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