mirror of
https://github.com/pgsty/minio.git
synced 2026-10-04 08:45:59 +03:00
Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e7963381ba | |||
| 37c0edc7ca | |||
| a2fe70424e | |||
| 027d43b4eb | |||
| f99ed829b5 | |||
| 956a1a8e33 | |||
| 82f0a9828e |
@@ -96,6 +96,12 @@ jobs:
|
|||||||
- name: Run multipart listing and cancellation tests under race detector
|
- 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
|
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:
|
crosscompile:
|
||||||
name: Cross Compile
|
name: Cross Compile
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
|
|||||||
@@ -132,9 +132,9 @@ jobs:
|
|||||||
cosign-release: v3.1.2
|
cosign-release: v3.1.2
|
||||||
|
|
||||||
- name: Build Draft release with GoReleaser
|
- name: Build Draft release with GoReleaser
|
||||||
uses: goreleaser/goreleaser-action@v7
|
uses: goreleaser/goreleaser-action@f06c13b6b1a9625abc9e6e439d9c05a8f2190e94 # v7.2.3
|
||||||
with:
|
with:
|
||||||
version: "~> v2"
|
version: v2.18.1
|
||||||
args: release --clean --skip=validate --config .github/goreleaser.yml
|
args: release --clean --skip=validate --config .github/goreleaser.yml
|
||||||
env:
|
env:
|
||||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||||
|
|||||||
@@ -4,6 +4,8 @@ on:
|
|||||||
workflow_dispatch:
|
workflow_dispatch:
|
||||||
pull_request:
|
pull_request:
|
||||||
paths:
|
paths:
|
||||||
|
- "go.mod"
|
||||||
|
- "go.sum"
|
||||||
- ".github/goreleaser.yml"
|
- ".github/goreleaser.yml"
|
||||||
- ".github/nfpm.yml"
|
- ".github/nfpm.yml"
|
||||||
- "Dockerfile.goreleaser"
|
- "Dockerfile.goreleaser"
|
||||||
@@ -93,9 +95,9 @@ jobs:
|
|||||||
echo "LDFLAGS: ${LDFLAGS}"
|
echo "LDFLAGS: ${LDFLAGS}"
|
||||||
|
|
||||||
- name: GoReleaser config check
|
- name: GoReleaser config check
|
||||||
uses: goreleaser/goreleaser-action@v7
|
uses: goreleaser/goreleaser-action@f06c13b6b1a9625abc9e6e439d9c05a8f2190e94 # v7.2.3
|
||||||
with:
|
with:
|
||||||
version: "~> v2"
|
version: v2.18.1
|
||||||
args: check --config .github/goreleaser.yml
|
args: check --config .github/goreleaser.yml
|
||||||
|
|
||||||
- name: Validate Helm chart and legacy upgrade identity
|
- name: Validate Helm chart and legacy upgrade identity
|
||||||
@@ -107,9 +109,9 @@ jobs:
|
|||||||
syft-version: v1.50.0
|
syft-version: v1.50.0
|
||||||
|
|
||||||
- name: Build snapshot artifacts
|
- name: Build snapshot artifacts
|
||||||
uses: goreleaser/goreleaser-action@v7
|
uses: goreleaser/goreleaser-action@f06c13b6b1a9625abc9e6e439d9c05a8f2190e94 # v7.2.3
|
||||||
with:
|
with:
|
||||||
version: "~> v2"
|
version: v2.18.1
|
||||||
# A pull-request snapshot has no trusted release identity. Exercise
|
# A pull-request snapshot has no trusted release identity. Exercise
|
||||||
# the SBOM/checksum pipeline here, and reserve keyless signing for
|
# the SBOM/checksum pipeline here, and reserve keyless signing for
|
||||||
# the tag-triggered release workflow with GitHub OIDC.
|
# the tag-triggered release workflow with GitHub OIDC.
|
||||||
|
|||||||
+44
-25
@@ -2,13 +2,18 @@
|
|||||||
|
|
||||||
## Unreleased
|
## 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
|
**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/)
|
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).
|
and [complete commit range](https://github.com/pgsty/silo/compare/RELEASE.2026-09-03T13-18-01Z...main).
|
||||||
|
|
||||||
### Authorization and security
|
### 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
|
- Restrict embedded Console's anonymous sharing proxy to object-content GETs
|
||||||
at the configured S3 origin, and reject every redirect. Internal metrics,
|
at the configured S3 origin, and reject every redirect. Internal metrics,
|
||||||
system paths and non-download S3 operations cannot be reached through it.
|
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
|
quorum-written `xl.meta`; completion removes those upload-only fields. Native
|
||||||
markers remain usable after their upload is completed or canceled. Strict
|
markers remain usable after their upload is completed or canceled. Strict
|
||||||
listing returns a diagnostic 503 for legacy uploads or uncertain coverage;
|
listing returns a diagnostic 503 for legacy uploads or uncertain coverage;
|
||||||
`api multipart_listing=legacy` is an explicit temporary migration mode.
|
the default remains the released exact-key/cache-based `legacy` behavior.
|
||||||
Upgrade every writer, drain old uploads and check the read-only admin
|
Opt into strict mode only through `MINIO_API_MULTIPART_LISTING=strict`, after
|
||||||
`multipart-preflight` report before relying on strict listing. Per-process
|
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;
|
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)
|
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/).
|
and its [design record](https://silo.pgsty.com/blog/design/list-multipart-uploads/).
|
||||||
Thanks to mr javad seydi (@mrjavadseydi) for the original implementation.
|
Thanks to mr javad seydi (@mrjavadseydi) for the original implementation.
|
||||||
- Confirm multipart cancellation on a strict majority of each relevant set,
|
- Retain released read-quorum and best-effort multipart cancellation in default
|
||||||
and allow retries after partial deletion. Uncertain pools or insufficient
|
legacy mode. Strict mode requires majority deletion acknowledgements and
|
||||||
confirmations return 503 rather than acknowledging a cancellation whose
|
permits retries below read quorum. When most drives were already empty,
|
||||||
static remnants can later become readable. **Known boundary:** creation
|
failed deletion of an observed remnant now returns 503 instead of being
|
||||||
writes that finish after a storage timeout can still restore an upload after
|
masked by empty-drive successes. Wrong-key or wrong-bucket cancellation
|
||||||
successful cancellation; this change does not add a durable creation fence.
|
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
|
- Preserve object tags during multi-pool metadata reconciliation by reading the
|
||||||
resolved tag field together with its revision (#189). Previously, reconciliation
|
resolved tag field together with its revision (#189). Previously, reconciliation
|
||||||
could replace existing tags with an empty value.
|
could replace existing tags with an empty value.
|
||||||
@@ -158,23 +172,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
|
- Restore embedded Console login over loopback TLS, trusted-proxy handling and
|
||||||
all four WebSocket connection limits. Preserve Go TLS defaults across transports.
|
all four WebSocket connection limits. Preserve Go TLS defaults across transports.
|
||||||
- Directly require `github.com/pgsty/silo-pkg/v3` v3.14.0; select Console
|
- Directly require `github.com/pgsty/silo-pkg/v3` v3.14.1; select released Console
|
||||||
`v0.0.0-20260916034812-56dfe455ac2f` and MC
|
v2.4.1 (`v0.0.0-20260916075814-1360e26d976d`) and mcli 20260916
|
||||||
`v0.0.0-20260913012246-4f609a4da3bb` with explicit PGSTY replacements.
|
(`v0.0.0-20260916070421-e952aa78f10a`) with explicit PGSTY replacements.
|
||||||
- Pin upstream minio-go `v7.3.1-0.20260910142817-60bd07042d49`; refresh Go x/*
|
The embedded frontend identifies itself as Console v2.4.1.
|
||||||
modules and security fixes including bounded AMQP frame handling. Keep Go
|
- Pin upstream minio-go `v7.3.1-0.20260915093545-32e1f32cb176` to handle
|
||||||
1.27.1 and go-systemd v22.6.0's NetBSD compatibility replacement.
|
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
|
- 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
|
source for both Linux architectures. Pin the published mcli 20260916 archives
|
||||||
hashes. Helm's client image follows that release; its Server image still names
|
and hashes in the container and update the client installer default.
|
||||||
the latest published Server 20260903.
|
- 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
|
Validation of earlier source revisions does not establish acceptance of this
|
||||||
Release workflows; native curl builds passed on both architectures. A local
|
candidate. Final source, package, image and multi-process checks are tracked
|
||||||
ARM64 image passed startup, health, S3 transfer and embedded Console checks.
|
separately in [#203](https://github.com/pgsty/silo/issues/203). No Server release
|
||||||
These checks do not publish a Server tag or production image and do not replace
|
or production rollout is implied by this preparation target.
|
||||||
cluster upgrade/rollback acceptance for the next release. Dated investigations
|
|
||||||
retain the exact source and runtime boundaries they tested.
|
|
||||||
|
|
||||||
## RELEASE.2026-09-03T13-18-01Z
|
## RELEASE.2026-09-03T13-18-01Z
|
||||||
|
|
||||||
|
|||||||
@@ -17,9 +17,9 @@ ENV GOPATH=/go
|
|||||||
ENV CGO_ENABLED=0
|
ENV CGO_ENABLED=0
|
||||||
|
|
||||||
ARG MC_REPO=pgsty/mc
|
ARG MC_REPO=pgsty/mc
|
||||||
ARG MC_VERSION=RELEASE.2026-09-13T00-00-00Z
|
ARG MC_VERSION=RELEASE.2026-09-16T00-00-00Z
|
||||||
ARG MC_AMD64_SHA256=9d2a92de9c7b887d9b944fe9ddce68d23f1b6df3415092e737594e56593f5e2b
|
ARG MC_AMD64_SHA256=4ba2814fd5507fbe6b4d237c359750b9119d28d7217495fa5b48002fcbd397ef
|
||||||
ARG MC_ARM64_SHA256=3d82e9ea6c601c4cb44fe5dd5f2ad1b7d7d64369378110f9ada9c524688a452a
|
ARG MC_ARM64_SHA256=b7008ca2a1bc5735b6585981c59a3640a0daa152dc789df3d1a0d0438787de82
|
||||||
|
|
||||||
RUN apk add -U --no-cache \
|
RUN apk add -U --no-cache \
|
||||||
ca-certificates \
|
ca-certificates \
|
||||||
|
|||||||
@@ -40,7 +40,7 @@ if [ -n "${MCLI_BIN:-}" ]; then
|
|||||||
exit 0
|
exit 0
|
||||||
fi
|
fi
|
||||||
|
|
||||||
release=${MCLI_RELEASE:-RELEASE.2026-09-13T00-00-00Z}
|
release=${MCLI_RELEASE:-RELEASE.2026-09-16T00-00-00Z}
|
||||||
version_hyphen=${release#RELEASE.}
|
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/')
|
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
|
if [ "${package_version}" = "${version_hyphen}" ]; then
|
||||||
|
|||||||
@@ -571,6 +571,14 @@
|
|||||||
"minio_resource_stats",
|
"minio_resource_stats",
|
||||||
"minio_s3",
|
"minio_s3",
|
||||||
"minio_stats",
|
"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"
|
"minio_test"
|
||||||
],
|
],
|
||||||
"headers": [
|
"headers": [
|
||||||
|
|||||||
@@ -113,7 +113,7 @@ helm_run template my-release "${new_chart}" \
|
|||||||
go run ./buildscripts/helm-migration-guard "${old_render}" "${new_render}"
|
go run ./buildscripts/helm-migration-guard "${old_render}" "${new_render}"
|
||||||
|
|
||||||
helm_run package "${new_chart}" --destination "${output_dir}" >/dev/null
|
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
|
if find "${work_dir}" -maxdepth 1 -type f -name 'minio-*.tgz' | grep -q .; then
|
||||||
echo "Helm packaging emitted a legacy MinIO chart name" >&2
|
echo "Helm packaging emitted a legacy MinIO chart name" >&2
|
||||||
exit 1
|
exit 1
|
||||||
|
|||||||
@@ -31,6 +31,7 @@ func TestMultipartListingAbortBetweenHTTPPages(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
setMultipartListingTestMode(t, false)
|
||||||
var firstID, secondID string
|
var firstID, secondID string
|
||||||
for attempt := 0; attempt < 32; attempt++ {
|
for attempt := 0; attempt < 32; attempt++ {
|
||||||
one, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
one, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
||||||
@@ -147,6 +148,22 @@ func TestMultipartListingLegacyPreflight(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
z.mpCache.Clear()
|
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)
|
_, err = z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 10)
|
||||||
if !errors.Is(err, errMultipartListingLegacy) {
|
if !errors.Is(err, errMultipartListingLegacy) {
|
||||||
t.Fatalf("old upload in another bucket: %v", err)
|
t.Fatalf("old upload in another bucket: %v", err)
|
||||||
@@ -214,6 +231,7 @@ func TestMultipartListingIdentityFallback(t *testing.T) {
|
|||||||
|
|
||||||
func TestMultipartAbortPoolsAndRetry(t *testing.T) {
|
func TestMultipartAbortPoolsAndRetry(t *testing.T) {
|
||||||
z, bucket := consistencyPools(t)
|
z, bucket := consistencyPools(t)
|
||||||
|
setMultipartListingTestMode(t, false)
|
||||||
mp, err := z.serverPools[1].NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
mp, err := z.serverPools[1].NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
@@ -286,6 +304,7 @@ func TestMultipartListingMarkerHTTP(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
setMultipartListingTestMode(t, false)
|
||||||
for _, tc := range []struct {
|
for _, tc := range []struct {
|
||||||
key, marker string
|
key, marker string
|
||||||
status int
|
status int
|
||||||
@@ -369,6 +388,7 @@ func TestMultipartAbortLateCreateBoundary(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
z := obj.(*erasureServerPools)
|
z := obj.(*erasureServerPools)
|
||||||
|
setMultipartListingTestMode(t, false)
|
||||||
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
|
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
|
||||||
saved := globalStorageClass
|
saved := globalStorageClass
|
||||||
globalStorageClass.Update(storageclass.Config{Standard: storageclass.StorageClass{Parity: 8}})
|
globalStorageClass.Update(storageclass.Config{Standard: storageclass.StorageClass{Parity: 8}})
|
||||||
@@ -454,6 +474,7 @@ func TestMultipartListingScanCosts(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
z := obj.(*erasureServerPools)
|
z := obj.(*erasureServerPools)
|
||||||
|
setMultipartListingTestMode(t, false)
|
||||||
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
|
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
|
||||||
const bucket, otherBucket = "r9-scan-target", "r9-scan-unrelated"
|
const bucket, otherBucket = "r9-scan-target", "r9-scan-unrelated"
|
||||||
for _, name := range []string{bucket, otherBucket} {
|
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 {
|
if err := z.MakeBucket(t.Context(), bucket, MakeBucketOptions{}); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
setMultipartListingTestMode(t, false)
|
||||||
return z, z.serverPools[0].getHashedSet("a"), bucket
|
return z, z.serverPools[0].getHashedSet("a"), bucket
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -673,6 +695,7 @@ func TestMultipartListingAdmissionHTTP(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
setMultipartListingTestMode(t, false)
|
||||||
for range cap(multipartScanSlots) {
|
for range cap(multipartScanSlots) {
|
||||||
scan, err := startMultipartScan(t.Context(), false)
|
scan, err := startMultipartScan(t.Context(), false)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -0,0 +1,331 @@
|
|||||||
|
// Copyright (c) 2026 Ruohang Feng
|
||||||
|
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||||
|
|
||||||
|
package cmd
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/base64"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"strings"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/minio/minio/internal/config/api"
|
||||||
|
"github.com/minio/minio/internal/config/storageclass"
|
||||||
|
"github.com/minio/minio/internal/dsync"
|
||||||
|
"github.com/minio/minio/internal/grid"
|
||||||
|
xnet "github.com/pgsty/silo-pkg/v3/net"
|
||||||
|
)
|
||||||
|
|
||||||
|
func setMultipartListingTestMode(t *testing.T, legacy bool) {
|
||||||
|
t.Helper()
|
||||||
|
globalAPIConfig.mu.Lock()
|
||||||
|
previous := globalAPIConfig.multipartListingStrict
|
||||||
|
globalAPIConfig.multipartListingStrict = !legacy
|
||||||
|
globalAPIConfig.mu.Unlock()
|
||||||
|
t.Cleanup(func() {
|
||||||
|
globalAPIConfig.mu.Lock()
|
||||||
|
globalAPIConfig.multipartListingStrict = previous
|
||||||
|
globalAPIConfig.mu.Unlock()
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestMultipartListingDefaultMode(t *testing.T) {
|
||||||
|
var uninitialized apiConfig
|
||||||
|
if !uninitialized.getMultipartListingLegacy() {
|
||||||
|
t.Fatal("runtime zero value enabled strict mode before config initialization")
|
||||||
|
}
|
||||||
|
for _, mode := range []string{"", "legacy", "strict"} {
|
||||||
|
var local apiConfig
|
||||||
|
local.init(api.Config{RequestsMax: 1, MultipartListing: mode}, []int{4}, false)
|
||||||
|
if local.getMultipartListingLegacy() != (mode != "strict") {
|
||||||
|
t.Fatalf("unexpected mode for %q", mode)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestMultipartAbortLegacyAvailability(t *testing.T) {
|
||||||
|
for _, drives := range []int{4, 16} {
|
||||||
|
for _, offline := range []int{0, 1, drives / 2} {
|
||||||
|
t.Run(fmt.Sprintf("drives=%d/offline=%d", drives, offline), func(t *testing.T) {
|
||||||
|
obj, dirs, err := prepareErasure(t.Context(), drives)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
z := obj.(*erasureServerPools)
|
||||||
|
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
|
||||||
|
setMultipartListingTestMode(t, true)
|
||||||
|
previous := globalStorageClass
|
||||||
|
globalStorageClass.Update(storageclass.Config{Standard: storageclass.StorageClass{Parity: drives / 2}})
|
||||||
|
t.Cleanup(func() { globalStorageClass.Update(previous) })
|
||||||
|
const bucket, key = "multipart-legacy-availability", "keep-available"
|
||||||
|
if err := z.MakeBucket(t.Context(), bucket, MakeBucketOptions{}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
mp, err := z.NewMultipartUpload(t.Context(), bucket, key, ObjectOptions{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
set := z.serverPools[0].getHashedSet(key)
|
||||||
|
original := set.getDisks
|
||||||
|
disks := original()
|
||||||
|
set.getDisks = func() []StorageAPI {
|
||||||
|
visible := append([]StorageAPI(nil), disks...)
|
||||||
|
clear(visible[:offline])
|
||||||
|
return visible
|
||||||
|
}
|
||||||
|
t.Cleanup(func() { set.getDisks = original })
|
||||||
|
if err := z.AbortMultipartUpload(t.Context(), bucket, key, mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatalf("released read-quorum availability regressed: %v", err)
|
||||||
|
}
|
||||||
|
for _, disk := range disks[offline:] {
|
||||||
|
_, err := disk.ReadVersion(t.Context(), bucket, minioMetaMultipartBucket, set.getUploadIDDir(bucket, key, mp.UploadID), "", ReadOptions{})
|
||||||
|
if !errors.Is(err, errFileNotFound) {
|
||||||
|
t.Fatalf("online replica not cleaned: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestMultipartAbortLegacyPoolOrder(t *testing.T) {
|
||||||
|
for _, owner := range []int{0, 1} {
|
||||||
|
t.Run(fmt.Sprintf("owner=%d", owner), func(t *testing.T) {
|
||||||
|
z, bucket := consistencyPools(t)
|
||||||
|
setMultipartListingTestMode(t, true)
|
||||||
|
mp, err := z.serverPools[owner].NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
emptySet := z.serverPools[1-owner].getHashedSet("a")
|
||||||
|
original := emptySet.getDisks
|
||||||
|
t.Cleanup(func() { emptySet.getDisks = original })
|
||||||
|
emptySet.getDisks = func() []StorageAPI { return make([]StorageAPI, emptySet.setDriveCount) }
|
||||||
|
err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{})
|
||||||
|
if owner == 0 && err != nil {
|
||||||
|
t.Fatalf("unrelated later pool blocked legacy cancellation: %v", err)
|
||||||
|
}
|
||||||
|
if owner == 1 {
|
||||||
|
if err == nil || toAPIError(t.Context(), err).HTTPStatusCode != 503 {
|
||||||
|
t.Fatalf("unknown earlier pool must retain released error behavior: %v", err)
|
||||||
|
}
|
||||||
|
if _, err := z.serverPools[owner].GetMultipartInfo(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatalf("later pool visited despite earlier error: %v", err)
|
||||||
|
}
|
||||||
|
emptySet.getDisks = original
|
||||||
|
if err := z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestMultipartAbortMinorityLegacyCleanup(t *testing.T) {
|
||||||
|
for _, failDelete := range []bool{false, true} {
|
||||||
|
t.Run(fmt.Sprintf("delete-fails=%v", failDelete), func(t *testing.T) {
|
||||||
|
z, set, bucket := multipartListingFixture(t)
|
||||||
|
mp, err := z.NewMultipartUpload(t.Context(), bucket, "old", ObjectOptions{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
path := set.getUploadIDDir(bucket, "old", mp.UploadID)
|
||||||
|
original := set.getDisks
|
||||||
|
disks := original()
|
||||||
|
fi, err := disks[0].ReadVersion(t.Context(), bucket, minioMetaMultipartBucket, path, "", ReadOptions{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
delete(fi.Metadata, multipartMetaBucket)
|
||||||
|
delete(fi.Metadata, multipartMetaObject)
|
||||||
|
if err := disks[0].WriteMetadata(t.Context(), bucket, minioMetaMultipartBucket, path, fi); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
for _, d := range disks[1:] {
|
||||||
|
if err := d.Delete(t.Context(), minioMetaMultipartBucket, path, DeleteOptions{Recursive: true}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
report, err := z.multipartPreflight(t.Context())
|
||||||
|
if err != nil || report.Ready || report.LegacyUploads != 1 {
|
||||||
|
t.Fatalf("invalid legacy fixture: %+v %v", report, err)
|
||||||
|
}
|
||||||
|
if failDelete {
|
||||||
|
set.getDisks = func() []StorageAPI {
|
||||||
|
wrapped := append([]StorageAPI(nil), disks...)
|
||||||
|
wrapped[0] = multipartListingFaultDisk{StorageAPI: disks[0], delete: func(context.Context, string, string, DeleteOptions) error { return errFaultyDisk }}
|
||||||
|
return wrapped
|
||||||
|
}
|
||||||
|
t.Cleanup(func() { set.getDisks = original })
|
||||||
|
err := z.AbortMultipartUpload(t.Context(), bucket, "old", mp.UploadID, ObjectOptions{})
|
||||||
|
if err == nil || toAPIError(t.Context(), err).HTTPStatusCode != 503 {
|
||||||
|
t.Fatalf("empty drives masked failed deletion of the observed replica: %v", err)
|
||||||
|
}
|
||||||
|
if _, present := z.mpCache.Load(mp.UploadID); !present {
|
||||||
|
t.Fatal("failed cancellation evicted upload from cache")
|
||||||
|
}
|
||||||
|
set.getDisks = original
|
||||||
|
}
|
||||||
|
if err := z.AbortMultipartUpload(t.Context(), bucket, "old", mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
report, err = z.multipartPreflight(t.Context())
|
||||||
|
if err != nil || !report.Ready || report.LegacyUploads != 0 {
|
||||||
|
t.Fatalf("acknowledged cleanup left online legacy remnants: %+v %v", report, err)
|
||||||
|
}
|
||||||
|
if _, present := z.mpCache.Load(mp.UploadID); present {
|
||||||
|
t.Fatal("successful cancellation retained cache entry")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestMultipartAbortWrongTargetKeepsCache(t *testing.T) {
|
||||||
|
for _, legacy := range []bool{true, false} {
|
||||||
|
t.Run(fmt.Sprintf("legacy=%v", legacy), func(t *testing.T) {
|
||||||
|
z, _, bucket := multipartListingFixture(t)
|
||||||
|
setMultipartListingTestMode(t, legacy)
|
||||||
|
const other = "multipart-other-bucket"
|
||||||
|
if err := z.MakeBucket(t.Context(), other, MakeBucketOptions{}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
mp, err := z.NewMultipartUpload(t.Context(), bucket, "valid-key", ObjectOptions{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
for _, target := range [][2]string{{bucket, "wrong-key"}, {other, "valid-key"}} {
|
||||||
|
err := z.AbortMultipartUpload(t.Context(), target[0], target[1], mp.UploadID, ObjectOptions{})
|
||||||
|
var invalid InvalidUploadID
|
||||||
|
if !errors.As(err, &invalid) {
|
||||||
|
t.Fatalf("wrong target result: %v", err)
|
||||||
|
}
|
||||||
|
if _, present := z.mpCache.Load(mp.UploadID); !present {
|
||||||
|
t.Fatal("wrong-target abort evicted another upload")
|
||||||
|
}
|
||||||
|
if _, err := z.GetMultipartInfo(t.Context(), bucket, "valid-key", mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatalf("wrong-target abort harmed valid upload: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if err := z.AbortMultipartUpload(t.Context(), bucket, "valid-key", mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
var invalid InvalidUploadID
|
||||||
|
if err := z.AbortMultipartUpload(t.Context(), bucket, "valid-key", mp.UploadID, ObjectOptions{}); !errors.As(err, &invalid) {
|
||||||
|
t.Fatalf("repeat abort: %v", err)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestMultipartAbortRejectsUnsafeID(t *testing.T) {
|
||||||
|
for _, legacy := range []bool{true, false} {
|
||||||
|
t.Run(fmt.Sprintf("legacy=%v", legacy), func(t *testing.T) {
|
||||||
|
z, _, bucket := multipartListingFixture(t)
|
||||||
|
setMultipartListingTestMode(t, legacy)
|
||||||
|
mp, err := z.NewMultipartUpload(t.Context(), bucket, "keep", ObjectOptions{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
for _, suffix := range []string{".", "..", "../other", "a/b", "a\\b", ""} {
|
||||||
|
id := base64.RawURLEncoding.EncodeToString([]byte("deployment." + suffix))
|
||||||
|
err := z.AbortMultipartUpload(t.Context(), bucket, "keep", id, ObjectOptions{})
|
||||||
|
var invalid InvalidUploadID
|
||||||
|
if !errors.As(err, &invalid) {
|
||||||
|
t.Fatalf("unsafe ID accepted: suffix=%q err=%v", suffix, err)
|
||||||
|
}
|
||||||
|
if _, err := z.GetMultipartInfo(t.Context(), bucket, "keep", mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatalf("valid upload harmed by rejected ID: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestMultipartAbortPeerNotification(t *testing.T) {
|
||||||
|
z, set, bucket := multipartListingFixture(t)
|
||||||
|
tg, err := grid.SetupTestGrid(2)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
t.Cleanup(tg.Cleanup)
|
||||||
|
var notifications atomic.Int32
|
||||||
|
if err := cleanupUploadIDCacheMetaRPC.Register(tg.Managers[1], func(*grid.MSS) (grid.NoPayload, *grid.RemoteErr) {
|
||||||
|
notifications.Add(1)
|
||||||
|
return grid.NoPayload{}, nil
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
host, err := xnet.ParseHost(strings.TrimPrefix(tg.Hosts[1], "http://"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
previous := globalNotificationSys
|
||||||
|
globalNotificationSys = &NotificationSys{peerClients: []*peerRESTClient{{
|
||||||
|
host: host,
|
||||||
|
gridConn: func() *grid.Connection { return tg.Managers[0].Connection(tg.Hosts[1]) },
|
||||||
|
}}}
|
||||||
|
t.Cleanup(func() { globalNotificationSys = previous })
|
||||||
|
// Use the real distributed lock implementation: Unlock cancels its derived
|
||||||
|
// context. Local locks do not, and would hide a broken notification context.
|
||||||
|
oldMutex, oldLockers := set.nsMutex, set.getLockers
|
||||||
|
lockers := []dsync.NetLocker{newLocker(), newLocker()}
|
||||||
|
set.nsMutex = newNSLock(true)
|
||||||
|
set.getLockers = func() ([]dsync.NetLocker, string) { return lockers, "multipart-test" }
|
||||||
|
t.Cleanup(func() { set.nsMutex, set.getLockers = oldMutex, oldLockers })
|
||||||
|
for _, legacy := range []bool{true, false} {
|
||||||
|
t.Run(fmt.Sprintf("legacy=%v", legacy), func(t *testing.T) {
|
||||||
|
setMultipartListingTestMode(t, legacy)
|
||||||
|
for range 4 {
|
||||||
|
mp, err := z.NewMultipartUpload(t.Context(), bucket, "valid", ObjectOptions{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
before := notifications.Load()
|
||||||
|
var invalid InvalidUploadID
|
||||||
|
if err := z.AbortMultipartUpload(t.Context(), bucket, "wrong", mp.UploadID, ObjectOptions{}); !errors.As(err, &invalid) {
|
||||||
|
t.Fatalf("wrong target: %v", err)
|
||||||
|
}
|
||||||
|
if notifications.Load() != before {
|
||||||
|
t.Fatal("wrong-target abort sent a destructive peer notification")
|
||||||
|
}
|
||||||
|
if err := z.AbortMultipartUpload(t.Context(), bucket, "valid", mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if notifications.Load() != before+1 {
|
||||||
|
t.Fatal("successful abort lost peer notification after unlocking")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestMultipartAbortStrictMajority(t *testing.T) {
|
||||||
|
z, set, bucket := multipartListingFixture(t)
|
||||||
|
mp, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
original := set.getDisks
|
||||||
|
disks := original()
|
||||||
|
set.getDisks = func() []StorageAPI {
|
||||||
|
wrapped := append([]StorageAPI(nil), disks...)
|
||||||
|
wrapped[0] = multipartListingFaultDisk{StorageAPI: disks[0], delete: func(context.Context, string, string, DeleteOptions) error { return errFaultyDisk }}
|
||||||
|
return wrapped
|
||||||
|
}
|
||||||
|
t.Cleanup(func() { set.getDisks = original })
|
||||||
|
if err := z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatalf("ordinary majority cancellation was tightened: %v", err)
|
||||||
|
}
|
||||||
|
// One failed deletion is tolerated in the ordinary majority path. Retrying
|
||||||
|
// after recovery must also clean the now-minority remnant.
|
||||||
|
if _, err := disks[0].ReadVersion(t.Context(), bucket, minioMetaMultipartBucket, set.getUploadIDDir(bucket, "a", mp.UploadID), "", ReadOptions{}); err != nil {
|
||||||
|
t.Fatalf("fault fixture did not leave a remnant: %v", err)
|
||||||
|
}
|
||||||
|
set.getDisks = original
|
||||||
|
if err := z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
+37
-12
@@ -1753,11 +1753,11 @@ func (er erasureObjects) CompleteMultipartUpload(ctx context.Context, bucket str
|
|||||||
return fi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended), nil
|
return fi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// abortMultipartUpload confirms absence on a strict majority of this set.
|
// abortMultipartUpload retains read-quorum validation and best-effort cleanup
|
||||||
// Unlike an existence read followed by best-effort deletion, it also permits
|
// in legacy mode. Strict mode requires majority deletion acknowledgements and
|
||||||
// retrying a partial deletion which no longer has a readable metadata quorum.
|
// permits retrying remnants below read quorum. Neither mode fences creation
|
||||||
// This does not fence creation writes still executing after a storage timeout.
|
// writes still executing after a storage timeout.
|
||||||
func (er erasureObjects) abortMultipartUpload(ctx context.Context, bucket, object, uploadID string, opts ObjectOptions) (bool, error) {
|
func (er erasureObjects) abortMultipartUpload(ctx context.Context, bucket, object, uploadID string, opts ObjectOptions, legacy bool) (bool, error) {
|
||||||
if !opts.NoAuditLog {
|
if !opts.NoAuditLog {
|
||||||
auditObjectErasureSet(ctx, "AbortMultipartUpload", object, &er)
|
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, "/\\") {
|
if !ok || internalID == "" || internalID == "." || internalID == ".." || strings.ContainsAny(internalID, "/\\") {
|
||||||
return false, InvalidUploadID{Bucket: bucket, Object: object, UploadID: uploadID}
|
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()
|
disks := er.getDisks()
|
||||||
uploadPath := er.getUploadIDDir(bucket, object, uploadID)
|
uploadPath := er.getUploadIDDir(bucket, object, uploadID)
|
||||||
_, errs := readAllFileInfo(ctx, disks, bucket, minioMetaMultipartBucket, uploadPath, "", false, false)
|
_, errs := readAllFileInfo(ctx, disks, bucket, minioMetaMultipartBucket, uploadPath, "", false, false)
|
||||||
quorum := er.setDriveCount/2 + 1
|
quorum := er.setDriveCount/2 + 1
|
||||||
found, absent := false, 0
|
found, absent := false, 0
|
||||||
for _, err := range errs {
|
observed := make([]bool, len(disks))
|
||||||
|
for i, err := range errs {
|
||||||
switch {
|
switch {
|
||||||
case err == nil, errors.Is(err, errFileCorrupt):
|
case err == nil, errors.Is(err, errFileCorrupt):
|
||||||
found = true
|
found = true
|
||||||
|
observed[i] = true
|
||||||
case errors.Is(err, errFileNotFound), errors.Is(err, errFileVersionNotFound):
|
case errors.Is(err, errFileNotFound), errors.Is(err, errFileVersionNotFound):
|
||||||
absent++
|
absent++
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if absent >= quorum {
|
if absent >= quorum && !found {
|
||||||
return found, nil
|
return false, nil
|
||||||
}
|
}
|
||||||
if !found {
|
if !found {
|
||||||
return false, toObjectErr(errErasureReadQuorum, bucket, object, uploadID)
|
return false, toObjectErr(errErasureReadQuorum, bucket, object, uploadID)
|
||||||
@@ -1801,13 +1814,25 @@ func (er erasureObjects) abortMultipartUpload(ctx context.Context, bucket, objec
|
|||||||
return err
|
return err
|
||||||
}, i)
|
}, 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
|
// AbortMultipartUpload cancels an upload using the configured mode. Offline
|
||||||
// still need stale-upload cleanup after its drives return.
|
// 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) {
|
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 {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2160,13 +2160,14 @@ func (z *erasureServerPools) AbortMultipartUpload(ctx context.Context, bucket, o
|
|||||||
return toObjectErr(err, bucket)
|
return toObjectErr(err, bucket)
|
||||||
}
|
}
|
||||||
|
|
||||||
defer func() {
|
// Unlock cancels the derived lock context before this notification runs.
|
||||||
_, absent := err.(InvalidUploadID)
|
// Keep the request context so successful cancellation reaches peer caches.
|
||||||
if err == nil || absent {
|
defer func(ctx context.Context) {
|
||||||
|
if err == nil {
|
||||||
z.mpCache.Delete(uploadID)
|
z.mpCache.Delete(uploadID)
|
||||||
globalNotificationSys.DeleteUploadID(ctx, uploadID)
|
globalNotificationSys.DeleteUploadID(ctx, uploadID)
|
||||||
}
|
}
|
||||||
}()
|
}(ctx)
|
||||||
|
|
||||||
lk := z.NewNSLock(bucket, pathJoin(object, uploadID))
|
lk := z.NewNSLock(bucket, pathJoin(object, uploadID))
|
||||||
lkctx, err := lk.GetLock(ctx, globalOperationTimeout)
|
lkctx, err := lk.GetLock(ctx, globalOperationTimeout)
|
||||||
@@ -2176,13 +2177,18 @@ func (z *erasureServerPools) AbortMultipartUpload(ctx context.Context, bucket, o
|
|||||||
ctx = lkctx.Context()
|
ctx = lkctx.Context()
|
||||||
defer lk.Unlock(lkctx)
|
defer lk.Unlock(lkctx)
|
||||||
|
|
||||||
|
legacy := globalAPIConfig.getMultipartListingLegacy()
|
||||||
found := false
|
found := false
|
||||||
var firstErr error
|
var firstErr error
|
||||||
for idx, pool := range z.serverPools {
|
for idx, pool := range z.serverPools {
|
||||||
if z.IsSuspended(idx) {
|
if z.IsSuspended(idx) {
|
||||||
continue
|
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
|
found = found || poolFound
|
||||||
if err != nil && firstErr == nil {
|
if err != nil && firstErr == nil {
|
||||||
firstErr = err
|
firstErr = err
|
||||||
|
|||||||
+3
-3
@@ -50,7 +50,7 @@ type apiConfig struct {
|
|||||||
transitionWorkers int
|
transitionWorkers int
|
||||||
|
|
||||||
staleUploadsExpiry time.Duration
|
staleUploadsExpiry time.Duration
|
||||||
multipartListingLegacy bool
|
multipartListingStrict bool
|
||||||
staleUploadsCleanupInterval time.Duration
|
staleUploadsCleanupInterval time.Duration
|
||||||
deleteCleanupInterval time.Duration
|
deleteCleanupInterval time.Duration
|
||||||
enableODirect bool
|
enableODirect bool
|
||||||
@@ -182,7 +182,7 @@ func (t *apiConfig) init(cfg api.Config, setDriveCounts []int, legacy bool) {
|
|||||||
t.transitionWorkers = cfg.TransitionWorkers
|
t.transitionWorkers = cfg.TransitionWorkers
|
||||||
|
|
||||||
t.staleUploadsExpiry = cfg.StaleUploadsExpiry
|
t.staleUploadsExpiry = cfg.StaleUploadsExpiry
|
||||||
t.multipartListingLegacy = cfg.MultipartListing == "legacy"
|
t.multipartListingStrict = cfg.MultipartListing == "strict"
|
||||||
t.deleteCleanupInterval = cfg.DeleteCleanupInterval
|
t.deleteCleanupInterval = cfg.DeleteCleanupInterval
|
||||||
t.enableODirect = cfg.EnableODirect
|
t.enableODirect = cfg.EnableODirect
|
||||||
t.gzipObjects = cfg.GzipObjects
|
t.gzipObjects = cfg.GzipObjects
|
||||||
@@ -211,7 +211,7 @@ func (t *apiConfig) odirectEnabled() bool {
|
|||||||
func (t *apiConfig) getMultipartListingLegacy() bool {
|
func (t *apiConfig) getMultipartListingLegacy() bool {
|
||||||
t.mu.RLock()
|
t.mu.RLock()
|
||||||
defer t.mu.RUnlock()
|
defer t.mu.RUnlock()
|
||||||
return t.multipartListingLegacy
|
return !t.multipartListingStrict
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *apiConfig) shouldGzipObjects() bool {
|
func (t *apiConfig) shouldGzipObjects() bool {
|
||||||
|
|||||||
@@ -47,6 +47,7 @@ func TestListMultipartUploadsS3Compatibility(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
z := obj.(*erasureServerPools)
|
z := obj.(*erasureServerPools)
|
||||||
|
setMultipartListingTestMode(t, false)
|
||||||
t.Cleanup(func() {
|
t.Cleanup(func() {
|
||||||
z.Shutdown(t.Context())
|
z.Shutdown(t.Context())
|
||||||
removeRoots(dirs)
|
removeRoots(dirs)
|
||||||
@@ -201,15 +202,7 @@ func TestListMultipartUploadsS3Compatibility(t *testing.T) {
|
|||||||
if !errors.Is(err, errMultipartListingLegacy) {
|
if !errors.Is(err, errMultipartListingLegacy) {
|
||||||
t.Fatalf("legacy strict listing: %v", err)
|
t.Fatalf("legacy strict listing: %v", err)
|
||||||
}
|
}
|
||||||
globalAPIConfig.mu.Lock()
|
setMultipartListingTestMode(t, true)
|
||||||
oldLegacy := globalAPIConfig.multipartListingLegacy
|
|
||||||
globalAPIConfig.multipartListingLegacy = true
|
|
||||||
globalAPIConfig.mu.Unlock()
|
|
||||||
t.Cleanup(func() {
|
|
||||||
globalAPIConfig.mu.Lock()
|
|
||||||
globalAPIConfig.multipartListingLegacy = oldLegacy
|
|
||||||
globalAPIConfig.mu.Unlock()
|
|
||||||
})
|
|
||||||
legacy, err := z.ListMultipartUploads(t.Context(), bucket, legacyObject, "", "", "", 100)
|
legacy, err := z.ListMultipartUploads(t.Context(), bucket, legacyObject, "", "", "", 100)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
@@ -267,6 +260,7 @@ func TestPaginateMultipartUploads(t *testing.T) {
|
|||||||
|
|
||||||
func TestListMultipartUploadsGlobalPageAcrossPools(t *testing.T) {
|
func TestListMultipartUploadsGlobalPageAcrossPools(t *testing.T) {
|
||||||
z, bucket := consistencyPools(t)
|
z, bucket := consistencyPools(t)
|
||||||
|
setMultipartListingTestMode(t, false)
|
||||||
objects := []string{"a/one", "b/two", "c/three", "d/four"}
|
objects := []string{"a/one", "b/two", "c/three", "d/four"}
|
||||||
for i, object := range objects {
|
for i, object := range objects {
|
||||||
if _, err := z.serverPools[i%len(z.serverPools)].NewMultipartUpload(t.Context(), bucket, object, ObjectOptions{}); err != nil {
|
if _, err := z.serverPools[i%len(z.serverPools)].NewMultipartUpload(t.Context(), bucket, object, ObjectOptions{}); err != nil {
|
||||||
|
|||||||
@@ -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
|
// 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.
|
// 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]
|
cpuResourceMetrics, found := resourceMetricsMap[cpuSubsystem]
|
||||||
if found {
|
if found {
|
||||||
if cpuIdleMetric, ok := cpuResourceMetrics[getResourceKey(cpuIdle, nil)]; ok {
|
if cpuIdleMetric, ok := cpuResourceMetrics[getResourceKey(cpuIdle, nil)]; ok {
|
||||||
|
|||||||
@@ -0,0 +1,185 @@
|
|||||||
|
// Copyright (c) 2026 Ruohang Feng
|
||||||
|
//
|
||||||
|
// This program is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Affero General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// This program is distributed in the hope that it will be useful
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Affero General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Affero General Public License
|
||||||
|
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package cmd
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"runtime"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/minio/madmin-go/v3"
|
||||||
|
"github.com/minio/minio/internal/cachevalue"
|
||||||
|
"github.com/prometheus/client_golang/prometheus"
|
||||||
|
dto "github.com/prometheus/client_model/go"
|
||||||
|
cpustats "github.com/shirou/gopsutil/v3/cpu"
|
||||||
|
"github.com/shirou/gopsutil/v3/load"
|
||||||
|
)
|
||||||
|
|
||||||
|
// These tests replace global resource metrics and must not run in parallel.
|
||||||
|
// Concurrent writers must finish before the fixture is restored.
|
||||||
|
func setCPUResourceMetricsForTest(t *testing.T, value map[MetricSubsystem]ResourceMetrics) {
|
||||||
|
t.Helper()
|
||||||
|
resourceMetricsMapMu.Lock()
|
||||||
|
saved := resourceMetricsMap
|
||||||
|
resourceMetricsMap = value
|
||||||
|
resourceMetricsMapMu.Unlock()
|
||||||
|
t.Cleanup(func() {
|
||||||
|
resourceMetricsMapMu.Lock()
|
||||||
|
resourceMetricsMap = saved
|
||||||
|
resourceMetricsMapMu.Unlock()
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func newCPUMetricsTestRegistry() *prometheus.Registry {
|
||||||
|
c := &metricsCache{cpuMetrics: cachevalue.NewFromFunc(time.Hour, cachevalue.Opts{},
|
||||||
|
func(context.Context) (madmin.CPUMetrics, error) {
|
||||||
|
return madmin.CPUMetrics{
|
||||||
|
CPUCount: 4,
|
||||||
|
LoadStat: &load.AvgStat{Load1: 2},
|
||||||
|
TimesStat: &cpustats.TimesStat{
|
||||||
|
User: 10, System: 20, Idle: 60, Iowait: 5, Nice: 3, Steal: 2,
|
||||||
|
},
|
||||||
|
}, nil
|
||||||
|
})}
|
||||||
|
_, _ = c.cpuMetrics.Get() // Warm the cache: the cache does not protect the resource map.
|
||||||
|
g := NewMetricsGroup(systemCPUCollectorPath, []MetricDescriptor{
|
||||||
|
sysCPUAvgIdleMD, sysCPUAvgIOWaitMD, sysCPULoadMD, sysCPULoadPercMD,
|
||||||
|
sysCPUNiceMD, sysCPUStealMD, sysCPUSystemMD, sysCPUUserMD,
|
||||||
|
}, loadCPUMetrics)
|
||||||
|
g.SetCache(c)
|
||||||
|
r := prometheus.NewPedanticRegistry()
|
||||||
|
r.MustRegister(g)
|
||||||
|
return r
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestLoadCPUMetricsValues(t *testing.T) {
|
||||||
|
for _, tc := range []struct {
|
||||||
|
name string
|
||||||
|
data map[MetricSubsystem]ResourceMetrics
|
||||||
|
want map[string]float64
|
||||||
|
}{
|
||||||
|
{name: "nil-map"},
|
||||||
|
{name: "empty-map", data: map[MetricSubsystem]ResourceMetrics{}},
|
||||||
|
{name: "no-cpu", data: map[MetricSubsystem]ResourceMetrics{memSubsystem: {}}},
|
||||||
|
{name: "nil-cpu", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: nil}},
|
||||||
|
{name: "empty-cpu", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {}}},
|
||||||
|
{name: "idle-only", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {
|
||||||
|
getResourceKey(cpuIdle, nil): {Avg: 87.654},
|
||||||
|
}}, want: map[string]float64{"minio_system_cpu_avg_idle": 87.65}},
|
||||||
|
{name: "iowait-only", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {
|
||||||
|
getResourceKey(cpuIOWait, nil): {Avg: 1.236},
|
||||||
|
}}, want: map[string]float64{"minio_system_cpu_avg_iowait": 1.24}},
|
||||||
|
{name: "both-rounded", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {
|
||||||
|
getResourceKey(cpuIdle, nil): {Avg: 87.654},
|
||||||
|
getResourceKey(cpuIOWait, nil): {Avg: 1.236},
|
||||||
|
}}, want: map[string]float64{"minio_system_cpu_avg_idle": 87.65, "minio_system_cpu_avg_iowait": 1.24}},
|
||||||
|
{name: "zero-keeps-existing-omission", data: map[MetricSubsystem]ResourceMetrics{cpuSubsystem: {
|
||||||
|
getResourceKey(cpuIdle, nil): {Avg: 0},
|
||||||
|
getResourceKey(cpuIOWait, nil): {Avg: 0},
|
||||||
|
}}},
|
||||||
|
} {
|
||||||
|
t.Run(tc.name, func(t *testing.T) {
|
||||||
|
setCPUResourceMetricsForTest(t, tc.data)
|
||||||
|
families, err := newCPUMetricsTestRegistry().Gather()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
want := map[string]float64{
|
||||||
|
"minio_system_cpu_load": 2, "minio_system_cpu_load_perc": 50,
|
||||||
|
"minio_system_cpu_nice": 3, "minio_system_cpu_steal": 2,
|
||||||
|
"minio_system_cpu_system": 20, "minio_system_cpu_user": 10,
|
||||||
|
}
|
||||||
|
for k, v := range tc.want {
|
||||||
|
want[k] = v
|
||||||
|
}
|
||||||
|
if len(families) != len(want) {
|
||||||
|
t.Fatalf("got %d families, want %d", len(families), len(want))
|
||||||
|
}
|
||||||
|
for _, family := range families {
|
||||||
|
v, ok := want[family.GetName()]
|
||||||
|
if !ok || family.GetType() != dto.MetricType_GAUGE || len(family.Metric) != 1 || family.Metric[0].GetGauge().GetValue() != v {
|
||||||
|
t.Errorf("unexpected family: %v (want %v)", family, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestLoadCPUMetricsConcurrentUpdate(t *testing.T) {
|
||||||
|
for _, subsystem := range []MetricSubsystem{cpuSubsystem, memSubsystem} {
|
||||||
|
for _, readers := range []int{1, 4} {
|
||||||
|
t.Run(fmt.Sprintf("writer-%s/readers-%d", subsystem, readers), func(t *testing.T) {
|
||||||
|
setCPUResourceMetricsForTest(t, map[MetricSubsystem]ResourceMetrics{})
|
||||||
|
updateResourceMetrics(cpuSubsystem, cpuIdle, 80, nil, false)
|
||||||
|
updateResourceMetrics(cpuSubsystem, cpuIOWait, 5, nil, false)
|
||||||
|
r := newCPUMetricsTestRegistry()
|
||||||
|
stop := make(chan struct{})
|
||||||
|
done := make(chan struct{})
|
||||||
|
ready := make(chan struct{})
|
||||||
|
var updates atomic.Uint64
|
||||||
|
go func() {
|
||||||
|
defer close(done)
|
||||||
|
if subsystem == cpuSubsystem {
|
||||||
|
updateResourceMetrics(subsystem, cpuIdle, 80, nil, false)
|
||||||
|
} else {
|
||||||
|
updateResourceMetrics(subsystem, memUsed, 1024, nil, false)
|
||||||
|
}
|
||||||
|
close(ready)
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-stop:
|
||||||
|
return
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
if subsystem == cpuSubsystem {
|
||||||
|
updateResourceMetrics(subsystem, cpuIdle, 80, nil, false)
|
||||||
|
updateResourceMetrics(subsystem, cpuIOWait, 5, nil, false)
|
||||||
|
} else {
|
||||||
|
updateResourceMetrics(subsystem, memUsed, 1024, nil, false)
|
||||||
|
}
|
||||||
|
updates.Add(1)
|
||||||
|
runtime.Gosched()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
<-ready
|
||||||
|
defer func() { close(stop); <-done }()
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
for range readers {
|
||||||
|
wg.Add(1)
|
||||||
|
go func() {
|
||||||
|
defer wg.Done()
|
||||||
|
for range 1000 {
|
||||||
|
families, err := r.Gather()
|
||||||
|
if err != nil || len(families) != 8 {
|
||||||
|
t.Errorf("Gather: families=%d, error=%v", len(families), err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
|
if updates.Load() == 0 {
|
||||||
|
t.Fatal("writer made no progress")
|
||||||
|
}
|
||||||
|
t.Logf("completed %d gathers with %d concurrent update iterations", readers*1000, updates.Load())
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -41,15 +41,23 @@ const r5TagStamp = ReservedMetadataPrefixLower + TaggingTimestamp
|
|||||||
func r5Capacity(z *erasureServerPools) func() {
|
func r5Capacity(z *erasureServerPools) func() {
|
||||||
var restores []func()
|
var restores []func()
|
||||||
for _, pool := range z.serverPools {
|
for _, pool := range z.serverPools {
|
||||||
|
pool.erasureDisksMu.Lock()
|
||||||
for _, set := range pool.sets {
|
for _, set := range pool.sets {
|
||||||
old := set.getDisks
|
old := pool.erasureDisks[set.setIndex]
|
||||||
disks := append([]StorageAPI(nil), old()...)
|
disks := append([]StorageAPI(nil), old...)
|
||||||
for i := range disks {
|
for i := range disks {
|
||||||
disks[i] = tagTestCapacityDisk{StorageAPI: disks[i]}
|
disks[i] = tagTestCapacityDisk{StorageAPI: disks[i]}
|
||||||
}
|
}
|
||||||
set.getDisks = func() []StorageAPI { return disks }
|
// GetDisks copies this list under the same mutex. Keep its function
|
||||||
restores = append(restores, func() { set.getDisks = old })
|
// 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() {
|
return func() {
|
||||||
for _, restore := range restores {
|
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 {
|
func r5Request(t *testing.T, router http.Handler, cred auth.Credentials, method, path, body string, headers map[string]string) *httptest.ResponseRecorder {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
r, err := newTestSignedRequestV4(method, path, int64(len(body)), strings.NewReader(body), cred.AccessKey, cred.SecretKey, headers)
|
r, err := newTestSignedRequestV4(method, path, int64(len(body)), strings.NewReader(body), cred.AccessKey, cred.SecretKey, headers)
|
||||||
|
|||||||
@@ -0,0 +1,101 @@
|
|||||||
|
// Copyright (c) 2026 Feng Ruohang
|
||||||
|
//
|
||||||
|
// This program is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Affero General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// This program is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Affero General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Affero General Public License
|
||||||
|
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package cmd
|
||||||
|
|
||||||
|
import (
|
||||||
|
"archive/tar"
|
||||||
|
"bytes"
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"os"
|
||||||
|
"sync"
|
||||||
|
"testing"
|
||||||
|
"testing/iotest"
|
||||||
|
|
||||||
|
"github.com/pierrec/lz4/v4"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestUntarLZ4(t *testing.T) {
|
||||||
|
for _, empty := range []bool{false, true} {
|
||||||
|
t.Run(fmt.Sprintf("empty=%t", empty), func(t *testing.T) {
|
||||||
|
want := map[string][]byte{}
|
||||||
|
if !empty {
|
||||||
|
want["small.txt"] = bytes.Repeat([]byte("small object\n"), 10)
|
||||||
|
want["large.txt"] = bytes.Repeat([]byte("large object\n"), 32768)
|
||||||
|
}
|
||||||
|
var compressed bytes.Buffer
|
||||||
|
compressor := lz4.NewWriter(&compressed)
|
||||||
|
archive := tar.NewWriter(compressor)
|
||||||
|
for name, data := range want {
|
||||||
|
if err := archive.WriteHeader(&tar.Header{Name: name, Mode: 0o600, Size: int64(len(data))}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := archive.Write(data); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if err := archive.Close(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := compressor.Close(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
var mu sync.Mutex
|
||||||
|
got := map[string][]byte{}
|
||||||
|
err := untar(t.Context(), iotest.HalfReader(bytes.NewReader(compressed.Bytes())), func(r io.Reader, info os.FileInfo, name string) error {
|
||||||
|
// Exercise a short first read followed by streaming the remaining content.
|
||||||
|
var output bytes.Buffer
|
||||||
|
prefix := make([]byte, 7)
|
||||||
|
if _, err := io.ReadFull(r, prefix); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
output.Write(prefix)
|
||||||
|
if _, err := io.Copy(&output, r); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if int64(output.Len()) != info.Size() {
|
||||||
|
return fmt.Errorf("size mismatch for %s", name)
|
||||||
|
}
|
||||||
|
mu.Lock()
|
||||||
|
got[name] = output.Bytes()
|
||||||
|
mu.Unlock()
|
||||||
|
return nil
|
||||||
|
}, untarOptions{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if len(got) != len(want) {
|
||||||
|
t.Fatalf("objects=%d, want %d", len(got), len(want))
|
||||||
|
}
|
||||||
|
for name, data := range want {
|
||||||
|
if !bytes.Equal(got[name], data) {
|
||||||
|
t.Errorf("content mismatch for %s", name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
t.Run("corrupt-header", func(t *testing.T) {
|
||||||
|
damaged := bytes.Clone(compressed.Bytes())
|
||||||
|
damaged[6] ^= 0xff
|
||||||
|
err := untar(t.Context(), bytes.NewReader(damaged), func(io.Reader, os.FileInfo, string) error {
|
||||||
|
t.Error("corrupt header must be rejected before uploading objects")
|
||||||
|
return nil
|
||||||
|
}, untarOptions{})
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("corrupt LZ4 header accepted")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -4,9 +4,9 @@ go 1.27.1
|
|||||||
|
|
||||||
// Console and MC retain their historical module paths for best-effort upstream
|
// Console and MC retain their historical module paths for best-effort upstream
|
||||||
// compatibility. Pin the maintained PGSTY implementations used by SILO.
|
// 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
|
// v22.7.0 does not compile on NetBSD because its unix implementation uses
|
||||||
// CLOCK_MONOTONIC, which is unavailable there. Keep the last portable release
|
// 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/kes v0.3.1
|
||||||
github.com/minio/kms-go/kms v0.6.0
|
github.com/minio/kms-go/kms v0.6.0
|
||||||
github.com/minio/madmin-go/v3 v3.0.110
|
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/mux v1.10.1
|
||||||
github.com/minio/selfupdate v0.6.0
|
github.com/minio/selfupdate v0.6.0
|
||||||
github.com/minio/simdjson-go v0.4.5
|
github.com/minio/simdjson-go v0.4.5
|
||||||
@@ -81,9 +81,9 @@ require (
|
|||||||
github.com/nats-io/stan.go v0.10.4
|
github.com/nats-io/stan.go v0.10.4
|
||||||
github.com/ncw/directio v1.0.5
|
github.com/ncw/directio v1.0.5
|
||||||
github.com/nsqio/go-nsq v1.1.0
|
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/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/errors v0.9.1
|
||||||
github.com/pkg/sftp v1.13.11
|
github.com/pkg/sftp v1.13.11
|
||||||
github.com/pkg/xattr v0.4.12
|
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 v0.33.1 // indirect
|
||||||
github.com/go-openapi/runtime/server-middleware 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/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 v0.29.1 // indirect
|
||||||
github.com/go-openapi/swag/cmdutils v0.29.1 // indirect
|
github.com/go-openapi/swag/cmdutils v0.29.1 // indirect
|
||||||
github.com/go-openapi/swag/conv 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/dsig-secp256k1 v1.0.0 // indirect
|
||||||
github.com/lestrrat-go/httpcc v1.0.1 // indirect
|
github.com/lestrrat-go/httpcc v1.0.1 // indirect
|
||||||
github.com/lestrrat-go/httprc/v3 v3.0.6 // 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/lestrrat-go/option/v2 v2.0.0 // indirect
|
||||||
github.com/lucasb-eyer/go-colorful v1.4.1 // indirect
|
github.com/lucasb-eyer/go-colorful v1.4.1 // indirect
|
||||||
github.com/lufia/plan9stats v0.0.0-20260802145828-341c2f0c90b5 // indirect
|
github.com/lufia/plan9stats v0.0.0-20260802145828-341c2f0c90b5 // indirect
|
||||||
|
|||||||
@@ -214,8 +214,8 @@ github.com/go-openapi/runtime/server-middleware v0.33.1 h1:IAeKbwWnBnpsYTpuPVS8t
|
|||||||
github.com/go-openapi/runtime/server-middleware v0.33.1/go.mod h1:2Gej5fDxqeJxY+w38vxXYW0BgFASfgBsJ5rXwN1Fseg=
|
github.com/go-openapi/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 h1:JtB/GHOj+eetjse6YvxqLze88oEekl/4uPBethvzRrA=
|
||||||
github.com/go-openapi/spec v1.0.0/go.mod h1:boj1PRhqS0x5jylgcNp9BRWxenphrh7vVNYZ90IRCoo=
|
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.2 h1:SG32SlbwNy92s0KJiVxt2joJeFdqIYHvwrA0OU6HqzQ=
|
||||||
github.com/go-openapi/strfmt v0.27.0/go.mod h1:s/qhDqfY72irigXUGJmtgid2Rm+3tnz3k8hZaRmvWYc=
|
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 h1:C6EeWzUwQtcWEhE9eqBdUubGXxhWY4PlzHMLD7kLaiQ=
|
||||||
github.com/go-openapi/swag v0.29.1/go.mod h1:BzxEXKiPlSXRsRTv1KSBF/BpGKHxA/YciCnr4tv9bvA=
|
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=
|
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/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 h1:4FpLQ18KK/ypPbVU3NLWJNRvH3kcYiqKqWfKGqNWxxI=
|
||||||
github.com/lestrrat-go/httprc/v3 v3.0.6/go.mod h1:mSMtkZW92Z98M5YoNNztbRGxbXHql7tSitCvaxvo9l0=
|
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.3.0 h1:OXcYvQOQ7cxWzeZ/Q9sYk8ABe/kCSI371WmuACiCT+4=
|
||||||
github.com/lestrrat-go/jwx/v3 v3.2.0/go.mod h1:38vQ8iWKq3qRSbilbzvzdQPuywhowwuR03lhkYskyrw=
|
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 h1:XxrcaJESE1fokHy3FpaQ/cXW8ZsIdWcdFzzLOcID3Ss=
|
||||||
github.com/lestrrat-go/option/v2 v2.0.0/go.mod h1:oSySsmzMoR0iRzCDCaUfsCzxQHUEuhOViQObyy7S6Vg=
|
github.com/lestrrat-go/option/v2 v2.0.0/go.mod h1:oSySsmzMoR0iRzCDCaUfsCzxQHUEuhOViQObyy7S6Vg=
|
||||||
github.com/lib/pq v1.10.4/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o=
|
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/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 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34=
|
||||||
github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM=
|
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.20260915093545-32e1f32cb176 h1:TIrbAkJoMN/kbtV4LAitf/eidrbvvYcujXPbSe+4j/0=
|
||||||
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/go.mod h1:KUPWdecEO1LWyUz+sTGXAuf2jZHrPh5fCsRH86QbPfk=
|
||||||
github.com/minio/mux v1.10.1 h1:grrK8SwRKbkNFE6qG7WAvFGH09bB46d5teOOtKfQ14s=
|
github.com/minio/mux v1.10.1 h1:grrK8SwRKbkNFE6qG7WAvFGH09bB46d5teOOtKfQ14s=
|
||||||
github.com/minio/mux v1.10.1/go.mod h1:INYT4sMSTJy0QWUEA/E2DZNxJ5sAxIwbnyZjkzNFRfE=
|
github.com/minio/mux v1.10.1/go.mod h1:INYT4sMSTJy0QWUEA/E2DZNxJ5sAxIwbnyZjkzNFRfE=
|
||||||
github.com/minio/pkg/v3 v3.6.1 h1:gaNT80BS/iuIany5ylTkVmfN4s6UYY30OtImFv4GQA8=
|
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/pascaldekloe/goe v0.1.0/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc=
|
||||||
github.com/pborman/getopt v0.0.0-20170112200414-7148bc3a4c30/go.mod h1:85jBQOZwpVEaDAr341tbn15RS4fCAsIst0qp7i8ex1o=
|
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/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-20260916070421-e952aa78f10a h1:FY1vcTQT67HGBe23ozlhOiQnWCzUCGmeSN5x6DT/kFs=
|
||||||
github.com/pgsty/mc v0.0.0-20260913012246-4f609a4da3bb/go.mod h1:kJN7dsWtSUhXd2vNPXdjDi/or6lJKlSLfb56cBfMckA=
|
github.com/pgsty/mc v0.0.0-20260916070421-e952aa78f10a/go.mod h1:6z9lX8kOvaWz2vLlP7BjjF3JWHnIH86XDA2HwcmpZH8=
|
||||||
github.com/pgsty/silo-console v0.0.0-20260916034812-56dfe455ac2f h1:ull9m/nXOEMfggKtMQP2idYXPh6l6chCUJ9hJHYyzDc=
|
github.com/pgsty/silo-console v0.0.0-20260916075814-1360e26d976d h1:A9ou1GkUckGZd3BOgEXeZ393nbYST967SKcI4Lgr/Vs=
|
||||||
github.com/pgsty/silo-console v0.0.0-20260916034812-56dfe455ac2f/go.mod h1:YtRQZ6jYXRUE03oPA+IMGtflQN6nCvDWdKtroA7tfKo=
|
github.com/pgsty/silo-console v0.0.0-20260916075814-1360e26d976d/go.mod h1:JS165CF2yoTBIx1X8k1IQ6vunJcZFG3pRZlWHJ8UT8o=
|
||||||
github.com/pgsty/silo-pkg/v3 v3.14.0 h1:RCuVkzr6mdjbkV/Rv4OVO/XgdMBE0XYvUnT6GFuihFk=
|
github.com/pgsty/silo-pkg/v3 v3.14.1 h1:8MONT3Hky9EOnsuPk29590pEBE6MILaOnJtaejPPi7M=
|
||||||
github.com/pgsty/silo-pkg/v3 v3.14.0/go.mod h1:c26IoMVITlP1+Sirl0AsHwoijlsSC9wPpyCuvhb3Axc=
|
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 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM=
|
||||||
github.com/philhofer/fwd v1.2.0/go.mod h1:RqIHx9QI14HlwKwm98g9Re5prTQ6LdeRQn+gXJFxsJM=
|
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.30 h1:cchX8N2DVP668WkElI9QMwVyoNabLkq1LofDHFeIrdg=
|
||||||
github.com/pierrec/lz4/v4 v4.1.29/go.mod h1:EoQMVJgeeEOMsCqCzqFm2O0cJvljX2nGZjcRIPL34O4=
|
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 h1:+mdjkGKdHQG3305AYmdv1U2eRNDiU2ErMBj1gwrq8eQ=
|
||||||
github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c/go.mod h1:7rwL4CYBLnjLxUqIJNnCWiEdr3bn6IUYi15bNlnbCCU=
|
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=
|
github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
|
||||||
|
|||||||
@@ -1,8 +1,8 @@
|
|||||||
apiVersion: v1
|
apiVersion: v1
|
||||||
description: S3-Interface Libre Object Storage
|
description: S3-Interface Libre Object Storage
|
||||||
name: silo
|
name: silo
|
||||||
version: 7.0.2
|
version: 7.0.3
|
||||||
appVersion: RELEASE.2026-09-03T13-18-01Z
|
appVersion: RELEASE.2026-09-16T00-00-00Z
|
||||||
keywords:
|
keywords:
|
||||||
- silo
|
- silo
|
||||||
- storage
|
- storage
|
||||||
|
|||||||
@@ -15,7 +15,7 @@ clusterDomain: cluster.local
|
|||||||
##
|
##
|
||||||
image:
|
image:
|
||||||
repository: pgsty/silo
|
repository: pgsty/silo
|
||||||
tag: RELEASE.2026-09-03T13-18-01Z
|
tag: RELEASE.2026-09-16T00-00-00Z
|
||||||
pullPolicy: IfNotPresent
|
pullPolicy: IfNotPresent
|
||||||
|
|
||||||
imagePullSecrets: []
|
imagePullSecrets: []
|
||||||
@@ -25,7 +25,7 @@ imagePullSecrets: []
|
|||||||
##
|
##
|
||||||
mcImage:
|
mcImage:
|
||||||
repository: pgsty/mc
|
repository: pgsty/mc
|
||||||
tag: RELEASE.2026-09-13T00-00-00Z
|
tag: RELEASE.2026-09-16T00-00-00Z
|
||||||
pullPolicy: IfNotPresent
|
pullPolicy: IfNotPresent
|
||||||
|
|
||||||
## Silo mode, i.e. standalone or distributed.
|
## Silo mode, i.e. standalone or distributed.
|
||||||
|
|||||||
@@ -135,7 +135,6 @@ var (
|
|||||||
Key: apiStaleUploadsExpiry,
|
Key: apiStaleUploadsExpiry,
|
||||||
Value: "24h",
|
Value: "24h",
|
||||||
},
|
},
|
||||||
config.KV{Key: apiMultipartListing, Value: "strict"},
|
|
||||||
config.KV{
|
config.KV{
|
||||||
Key: apiDeleteCleanupInterval,
|
Key: apiDeleteCleanupInterval,
|
||||||
Value: "5m",
|
Value: "5m",
|
||||||
@@ -210,6 +209,7 @@ func LookupConfig(kvs config.KVS) (cfg Config, err error) {
|
|||||||
apiReplicationWorkers,
|
apiReplicationWorkers,
|
||||||
apiReplicationFailedWorkers,
|
apiReplicationFailedWorkers,
|
||||||
"expiry_workers",
|
"expiry_workers",
|
||||||
|
apiMultipartListing, // Ignore the retired shared-config key during migration.
|
||||||
}
|
}
|
||||||
|
|
||||||
disableODirect := env.Get(EnvAPIDisableODirect, kvs.Get(apiDisableODirect)) == config.EnableOn
|
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
|
return cfg, err
|
||||||
}
|
}
|
||||||
cfg.StaleUploadsExpiry = staleUploadsExpiry
|
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
|
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.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
|
return cfg, nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,38 +4,111 @@
|
|||||||
package api
|
package api
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"encoding/json"
|
||||||
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/minio/minio/internal/config"
|
"github.com/minio/minio/internal/config"
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestMultipartListingMigrationMode(t *testing.T) {
|
func TestMultipartListingMigrationMode(t *testing.T) {
|
||||||
|
t.Setenv(EnvAPIRequestsMax, "")
|
||||||
|
t.Setenv(EnvAPISyncEvents, "")
|
||||||
|
t.Setenv(EnvAPIObjectMaxVersions, "")
|
||||||
for _, tc := range []struct {
|
for _, tc := range []struct {
|
||||||
name, stored, override, want string
|
name, stored, override, want string
|
||||||
invalid bool
|
invalid bool
|
||||||
}{
|
}{
|
||||||
{name: "default", want: "strict"},
|
{name: "default", want: "legacy"},
|
||||||
{name: "explicit-migration", stored: "legacy", want: "legacy"},
|
{name: "explicit-legacy", override: "legacy", want: "legacy"},
|
||||||
{name: "environment-override", stored: "strict", override: "legacy", want: "legacy"},
|
{name: "explicit-strict", override: "strict", want: "strict"},
|
||||||
{name: "invalid-stored", stored: "automatic", invalid: true},
|
{name: "retired-stored-strict", stored: "strict", want: "legacy"},
|
||||||
{name: "invalid-environment", stored: "strict", override: "automatic", invalid: true},
|
{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.Run(tc.name, func(t *testing.T) {
|
||||||
t.Setenv(EnvAPIMultipartListing, tc.override)
|
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 != "" {
|
if tc.stored != "" {
|
||||||
kvs = config.KVS{{Key: apiMultipartListing, Value: tc.stored}}
|
kvs.Set(apiMultipartListing, tc.stored)
|
||||||
}
|
}
|
||||||
cfg, err := LookupConfig(kvs)
|
cfg, err := LookupConfig(kvs)
|
||||||
if tc.invalid {
|
if (err != nil) != tc.invalid || cfg.MultipartListing != tc.want {
|
||||||
if err == nil {
|
t.Fatalf("mode=%q err=%v, want %q invalid=%v", cfg.MultipartListing, err, tc.want, tc.invalid)
|
||||||
t.Fatal("invalid migration mode silently accepted")
|
|
||||||
}
|
}
|
||||||
return
|
if cfg.RequestsMax != 17 || !cfg.SyncEvents || cfg.ObjectMaxVersions != 71 {
|
||||||
}
|
t.Fatalf("migration setting discarded unrelated config: %+v", cfg)
|
||||||
if err != nil || cfg.MultipartListing != tc.want {
|
|
||||||
t.Fatalf("mode=%q err=%v, want %q", cfg.MultipartListing, err, tc.want)
|
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestMultipartListingConfigPersistence(t *testing.T) {
|
||||||
|
t.Setenv(EnvAPIRequestsMax, "")
|
||||||
|
t.Setenv(EnvAPISyncEvents, "")
|
||||||
|
t.Setenv(EnvAPIObjectMaxVersions, "")
|
||||||
|
t.Setenv(EnvAPIMultipartListing, "strict")
|
||||||
|
previous := config.DefaultKVS
|
||||||
|
config.DefaultKVS = map[string]config.KVS{config.APISubSys: DefaultKVS}
|
||||||
|
t.Cleanup(func() { config.DefaultKVS = previous })
|
||||||
|
for _, help := range Help {
|
||||||
|
if help.Key == apiMultipartListing {
|
||||||
|
t.Fatal("process-only setting advertised as shared config")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for _, input := range []string{
|
||||||
|
"",
|
||||||
|
"api requests_max=17 sync_events=on object_max_versions=71",
|
||||||
|
`api requests_max=17 comment="later disable multipart_listing=strict"`,
|
||||||
|
} {
|
||||||
|
t.Run(input, func(t *testing.T) {
|
||||||
|
c := config.New().Clone()
|
||||||
|
if _, err := c.ReadConfig(strings.NewReader(input)); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
// Exercise the JSON save/load and Merge paths used for shared config.
|
||||||
|
data, err := json.Marshal(c)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
var loaded config.Config
|
||||||
|
if err = json.Unmarshal(data, &loaded); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
kvs := loaded.Merge()[config.APISubSys][config.Default]
|
||||||
|
if _, present := kvs.Lookup(apiMultipartListing); present {
|
||||||
|
t.Fatal("process-only setting persisted in shared config")
|
||||||
|
}
|
||||||
|
cfg, err := LookupConfig(kvs)
|
||||||
|
if err != nil || cfg.MultipartListing != "strict" {
|
||||||
|
t.Fatalf("environment mode lost: %+v %v", cfg, err)
|
||||||
|
}
|
||||||
|
if input != "" && cfg.RequestsMax != 17 {
|
||||||
|
t.Fatalf("unrelated setting lost: %+v", cfg)
|
||||||
|
}
|
||||||
|
if strings.Contains(input, "comment=") && kvs.Get(config.Comment) != "later disable multipart_listing=strict" {
|
||||||
|
t.Fatalf("literal comment changed: %q", kvs.Get(config.Comment))
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
// A key persisted by a previous development build is tolerated, retained for
|
||||||
|
// explicit cleanup, and removable without resetting other API settings.
|
||||||
|
c := config.New().Clone()
|
||||||
|
kvs := c[config.APISubSys][config.Default]
|
||||||
|
kvs.Set(apiMultipartListing, "strict")
|
||||||
|
kvs.Set(apiRequestsMax, "17")
|
||||||
|
c[config.APISubSys][config.Default] = kvs
|
||||||
|
c = c.Merge()
|
||||||
|
if err := c.DelKVS("api multipart_listing"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
kvs = c.Merge()[config.APISubSys][config.Default]
|
||||||
|
if _, present := kvs.Lookup(apiMultipartListing); present || kvs.Get(apiRequestsMax) != "17" {
|
||||||
|
t.Fatalf("targeted reset changed unrelated config or restored retired key: %v", kvs)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -26,12 +26,6 @@ var (
|
|||||||
|
|
||||||
// Help holds configuration keys and their default values for api subsystem.
|
// Help holds configuration keys and their default values for api subsystem.
|
||||||
Help = config.HelpKVS{
|
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{
|
config.HelpKV{
|
||||||
Key: apiRequestsMax,
|
Key: apiRequestsMax,
|
||||||
Description: `set the maximum number of concurrent requests (default: auto)`,
|
Description: `set the maximum number of concurrent requests (default: auto)`,
|
||||||
|
|||||||
@@ -0,0 +1,71 @@
|
|||||||
|
// Copyright (c) 2026 Feng Ruohang
|
||||||
|
//
|
||||||
|
// This program is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Affero General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// This program is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Affero General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Affero General Public License
|
||||||
|
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package s3select
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"io"
|
||||||
|
"testing"
|
||||||
|
"testing/iotest"
|
||||||
|
|
||||||
|
"github.com/pierrec/lz4/v4"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestLZ4Input(t *testing.T) {
|
||||||
|
request := []byte(`<SelectObjectContentRequest>
|
||||||
|
<Expression>SELECT * FROM S3Object</Expression><ExpressionType>SQL</ExpressionType>
|
||||||
|
<InputSerialization><CompressionType>LZ4</CompressionType><CSV><FileHeaderInfo>NONE</FileHeaderInfo></CSV></InputSerialization>
|
||||||
|
<OutputSerialization><CSV/></OutputSerialization></SelectObjectContentRequest>`)
|
||||||
|
for _, input := range []string{"", "one,1\ntwo,2\n"} {
|
||||||
|
t.Run(input, func(t *testing.T) {
|
||||||
|
var compressed bytes.Buffer
|
||||||
|
writer := lz4.NewWriter(&compressed)
|
||||||
|
if _, err := io.WriteString(writer, input); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := writer.Close(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
got, err := evaluateSelectForTest(t, request, compressed.Bytes())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if string(got) != input {
|
||||||
|
t.Fatalf("result=%q, want %q", got, input)
|
||||||
|
}
|
||||||
|
|
||||||
|
reader, err := newProgressReader(io.NopCloser(iotest.OneByteReader(bytes.NewReader(compressed.Bytes()))), lz4Type)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
t.Cleanup(func() { _ = reader.Close() })
|
||||||
|
got, err = io.ReadAll(reader)
|
||||||
|
if err != nil || string(got) != input {
|
||||||
|
t.Fatalf("fragmented input: result=%q, err=%v", got, err)
|
||||||
|
}
|
||||||
|
scanned, processed := reader.Stats()
|
||||||
|
if scanned != int64(compressed.Len()) || processed != int64(len(input)) {
|
||||||
|
t.Fatalf("scanned=%d processed=%d", scanned, processed)
|
||||||
|
}
|
||||||
|
if err := reader.Close(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := reader.Read(make([]byte, 1)); err == nil {
|
||||||
|
t.Fatal("read after Close succeeded")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user