Compare commits

...

17 Commits

Author SHA1 Message Date
Feng Ruohang 2a4d51406b test: build the purge-response fixture ARN at runtime
The rebrand compatibility guard tracks every literal policy value that
names the upstream brand. Construct the replication ARN the way the
other replication fixtures do instead of adding a new literal.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 23:50:09 +08:00
Feng Ruohang 0596685ae7 docs: record the delete-marker purge, listing quorum and migration tag repairs
Describe the three correctness repairs merged after the 20260916 candidate
was prepared, together with their known remaining limits (#217, #218).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 23:44:44 +08:00
Feng Ruohang 358ab38fb0 fix: report a purged version as a delete marker only for stored markers
The exact-version purge path copies the looked-up DeleteMarker flag
into the DELETE response so that removing a delete marker keeps its
x-amz-delete-marker header. That flag is also set for a data version
whose purge is pending, because the lookup exposes such a version as
deleted for visibility. Retrying the purge after the replication
configuration was removed therefore answered with
x-amz-delete-marker: true and raised ObjectRemovedDeleteMarkerCreated
for a data version.

Only a stored marker carries no erasure layout; use that to decide the
response identity. Found by the adversarial review of the purge change.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 22:53:12 +08:00
Feng Ruohang 254b19ac07 fix(replication): recheck queued delete-marker creations under the lock
A delete-marker 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. Tasks from several frontends serialize on the per-object
replication lock, so a task can run after the user has purged that
marker and the purge has already reached the targets. The target no
longer holds the marker, so the queued creation recreated it there with
the original VersionID and mtime. This is the late-create sequence seen
in the 2026-09-16 three-site runs: purge 204 on both targets, then about
85 ms later a replicated creation for the same VersionID.

Re-read the source version under the replication lock before sending a
creation. A missing version, a non-marker version, or a version under
purge makes the task stale; it is dropped without touching the source.
A read that cannot confirm either way is retried through MRF instead of
being treated as absence.

Creations already on the wire, replays from other sites, and cleanup of
minority residue after a crash are not covered by this check.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 22:35:08 +08:00
Feng Ruohang eb4f5e5b31 Fix exact-version purges and delete-marker metadata healing
Require write-quorum absence for missing purge retries, aggregate removed and already-absent votes only for physical purges, and resolve receiver purges by version across pools. Preserve marker metadata through creation and healing while retaining DELETE response semantics.

Add real-disk quorum, pool, callback, metadata, heal and outbound replication regressions. Synchronize capacity fixture installation and restoration with concurrent IAM readers.

Signed-off-by: Feng Ruohang <rh@vonng.com>
(cherry picked from commit 22ba4d426677aa07470927fc349a43afd87e19b6)
2026-09-16 22:24:27 +08:00
Feng Ruohang 8d06424b12 fix: retain quorate null versions after newer minorities
Recount only when the existing selection is below quorum and each original
nonempty stream has one ordinary null version with identical EC headers.
Reuse the existing header grouping and leave history pruning unchanged.

Add deterministic signed LIST, complete-metadata permutation, excluded-history,
object-read, scanner/heal, Walk and migration-consumer regressions.

Signed-off-by: Feng Ruohang <rh@vonng.com>
(cherry picked from commit 9804e4deb4d5a18cca640be72bff9b47a415f7de)
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 21:50:03 +08:00
Feng Ruohang fced863036 fix(storage): preserve tags during pool migration
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 20:22:37 +08:00
Feng Ruohang 5af0865aab Merge pull request #214 from pgsty/codex/release-prep-20260916
Prepare the SILO 20260916000000 source candidate with released component pins, CPU metrics synchronization, regression coverage and release tooling updates.

All 11 checks passed on e7963381ba. Signed release artifacts, image publication and deployment remain separate deliverables tracked by #203.

Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 19:30:07 +08:00
Feng Ruohang 83821f0f1f Merge pull request #216 from pgsty/codex/docs-published-design-links-20260916
docs: link retained guides to published design records
2026-09-16 18:50:28 +08:00
Feng Ruohang bf053386a4 docs: link compression and federation to maintained design records
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 18:26:03 +08:00
Feng Ruohang e7963381ba test: synchronize tag replication capacity fixtures
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 17:37:47 +08:00
Feng Ruohang 37c0edc7ca test: use canonical octal permissions in LZ4 fixture
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 17:33:58 +08:00
Feng Ruohang a2fe70424e release: prepare SILO 20260916000000 candidate
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 17:30:29 +08:00
Feng Ruohang 027d43b4eb fix: synchronize CPU resource metrics reads
Hold resourceMetricsMapMu.RLock while the v3 CPU collector reads both the
subsystem map and its nested idle/iowait metrics. The unsynchronized reads
were inherited from upstream commit f7b665347 (minio/minio#19560).

Add value and concurrent Prometheus Gather regression tests, run them in
the targeted race CI job, and record the existing CPU metric names in the
compatibility inventory.

Fixes #210

Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 17:09:19 +08:00
Feng Ruohang f99ed829b5 Merge pull request #213 from pgsty/codex/pr198-release-compat
fix: preserve multipart defaults and rollback compatibility
2026-09-16 16:42:56 +08:00
Feng Ruohang 956a1a8e33 Merge pull request #212 from pgsty/codex/docs-migration-cleanup
docs: consolidate repository entry points and retire migrated work records
2026-09-16 16:39:40 +08:00
Feng Ruohang 82f0a9828e fix: preserve multipart defaults and rollback compatibility
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 16:30:43 +08:00
46 changed files with 3825 additions and 172 deletions
+6
View File
@@ -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
+2 -2
View File
@@ -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 }}
+6 -4
View File
@@ -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
View File
@@ -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
+3 -3
View File
@@ -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 \
+1 -1
View File
@@ -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": [
+1 -1
View File
@@ -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
+64
View File
@@ -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])
}
}})
}
+524
View File
@@ -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())
}
}
+23
View File
@@ -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 {
+331
View File
@@ -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
View File
@@ -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
View File
@@ -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
}
+20
View File
@@ -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.
+2 -2
View File
@@ -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.
+33
View File
@@ -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
}
+592
View File
@@ -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)
}
})
}
}
}
+2 -2
View File
@@ -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.
+28 -7
View File
@@ -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
+87
View File
@@ -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
View File
@@ -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 {
+3 -9
View File
@@ -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 {
+433
View File
@@ -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)
}
+2
View File
@@ -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 {
+185
View File
@@ -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())
})
}
}
}
+41 -4
View File
@@ -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()
+264
View File
@@ -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{}
}
+40 -4
View File
@@ -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)
+101
View File
@@ -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)
}
}
}
+67
View File
@@ -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
}
+274
View File
@@ -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
View File
@@ -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
View File
@@ -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.
+2
View File
@@ -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)
+2
View File
@@ -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
+7 -7
View File
@@ -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
+14 -14
View File
@@ -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=
+2 -2
View File
@@ -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
+2 -2
View File
@@ -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.
+7 -5
View File
@@ -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
}
+87 -14
View File
@@ -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)
}
}
-6
View File
@@ -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)`,
+71
View File
@@ -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")
}
})
}
}