Compare commits

...

7 Commits

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

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

Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 19:30:07 +08:00
Feng Ruohang 83821f0f1f Merge pull request #216 from pgsty/codex/docs-published-design-links-20260916
docs: link retained guides to published design records
2026-09-16 18:50:28 +08:00
Feng Ruohang 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
21 changed files with 1106 additions and 60 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.
+26 -16
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.
@@ -167,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
+2 -2
View File
@@ -618,7 +618,7 @@ func (z *erasureServerPools) decommissionObject(ctx context.Context, idx int, bu
if objInfo.isMultipart() {
res, err := z.NewMultipartUpload(ctx, bucket, objInfo.Name, ObjectOptions{
VersionID: objInfo.VersionID,
UserDefined: objInfo.UserDefined,
UserDefined: migrationObjectMetadata(objInfo),
NoAuditLog: true,
SrcPoolIdx: idx,
DataMovement: true,
@@ -681,7 +681,7 @@ func (z *erasureServerPools) decommissionObject(ctx context.Context, idx int, bu
SrcPoolIdx: idx,
VersionID: objInfo.VersionID,
MTime: objInfo.ModTime,
UserDefined: objInfo.UserDefined,
UserDefined: migrationObjectMetadata(objInfo),
PreserveETag: objInfo.ETag, // Preserve original ETag to ensure same metadata.
IndexCB: func() []byte {
return objInfo.Parts[0].Index // Preserve part Index to ensure decompression works.
+33
View File
@@ -0,0 +1,33 @@
// Copyright (c) 2026 Feng Ruohang
//
// This file is part of Silo Object Storage stack
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package cmd
import xhttp "github.com/minio/minio/internal/http"
// migrationObjectMetadata restores tags removed from the read representation.
// Keep the original revision, including ordered empty states, for the existing
// locked reconciliation at PUT or multipart completion. Moving is not a new
// tagging mutation, and the write must not modify the reader's metadata map.
func migrationObjectMetadata(oi ObjectInfo) map[string]string {
metadata := cloneMSS(oi.UserDefined)
delete(metadata, xhttp.AmzObjectTagging)
if oi.UserTags != "" {
metadata[xhttp.AmzObjectTagging] = oi.UserTags
}
return metadata
}
+592
View File
@@ -0,0 +1,592 @@
// Copyright (c) 2026 Feng Ruohang
//
// This file is part of Silo Object Storage stack
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"maps"
"strings"
"sync"
"testing"
"time"
xhttp "github.com/minio/minio/internal/http"
)
const migrationTagRevision = ReservedMetadataPrefixLower + TaggingTimestamp
var migrationTagStates = []struct {
name, tags, revision string
hasRevision bool
}{
{name: "initial", tags: "team=storage&path=a%2Fb"},
{name: "empty-revision", tags: "team=storage", hasRevision: true},
{name: "invalid-revision", tags: "team=storage", revision: "invalid", hasRevision: true},
{name: "ordered", tags: "team=storage", revision: "2026-09-16T09:00:00Z", hasRevision: true},
{name: "cleared", revision: "2026-09-16T09:00:00Z", hasRevision: true},
{name: "never-tagged"},
{name: "empty-empty-revision", hasRevision: true},
{name: "empty-invalid-revision", revision: "invalid", hasRevision: true},
}
func TestMigrationObjectMetadata(t *testing.T) {
for _, state := range migrationTagStates {
for _, residual := range []bool{false, true} {
t.Run(fmt.Sprintf("%s/residual=%t", state.name, residual), func(t *testing.T) {
raw := map[string]string{
xhttp.AmzObjectTagging: state.tags,
"x-amz-meta-team": "storage",
"x-amz-meta-empty": "",
ReservedMetadataPrefix + "compression": "opaque-compression-state",
ReservedMetadataPrefix + "sealed-key": "opaque-key-state",
"x-amz-object-lock-retain-until-date": "2030-01-01T00:00:00Z",
"x-amz-object-lock-legal-hold": "ON",
"x-amz-tagging": "keep-noncanonical-metadata",
ReservedMetadataPrefix + "actual-size": "42",
ReservedMetadataPrefix + "replica-status": "opaque-replication-state",
}
if state.hasRevision {
raw[migrationTagRevision] = state.revision
}
oi := (FileInfo{Metadata: raw}).ToObjectInfo("bucket", "object", false)
if residual {
oi.UserDefined[xhttp.AmzObjectTagging] = "stale=raw"
}
before := maps.Clone(oi.UserDefined)
got := migrationObjectMetadata(oi)
if got == nil || !maps.Equal(before, oi.UserDefined) {
t.Fatal("helper must return a new non-nil map without changing the read snapshot")
}
value, exists := got[xhttp.AmzObjectTagging]
if value != state.tags || exists != (state.tags != "") {
t.Fatalf("raw tags=(%q,%t), want (%q,%t)", value, exists, state.tags, state.tags != "")
}
stamp, present := got[migrationTagRevision]
if stamp != state.revision || present != state.hasRevision {
t.Fatalf("revision=(%q,%t), want (%q,%t)", stamp, present, state.revision, state.hasRevision)
}
for key, value := range before {
if stored, exists := got[key]; key != xhttp.AmzObjectTagging && (!exists || stored != value) {
t.Errorf("lost metadata %q", key)
}
}
got["x-amz-meta-team"] = "changed"
if !maps.Equal(before, oi.UserDefined) {
t.Fatal("output aliases the read snapshot")
}
})
}
}
for _, metadata := range []map[string]string{nil, {}} {
got := migrationObjectMetadata(ObjectInfo{UserDefined: metadata})
if got == nil || len(got) != 0 {
t.Fatalf("empty input must yield a writable empty map: %#v", got)
}
got["new"] = "value"
if len(metadata) != 0 {
t.Fatal("empty map input was aliased")
}
}
}
// Allocation ignores the host's used-space percentage, as in the existing tag
// fixtures. All object data and metadata still use real erasure-storage disks.
func migrationTestPools(t *testing.T) (*erasureServerPools, string) {
t.Helper()
z, bucket := consistencyPools(t)
for _, pool := range z.serverPools {
for _, set := range pool.sets {
original := set.getDisks
disks := append([]StorageAPI(nil), original()...)
for i := range disks {
disks[i] = tagTestCapacityDisk{StorageAPI: disks[i]}
}
set.getDisks = func() []StorageAPI { return disks }
t.Cleanup(func() { set.getDisks = original })
}
}
return z, bucket
}
func migrationTestMover(t *testing.T, z *erasureServerPools, source int, kind string) func(context.Context, int, string, *GetObjectReader) error {
t.Helper()
if kind == "rebalance" {
z.rebalMu.Lock()
z.rebalMeta = &rebalanceMeta{PoolStats: []*rebalanceStats{{}, {}}}
z.rebalMeta.PoolStats[source] = &rebalanceStats{Participating: true, Info: rebalanceInfo{Status: rebalStarted}}
z.rebalMu.Unlock()
t.Cleanup(func() { z.rebalMu.Lock(); z.rebalMeta = nil; z.rebalMu.Unlock() })
return z.rebalanceObject
}
z.poolMetaMutex.Lock()
z.poolMeta.Pools[source].Decommission = &PoolDecommissionInfo{}
z.poolMetaMutex.Unlock()
t.Cleanup(func() { z.poolMetaMutex.Lock(); z.poolMeta.Pools[source].Decommission = nil; z.poolMetaMutex.Unlock() })
return z.decommissionObject
}
func migrationTestObject(t *testing.T, z *erasureServerPools, bucket, object string, source int, multipart bool, parts []string, opts ObjectOptions) ObjectInfo {
t.Helper()
if !multipart {
oi := putConsistencyObject(t, z, bucket, object, source, strings.Join(parts, ""), opts)
return migrationTestPersistedInfo(t, z, source, oi)
}
mp, err := z.serverPools[source].NewMultipartUpload(t.Context(), bucket, object, opts)
if err != nil {
t.Fatal(err)
}
completed := make([]CompletePart, len(parts))
for i, body := range parts {
part, err := z.serverPools[source].PutObjectPart(t.Context(), bucket, object, mp.UploadID, i+1,
mustGetPutObjReader(t, strings.NewReader(body), int64(len(body)), "", ""), ObjectOptions{})
if err != nil {
t.Fatal(err)
}
completed[i] = CompletePart{PartNumber: i + 1, ETag: part.ETag}
}
oi, err := z.serverPools[source].CompleteMultipartUpload(t.Context(), bucket, object, mp.UploadID, completed, opts)
if err != nil {
t.Fatal(err)
}
if !oi.isMultipart() {
t.Fatal("fixture did not create a multipart object")
}
return migrationTestPersistedInfo(t, z, source, oi)
}
func migrationTestPersistedInfo(t *testing.T, z *erasureServerPools, pool int, oi ObjectInfo) ObjectInfo {
t.Helper()
// PUT can return transient fields such as tier-free-versionID that are not
// stored. Take the same persisted read representation used by the movers.
got, err := z.serverPools[pool].GetObjectInfo(t.Context(), oi.Bucket, oi.Name, ObjectOptions{VersionID: migrationTestVersion(oi)})
if err != nil {
t.Fatal(err)
}
return got
}
func migrationTestVersion(oi ObjectInfo) string {
if oi.VersionID == "" {
return nullVersionID
}
return oi.VersionID
}
func migrationTestReader(t *testing.T, z *erasureServerPools, source int, oi ObjectInfo) *GetObjectReader {
t.Helper()
gr, err := z.serverPools[source].GetObjectNInfo(t.Context(), oi.Bucket, oi.Name, nil, nil,
ObjectOptions{VersionID: migrationTestVersion(oi), NoLock: true, NoDecryption: true})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { gr.Close() })
return gr
}
func migrationTestStored(t *testing.T, z *erasureServerPools, pool int, original ObjectInfo, tags, revision string, hasRevision bool, body string) {
t.Helper()
version := migrationTestVersion(original)
got, err := z.serverPools[pool].GetObjectInfo(t.Context(), original.Bucket, original.Name, ObjectOptions{VersionID: version})
if err != nil {
t.Fatal(err)
}
if got.UserTags != tags || got.ETag != original.ETag || !got.ModTime.Equal(original.ModTime) || got.VersionID != original.VersionID || got.Size != original.Size {
t.Errorf("pool %d changed tags/object identity: tags=%q want=%q version=%q/%q etag=%q/%q size=%d/%d mtime=%v/%v", pool, got.UserTags, tags, got.VersionID, original.VersionID, got.ETag, original.ETag, got.Size, original.Size, got.ModTime, original.ModTime)
}
infos, errs := readAllFileInfo(t.Context(), z.serverPools[pool].getHashedSet(original.Name).getDisks(), "", original.Bucket, original.Name, version, false, false)
for disk, fi := range infos {
if errs[disk] != nil {
t.Fatalf("pool %d disk %d: %v", pool, disk, errs[disk])
}
stamp, exists := fi.Metadata[migrationTagRevision]
if fi.Metadata[xhttp.AmzObjectTagging] != tags || stamp != revision || exists != hasRevision {
t.Errorf("pool %d disk %d: tags=%q revision=(%q,%t), want %q (%q,%t)", pool, disk, fi.Metadata[xhttp.AmzObjectTagging], stamp, exists, tags, revision, hasRevision)
}
for key, value := range original.UserDefined {
if key != migrationTagRevision && fi.Metadata[key] != value {
t.Errorf("pool %d disk %d lost metadata %q", pool, disk, key)
}
}
}
gr, err := z.serverPools[pool].GetObjectNInfo(t.Context(), original.Bucket, original.Name, nil, nil, ObjectOptions{VersionID: version})
if err != nil {
t.Fatal(err)
}
data, err := io.ReadAll(gr)
gr.Close()
if err != nil || !bytes.Equal(data, []byte(body)) {
t.Fatalf("pool %d content changed: bytes=%d err=%v", pool, len(data), err)
}
}
func TestPoolsMigrationPreservesTagState(t *testing.T) {
z, bucket := migrationTestPools(t)
for _, kind := range []string{"rebalance", "decommission"} {
for _, method := range []string{"put", "multipart"} {
for source := range 2 {
for _, state := range migrationTagStates {
t.Run(fmt.Sprintf("%s/%s/source=%d/%s", kind, method, source, state.name), func(t *testing.T) {
move := migrationTestMover(t, z, source, kind)
meta := map[string]string{xhttp.AmzObjectTagging: state.tags, "x-amz-meta-owner": "retained"}
if state.hasRevision {
meta[migrationTagRevision] = state.revision
}
oi := migrationTestObject(t, z, bucket, t.Name(), source, method == "multipart", []string{"payload"}, ObjectOptions{Versioned: true, UserDefined: meta})
migrationTestStored(t, z, source, oi, state.tags, state.revision, state.hasRevision, "payload")
gr := migrationTestReader(t, z, source, oi)
before := maps.Clone(gr.ObjInfo.UserDefined)
if err := move(t.Context(), source, bucket, gr); err != nil {
t.Fatal(err)
}
if !maps.Equal(before, gr.ObjInfo.UserDefined) {
t.Error("migration modified its source snapshot")
}
migrationTestStored(t, z, 1-source, oi, state.tags, state.revision, state.hasRevision, "payload")
})
}
}
}
}
}
type migrationTestGate struct {
entered, resume chan struct{}
once, release sync.Once
}
func newMigrationTestGate() *migrationTestGate {
return &migrationTestGate{entered: make(chan struct{}), resume: make(chan struct{})}
}
func (g *migrationTestGate) wait(ctx context.Context) error {
g.once.Do(func() { close(g.entered) })
select {
case <-g.resume:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func (g *migrationTestGate) unblock() { g.release.Do(func() { close(g.resume) }) }
type migrationTestGateReader struct {
io.Reader
ctx context.Context
gate *migrationTestGate
}
func (r migrationTestGateReader) Read(p []byte) (int, error) {
if err := r.gate.wait(r.ctx); err != nil {
return 0, err
}
return r.Reader.Read(p)
}
func TestPoolsMigrationRechecksTags(t *testing.T) {
z, bucket := migrationTestPools(t)
const recent = "2026-09-16T10:00:00Z"
for _, kind := range []string{"rebalance", "decommission"} {
for _, method := range []string{"put", "multipart"} {
for source := range 2 {
for _, tags := range []string{"state=after", ""} {
t.Run(fmt.Sprintf("%s/%s/source=%d/clear=%t", kind, method, source, tags == ""), func(t *testing.T) {
move := migrationTestMover(t, z, source, kind)
oi := migrationTestObject(t, z, bucket, t.Name(), source, method == "multipart", []string{"payload"}, ObjectOptions{
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "state=before"},
})
gr := migrationTestReader(t, z, source, oi)
before := maps.Clone(gr.ObjInfo.UserDefined)
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
defer cancel()
gate := newMigrationTestGate()
done, finished := make(chan error, 1), make(chan struct{})
if method == "multipart" {
// The first part read is after persisted upload initialization,
// but before completion takes the object lock and reconciles.
gr.Reader = migrationTestGateReader{Reader: gr.Reader, ctx: ctx, gate: gate}
go func() { defer close(finished); done <- move(ctx, source, bucket, gr) }()
defer func() { gate.unblock(); cancel(); <-finished }()
select {
case <-gate.entered:
case err := <-done:
t.Fatalf("migration finished before the part barrier: %v", err)
case <-ctx.Done():
t.Fatal(ctx.Err())
}
uploads, err := z.serverPools[1-source].listMultipartUploadsExact(ctx, bucket, oi.Name)
if err != nil || len(uploads.Uploads) != 1 {
t.Fatalf("upload not persisted at barrier: %+v, %v", uploads, err)
}
info, err := z.serverPools[1-source].GetMultipartInfo(ctx, bucket, oi.Name, uploads.Uploads[0].UploadID, ObjectOptions{})
if err != nil || info.UserDefined[xhttp.AmzObjectTagging] != "state=before" {
t.Fatalf("upload lost snapshot tags: %v, %v", info.UserDefined, err)
}
}
// For ordinary PUT this must finish before calling the mover:
// blocking its data Read would already hold the object lock.
opts := ObjectOptions{VersionID: migrationTestVersion(oi), UserDefined: map[string]string{migrationTagRevision: recent}}
var err error
if tags == "" {
_, err = z.DeleteObjectTags(ctx, bucket, oi.Name, opts)
} else {
_, err = z.PutObjectTags(ctx, bucket, oi.Name, tags, opts)
}
if err != nil {
t.Fatal(err)
}
if gr.ObjInfo.UserTags != "state=before" || !maps.Equal(before, gr.ObjInfo.UserDefined) {
t.Fatal("test must retain the old source snapshot")
}
migrationTestStored(t, z, source, oi, tags, recent, true, "payload")
if method == "multipart" {
gate.unblock()
err = <-done
} else {
err = move(ctx, source, bucket, gr)
}
if err != nil {
t.Fatal(err)
}
if !maps.Equal(before, gr.ObjInfo.UserDefined) {
t.Error("reconciliation changed the reader's map")
}
migrationTestStored(t, z, 1-source, oi, tags, recent, true, "payload")
})
}
}
}
}
}
type migrationTestMetadataGateDisk struct {
StorageAPI
bucket, object string
gate *migrationTestGate
}
func (d migrationTestMetadataGateDisk) UpdateMetadata(ctx context.Context, volume, path string, fi FileInfo, opts UpdateMetadataOpts) error {
if volume == d.bucket && path == d.object {
if err := d.gate.wait(ctx); err != nil {
return err
}
}
return d.StorageAPI.UpdateMetadata(ctx, volume, path, fi, opts)
}
func TestPoolsMigrationTagsDuringCleanup(t *testing.T) {
z, bucket := migrationTestPools(t)
const recent = "2026-09-16T10:00:00Z"
for _, kind := range []string{"rebalance", "decommission"} {
for source := range 2 {
for _, tags := range []string{"state=after", ""} {
for _, interrupt := range []bool{false, true} {
t.Run(fmt.Sprintf("%s/source=%d/clear=%t/interrupt=%t", kind, source, tags == "", interrupt), func(t *testing.T) {
move := migrationTestMover(t, z, source, kind)
oi := migrationTestObject(t, z, bucket, t.Name(), source, false, []string{"payload"}, ObjectOptions{
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "state=before"},
})
if err := move(t.Context(), source, bucket, migrationTestReader(t, z, source, oi)); err != nil {
t.Fatal(err)
}
migrationTestStored(t, z, 1-source, oi, "state=before", "", false, "payload")
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
defer cancel()
mutate := func() (ObjectInfo, error) {
opts := ObjectOptions{VersionID: migrationTestVersion(oi), UserDefined: map[string]string{migrationTagRevision: recent}}
if tags == "" {
return z.DeleteObjectTags(ctx, bucket, oi.Name, opts)
}
return z.PutObjectTags(ctx, bucket, oi.Name, tags, opts)
}
set := z.serverPools[source].getHashedSet(oi.Name)
cleanup := func() {
// Use the same source prefix cleanup as both outer movers.
if _, err := set.DeleteObject(ctx, bucket, oi.Name, ObjectOptions{DeletePrefix: true, DeletePrefixObject: true}); err != nil {
t.Fatal(err)
}
}
var updated ObjectInfo
if !interrupt {
var err error
updated, err = mutate()
if err != nil {
t.Fatal(err)
}
cleanup()
} else {
gate := newMigrationTestGate()
original := set.getDisks
disks := append([]StorageAPI(nil), original()...)
for i := range disks {
disks[i] = migrationTestMetadataGateDisk{StorageAPI: disks[i], bucket: bucket, object: oi.Name, gate: gate}
}
set.getDisks = func() []StorageAPI { return disks }
defer func() { set.getDisks = original }()
done, finished := make(chan error, 1), make(chan struct{})
go func() { defer close(finished); _, err := mutate(); done <- err }()
defer func() { gate.unblock(); cancel(); <-finished }()
select {
case <-gate.entered:
case err := <-done:
t.Fatalf("metadata update missed the barrier: %v", err)
case <-ctx.Done():
t.Fatal(ctx.Err())
}
cleanup()
gate.unblock()
if err := <-done; err == nil {
t.Fatal("update reported success after its source copy was removed")
}
<-finished
set.getDisks = original
var err error
updated, err = mutate()
if err != nil {
t.Fatalf("metadata retry did not converge: %v", err)
}
}
if _, err := z.serverPools[source].GetObjectInfo(ctx, bucket, oi.Name, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
t.Fatalf("source cleanup left a copy: %v", err)
}
migrationTestStored(t, z, 1-source, oi, tags, updated.UserDefined[migrationTagRevision], true, "payload")
})
}
}
}
}
}
func TestPoolsMigrationVersionHistory(t *testing.T) {
z, bucket := migrationTestPools(t)
for _, kind := range []string{"rebalance", "decommission"} {
for _, method := range []string{"put", "multipart"} {
for _, versioning := range []string{"unversioned", "null", "history"} {
t.Run(kind+"/"+method+"/"+versioning, func(t *testing.T) {
move := migrationTestMover(t, z, 0, kind)
opts := ObjectOptions{
Versioned: versioning == "history", VersionSuspended: versioning == "null",
MTime: time.Date(2026, 9, 16, 8, 0, 0, 0, time.UTC),
UserDefined: map[string]string{xhttp.AmzObjectTagging: "version=old"},
}
versions := []ObjectInfo{migrationTestObject(t, z, bucket, t.Name(), 0, method == "multipart", []string{"old"}, opts)}
if versioning == "history" {
opts.MTime = opts.MTime.Add(time.Minute)
versions = append(versions, migrationTestObject(t, z, bucket, t.Name(), 0, method == "multipart", []string{"old"}, opts))
}
var newest ObjectInfo
if versioning != "unversioned" {
newest = migrationTestObject(t, z, bucket, t.Name(), 0, false, []string{"latest"}, ObjectOptions{
Versioned: true, MTime: opts.MTime.Add(time.Minute),
UserDefined: map[string]string{xhttp.AmzObjectTagging: "version=new"},
})
}
for _, oi := range versions {
if err := move(t.Context(), 0, bucket, migrationTestReader(t, z, 0, oi)); err != nil {
t.Fatal(err)
}
migrationTestStored(t, z, 1, oi, "version=old", "", false, "old")
}
if versioning != "unversioned" {
migrationTestStored(t, z, 0, newest, "version=new", "", false, "latest")
if _, err := z.serverPools[1].GetObjectInfo(t.Context(), bucket, t.Name(), ObjectOptions{VersionID: newest.VersionID}); !isErrVersionNotFound(err) {
t.Fatalf("moving an addressed version affected the newer version: %v", err)
}
}
})
}
}
}
}
func TestPoolsMigrationMultipartParts(t *testing.T) {
z, bucket := migrationTestPools(t)
parts := []string{strings.Repeat("a", 5<<20), "tail-with-tags"}
body := strings.Join(parts, "")
for _, kind := range []string{"rebalance", "decommission"} {
for source := range 2 {
t.Run(fmt.Sprintf("%s/source=%d", kind, source), func(t *testing.T) {
move := migrationTestMover(t, z, source, kind)
oi := migrationTestObject(t, z, bucket, t.Name(), source, true, parts, ObjectOptions{
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "multipart=kept"},
})
if err := move(t.Context(), source, bucket, migrationTestReader(t, z, source, oi)); err != nil {
t.Fatal(err)
}
migrationTestStored(t, z, 1-source, oi, "multipart=kept", "", false, body)
got, err := z.serverPools[1-source].GetObjectInfo(t.Context(), bucket, oi.Name, ObjectOptions{VersionID: oi.VersionID})
if err != nil || len(got.Parts) != len(oi.Parts) {
t.Fatalf("parts changed: %+v, %v", got.Parts, err)
}
for i, part := range got.Parts {
old := oi.Parts[i]
if part.Number != old.Number || part.Size != old.Size || part.ActualSize != old.ActualSize || part.ETag != old.ETag || !bytes.Equal(part.Index, old.Index) {
t.Errorf("part %d changed: %+v / %+v", i, part, old)
}
}
start := int64(len(parts[0]) - 3)
gr, err := z.serverPools[1-source].GetObjectNInfo(t.Context(), bucket, oi.Name, &HTTPRangeSpec{Start: start, End: start + 7}, nil, ObjectOptions{VersionID: oi.VersionID})
if err != nil {
t.Fatal(err)
}
data, err := io.ReadAll(gr)
gr.Close()
if err != nil || string(data) != body[start:start+8] {
t.Fatalf("cross-part range changed: %q, %v", data, err)
}
})
}
}
}
func TestPoolsMigrationUnreadablePool(t *testing.T) {
z, bucket := migrationTestPools(t)
for _, kind := range []string{"rebalance", "decommission"} {
for _, method := range []string{"put", "multipart"} {
t.Run(kind+"/"+method, func(t *testing.T) {
move := migrationTestMover(t, z, 0, kind)
oi := migrationTestObject(t, z, bucket, t.Name(), 0, method == "multipart", []string{"payload"}, ObjectOptions{
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "keep=source"},
})
gr := migrationTestReader(t, z, 0, oi)
set := z.serverPools[0].getHashedSet(oi.Name)
original := set.getDisks
disks := append([]StorageAPI(nil), original()...)
for i := range disks {
disks[i] = consistencyReadFaultDisk{StorageAPI: disks[i], bucket: bucket, object: oi.Name}
}
set.getDisks = func() []StorageAPI { return disks }
defer func() { set.getDisks = original }()
err := move(t.Context(), 0, bucket, gr)
var quorum InsufficientReadQuorum
if !errors.As(err, &quorum) {
t.Fatalf("unreadable pool must fail migration with read quorum error: %v", err)
}
set.getDisks = original
migrationTestStored(t, z, 0, oi, "keep=source", "", false, "payload")
if _, err := z.serverPools[1].GetObjectInfo(t.Context(), bucket, oi.Name, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
t.Fatalf("failed migration committed a target version: %v", err)
}
})
}
}
}
+2 -2
View File
@@ -865,7 +865,7 @@ func (z *erasureServerPools) rebalanceObject(ctx context.Context, poolIdx int, b
if oi.isMultipart() {
res, err := z.NewMultipartUpload(ctx, bucket, oi.Name, ObjectOptions{
VersionID: oi.VersionID,
UserDefined: oi.UserDefined,
UserDefined: migrationObjectMetadata(oi),
NoAuditLog: true,
DataMovement: true,
SrcPoolIdx: poolIdx,
@@ -924,7 +924,7 @@ func (z *erasureServerPools) rebalanceObject(ctx context.Context, poolIdx int, b
DataMovement: true,
VersionID: oi.VersionID,
MTime: oi.ModTime,
UserDefined: oi.UserDefined,
UserDefined: migrationObjectMetadata(oi),
PreserveETag: oi.ETag, // Preserve original ETag to ensure same metadata.
IndexCB: func() []byte {
return oi.Parts[0].Index // Preserve part Index to ensure decompression works.
+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.
+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")
}
})
}
}