mirror of
https://github.com/pgsty/minio.git
synced 2026-09-30 23:05:59 +03:00
Compare commits
17 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 2a4d51406b | |||
| 0596685ae7 | |||
| 358ab38fb0 | |||
| 254b19ac07 | |||
| eb4f5e5b31 | |||
| 8d06424b12 | |||
| fced863036 | |||
| 5af0865aab | |||
| 83821f0f1f | |||
| bf053386a4 | |||
| e7963381ba | |||
| 37c0edc7ca | |||
| a2fe70424e | |||
| 027d43b4eb | |||
| f99ed829b5 | |||
| 956a1a8e33 | |||
| 82f0a9828e |
@@ -96,6 +96,12 @@ jobs:
|
||||
- name: Run multipart listing and cancellation tests under race detector
|
||||
run: go test -race ./cmd -run '^Test(MultipartListing|MultipartAbort|PaginateMultipartUploads|ListMultipartUploads)' -count=1 -timeout=5m
|
||||
|
||||
- name: Run CPU metrics tests under race detector
|
||||
run: go test -race ./cmd -run '^TestLoadCPUMetrics' -count=1 -timeout=5m
|
||||
|
||||
- name: Run tag replication tests under race detector
|
||||
run: go test -race ./cmd -run '^TestAPITagging' -count=1 -timeout=5m
|
||||
|
||||
crosscompile:
|
||||
name: Cross Compile
|
||||
runs-on: ubuntu-latest
|
||||
|
||||
@@ -132,9 +132,9 @@ jobs:
|
||||
cosign-release: v3.1.2
|
||||
|
||||
- name: Build Draft release with GoReleaser
|
||||
uses: goreleaser/goreleaser-action@v7
|
||||
uses: goreleaser/goreleaser-action@f06c13b6b1a9625abc9e6e439d9c05a8f2190e94 # v7.2.3
|
||||
with:
|
||||
version: "~> v2"
|
||||
version: v2.18.1
|
||||
args: release --clean --skip=validate --config .github/goreleaser.yml
|
||||
env:
|
||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
@@ -4,6 +4,8 @@ on:
|
||||
workflow_dispatch:
|
||||
pull_request:
|
||||
paths:
|
||||
- "go.mod"
|
||||
- "go.sum"
|
||||
- ".github/goreleaser.yml"
|
||||
- ".github/nfpm.yml"
|
||||
- "Dockerfile.goreleaser"
|
||||
@@ -93,9 +95,9 @@ jobs:
|
||||
echo "LDFLAGS: ${LDFLAGS}"
|
||||
|
||||
- name: GoReleaser config check
|
||||
uses: goreleaser/goreleaser-action@v7
|
||||
uses: goreleaser/goreleaser-action@f06c13b6b1a9625abc9e6e439d9c05a8f2190e94 # v7.2.3
|
||||
with:
|
||||
version: "~> v2"
|
||||
version: v2.18.1
|
||||
args: check --config .github/goreleaser.yml
|
||||
|
||||
- name: Validate Helm chart and legacy upgrade identity
|
||||
@@ -107,9 +109,9 @@ jobs:
|
||||
syft-version: v1.50.0
|
||||
|
||||
- name: Build snapshot artifacts
|
||||
uses: goreleaser/goreleaser-action@v7
|
||||
uses: goreleaser/goreleaser-action@f06c13b6b1a9625abc9e6e439d9c05a8f2190e94 # v7.2.3
|
||||
with:
|
||||
version: "~> v2"
|
||||
version: v2.18.1
|
||||
# A pull-request snapshot has no trusted release identity. Exercise
|
||||
# the SBOM/checksum pipeline here, and reserve keyless signing for
|
||||
# the tag-triggered release workflow with GitHub OIDC.
|
||||
|
||||
+73
-25
@@ -2,13 +2,18 @@
|
||||
|
||||
## Unreleased
|
||||
|
||||
The entries below describe source changes on main since the latest published Server.
|
||||
Preparation target: `RELEASE.2026-09-16T00-00-00Z` (package version
|
||||
`20260916000000.0.0`). The entries below describe the candidate changes since
|
||||
the latest published Server.
|
||||
**The latest published Server remains 20260903.** These changes are not in its
|
||||
binaries, packages or images. See the [component matrix](https://silo.pgsty.com/compatibility/versions/)
|
||||
and [complete commit range](https://github.com/pgsty/silo/compare/RELEASE.2026-09-03T13-18-01Z...main).
|
||||
|
||||
### Authorization and security
|
||||
|
||||
- Synchronize CPU metrics reads with resource-metrics updates (#210), preventing
|
||||
concurrent map access from terminating the server during Prometheus scraping.
|
||||
Metric names, values and authentication requirements are unchanged.
|
||||
- Restrict embedded Console's anonymous sharing proxy to object-content GETs
|
||||
at the configured S3 origin, and reject every redirect. Internal metrics,
|
||||
system paths and non-download S3 operations cannot be reached through it.
|
||||
@@ -56,19 +61,28 @@ and [complete commit range](https://github.com/pgsty/silo/compare/RELEASE.2026-0
|
||||
quorum-written `xl.meta`; completion removes those upload-only fields. Native
|
||||
markers remain usable after their upload is completed or canceled. Strict
|
||||
listing returns a diagnostic 503 for legacy uploads or uncertain coverage;
|
||||
`api multipart_listing=legacy` is an explicit temporary migration mode.
|
||||
Upgrade every writer, drain old uploads and check the read-only admin
|
||||
`multipart-preflight` report before relying on strict listing. Per-process
|
||||
the default remains the released exact-key/cache-based `legacy` behavior.
|
||||
Opt into strict mode only through `MINIO_API_MULTIPART_LISTING=strict`, after
|
||||
upgrading every writer, draining old uploads, checking the read-only admin
|
||||
`multipart-preflight` report and validating scan capacity. The process-only
|
||||
setting is not persisted into shared API configuration. Historical
|
||||
`multipart_listing` keys are ignored and can be removed with a targeted
|
||||
`mcli admin config reset ALIAS api multipart_listing` before rollback. Per-process
|
||||
admission, directory-entry, worker and time budgets bound scan scheduling;
|
||||
each page still scans durable state. See [issue #79](https://github.com/pgsty/silo/issues/79)
|
||||
and its [design record](https://silo.pgsty.com/blog/design/list-multipart-uploads/).
|
||||
Thanks to mr javad seydi (@mrjavadseydi) for the original implementation.
|
||||
- Confirm multipart cancellation on a strict majority of each relevant set,
|
||||
and allow retries after partial deletion. Uncertain pools or insufficient
|
||||
confirmations return 503 rather than acknowledging a cancellation whose
|
||||
static remnants can later become readable. **Known boundary:** creation
|
||||
writes that finish after a storage timeout can still restore an upload after
|
||||
successful cancellation; this change does not add a durable creation fence.
|
||||
- Retain released read-quorum and best-effort multipart cancellation in default
|
||||
legacy mode. Strict mode requires majority deletion acknowledgements and
|
||||
permits retries below read quorum. When most drives were already empty,
|
||||
failed deletion of an observed remnant now returns 503 instead of being
|
||||
masked by empty-drive successes. Wrong-key or wrong-bucket cancellation
|
||||
preserves the valid upload's cache entry; successful cancellation notifies
|
||||
peers with the request context after the distributed lock is released.
|
||||
The HTTP response for an absent
|
||||
upload remains 204; this is not proof of physical cleanup. **Known boundary:**
|
||||
delayed creation writes can still restore an upload after cancellation;
|
||||
this change does not add a durable creation fence.
|
||||
- Preserve object tags during multi-pool metadata reconciliation by reading the
|
||||
resolved tag field together with its revision (#189). Previously, reconciliation
|
||||
could replace existing tags with an empty value.
|
||||
@@ -94,6 +108,35 @@ and [complete commit range](https://github.com/pgsty/silo/compare/RELEASE.2026-0
|
||||
converged. Follow the [read-only audit procedure](https://silo.pgsty.com/operations/replication/replica-metadata-audit/)
|
||||
before planning any repair of stored state.
|
||||
|
||||
- Make exact-version delete-marker purges converge in replicated buckets
|
||||
(`eb4f5e5b3`, `254b19ac0`, `358ab38fb`). A purge no longer creates the
|
||||
marker on drives that lacked it; a retried purge of a missing version is
|
||||
acknowledged only when a write-quorum majority of drives report it absent;
|
||||
purge results merge removed and reliably absent replies; healing a marker
|
||||
preserves its stored replication and purge metadata; and a queued marker
|
||||
creation is re-checked against the source under the replication lock before
|
||||
it is sent, so a purge that already reached the targets is not undone by a
|
||||
stale task from another frontend, a GET/LIST heal, the scanner or MRF.
|
||||
Purging a data version whose earlier purge is still pending reports it as a
|
||||
data version. **Known limitations:** creations already in flight or replayed
|
||||
from another site, and minority marker copies left by a crash after a
|
||||
majority-acknowledged purge, are tracked in #217. See
|
||||
[the replication reliability record](https://silo.pgsty.com/blog/design/replication-reliability/).
|
||||
- Keep a null object version that has listing quorum when a newer minority of
|
||||
drives sorts first (`8d06424b1`). The resolver recounts per header only when
|
||||
the original selection lacks quorum, every non-empty drive stream holds
|
||||
exactly one ordinary null version and all share the same erasure layout;
|
||||
mixed histories keep their previous behavior. **Known limitation:** a
|
||||
successful ListObjects can still omit readable keys during rolling restarts
|
||||
with concurrent overwrites (#218). Do not run destination-deleting sync tools
|
||||
against a listing taken during a rolling restart; list again once the
|
||||
cluster is stable.
|
||||
- Carry object tags through rebalance and decommission for ordinary and
|
||||
multipart writes (`fced86303`). Both migration entry points restore the tags
|
||||
and their revision fields when rewriting the object in the destination pool.
|
||||
Tags dropped by earlier migrations are not recovered; audit tag-dependent
|
||||
lifecycle and policy rules for pools migrated with an older build.
|
||||
|
||||
- Evaluate conditional multipart completion against the logical current object
|
||||
across all pools while holding the existing object lock. A stale `If-Match`
|
||||
can no longer replace newer data in another pool, and the current ETag is no
|
||||
@@ -158,23 +201,28 @@ and [complete commit range](https://github.com/pgsty/silo/compare/RELEASE.2026-0
|
||||
|
||||
- Restore embedded Console login over loopback TLS, trusted-proxy handling and
|
||||
all four WebSocket connection limits. Preserve Go TLS defaults across transports.
|
||||
- Directly require `github.com/pgsty/silo-pkg/v3` v3.14.0; select Console
|
||||
`v0.0.0-20260916034812-56dfe455ac2f` and MC
|
||||
`v0.0.0-20260913012246-4f609a4da3bb` with explicit PGSTY replacements.
|
||||
- Pin upstream minio-go `v7.3.1-0.20260910142817-60bd07042d49`; refresh Go x/*
|
||||
modules and security fixes including bounded AMQP frame handling. Keep Go
|
||||
1.27.1 and go-systemd v22.6.0's NetBSD compatibility replacement.
|
||||
- Directly require `github.com/pgsty/silo-pkg/v3` v3.14.1; select released Console
|
||||
v2.4.1 (`v0.0.0-20260916075814-1360e26d976d`) and mcli 20260916
|
||||
(`v0.0.0-20260916070421-e952aa78f10a`) with explicit PGSTY replacements.
|
||||
The embedded frontend identifies itself as Console v2.4.1.
|
||||
- Pin upstream minio-go `v7.3.1-0.20260915093545-32e1f32cb176` to handle
|
||||
CopyObject errors embedded in HTTP 200 responses. Update JWX to v3.3.0 for
|
||||
JSON field-name escaping, strfmt to v0.27.2 for Go 1.27 hostname validation,
|
||||
and LZ4 to v4.1.30 for frame-reader, partial-read and concurrency fixes.
|
||||
Retain the earlier Go x/* and bounded AMQP frame updates, Go 1.27.1, and
|
||||
go-systemd v22.6.0's NetBSD compatibility replacement.
|
||||
- Refresh container base digests and build static curl 8.22.0 from verified
|
||||
source for both Linux architectures. Pin the actual mcli 20260913 archives and
|
||||
hashes. Helm's client image follows that release; its Server image still names
|
||||
the latest published Server 20260903.
|
||||
source for both Linux architectures. Pin the published mcli 20260916 archives
|
||||
and hashes in the container and update the client installer default.
|
||||
- Prepare Helm chart 7.0.3 with Server and client defaults for the September 16
|
||||
batch. Publish the chart only after the corresponding Server image exists.
|
||||
- Pin GoReleaser v2.18.1 and its action commit identically in snapshot and release
|
||||
workflows. Dependency-only PRs now run the Test Release Pipeline too.
|
||||
|
||||
The dependency update passed the final candidate's Go, vulnerability and Test
|
||||
Release workflows; native curl builds passed on both architectures. A local
|
||||
ARM64 image passed startup, health, S3 transfer and embedded Console checks.
|
||||
These checks do not publish a Server tag or production image and do not replace
|
||||
cluster upgrade/rollback acceptance for the next release. Dated investigations
|
||||
retain the exact source and runtime boundaries they tested.
|
||||
Validation of earlier source revisions does not establish acceptance of this
|
||||
candidate. Final source, package, image and multi-process checks are tracked
|
||||
separately in [#203](https://github.com/pgsty/silo/issues/203). No Server release
|
||||
or production rollout is implied by this preparation target.
|
||||
|
||||
## RELEASE.2026-09-03T13-18-01Z
|
||||
|
||||
|
||||
@@ -17,9 +17,9 @@ ENV GOPATH=/go
|
||||
ENV CGO_ENABLED=0
|
||||
|
||||
ARG MC_REPO=pgsty/mc
|
||||
ARG MC_VERSION=RELEASE.2026-09-13T00-00-00Z
|
||||
ARG MC_AMD64_SHA256=9d2a92de9c7b887d9b944fe9ddce68d23f1b6df3415092e737594e56593f5e2b
|
||||
ARG MC_ARM64_SHA256=3d82e9ea6c601c4cb44fe5dd5f2ad1b7d7d64369378110f9ada9c524688a452a
|
||||
ARG MC_VERSION=RELEASE.2026-09-16T00-00-00Z
|
||||
ARG MC_AMD64_SHA256=4ba2814fd5507fbe6b4d237c359750b9119d28d7217495fa5b48002fcbd397ef
|
||||
ARG MC_ARM64_SHA256=b7008ca2a1bc5735b6585981c59a3640a0daa152dc789df3d1a0d0438787de82
|
||||
|
||||
RUN apk add -U --no-cache \
|
||||
ca-certificates \
|
||||
|
||||
@@ -40,7 +40,7 @@ if [ -n "${MCLI_BIN:-}" ]; then
|
||||
exit 0
|
||||
fi
|
||||
|
||||
release=${MCLI_RELEASE:-RELEASE.2026-09-13T00-00-00Z}
|
||||
release=${MCLI_RELEASE:-RELEASE.2026-09-16T00-00-00Z}
|
||||
version_hyphen=${release#RELEASE.}
|
||||
package_version=$(printf '%s\n' "${version_hyphen}" | sed -E 's/^([0-9]{4})-([0-9]{2})-([0-9]{2})T([0-9]{2})-([0-9]{2})-([0-9]{2})Z$/\1\2\3\4\5\6.0.0/')
|
||||
if [ "${package_version}" = "${version_hyphen}" ]; then
|
||||
|
||||
@@ -571,6 +571,14 @@
|
||||
"minio_resource_stats",
|
||||
"minio_s3",
|
||||
"minio_stats",
|
||||
"minio_system_cpu_avg_idle",
|
||||
"minio_system_cpu_avg_iowait",
|
||||
"minio_system_cpu_load",
|
||||
"minio_system_cpu_load_perc",
|
||||
"minio_system_cpu_nice",
|
||||
"minio_system_cpu_steal",
|
||||
"minio_system_cpu_system",
|
||||
"minio_system_cpu_user",
|
||||
"minio_test"
|
||||
],
|
||||
"headers": [
|
||||
|
||||
@@ -113,7 +113,7 @@ helm_run template my-release "${new_chart}" \
|
||||
go run ./buildscripts/helm-migration-guard "${old_render}" "${new_render}"
|
||||
|
||||
helm_run package "${new_chart}" --destination "${output_dir}" >/dev/null
|
||||
test -s "${work_dir}/silo-7.0.2.tgz"
|
||||
test -s "${work_dir}/silo-7.0.3.tgz"
|
||||
if find "${work_dir}" -maxdepth 1 -type f -name 'minio-*.tgz' | grep -q .; then
|
||||
echo "Helm packaging emitted a legacy MinIO chart name" >&2
|
||||
exit 1
|
||||
|
||||
@@ -504,6 +504,35 @@ func replicateDelete(ctx context.Context, dobj DeletedObjectReplicationInfo, obj
|
||||
ctx = lkctx.Context()
|
||||
defer lk.Unlock(lkctx)
|
||||
|
||||
if !isPurge && dobj.DeleteMarkerVersionID != "" {
|
||||
// A creation task carries the marker as it looked when it was queued
|
||||
// by the DELETE handler, a GET/HEAD/LIST heal, the scanner or MRF.
|
||||
// While it waited for this lock, another frontend may have purged the
|
||||
// marker and replicated that purge. The targets no longer hold the
|
||||
// marker, so the queued creation would recreate it from the stale
|
||||
// snapshot. Confirm the source version under the lock first.
|
||||
switch deleteMarkerCreationState(ctx, objectAPI, dobj) {
|
||||
case creationStale:
|
||||
return replicatedInfos{}
|
||||
case creationUnverified:
|
||||
dobj.RetryCount++
|
||||
globalReplicationPool.Get().queueMRFSave(dobj.ToMRFEntry())
|
||||
sendEvent(eventArgs{
|
||||
BucketName: bucket,
|
||||
Object: ObjectInfo{
|
||||
Bucket: bucket,
|
||||
Name: dobj.ObjectName,
|
||||
VersionID: versionID,
|
||||
DeleteMarker: dobj.DeleteMarker,
|
||||
},
|
||||
UserAgent: "Internal: [Replication]",
|
||||
Host: globalLocalNodeName,
|
||||
EventName: event.ObjectReplicationNotTracked,
|
||||
})
|
||||
return replicatedInfos{}
|
||||
}
|
||||
}
|
||||
|
||||
rinfos := replicatedInfos{Targets: make([]replicatedTargetInfo, 0, len(dsc.targetsMap))}
|
||||
var wg sync.WaitGroup
|
||||
var mu sync.Mutex
|
||||
@@ -1955,6 +1984,41 @@ func (di DeletedObjectReplicationInfo) isVersionPurge() bool {
|
||||
return di.VersionID != "" || di.DeleteMarkerVersionID != "" && !di.VersionPurgeStatus().Empty()
|
||||
}
|
||||
|
||||
// creationState is the outcome of re-reading a queued delete-marker creation
|
||||
// against the source.
|
||||
type creationState int
|
||||
|
||||
const (
|
||||
creationCurrent creationState = iota
|
||||
creationStale
|
||||
creationUnverified
|
||||
)
|
||||
|
||||
// deleteMarkerCreationState re-reads the marker a queued creation task refers
|
||||
// to. The version being absent, no longer a marker, or under a version purge
|
||||
// makes the creation stale. A failed read is not absence: the caller retries
|
||||
// later instead of guessing.
|
||||
func deleteMarkerCreationState(ctx context.Context, objectAPI ObjectLayer, dobj DeletedObjectReplicationInfo) creationState {
|
||||
oi, err := objectAPI.GetObjectInfo(ctx, dobj.Bucket, dobj.ObjectName, ObjectOptions{
|
||||
VersionID: dobj.DeleteMarkerVersionID,
|
||||
Versioned: globalBucketVersioningSys.PrefixEnabled(dobj.Bucket, dobj.ObjectName),
|
||||
VersionSuspended: globalBucketVersioningSys.Suspended(dobj.Bucket),
|
||||
})
|
||||
switch {
|
||||
case isErrObjectNotFound(err), isErrVersionNotFound(err):
|
||||
return creationStale
|
||||
case err != nil && !isErrMethodNotAllowed(err):
|
||||
return creationUnverified
|
||||
}
|
||||
if !oi.DeleteMarker || oi.VersionID != dobj.DeleteMarkerVersionID {
|
||||
return creationStale
|
||||
}
|
||||
if !oi.VersionPurgeStatus.Empty() || oi.VersionPurgeStatusInternal != "" {
|
||||
return creationStale
|
||||
}
|
||||
return creationCurrent
|
||||
}
|
||||
|
||||
// Purge metadata uses COMPLETE; operation statistics and audit use COMPLETED.
|
||||
func purgeReplicationStatus(status VersionPurgeStatusType) replication.StatusType {
|
||||
if replication.StatusType(status) == replication.CompletedLegacy {
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
|
||||
"github.com/minio/minio/internal/auth"
|
||||
"github.com/minio/minio/internal/bucket/replication"
|
||||
xhttp "github.com/minio/minio/internal/http"
|
||||
)
|
||||
|
||||
// An exact-version DELETE without a replication decision physically purges
|
||||
// the version. Its response must describe the stored version: a data version
|
||||
// is not a delete marker even while pending-purge metadata makes the lookup
|
||||
// expose it as deleted, and a stored marker stays a marker.
|
||||
func TestDeleteMarkerPurgeResponseIdentity(t *testing.T) {
|
||||
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, endpoints: []string{"DeleteObject"}, objAPITest: func(obj ObjectLayer, backend, bucket string, router http.Handler, creds auth.Credentials, t *testing.T) {
|
||||
defer replicationTestCapacity(obj)()
|
||||
ctx := t.Context()
|
||||
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, enabledBucketVersioningConfig); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Seed the purge state the DELETE handler would record for a target.
|
||||
arn := "arn:minio:replication::" + mustGetUUID() + ":bucket"
|
||||
pendingPurge := ReplicationState{
|
||||
VersionPurgeStatusInternal: arn + "=PENDING;",
|
||||
PurgeTargets: map[string]VersionPurgeStatusType{arn: replication.VersionPurgePending},
|
||||
}
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
marker, purged bool
|
||||
}{
|
||||
{name: "data"},
|
||||
{name: "data-pending-purge", purged: true},
|
||||
{name: "marker", marker: true},
|
||||
{name: "marker-pending-purge", marker: true, purged: true},
|
||||
} {
|
||||
name := "purge-response-" + mustGetUUID()
|
||||
oi, err := obj.PutObject(ctx, bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
version := oi.VersionID
|
||||
if tc.marker {
|
||||
opts := ObjectOptions{Versioned: true, VersionID: mustGetUUID(), DeleteMarker: true, MTime: UTCNow(), ReplicationRequest: true}
|
||||
opts.SetReplicaStatus(replication.Replica)
|
||||
if _, err := obj.DeleteObject(ctx, bucket, name, opts); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
version = opts.VersionID
|
||||
}
|
||||
if tc.purged {
|
||||
if _, err := obj.DeleteObject(ctx, bucket, name, ObjectOptions{Versioned: true, VersionID: version, DeleteReplication: pendingPurge}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
before, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{Versioned: true, VersionID: version})
|
||||
if tc.marker || tc.purged {
|
||||
if !isErrMethodNotAllowed(err) || !before.DeleteMarker {
|
||||
t.Fatalf("%s: seeded version not exposed as deleted: %+v %v", tc.name, before, err)
|
||||
}
|
||||
} else if err != nil || before.DeleteMarker {
|
||||
t.Fatalf("%s: seeded data version: %+v %v", tc.name, before, err)
|
||||
}
|
||||
if tc.purged && before.VersionPurgeStatus != replication.VersionPurgePending {
|
||||
t.Fatalf("%s: purge state not seeded: %+v", tc.name, before)
|
||||
}
|
||||
req, err := newTestSignedRequestV4(http.MethodDelete, "/"+bucket+"/"+name+"?versionId="+version, 0, nil, creds.AccessKey, creds.SecretKey, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rec := httptest.NewRecorder()
|
||||
router.ServeHTTP(rec, req)
|
||||
if rec.Code != http.StatusNoContent {
|
||||
t.Fatalf("%s: status=%d: %s", tc.name, rec.Code, rec.Body.String())
|
||||
}
|
||||
// The handler sets these headers by map key, not canonical name.
|
||||
if got := rec.Header()[xhttp.AmzVersionID]; len(got) != 1 || got[0] != version {
|
||||
t.Fatalf("%s: response version %q, want %q", tc.name, got, version)
|
||||
}
|
||||
if got := len(rec.Header()[xhttp.AmzDeleteMarker]) == 1 && rec.Header()[xhttp.AmzDeleteMarker][0] == "true"; got != tc.marker {
|
||||
t.Fatalf("%s: x-amz-delete-marker=%v for a stored marker=%v", tc.name, got, tc.marker)
|
||||
}
|
||||
if _, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{Versioned: true, VersionID: version}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
t.Fatalf("%s: version not purged: %v", tc.name, err)
|
||||
}
|
||||
t.Logf("%s %s: purged with x-amz-delete-marker=%q", backend, tc.name, rec.Header()[xhttp.AmzDeleteMarker])
|
||||
}
|
||||
}})
|
||||
}
|
||||
@@ -0,0 +1,524 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"maps"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/minio/madmin-go/v3"
|
||||
"github.com/minio/minio-go/v7"
|
||||
"github.com/minio/minio/internal/bucket/replication"
|
||||
xhttp "github.com/minio/minio/internal/http"
|
||||
"github.com/minio/minio/internal/once"
|
||||
)
|
||||
|
||||
func TestDeleteMarkerPurgeIntent(t *testing.T) {
|
||||
version := mustGetUUID()
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
edit func(*ObjectOptions)
|
||||
want bool
|
||||
}{
|
||||
{"ordinary", func(*ObjectOptions) {}, true},
|
||||
{"receiver", func(o *ObjectOptions) { o.ReplicationRequest = true; o.SetReplicaStatus(replication.Replica) }, true},
|
||||
{"untrusted-replica", func(o *ObjectOptions) { o.SetReplicaStatus(replication.Replica) }, false},
|
||||
{"complete-purge", func(o *ObjectOptions) {
|
||||
o.DeleteReplication.VersionPurgeStatusInternal = string(replication.VersionPurgeComplete)
|
||||
}, true},
|
||||
{"pending-purge", func(o *ObjectOptions) { o.DeleteReplication.VersionPurgeStatusInternal = "PENDING" }, false},
|
||||
{"failed-purge", func(o *ObjectOptions) { o.DeleteReplication.VersionPurgeStatusInternal = "FAILED" }, false},
|
||||
{"unknown-purge", func(o *ObjectOptions) { o.DeleteReplication.VersionPurgeStatusInternal = "future-state" }, false},
|
||||
{"pending-creation", func(o *ObjectOptions) { o.DeleteReplication.ReplicationStatusInternal = "PENDING" }, false},
|
||||
{"failed-creation", func(o *ObjectOptions) { o.DeleteReplication.ReplicationStatusInternal = "FAILED" }, false},
|
||||
{"completed-creation", func(o *ObjectOptions) { o.DeleteReplication.ReplicationStatusInternal = "COMPLETED" }, false},
|
||||
{"create-marker", func(o *ObjectOptions) { o.DeleteMarker = true }, false},
|
||||
{"empty", func(o *ObjectOptions) { o.VersionID = "" }, false},
|
||||
{"null", func(o *ObjectOptions) { o.VersionID = nullVersionID }, false},
|
||||
{"invalid", func(o *ObjectOptions) { o.VersionID = "invalid" }, false},
|
||||
{"zero-uuid", func(o *ObjectOptions) { o.VersionID = emptyUUID }, false},
|
||||
{"movement", func(o *ObjectOptions) { o.DataMovement = true }, false},
|
||||
{"free-version", func(o *ObjectOptions) { o.InclFreeVersions = true }, false},
|
||||
{"expiration", func(o *ObjectOptions) { o.Expiration.Expire = true }, false},
|
||||
{"transition", func(o *ObjectOptions) { o.Transition.Status = "complete" }, false},
|
||||
{"restored-expiration", func(o *ObjectOptions) { o.Transition.ExpireRestored = true }, false},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
o := ObjectOptions{VersionID: version, Versioned: true}
|
||||
tc.edit(&o)
|
||||
if got := o.isVersionPurge(); got != tc.want {
|
||||
t.Fatalf("purge=%v want=%v", got, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Every other StorageAPI method panics through the nil embedding: a proof may
|
||||
// read metadata, but must never delete, heal, or ask a different storage API.
|
||||
type absenceProofDisk struct {
|
||||
StorageAPI
|
||||
t *testing.T
|
||||
err error
|
||||
reads *atomic.Int32
|
||||
deletes *atomic.Int32
|
||||
}
|
||||
|
||||
func (d absenceProofDisk) ReadVersion(_ context.Context, _, _, _, _ string, opts ReadOptions) (FileInfo, error) {
|
||||
if opts.ReadData || opts.Healing {
|
||||
d.t.Error("absence proof requested data/healing")
|
||||
}
|
||||
d.reads.Add(1)
|
||||
return FileInfo{}, d.err
|
||||
}
|
||||
|
||||
func (d absenceProofDisk) DeleteVersion(_ context.Context, _, _ string, _ FileInfo, _ bool, _ DeleteOptions) error {
|
||||
if d.deletes == nil {
|
||||
d.t.Error("read-only absence proof attempted a deletion")
|
||||
} else {
|
||||
d.deletes.Add(1)
|
||||
}
|
||||
return d.err
|
||||
}
|
||||
|
||||
func TestVersionPurgeAbsenceProof(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
errs []error
|
||||
quorum, all bool
|
||||
}{
|
||||
{"all-absent", []error{errFileNotFound, errFileVersionNotFound, errFileNotFound, errFileVersionNotFound}, true, true},
|
||||
{"majority-absent", []error{errFileNotFound, errFileVersionNotFound, errFileNotFound, nil}, true, false},
|
||||
{"half-absent", []error{errFileNotFound, errFileVersionNotFound, errDiskNotFound, errDiskNotFound}, false, false},
|
||||
{"corrupt", []error{errFileNotFound, errFileVersionNotFound, errFileCorrupt, errFileCorrupt}, false, false},
|
||||
{"permission", []error{errFileNotFound, errFileVersionNotFound, errDiskAccessDenied, errVolumeAccessDenied}, false, false},
|
||||
{"volume-missing", []error{errFileNotFound, errFileVersionNotFound, errVolumeNotFound, errVolumeNotFound}, false, false},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
var reads atomic.Int32
|
||||
disks := make([]StorageAPI, len(tc.errs))
|
||||
for i, err := range tc.errs {
|
||||
disks[i] = absenceProofDisk{t: t, err: err, reads: &reads}
|
||||
}
|
||||
er := erasureObjects{getDisks: func() []StorageAPI { return disks }}
|
||||
oldQueue := globalMRFState.opCh
|
||||
globalMRFState.opCh = make(chan PartialOperation, 1)
|
||||
defer func() { globalMRFState.opCh = oldQueue }()
|
||||
all, err := er.confirmVersionAbsent(t.Context(), "bucket", "object", mustGetUUID())
|
||||
if (err == nil) != tc.quorum || all != tc.all || reads.Load() != int32(len(disks)) || len(globalMRFState.opCh) != 0 {
|
||||
t.Fatalf("proof all=%v error=%v reads=%d queued=%d", all, err, reads.Load(), len(globalMRFState.opCh))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestVersionPurgeAggregation(t *testing.T) {
|
||||
for _, pure := range []bool{false, true} {
|
||||
for _, fault := range []error{nil, errFileCorrupt, errDiskAccessDenied} {
|
||||
t.Run(fmt.Sprintf("purge_%v_fault_%v", pure, fault), func(t *testing.T) {
|
||||
var calls atomic.Int32
|
||||
disks := make([]StorageAPI, 4)
|
||||
for i, err := range []error{nil, errFileNotFound, errFileVersionNotFound, fault} {
|
||||
disks[i] = absenceProofDisk{t: t, err: err, deletes: &calls}
|
||||
}
|
||||
er := erasureObjects{getDisks: func() []StorageAPI { return disks }}
|
||||
err := er.deleteObjectVersion(t.Context(), "bucket", "object", FileInfo{VersionID: mustGetUUID()}, false, pure)
|
||||
if (err == nil) != pure || calls.Load() != 4 {
|
||||
t.Fatalf("aggregation error=%v calls=%d", err, calls.Load())
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteMarkerPurgeOmittedPool(t *testing.T) {
|
||||
z, bucket := consistencyPools(t)
|
||||
defer replicationTestCapacity(z)()
|
||||
const name = "partially-absent"
|
||||
er := z.serverPools[1].getHashedSet(name)
|
||||
_, opts := seedPurgeMarker(t, er, bucket, name, false)
|
||||
opts.ReplicationRequest = false
|
||||
opts.DeleteReplication = ReplicationState{}
|
||||
// Pool 0 appears absent at read quorum, but half of it cannot be read.
|
||||
missing := z.serverPools[0].getHashedSet(name)
|
||||
original := missing.getDisks
|
||||
disks := append([]StorageAPI(nil), original()...)
|
||||
clear(disks[len(disks)/2:])
|
||||
missing.getDisks = func() []StorageAPI { return disks }
|
||||
defer func() { missing.getDisks = original }()
|
||||
_, err := z.DeleteObject(t.Context(), bucket, name, opts)
|
||||
if !isErrWriteQuorum(err) {
|
||||
t.Errorf("omitted pool acknowledged deletion: %T %v", err, err)
|
||||
}
|
||||
if countPurgeMarkers(t, er.getDisks(), bucket, name, opts.VersionID) != 16 {
|
||||
t.Fatal("mutated known copies before checking omitted pool")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteMarkerPurgeReceivingPool(t *testing.T) {
|
||||
z, bucket := consistencyPools(t)
|
||||
defer replicationTestCapacity(z)()
|
||||
for _, duplicate := range []bool{false, true} {
|
||||
t.Run(fmt.Sprintf("duplicate_%v", duplicate), func(t *testing.T) {
|
||||
name := "receiver-" + mustGetUUID()
|
||||
_, opts := seedPurgeMarker(t, z.serverPools[0].getHashedSet(name), bucket, name, false)
|
||||
if duplicate {
|
||||
create := opts
|
||||
create.DeleteMarker = true
|
||||
if _, err := z.serverPools[1].DeleteObject(t.Context(), bucket, name, create); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
} else {
|
||||
// Latest-key routing selects pool 1, but the addressed marker is
|
||||
// in pool 0. A replica purge addresses a version, not the latest.
|
||||
putConsistencyObject(t, z, bucket, name, 1, "newer-data", ObjectOptions{Versioned: true})
|
||||
}
|
||||
if _, err := z.DeleteObject(t.Context(), bucket, name, opts); err != nil {
|
||||
t.Errorf("replica purge did not find its addressed version: %v", err)
|
||||
}
|
||||
for i, pool := range z.serverPools {
|
||||
if got := countPurgeMarkers(t, pool.getHashedSet(name).getDisks(), bucket, name, opts.VersionID); got != 0 {
|
||||
t.Errorf("replica purge left %d marker copies in pool %d", got, i)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteMarkerPurgeCallbackIntent(t *testing.T) {
|
||||
for _, loseCopies := range []bool{false, true} {
|
||||
t.Run(fmt.Sprintf("missing_after_callback_%v", loseCopies), func(t *testing.T) {
|
||||
z, bucket := consistencyPools(t)
|
||||
defer replicationTestCapacity(z)()
|
||||
const name = "callback-pending"
|
||||
er := z.serverPools[0].getHashedSet(name)
|
||||
_, opts := seedPurgeMarker(t, er, bucket, name, false)
|
||||
original := er.getDisks
|
||||
all := append([]StorageAPI(nil), original()...)
|
||||
defer func() { er.getDisks = original }()
|
||||
opts.ReplicationRequest = false
|
||||
opts.DeleteReplication = ReplicationState{}
|
||||
metadata, retention := 0, 0
|
||||
opts.EvalRetentionBypassFn = func(oi ObjectInfo, err error) error {
|
||||
retention++
|
||||
if !oi.DeleteMarker || !isErrMethodNotAllowed(err) {
|
||||
t.Fatalf("retention callback lost marker: %+v %v", oi, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
opts.EvalMetadataFn = func(*ObjectInfo, error) (ReplicateDecision, error) {
|
||||
metadata++
|
||||
if loseCopies {
|
||||
// Inject a disk-state change after preflight, before the set's
|
||||
// metadata-update read. It returns NotFound at read quorum.
|
||||
for _, disk := range all[:8] {
|
||||
if err := disk.DeleteVersion(t.Context(), bucket, name, FileInfo{VersionID: opts.VersionID}, false, DeleteOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
partial := make([]StorageAPI, len(all))
|
||||
copy(partial, all[:8])
|
||||
er.getDisks = func() []StorageAPI { return partial }
|
||||
}
|
||||
decision := ReplicateDecision{}
|
||||
decision.Set(newReplicateTargetDecision("arn1", true, false))
|
||||
return decision, nil
|
||||
}
|
||||
_, err := z.DeleteObject(t.Context(), bucket, name, opts)
|
||||
if loseCopies {
|
||||
if !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
t.Errorf("pending update misclassified as physical purge: %T %v", err, err)
|
||||
}
|
||||
} else if err != nil {
|
||||
t.Errorf("pending update failed: %v", err)
|
||||
}
|
||||
want := 16
|
||||
if loseCopies {
|
||||
want = 8
|
||||
}
|
||||
if got := countPurgeMarkers(t, all, bucket, name, opts.VersionID); got != want || metadata != 1 || retention != 1 {
|
||||
t.Errorf("pending update: copies=%d want=%d metadata=%d retention=%d", got, want, metadata, retention)
|
||||
}
|
||||
if !loseCopies {
|
||||
fi, err := all[0].ReadVersion(t.Context(), "", bucket, name, opts.VersionID, ReadOptions{})
|
||||
if err != nil || fi.VersionPurgeStatus() != replication.VersionPurgePending {
|
||||
t.Errorf("pending state not persisted: %+v %v", fi.ReplicationState, err)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteMarkerReceivingMetadataUpdate(t *testing.T) {
|
||||
z, bucket := consistencyPools(t)
|
||||
defer replicationTestCapacity(z)()
|
||||
for _, state := range []VersionPurgeStatusType{replication.VersionPurgePending, replication.VersionPurgeFailed} {
|
||||
t.Run(string(state), func(t *testing.T) {
|
||||
name := "update-" + mustGetUUID()
|
||||
_, opts := seedPurgeMarker(t, z.serverPools[0].getHashedSet(name), bucket, name, false)
|
||||
create := opts
|
||||
create.DeleteMarker = true
|
||||
if _, err := z.serverPools[1].DeleteObject(t.Context(), bucket, name, create); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
opts.DeleteReplication = ReplicationState{VersionPurgeStatusInternal: string(state)}
|
||||
if _, err := z.DeleteObject(t.Context(), bucket, name, opts); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for i, pool := range z.serverPools {
|
||||
if got := countPurgeMarkers(t, pool.getHashedSet(name).getDisks(), bucket, name, opts.VersionID); got != 16 {
|
||||
t.Errorf("metadata update removed marker copies in pool %d: %d", i, got)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestVersionPurgeRetentionGate(t *testing.T) {
|
||||
z, er, all, bucket := markerPurgeFixture(t, 4)
|
||||
data, _ := seedPurgeMarker(t, er, bucket, "protected-data", true)
|
||||
denied := errors.New("retention denied")
|
||||
opts := ObjectOptions{Versioned: true, VersionID: data, EvalRetentionBypassFn: func(oi ObjectInfo, err error) error {
|
||||
if err != nil || oi.VersionID != data || oi.DeleteMarker {
|
||||
t.Errorf("wrong retained version: %+v error=%v", oi, err)
|
||||
}
|
||||
return denied
|
||||
}}
|
||||
_, err := z.DeleteObject(t.Context(), bucket, "protected-data", opts)
|
||||
if !errors.Is(err, denied) {
|
||||
t.Fatalf("retention gate: %v", err)
|
||||
}
|
||||
for _, disk := range all {
|
||||
if _, err := disk.ReadVersion(t.Context(), "", bucket, "protected-data", data, ReadOptions{}); err != nil {
|
||||
t.Errorf("protected version changed: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func markerPurgeFixture(t *testing.T, n int) (*erasureServerPools, *erasureObjects, []StorageAPI, string) {
|
||||
t.Helper()
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
obj, dirs, err := prepareErasure(ctx, n)
|
||||
if err != nil {
|
||||
cancel()
|
||||
t.Fatal(err)
|
||||
}
|
||||
restore := replicationTestCapacity(obj)
|
||||
t.Cleanup(func() {
|
||||
cancel()
|
||||
restore()
|
||||
obj.Shutdown(context.Background())
|
||||
removeRoots(dirs)
|
||||
})
|
||||
bucket := getRandomBucketName()
|
||||
if err := obj.MakeBucket(ctx, bucket, MakeBucketOptions{VersioningEnabled: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
z := obj.(*erasureServerPools)
|
||||
er := z.serverPools[0].sets[0]
|
||||
original := er.getDisks
|
||||
t.Cleanup(func() { er.getDisks = original })
|
||||
return z, er, append([]StorageAPI(nil), original()...), bucket
|
||||
}
|
||||
|
||||
func seedPurgeMarker(t *testing.T, er *erasureObjects, bucket, name string, withData bool) (string, ObjectOptions) {
|
||||
t.Helper()
|
||||
dataVersion := ""
|
||||
if withData {
|
||||
oi, err := er.PutObject(t.Context(), bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
dataVersion = oi.VersionID
|
||||
}
|
||||
opts := ObjectOptions{Versioned: true, VersionID: mustGetUUID(), DeleteMarker: true, ReplicationRequest: true, MTime: UTCNow().Add(-time.Hour), NoAuditLog: true}
|
||||
opts.SetReplicaStatus(replication.Replica)
|
||||
if _, err := er.DeleteObject(t.Context(), bucket, name, opts); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
opts.DeleteMarker = false
|
||||
return dataVersion, opts
|
||||
}
|
||||
|
||||
func countPurgeMarkers(t *testing.T, disks []StorageAPI, bucket, name, version string) int {
|
||||
t.Helper()
|
||||
n := 0
|
||||
for _, disk := range disks {
|
||||
fi, err := disk.ReadVersion(t.Context(), "", bucket, name, version, ReadOptions{})
|
||||
if err == nil && fi.Deleted {
|
||||
n++
|
||||
} else if err != errFileNotFound && err != errFileVersionNotFound {
|
||||
t.Fatalf("unexpected disk version: deleted=%v error=%v", fi.Deleted, err)
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
func TestDeleteMarkerPurgeQuorum(t *testing.T) {
|
||||
for _, tc := range []struct{ disks, online int }{{4, 2}, {16, 7}, {16, 8}, {16, 9}} {
|
||||
t.Run(fmt.Sprintf("%d_disks_%d_online", tc.disks, tc.online), func(t *testing.T) {
|
||||
z, er, all, bucket := markerPurgeFixture(t, tc.disks)
|
||||
const name = "marker"
|
||||
dataVersion, opts := seedPurgeMarker(t, er, bucket, name, true)
|
||||
partial := make([]StorageAPI, len(all))
|
||||
copy(partial, all[:tc.online])
|
||||
er.getDisks = func() []StorageAPI { return partial }
|
||||
_, first := z.DeleteObject(t.Context(), bucket, name, opts)
|
||||
_, retry := z.DeleteObject(t.Context(), bucket, name, opts)
|
||||
if tc.online <= tc.disks/2 {
|
||||
if !isErrWriteQuorum(first) || !isErrWriteQuorum(retry) {
|
||||
t.Errorf("below write quorum: first=%T %v, retry=%T %v", first, first, retry, retry)
|
||||
}
|
||||
} else if first != nil || (retry != nil && !isErrVersionNotFound(retry) && !isErrObjectNotFound(retry)) {
|
||||
t.Errorf("write quorum available: first=%v retry=%v", first, retry)
|
||||
}
|
||||
want := tc.disks - tc.online
|
||||
if tc.online < tc.disks/2 {
|
||||
want = tc.disks
|
||||
}
|
||||
if got := countPurgeMarkers(t, all, bucket, name, opts.VersionID); got != want {
|
||||
t.Errorf("partial purge markers=%d want=%d", got, want)
|
||||
}
|
||||
er.getDisks = func() []StorageAPI { return all }
|
||||
_, err := z.DeleteObject(t.Context(), bucket, name, opts)
|
||||
if err != nil && !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
t.Errorf("full-online retry: %v", err)
|
||||
}
|
||||
if tc.online <= tc.disks/2 {
|
||||
if got := countPurgeMarkers(t, all, bucket, name, opts.VersionID); got != 0 {
|
||||
t.Errorf("full-online retry recreated/retained %d marker copies", got)
|
||||
}
|
||||
}
|
||||
// A quorum-confirmed missing retry can leave minority residue. Exercise
|
||||
// the real dangling-version healer after every disk has returned.
|
||||
_, _ = er.HealObject(t.Context(), bucket, name, opts.VersionID, madmin.HealOpts{Remove: true})
|
||||
if got := countPurgeMarkers(t, all, bucket, name, opts.VersionID); got != 0 {
|
||||
t.Errorf("marker copies after retry and heal=%d", got)
|
||||
}
|
||||
oi, err := z.GetObjectInfo(t.Context(), bucket, name, ObjectOptions{Versioned: true})
|
||||
if err != nil || oi.DeleteMarker || oi.VersionID != dataVersion {
|
||||
t.Errorf("underlying version not visible: %+v error=%v", oi, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteMarkerPurgeMissingKey(t *testing.T) {
|
||||
for _, pooled := range []bool{false, true} {
|
||||
t.Run(fmt.Sprintf("pooled_%v", pooled), func(t *testing.T) {
|
||||
z, er, all, bucket := markerPurgeFixture(t, 4)
|
||||
_, opts := seedPurgeMarker(t, er, bucket, "marker-only", false)
|
||||
er.getDisks = func() []StorageAPI { return []StorageAPI{all[0], all[1], nil, nil} }
|
||||
deleteFn := er.DeleteObject
|
||||
if pooled {
|
||||
deleteFn = z.DeleteObject
|
||||
}
|
||||
_, _ = deleteFn(t.Context(), bucket, "marker-only", opts)
|
||||
_, err := deleteFn(t.Context(), bucket, "marker-only", opts)
|
||||
if !isErrWriteQuorum(err) {
|
||||
t.Errorf("missing retry with only 2/4 absence votes returned %T %v", err, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteMarkerHealReplicationIdentity(t *testing.T) {
|
||||
obj, er, all, bucket := markerPurgeFixture(t, 4)
|
||||
const name = "healed-marker"
|
||||
_, opts := seedPurgeMarker(t, er, bucket, name, true)
|
||||
original, err := all[3].ReadVersion(t.Context(), "", bucket, name, opts.VersionID, ReadOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, disk := range all[:2] {
|
||||
if err := disk.DeleteVersion(t.Context(), bucket, name, FileInfo{VersionID: opts.VersionID}, false, DeleteOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if _, err := er.HealObject(t.Context(), bucket, name, opts.VersionID, madmin.HealOpts{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for i, disk := range all {
|
||||
fi, err := disk.ReadVersion(t.Context(), "", bucket, name, opts.VersionID, ReadOptions{})
|
||||
if err != nil || !maps.Equal(fi.Metadata, original.Metadata) {
|
||||
t.Errorf("disk %d lost marker metadata: before=%v after=%v error=%v", i, original.Metadata, fi.Metadata, err)
|
||||
}
|
||||
}
|
||||
// Only the healed half remains readable. Read actual disk state before
|
||||
// invoking scanner replication, rather than constructing an ObjectInfo.
|
||||
er.getDisks = func() []StorageAPI { return []StorageAPI{all[0], all[1], nil, nil} }
|
||||
oi, err := er.GetObjectInfo(t.Context(), bucket, name, ObjectOptions{VersionID: opts.VersionID, Versioned: true})
|
||||
if !oi.DeleteMarker || !isErrMethodNotAllowed(err) {
|
||||
t.Fatalf("healed marker unreadable: %+v %v", oi, err)
|
||||
}
|
||||
er.getDisks = func() []StorageAPI { return all }
|
||||
var outbound atomic.Int32
|
||||
remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method == http.MethodDelete {
|
||||
outbound.Add(1)
|
||||
if r.Header.Get(xhttp.MinIOSourceDeleteMarker) != "true" || r.URL.Query().Get("versionId") != opts.VersionID {
|
||||
t.Errorf("unexpected outbound request: %s %s", r.Method, r.URL)
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusNotFound)
|
||||
}))
|
||||
defer remote.Close()
|
||||
client, err := minio.New(strings.TrimPrefix(remote.URL, "http://"), &minio.Options{Region: "us-east-1", MaxRetries: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
arn := "arn:minio:replication::" + mustGetUUID() + ":bucket"
|
||||
target := &TargetClient{Client: client, ARN: arn, Bucket: "target"}
|
||||
oldTargets := globalBucketTargetSys
|
||||
globalBucketTargetSys = &BucketTargetSys{arnRemotesMap: map[string]arnTarget{arn: {Client: target, lastRefresh: UTCNow()}}, targetsMap: map[string][]madmin.BucketTarget{bucket: {{Arn: arn, TargetBucket: "target"}}}, hc: map[string]epHealth{client.EndpointURL().Host: {Online: true}}}
|
||||
defer func() { globalBucketTargetSys = oldTargets }()
|
||||
cfg := configs[0]
|
||||
cfg.RoleArn = arn
|
||||
meta, err := globalBucketMetadataSys.Get(bucket)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
meta.replicationConfig = &cfg
|
||||
globalBucketMetadataSys.Set(bucket, meta)
|
||||
worker := make(chan ReplicationWorkerOperation, 1)
|
||||
oldPool := globalReplicationPool
|
||||
globalReplicationPool = once.NewSingleton[ReplicationPool]()
|
||||
globalReplicationPool.Set(&ReplicationPool{ctx: t.Context(), objLayer: obj, workers: []chan ReplicationWorkerOperation{worker}, stats: globalReplicationStats.Load(), mrfSaveCh: make(chan MRFReplicateEntry, 1)})
|
||||
defer func() { globalReplicationPool = oldPool }()
|
||||
roi := queueReplicationHeal(t.Context(), bucket, oi, replicationConfig{Config: &cfg, remotes: &madmin.BucketTargets{Targets: []madmin.BucketTarget{{Arn: arn, TargetBucket: "target"}}}}, 0)
|
||||
select {
|
||||
case op := <-worker:
|
||||
d := op.(DeletedObjectReplicationInfo)
|
||||
result := replicateDeleteToTarget(t.Context(), d, target)
|
||||
t.Errorf("healed replica scheduled for creation: version=%s existing=%v result=%+v", d.DeleteMarkerVersionID, roi.ExistingObjResync.mustResync(), result)
|
||||
default:
|
||||
}
|
||||
if outbound.Load() != 0 {
|
||||
t.Errorf("healed replica sent %d new marker creations", outbound.Load())
|
||||
}
|
||||
}
|
||||
@@ -31,6 +31,7 @@ func TestMultipartListingAbortBetweenHTTPPages(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
setMultipartListingTestMode(t, false)
|
||||
var firstID, secondID string
|
||||
for attempt := 0; attempt < 32; attempt++ {
|
||||
one, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
||||
@@ -147,6 +148,22 @@ func TestMultipartListingLegacyPreflight(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
z.mpCache.Clear()
|
||||
t.Run("legacy-upgrade", func(t *testing.T) {
|
||||
setMultipartListingTestMode(t, true)
|
||||
exact, err := z.ListMultipartUploads(t.Context(), other, "old", "", "", "", 10)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
requireMultipartUploadKeys(t, exact, "old")
|
||||
const empty = "multipart-legacy-empty"
|
||||
if err := z.MakeBucket(t.Context(), empty, MakeBucketOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got, err := z.ListMultipartUploads(t.Context(), empty, "", "", "", "", 10)
|
||||
if err != nil || len(got.Uploads) != 0 {
|
||||
t.Fatalf("old upload in another bucket broke legacy listing: %+v %v", got, err)
|
||||
}
|
||||
})
|
||||
_, err = z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 10)
|
||||
if !errors.Is(err, errMultipartListingLegacy) {
|
||||
t.Fatalf("old upload in another bucket: %v", err)
|
||||
@@ -214,6 +231,7 @@ func TestMultipartListingIdentityFallback(t *testing.T) {
|
||||
|
||||
func TestMultipartAbortPoolsAndRetry(t *testing.T) {
|
||||
z, bucket := consistencyPools(t)
|
||||
setMultipartListingTestMode(t, false)
|
||||
mp, err := z.serverPools[1].NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -286,6 +304,7 @@ func TestMultipartListingMarkerHTTP(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
setMultipartListingTestMode(t, false)
|
||||
for _, tc := range []struct {
|
||||
key, marker string
|
||||
status int
|
||||
@@ -369,6 +388,7 @@ func TestMultipartAbortLateCreateBoundary(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
z := obj.(*erasureServerPools)
|
||||
setMultipartListingTestMode(t, false)
|
||||
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
|
||||
saved := globalStorageClass
|
||||
globalStorageClass.Update(storageclass.Config{Standard: storageclass.StorageClass{Parity: 8}})
|
||||
@@ -454,6 +474,7 @@ func TestMultipartListingScanCosts(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
z := obj.(*erasureServerPools)
|
||||
setMultipartListingTestMode(t, false)
|
||||
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
|
||||
const bucket, otherBucket = "r9-scan-target", "r9-scan-unrelated"
|
||||
for _, name := range []string{bucket, otherBucket} {
|
||||
@@ -523,6 +544,7 @@ func multipartListingFixture(t *testing.T) (*erasureServerPools, *erasureObjects
|
||||
if err := z.MakeBucket(t.Context(), bucket, MakeBucketOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
setMultipartListingTestMode(t, false)
|
||||
return z, z.serverPools[0].getHashedSet("a"), bucket
|
||||
}
|
||||
|
||||
@@ -673,6 +695,7 @@ func TestMultipartListingAdmissionHTTP(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
setMultipartListingTestMode(t, false)
|
||||
for range cap(multipartScanSlots) {
|
||||
scan, err := startMultipartScan(t.Context(), false)
|
||||
if err != nil {
|
||||
|
||||
@@ -0,0 +1,331 @@
|
||||
// Copyright (c) 2026 Ruohang Feng
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
|
||||
"github.com/minio/minio/internal/config/api"
|
||||
"github.com/minio/minio/internal/config/storageclass"
|
||||
"github.com/minio/minio/internal/dsync"
|
||||
"github.com/minio/minio/internal/grid"
|
||||
xnet "github.com/pgsty/silo-pkg/v3/net"
|
||||
)
|
||||
|
||||
func setMultipartListingTestMode(t *testing.T, legacy bool) {
|
||||
t.Helper()
|
||||
globalAPIConfig.mu.Lock()
|
||||
previous := globalAPIConfig.multipartListingStrict
|
||||
globalAPIConfig.multipartListingStrict = !legacy
|
||||
globalAPIConfig.mu.Unlock()
|
||||
t.Cleanup(func() {
|
||||
globalAPIConfig.mu.Lock()
|
||||
globalAPIConfig.multipartListingStrict = previous
|
||||
globalAPIConfig.mu.Unlock()
|
||||
})
|
||||
}
|
||||
|
||||
func TestMultipartListingDefaultMode(t *testing.T) {
|
||||
var uninitialized apiConfig
|
||||
if !uninitialized.getMultipartListingLegacy() {
|
||||
t.Fatal("runtime zero value enabled strict mode before config initialization")
|
||||
}
|
||||
for _, mode := range []string{"", "legacy", "strict"} {
|
||||
var local apiConfig
|
||||
local.init(api.Config{RequestsMax: 1, MultipartListing: mode}, []int{4}, false)
|
||||
if local.getMultipartListingLegacy() != (mode != "strict") {
|
||||
t.Fatalf("unexpected mode for %q", mode)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultipartAbortLegacyAvailability(t *testing.T) {
|
||||
for _, drives := range []int{4, 16} {
|
||||
for _, offline := range []int{0, 1, drives / 2} {
|
||||
t.Run(fmt.Sprintf("drives=%d/offline=%d", drives, offline), func(t *testing.T) {
|
||||
obj, dirs, err := prepareErasure(t.Context(), drives)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
z := obj.(*erasureServerPools)
|
||||
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
|
||||
setMultipartListingTestMode(t, true)
|
||||
previous := globalStorageClass
|
||||
globalStorageClass.Update(storageclass.Config{Standard: storageclass.StorageClass{Parity: drives / 2}})
|
||||
t.Cleanup(func() { globalStorageClass.Update(previous) })
|
||||
const bucket, key = "multipart-legacy-availability", "keep-available"
|
||||
if err := z.MakeBucket(t.Context(), bucket, MakeBucketOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
mp, err := z.NewMultipartUpload(t.Context(), bucket, key, ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
set := z.serverPools[0].getHashedSet(key)
|
||||
original := set.getDisks
|
||||
disks := original()
|
||||
set.getDisks = func() []StorageAPI {
|
||||
visible := append([]StorageAPI(nil), disks...)
|
||||
clear(visible[:offline])
|
||||
return visible
|
||||
}
|
||||
t.Cleanup(func() { set.getDisks = original })
|
||||
if err := z.AbortMultipartUpload(t.Context(), bucket, key, mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatalf("released read-quorum availability regressed: %v", err)
|
||||
}
|
||||
for _, disk := range disks[offline:] {
|
||||
_, err := disk.ReadVersion(t.Context(), bucket, minioMetaMultipartBucket, set.getUploadIDDir(bucket, key, mp.UploadID), "", ReadOptions{})
|
||||
if !errors.Is(err, errFileNotFound) {
|
||||
t.Fatalf("online replica not cleaned: %v", err)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultipartAbortLegacyPoolOrder(t *testing.T) {
|
||||
for _, owner := range []int{0, 1} {
|
||||
t.Run(fmt.Sprintf("owner=%d", owner), func(t *testing.T) {
|
||||
z, bucket := consistencyPools(t)
|
||||
setMultipartListingTestMode(t, true)
|
||||
mp, err := z.serverPools[owner].NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
emptySet := z.serverPools[1-owner].getHashedSet("a")
|
||||
original := emptySet.getDisks
|
||||
t.Cleanup(func() { emptySet.getDisks = original })
|
||||
emptySet.getDisks = func() []StorageAPI { return make([]StorageAPI, emptySet.setDriveCount) }
|
||||
err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{})
|
||||
if owner == 0 && err != nil {
|
||||
t.Fatalf("unrelated later pool blocked legacy cancellation: %v", err)
|
||||
}
|
||||
if owner == 1 {
|
||||
if err == nil || toAPIError(t.Context(), err).HTTPStatusCode != 503 {
|
||||
t.Fatalf("unknown earlier pool must retain released error behavior: %v", err)
|
||||
}
|
||||
if _, err := z.serverPools[owner].GetMultipartInfo(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatalf("later pool visited despite earlier error: %v", err)
|
||||
}
|
||||
emptySet.getDisks = original
|
||||
if err := z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultipartAbortMinorityLegacyCleanup(t *testing.T) {
|
||||
for _, failDelete := range []bool{false, true} {
|
||||
t.Run(fmt.Sprintf("delete-fails=%v", failDelete), func(t *testing.T) {
|
||||
z, set, bucket := multipartListingFixture(t)
|
||||
mp, err := z.NewMultipartUpload(t.Context(), bucket, "old", ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
path := set.getUploadIDDir(bucket, "old", mp.UploadID)
|
||||
original := set.getDisks
|
||||
disks := original()
|
||||
fi, err := disks[0].ReadVersion(t.Context(), bucket, minioMetaMultipartBucket, path, "", ReadOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
delete(fi.Metadata, multipartMetaBucket)
|
||||
delete(fi.Metadata, multipartMetaObject)
|
||||
if err := disks[0].WriteMetadata(t.Context(), bucket, minioMetaMultipartBucket, path, fi); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, d := range disks[1:] {
|
||||
if err := d.Delete(t.Context(), minioMetaMultipartBucket, path, DeleteOptions{Recursive: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
report, err := z.multipartPreflight(t.Context())
|
||||
if err != nil || report.Ready || report.LegacyUploads != 1 {
|
||||
t.Fatalf("invalid legacy fixture: %+v %v", report, err)
|
||||
}
|
||||
if failDelete {
|
||||
set.getDisks = func() []StorageAPI {
|
||||
wrapped := append([]StorageAPI(nil), disks...)
|
||||
wrapped[0] = multipartListingFaultDisk{StorageAPI: disks[0], delete: func(context.Context, string, string, DeleteOptions) error { return errFaultyDisk }}
|
||||
return wrapped
|
||||
}
|
||||
t.Cleanup(func() { set.getDisks = original })
|
||||
err := z.AbortMultipartUpload(t.Context(), bucket, "old", mp.UploadID, ObjectOptions{})
|
||||
if err == nil || toAPIError(t.Context(), err).HTTPStatusCode != 503 {
|
||||
t.Fatalf("empty drives masked failed deletion of the observed replica: %v", err)
|
||||
}
|
||||
if _, present := z.mpCache.Load(mp.UploadID); !present {
|
||||
t.Fatal("failed cancellation evicted upload from cache")
|
||||
}
|
||||
set.getDisks = original
|
||||
}
|
||||
if err := z.AbortMultipartUpload(t.Context(), bucket, "old", mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
report, err = z.multipartPreflight(t.Context())
|
||||
if err != nil || !report.Ready || report.LegacyUploads != 0 {
|
||||
t.Fatalf("acknowledged cleanup left online legacy remnants: %+v %v", report, err)
|
||||
}
|
||||
if _, present := z.mpCache.Load(mp.UploadID); present {
|
||||
t.Fatal("successful cancellation retained cache entry")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultipartAbortWrongTargetKeepsCache(t *testing.T) {
|
||||
for _, legacy := range []bool{true, false} {
|
||||
t.Run(fmt.Sprintf("legacy=%v", legacy), func(t *testing.T) {
|
||||
z, _, bucket := multipartListingFixture(t)
|
||||
setMultipartListingTestMode(t, legacy)
|
||||
const other = "multipart-other-bucket"
|
||||
if err := z.MakeBucket(t.Context(), other, MakeBucketOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
mp, err := z.NewMultipartUpload(t.Context(), bucket, "valid-key", ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, target := range [][2]string{{bucket, "wrong-key"}, {other, "valid-key"}} {
|
||||
err := z.AbortMultipartUpload(t.Context(), target[0], target[1], mp.UploadID, ObjectOptions{})
|
||||
var invalid InvalidUploadID
|
||||
if !errors.As(err, &invalid) {
|
||||
t.Fatalf("wrong target result: %v", err)
|
||||
}
|
||||
if _, present := z.mpCache.Load(mp.UploadID); !present {
|
||||
t.Fatal("wrong-target abort evicted another upload")
|
||||
}
|
||||
if _, err := z.GetMultipartInfo(t.Context(), bucket, "valid-key", mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatalf("wrong-target abort harmed valid upload: %v", err)
|
||||
}
|
||||
}
|
||||
if err := z.AbortMultipartUpload(t.Context(), bucket, "valid-key", mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var invalid InvalidUploadID
|
||||
if err := z.AbortMultipartUpload(t.Context(), bucket, "valid-key", mp.UploadID, ObjectOptions{}); !errors.As(err, &invalid) {
|
||||
t.Fatalf("repeat abort: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultipartAbortRejectsUnsafeID(t *testing.T) {
|
||||
for _, legacy := range []bool{true, false} {
|
||||
t.Run(fmt.Sprintf("legacy=%v", legacy), func(t *testing.T) {
|
||||
z, _, bucket := multipartListingFixture(t)
|
||||
setMultipartListingTestMode(t, legacy)
|
||||
mp, err := z.NewMultipartUpload(t.Context(), bucket, "keep", ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, suffix := range []string{".", "..", "../other", "a/b", "a\\b", ""} {
|
||||
id := base64.RawURLEncoding.EncodeToString([]byte("deployment." + suffix))
|
||||
err := z.AbortMultipartUpload(t.Context(), bucket, "keep", id, ObjectOptions{})
|
||||
var invalid InvalidUploadID
|
||||
if !errors.As(err, &invalid) {
|
||||
t.Fatalf("unsafe ID accepted: suffix=%q err=%v", suffix, err)
|
||||
}
|
||||
if _, err := z.GetMultipartInfo(t.Context(), bucket, "keep", mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatalf("valid upload harmed by rejected ID: %v", err)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultipartAbortPeerNotification(t *testing.T) {
|
||||
z, set, bucket := multipartListingFixture(t)
|
||||
tg, err := grid.SetupTestGrid(2)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(tg.Cleanup)
|
||||
var notifications atomic.Int32
|
||||
if err := cleanupUploadIDCacheMetaRPC.Register(tg.Managers[1], func(*grid.MSS) (grid.NoPayload, *grid.RemoteErr) {
|
||||
notifications.Add(1)
|
||||
return grid.NoPayload{}, nil
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
host, err := xnet.ParseHost(strings.TrimPrefix(tg.Hosts[1], "http://"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
previous := globalNotificationSys
|
||||
globalNotificationSys = &NotificationSys{peerClients: []*peerRESTClient{{
|
||||
host: host,
|
||||
gridConn: func() *grid.Connection { return tg.Managers[0].Connection(tg.Hosts[1]) },
|
||||
}}}
|
||||
t.Cleanup(func() { globalNotificationSys = previous })
|
||||
// Use the real distributed lock implementation: Unlock cancels its derived
|
||||
// context. Local locks do not, and would hide a broken notification context.
|
||||
oldMutex, oldLockers := set.nsMutex, set.getLockers
|
||||
lockers := []dsync.NetLocker{newLocker(), newLocker()}
|
||||
set.nsMutex = newNSLock(true)
|
||||
set.getLockers = func() ([]dsync.NetLocker, string) { return lockers, "multipart-test" }
|
||||
t.Cleanup(func() { set.nsMutex, set.getLockers = oldMutex, oldLockers })
|
||||
for _, legacy := range []bool{true, false} {
|
||||
t.Run(fmt.Sprintf("legacy=%v", legacy), func(t *testing.T) {
|
||||
setMultipartListingTestMode(t, legacy)
|
||||
for range 4 {
|
||||
mp, err := z.NewMultipartUpload(t.Context(), bucket, "valid", ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
before := notifications.Load()
|
||||
var invalid InvalidUploadID
|
||||
if err := z.AbortMultipartUpload(t.Context(), bucket, "wrong", mp.UploadID, ObjectOptions{}); !errors.As(err, &invalid) {
|
||||
t.Fatalf("wrong target: %v", err)
|
||||
}
|
||||
if notifications.Load() != before {
|
||||
t.Fatal("wrong-target abort sent a destructive peer notification")
|
||||
}
|
||||
if err := z.AbortMultipartUpload(t.Context(), bucket, "valid", mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if notifications.Load() != before+1 {
|
||||
t.Fatal("successful abort lost peer notification after unlocking")
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultipartAbortStrictMajority(t *testing.T) {
|
||||
z, set, bucket := multipartListingFixture(t)
|
||||
mp, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
original := set.getDisks
|
||||
disks := original()
|
||||
set.getDisks = func() []StorageAPI {
|
||||
wrapped := append([]StorageAPI(nil), disks...)
|
||||
wrapped[0] = multipartListingFaultDisk{StorageAPI: disks[0], delete: func(context.Context, string, string, DeleteOptions) error { return errFaultyDisk }}
|
||||
return wrapped
|
||||
}
|
||||
t.Cleanup(func() { set.getDisks = original })
|
||||
if err := z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatalf("ordinary majority cancellation was tightened: %v", err)
|
||||
}
|
||||
// One failed deletion is tolerated in the ordinary majority path. Retrying
|
||||
// after recovery must also clean the now-minority remnant.
|
||||
if _, err := disks[0].ReadVersion(t.Context(), bucket, minioMetaMultipartBucket, set.getUploadIDDir(bucket, "a", mp.UploadID), "", ReadOptions{}); err != nil {
|
||||
t.Fatalf("fault fixture did not leave a remnant: %v", err)
|
||||
}
|
||||
set.getDisks = original
|
||||
if err := z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
+37
-12
@@ -1753,11 +1753,11 @@ func (er erasureObjects) CompleteMultipartUpload(ctx context.Context, bucket str
|
||||
return fi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended), nil
|
||||
}
|
||||
|
||||
// abortMultipartUpload confirms absence on a strict majority of this set.
|
||||
// Unlike an existence read followed by best-effort deletion, it also permits
|
||||
// retrying a partial deletion which no longer has a readable metadata quorum.
|
||||
// This does not fence creation writes still executing after a storage timeout.
|
||||
func (er erasureObjects) abortMultipartUpload(ctx context.Context, bucket, object, uploadID string, opts ObjectOptions) (bool, error) {
|
||||
// abortMultipartUpload retains read-quorum validation and best-effort cleanup
|
||||
// in legacy mode. Strict mode requires majority deletion acknowledgements and
|
||||
// permits retrying remnants below read quorum. Neither mode fences creation
|
||||
// writes still executing after a storage timeout.
|
||||
func (er erasureObjects) abortMultipartUpload(ctx context.Context, bucket, object, uploadID string, opts ObjectOptions, legacy bool) (bool, error) {
|
||||
if !opts.NoAuditLog {
|
||||
auditObjectErasureSet(ctx, "AbortMultipartUpload", object, &er)
|
||||
}
|
||||
@@ -1769,21 +1769,34 @@ func (er erasureObjects) abortMultipartUpload(ctx context.Context, bucket, objec
|
||||
if !ok || internalID == "" || internalID == "." || internalID == ".." || strings.ContainsAny(internalID, "/\\") {
|
||||
return false, InvalidUploadID{Bucket: bucket, Object: object, UploadID: uploadID}
|
||||
}
|
||||
if legacy {
|
||||
// Keep the released read-quorum and best-effort cleanup behavior.
|
||||
// The upload ID safety check above applies to both modes.
|
||||
defer er.deleteAll(ctx, minioMetaMultipartBucket, er.getUploadIDDir(bucket, object, uploadID))
|
||||
_, _, err := er.checkUploadIDExists(ctx, bucket, object, uploadID, false)
|
||||
err = toObjectErr(err, bucket, object, uploadID)
|
||||
if _, absent := err.(InvalidUploadID); absent {
|
||||
return false, nil
|
||||
}
|
||||
return err == nil, err
|
||||
}
|
||||
disks := er.getDisks()
|
||||
uploadPath := er.getUploadIDDir(bucket, object, uploadID)
|
||||
_, errs := readAllFileInfo(ctx, disks, bucket, minioMetaMultipartBucket, uploadPath, "", false, false)
|
||||
quorum := er.setDriveCount/2 + 1
|
||||
found, absent := false, 0
|
||||
for _, err := range errs {
|
||||
observed := make([]bool, len(disks))
|
||||
for i, err := range errs {
|
||||
switch {
|
||||
case err == nil, errors.Is(err, errFileCorrupt):
|
||||
found = true
|
||||
observed[i] = true
|
||||
case errors.Is(err, errFileNotFound), errors.Is(err, errFileVersionNotFound):
|
||||
absent++
|
||||
}
|
||||
}
|
||||
if absent >= quorum {
|
||||
return found, nil
|
||||
if absent >= quorum && !found {
|
||||
return false, nil
|
||||
}
|
||||
if !found {
|
||||
return false, toObjectErr(errErasureReadQuorum, bucket, object, uploadID)
|
||||
@@ -1801,13 +1814,25 @@ func (er erasureObjects) abortMultipartUpload(ctx context.Context, bucket, objec
|
||||
return err
|
||||
}, i)
|
||||
}
|
||||
return true, toObjectErr(reduceWriteQuorumErrs(ctx, g.Wait(), nil, quorum), bucket, object, uploadID)
|
||||
deleteErrs := g.Wait()
|
||||
if err := reduceWriteQuorumErrs(ctx, deleteErrs, nil, quorum); err != nil {
|
||||
return true, toObjectErr(err, bucket, object, uploadID)
|
||||
}
|
||||
if absent >= quorum {
|
||||
// Missing disks must not mask a failed cleanup of known remnants.
|
||||
for i, found := range observed {
|
||||
if found && deleteErrs[i] != nil {
|
||||
return true, toObjectErr(errErasureWriteQuorum, bucket, object, uploadID)
|
||||
}
|
||||
}
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// AbortMultipartUpload confirms logical cancellation. Offline part data may
|
||||
// still need stale-upload cleanup after its drives return.
|
||||
// AbortMultipartUpload cancels an upload using the configured mode. Offline
|
||||
// part data may still need stale-upload cleanup after its drives return.
|
||||
func (er erasureObjects) AbortMultipartUpload(ctx context.Context, bucket, object, uploadID string, opts ObjectOptions) (err error) {
|
||||
found, err := er.abortMultipartUpload(ctx, bucket, object, uploadID, opts)
|
||||
found, err := er.abortMultipartUpload(ctx, bucket, object, uploadID, opts, globalAPIConfig.getMultipartListingLegacy())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
+32
-6
@@ -1657,7 +1657,7 @@ func (er erasureObjects) putObject(ctx context.Context, bucket string, object st
|
||||
return fi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended), nil
|
||||
}
|
||||
|
||||
func (er erasureObjects) deleteObjectVersion(ctx context.Context, bucket, object string, fi FileInfo, forceDelMarker bool) error {
|
||||
func (er erasureObjects) deleteObjectVersion(ctx context.Context, bucket, object string, fi FileInfo, forceDelMarker, purge bool) error {
|
||||
disks := er.getDisks()
|
||||
// Assume (N/2 + 1) quorum for Delete()
|
||||
// this is a theoretical assumption such that
|
||||
@@ -1673,7 +1673,13 @@ func (er erasureObjects) deleteObjectVersion(ctx context.Context, bucket, object
|
||||
if disks[index] == nil {
|
||||
return errDiskNotFound
|
||||
}
|
||||
return disks[index].DeleteVersion(ctx, bucket, object, fi, forceDelMarker, DeleteOptions{})
|
||||
err := disks[index].DeleteVersion(ctx, bucket, object, fi, forceDelMarker, DeleteOptions{})
|
||||
// Physical removal and reliable already-absent replies are the same
|
||||
// outcome. Creation and metadata updates must not use this quorum.
|
||||
if purge && (err == errFileNotFound || err == errFileVersionNotFound) {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}, index)
|
||||
}
|
||||
// return errors if any during deletion
|
||||
@@ -2023,6 +2029,11 @@ func (er erasureObjects) DeleteObject(ctx context.Context, bucket, object string
|
||||
if opts.DeleteMarker {
|
||||
versionFound = false
|
||||
} else if !tryDel {
|
||||
if opts.isVersionPurge() && (isErrObjectNotFound(gerr) || isErrVersionNotFound(gerr)) {
|
||||
if err := er.checkPurgeAbsent(ctx, bucket, object, opts.VersionID); err != nil {
|
||||
return objInfo, err
|
||||
}
|
||||
}
|
||||
return objInfo, gerr
|
||||
}
|
||||
}
|
||||
@@ -2105,6 +2116,13 @@ func (er erasureObjects) DeleteObject(ctx context.Context, bucket, object string
|
||||
}
|
||||
}
|
||||
|
||||
purge := opts.isVersionPurge()
|
||||
if purge {
|
||||
// A delete marker describes the version being removed; it is not an
|
||||
// instruction to create that marker on disks which already lack it.
|
||||
markDelete, deleteMarker = false, false
|
||||
}
|
||||
|
||||
modTime := opts.MTime
|
||||
if opts.MTime.IsZero() {
|
||||
modTime = UTCNow()
|
||||
@@ -2151,7 +2169,7 @@ func (er erasureObjects) DeleteObject(ctx context.Context, bucket, object string
|
||||
// delete marker. Add delete marker, since we don't have
|
||||
// any version specified explicitly. Or if a particular
|
||||
// version id needs to be replicated.
|
||||
if err = er.deleteObjectVersion(ctx, bucket, object, fi, opts.DeleteMarker); err != nil {
|
||||
if err = er.deleteObjectVersion(ctx, bucket, object, fi, opts.DeleteMarker, false); err != nil {
|
||||
return objInfo, toObjectErr(err, bucket, object)
|
||||
}
|
||||
oi := fi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended)
|
||||
@@ -2174,11 +2192,19 @@ func (er erasureObjects) DeleteObject(ctx context.Context, bucket, object string
|
||||
if opts.SkipFreeVersion {
|
||||
dfi.SetSkipTierFreeVersion()
|
||||
}
|
||||
if err = er.deleteObjectVersion(ctx, bucket, object, dfi, opts.DeleteMarker); err != nil {
|
||||
if err = er.deleteObjectVersion(ctx, bucket, object, dfi, opts.DeleteMarker, purge); err != nil {
|
||||
return objInfo, toObjectErr(err, bucket, object)
|
||||
}
|
||||
|
||||
return dfi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended), nil
|
||||
oi := dfi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended)
|
||||
if purge {
|
||||
// Preserve the DELETE response's identity without reusing it as a disk
|
||||
// instruction to create a marker. The lookup also exposes a data
|
||||
// version pending purge as deleted for visibility; only a stored
|
||||
// marker, which carries no erasure layout, is reported as one.
|
||||
oi.DeleteMarker = goi.DeleteMarker && goi.DataBlocks == 0
|
||||
}
|
||||
return oi, nil
|
||||
}
|
||||
|
||||
// Send the successful but partial upload/delete, however ignore
|
||||
@@ -2467,7 +2493,7 @@ func (er erasureObjects) TransitionObject(ctx context.Context, bucket, object st
|
||||
|
||||
storageDisks := er.getDisks()
|
||||
|
||||
if err = er.deleteObjectVersion(ctx, bucket, object, fi, false); err != nil {
|
||||
if err = er.deleteObjectVersion(ctx, bucket, object, fi, false, false); err != nil {
|
||||
eventName = event.ObjectTransitionFailed
|
||||
}
|
||||
|
||||
|
||||
@@ -316,6 +316,11 @@ func (z *erasureServerPools) retireReplicaCopies(ctx context.Context, bucket, ob
|
||||
func (z *erasureServerPools) deleteObjectReconciled(ctx context.Context, bucket, object string, opts ObjectOptions) (ObjectInfo, error) {
|
||||
copies, err := z.objectPoolInfos(ctx, bucket, object, opts)
|
||||
if err != nil {
|
||||
if opts.isVersionPurge() && (isErrObjectNotFound(err) || isErrVersionNotFound(err)) {
|
||||
if quorumErr := z.checkPurgeAbsent(ctx, bucket, object, opts.VersionID); quorumErr != nil {
|
||||
return ObjectInfo{}, quorumErr
|
||||
}
|
||||
}
|
||||
return ObjectInfo{}, err
|
||||
}
|
||||
primary := copies[0]
|
||||
@@ -365,6 +370,21 @@ func (z *erasureServerPools) deleteObjectReconciled(ctx context.Context, bucket,
|
||||
opts.EvalMetadataFn = nil
|
||||
}
|
||||
}
|
||||
if opts.isVersionPurge() {
|
||||
// A missing pool may have only read-quorum absence. Verify every pool
|
||||
// omitted by lookup before mutating the known copies.
|
||||
for i, pool := range z.serverPools {
|
||||
found := false
|
||||
for _, copy := range copies {
|
||||
found = found || copy.Index == i
|
||||
}
|
||||
if !found {
|
||||
if err := pool.getHashedSet(object).checkPurgeAbsent(ctx, bucket, object, opts.VersionID); err != nil {
|
||||
return ObjectInfo{}, err
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if opts.VersionID == "" && (opts.Versioned || opts.VersionSuspended) {
|
||||
// A single new delete marker hides the current version. Older versions
|
||||
// remain history and must not be removed by an unqualified DELETE.
|
||||
|
||||
@@ -618,7 +618,7 @@ func (z *erasureServerPools) decommissionObject(ctx context.Context, idx int, bu
|
||||
if objInfo.isMultipart() {
|
||||
res, err := z.NewMultipartUpload(ctx, bucket, objInfo.Name, ObjectOptions{
|
||||
VersionID: objInfo.VersionID,
|
||||
UserDefined: objInfo.UserDefined,
|
||||
UserDefined: migrationObjectMetadata(objInfo),
|
||||
NoAuditLog: true,
|
||||
SrcPoolIdx: idx,
|
||||
DataMovement: true,
|
||||
@@ -681,7 +681,7 @@ func (z *erasureServerPools) decommissionObject(ctx context.Context, idx int, bu
|
||||
SrcPoolIdx: idx,
|
||||
VersionID: objInfo.VersionID,
|
||||
MTime: objInfo.ModTime,
|
||||
UserDefined: objInfo.UserDefined,
|
||||
UserDefined: migrationObjectMetadata(objInfo),
|
||||
PreserveETag: objInfo.ETag, // Preserve original ETag to ensure same metadata.
|
||||
IndexCB: func() []byte {
|
||||
return objInfo.Parts[0].Index // Preserve part Index to ensure decompression works.
|
||||
|
||||
@@ -0,0 +1,33 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import xhttp "github.com/minio/minio/internal/http"
|
||||
|
||||
// migrationObjectMetadata restores tags removed from the read representation.
|
||||
// Keep the original revision, including ordered empty states, for the existing
|
||||
// locked reconciliation at PUT or multipart completion. Moving is not a new
|
||||
// tagging mutation, and the write must not modify the reader's metadata map.
|
||||
func migrationObjectMetadata(oi ObjectInfo) map[string]string {
|
||||
metadata := cloneMSS(oi.UserDefined)
|
||||
delete(metadata, xhttp.AmzObjectTagging)
|
||||
if oi.UserTags != "" {
|
||||
metadata[xhttp.AmzObjectTagging] = oi.UserTags
|
||||
}
|
||||
return metadata
|
||||
}
|
||||
@@ -0,0 +1,592 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"maps"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
xhttp "github.com/minio/minio/internal/http"
|
||||
)
|
||||
|
||||
const migrationTagRevision = ReservedMetadataPrefixLower + TaggingTimestamp
|
||||
|
||||
var migrationTagStates = []struct {
|
||||
name, tags, revision string
|
||||
hasRevision bool
|
||||
}{
|
||||
{name: "initial", tags: "team=storage&path=a%2Fb"},
|
||||
{name: "empty-revision", tags: "team=storage", hasRevision: true},
|
||||
{name: "invalid-revision", tags: "team=storage", revision: "invalid", hasRevision: true},
|
||||
{name: "ordered", tags: "team=storage", revision: "2026-09-16T09:00:00Z", hasRevision: true},
|
||||
{name: "cleared", revision: "2026-09-16T09:00:00Z", hasRevision: true},
|
||||
{name: "never-tagged"},
|
||||
{name: "empty-empty-revision", hasRevision: true},
|
||||
{name: "empty-invalid-revision", revision: "invalid", hasRevision: true},
|
||||
}
|
||||
|
||||
func TestMigrationObjectMetadata(t *testing.T) {
|
||||
for _, state := range migrationTagStates {
|
||||
for _, residual := range []bool{false, true} {
|
||||
t.Run(fmt.Sprintf("%s/residual=%t", state.name, residual), func(t *testing.T) {
|
||||
raw := map[string]string{
|
||||
xhttp.AmzObjectTagging: state.tags,
|
||||
"x-amz-meta-team": "storage",
|
||||
"x-amz-meta-empty": "",
|
||||
ReservedMetadataPrefix + "compression": "opaque-compression-state",
|
||||
ReservedMetadataPrefix + "sealed-key": "opaque-key-state",
|
||||
"x-amz-object-lock-retain-until-date": "2030-01-01T00:00:00Z",
|
||||
"x-amz-object-lock-legal-hold": "ON",
|
||||
"x-amz-tagging": "keep-noncanonical-metadata",
|
||||
ReservedMetadataPrefix + "actual-size": "42",
|
||||
ReservedMetadataPrefix + "replica-status": "opaque-replication-state",
|
||||
}
|
||||
if state.hasRevision {
|
||||
raw[migrationTagRevision] = state.revision
|
||||
}
|
||||
oi := (FileInfo{Metadata: raw}).ToObjectInfo("bucket", "object", false)
|
||||
if residual {
|
||||
oi.UserDefined[xhttp.AmzObjectTagging] = "stale=raw"
|
||||
}
|
||||
before := maps.Clone(oi.UserDefined)
|
||||
got := migrationObjectMetadata(oi)
|
||||
if got == nil || !maps.Equal(before, oi.UserDefined) {
|
||||
t.Fatal("helper must return a new non-nil map without changing the read snapshot")
|
||||
}
|
||||
value, exists := got[xhttp.AmzObjectTagging]
|
||||
if value != state.tags || exists != (state.tags != "") {
|
||||
t.Fatalf("raw tags=(%q,%t), want (%q,%t)", value, exists, state.tags, state.tags != "")
|
||||
}
|
||||
stamp, present := got[migrationTagRevision]
|
||||
if stamp != state.revision || present != state.hasRevision {
|
||||
t.Fatalf("revision=(%q,%t), want (%q,%t)", stamp, present, state.revision, state.hasRevision)
|
||||
}
|
||||
for key, value := range before {
|
||||
if stored, exists := got[key]; key != xhttp.AmzObjectTagging && (!exists || stored != value) {
|
||||
t.Errorf("lost metadata %q", key)
|
||||
}
|
||||
}
|
||||
got["x-amz-meta-team"] = "changed"
|
||||
if !maps.Equal(before, oi.UserDefined) {
|
||||
t.Fatal("output aliases the read snapshot")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
for _, metadata := range []map[string]string{nil, {}} {
|
||||
got := migrationObjectMetadata(ObjectInfo{UserDefined: metadata})
|
||||
if got == nil || len(got) != 0 {
|
||||
t.Fatalf("empty input must yield a writable empty map: %#v", got)
|
||||
}
|
||||
got["new"] = "value"
|
||||
if len(metadata) != 0 {
|
||||
t.Fatal("empty map input was aliased")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Allocation ignores the host's used-space percentage, as in the existing tag
|
||||
// fixtures. All object data and metadata still use real erasure-storage disks.
|
||||
func migrationTestPools(t *testing.T) (*erasureServerPools, string) {
|
||||
t.Helper()
|
||||
z, bucket := consistencyPools(t)
|
||||
for _, pool := range z.serverPools {
|
||||
for _, set := range pool.sets {
|
||||
original := set.getDisks
|
||||
disks := append([]StorageAPI(nil), original()...)
|
||||
for i := range disks {
|
||||
disks[i] = tagTestCapacityDisk{StorageAPI: disks[i]}
|
||||
}
|
||||
set.getDisks = func() []StorageAPI { return disks }
|
||||
t.Cleanup(func() { set.getDisks = original })
|
||||
}
|
||||
}
|
||||
return z, bucket
|
||||
}
|
||||
|
||||
func migrationTestMover(t *testing.T, z *erasureServerPools, source int, kind string) func(context.Context, int, string, *GetObjectReader) error {
|
||||
t.Helper()
|
||||
if kind == "rebalance" {
|
||||
z.rebalMu.Lock()
|
||||
z.rebalMeta = &rebalanceMeta{PoolStats: []*rebalanceStats{{}, {}}}
|
||||
z.rebalMeta.PoolStats[source] = &rebalanceStats{Participating: true, Info: rebalanceInfo{Status: rebalStarted}}
|
||||
z.rebalMu.Unlock()
|
||||
t.Cleanup(func() { z.rebalMu.Lock(); z.rebalMeta = nil; z.rebalMu.Unlock() })
|
||||
return z.rebalanceObject
|
||||
}
|
||||
z.poolMetaMutex.Lock()
|
||||
z.poolMeta.Pools[source].Decommission = &PoolDecommissionInfo{}
|
||||
z.poolMetaMutex.Unlock()
|
||||
t.Cleanup(func() { z.poolMetaMutex.Lock(); z.poolMeta.Pools[source].Decommission = nil; z.poolMetaMutex.Unlock() })
|
||||
return z.decommissionObject
|
||||
}
|
||||
|
||||
func migrationTestObject(t *testing.T, z *erasureServerPools, bucket, object string, source int, multipart bool, parts []string, opts ObjectOptions) ObjectInfo {
|
||||
t.Helper()
|
||||
if !multipart {
|
||||
oi := putConsistencyObject(t, z, bucket, object, source, strings.Join(parts, ""), opts)
|
||||
return migrationTestPersistedInfo(t, z, source, oi)
|
||||
}
|
||||
mp, err := z.serverPools[source].NewMultipartUpload(t.Context(), bucket, object, opts)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
completed := make([]CompletePart, len(parts))
|
||||
for i, body := range parts {
|
||||
part, err := z.serverPools[source].PutObjectPart(t.Context(), bucket, object, mp.UploadID, i+1,
|
||||
mustGetPutObjReader(t, strings.NewReader(body), int64(len(body)), "", ""), ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
completed[i] = CompletePart{PartNumber: i + 1, ETag: part.ETag}
|
||||
}
|
||||
oi, err := z.serverPools[source].CompleteMultipartUpload(t.Context(), bucket, object, mp.UploadID, completed, opts)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !oi.isMultipart() {
|
||||
t.Fatal("fixture did not create a multipart object")
|
||||
}
|
||||
return migrationTestPersistedInfo(t, z, source, oi)
|
||||
}
|
||||
|
||||
func migrationTestPersistedInfo(t *testing.T, z *erasureServerPools, pool int, oi ObjectInfo) ObjectInfo {
|
||||
t.Helper()
|
||||
// PUT can return transient fields such as tier-free-versionID that are not
|
||||
// stored. Take the same persisted read representation used by the movers.
|
||||
got, err := z.serverPools[pool].GetObjectInfo(t.Context(), oi.Bucket, oi.Name, ObjectOptions{VersionID: migrationTestVersion(oi)})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return got
|
||||
}
|
||||
|
||||
func migrationTestVersion(oi ObjectInfo) string {
|
||||
if oi.VersionID == "" {
|
||||
return nullVersionID
|
||||
}
|
||||
return oi.VersionID
|
||||
}
|
||||
|
||||
func migrationTestReader(t *testing.T, z *erasureServerPools, source int, oi ObjectInfo) *GetObjectReader {
|
||||
t.Helper()
|
||||
gr, err := z.serverPools[source].GetObjectNInfo(t.Context(), oi.Bucket, oi.Name, nil, nil,
|
||||
ObjectOptions{VersionID: migrationTestVersion(oi), NoLock: true, NoDecryption: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { gr.Close() })
|
||||
return gr
|
||||
}
|
||||
|
||||
func migrationTestStored(t *testing.T, z *erasureServerPools, pool int, original ObjectInfo, tags, revision string, hasRevision bool, body string) {
|
||||
t.Helper()
|
||||
version := migrationTestVersion(original)
|
||||
got, err := z.serverPools[pool].GetObjectInfo(t.Context(), original.Bucket, original.Name, ObjectOptions{VersionID: version})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got.UserTags != tags || got.ETag != original.ETag || !got.ModTime.Equal(original.ModTime) || got.VersionID != original.VersionID || got.Size != original.Size {
|
||||
t.Errorf("pool %d changed tags/object identity: tags=%q want=%q version=%q/%q etag=%q/%q size=%d/%d mtime=%v/%v", pool, got.UserTags, tags, got.VersionID, original.VersionID, got.ETag, original.ETag, got.Size, original.Size, got.ModTime, original.ModTime)
|
||||
}
|
||||
infos, errs := readAllFileInfo(t.Context(), z.serverPools[pool].getHashedSet(original.Name).getDisks(), "", original.Bucket, original.Name, version, false, false)
|
||||
for disk, fi := range infos {
|
||||
if errs[disk] != nil {
|
||||
t.Fatalf("pool %d disk %d: %v", pool, disk, errs[disk])
|
||||
}
|
||||
stamp, exists := fi.Metadata[migrationTagRevision]
|
||||
if fi.Metadata[xhttp.AmzObjectTagging] != tags || stamp != revision || exists != hasRevision {
|
||||
t.Errorf("pool %d disk %d: tags=%q revision=(%q,%t), want %q (%q,%t)", pool, disk, fi.Metadata[xhttp.AmzObjectTagging], stamp, exists, tags, revision, hasRevision)
|
||||
}
|
||||
for key, value := range original.UserDefined {
|
||||
if key != migrationTagRevision && fi.Metadata[key] != value {
|
||||
t.Errorf("pool %d disk %d lost metadata %q", pool, disk, key)
|
||||
}
|
||||
}
|
||||
}
|
||||
gr, err := z.serverPools[pool].GetObjectNInfo(t.Context(), original.Bucket, original.Name, nil, nil, ObjectOptions{VersionID: version})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
data, err := io.ReadAll(gr)
|
||||
gr.Close()
|
||||
if err != nil || !bytes.Equal(data, []byte(body)) {
|
||||
t.Fatalf("pool %d content changed: bytes=%d err=%v", pool, len(data), err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPoolsMigrationPreservesTagState(t *testing.T) {
|
||||
z, bucket := migrationTestPools(t)
|
||||
for _, kind := range []string{"rebalance", "decommission"} {
|
||||
for _, method := range []string{"put", "multipart"} {
|
||||
for source := range 2 {
|
||||
for _, state := range migrationTagStates {
|
||||
t.Run(fmt.Sprintf("%s/%s/source=%d/%s", kind, method, source, state.name), func(t *testing.T) {
|
||||
move := migrationTestMover(t, z, source, kind)
|
||||
meta := map[string]string{xhttp.AmzObjectTagging: state.tags, "x-amz-meta-owner": "retained"}
|
||||
if state.hasRevision {
|
||||
meta[migrationTagRevision] = state.revision
|
||||
}
|
||||
oi := migrationTestObject(t, z, bucket, t.Name(), source, method == "multipart", []string{"payload"}, ObjectOptions{Versioned: true, UserDefined: meta})
|
||||
migrationTestStored(t, z, source, oi, state.tags, state.revision, state.hasRevision, "payload")
|
||||
gr := migrationTestReader(t, z, source, oi)
|
||||
before := maps.Clone(gr.ObjInfo.UserDefined)
|
||||
if err := move(t.Context(), source, bucket, gr); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !maps.Equal(before, gr.ObjInfo.UserDefined) {
|
||||
t.Error("migration modified its source snapshot")
|
||||
}
|
||||
migrationTestStored(t, z, 1-source, oi, state.tags, state.revision, state.hasRevision, "payload")
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type migrationTestGate struct {
|
||||
entered, resume chan struct{}
|
||||
once, release sync.Once
|
||||
}
|
||||
|
||||
func newMigrationTestGate() *migrationTestGate {
|
||||
return &migrationTestGate{entered: make(chan struct{}), resume: make(chan struct{})}
|
||||
}
|
||||
|
||||
func (g *migrationTestGate) wait(ctx context.Context) error {
|
||||
g.once.Do(func() { close(g.entered) })
|
||||
select {
|
||||
case <-g.resume:
|
||||
return nil
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
func (g *migrationTestGate) unblock() { g.release.Do(func() { close(g.resume) }) }
|
||||
|
||||
type migrationTestGateReader struct {
|
||||
io.Reader
|
||||
ctx context.Context
|
||||
gate *migrationTestGate
|
||||
}
|
||||
|
||||
func (r migrationTestGateReader) Read(p []byte) (int, error) {
|
||||
if err := r.gate.wait(r.ctx); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return r.Reader.Read(p)
|
||||
}
|
||||
|
||||
func TestPoolsMigrationRechecksTags(t *testing.T) {
|
||||
z, bucket := migrationTestPools(t)
|
||||
const recent = "2026-09-16T10:00:00Z"
|
||||
for _, kind := range []string{"rebalance", "decommission"} {
|
||||
for _, method := range []string{"put", "multipart"} {
|
||||
for source := range 2 {
|
||||
for _, tags := range []string{"state=after", ""} {
|
||||
t.Run(fmt.Sprintf("%s/%s/source=%d/clear=%t", kind, method, source, tags == ""), func(t *testing.T) {
|
||||
move := migrationTestMover(t, z, source, kind)
|
||||
oi := migrationTestObject(t, z, bucket, t.Name(), source, method == "multipart", []string{"payload"}, ObjectOptions{
|
||||
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "state=before"},
|
||||
})
|
||||
gr := migrationTestReader(t, z, source, oi)
|
||||
before := maps.Clone(gr.ObjInfo.UserDefined)
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
|
||||
defer cancel()
|
||||
gate := newMigrationTestGate()
|
||||
done, finished := make(chan error, 1), make(chan struct{})
|
||||
if method == "multipart" {
|
||||
// The first part read is after persisted upload initialization,
|
||||
// but before completion takes the object lock and reconciles.
|
||||
gr.Reader = migrationTestGateReader{Reader: gr.Reader, ctx: ctx, gate: gate}
|
||||
go func() { defer close(finished); done <- move(ctx, source, bucket, gr) }()
|
||||
defer func() { gate.unblock(); cancel(); <-finished }()
|
||||
select {
|
||||
case <-gate.entered:
|
||||
case err := <-done:
|
||||
t.Fatalf("migration finished before the part barrier: %v", err)
|
||||
case <-ctx.Done():
|
||||
t.Fatal(ctx.Err())
|
||||
}
|
||||
uploads, err := z.serverPools[1-source].listMultipartUploadsExact(ctx, bucket, oi.Name)
|
||||
if err != nil || len(uploads.Uploads) != 1 {
|
||||
t.Fatalf("upload not persisted at barrier: %+v, %v", uploads, err)
|
||||
}
|
||||
info, err := z.serverPools[1-source].GetMultipartInfo(ctx, bucket, oi.Name, uploads.Uploads[0].UploadID, ObjectOptions{})
|
||||
if err != nil || info.UserDefined[xhttp.AmzObjectTagging] != "state=before" {
|
||||
t.Fatalf("upload lost snapshot tags: %v, %v", info.UserDefined, err)
|
||||
}
|
||||
}
|
||||
// For ordinary PUT this must finish before calling the mover:
|
||||
// blocking its data Read would already hold the object lock.
|
||||
opts := ObjectOptions{VersionID: migrationTestVersion(oi), UserDefined: map[string]string{migrationTagRevision: recent}}
|
||||
var err error
|
||||
if tags == "" {
|
||||
_, err = z.DeleteObjectTags(ctx, bucket, oi.Name, opts)
|
||||
} else {
|
||||
_, err = z.PutObjectTags(ctx, bucket, oi.Name, tags, opts)
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if gr.ObjInfo.UserTags != "state=before" || !maps.Equal(before, gr.ObjInfo.UserDefined) {
|
||||
t.Fatal("test must retain the old source snapshot")
|
||||
}
|
||||
migrationTestStored(t, z, source, oi, tags, recent, true, "payload")
|
||||
if method == "multipart" {
|
||||
gate.unblock()
|
||||
err = <-done
|
||||
} else {
|
||||
err = move(ctx, source, bucket, gr)
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !maps.Equal(before, gr.ObjInfo.UserDefined) {
|
||||
t.Error("reconciliation changed the reader's map")
|
||||
}
|
||||
migrationTestStored(t, z, 1-source, oi, tags, recent, true, "payload")
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type migrationTestMetadataGateDisk struct {
|
||||
StorageAPI
|
||||
bucket, object string
|
||||
gate *migrationTestGate
|
||||
}
|
||||
|
||||
func (d migrationTestMetadataGateDisk) UpdateMetadata(ctx context.Context, volume, path string, fi FileInfo, opts UpdateMetadataOpts) error {
|
||||
if volume == d.bucket && path == d.object {
|
||||
if err := d.gate.wait(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return d.StorageAPI.UpdateMetadata(ctx, volume, path, fi, opts)
|
||||
}
|
||||
|
||||
func TestPoolsMigrationTagsDuringCleanup(t *testing.T) {
|
||||
z, bucket := migrationTestPools(t)
|
||||
const recent = "2026-09-16T10:00:00Z"
|
||||
for _, kind := range []string{"rebalance", "decommission"} {
|
||||
for source := range 2 {
|
||||
for _, tags := range []string{"state=after", ""} {
|
||||
for _, interrupt := range []bool{false, true} {
|
||||
t.Run(fmt.Sprintf("%s/source=%d/clear=%t/interrupt=%t", kind, source, tags == "", interrupt), func(t *testing.T) {
|
||||
move := migrationTestMover(t, z, source, kind)
|
||||
oi := migrationTestObject(t, z, bucket, t.Name(), source, false, []string{"payload"}, ObjectOptions{
|
||||
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "state=before"},
|
||||
})
|
||||
if err := move(t.Context(), source, bucket, migrationTestReader(t, z, source, oi)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
migrationTestStored(t, z, 1-source, oi, "state=before", "", false, "payload")
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
|
||||
defer cancel()
|
||||
mutate := func() (ObjectInfo, error) {
|
||||
opts := ObjectOptions{VersionID: migrationTestVersion(oi), UserDefined: map[string]string{migrationTagRevision: recent}}
|
||||
if tags == "" {
|
||||
return z.DeleteObjectTags(ctx, bucket, oi.Name, opts)
|
||||
}
|
||||
return z.PutObjectTags(ctx, bucket, oi.Name, tags, opts)
|
||||
}
|
||||
set := z.serverPools[source].getHashedSet(oi.Name)
|
||||
cleanup := func() {
|
||||
// Use the same source prefix cleanup as both outer movers.
|
||||
if _, err := set.DeleteObject(ctx, bucket, oi.Name, ObjectOptions{DeletePrefix: true, DeletePrefixObject: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
var updated ObjectInfo
|
||||
if !interrupt {
|
||||
var err error
|
||||
updated, err = mutate()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
cleanup()
|
||||
} else {
|
||||
gate := newMigrationTestGate()
|
||||
original := set.getDisks
|
||||
disks := append([]StorageAPI(nil), original()...)
|
||||
for i := range disks {
|
||||
disks[i] = migrationTestMetadataGateDisk{StorageAPI: disks[i], bucket: bucket, object: oi.Name, gate: gate}
|
||||
}
|
||||
set.getDisks = func() []StorageAPI { return disks }
|
||||
defer func() { set.getDisks = original }()
|
||||
done, finished := make(chan error, 1), make(chan struct{})
|
||||
go func() { defer close(finished); _, err := mutate(); done <- err }()
|
||||
defer func() { gate.unblock(); cancel(); <-finished }()
|
||||
select {
|
||||
case <-gate.entered:
|
||||
case err := <-done:
|
||||
t.Fatalf("metadata update missed the barrier: %v", err)
|
||||
case <-ctx.Done():
|
||||
t.Fatal(ctx.Err())
|
||||
}
|
||||
cleanup()
|
||||
gate.unblock()
|
||||
if err := <-done; err == nil {
|
||||
t.Fatal("update reported success after its source copy was removed")
|
||||
}
|
||||
<-finished
|
||||
set.getDisks = original
|
||||
var err error
|
||||
updated, err = mutate()
|
||||
if err != nil {
|
||||
t.Fatalf("metadata retry did not converge: %v", err)
|
||||
}
|
||||
}
|
||||
if _, err := z.serverPools[source].GetObjectInfo(ctx, bucket, oi.Name, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
t.Fatalf("source cleanup left a copy: %v", err)
|
||||
}
|
||||
migrationTestStored(t, z, 1-source, oi, tags, updated.UserDefined[migrationTagRevision], true, "payload")
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestPoolsMigrationVersionHistory(t *testing.T) {
|
||||
z, bucket := migrationTestPools(t)
|
||||
for _, kind := range []string{"rebalance", "decommission"} {
|
||||
for _, method := range []string{"put", "multipart"} {
|
||||
for _, versioning := range []string{"unversioned", "null", "history"} {
|
||||
t.Run(kind+"/"+method+"/"+versioning, func(t *testing.T) {
|
||||
move := migrationTestMover(t, z, 0, kind)
|
||||
opts := ObjectOptions{
|
||||
Versioned: versioning == "history", VersionSuspended: versioning == "null",
|
||||
MTime: time.Date(2026, 9, 16, 8, 0, 0, 0, time.UTC),
|
||||
UserDefined: map[string]string{xhttp.AmzObjectTagging: "version=old"},
|
||||
}
|
||||
versions := []ObjectInfo{migrationTestObject(t, z, bucket, t.Name(), 0, method == "multipart", []string{"old"}, opts)}
|
||||
if versioning == "history" {
|
||||
opts.MTime = opts.MTime.Add(time.Minute)
|
||||
versions = append(versions, migrationTestObject(t, z, bucket, t.Name(), 0, method == "multipart", []string{"old"}, opts))
|
||||
}
|
||||
var newest ObjectInfo
|
||||
if versioning != "unversioned" {
|
||||
newest = migrationTestObject(t, z, bucket, t.Name(), 0, false, []string{"latest"}, ObjectOptions{
|
||||
Versioned: true, MTime: opts.MTime.Add(time.Minute),
|
||||
UserDefined: map[string]string{xhttp.AmzObjectTagging: "version=new"},
|
||||
})
|
||||
}
|
||||
for _, oi := range versions {
|
||||
if err := move(t.Context(), 0, bucket, migrationTestReader(t, z, 0, oi)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
migrationTestStored(t, z, 1, oi, "version=old", "", false, "old")
|
||||
}
|
||||
if versioning != "unversioned" {
|
||||
migrationTestStored(t, z, 0, newest, "version=new", "", false, "latest")
|
||||
if _, err := z.serverPools[1].GetObjectInfo(t.Context(), bucket, t.Name(), ObjectOptions{VersionID: newest.VersionID}); !isErrVersionNotFound(err) {
|
||||
t.Fatalf("moving an addressed version affected the newer version: %v", err)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestPoolsMigrationMultipartParts(t *testing.T) {
|
||||
z, bucket := migrationTestPools(t)
|
||||
parts := []string{strings.Repeat("a", 5<<20), "tail-with-tags"}
|
||||
body := strings.Join(parts, "")
|
||||
for _, kind := range []string{"rebalance", "decommission"} {
|
||||
for source := range 2 {
|
||||
t.Run(fmt.Sprintf("%s/source=%d", kind, source), func(t *testing.T) {
|
||||
move := migrationTestMover(t, z, source, kind)
|
||||
oi := migrationTestObject(t, z, bucket, t.Name(), source, true, parts, ObjectOptions{
|
||||
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "multipart=kept"},
|
||||
})
|
||||
if err := move(t.Context(), source, bucket, migrationTestReader(t, z, source, oi)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
migrationTestStored(t, z, 1-source, oi, "multipart=kept", "", false, body)
|
||||
got, err := z.serverPools[1-source].GetObjectInfo(t.Context(), bucket, oi.Name, ObjectOptions{VersionID: oi.VersionID})
|
||||
if err != nil || len(got.Parts) != len(oi.Parts) {
|
||||
t.Fatalf("parts changed: %+v, %v", got.Parts, err)
|
||||
}
|
||||
for i, part := range got.Parts {
|
||||
old := oi.Parts[i]
|
||||
if part.Number != old.Number || part.Size != old.Size || part.ActualSize != old.ActualSize || part.ETag != old.ETag || !bytes.Equal(part.Index, old.Index) {
|
||||
t.Errorf("part %d changed: %+v / %+v", i, part, old)
|
||||
}
|
||||
}
|
||||
start := int64(len(parts[0]) - 3)
|
||||
gr, err := z.serverPools[1-source].GetObjectNInfo(t.Context(), bucket, oi.Name, &HTTPRangeSpec{Start: start, End: start + 7}, nil, ObjectOptions{VersionID: oi.VersionID})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
data, err := io.ReadAll(gr)
|
||||
gr.Close()
|
||||
if err != nil || string(data) != body[start:start+8] {
|
||||
t.Fatalf("cross-part range changed: %q, %v", data, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestPoolsMigrationUnreadablePool(t *testing.T) {
|
||||
z, bucket := migrationTestPools(t)
|
||||
for _, kind := range []string{"rebalance", "decommission"} {
|
||||
for _, method := range []string{"put", "multipart"} {
|
||||
t.Run(kind+"/"+method, func(t *testing.T) {
|
||||
move := migrationTestMover(t, z, 0, kind)
|
||||
oi := migrationTestObject(t, z, bucket, t.Name(), 0, method == "multipart", []string{"payload"}, ObjectOptions{
|
||||
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "keep=source"},
|
||||
})
|
||||
gr := migrationTestReader(t, z, 0, oi)
|
||||
set := z.serverPools[0].getHashedSet(oi.Name)
|
||||
original := set.getDisks
|
||||
disks := append([]StorageAPI(nil), original()...)
|
||||
for i := range disks {
|
||||
disks[i] = consistencyReadFaultDisk{StorageAPI: disks[i], bucket: bucket, object: oi.Name}
|
||||
}
|
||||
set.getDisks = func() []StorageAPI { return disks }
|
||||
defer func() { set.getDisks = original }()
|
||||
err := move(t.Context(), 0, bucket, gr)
|
||||
var quorum InsufficientReadQuorum
|
||||
if !errors.As(err, &quorum) {
|
||||
t.Fatalf("unreadable pool must fail migration with read quorum error: %v", err)
|
||||
}
|
||||
set.getDisks = original
|
||||
migrationTestStored(t, z, 0, oi, "keep=source", "", false, "payload")
|
||||
if _, err := z.serverPools[1].GetObjectInfo(t.Context(), bucket, oi.Name, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
t.Fatalf("failed migration committed a target version: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -865,7 +865,7 @@ func (z *erasureServerPools) rebalanceObject(ctx context.Context, poolIdx int, b
|
||||
if oi.isMultipart() {
|
||||
res, err := z.NewMultipartUpload(ctx, bucket, oi.Name, ObjectOptions{
|
||||
VersionID: oi.VersionID,
|
||||
UserDefined: oi.UserDefined,
|
||||
UserDefined: migrationObjectMetadata(oi),
|
||||
NoAuditLog: true,
|
||||
DataMovement: true,
|
||||
SrcPoolIdx: poolIdx,
|
||||
@@ -924,7 +924,7 @@ func (z *erasureServerPools) rebalanceObject(ctx context.Context, poolIdx int, b
|
||||
DataMovement: true,
|
||||
VersionID: oi.VersionID,
|
||||
MTime: oi.ModTime,
|
||||
UserDefined: oi.UserDefined,
|
||||
UserDefined: migrationObjectMetadata(oi),
|
||||
PreserveETag: oi.ETag, // Preserve original ETag to ensure same metadata.
|
||||
IndexCB: func() []byte {
|
||||
return oi.Parts[0].Index // Preserve part Index to ensure decompression works.
|
||||
|
||||
@@ -1239,9 +1239,12 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob
|
||||
return ObjectInfo{}, z.deletePrefix(ctx, bucket, object)
|
||||
}
|
||||
|
||||
// Reconcile ordinary addressed-version deletes independently of pool movement.
|
||||
// Resolve a physical purge by its addressed version even on a replica
|
||||
// receiver: latest-key routing can select a different pool or miss copies.
|
||||
// Marker creation and specialized movement/scanner operations retain their
|
||||
// existing routing.
|
||||
reconcileVersion := opts.VersionID != "" && !opts.DataMovement &&
|
||||
!opts.ReplicationRequest && !opts.Expiration.Expire && !opts.InclFreeVersions
|
||||
(!opts.ReplicationRequest || opts.isVersionPurge()) && !opts.Expiration.Expire && !opts.InclFreeVersions
|
||||
if !z.SinglePool() && (opts.CheckPrecondFn != nil || reconcileVersion) {
|
||||
return z.deleteObjectReconciled(ctx, bucket, object, opts)
|
||||
}
|
||||
@@ -1254,6 +1257,13 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob
|
||||
if _, ok := err.(InsufficientReadQuorum); ok {
|
||||
return objInfo, InsufficientWriteQuorum{}
|
||||
}
|
||||
// Lookup can return before the set's purge confirmation. Check here,
|
||||
// before any callback can change the request into a metadata update.
|
||||
if opts.isVersionPurge() && (isErrObjectNotFound(err) || isErrVersionNotFound(err)) {
|
||||
if quorumErr := z.checkPurgeAbsent(ctx, bucket, object, opts.VersionID); quorumErr != nil {
|
||||
return objInfo, quorumErr
|
||||
}
|
||||
}
|
||||
// A conditional (If-Match) delete addressing a specific version treats an
|
||||
// absent key as an absent version. getPoolInfoExistingWithOpts strips
|
||||
// VersionID, so a missing key surfaces ObjectNotFound here even for a
|
||||
@@ -1284,6 +1294,11 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob
|
||||
if verr != nil && (!isErrMethodNotAllowed(verr) || !vi.DeleteMarker) {
|
||||
// Genuine read failure for the addressed version: a missing
|
||||
// version -> VersionNotFound (NoSuchVersion), read-quorum loss, etc.
|
||||
if opts.isVersionPurge() && (isErrObjectNotFound(verr) || isErrVersionNotFound(verr)) {
|
||||
if quorumErr := z.checkPurgeAbsent(ctx, bucket, object, opts.VersionID); quorumErr != nil {
|
||||
return objInfo, quorumErr
|
||||
}
|
||||
}
|
||||
return objInfo, verr
|
||||
}
|
||||
// verr is nil for a live version, or MethodNotAllowed with a populated
|
||||
@@ -2160,13 +2175,14 @@ func (z *erasureServerPools) AbortMultipartUpload(ctx context.Context, bucket, o
|
||||
return toObjectErr(err, bucket)
|
||||
}
|
||||
|
||||
defer func() {
|
||||
_, absent := err.(InvalidUploadID)
|
||||
if err == nil || absent {
|
||||
// Unlock cancels the derived lock context before this notification runs.
|
||||
// Keep the request context so successful cancellation reaches peer caches.
|
||||
defer func(ctx context.Context) {
|
||||
if err == nil {
|
||||
z.mpCache.Delete(uploadID)
|
||||
globalNotificationSys.DeleteUploadID(ctx, uploadID)
|
||||
}
|
||||
}()
|
||||
}(ctx)
|
||||
|
||||
lk := z.NewNSLock(bucket, pathJoin(object, uploadID))
|
||||
lkctx, err := lk.GetLock(ctx, globalOperationTimeout)
|
||||
@@ -2176,13 +2192,18 @@ func (z *erasureServerPools) AbortMultipartUpload(ctx context.Context, bucket, o
|
||||
ctx = lkctx.Context()
|
||||
defer lk.Unlock(lkctx)
|
||||
|
||||
legacy := globalAPIConfig.getMultipartListingLegacy()
|
||||
found := false
|
||||
var firstErr error
|
||||
for idx, pool := range z.serverPools {
|
||||
if z.IsSuspended(idx) {
|
||||
continue
|
||||
}
|
||||
poolFound, err := pool.getHashedSet(object).abortMultipartUpload(ctx, bucket, object, uploadID, opts)
|
||||
poolFound, err := pool.getHashedSet(object).abortMultipartUpload(ctx, bucket, object, uploadID, opts, legacy)
|
||||
if legacy && (poolFound || err != nil) {
|
||||
// Match the released first-matching-pool behavior.
|
||||
return err
|
||||
}
|
||||
found = found || poolFound
|
||||
if err != nil && firstErr == nil {
|
||||
firstErr = err
|
||||
|
||||
@@ -0,0 +1,87 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/minio/minio/internal/bucket/replication"
|
||||
)
|
||||
|
||||
// isVersionPurge distinguishes physical removal from marker creation and
|
||||
// replication-state updates. Call it again after metadata callbacks, which
|
||||
// can change an ordinary deletion into a pending purge update.
|
||||
func (o ObjectOptions) isVersionPurge() bool {
|
||||
if o.VersionID == "" || o.VersionID == nullVersionID || o.DeleteMarker ||
|
||||
o.DeletePrefix || o.DataMovement || o.InclFreeVersions || o.Expiration.Expire ||
|
||||
o.Transition != (TransitionOptions{}) {
|
||||
return false
|
||||
}
|
||||
id, err := uuid.Parse(o.VersionID)
|
||||
if err != nil || id == uuid.Nil {
|
||||
return false
|
||||
}
|
||||
purge := o.VersionPurgeStatus()
|
||||
if purge == replication.VersionPurgeComplete {
|
||||
return true
|
||||
}
|
||||
if !purge.Empty() || o.DeleteReplication.VersionPurgeStatusInternal != "" {
|
||||
return false
|
||||
}
|
||||
status := o.DeleteMarkerReplicationStatus()
|
||||
return (status.Empty() && o.DeleteReplication.ReplicationStatusInternal == "") ||
|
||||
(status == replication.Replica && o.ReplicationRequest)
|
||||
}
|
||||
|
||||
// confirmVersionAbsent is a read-only proof for an already-missing purge.
|
||||
// Read quorum (or a synthesized NotFound) cannot acknowledge a write. Do not
|
||||
// delete here: unreadable minority copies may carry retention we cannot check.
|
||||
// The caller holds the object lock. This function never heals or enqueues work.
|
||||
func (er erasureObjects) confirmVersionAbsent(ctx context.Context, bucket, object, versionID string) (allAbsent bool, err error) {
|
||||
disks := er.getDisks()
|
||||
_, errs := readAllFileInfo(ctx, disks, "", bucket, object, versionID, false, false)
|
||||
absent := 0
|
||||
for _, err := range errs {
|
||||
if err == errFileNotFound || err == errFileVersionNotFound {
|
||||
absent++
|
||||
}
|
||||
}
|
||||
if absent < len(disks)/2+1 {
|
||||
return false, InsufficientWriteQuorum{}
|
||||
}
|
||||
return absent == len(disks), nil
|
||||
}
|
||||
|
||||
// checkPurgeAbsent schedules recovery explicitly, outside the read-only proof.
|
||||
func (er erasureObjects) checkPurgeAbsent(ctx context.Context, bucket, object, versionID string) error {
|
||||
allAbsent, err := er.confirmVersionAbsent(ctx, bucket, object, versionID)
|
||||
if !allAbsent {
|
||||
er.addPartial(bucket, object, versionID)
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (z *erasureServerPools) checkPurgeAbsent(ctx context.Context, bucket, object, versionID string) error {
|
||||
for _, pool := range z.serverPools {
|
||||
if err := pool.getHashedSet(object).checkPurgeAbsent(ctx, bucket, object, versionID); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
+3
-3
@@ -50,7 +50,7 @@ type apiConfig struct {
|
||||
transitionWorkers int
|
||||
|
||||
staleUploadsExpiry time.Duration
|
||||
multipartListingLegacy bool
|
||||
multipartListingStrict bool
|
||||
staleUploadsCleanupInterval time.Duration
|
||||
deleteCleanupInterval time.Duration
|
||||
enableODirect bool
|
||||
@@ -182,7 +182,7 @@ func (t *apiConfig) init(cfg api.Config, setDriveCounts []int, legacy bool) {
|
||||
t.transitionWorkers = cfg.TransitionWorkers
|
||||
|
||||
t.staleUploadsExpiry = cfg.StaleUploadsExpiry
|
||||
t.multipartListingLegacy = cfg.MultipartListing == "legacy"
|
||||
t.multipartListingStrict = cfg.MultipartListing == "strict"
|
||||
t.deleteCleanupInterval = cfg.DeleteCleanupInterval
|
||||
t.enableODirect = cfg.EnableODirect
|
||||
t.gzipObjects = cfg.GzipObjects
|
||||
@@ -211,7 +211,7 @@ func (t *apiConfig) odirectEnabled() bool {
|
||||
func (t *apiConfig) getMultipartListingLegacy() bool {
|
||||
t.mu.RLock()
|
||||
defer t.mu.RUnlock()
|
||||
return t.multipartListingLegacy
|
||||
return !t.multipartListingStrict
|
||||
}
|
||||
|
||||
func (t *apiConfig) shouldGzipObjects() bool {
|
||||
|
||||
@@ -47,6 +47,7 @@ func TestListMultipartUploadsS3Compatibility(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
z := obj.(*erasureServerPools)
|
||||
setMultipartListingTestMode(t, false)
|
||||
t.Cleanup(func() {
|
||||
z.Shutdown(t.Context())
|
||||
removeRoots(dirs)
|
||||
@@ -201,15 +202,7 @@ func TestListMultipartUploadsS3Compatibility(t *testing.T) {
|
||||
if !errors.Is(err, errMultipartListingLegacy) {
|
||||
t.Fatalf("legacy strict listing: %v", err)
|
||||
}
|
||||
globalAPIConfig.mu.Lock()
|
||||
oldLegacy := globalAPIConfig.multipartListingLegacy
|
||||
globalAPIConfig.multipartListingLegacy = true
|
||||
globalAPIConfig.mu.Unlock()
|
||||
t.Cleanup(func() {
|
||||
globalAPIConfig.mu.Lock()
|
||||
globalAPIConfig.multipartListingLegacy = oldLegacy
|
||||
globalAPIConfig.mu.Unlock()
|
||||
})
|
||||
setMultipartListingTestMode(t, true)
|
||||
legacy, err := z.ListMultipartUploads(t.Context(), bucket, legacyObject, "", "", "", 100)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -267,6 +260,7 @@ func TestPaginateMultipartUploads(t *testing.T) {
|
||||
|
||||
func TestListMultipartUploadsGlobalPageAcrossPools(t *testing.T) {
|
||||
z, bucket := consistencyPools(t)
|
||||
setMultipartListingTestMode(t, false)
|
||||
objects := []string{"a/one", "b/two", "c/three", "d/four"}
|
||||
for i, object := range objects {
|
||||
if _, err := z.serverPools[i%len(z.serverPools)].NewMultipartUpload(t.Context(), bucket, object, ObjectOptions{}); err != nil {
|
||||
|
||||
@@ -0,0 +1,433 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/xml"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"slices"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
madmin "github.com/minio/madmin-go/v3"
|
||||
)
|
||||
|
||||
// Hold nine completed disk walks while a PUT overwrites one name. The other
|
||||
// seven walks start after that PUT completes. Both generations come from real
|
||||
// disk walkers and valid PUT metadata; only the reader schedule is controlled.
|
||||
type nullQuorumWalkSchedule struct {
|
||||
bucket string
|
||||
oldReady chan struct{}
|
||||
resume chan struct{}
|
||||
mu sync.Mutex
|
||||
snapshots [16][]byte
|
||||
}
|
||||
|
||||
type nullQuorumWalkDisk struct {
|
||||
StorageAPI
|
||||
schedule *nullQuorumWalkSchedule
|
||||
index int
|
||||
}
|
||||
|
||||
func (d *nullQuorumWalkDisk) WalkDir(ctx context.Context, opts WalkDirOptions, out io.Writer) error {
|
||||
if opts.Bucket != d.schedule.bucket {
|
||||
return d.StorageAPI.WalkDir(ctx, opts, out)
|
||||
}
|
||||
var stream bytes.Buffer
|
||||
if d.index < 9 {
|
||||
if err := d.StorageAPI.WalkDir(ctx, opts, &stream); err != nil {
|
||||
return err
|
||||
}
|
||||
select {
|
||||
case d.schedule.oldReady <- struct{}{}:
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
}
|
||||
select {
|
||||
case <-d.schedule.resume:
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
if d.index >= 9 {
|
||||
if err := d.StorageAPI.WalkDir(ctx, opts, &stream); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
d.schedule.mu.Lock()
|
||||
d.schedule.snapshots[d.index] = bytes.Clone(stream.Bytes())
|
||||
d.schedule.mu.Unlock()
|
||||
_, err := out.Write(stream.Bytes())
|
||||
return err
|
||||
}
|
||||
|
||||
func nullQuorumBackend(t *testing.T) (*erasureServerPools, string, http.Handler) {
|
||||
t.Helper()
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
obj, dirs, err := prepareErasure16(ctx)
|
||||
if err != nil {
|
||||
cancel()
|
||||
t.Fatal(err)
|
||||
}
|
||||
z := obj.(*erasureServerPools)
|
||||
previous := newObjectLayerFn()
|
||||
setObjectLayer(z)
|
||||
t.Cleanup(func() {
|
||||
cancel()
|
||||
z.Shutdown(context.Background())
|
||||
removeRoots(dirs)
|
||||
setObjectLayer(previous)
|
||||
})
|
||||
bucket, router, err := initAPIHandlerTest(ctx, z, nil, MakeBucketOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return z, bucket, router
|
||||
}
|
||||
|
||||
func nullQuorumRequest(t *testing.T, router http.Handler, method, target string) *httptest.ResponseRecorder {
|
||||
t.Helper()
|
||||
req, err := newTestSignedRequestV4(method, target, 0, nil, globalActiveCred.AccessKey, globalActiveCred.SecretKey, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rec := httptest.NewRecorder()
|
||||
router.ServeHTTP(rec, req.WithContext(t.Context()))
|
||||
return rec
|
||||
}
|
||||
|
||||
func TestListObjectsSingleNullQuorumHTTP(t *testing.T) {
|
||||
z, bucket, router := nullQuorumBackend(t)
|
||||
oldBody := bytes.Repeat([]byte{'a'}, 8192)
|
||||
newBody := bytes.Repeat([]byte{'b'}, 8192)
|
||||
oldTime := time.Now().UTC().Add(-time.Hour)
|
||||
newTime := oldTime.Add(time.Minute)
|
||||
var want []string
|
||||
for i := range 16 {
|
||||
name := fmt.Sprintf("fixed/%02d", i)
|
||||
want = append(want, name)
|
||||
_, err := z.PutObject(t.Context(), bucket, name, mustGetPutObjReader(t, bytes.NewReader(oldBody), int64(len(oldBody)), "", ""), ObjectOptions{MTime: oldTime})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
const overwritten = "fixed/10"
|
||||
schedule := &nullQuorumWalkSchedule{bucket: bucket, oldReady: make(chan struct{}, 16), resume: make(chan struct{})}
|
||||
set := z.serverPools[0].sets[0]
|
||||
disks := set.getDisks()
|
||||
wrapped := make([]StorageAPI, len(disks))
|
||||
for i, disk := range disks {
|
||||
wrapped[i] = &nullQuorumWalkDisk{StorageAPI: disk, schedule: schedule, index: i}
|
||||
}
|
||||
set.getDisks = func() []StorageAPI { return wrapped }
|
||||
if globalAPIConfig.getListQuorum() != "strict" {
|
||||
t.Fatal("test requires strict listing across all sixteen disks")
|
||||
}
|
||||
writeDone := make(chan error, 1)
|
||||
putReader := mustGetPutObjReader(t, bytes.NewReader(newBody), int64(len(newBody)), "", "")
|
||||
go func() {
|
||||
defer close(schedule.resume)
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
|
||||
defer cancel()
|
||||
for range 9 {
|
||||
select {
|
||||
case <-schedule.oldReady:
|
||||
case <-ctx.Done():
|
||||
writeDone <- ctx.Err()
|
||||
return
|
||||
}
|
||||
}
|
||||
_, err := z.PutObject(ctx, bucket, overwritten, putReader, ObjectOptions{MTime: newTime})
|
||||
writeDone <- err
|
||||
}()
|
||||
rec := nullQuorumRequest(t, router, http.MethodGet, getListObjectsV2URL("", bucket, "fixed/", "1000", "", "", ""))
|
||||
if err := <-writeDone; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var list ListObjectsV2Response
|
||||
if rec.Code != http.StatusOK || xml.Unmarshal(rec.Body.Bytes(), &list) != nil {
|
||||
t.Fatalf("LIST: %d %s", rec.Code, rec.Body.String())
|
||||
}
|
||||
var got []string
|
||||
for _, object := range list.Contents {
|
||||
got = append(got, object.Key)
|
||||
}
|
||||
t.Logf("LIST HTTP=%d KeyCount=%d IsTruncated=%v keys=%v", rec.Code, list.KeyCount, list.IsTruncated, got)
|
||||
if list.KeyCount != 16 || list.IsTruncated || !slices.Equal(got, want) {
|
||||
t.Errorf("LIST omitted a name during overwrite: got %v, want %v", got, want)
|
||||
}
|
||||
|
||||
// Verify the actual reader inputs, including shape and EC, instead of
|
||||
// assuming the scheduling barrier produced the intended resolver case.
|
||||
schedule.mu.Lock()
|
||||
snapshots := schedule.snapshots
|
||||
schedule.mu.Unlock()
|
||||
for i, snapshot := range snapshots {
|
||||
reader := newMetacacheReader(bytes.NewReader(snapshot))
|
||||
var found bool
|
||||
for {
|
||||
entry, err := reader.next()
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if entry.name != overwritten {
|
||||
continue
|
||||
}
|
||||
xl, err := entry.xlmeta()
|
||||
if err != nil || len(xl.versions) != 1 {
|
||||
t.Fatalf("disk %d: invalid input shape: %v", i, err)
|
||||
}
|
||||
header := xl.versions[0].header
|
||||
wantTime := oldTime
|
||||
if i >= 9 {
|
||||
wantTime = newTime
|
||||
}
|
||||
if header.ModTime != wantTime.UnixNano() || header.VersionID != [16]byte{} || header.Type != ObjectType || header.FreeVersion() {
|
||||
t.Fatalf("disk %d: unexpected header %v", i, header)
|
||||
}
|
||||
t.Logf("disk=%d versions=1 header=%v", i, header)
|
||||
found = true
|
||||
}
|
||||
reader.Close()
|
||||
if !found {
|
||||
t.Fatalf("disk %d did not emit the overwritten name", i)
|
||||
}
|
||||
}
|
||||
get := nullQuorumRequest(t, router, http.MethodGet, getGetObjectURL("", bucket, overwritten))
|
||||
head := nullQuorumRequest(t, router, http.MethodHead, getGetObjectURL("", bucket, overwritten))
|
||||
t.Logf("completed PUT readback: GET=%d bytes=%d HEAD=%d length=%s", get.Code, get.Body.Len(), head.Code, head.Header().Get("Content-Length"))
|
||||
if get.Code != http.StatusOK || !bytes.Equal(get.Body.Bytes(), newBody) || head.Code != http.StatusOK || head.Header().Get("Content-Length") != "8192" {
|
||||
t.Fatalf("completed PUT was not readable: GET=%d HEAD=%d", get.Code, head.Code)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSingleNullQuorumReadAndConsumers(t *testing.T) {
|
||||
z, bucket, router := nullQuorumBackend(t)
|
||||
set := z.serverPools[0].sets[0]
|
||||
disks := set.getDisks()
|
||||
const object = "object"
|
||||
oldBody, newBody := bytes.Repeat([]byte{'a'}, 8192), bytes.Repeat([]byte{'b'}, 8192)
|
||||
oldTime := time.Now().UTC().Add(-time.Hour)
|
||||
put := func(body []byte, modTime time.Time) {
|
||||
t.Helper()
|
||||
_, err := z.PutObject(t.Context(), bucket, object, mustGetPutObjReader(t, bytes.NewReader(body), int64(len(body)), "", ""), ObjectOptions{MTime: modTime})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
put(oldBody, oldTime)
|
||||
oldMeta := make([][]byte, len(disks))
|
||||
for i, disk := range disks {
|
||||
var err error
|
||||
oldMeta[i], err = disk.ReadAll(t.Context(), bucket, object+"/"+xlStorageFormatFile)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
put(newBody, oldTime.Add(time.Minute))
|
||||
newMinorityMeta := mustReadNullQuorumMeta(t, disks[14], bucket, object)
|
||||
// Model an incomplete overwrite using real inline shards: fifteen old
|
||||
// copies retain a read quorum; the final disk contains the newer minority.
|
||||
for i := range 15 {
|
||||
if err := disks[i].WriteAll(t.Context(), bucket, object+"/"+xlStorageFormatFile, oldMeta[i]); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
checkRead := func(body []byte, modTime time.Time) {
|
||||
t.Helper()
|
||||
get := nullQuorumRequest(t, router, http.MethodGet, getGetObjectURL("", bucket, object))
|
||||
head := nullQuorumRequest(t, router, http.MethodHead, getGetObjectURL("", bucket, object))
|
||||
if get.Code != 200 || !bytes.Equal(get.Body.Bytes(), body) || head.Code != 200 || head.Header().Get("Content-Length") != "8192" {
|
||||
t.Fatalf("GET/HEAD lost readable quorum: GET=%d HEAD=%d body=%q", get.Code, head.Code, get.Body.String())
|
||||
}
|
||||
oi, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{})
|
||||
if err != nil || !oi.ModTime.Equal(modTime) {
|
||||
t.Fatalf("wrong readable generation: %v, %v", oi.ModTime, err)
|
||||
}
|
||||
}
|
||||
checkRead(oldBody, oldTime)
|
||||
// Reading must not reinterpret the minority as absent or delete it.
|
||||
minority, err := disks[15].ReadVersion(t.Context(), "", bucket, object, "", ReadOptions{})
|
||||
if err != nil || !minority.ModTime.Equal(oldTime.Add(time.Minute)) {
|
||||
t.Fatalf("read mutated the minority: %+v, %v", minority, err)
|
||||
}
|
||||
for _, kind := range []string{"rebalance", "decommission"} {
|
||||
t.Run(kind, func(t *testing.T) {
|
||||
var names []string
|
||||
consume := func(entry metaCacheEntry) {
|
||||
versions, err := entry.fileInfoVersions(bucket)
|
||||
if err != nil || len(versions.Versions) != 1 || !versions.Versions[0].ModTime.Equal(oldTime) {
|
||||
t.Errorf("invalid migration metadata: %+v, %v", versions, err)
|
||||
return
|
||||
}
|
||||
// Migration reopens the selected name/version on all drives;
|
||||
// it must not treat the selected listing metadata as disk agreement.
|
||||
reader, err := set.GetObjectNInfo(t.Context(), bucket, entry.name, nil, nil, ObjectOptions{VersionID: nullVersionID, NoLock: true, NoDecryption: true})
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
return
|
||||
}
|
||||
body, err := io.ReadAll(reader)
|
||||
reader.Close()
|
||||
if err != nil || !bytes.Equal(body, oldBody) {
|
||||
t.Errorf("migration could not read selected object data: %v", err)
|
||||
}
|
||||
names = append(names, entry.name)
|
||||
}
|
||||
var err error
|
||||
if kind == "rebalance" {
|
||||
err = set.listObjectsToRebalance(t.Context(), bucket, consume)
|
||||
} else {
|
||||
err = set.listObjectsToDecommission(t.Context(), decomBucketInfo{Name: bucket}, consume)
|
||||
}
|
||||
if err != nil || !slices.Equal(names, []string{object}) {
|
||||
t.Fatalf("migration listing dropped the quorum: %v, %v", names, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
for _, latestOnly := range []bool{false, true} {
|
||||
results := make(chan itemOrErr[ObjectInfo], 16)
|
||||
if err := z.Walk(t.Context(), bucket, "", results, WalkOptions{AskDisks: "strict", LatestOnly: latestOnly}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var names []string
|
||||
for result := range results {
|
||||
if result.Err != nil || !result.Item.ModTime.Equal(oldTime) {
|
||||
t.Fatalf("Walk returned invalid metadata: %+v", result)
|
||||
}
|
||||
names = append(names, result.Item.Name)
|
||||
}
|
||||
if !slices.Equal(names, []string{object}) {
|
||||
t.Fatalf("Walk dropped the quorum: %v", names)
|
||||
}
|
||||
}
|
||||
|
||||
// Exercise the scanner's actual abandoned-child path. Make the scan drive
|
||||
// miss the object, while keeping an older quorum and a newer minority on
|
||||
// the remaining disks. Its queued heal must read all disks and repair both.
|
||||
if err := disks[14].WriteAll(t.Context(), bucket, object+"/"+xlStorageFormatFile, newMinorityMeta); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.RemoveAll(filepath.Join(disks[15].Endpoint().Path, bucket, object)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
scanNullQuorumAbandoned(t, z, bucket, object, oldTime)
|
||||
for i, disk := range disks {
|
||||
fi, err := disk.ReadVersion(t.Context(), "", bucket, object, "", ReadOptions{Healing: true})
|
||||
if err != nil || !fi.ModTime.Equal(oldTime) {
|
||||
t.Fatalf("scanner heal did not reconcile disk %d: %v, %v", i, fi.ModTime, err)
|
||||
}
|
||||
}
|
||||
checkRead(oldBody, oldTime)
|
||||
put(newBody, oldTime.Add(2*time.Minute))
|
||||
checkRead(newBody, oldTime.Add(2*time.Minute))
|
||||
}
|
||||
|
||||
func mustReadNullQuorumMeta(t *testing.T, disk StorageAPI, bucket, object string) []byte {
|
||||
t.Helper()
|
||||
data, err := disk.ReadAll(t.Context(), bucket, object+"/"+xlStorageFormatFile)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return data
|
||||
}
|
||||
|
||||
func scanNullQuorumAbandoned(t *testing.T, z *erasureServerPools, bucket, object string, oldTime time.Time) {
|
||||
t.Helper()
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
defer cancel()
|
||||
disks := z.serverPools[0].sets[0].getDisks()
|
||||
seq := newBgHealSequence()
|
||||
defer seq.cancelCtx()
|
||||
previousState, previousRoutine := globalBackgroundHealState, globalBackgroundHealRoutine
|
||||
globalBackgroundHealState = &allHealState{healSeqMap: map[string]*healSequence{"test": seq}}
|
||||
routine := &healRoutine{tasks: make(chan healTask)}
|
||||
globalBackgroundHealRoutine = routine
|
||||
defer func() { globalBackgroundHealState, globalBackgroundHealRoutine = previousState, previousRoutine }()
|
||||
workerDone := make(chan struct{})
|
||||
var healed []string
|
||||
go func() {
|
||||
defer close(workerDone)
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case task := <-routine.tasks:
|
||||
var result madmin.HealResultItem
|
||||
var err error
|
||||
if task.object == "" {
|
||||
result, err = z.HealBucket(ctx, task.bucket, task.opts)
|
||||
} else {
|
||||
healed = append(healed, task.object+"/"+task.versionID)
|
||||
result, err = z.HealObject(ctx, task.bucket, task.object, task.versionID, task.opts)
|
||||
}
|
||||
task.respCh <- healResult{result: result, err: err}
|
||||
}
|
||||
}
|
||||
}()
|
||||
cache := dataUsageCache{Info: dataUsageCacheInfo{Name: bucket}}
|
||||
cache.replace(bucket, "", dataUsageEntry{})
|
||||
cache.replace(bucket+"/"+object, bucket, dataUsageEntry{Size: 8192, Objects: 1, Versions: 1})
|
||||
scanner := folderScanner{
|
||||
root: disks[15].Endpoint().Path, oldCache: cache,
|
||||
newCache: dataUsageCache{Info: cache.Info}, updateCache: dataUsageCache{Info: cache.Info},
|
||||
disks: disks, disksQuorum: 8, healObjectSelect: 1,
|
||||
weSleep: func() bool { return false }, shouldHeal: func() bool { return true }, updateCurrentPath: func(string) {},
|
||||
getSize: func(item scannerItem) (sizeSummary, error) {
|
||||
if filepath.Base(item.Path) != xlStorageFormatFile {
|
||||
return sizeSummary{}, errSkipFile
|
||||
}
|
||||
data, err := os.ReadFile(item.Path)
|
||||
if err != nil {
|
||||
return sizeSummary{}, err
|
||||
}
|
||||
var xl xlMetaV2
|
||||
if err := xl.Load(data); err != nil {
|
||||
return sizeSummary{}, err
|
||||
}
|
||||
fi, err := xl.ToFileInfo(bucket, object, "", false, false)
|
||||
if err != nil || !fi.ModTime.Equal(oldTime) {
|
||||
return sizeSummary{}, fmt.Errorf("scanner re-read wrong generation: %v, %v", fi.ModTime, err)
|
||||
}
|
||||
return sizeSummary{totalSize: fi.Size, versions: 1}, nil
|
||||
},
|
||||
}
|
||||
var usage dataUsageEntry
|
||||
err := scanner.scanFolder(ctx, cachedFolder{name: bucket, objectHealProbDiv: 1}, &usage)
|
||||
cancel()
|
||||
<-workerDone
|
||||
if err != nil || len(healed) != 1 || healed[0] != object+"/"+nullVersionID && healed[0] != object+"/" {
|
||||
t.Fatalf("scanner did not queue a name-based heal: %v, %v", healed, err)
|
||||
}
|
||||
flat := scanner.newCache.sizeRecursive(bucket)
|
||||
if flat == nil || flat.Objects != 1 || flat.Size != 8192 {
|
||||
t.Fatalf("scanner lost the healed object from usage: %+v", flat)
|
||||
}
|
||||
t.Logf("scanner queued %v; healed old quorum, repaired minority/missing copies, and retained one 8192-byte object", healed)
|
||||
}
|
||||
@@ -69,6 +69,8 @@ func loadCPUMetrics(ctx context.Context, m MetricValues, c *metricsCache) error
|
||||
|
||||
// metrics-resource.go runs a job to collect resource metrics including their Avg values and
|
||||
// stores them in resourceMetricsMap. We can use it to get the Avg values of CPU idle and IOWait.
|
||||
resourceMetricsMapMu.RLock()
|
||||
defer resourceMetricsMapMu.RUnlock()
|
||||
cpuResourceMetrics, found := resourceMetricsMap[cpuSubsystem]
|
||||
if found {
|
||||
if cpuIdleMetric, ok := cpuResourceMetrics[getResourceKey(cpuIdle, nil)]; ok {
|
||||
|
||||
@@ -0,0 +1,185 @@
|
||||
// Copyright (c) 2026 Ruohang Feng
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"runtime"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/minio/madmin-go/v3"
|
||||
"github.com/minio/minio/internal/cachevalue"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
dto "github.com/prometheus/client_model/go"
|
||||
cpustats "github.com/shirou/gopsutil/v3/cpu"
|
||||
"github.com/shirou/gopsutil/v3/load"
|
||||
)
|
||||
|
||||
// These tests replace global resource metrics and must not run in parallel.
|
||||
// Concurrent writers must finish before the fixture is restored.
|
||||
func setCPUResourceMetricsForTest(t *testing.T, value map[MetricSubsystem]ResourceMetrics) {
|
||||
t.Helper()
|
||||
resourceMetricsMapMu.Lock()
|
||||
saved := resourceMetricsMap
|
||||
resourceMetricsMap = value
|
||||
resourceMetricsMapMu.Unlock()
|
||||
t.Cleanup(func() {
|
||||
resourceMetricsMapMu.Lock()
|
||||
resourceMetricsMap = saved
|
||||
resourceMetricsMapMu.Unlock()
|
||||
})
|
||||
}
|
||||
|
||||
func newCPUMetricsTestRegistry() *prometheus.Registry {
|
||||
c := &metricsCache{cpuMetrics: cachevalue.NewFromFunc(time.Hour, cachevalue.Opts{},
|
||||
func(context.Context) (madmin.CPUMetrics, error) {
|
||||
return madmin.CPUMetrics{
|
||||
CPUCount: 4,
|
||||
LoadStat: &load.AvgStat{Load1: 2},
|
||||
TimesStat: &cpustats.TimesStat{
|
||||
User: 10, System: 20, Idle: 60, Iowait: 5, Nice: 3, Steal: 2,
|
||||
},
|
||||
}, nil
|
||||
})}
|
||||
_, _ = c.cpuMetrics.Get() // Warm the cache: the cache does not protect the resource map.
|
||||
g := NewMetricsGroup(systemCPUCollectorPath, []MetricDescriptor{
|
||||
sysCPUAvgIdleMD, sysCPUAvgIOWaitMD, sysCPULoadMD, sysCPULoadPercMD,
|
||||
sysCPUNiceMD, sysCPUStealMD, sysCPUSystemMD, sysCPUUserMD,
|
||||
}, loadCPUMetrics)
|
||||
g.SetCache(c)
|
||||
r := prometheus.NewPedanticRegistry()
|
||||
r.MustRegister(g)
|
||||
return r
|
||||
}
|
||||
|
||||
func TestLoadCPUMetricsValues(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
data map[MetricSubsystem]ResourceMetrics
|
||||
want map[string]float64
|
||||
}{
|
||||
{name: "nil-map"},
|
||||
{name: "empty-map", data: map[MetricSubsystem]ResourceMetrics{}},
|
||||
{name: "no-cpu", data: map[MetricSubsystem]ResourceMetrics{memSubsystem: {}}},
|
||||
{name: "nil-cpu", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: nil}},
|
||||
{name: "empty-cpu", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {}}},
|
||||
{name: "idle-only", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {
|
||||
getResourceKey(cpuIdle, nil): {Avg: 87.654},
|
||||
}}, want: map[string]float64{"minio_system_cpu_avg_idle": 87.65}},
|
||||
{name: "iowait-only", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {
|
||||
getResourceKey(cpuIOWait, nil): {Avg: 1.236},
|
||||
}}, want: map[string]float64{"minio_system_cpu_avg_iowait": 1.24}},
|
||||
{name: "both-rounded", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {
|
||||
getResourceKey(cpuIdle, nil): {Avg: 87.654},
|
||||
getResourceKey(cpuIOWait, nil): {Avg: 1.236},
|
||||
}}, want: map[string]float64{"minio_system_cpu_avg_idle": 87.65, "minio_system_cpu_avg_iowait": 1.24}},
|
||||
{name: "zero-keeps-existing-omission", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {
|
||||
getResourceKey(cpuIdle, nil): {Avg: 0},
|
||||
getResourceKey(cpuIOWait, nil): {Avg: 0},
|
||||
}}},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
setCPUResourceMetricsForTest(t, tc.data)
|
||||
families, err := newCPUMetricsTestRegistry().Gather()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
want := map[string]float64{
|
||||
"minio_system_cpu_load": 2, "minio_system_cpu_load_perc": 50,
|
||||
"minio_system_cpu_nice": 3, "minio_system_cpu_steal": 2,
|
||||
"minio_system_cpu_system": 20, "minio_system_cpu_user": 10,
|
||||
}
|
||||
for k, v := range tc.want {
|
||||
want[k] = v
|
||||
}
|
||||
if len(families) != len(want) {
|
||||
t.Fatalf("got %d families, want %d", len(families), len(want))
|
||||
}
|
||||
for _, family := range families {
|
||||
v, ok := want[family.GetName()]
|
||||
if !ok || family.GetType() != dto.MetricType_GAUGE || len(family.Metric) != 1 || family.Metric[0].GetGauge().GetValue() != v {
|
||||
t.Errorf("unexpected family: %v (want %v)", family, want)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadCPUMetricsConcurrentUpdate(t *testing.T) {
|
||||
for _, subsystem := range []MetricSubsystem{cpuSubsystem, memSubsystem} {
|
||||
for _, readers := range []int{1, 4} {
|
||||
t.Run(fmt.Sprintf("writer-%s/readers-%d", subsystem, readers), func(t *testing.T) {
|
||||
setCPUResourceMetricsForTest(t, map[MetricSubsystem]ResourceMetrics{})
|
||||
updateResourceMetrics(cpuSubsystem, cpuIdle, 80, nil, false)
|
||||
updateResourceMetrics(cpuSubsystem, cpuIOWait, 5, nil, false)
|
||||
r := newCPUMetricsTestRegistry()
|
||||
stop := make(chan struct{})
|
||||
done := make(chan struct{})
|
||||
ready := make(chan struct{})
|
||||
var updates atomic.Uint64
|
||||
go func() {
|
||||
defer close(done)
|
||||
if subsystem == cpuSubsystem {
|
||||
updateResourceMetrics(subsystem, cpuIdle, 80, nil, false)
|
||||
} else {
|
||||
updateResourceMetrics(subsystem, memUsed, 1024, nil, false)
|
||||
}
|
||||
close(ready)
|
||||
for {
|
||||
select {
|
||||
case <-stop:
|
||||
return
|
||||
default:
|
||||
}
|
||||
if subsystem == cpuSubsystem {
|
||||
updateResourceMetrics(subsystem, cpuIdle, 80, nil, false)
|
||||
updateResourceMetrics(subsystem, cpuIOWait, 5, nil, false)
|
||||
} else {
|
||||
updateResourceMetrics(subsystem, memUsed, 1024, nil, false)
|
||||
}
|
||||
updates.Add(1)
|
||||
runtime.Gosched()
|
||||
}
|
||||
}()
|
||||
<-ready
|
||||
defer func() { close(stop); <-done }()
|
||||
var wg sync.WaitGroup
|
||||
for range readers {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
for range 1000 {
|
||||
families, err := r.Gather()
|
||||
if err != nil || len(families) != 8 {
|
||||
t.Errorf("Gather: families=%d, error=%v", len(families), err)
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
if updates.Load() == 0 {
|
||||
t.Fatal("writer made no progress")
|
||||
}
|
||||
t.Logf("completed %d gathers with %d concurrent update iterations", readers*1000, updates.Load())
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -82,17 +82,26 @@ type markerRecoveryTarget struct {
|
||||
func replicationTestCapacity(obj ObjectLayer) func() {
|
||||
var restore []func()
|
||||
for _, pool := range obj.(*erasureServerPools).serverPools {
|
||||
pool.erasureDisksMu.Lock()
|
||||
for _, set := range pool.sets {
|
||||
original := set.getDisks
|
||||
disks := append([]StorageAPI(nil), original()...)
|
||||
original := pool.erasureDisks[set.setIndex]
|
||||
disks := append([]StorageAPI(nil), original...)
|
||||
for i, disk := range disks {
|
||||
if disk != nil {
|
||||
disks[i] = tagTestCapacityDisk{StorageAPI: disk}
|
||||
}
|
||||
}
|
||||
set.getDisks = func() []StorageAPI { return disks }
|
||||
restore = append(restore, func() { set.getDisks = original })
|
||||
// GetDisks copies this list under the same mutex. Keep its function
|
||||
// stable while background IAM scans use the fixture, both when
|
||||
// installing the adapter and when restoring the original disks.
|
||||
pool.erasureDisks[set.setIndex] = disks
|
||||
restore = append(restore, func() {
|
||||
pool.erasureDisksMu.Lock()
|
||||
pool.erasureDisks[set.setIndex] = original
|
||||
pool.erasureDisksMu.Unlock()
|
||||
})
|
||||
}
|
||||
pool.erasureDisksMu.Unlock()
|
||||
}
|
||||
return func() {
|
||||
for _, fn := range restore {
|
||||
@@ -101,6 +110,34 @@ func replicationTestCapacity(obj ObjectLayer) func() {
|
||||
}
|
||||
}
|
||||
|
||||
func TestReplicationCapacityConcurrentIAM(t *testing.T) {
|
||||
z, _ := consistencyPools(t)
|
||||
if _, _, err := initAPIHandlerTest(t.Context(), z, nil, MakeBucketOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
|
||||
defer cancel()
|
||||
iam := globalIAMSys
|
||||
started, finished := make(chan struct{}), make(chan error, 1)
|
||||
go func() {
|
||||
close(started)
|
||||
for range 20 {
|
||||
if err := iam.Load(ctx, false); err != nil {
|
||||
finished <- err
|
||||
return
|
||||
}
|
||||
}
|
||||
finished <- nil
|
||||
}()
|
||||
<-started
|
||||
for range 5000 {
|
||||
replicationTestCapacity(z)()
|
||||
}
|
||||
if err := <-finished; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func testReplicationMRFMarkerRecovery(t *testing.T, obj ObjectLayer, backend, bucket string, router http.Handler, creds auth.Credentials, tc markerRecoveryCase) {
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
defer cancel()
|
||||
|
||||
@@ -0,0 +1,264 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/minio/madmin-go/v3"
|
||||
"github.com/minio/minio-go/v7"
|
||||
"github.com/minio/minio/internal/auth"
|
||||
"github.com/minio/minio/internal/bucket/replication"
|
||||
xhttp "github.com/minio/minio/internal/http"
|
||||
"github.com/minio/minio/internal/once"
|
||||
)
|
||||
|
||||
type staleCreationCase struct {
|
||||
name string
|
||||
purged, pendingPurge, unreadable bool
|
||||
}
|
||||
|
||||
// A queued delete-marker creation must be checked against the source marker
|
||||
// under the replication lock before it is sent. Between queueing (DELETE
|
||||
// handler, GET/HEAD/LIST heal, scanner, MRF) and sending, a user purge of the
|
||||
// same marker can complete and reach the targets; the stale creation would
|
||||
// then recreate the marker there.
|
||||
func TestReplicationDeleteMarkerCreationRevalidated(t *testing.T) {
|
||||
for _, tc := range []staleCreationCase{
|
||||
{name: "current"},
|
||||
{name: "purged", purged: true},
|
||||
{name: "purge-pending", pendingPurge: true},
|
||||
{name: "unreadable", unreadable: true},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, endpoints: []string{"DeleteObject"}, objAPITest: func(obj ObjectLayer, backend, bucket string, router http.Handler, creds auth.Credentials, t *testing.T) {
|
||||
testReplicationDeleteMarkerCreationRevalidated(t, obj, backend, bucket, router, creds, tc)
|
||||
}})
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
type staleCreationTarget struct {
|
||||
arn, bucket string
|
||||
creations atomic.Int32
|
||||
purges atomic.Int32
|
||||
}
|
||||
|
||||
func testReplicationDeleteMarkerCreationRevalidated(t *testing.T, obj ObjectLayer, backend, bucket string, router http.Handler, creds auth.Credentials, tc staleCreationCase) {
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
defer cancel()
|
||||
defer replicationTestCapacity(obj)()
|
||||
stats := NewReplicationStats(ctx, nil)
|
||||
oldStats := globalReplicationStats.Swap(stats)
|
||||
defer globalReplicationStats.Store(oldStats)
|
||||
oldPool := globalReplicationPool
|
||||
defer func() { globalReplicationPool = oldPool }()
|
||||
const name = "marker"
|
||||
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, enabledBucketVersioningConfig); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := obj.PutObject(ctx, bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
target := &staleCreationTarget{arn: "arn:minio:replication::" + mustGetUUID() + ":bucket", bucket: getRandomBucketName()}
|
||||
if err := obj.MakeBucket(ctx, target.bucket, MakeBucketOptions{VersioningEnabled: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := obj.PutObject(ctx, target.bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// The remote behaves like a real receiver: a marker lookup and a
|
||||
// replicated DELETE applied to the target bucket on the same fixture.
|
||||
remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
opts := ObjectOptions{VersionID: r.URL.Query().Get("versionId"), Versioned: true}
|
||||
switch r.Method {
|
||||
case http.MethodHead:
|
||||
oi, err := obj.GetObjectInfo(r.Context(), target.bucket, name, opts)
|
||||
if oi.DeleteMarker {
|
||||
w.Header().Set(xhttp.AmzDeleteMarker, "true")
|
||||
w.Header().Set(xhttp.AmzVersionID, oi.VersionID)
|
||||
}
|
||||
if err != nil {
|
||||
writeErrorResponseHeadersOnly(w, toAPIError(r.Context(), err))
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusOK)
|
||||
case http.MethodDelete:
|
||||
opts.DeleteMarker = r.Header.Get(xhttp.MinIOSourceDeleteMarker) == "true"
|
||||
if opts.DeleteMarker {
|
||||
target.creations.Add(1)
|
||||
} else {
|
||||
target.purges.Add(1)
|
||||
}
|
||||
opts.ReplicationRequest = true
|
||||
opts.SetReplicaStatus(replication.Replica)
|
||||
if _, err := obj.DeleteObject(r.Context(), target.bucket, name, opts); err != nil && !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
writeErrorResponse(r.Context(), w, toAPIError(r.Context(), err), r.URL)
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
default:
|
||||
t.Errorf("unexpected remote method %s", r.Method)
|
||||
}
|
||||
}))
|
||||
defer remote.Close()
|
||||
client, err := minio.New(strings.TrimPrefix(remote.URL, "http://"), &minio.Options{Region: "us-east-1", MaxRetries: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
globalBucketTargetSys.Lock()
|
||||
globalBucketTargetSys.arnRemotesMap[target.arn] = arnTarget{Client: &TargetClient{Client: client, ARN: target.arn, Bucket: target.bucket}, lastRefresh: UTCNow()}
|
||||
globalBucketTargetSys.targetsMap[bucket] = append(globalBucketTargetSys.targetsMap[bucket], madmin.BucketTarget{Arn: target.arn, TargetBucket: target.bucket})
|
||||
globalBucketTargetSys.Unlock()
|
||||
globalBucketTargetSys.hMutex.Lock()
|
||||
globalBucketTargetSys.hc[client.EndpointURL().Host] = epHealth{Online: true}
|
||||
globalBucketTargetSys.hMutex.Unlock()
|
||||
rule := configs[0].Rules[0]
|
||||
rule.Destination = replication.Destination{ARN: target.arn, Bucket: target.bucket}
|
||||
cfg := replication.Config{RoleArn: target.arn, Rules: []replication.Rule{rule}}
|
||||
meta, err := globalBucketMetadataSys.Get(bucket)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
meta.replicationConfig = &cfg
|
||||
globalBucketMetadataSys.Set(bucket, meta)
|
||||
p := &ReplicationPool{ctx: ctx, objLayer: obj, workers: []chan ReplicationWorkerOperation{make(chan ReplicationWorkerOperation, 8)}, stats: stats, mrfSaveCh: make(chan MRFReplicateEntry, 8)}
|
||||
globalReplicationPool = once.NewSingleton[ReplicationPool]()
|
||||
globalReplicationPool.Set(p)
|
||||
|
||||
// The real DELETE handler creates the pending marker and queues the
|
||||
// creation task that a worker would later run.
|
||||
req, err := newTestSignedRequestV4(http.MethodDelete, "/"+bucket+"/"+name, 0, nil, creds.AccessKey, creds.SecretKey, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
w := httptest.NewRecorder()
|
||||
router.ServeHTTP(w, req)
|
||||
if w.Code != http.StatusNoContent {
|
||||
t.Fatalf("source DELETE status %d: %s", w.Code, w.Body)
|
||||
}
|
||||
var creation DeletedObjectReplicationInfo
|
||||
select {
|
||||
case op := <-p.workers[0]:
|
||||
creation = op.(DeletedObjectReplicationInfo)
|
||||
case <-time.After(3 * time.Second):
|
||||
t.Fatal("handler queued no creation task")
|
||||
}
|
||||
version := creation.DeleteMarkerVersionID
|
||||
if version == "" || creation.VersionID != "" || !creation.DeleteMarker {
|
||||
t.Fatalf("handler queued a non-creation task: %+v", creation)
|
||||
}
|
||||
before, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true})
|
||||
if !isErrMethodNotAllowed(err) || !before.DeleteMarker || before.ReplicationStatus != replication.Pending {
|
||||
t.Fatalf("source marker not pending: %+v %v", before, err)
|
||||
}
|
||||
|
||||
// Meanwhile the user purges the marker through another path. The purge is
|
||||
// either complete (version gone) or still pending on the source.
|
||||
switch {
|
||||
case tc.purged:
|
||||
if _, err := obj.DeleteObject(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
t.Fatalf("source purge left the marker: %v", err)
|
||||
}
|
||||
case tc.pendingPurge:
|
||||
opts := ObjectOptions{VersionID: version, Versioned: true, DeleteReplication: ReplicationState{
|
||||
ReplicateDecisionStr: creation.ReplicationState.ReplicateDecisionStr,
|
||||
VersionPurgeStatusInternal: target.arn + "=PENDING;",
|
||||
PurgeTargets: map[string]VersionPurgeStatusType{target.arn: replication.VersionPurgePending},
|
||||
}}
|
||||
if _, err := obj.DeleteObject(ctx, bucket, name, opts); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
oi, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true})
|
||||
if !isErrMethodNotAllowed(err) || !oi.DeleteMarker || oi.VersionPurgeStatus != replication.VersionPurgePending {
|
||||
t.Fatalf("source marker not pending purge: %+v %v", oi, err)
|
||||
}
|
||||
}
|
||||
source := obj
|
||||
if tc.unreadable {
|
||||
source = staleLookupLayer{ObjectLayer: obj}
|
||||
}
|
||||
pending, _ := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true})
|
||||
|
||||
result := replicateDelete(ctx, creation, source)
|
||||
|
||||
after, aerr := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true})
|
||||
_, terr := obj.GetObjectInfo(ctx, target.bucket, name, ObjectOptions{VersionID: version, Versioned: true})
|
||||
targetHasMarker := isErrMethodNotAllowed(terr)
|
||||
switch {
|
||||
case tc.purged, tc.pendingPurge, tc.unreadable:
|
||||
if len(result.Targets) != 0 || target.creations.Load() != 0 || target.purges.Load() != 0 {
|
||||
t.Fatalf("%s: stale creation was sent: result=%+v creations=%d purges=%d", tc.name, result, target.creations.Load(), target.purges.Load())
|
||||
}
|
||||
if targetHasMarker {
|
||||
t.Fatalf("%s: target marker recreated from the stale creation", tc.name)
|
||||
}
|
||||
if tc.purged {
|
||||
if !isErrVersionNotFound(aerr) && !isErrObjectNotFound(aerr) {
|
||||
t.Fatalf("stale creation resurrected the source marker: %+v %v", after, aerr)
|
||||
}
|
||||
} else if !isErrMethodNotAllowed(aerr) || !after.DeleteMarker ||
|
||||
after.ReplicationStatusInternal != pending.ReplicationStatusInternal ||
|
||||
after.VersionPurgeStatusInternal != pending.VersionPurgeStatusInternal ||
|
||||
after.UserDefined[ReservedMetadataPrefixLower+ReplicationTimestamp] != pending.UserDefined[ReservedMetadataPrefixLower+ReplicationTimestamp] {
|
||||
t.Fatalf("skipped creation rewrote the source marker: before=%+v after=%+v %v", pending, after, aerr)
|
||||
}
|
||||
if tc.unreadable {
|
||||
select {
|
||||
case entry := <-p.mrfSaveCh:
|
||||
if entry.versionID != version || entry.RetryCount != 1 || entry.Bucket != bucket || entry.Object != name {
|
||||
t.Fatalf("unverified creation queued wrong MRF entry: %+v", entry)
|
||||
}
|
||||
default:
|
||||
t.Fatal("unverified source read did not queue a retry")
|
||||
}
|
||||
}
|
||||
if len(p.mrfSaveCh) != 0 {
|
||||
t.Fatalf("%s: unexpected MRF entries queued: %d", tc.name, len(p.mrfSaveCh))
|
||||
}
|
||||
default:
|
||||
if result.ReplicationStatus() != replication.Completed || target.creations.Load() != 1 || !targetHasMarker {
|
||||
t.Fatalf("current creation not replicated: result=%+v creations=%d targetMarker=%v", result, target.creations.Load(), targetHasMarker)
|
||||
}
|
||||
if !isErrMethodNotAllowed(aerr) || !after.DeleteMarker || replicationStatusesMap(after.ReplicationStatusInternal)[target.arn] != replication.Completed {
|
||||
t.Fatalf("source creation state not completed: %+v %v", after, aerr)
|
||||
}
|
||||
if len(p.mrfSaveCh) != 0 {
|
||||
t.Fatal("completed creation queued MRF work")
|
||||
}
|
||||
}
|
||||
t.Logf("%s: %s checked; creations=%d purges=%d", backend, tc.name, target.creations.Load(), target.purges.Load())
|
||||
}
|
||||
|
||||
// staleLookupLayer cannot confirm the source marker: the read fails without
|
||||
// saying whether the version is present.
|
||||
type staleLookupLayer struct{ ObjectLayer }
|
||||
|
||||
func (staleLookupLayer) GetObjectInfo(context.Context, string, string, ObjectOptions) (ObjectInfo, error) {
|
||||
return ObjectInfo{}, InsufficientReadQuorum{}
|
||||
}
|
||||
@@ -41,15 +41,23 @@ const r5TagStamp = ReservedMetadataPrefixLower + TaggingTimestamp
|
||||
func r5Capacity(z *erasureServerPools) func() {
|
||||
var restores []func()
|
||||
for _, pool := range z.serverPools {
|
||||
pool.erasureDisksMu.Lock()
|
||||
for _, set := range pool.sets {
|
||||
old := set.getDisks
|
||||
disks := append([]StorageAPI(nil), old()...)
|
||||
old := pool.erasureDisks[set.setIndex]
|
||||
disks := append([]StorageAPI(nil), old...)
|
||||
for i := range disks {
|
||||
disks[i] = tagTestCapacityDisk{StorageAPI: disks[i]}
|
||||
}
|
||||
set.getDisks = func() []StorageAPI { return disks }
|
||||
restores = append(restores, func() { set.getDisks = old })
|
||||
// GetDisks copies this list under the same mutex. Keep its function
|
||||
// stable while background IAM scans are using the fixture.
|
||||
pool.erasureDisks[set.setIndex] = disks
|
||||
restores = append(restores, func() {
|
||||
pool.erasureDisksMu.Lock()
|
||||
pool.erasureDisks[set.setIndex] = old
|
||||
pool.erasureDisksMu.Unlock()
|
||||
})
|
||||
}
|
||||
pool.erasureDisksMu.Unlock()
|
||||
}
|
||||
return func() {
|
||||
for _, restore := range restores {
|
||||
@@ -58,6 +66,34 @@ func r5Capacity(z *erasureServerPools) func() {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAPITaggingCapacityConcurrentIAM(t *testing.T) {
|
||||
z, _ := consistencyPools(t)
|
||||
if _, _, err := initAPIHandlerTest(t.Context(), z, nil, MakeBucketOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
|
||||
defer cancel()
|
||||
iam := globalIAMSys
|
||||
started, finished := make(chan struct{}), make(chan error, 1)
|
||||
go func() {
|
||||
close(started)
|
||||
for range 20 {
|
||||
if err := iam.Load(ctx, false); err != nil {
|
||||
finished <- err
|
||||
return
|
||||
}
|
||||
}
|
||||
finished <- nil
|
||||
}()
|
||||
<-started
|
||||
for range 5000 {
|
||||
r5Capacity(z)()
|
||||
}
|
||||
if err := <-finished; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func r5Request(t *testing.T, router http.Handler, cred auth.Credentials, method, path, body string, headers map[string]string) *httptest.ResponseRecorder {
|
||||
t.Helper()
|
||||
r, err := newTestSignedRequestV4(method, path, int64(len(body)), strings.NewReader(body), cred.AccessKey, cred.SecretKey, headers)
|
||||
|
||||
@@ -0,0 +1,101 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful,
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"archive/tar"
|
||||
"bytes"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"sync"
|
||||
"testing"
|
||||
"testing/iotest"
|
||||
|
||||
"github.com/pierrec/lz4/v4"
|
||||
)
|
||||
|
||||
func TestUntarLZ4(t *testing.T) {
|
||||
for _, empty := range []bool{false, true} {
|
||||
t.Run(fmt.Sprintf("empty=%t", empty), func(t *testing.T) {
|
||||
want := map[string][]byte{}
|
||||
if !empty {
|
||||
want["small.txt"] = bytes.Repeat([]byte("small object\n"), 10)
|
||||
want["large.txt"] = bytes.Repeat([]byte("large object\n"), 32768)
|
||||
}
|
||||
var compressed bytes.Buffer
|
||||
compressor := lz4.NewWriter(&compressed)
|
||||
archive := tar.NewWriter(compressor)
|
||||
for name, data := range want {
|
||||
if err := archive.WriteHeader(&tar.Header{Name: name, Mode: 0o600, Size: int64(len(data))}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := archive.Write(data); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if err := archive.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := compressor.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var mu sync.Mutex
|
||||
got := map[string][]byte{}
|
||||
err := untar(t.Context(), iotest.HalfReader(bytes.NewReader(compressed.Bytes())), func(r io.Reader, info os.FileInfo, name string) error {
|
||||
// Exercise a short first read followed by streaming the remaining content.
|
||||
var output bytes.Buffer
|
||||
prefix := make([]byte, 7)
|
||||
if _, err := io.ReadFull(r, prefix); err != nil {
|
||||
return err
|
||||
}
|
||||
output.Write(prefix)
|
||||
if _, err := io.Copy(&output, r); err != nil {
|
||||
return err
|
||||
}
|
||||
if int64(output.Len()) != info.Size() {
|
||||
return fmt.Errorf("size mismatch for %s", name)
|
||||
}
|
||||
mu.Lock()
|
||||
got[name] = output.Bytes()
|
||||
mu.Unlock()
|
||||
return nil
|
||||
}, untarOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(got) != len(want) {
|
||||
t.Fatalf("objects=%d, want %d", len(got), len(want))
|
||||
}
|
||||
for name, data := range want {
|
||||
if !bytes.Equal(got[name], data) {
|
||||
t.Errorf("content mismatch for %s", name)
|
||||
}
|
||||
}
|
||||
t.Run("corrupt-header", func(t *testing.T) {
|
||||
damaged := bytes.Clone(compressed.Bytes())
|
||||
damaged[6] ^= 0xff
|
||||
err := untar(t.Context(), bytes.NewReader(damaged), func(io.Reader, os.FileInfo, string) error {
|
||||
t.Error("corrupt header must be rejected before uploading objects")
|
||||
return nil
|
||||
}, untarOptions{})
|
||||
if err == nil {
|
||||
t.Fatal("corrupt LZ4 header accepted")
|
||||
}
|
||||
})
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,148 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"io"
|
||||
"maps"
|
||||
"net/http"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/minio/minio/internal/bucket/replication"
|
||||
)
|
||||
|
||||
func TestDeleteMarkerReadDataControls(t *testing.T) {
|
||||
obj, er, disks, bucket := markerPurgeFixture(t, 4)
|
||||
for _, body := range []string{"", "inline-data-control"} {
|
||||
name := "data-" + mustGetUUID()
|
||||
oi, err := er.PutObject(t.Context(), bucket, name, mustGetPutObjReader(t, bytes.NewBufferString(body), int64(len(body)), "", ""), ObjectOptions{Versioned: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, disk := range disks {
|
||||
fi, err := disk.ReadVersion(t.Context(), "", bucket, name, oi.VersionID, ReadOptions{ReadData: true})
|
||||
if err != nil || fi.Deleted || fi.Size != int64(len(body)) || !fi.InlineData() || (body != "" && len(fi.Data) == 0) {
|
||||
t.Fatalf("data read: size=%d inline=%v bytes=%d error=%v", fi.Size, fi.InlineData(), len(fi.Data), err)
|
||||
}
|
||||
// Legacy inline data without its annotation still gets the original
|
||||
// ReadData behavior; the new marker guard must not affect it.
|
||||
delete(fi.Metadata, ReservedMetadataPrefixLower+"inline-data")
|
||||
if err := disk.WriteMetadata(t.Context(), "", bucket, name, fi); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fi, err = disk.ReadVersion(t.Context(), "", bucket, name, oi.VersionID, ReadOptions{ReadData: true})
|
||||
if err != nil || !fi.InlineData() {
|
||||
t.Fatalf("legacy inline read: %+v error=%v", fi, err)
|
||||
}
|
||||
}
|
||||
reader, err := obj.GetObjectNInfo(t.Context(), bucket, name, nil, http.Header{}, ObjectOptions{VersionID: oi.VersionID})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got, err := io.ReadAll(reader)
|
||||
reader.Close()
|
||||
if err != nil || string(got) != body {
|
||||
t.Fatalf("payload changed: %q error=%v", got, err)
|
||||
}
|
||||
}
|
||||
_, opts := seedPurgeMarker(t, er, bucket, "stored-marker-key", false)
|
||||
for _, disk := range disks {
|
||||
fi, err := disk.ReadVersion(t.Context(), "", bucket, "stored-marker-key", opts.VersionID, ReadOptions{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Preserve even an existing unusual marker key; suppress only the
|
||||
// manufacture of a new inline annotation by the data-read branch.
|
||||
fi.Metadata[ReservedMetadataPrefixLower+"inline-data"] = "original"
|
||||
if err := disk.WriteMetadata(t.Context(), "", bucket, "stored-marker-key", fi); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fi, err = disk.ReadVersion(t.Context(), "", bucket, "stored-marker-key", opts.VersionID, ReadOptions{ReadData: true})
|
||||
if err != nil || fi.Metadata[ReservedMetadataPrefixLower+"inline-data"] != "original" {
|
||||
t.Errorf("existing marker key changed: metadata=%v error=%v", fi.Metadata, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteMarkerMetadataRoundTrip(t *testing.T) {
|
||||
stamp := time.Date(2026, 9, 1, 12, 1, 2, 345, time.UTC)
|
||||
for _, free := range []bool{false, true} {
|
||||
t.Run(map[bool]string{false: "replication", true: "free-version"}[free], func(t *testing.T) {
|
||||
metadata := map[string]string{
|
||||
ReservedMetadataPrefixLower + ReplicaStatus: "REPLICA",
|
||||
ReservedMetadataPrefixLower + ReplicaTimestamp: stamp.Format(time.RFC3339Nano),
|
||||
ReservedMetadataPrefixLower + ReplicationStatus: "arn1=COMPLETED;arn2=FAILED;",
|
||||
ReservedMetadataPrefixLower + ReplicationTimestamp: stamp.Add(-time.Hour).Format(time.RFC3339Nano),
|
||||
VersionPurgeStatusKey: "arn1=PENDING;arn2=FAILED;",
|
||||
targetResetHeader("arn1"): "original;reset1",
|
||||
targetResetHeader("arn2"): "original;reset2",
|
||||
ReservedMetadataPrefixLower + "unknown": "",
|
||||
}
|
||||
if free {
|
||||
metadata[ReservedMetadataPrefixLower+freeVersion] = ""
|
||||
metadata[metaTierName] = "tier"
|
||||
metadata[metaTierObjName] = "remote-object"
|
||||
metadata[metaTierVersionID] = "remote-version"
|
||||
}
|
||||
fi := FileInfo{VersionID: mustGetUUID(), Deleted: true, ModTime: stamp, Metadata: maps.Clone(metadata)}
|
||||
// Deliberately conflicting parsed state must never replace stored
|
||||
// values or add a key missing from the raw metadata.
|
||||
fi.ReplicationState = ReplicationState{ReplicaStatus: replication.Failed, ReplicaTimeStamp: stamp.Add(time.Hour), ResetStatusesMap: map[string]string{"new-arn": "invented"}}
|
||||
fi.SetHealing()
|
||||
fi.SetDataMov()
|
||||
fi.SetTierFreeVersionID(mustGetUUID())
|
||||
fi.SetTierFreeVersion()
|
||||
fi.SetSkipTierFreeVersion()
|
||||
xl := xlMetaV2{}
|
||||
if err := xl.AddVersion(fi); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
stored, err := xl.getIdx(0)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got := make(map[string]string)
|
||||
for k, v := range stored.DeleteMarker.MetaSys {
|
||||
got[k] = string(v)
|
||||
}
|
||||
if !maps.Equal(got, metadata) {
|
||||
t.Errorf("stored metadata=%v want=%v", got, metadata)
|
||||
}
|
||||
decoded, err := xl.ToFileInfo("bucket", "marker", fi.VersionID, true, true)
|
||||
if err != nil || !decoded.ModTime.Equal(fi.ModTime) || decoded.TierFreeVersion() != free {
|
||||
t.Errorf("round trip identity: %+v error=%v", decoded, err)
|
||||
}
|
||||
if decoded.ReplicationState.ReplicaStatus != replication.Replica || decoded.ReplicationState.PurgeTargets["arn2"] != replication.VersionPurgeFailed || decoded.ReplicationState.Targets["arn1"] != replication.Completed {
|
||||
t.Errorf("lost parsed replication state: %+v", decoded.ReplicationState)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteMarkerCreationMetadata(t *testing.T) {
|
||||
_, er, disks, bucket := markerPurgeFixture(t, 4)
|
||||
_, opts := seedPurgeMarker(t, er, bucket, "empty-key", false)
|
||||
for i, disk := range disks {
|
||||
fi, err := disk.ReadVersion(t.Context(), "", bucket, "empty-key", opts.VersionID, ReadOptions{})
|
||||
if err != nil || fi.Metadata[ReservedMetadataPrefixLower+ReplicaStatus] != "REPLICA" || fi.Metadata[ReservedMetadataPrefixLower+ReplicaTimestamp] != opts.DeleteReplication.ReplicaTimeStamp.Format(time.RFC3339Nano) {
|
||||
t.Errorf("disk %d lost new marker replica identity: metadata=%v error=%v", i, fi.Metadata, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// deleteMarkerMetadata preserves the complete stored marker state during
|
||||
// healing. The parsed ReplicationState is lossy; even absent raw keys must not
|
||||
// be reconstructed from it. Only creation writers without stored metadata
|
||||
// use typed state, and never invent timestamps.
|
||||
func deleteMarkerMetadata(fi FileInfo) map[string][]byte {
|
||||
meta := make(map[string][]byte, len(fi.Metadata))
|
||||
for k, v := range fi.Metadata {
|
||||
switch k {
|
||||
case xMinIOHealing, xMinIODataMov,
|
||||
ReservedMetadataPrefixLower + tierFVID,
|
||||
ReservedMetadataPrefixLower + tierFVMarker,
|
||||
ReservedMetadataPrefixLower + tierSkipFVID:
|
||||
continue
|
||||
}
|
||||
meta[k] = []byte(v)
|
||||
}
|
||||
if len(meta) != 0 || fi.Healing() || fi.DataMov() {
|
||||
return meta
|
||||
}
|
||||
rs := fi.ReplicationState
|
||||
if !rs.ReplicaStatus.Empty() {
|
||||
meta[ReservedMetadataPrefixLower+ReplicaStatus] = []byte(rs.ReplicaStatus)
|
||||
if !rs.ReplicaTimeStamp.IsZero() {
|
||||
meta[ReservedMetadataPrefixLower+ReplicaTimestamp] = []byte(rs.ReplicaTimeStamp.UTC().Format(time.RFC3339Nano))
|
||||
}
|
||||
}
|
||||
if rs.ReplicationStatusInternal != "" {
|
||||
meta[ReservedMetadataPrefixLower+ReplicationStatus] = []byte(rs.ReplicationStatusInternal)
|
||||
if !rs.ReplicationTimeStamp.IsZero() {
|
||||
meta[ReservedMetadataPrefixLower+ReplicationTimestamp] = []byte(rs.ReplicationTimeStamp.UTC().Format(time.RFC3339Nano))
|
||||
}
|
||||
}
|
||||
if rs.VersionPurgeStatusInternal != "" {
|
||||
meta[VersionPurgeStatusKey] = []byte(rs.VersionPurgeStatusInternal)
|
||||
}
|
||||
for k, v := range rs.ResetStatusesMap {
|
||||
if !strings.HasPrefix(k, ReservedMetadataPrefixLower+ReplicationReset) {
|
||||
k = targetResetHeader(k)
|
||||
}
|
||||
meta[k] = []byte(v)
|
||||
}
|
||||
return meta
|
||||
}
|
||||
@@ -0,0 +1,274 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This file is part of Silo Object Storage stack
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"reflect"
|
||||
"slices"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Use complete metadata so the selected version can also be decoded by LIST
|
||||
// and the object read path, rather than testing only synthetic shallow headers.
|
||||
func quorumNullVersion(t testing.TB, generation int64) xlMetaV2ShallowVersion {
|
||||
t.Helper()
|
||||
fi := newFileInfo("fixed/object", 12, 4)
|
||||
fi.Erasure.Index = 1
|
||||
fi.DataDir = "11111111-1111-1111-1111-111111111111"
|
||||
fi.ModTime = time.Unix(generation, 0).UTC()
|
||||
fi.Size = 8192
|
||||
fi.Parts = []ObjectPartInfo{{Number: 1, Size: fi.Size, ActualSize: fi.Size}}
|
||||
fi.Metadata = map[string]string{"etag": fmt.Sprintf("%032d", generation)}
|
||||
var xl xlMetaV2
|
||||
if err := xl.AddVersion(fi); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return xl.versions[0]
|
||||
}
|
||||
|
||||
// Start with forward, reverse and interleaved orders, then fixed-seed shuffles.
|
||||
func quorumVersionOrders(input [][]xlMetaV2ShallowVersion) [][][]xlMetaV2ShallowVersion {
|
||||
orders := [][][]xlMetaV2ShallowVersion{slices.Clone(input), slices.Clone(input)}
|
||||
slices.Reverse(orders[1])
|
||||
interleaved := make([][]xlMetaV2ShallowVersion, 0, len(input))
|
||||
for i, j := 0, len(input)-1; i <= j; i, j = i+1, j-1 {
|
||||
interleaved = append(interleaved, input[i])
|
||||
if i != j {
|
||||
interleaved = append(interleaved, input[j])
|
||||
}
|
||||
}
|
||||
orders = append(orders, interleaved)
|
||||
for seed := range int64(32) {
|
||||
order := slices.Clone(input)
|
||||
rand.New(rand.NewSource(seed)).Shuffle(len(order), func(i, j int) {
|
||||
order[i], order[j] = order[j], order[i]
|
||||
})
|
||||
orders = append(orders, order)
|
||||
}
|
||||
return orders
|
||||
}
|
||||
|
||||
func TestMergeXLV2SingleNullQuorum(t *testing.T) {
|
||||
v := []xlMetaV2ShallowVersion{quorumNullVersion(t, 100), quorumNullVersion(t, 200), quorumNullVersion(t, 300)}
|
||||
for _, tt := range []struct {
|
||||
name string
|
||||
counts [3]int
|
||||
empty int
|
||||
quorum int
|
||||
want int
|
||||
}{
|
||||
{"consistent", [3]int{16}, 0, 8, 0},
|
||||
{"old9-new7", [3]int{9, 7}, 0, 8, 0},
|
||||
{"old15-new1", [3]int{15, 1}, 0, 8, 0},
|
||||
{"old3-new1", [3]int{3, 1}, 0, 3, 0},
|
||||
{"two-quorums", [3]int{8, 8}, 0, 8, 1},
|
||||
{"empty-streams", [3]int{9, 1}, 6, 8, 0},
|
||||
{"two-subquorums", [3]int{7, 7}, 2, 8, -1},
|
||||
{"three-subquorums", [3]int{6, 5, 5}, 0, 8, -1},
|
||||
} {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
var input [][]xlMetaV2ShallowVersion
|
||||
for g, n := range tt.counts {
|
||||
for range n {
|
||||
input = append(input, []xlMetaV2ShallowVersion{v[g]})
|
||||
}
|
||||
}
|
||||
input = append(input, make([][]xlMetaV2ShallowVersion, tt.empty)...)
|
||||
want := []xlMetaV2ShallowVersion{}
|
||||
if tt.want >= 0 {
|
||||
want = append(want, v[tt.want])
|
||||
}
|
||||
for order, versions := range quorumVersionOrders(input) {
|
||||
for _, strict := range []bool{false, true} {
|
||||
for _, requested := range []int{0, 1} {
|
||||
got := mergeXLV2Versions(tt.quorum, strict, requested, versions...)
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("order=%d strict=%v requested=%d: got %#v, want %#v", order, strict, requested, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestMergeXLV2SingleNullHeaderGroups(t *testing.T) {
|
||||
old := quorumNullVersion(t, 100)
|
||||
newer := quorumNullVersion(t, 200)
|
||||
differentSig := old
|
||||
differentSig.header.Signature[0]++
|
||||
differentFlags := old
|
||||
differentFlags.header.Flags ^= xlFlagInlineData
|
||||
for _, tt := range []struct {
|
||||
name string
|
||||
input [][]xlMetaV2ShallowVersion
|
||||
strict bool
|
||||
want []xlMetaV2ShallowVersion
|
||||
}{
|
||||
{"signature-nonstrict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentSig}, {newer}}, false, []xlMetaV2ShallowVersion{differentSig}},
|
||||
{"signature-strict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentSig}, {newer}}, true, []xlMetaV2ShallowVersion{}},
|
||||
{"flags-nonstrict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentFlags}, {newer}}, false, []xlMetaV2ShallowVersion{}},
|
||||
{"flags-strict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentFlags}, {newer}}, true, []xlMetaV2ShallowVersion{}},
|
||||
} {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
got := mergeXLV2Versions(3, tt.strict, 0, tt.input...)
|
||||
if !reflect.DeepEqual(got, tt.want) {
|
||||
t.Fatalf("got %#v, want %#v", got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestResolveSingleNullQuorum(t *testing.T) {
|
||||
old, newer := quorumNullVersion(t, 100), quorumNullVersion(t, 200)
|
||||
input := make([][]xlMetaV2ShallowVersion, 16)
|
||||
for i := range input {
|
||||
input[i] = []xlMetaV2ShallowVersion{old}
|
||||
if i >= 9 {
|
||||
input[i] = []xlMetaV2ShallowVersion{newer}
|
||||
}
|
||||
}
|
||||
for order, versions := range quorumVersionOrders(input) {
|
||||
entries := make(metaCacheEntries, len(versions))
|
||||
for i := range versions {
|
||||
xl := &xlMetaV2{versions: versions[i]}
|
||||
metadata, err := xl.AppendTo(nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
entries[i] = metaCacheEntry{name: "fixed/object", metadata: metadata}
|
||||
}
|
||||
selected, ok := entries.resolve(&metadataResolutionParams{objQuorum: 8, dirQuorum: 8, requestedVersions: 1})
|
||||
if !ok {
|
||||
t.Fatalf("order %d: a complete null object version has quorum but was omitted", order)
|
||||
}
|
||||
xl, err := selected.xlmeta()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !reflect.DeepEqual(xl.versions, []xlMetaV2ShallowVersion{old}) {
|
||||
t.Fatalf("order %d: incorrect complete selected version: %#v", order, xl.versions)
|
||||
}
|
||||
fi, err := xl.ToFileInfo("bucket", "fixed/object", "", false, false)
|
||||
if err != nil || !fi.IsValid() || fi.Size != 8192 || !fi.ModTime.Equal(time.Unix(100, 0)) {
|
||||
t.Fatalf("order %d: selected metadata cannot be read: %+v, %v", order, fi, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// These are compatibility assertions for inputs outside the narrow repair.
|
||||
// In particular, the mixed-version outputs below are not a new correctness
|
||||
// contract for version enumeration; that behavior requires a separate change.
|
||||
func TestMergeXLV2NullHistoriesUnchanged(t *testing.T) {
|
||||
old, newer := quorumNullVersion(t, 100), quorumNullVersion(t, 300)
|
||||
middle := quorumNullVersion(t, 200)
|
||||
middle.header.VersionID = [16]byte{1}
|
||||
highest := quorumNullVersion(t, 400)
|
||||
highest.header.VersionID = [16]byte{2}
|
||||
for _, tt := range []struct {
|
||||
name string
|
||||
input [][]xlMetaV2ShallowVersion
|
||||
want []xlMetaV2ShallowVersion
|
||||
}{
|
||||
{"mixed-forward", [][]xlMetaV2ShallowVersion{{old}, {old}, {old}, {newer, middle}, {middle}, {middle}}, []xlMetaV2ShallowVersion{middle}},
|
||||
{"mixed-reverse", [][]xlMetaV2ShallowVersion{{middle}, {middle}, {newer, middle}, {old}, {old}, {old}}, []xlMetaV2ShallowVersion{old}},
|
||||
{"history-pruned-first", [][]xlMetaV2ShallowVersion{{highest, old}, {highest, old}, {highest, old}, {highest, newer}}, []xlMetaV2ShallowVersion{highest}},
|
||||
} {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
got := mergeXLV2Versions(3, false, 0, tt.input...)
|
||||
if !reflect.DeepEqual(got, tt.want) {
|
||||
t.Fatalf("excluded history changed: got %#v, want %#v", got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
for _, kind := range []string{"delete", "free", "legacy", "mixed-ec", "legacy-modern-ec", "nonzero-strict"} {
|
||||
t.Run(kind, func(t *testing.T) {
|
||||
a, b := old, newer
|
||||
switch kind {
|
||||
case "delete":
|
||||
b.header.Type = DeleteType
|
||||
case "free":
|
||||
b.header.Flags |= xlFlagFreeVersion
|
||||
case "legacy":
|
||||
b.header.Type = LegacyType
|
||||
case "mixed-ec":
|
||||
b.header.EcN, b.header.EcM = 8, 8
|
||||
case "legacy-modern-ec":
|
||||
b.header.EcN, b.header.EcM = 0, 0
|
||||
case "nonzero-strict":
|
||||
a.header.VersionID, b.header.VersionID = [16]byte{1}, [16]byte{1}
|
||||
}
|
||||
for _, requested := range []int{0, 1} {
|
||||
got := mergeXLV2Versions(3, true, requested, []xlMetaV2ShallowVersion{a}, []xlMetaV2ShallowVersion{a}, []xlMetaV2ShallowVersion{a}, []xlMetaV2ShallowVersion{b})
|
||||
if len(got) != 0 {
|
||||
t.Fatalf("excluded input changed: requested=%d got %#v", requested, got)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
// All-legacy EC headers are eligible, unlike mixing legacy and modern EC.
|
||||
old.header.EcN, old.header.EcM = 0, 0
|
||||
newer.header.EcN, newer.header.EcM = 0, 0
|
||||
got := mergeXLV2Versions(3, false, 0, []xlMetaV2ShallowVersion{old}, []xlMetaV2ShallowVersion{old}, []xlMetaV2ShallowVersion{old}, []xlMetaV2ShallowVersion{newer})
|
||||
if !reflect.DeepEqual(got, []xlMetaV2ShallowVersion{old}) {
|
||||
t.Fatalf("all-legacy EC lost the old quorum: %#v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkResolveSingleNullQuorumPage(b *testing.B) {
|
||||
old, newer := quorumNullVersion(b, 100), quorumNullVersion(b, 200)
|
||||
for _, kind := range []string{"consistent", "old-first", "new-first"} {
|
||||
b.Run(kind, func(b *testing.B) {
|
||||
var templates [16]metaCacheEntry
|
||||
for i := range templates {
|
||||
v := old
|
||||
if kind == "old-first" && i >= 9 || kind == "new-first" && i < 7 {
|
||||
v = newer
|
||||
}
|
||||
xl := &xlMetaV2{versions: []xlMetaV2ShallowVersion{v}}
|
||||
data, err := xl.AppendTo(nil)
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
templates[i] = metaCacheEntry{metadata: data}
|
||||
}
|
||||
var names [1000]string
|
||||
for i := range names {
|
||||
names[i] = fmt.Sprintf("fixed/%04d", i)
|
||||
}
|
||||
resolver := metadataResolutionParams{objQuorum: 8, dirQuorum: 8, requestedVersions: 1}
|
||||
b.ReportAllocs()
|
||||
b.ResetTimer()
|
||||
for b.Loop() {
|
||||
for _, name := range names {
|
||||
entries := templates
|
||||
for i := range entries {
|
||||
entries[i].name = name
|
||||
}
|
||||
selected, ok := metaCacheEntries(entries[:]).resolve(&resolver)
|
||||
if ok && selected.reusable {
|
||||
metaDataPoolPut(selected.metadata)
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
+83
-36
@@ -1630,7 +1630,7 @@ func (x *xlMetaV2) AddVersion(fi FileInfo) error {
|
||||
ventry.DeleteMarker = &xlMetaV2DeleteMarker{
|
||||
VersionID: uv,
|
||||
ModTime: fi.ModTime.UnixNano(),
|
||||
MetaSys: make(map[string][]byte),
|
||||
MetaSys: deleteMarkerMetadata(fi),
|
||||
}
|
||||
} else {
|
||||
ventry.Type = ObjectType
|
||||
@@ -1941,6 +1941,10 @@ func mergeXLV2Versions(quorum int, strict bool, requestedVersions int, versions
|
||||
// No need for non-strict checks if quorum is 1.
|
||||
strict = true
|
||||
}
|
||||
// Keep the original stream shapes: pruning must not make a versioned object
|
||||
// eligible for the single null-version recount below.
|
||||
originalVersions := versions
|
||||
var checkedSingleNull, singleNull bool
|
||||
// Shallow copy input
|
||||
versions = append(make([][]xlMetaV2ShallowVersion, 0, len(versions)), versions...)
|
||||
|
||||
@@ -2015,44 +2019,23 @@ func mergeXLV2Versions(quorum int, strict bool, requestedVersions int, versions
|
||||
// Version IDs match, but otherwise unable to resolve.
|
||||
// We are either strict, or don't have enough information to match.
|
||||
// Switch to a pure counting algo.
|
||||
x := make(map[xlMetaV2VersionHeader]int, len(tops))
|
||||
for _, a := range tops {
|
||||
if a.header.VersionID != ver.header.VersionID {
|
||||
continue
|
||||
}
|
||||
if !strict {
|
||||
// we must match EC, when we are not strict.
|
||||
if !a.header.matchesEC(ver.header) {
|
||||
continue
|
||||
}
|
||||
|
||||
a.header.Signature = [4]byte{}
|
||||
}
|
||||
x[a.header]++
|
||||
}
|
||||
latestCount = 0
|
||||
for k, v := range x {
|
||||
if v < latestCount {
|
||||
continue
|
||||
}
|
||||
if v == latestCount && latest.header.sortsBefore(k) {
|
||||
// Tiebreak, use sort.
|
||||
continue
|
||||
}
|
||||
for _, a := range tops {
|
||||
hdr := a.header
|
||||
if !strict {
|
||||
hdr.Signature = [4]byte{}
|
||||
}
|
||||
if hdr == k {
|
||||
latest = a
|
||||
}
|
||||
}
|
||||
latestCount = v
|
||||
}
|
||||
latest, latestCount = countXLV2Versions(tops, ver.header, latest, strict)
|
||||
break
|
||||
}
|
||||
}
|
||||
if latestCount < quorum {
|
||||
if !checkedSingleNull {
|
||||
singleNull = singleNullVersionStreams(originalVersions)
|
||||
checkedSingleNull = true
|
||||
}
|
||||
if singleNull {
|
||||
// A newer minority at the end can hide an older quorum from
|
||||
// the selection loop. Recount before discarding the null ID.
|
||||
if candidate, count := countXLV2Versions(tops, latest.header, latest, strict); count >= quorum {
|
||||
latest, latestCount = candidate, count
|
||||
}
|
||||
}
|
||||
}
|
||||
if latestCount >= quorum {
|
||||
merged = append(merged, latest)
|
||||
|
||||
@@ -2113,6 +2096,70 @@ func mergeXLV2Versions(quorum int, strict bool, requestedVersions int, versions
|
||||
return merged
|
||||
}
|
||||
|
||||
// singleNullVersionStreams excludes histories and other version types from the
|
||||
// additional recount. Check the original inputs, before any stream is pruned.
|
||||
func singleNullVersionStreams(versions [][]xlMetaV2ShallowVersion) bool {
|
||||
var ec xlMetaV2VersionHeader
|
||||
var haveEC bool
|
||||
for _, stream := range versions {
|
||||
if len(stream) == 0 {
|
||||
continue
|
||||
}
|
||||
if len(stream) != 1 {
|
||||
return false
|
||||
}
|
||||
h := stream[0].header
|
||||
if h.VersionID != [16]byte{} || h.Type != ObjectType || h.FreeVersion() {
|
||||
return false
|
||||
}
|
||||
if haveEC && (h.EcN != ec.EcN || h.EcM != ec.EcM) {
|
||||
return false
|
||||
}
|
||||
ec, haveEC = h, true
|
||||
}
|
||||
return haveEC
|
||||
}
|
||||
|
||||
// countXLV2Versions selects the most frequent compatible header for reference's
|
||||
// VersionID, retaining the existing sort tiebreak and last matching entry.
|
||||
func countXLV2Versions(tops []xlMetaV2ShallowVersion, reference xlMetaV2VersionHeader, latest xlMetaV2ShallowVersion, strict bool) (xlMetaV2ShallowVersion, int) {
|
||||
x := make(map[xlMetaV2VersionHeader]int, len(tops))
|
||||
for _, a := range tops {
|
||||
if a.header.VersionID != reference.VersionID {
|
||||
continue
|
||||
}
|
||||
if !strict {
|
||||
// we must match EC, when we are not strict.
|
||||
if !a.header.matchesEC(reference) {
|
||||
continue
|
||||
}
|
||||
a.header.Signature = [4]byte{}
|
||||
}
|
||||
x[a.header]++
|
||||
}
|
||||
var latestCount int
|
||||
for k, v := range x {
|
||||
if v < latestCount {
|
||||
continue
|
||||
}
|
||||
if v == latestCount && latest.header.sortsBefore(k) {
|
||||
// Tiebreak, use sort.
|
||||
continue
|
||||
}
|
||||
for _, a := range tops {
|
||||
hdr := a.header
|
||||
if !strict {
|
||||
hdr.Signature = [4]byte{}
|
||||
}
|
||||
if hdr == k {
|
||||
latest = a
|
||||
}
|
||||
}
|
||||
latestCount = v
|
||||
}
|
||||
return latest, latestCount
|
||||
}
|
||||
|
||||
type xlMetaBuf []byte
|
||||
|
||||
// ToFileInfo converts xlMetaV2 into a common FileInfo datastructure
|
||||
|
||||
+3
-1
@@ -1719,7 +1719,9 @@ func (s *xlStorage) ReadVersion(ctx context.Context, origvolume, volume, path, v
|
||||
defer metaDataPoolPut(buf)
|
||||
}
|
||||
|
||||
if readData {
|
||||
// Delete markers have no payload. In particular, do not manufacture an
|
||||
// inline-data metadata key when healing their zero-size FileInfo.
|
||||
if readData && !fi.Deleted {
|
||||
if len(fi.Data) > 0 || fi.Size == 0 {
|
||||
if fi.InlineData() {
|
||||
// If written with header we are fine.
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
# Compression Guide
|
||||
|
||||
For SILO's current SSE-C compression and historical-object boundary, see the [maintained encryption guide](https://silo.pgsty.com/administration/server-side-encryption/server-side-encryption-sse-c/) and [SSE-C replica design](https://silo.pgsty.com/blog/design/ssec-replica-integrity/). Check the documented release boundary before applying main-branch behavior to an older binary.
|
||||
|
||||
Silo server allows streaming compression to ensure efficient disk space usage.
|
||||
Compression happens inflight, i.e objects are compressed before being written to disk(s).
|
||||
Silo uses [`klauspost/compress/s2`](https://github.com/klauspost/compress/tree/master/s2)
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
# Federation Quickstart Guide *Federation feature is deprecated and should be avoided for future deployments*
|
||||
|
||||
The maintained [federated CopyObject design](https://silo.pgsty.com/blog/design/federated-copy-object/) records the destination encryption, checksum, Object Lock and committed-response contract, including its Server release boundary.
|
||||
|
||||
This document explains how to configure Silo with `Bucket lookup from DNS` style federation.
|
||||
|
||||
## Cross-deployment copy behavior
|
||||
|
||||
@@ -4,9 +4,9 @@ go 1.27.1
|
||||
|
||||
// Console and MC retain their historical module paths for best-effort upstream
|
||||
// compatibility. Pin the maintained PGSTY implementations used by SILO.
|
||||
replace github.com/minio/console => github.com/pgsty/silo-console v0.0.0-20260916034812-56dfe455ac2f
|
||||
replace github.com/minio/console => github.com/pgsty/silo-console v0.0.0-20260916075814-1360e26d976d
|
||||
|
||||
replace github.com/minio/mc => github.com/pgsty/mc v0.0.0-20260913012246-4f609a4da3bb
|
||||
replace github.com/minio/mc => github.com/pgsty/mc v0.0.0-20260916070421-e952aa78f10a
|
||||
|
||||
// v22.7.0 does not compile on NetBSD because its unix implementation uses
|
||||
// CLOCK_MONOTONIC, which is unavailable there. Keep the last portable release
|
||||
@@ -68,7 +68,7 @@ require (
|
||||
github.com/minio/kms-go/kes v0.3.1
|
||||
github.com/minio/kms-go/kms v0.6.0
|
||||
github.com/minio/madmin-go/v3 v3.0.110
|
||||
github.com/minio/minio-go/v7 v7.3.1-0.20260910142817-60bd07042d49
|
||||
github.com/minio/minio-go/v7 v7.3.1-0.20260915093545-32e1f32cb176
|
||||
github.com/minio/mux v1.10.1
|
||||
github.com/minio/selfupdate v0.6.0
|
||||
github.com/minio/simdjson-go v0.4.5
|
||||
@@ -81,9 +81,9 @@ require (
|
||||
github.com/nats-io/stan.go v0.10.4
|
||||
github.com/ncw/directio v1.0.5
|
||||
github.com/nsqio/go-nsq v1.1.0
|
||||
github.com/pgsty/silo-pkg/v3 v3.14.0
|
||||
github.com/pgsty/silo-pkg/v3 v3.14.1
|
||||
github.com/philhofer/fwd v1.2.0
|
||||
github.com/pierrec/lz4/v4 v4.1.29
|
||||
github.com/pierrec/lz4/v4 v4.1.30
|
||||
github.com/pkg/errors v0.9.1
|
||||
github.com/pkg/sftp v1.13.11
|
||||
github.com/pkg/xattr v0.4.12
|
||||
@@ -175,7 +175,7 @@ require (
|
||||
github.com/go-openapi/runtime v0.33.1 // indirect
|
||||
github.com/go-openapi/runtime/server-middleware v0.33.1 // indirect
|
||||
github.com/go-openapi/spec v1.0.0 // indirect
|
||||
github.com/go-openapi/strfmt v0.27.0 // indirect
|
||||
github.com/go-openapi/strfmt v0.27.2 // indirect
|
||||
github.com/go-openapi/swag v0.29.1 // indirect
|
||||
github.com/go-openapi/swag/cmdutils v0.29.1 // indirect
|
||||
github.com/go-openapi/swag/conv v0.29.1 // indirect
|
||||
@@ -225,7 +225,7 @@ require (
|
||||
github.com/lestrrat-go/dsig-secp256k1 v1.0.0 // indirect
|
||||
github.com/lestrrat-go/httpcc v1.0.1 // indirect
|
||||
github.com/lestrrat-go/httprc/v3 v3.0.6 // indirect
|
||||
github.com/lestrrat-go/jwx/v3 v3.2.0 // indirect
|
||||
github.com/lestrrat-go/jwx/v3 v3.3.0 // indirect
|
||||
github.com/lestrrat-go/option/v2 v2.0.0 // indirect
|
||||
github.com/lucasb-eyer/go-colorful v1.4.1 // indirect
|
||||
github.com/lufia/plan9stats v0.0.0-20260802145828-341c2f0c90b5 // indirect
|
||||
|
||||
@@ -214,8 +214,8 @@ github.com/go-openapi/runtime/server-middleware v0.33.1 h1:IAeKbwWnBnpsYTpuPVS8t
|
||||
github.com/go-openapi/runtime/server-middleware v0.33.1/go.mod h1:2Gej5fDxqeJxY+w38vxXYW0BgFASfgBsJ5rXwN1Fseg=
|
||||
github.com/go-openapi/spec v1.0.0 h1:JtB/GHOj+eetjse6YvxqLze88oEekl/4uPBethvzRrA=
|
||||
github.com/go-openapi/spec v1.0.0/go.mod h1:boj1PRhqS0x5jylgcNp9BRWxenphrh7vVNYZ90IRCoo=
|
||||
github.com/go-openapi/strfmt v0.27.0 h1:kbcTeaD9TXuXD0hhMXzuYa1sdTo6+dWGvwjW93E80IM=
|
||||
github.com/go-openapi/strfmt v0.27.0/go.mod h1:s/qhDqfY72irigXUGJmtgid2Rm+3tnz3k8hZaRmvWYc=
|
||||
github.com/go-openapi/strfmt v0.27.2 h1:SG32SlbwNy92s0KJiVxt2joJeFdqIYHvwrA0OU6HqzQ=
|
||||
github.com/go-openapi/strfmt v0.27.2/go.mod h1:M4CKsMO0Fb8qR10+1Ra75wCKNNquy+Vj+4LWZrhTo2E=
|
||||
github.com/go-openapi/swag v0.29.1 h1:C6EeWzUwQtcWEhE9eqBdUubGXxhWY4PlzHMLD7kLaiQ=
|
||||
github.com/go-openapi/swag v0.29.1/go.mod h1:BzxEXKiPlSXRsRTv1KSBF/BpGKHxA/YciCnr4tv9bvA=
|
||||
github.com/go-openapi/swag/cmdutils v0.29.1 h1:3DorPGfUdE80BogKY22EzoHBcHMrkVomZMoV7kS4ANY=
|
||||
@@ -411,8 +411,8 @@ github.com/lestrrat-go/httpcc v1.0.1 h1:ydWCStUeJLkpYyjLDHihupbn2tYmZ7m22BGkcvZZ
|
||||
github.com/lestrrat-go/httpcc v1.0.1/go.mod h1:qiltp3Mt56+55GPVCbTdM9MlqhvzyuL6W/NMDA8vA5E=
|
||||
github.com/lestrrat-go/httprc/v3 v3.0.6 h1:4FpLQ18KK/ypPbVU3NLWJNRvH3kcYiqKqWfKGqNWxxI=
|
||||
github.com/lestrrat-go/httprc/v3 v3.0.6/go.mod h1:mSMtkZW92Z98M5YoNNztbRGxbXHql7tSitCvaxvo9l0=
|
||||
github.com/lestrrat-go/jwx/v3 v3.2.0 h1:Jb3zBASTSZXz7gzzSAfYqxXF8KejvKC4xWoePLQqXCA=
|
||||
github.com/lestrrat-go/jwx/v3 v3.2.0/go.mod h1:38vQ8iWKq3qRSbilbzvzdQPuywhowwuR03lhkYskyrw=
|
||||
github.com/lestrrat-go/jwx/v3 v3.3.0 h1:OXcYvQOQ7cxWzeZ/Q9sYk8ABe/kCSI371WmuACiCT+4=
|
||||
github.com/lestrrat-go/jwx/v3 v3.3.0/go.mod h1:eIJhDcKHBwcgxqv8RiIylV67TVl1wJp/265IAHY1Db8=
|
||||
github.com/lestrrat-go/option/v2 v2.0.0 h1:XxrcaJESE1fokHy3FpaQ/cXW8ZsIdWcdFzzLOcID3Ss=
|
||||
github.com/lestrrat-go/option/v2 v2.0.0/go.mod h1:oSySsmzMoR0iRzCDCaUfsCzxQHUEuhOViQObyy7S6Vg=
|
||||
github.com/lib/pq v1.10.4/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o=
|
||||
@@ -474,8 +474,8 @@ github.com/minio/madmin-go/v3 v3.0.110 h1:FIYekj7YPc430ffpXFWiUtyut3qBt/unIAcDzJ
|
||||
github.com/minio/madmin-go/v3 v3.0.110/go.mod h1:WOe2kYmYl1OIlY2DSRHVQ8j1v4OItARQ6jGyQqcCud8=
|
||||
github.com/minio/md5-simd v1.1.2 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34=
|
||||
github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM=
|
||||
github.com/minio/minio-go/v7 v7.3.1-0.20260910142817-60bd07042d49 h1:xDONRymps7J67VLaNzMqq25qYCHHm3JgXnKPRf1NoFs=
|
||||
github.com/minio/minio-go/v7 v7.3.1-0.20260910142817-60bd07042d49/go.mod h1:KUPWdecEO1LWyUz+sTGXAuf2jZHrPh5fCsRH86QbPfk=
|
||||
github.com/minio/minio-go/v7 v7.3.1-0.20260915093545-32e1f32cb176 h1:TIrbAkJoMN/kbtV4LAitf/eidrbvvYcujXPbSe+4j/0=
|
||||
github.com/minio/minio-go/v7 v7.3.1-0.20260915093545-32e1f32cb176/go.mod h1:KUPWdecEO1LWyUz+sTGXAuf2jZHrPh5fCsRH86QbPfk=
|
||||
github.com/minio/mux v1.10.1 h1:grrK8SwRKbkNFE6qG7WAvFGH09bB46d5teOOtKfQ14s=
|
||||
github.com/minio/mux v1.10.1/go.mod h1:INYT4sMSTJy0QWUEA/E2DZNxJ5sAxIwbnyZjkzNFRfE=
|
||||
github.com/minio/pkg/v3 v3.6.1 h1:gaNT80BS/iuIany5ylTkVmfN4s6UYY30OtImFv4GQA8=
|
||||
@@ -545,16 +545,16 @@ github.com/orisano/pixelmatch v0.0.0-20220722002657-fb0b55479cde/go.mod h1:nZgzb
|
||||
github.com/pascaldekloe/goe v0.1.0/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc=
|
||||
github.com/pborman/getopt v0.0.0-20170112200414-7148bc3a4c30/go.mod h1:85jBQOZwpVEaDAr341tbn15RS4fCAsIst0qp7i8ex1o=
|
||||
github.com/pelletier/go-toml v1.2.0/go.mod h1:5z9KED0ma1S8pY6P1sdut58dfprrGBbd/94hg7ilaic=
|
||||
github.com/pgsty/mc v0.0.0-20260913012246-4f609a4da3bb h1:S7IAYBoKvFRqrw5MUsQE34/q4cKC4K7gUmnlEGxWtqo=
|
||||
github.com/pgsty/mc v0.0.0-20260913012246-4f609a4da3bb/go.mod h1:kJN7dsWtSUhXd2vNPXdjDi/or6lJKlSLfb56cBfMckA=
|
||||
github.com/pgsty/silo-console v0.0.0-20260916034812-56dfe455ac2f h1:ull9m/nXOEMfggKtMQP2idYXPh6l6chCUJ9hJHYyzDc=
|
||||
github.com/pgsty/silo-console v0.0.0-20260916034812-56dfe455ac2f/go.mod h1:YtRQZ6jYXRUE03oPA+IMGtflQN6nCvDWdKtroA7tfKo=
|
||||
github.com/pgsty/silo-pkg/v3 v3.14.0 h1:RCuVkzr6mdjbkV/Rv4OVO/XgdMBE0XYvUnT6GFuihFk=
|
||||
github.com/pgsty/silo-pkg/v3 v3.14.0/go.mod h1:c26IoMVITlP1+Sirl0AsHwoijlsSC9wPpyCuvhb3Axc=
|
||||
github.com/pgsty/mc v0.0.0-20260916070421-e952aa78f10a h1:FY1vcTQT67HGBe23ozlhOiQnWCzUCGmeSN5x6DT/kFs=
|
||||
github.com/pgsty/mc v0.0.0-20260916070421-e952aa78f10a/go.mod h1:6z9lX8kOvaWz2vLlP7BjjF3JWHnIH86XDA2HwcmpZH8=
|
||||
github.com/pgsty/silo-console v0.0.0-20260916075814-1360e26d976d h1:A9ou1GkUckGZd3BOgEXeZ393nbYST967SKcI4Lgr/Vs=
|
||||
github.com/pgsty/silo-console v0.0.0-20260916075814-1360e26d976d/go.mod h1:JS165CF2yoTBIx1X8k1IQ6vunJcZFG3pRZlWHJ8UT8o=
|
||||
github.com/pgsty/silo-pkg/v3 v3.14.1 h1:8MONT3Hky9EOnsuPk29590pEBE6MILaOnJtaejPPi7M=
|
||||
github.com/pgsty/silo-pkg/v3 v3.14.1/go.mod h1:jxKWxi52jSsDaDhXXPqjOyPDYdxx2iqdBXyLfM/xSIw=
|
||||
github.com/philhofer/fwd v1.2.0 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM=
|
||||
github.com/philhofer/fwd v1.2.0/go.mod h1:RqIHx9QI14HlwKwm98g9Re5prTQ6LdeRQn+gXJFxsJM=
|
||||
github.com/pierrec/lz4/v4 v4.1.29 h1:CDQY6qZOLI4DW0Nx6R1vRrifrCeQHnNXkMb0hZWXFjg=
|
||||
github.com/pierrec/lz4/v4 v4.1.29/go.mod h1:EoQMVJgeeEOMsCqCzqFm2O0cJvljX2nGZjcRIPL34O4=
|
||||
github.com/pierrec/lz4/v4 v4.1.30 h1:cchX8N2DVP668WkElI9QMwVyoNabLkq1LofDHFeIrdg=
|
||||
github.com/pierrec/lz4/v4 v4.1.30/go.mod h1:EoQMVJgeeEOMsCqCzqFm2O0cJvljX2nGZjcRIPL34O4=
|
||||
github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c h1:+mdjkGKdHQG3305AYmdv1U2eRNDiU2ErMBj1gwrq8eQ=
|
||||
github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c/go.mod h1:7rwL4CYBLnjLxUqIJNnCWiEdr3bn6IUYi15bNlnbCCU=
|
||||
github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
apiVersion: v1
|
||||
description: S3-Interface Libre Object Storage
|
||||
name: silo
|
||||
version: 7.0.2
|
||||
appVersion: RELEASE.2026-09-03T13-18-01Z
|
||||
version: 7.0.3
|
||||
appVersion: RELEASE.2026-09-16T00-00-00Z
|
||||
keywords:
|
||||
- silo
|
||||
- storage
|
||||
|
||||
@@ -15,7 +15,7 @@ clusterDomain: cluster.local
|
||||
##
|
||||
image:
|
||||
repository: pgsty/silo
|
||||
tag: RELEASE.2026-09-03T13-18-01Z
|
||||
tag: RELEASE.2026-09-16T00-00-00Z
|
||||
pullPolicy: IfNotPresent
|
||||
|
||||
imagePullSecrets: []
|
||||
@@ -25,7 +25,7 @@ imagePullSecrets: []
|
||||
##
|
||||
mcImage:
|
||||
repository: pgsty/mc
|
||||
tag: RELEASE.2026-09-13T00-00-00Z
|
||||
tag: RELEASE.2026-09-16T00-00-00Z
|
||||
pullPolicy: IfNotPresent
|
||||
|
||||
## Silo mode, i.e. standalone or distributed.
|
||||
|
||||
@@ -135,7 +135,6 @@ var (
|
||||
Key: apiStaleUploadsExpiry,
|
||||
Value: "24h",
|
||||
},
|
||||
config.KV{Key: apiMultipartListing, Value: "strict"},
|
||||
config.KV{
|
||||
Key: apiDeleteCleanupInterval,
|
||||
Value: "5m",
|
||||
@@ -210,6 +209,7 @@ func LookupConfig(kvs config.KVS) (cfg Config, err error) {
|
||||
apiReplicationWorkers,
|
||||
apiReplicationFailedWorkers,
|
||||
"expiry_workers",
|
||||
apiMultipartListing, // Ignore the retired shared-config key during migration.
|
||||
}
|
||||
|
||||
disableODirect := env.Get(EnvAPIDisableODirect, kvs.Get(apiDisableODirect)) == config.EnableOn
|
||||
@@ -324,10 +324,6 @@ func LookupConfig(kvs config.KVS) (cfg Config, err error) {
|
||||
return cfg, err
|
||||
}
|
||||
cfg.StaleUploadsExpiry = staleUploadsExpiry
|
||||
cfg.MultipartListing = env.Get(EnvAPIMultipartListing, kvs.GetWithDefault(apiMultipartListing, DefaultKVS))
|
||||
if cfg.MultipartListing != "strict" && cfg.MultipartListing != "legacy" {
|
||||
return cfg, fmt.Errorf("%s must be strict or legacy", apiMultipartListing)
|
||||
}
|
||||
|
||||
cfg.SyncEvents = env.Get(EnvAPISyncEvents, kvs.Get(apiSyncEvents)) == config.EnableOn
|
||||
|
||||
@@ -348,5 +344,11 @@ func LookupConfig(kvs config.KVS) (cfg Config, err error) {
|
||||
cfg.ObjectMaxVersions = math.MaxInt64
|
||||
}
|
||||
|
||||
cfg.MultipartListing = env.Get(EnvAPIMultipartListing, "legacy")
|
||||
if cfg.MultipartListing != "strict" && cfg.MultipartListing != "legacy" {
|
||||
cfg.MultipartListing = "legacy"
|
||||
return cfg, fmt.Errorf("%s must be strict or legacy; using legacy", EnvAPIMultipartListing)
|
||||
}
|
||||
|
||||
return cfg, nil
|
||||
}
|
||||
|
||||
@@ -4,38 +4,111 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/minio/minio/internal/config"
|
||||
)
|
||||
|
||||
func TestMultipartListingMigrationMode(t *testing.T) {
|
||||
t.Setenv(EnvAPIRequestsMax, "")
|
||||
t.Setenv(EnvAPISyncEvents, "")
|
||||
t.Setenv(EnvAPIObjectMaxVersions, "")
|
||||
for _, tc := range []struct {
|
||||
name, stored, override, want string
|
||||
invalid bool
|
||||
}{
|
||||
{name: "default", want: "strict"},
|
||||
{name: "explicit-migration", stored: "legacy", want: "legacy"},
|
||||
{name: "environment-override", stored: "strict", override: "legacy", want: "legacy"},
|
||||
{name: "invalid-stored", stored: "automatic", invalid: true},
|
||||
{name: "invalid-environment", stored: "strict", override: "automatic", invalid: true},
|
||||
{name: "default", want: "legacy"},
|
||||
{name: "explicit-legacy", override: "legacy", want: "legacy"},
|
||||
{name: "explicit-strict", override: "strict", want: "strict"},
|
||||
{name: "retired-stored-strict", stored: "strict", want: "legacy"},
|
||||
{name: "retired-stored-legacy", stored: "legacy", want: "legacy"},
|
||||
{name: "retired-stored-invalid", stored: "automatic", want: "legacy"},
|
||||
{name: "environment-override", stored: "legacy", override: "strict", want: "strict"},
|
||||
{name: "invalid-environment", stored: "strict", override: "automatic", want: "legacy", invalid: true},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Setenv(EnvAPIMultipartListing, tc.override)
|
||||
var kvs config.KVS
|
||||
kvs := DefaultKVS.Clone()
|
||||
kvs.Set(apiRequestsMax, "17")
|
||||
kvs.Set(apiSyncEvents, config.EnableOn)
|
||||
kvs.Set(apiObjectMaxVersions, "71")
|
||||
if tc.stored != "" {
|
||||
kvs = config.KVS{{Key: apiMultipartListing, Value: tc.stored}}
|
||||
kvs.Set(apiMultipartListing, tc.stored)
|
||||
}
|
||||
cfg, err := LookupConfig(kvs)
|
||||
if tc.invalid {
|
||||
if err == nil {
|
||||
t.Fatal("invalid migration mode silently accepted")
|
||||
}
|
||||
return
|
||||
if (err != nil) != tc.invalid || cfg.MultipartListing != tc.want {
|
||||
t.Fatalf("mode=%q err=%v, want %q invalid=%v", cfg.MultipartListing, err, tc.want, tc.invalid)
|
||||
}
|
||||
if err != nil || cfg.MultipartListing != tc.want {
|
||||
t.Fatalf("mode=%q err=%v, want %q", cfg.MultipartListing, err, tc.want)
|
||||
if cfg.RequestsMax != 17 || !cfg.SyncEvents || cfg.ObjectMaxVersions != 71 {
|
||||
t.Fatalf("migration setting discarded unrelated config: %+v", cfg)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultipartListingConfigPersistence(t *testing.T) {
|
||||
t.Setenv(EnvAPIRequestsMax, "")
|
||||
t.Setenv(EnvAPISyncEvents, "")
|
||||
t.Setenv(EnvAPIObjectMaxVersions, "")
|
||||
t.Setenv(EnvAPIMultipartListing, "strict")
|
||||
previous := config.DefaultKVS
|
||||
config.DefaultKVS = map[string]config.KVS{config.APISubSys: DefaultKVS}
|
||||
t.Cleanup(func() { config.DefaultKVS = previous })
|
||||
for _, help := range Help {
|
||||
if help.Key == apiMultipartListing {
|
||||
t.Fatal("process-only setting advertised as shared config")
|
||||
}
|
||||
}
|
||||
for _, input := range []string{
|
||||
"",
|
||||
"api requests_max=17 sync_events=on object_max_versions=71",
|
||||
`api requests_max=17 comment="later disable multipart_listing=strict"`,
|
||||
} {
|
||||
t.Run(input, func(t *testing.T) {
|
||||
c := config.New().Clone()
|
||||
if _, err := c.ReadConfig(strings.NewReader(input)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Exercise the JSON save/load and Merge paths used for shared config.
|
||||
data, err := json.Marshal(c)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var loaded config.Config
|
||||
if err = json.Unmarshal(data, &loaded); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
kvs := loaded.Merge()[config.APISubSys][config.Default]
|
||||
if _, present := kvs.Lookup(apiMultipartListing); present {
|
||||
t.Fatal("process-only setting persisted in shared config")
|
||||
}
|
||||
cfg, err := LookupConfig(kvs)
|
||||
if err != nil || cfg.MultipartListing != "strict" {
|
||||
t.Fatalf("environment mode lost: %+v %v", cfg, err)
|
||||
}
|
||||
if input != "" && cfg.RequestsMax != 17 {
|
||||
t.Fatalf("unrelated setting lost: %+v", cfg)
|
||||
}
|
||||
if strings.Contains(input, "comment=") && kvs.Get(config.Comment) != "later disable multipart_listing=strict" {
|
||||
t.Fatalf("literal comment changed: %q", kvs.Get(config.Comment))
|
||||
}
|
||||
})
|
||||
}
|
||||
// A key persisted by a previous development build is tolerated, retained for
|
||||
// explicit cleanup, and removable without resetting other API settings.
|
||||
c := config.New().Clone()
|
||||
kvs := c[config.APISubSys][config.Default]
|
||||
kvs.Set(apiMultipartListing, "strict")
|
||||
kvs.Set(apiRequestsMax, "17")
|
||||
c[config.APISubSys][config.Default] = kvs
|
||||
c = c.Merge()
|
||||
if err := c.DelKVS("api multipart_listing"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
kvs = c.Merge()[config.APISubSys][config.Default]
|
||||
if _, present := kvs.Lookup(apiMultipartListing); present || kvs.Get(apiRequestsMax) != "17" {
|
||||
t.Fatalf("targeted reset changed unrelated config or restored retired key: %v", kvs)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -26,12 +26,6 @@ var (
|
||||
|
||||
// Help holds configuration keys and their default values for api subsystem.
|
||||
Help = config.HelpKVS{
|
||||
config.HelpKV{
|
||||
Key: apiMultipartListing,
|
||||
Description: "multipart listing mode: strict, or temporary legacy mode during coordinated upgrade" + defaultHelpPostfix(apiMultipartListing),
|
||||
Type: "string",
|
||||
Optional: true,
|
||||
},
|
||||
config.HelpKV{
|
||||
Key: apiRequestsMax,
|
||||
Description: `set the maximum number of concurrent requests (default: auto)`,
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
// Copyright (c) 2026 Feng Ruohang
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License as published by
|
||||
// the Free Software Foundation, either version 3 of the License, or
|
||||
// (at your option) any later version.
|
||||
//
|
||||
// This program is distributed in the hope that it will be useful,
|
||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
// GNU Affero General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Affero General Public License
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
package s3select
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"io"
|
||||
"testing"
|
||||
"testing/iotest"
|
||||
|
||||
"github.com/pierrec/lz4/v4"
|
||||
)
|
||||
|
||||
func TestLZ4Input(t *testing.T) {
|
||||
request := []byte(`<SelectObjectContentRequest>
|
||||
<Expression>SELECT * FROM S3Object</Expression><ExpressionType>SQL</ExpressionType>
|
||||
<InputSerialization><CompressionType>LZ4</CompressionType><CSV><FileHeaderInfo>NONE</FileHeaderInfo></CSV></InputSerialization>
|
||||
<OutputSerialization><CSV/></OutputSerialization></SelectObjectContentRequest>`)
|
||||
for _, input := range []string{"", "one,1\ntwo,2\n"} {
|
||||
t.Run(input, func(t *testing.T) {
|
||||
var compressed bytes.Buffer
|
||||
writer := lz4.NewWriter(&compressed)
|
||||
if _, err := io.WriteString(writer, input); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := writer.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got, err := evaluateSelectForTest(t, request, compressed.Bytes())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if string(got) != input {
|
||||
t.Fatalf("result=%q, want %q", got, input)
|
||||
}
|
||||
|
||||
reader, err := newProgressReader(io.NopCloser(iotest.OneByteReader(bytes.NewReader(compressed.Bytes()))), lz4Type)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = reader.Close() })
|
||||
got, err = io.ReadAll(reader)
|
||||
if err != nil || string(got) != input {
|
||||
t.Fatalf("fragmented input: result=%q, err=%v", got, err)
|
||||
}
|
||||
scanned, processed := reader.Stats()
|
||||
if scanned != int64(compressed.Len()) || processed != int64(len(input)) {
|
||||
t.Fatalf("scanned=%d processed=%d", scanned, processed)
|
||||
}
|
||||
if err := reader.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := reader.Read(make([]byte, 1)); err == nil {
|
||||
t.Fatal("read after Close succeeded")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user