Compare commits

..

7 Commits

Author SHA1 Message Date
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
26 changed files with 997 additions and 119 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.
+44 -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.
@@ -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
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
+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
}
+11 -5
View File
@@ -2160,13 +2160,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 +2177,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
+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 {
+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())
})
}
}
}
+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")
}
})
})
}
}
+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")
}
})
}
}