Merge remote-tracking branch 'origin/main' into codex/pin-prerelease-components

This commit is contained in:
Feng Ruohang
2026-08-29 20:33:44 +08:00
22 changed files with 1456 additions and 110 deletions
+3 -1
View File
@@ -77,8 +77,10 @@ release:
name: silo
draft: true
prerelease: false
mode: append
replace_existing_draft: true
replace_existing_artifacts: false
mode: replace
# Draft replacement matches the release name; keep it identical to the tag.
name_template: "{{ .Tag }}"
changelog:
+1 -1
View File
@@ -23,7 +23,7 @@ permissions:
artifact-metadata: write
concurrency:
group: finalize-release
group: release-${{ inputs.tag }}
cancel-in-progress: false
jobs:
+34
View File
@@ -1,5 +1,10 @@
name: Release
# Retry contract: an absent or single unfinalized Draft may be rebuilt from
# scratch; a published release or a Draft carrying finalize's GPG-derived
# provenance marker is terminal for this lane. The per-tag lock serializes
# workflows, but a maintainer must not publish the Draft while this job runs.
on:
push:
tags:
@@ -16,6 +21,10 @@ permissions:
attestations: write
artifact-metadata: write
concurrency:
group: release-${{ github.event_name == 'workflow_dispatch' && inputs.tag || github.ref_name }}
cancel-in-progress: false
jobs:
release:
runs-on: ubuntu-latest
@@ -78,6 +87,15 @@ jobs:
echo "Invalid release tag format: ${TAG}" >&2
exit 1
fi
if ! TAG_COMMIT="$(git rev-parse "${TAG}^{commit}" 2>/dev/null)"; then
echo "Release tag ${TAG} does not resolve to a commit" >&2
exit 1
fi
HEAD_COMMIT="$(git rev-parse HEAD)"
if [ "${TAG_COMMIT}" != "${HEAD_COMMIT}" ]; then
echo "Release tag ${TAG} resolves to ${TAG_COMMIT}, checkout is ${HEAD_COMMIT}" >&2
exit 1
fi
VERSION_HYPHEN="${TAG#RELEASE.}"
PKG_VERSION="$(echo "${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/')"
VERSION_COLON="$(echo "${VERSION_HYPHEN}" | sed -E 's/T([0-9]{2})-([0-9]{2})-([0-9]{2})Z$/T\1:\2:\3Z/')"
@@ -93,6 +111,13 @@ jobs:
echo "Package version: ${PKG_VERSION}"
echo "LDFLAGS: ${LDFLAGS}"
- name: Check existing release state
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
run: |
set -euo pipefail
buildscripts/check-release-state.sh "${RELEASE_TAG}"
# Both installer actions are pinned to immutable commits. The explicit
# tool versions keep the release format reproducible across workflow
# reruns while the installers verify the downloaded executables.
@@ -113,6 +138,7 @@ jobs:
args: release --clean --skip=validate --config .github/goreleaser.yml
env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
GORELEASER_CURRENT_TAG: ${{ env.RELEASE_TAG }}
LDFLAGS: ${{ env.LDFLAGS }}
PKG_VERSION: ${{ env.PKG_VERSION }}
@@ -189,6 +215,14 @@ jobs:
test -s "${BUNDLE_PATH}"
cp "${BUNDLE_PATH}" "dist/silo_${PKG_VERSION}_provenance.sigstore.json"
- name: Confirm unfinalized Draft release state
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
REQUIRE_DRAFT: "true"
run: |
set -euo pipefail
buildscripts/check-release-state.sh "${RELEASE_TAG}"
- name: Upload nFPM packages to Draft release
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
+7
View File
@@ -23,6 +23,8 @@ on:
- "buildscripts/minio-upgrade.sh"
- "buildscripts/sign-release-rpms.sh"
- "buildscripts/verify-build-provenance.sh"
- "buildscripts/check-release-state.sh"
- "buildscripts/check-release-state_test.sh"
- "buildscripts/verify-rebrand.sh"
- "buildscripts/verify-helm-migration.sh"
- "buildscripts/helm-migration-guard/**"
@@ -479,6 +481,8 @@ jobs:
bash -n buildscripts/minio-upgrade.sh
bash -n buildscripts/sign-release-rpms.sh
bash -n buildscripts/verify-build-provenance.sh
bash -n buildscripts/check-release-state.sh
bash -n buildscripts/check-release-state_test.sh
bash -n buildscripts/verify-rebrand.sh
bash -n buildscripts/verify-helm-migration.sh
sh -n buildscripts/package/postinstall.sh
@@ -493,9 +497,12 @@ jobs:
test -x buildscripts/package-release.sh
test -x buildscripts/sign-release-rpms.sh
test -x buildscripts/verify-build-provenance.sh
test -x buildscripts/check-release-state.sh
test -x buildscripts/check-release-state_test.sh
test -x buildscripts/verify-rebrand.sh
test -x buildscripts/verify-helm-migration.sh
test -x buildscripts/package/postinstall.sh
test -x buildscripts/package/preremove.sh
test -x buildscripts/package/lifecycle_test.sh
test -x dockerscripts/docker-entrypoint_test.sh
buildscripts/check-release-state_test.sh
+76
View File
@@ -0,0 +1,76 @@
#!/usr/bin/env bash
# Fail closed before a release job can replace published or finalized assets.
# An ordinary Draft is retry state; a finalized Draft contains GPG-derived
# materials and must never be replaced by the build lane.
set -euo pipefail
release_tag="${1:-}"
fixture="${2:-}"
repository="${GITHUB_REPOSITORY:-pgsty/silo}"
require_draft="${REQUIRE_DRAFT:-false}"
if ! command -v jq >/dev/null 2>&1; then
echo "jq is required to inspect GitHub release state" >&2
exit 1
fi
if [[ ! "${release_tag}" =~ ^RELEASE\.[0-9]{4}-[0-9]{2}-[0-9]{2}T[0-9]{2}-[0-9]{2}-[0-9]{2}Z$ ]]; then
echo "Invalid release tag format: ${release_tag:-<empty>}" >&2
exit 1
fi
if [ -n "${fixture}" ]; then
release_json="$(<"${fixture}")"
else
error_file="$(mktemp)"
trap 'rm -f "${error_file}"' EXIT
if ! release_json="$(
gh api --paginate "repos/${repository}/releases?per_page=100" --jq '.[]' 2>"${error_file}" |
jq --arg tag "${release_tag}" -s '[.[] | select(.tag_name == $tag)]'
)"; then
cat "${error_file}" >&2
exit 1
fi
fi
if ! jq -e 'type == "array" and all(.[]; type == "object" and (.tag_name | type == "string") and (.draft | type == "boolean"))' \
<<<"${release_json}" >/dev/null 2>&1; then
echo "Invalid release state response for ${release_tag}" >&2
exit 1
fi
if ! jq -e --arg tag "${release_tag}" 'all(.[]; .tag_name == $tag)' \
<<<"${release_json}" >/dev/null 2>&1; then
echo "Release state returned a tag other than ${release_tag}" >&2
exit 1
fi
release_count="$(jq 'length' <<<"${release_json}")"
if [ "${release_count}" -eq 0 ]; then
if [ "${require_draft}" = "true" ]; then
echo "Expected one Draft release for ${release_tag}, found none" >&2
exit 1
fi
echo "No existing release for ${release_tag}."
exit 0
fi
if [ "${release_count}" -ne 1 ]; then
echo "Refusing to choose among ${release_count} releases for ${release_tag}; clean duplicate Drafts first" >&2
exit 1
fi
if [ "$(jq -r '.[0].draft' <<<"${release_json}")" != "true" ]; then
echo "Refusing to overwrite published release ${release_tag}" >&2
exit 1
fi
finalize_markers="$(jq '[.[0].assets[]? | select(.name | endswith("_packages_provenance.sigstore.json"))] | length' <<<"${release_json}")"
if [ "${finalize_markers}" -ne 0 ]; then
echo "Refusing to replace finalized Draft ${release_tag}" >&2
exit 1
fi
echo "Existing unfinalized Draft ${release_tag} will be replaced from scratch."
+68
View File
@@ -0,0 +1,68 @@
#!/usr/bin/env bash
set -euo pipefail
script_dir="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
checker="${script_dir}/check-release-state.sh"
tag="RELEASE.2026-08-29T00-00-00Z"
fixture="$(mktemp)"
stdout_file="$(mktemp)"
stderr_file="$(mktemp)"
trap 'rm -f "${fixture}" "${stdout_file}" "${stderr_file}"' EXIT
expect_success() {
if ! "${checker}" "$@" >"${stdout_file}" 2>"${stderr_file}"; then
cat "${stderr_file}" >&2
return 1
fi
}
expect_failure() {
if "${checker}" "$@" >"${stdout_file}" 2>"${stderr_file}"; then
echo "Expected release-state check to fail: $*" >&2
return 1
fi
}
printf '[]\n' >"${fixture}"
expect_success "${tag}" "${fixture}"
grep -qF "No existing release for ${tag}." "${stdout_file}"
if REQUIRE_DRAFT=true "${checker}" "${tag}" "${fixture}" >"${stdout_file}" 2>"${stderr_file}"; then
echo "Expected required-Draft check to fail when no release exists" >&2
exit 1
fi
grep -qF "Expected one Draft release for ${tag}, found none" "${stderr_file}"
printf '[{"tag_name":"%s","draft":true,"assets":[]}]\n' "${tag}" >"${fixture}"
expect_success "${tag}" "${fixture}"
grep -qF "Existing unfinalized Draft ${tag} will be replaced from scratch." "${stdout_file}"
if ! REQUIRE_DRAFT=true "${checker}" "${tag}" "${fixture}" >"${stdout_file}" 2>"${stderr_file}"; then
cat "${stderr_file}" >&2
exit 1
fi
printf '[{"tag_name":"%s","draft":true,"assets":[{"name":"silo_20260829000000.0.0_packages_provenance.sigstore.json"}]}]\n' "${tag}" >"${fixture}"
expect_failure "${tag}" "${fixture}"
grep -qF "Refusing to replace finalized Draft ${tag}" "${stderr_file}"
printf '[{"tag_name":"%s","draft":false}]\n' "${tag}" >"${fixture}"
expect_failure "${tag}" "${fixture}"
grep -qF "Refusing to overwrite published release ${tag}" "${stderr_file}"
printf '[{"tag_name":"%s","draft":true},{"tag_name":"%s","draft":true}]\n' "${tag}" "${tag}" >"${fixture}"
expect_failure "${tag}" "${fixture}"
grep -qF "Refusing to choose among 2 releases" "${stderr_file}"
printf '[{"tag_name":"RELEASE.2026-08-28T00-00-00Z","draft":true}]\n' >"${fixture}"
expect_failure "${tag}" "${fixture}"
grep -qF "other than ${tag}" "${stderr_file}"
printf '{not-json}\n' >"${fixture}"
expect_failure "${tag}" "${fixture}"
grep -qF "Invalid release state response for ${tag}" "${stderr_file}"
expect_failure "not-a-release-tag" "${fixture}"
grep -qF "Invalid release tag format" "${stderr_file}"
echo "release-state decision tests passed"
+1 -1
View File
@@ -255,7 +255,7 @@ func (a adminAPIHandlers) SRPeerReplicateBucketItem(w http.ResponseWriter, r *ht
case madmin.SRBucketMetaTypeTags:
err = globalSiteReplicationSys.PeerBucketTaggingHandler(ctx, item.Bucket, item.Tags, item.UpdatedAt)
case madmin.SRBucketMetaTypeObjectLockConfig:
err = globalSiteReplicationSys.PeerBucketObjectLockConfigHandler(ctx, item.Bucket, item.ObjectLockConfig, item.UpdatedAt)
err = globalSiteReplicationSys.peerBucketObjectLockConfigItem(ctx, item)
case madmin.SRBucketMetaTypeSSEConfig:
err = globalSiteReplicationSys.PeerBucketSSEConfigHandler(ctx, item.Bucket, item.SSEConfig, item.UpdatedAt)
case madmin.SRBucketMetaTypeCorsConfig:
+13 -11
View File
@@ -785,19 +785,21 @@ func corsHandler(handler http.Handler) http.Handler {
}
globalCors := cors.New(opts).Handler(handler)
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if bucket, _ := request2BucketObjectName(r); bucket != "" && globalBucketMetadataSys != nil {
cfg, _, err := globalBucketMetadataSys.GetCorsConfig(bucket)
if err == nil && cfg != nil {
if applyBucketCors(w, r, cfg) {
if r.Header.Get("Origin") != "" {
if bucket, _ := request2BucketObjectName(r); bucket != "" && globalBucketMetadataSys != nil {
cfg, _, err := globalBucketMetadataSys.GetCorsConfig(bucket)
if err == nil && cfg != nil {
if applyBucketCors(w, r, cfg) {
return
}
handler.ServeHTTP(w, r)
return
}
if err != nil && !errors.Is(err, errConfigNotFound) {
internalLogOnceIf(r.Context(), err, "bucket-cors-metadata")
handler.ServeHTTP(w, r)
return
}
handler.ServeHTTP(w, r)
return
}
if err != nil && !errors.Is(err, errConfigNotFound) && r.Header.Get("Origin") != "" {
internalLogOnceIf(r.Context(), err, "bucket-cors-metadata")
handler.ServeHTTP(w, r)
return
}
}
globalCors.ServeHTTP(w, r)
+95 -1
View File
@@ -18,15 +18,27 @@
package cmd
import (
"context"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"github.com/minio/minio/internal/auth"
"github.com/minio/minio/internal/bucket/cors"
)
type corsLookupCountingObjectLayer struct {
ObjectLayer
getObjectNInfoCalls atomic.Int64
}
func (o *corsLookupCountingObjectLayer) GetObjectNInfo(ctx context.Context, bucket, object string, rs *HTTPRangeSpec, h http.Header, opts ObjectOptions) (*GetObjectReader, error) {
o.getObjectNInfoCalls.Add(1)
return o.ObjectLayer.GetObjectNInfo(ctx, bucket, object, rs, h, opts)
}
func TestPerBucketCorsPreflight(t *testing.T) {
cfg := &cors.Config{CORSRules: []cors.Rule{{
AllowedOrigins: []string{"http://example.com"},
@@ -245,8 +257,13 @@ func TestPerBucketCorsOriginPatternResponse(t *testing.T) {
func TestBucketCorsMetadataErrorFailsClosed(t *testing.T) {
oldObjectAPI := newObjectLayerFn()
oldMetadataSys := globalBucketMetadataSys
setObjectLayer(nil)
defer setObjectLayer(oldObjectAPI)
globalBucketMetadataSys = NewBucketMetadataSys()
defer func() {
setObjectLayer(oldObjectAPI)
globalBucketMetadataSys = oldMetadataSys
}()
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNoContent)
@@ -273,6 +290,55 @@ func TestBucketCorsMetadataErrorFailsClosed(t *testing.T) {
}
}
func TestBucketCorsSkipsMetadataLookupWithoutOrigin(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testBucketCorsSkipsMetadataLookupWithoutOrigin,
endpoints: []string{"GetObject"},
})
}
func testBucketCorsSkipsMetadataLookupWithoutOrigin(obj ObjectLayer, _ string, _ string, _ http.Handler, _ auth.Credentials, t *testing.T) {
oldObjectAPI := newObjectLayerFn()
counting := &corsLookupCountingObjectLayer{ObjectLayer: obj}
setObjectLayer(counting)
defer setObjectLayer(oldObjectAPI)
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNoContent)
}))
rec := httptest.NewRecorder()
wrapped.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, getGetObjectURL("", "api", "v1/login"), nil))
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusNoContent)
}
requireCorsOriginVary(t, rec.Header())
if got := counting.getObjectNInfoCalls.Load(); got != 0 {
t.Fatalf("request without Origin performed %d bucket metadata reads", got)
}
}
func TestBucketCorsOriginlessPreflightShapeUsesGlobalHandler(t *testing.T) {
nextCalled := false
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
nextCalled = true
w.WriteHeader(http.StatusTeapot)
}))
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodOptions, getGetObjectURL("", "api", "v1/login"), nil)
req.Header.Set("Access-Control-Request-Method", http.MethodGet)
wrapped.ServeHTTP(rec, req)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusNoContent)
}
if nextCalled {
t.Fatal("originless preflight-shaped OPTIONS reached the application handler")
}
requireCorsOriginVary(t, rec.Header())
}
func TestBucketCorsNoConfigUsesGlobalFallback(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
@@ -281,6 +347,34 @@ func TestBucketCorsNoConfigUsesGlobalFallback(t *testing.T) {
})
}
func TestBucketCorsMissingBucketUsesGlobalFallback(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testBucketCorsMissingBucketUsesGlobalFallback,
endpoints: []string{"GetBucketCors"},
})
}
func testBucketCorsMissingBucketUsesGlobalFallback(_ ObjectLayer, _ string, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) {
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNotFound)
}))
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, getGetObjectURL("", bucket+"-missing", "object"), nil)
req.Header.Set("Origin", "https://app.example.com")
wrapped.ServeHTTP(rec, req)
if rec.Code != http.StatusNotFound {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusNotFound)
}
if got := rec.Header().Get("Access-Control-Allow-Origin"); got != "https://app.example.com" {
t.Fatalf("allow-origin = %q", got)
}
if got := rec.Header().Get("Access-Control-Allow-Credentials"); got != "true" {
t.Fatalf("allow-credentials = %q", got)
}
}
func testBucketCorsNoConfigUsesGlobalFallback(_ ObjectLayer, _ string, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) {
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNoContent)
+1 -6
View File
@@ -1844,12 +1844,7 @@ func (api objectAPIHandlers) PutBucketObjectLockConfigHandler(w http.ResponseWri
// We encode the xml bytes as base64 to ensure there are no encoding
// errors.
cfgStr := base64.StdEncoding.EncodeToString(configData)
replLogIf(ctx, globalSiteReplicationSys.BucketMetaHook(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeObjectLockConfig,
Bucket: bucket,
ObjectLockConfig: &cfgStr,
UpdatedAt: updatedAt,
}))
replLogIf(ctx, globalSiteReplicationSys.BucketMetaHook(ctx, newSRBucketObjectLockMeta(bucket, &cfgStr, updatedAt)))
// Write success response.
writeSuccessResponseHeadersOnly(w)
+4 -1
View File
@@ -344,7 +344,10 @@ func (b *BucketMetadata) parseAllConfigs(ctx context.Context, objectAPI ObjectLa
}
if bytes.Equal(b.ObjectLockConfigXML, enabledBucketObjectLockConfig) {
b.VersioningConfigXML = enabledBucketVersioningConfig
config, versioningErr := versioning.ParseConfig(bytes.NewReader(b.VersioningConfigXML))
if versioningErr != nil || !config.Enabled() {
b.VersioningConfigXML = enabledBucketVersioningConfig
}
}
if len(b.ObjectLockConfigXML) != 0 {
+23 -7
View File
@@ -516,17 +516,33 @@ func testAPICompleteMultipartChecksumTypeMismatch(obj ObjectLayer, instanceType,
}
})
t.Run("crc64nvme-composite-completion-is-rejected", func(t *testing.T) {
t.Run("crc64nvme-composite-is-rejected", func(t *testing.T) {
crc64Type := hash.ChecksumCRC64NVME
objectName := "type-mismatch/crc64nvme-composite"
uploadID := newMultipartUploadHTTP(t, apiRouter, credentials, bucketName, objectName,
crc64Type.String(), xhttp.AmzChecksumTypeComposite)
etags := uploadPartsHTTP(t, apiRouter, credentials, bucketName, objectName, uploadID, crc64Type, partData)
rec := completeMultipartUploadHTTP(t, apiRouter, credentials, bucketName, objectName, uploadID, etags, nil,
map[string]string{
crc64Type.Key(): mustChecksum(t, crc64Type, full),
req, err := newTestSignedRequestV4(http.MethodPost, getNewMultipartURL("", bucketName, objectName),
0, nil, credentials.AccessKey, credentials.SecretKey, map[string]string{
xhttp.AmzChecksumAlgo: crc64Type.String(),
xhttp.AmzChecksumType: xhttp.AmzChecksumTypeComposite,
})
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
if rec.Code != http.StatusBadRequest || apiErrorCode(t, rec) != "InvalidArgument" {
t.Fatalf("%s: CRC64NVME/COMPOSITE returned %d %s", instanceType, rec.Code, rec.Body.String())
}
})
t.Run("crc64nvme-composite-completion-is-rejected", func(t *testing.T) {
crc64Type := hash.ChecksumCRC64NVME
objectName := "type-mismatch/crc64nvme-composite-completion"
uploadID := newMultipartUploadHTTP(t, apiRouter, credentials, bucketName, objectName,
crc64Type.String(), xhttp.AmzChecksumTypeFullObject)
etags := uploadPartsHTTP(t, apiRouter, credentials, bucketName, objectName, uploadID, crc64Type, partData)
partCS := []string{mustChecksum(t, crc64Type, partData[0]), mustChecksum(t, crc64Type, partData[1])}
rec := completeMultipartUploadHTTP(t, apiRouter, credentials, bucketName, objectName, uploadID, etags, partCS,
map[string]string{xhttp.AmzChecksumType: xhttp.AmzChecksumTypeComposite})
if rec.Code != http.StatusBadRequest || apiErrorCode(t, rec) != "BadDigest" {
t.Fatalf("%s: CRC64NVME composite completion returned %d %s", instanceType, rec.Code, rec.Body.String())
}
+3
View File
@@ -439,6 +439,9 @@ func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[strin
// get ObjectOptions for Copy calls with encryption headers provided on the target side and source side metadata
func copyDstOpts(ctx context.Context, r *http.Request, bucket, object string, metadata map[string]string) (opts ObjectOptions, err error) {
if _, err := hash.GetContentChecksum(r.Header); err != nil {
return opts, err
}
return putOptsFromReq(ctx, r, bucket, object, metadata)
}
+137
View File
@@ -0,0 +1,137 @@
// Copyright (c) 2015-2026 MinIO, Inc.
//
// This file is part of MinIO 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"
"encoding/base64"
"encoding/xml"
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/minio/minio/internal/auth"
xhttp "github.com/minio/minio/internal/http"
)
func TestAPIRejectsUnsupportedChecksumHeaders(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPIRejectsUnsupportedChecksumHeaders,
endpoints: []string{"CopyObject", "NewMultipart", "PutObject", "PutObjectPart"},
})
}
func testAPIRejectsUnsupportedChecksumHeaders(obj ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
) {
data := []byte("unsupported-checksum")
unsupportedValue := base64.StdEncoding.EncodeToString(make([]byte, 64))
put := func(object string, headers map[string]string) *httptest.ResponseRecorder {
t.Helper()
req, err := newTestSignedRequestV4(http.MethodPut, getPutObjectURL("", bucketName, object),
int64(len(data)), bytes.NewReader(data), credentials.AccessKey, credentials.SecretKey, headers)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
return rec
}
assertRejected := func(name string, rec *httptest.ResponseRecorder) {
t.Helper()
if rec.Code != http.StatusBadRequest || !strings.Contains(rec.Body.String(), "<Code>InvalidArgument</Code>") {
t.Fatalf("%s: %s returned %d, want InvalidArgument: %s", instanceType, name, rec.Code, rec.Body.String())
}
}
for _, algorithm := range []string{"md5", "sha512", "xxhash64", "xxhash3", "xxhash128", "future"} {
object := "checksums/unsupported-" + algorithm
assertRejected(algorithm, put(object, map[string]string{
"x-amz-sdk-checksum-algorithm": "SHA512",
"x-amz-checksum-" + algorithm: unsupportedValue,
}))
if _, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{}); !isErrObjectNotFound(err) {
t.Fatalf("%s: rejected %s checksum stored an object: %v", instanceType, algorithm, err)
}
}
assertRejected("unsupported trailer", put("checksums/unsupported-trailer", map[string]string{
xhttp.AmzTrailer: "x-amz-checksum-sha512",
}))
newMultipart := func(name string, headers map[string]string) *httptest.ResponseRecorder {
t.Helper()
req, err := newTestSignedRequestV4(http.MethodPost, getNewMultipartURL("", bucketName, name),
0, nil, credentials.AccessKey, credentials.SecretKey, headers)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
return rec
}
assertRejected("NewMultipartUpload value header", newMultipart("checksums/mp-value", map[string]string{
"x-amz-checksum-sha512": unsupportedValue,
}))
assertRejected("NewMultipartUpload trailer", newMultipart("checksums/mp-trailer", map[string]string{
xhttp.AmzTrailer: "x-amz-checksum-sha512",
}))
rec := newMultipart("checksums/mp-part", nil)
if rec.Code != http.StatusOK {
t.Fatalf("%s: NewMultipartUpload setup returned %d: %s", instanceType, rec.Code, rec.Body.String())
}
var initiated InitiateMultipartUploadResponse
if err := xml.Unmarshal(rec.Body.Bytes(), &initiated); err != nil {
t.Fatal(err)
}
req, err := newTestSignedRequestV4(http.MethodPut,
getPutObjectPartURL("", bucketName, "checksums/mp-part", initiated.UploadID, "1"),
int64(len(data)), bytes.NewReader(data), credentials.AccessKey, credentials.SecretKey,
map[string]string{"x-amz-checksum-sha512": unsupportedValue})
if err != nil {
t.Fatal(err)
}
rec = httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
assertRejected("UploadPart", rec)
parts, err := obj.ListObjectParts(t.Context(), bucketName, "checksums/mp-part", initiated.UploadID, 0, 1000, ObjectOptions{})
if err != nil {
t.Fatal(err)
}
if len(parts.Parts) != 0 {
t.Fatalf("%s: rejected UploadPart stored %d parts", instanceType, len(parts.Parts))
}
if err := obj.AbortMultipartUpload(t.Context(), bucketName, "checksums/mp-part", initiated.UploadID, ObjectOptions{}); err != nil {
t.Fatal(err)
}
source := "checksums/source"
putCopyChecksumSource(t, apiRouter, credentials, bucketName, source, data, nil)
rec = copyChecksumRequest(t, apiRouter, credentials, bucketName, source, "checksums/copy", map[string]string{
"x-amz-checksum-sha512": unsupportedValue,
})
assertRejected("CopyObject", rec)
if _, err := obj.GetObjectInfo(t.Context(), bucketName, "checksums/copy", ObjectOptions{}); !isErrObjectNotFound(err) {
t.Fatalf("%s: rejected CopyObject stored a destination: %v", instanceType, err)
}
}
+79
View File
@@ -0,0 +1,79 @@
// Copyright (c) 2015-2026 MinIO, Inc.
//
// This file is part of MinIO Object Storage stack
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"net/http"
"net/http/httptest"
"testing"
"github.com/minio/minio/internal/auth"
"github.com/minio/minio/internal/hash"
xhttp "github.com/minio/minio/internal/http"
)
func TestAPIPutObjectRejectsCRC64Composite(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPIPutObjectRejectsCRC64Composite,
endpoints: []string{"PutObject"},
})
}
func testAPIPutObjectRejectsCRC64Composite(obj ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
) {
data := []byte("crc64-composite")
object := "checksums/crc64-composite"
req, err := newTestSignedRequestV4(http.MethodPut, getPutObjectURL("", bucketName, object),
int64(len(data)), bytes.NewReader(data), credentials.AccessKey, credentials.SecretKey, map[string]string{
xhttp.AmzChecksumCRC64NVME: mustChecksum(t, hash.ChecksumCRC64NVME, data),
xhttp.AmzChecksumType: xhttp.AmzChecksumTypeComposite,
})
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
if rec.Code != http.StatusBadRequest || apiErrorCode(t, rec) != "InvalidArgument" {
t.Fatalf("%s: CRC64NVME/COMPOSITE PutObject returned %d %s", instanceType, rec.Code, rec.Body.String())
}
if _, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{}); !isErrObjectNotFound(err) {
t.Fatalf("%s: rejected PutObject stored an object: %v", instanceType, err)
}
trailerObject := "checksums/crc64-composite-trailer"
req, err = newTestSignedRequestV4(http.MethodPut, getPutObjectURL("", bucketName, trailerObject),
int64(len(data)), bytes.NewReader(data), credentials.AccessKey, credentials.SecretKey, map[string]string{
xhttp.AmzTrailer: xhttp.AmzChecksumCRC64NVME,
xhttp.AmzChecksumType: xhttp.AmzChecksumTypeComposite,
})
if err != nil {
t.Fatal(err)
}
rec = httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
if rec.Code != http.StatusBadRequest || apiErrorCode(t, rec) != "InvalidArgument" {
t.Fatalf("%s: trailing CRC64NVME/COMPOSITE PutObject returned %d %s", instanceType, rec.Code, rec.Body.String())
}
if _, err := obj.GetObjectInfo(t.Context(), bucketName, trailerObject, ObjectOptions{}); !isErrObjectNotFound(err) {
t.Fatalf("%s: rejected trailing PutObject stored an object: %v", instanceType, err)
}
}
+4
View File
@@ -309,6 +309,10 @@ func (api objectAPIHandlers) NewMultipartUploadHandler(w http.ResponseWriter, r
}
}
if _, err := hash.GetContentChecksum(r.Header); err != nil {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrInvalidChecksum), r.URL)
return
}
checksumType := hash.NewChecksumHeader(r.Header)
if checksumType.Is(hash.ChecksumInvalid) {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrInvalidChecksum), r.URL)
@@ -0,0 +1,190 @@
// Copyright (c) 2015-2026 MinIO, Inc.
//
// This file is part of MinIO Object Storage stack
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"net/http"
"testing"
"time"
"github.com/minio/minio/internal/auth"
)
func TestPeerBucketAdoptionPreservesLockAndVersioningConfigs(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testPeerBucketAdoptionPreservesLockAndVersioningConfigs,
makeBucketOptions: MakeBucketOptions{LockEnabled: true},
})
}
func testPeerBucketAdoptionPreservesLockAndVersioningConfigs(_ ObjectLayer, instanceType, bucketName string,
_ http.Handler, _ auth.Credentials, t *testing.T,
) {
objectLockXML := []byte(`<ObjectLockConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><ObjectLockEnabled>Enabled</ObjectLockEnabled><Rule><DefaultRetention><Mode>GOVERNANCE</Mode><Days>30</Days></DefaultRetention></Rule></ObjectLockConfiguration>`)
versioningXML := []byte(`<VersioningConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Status>Enabled</Status><ExcludeFolders>true</ExcludeFolders><ExcludedPrefixes><Prefix>temporary/</Prefix></ExcludedPrefixes></VersioningConfiguration>`)
if _, err := globalBucketMetadataSys.Update(t.Context(), bucketName, objectLockConfig, objectLockXML); err != nil {
t.Fatal(err)
}
if _, err := globalBucketMetadataSys.Update(t.Context(), bucketName, bucketVersioningConfig, versioningXML); err != nil {
t.Fatal(err)
}
before, err := globalBucketMetadataSys.Get(bucketName)
if err != nil {
t.Fatal(err)
}
if err := globalSiteReplicationSys.PeerBucketMakeWithVersioningHandler(t.Context(), bucketName, MakeBucketOptions{
CreatedAt: before.Created.Add(-time.Hour),
LockEnabled: true,
}); err != nil {
t.Fatalf("%s: adopting existing bucket failed: %v", instanceType, err)
}
after, err := globalBucketMetadataSys.Get(bucketName)
if err != nil {
t.Fatal(err)
}
if !bytes.Equal(after.ObjectLockConfigXML, before.ObjectLockConfigXML) || !after.ObjectLockConfigUpdatedAt.Equal(before.ObjectLockConfigUpdatedAt) {
t.Fatalf("%s: Object Lock config changed during adoption", instanceType)
}
if !bytes.Equal(after.VersioningConfigXML, before.VersioningConfigXML) || !after.VersioningConfigUpdatedAt.Equal(before.VersioningConfigUpdatedAt) {
t.Fatalf("%s: versioning config changed during adoption", instanceType)
}
}
func TestPeerBucketAdoptionBootstrapsMissingConfigs(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testPeerBucketAdoptionBootstrapsMissingConfigs,
})
}
func TestPeerBucketAdoptionPreservesCustomVersioningWhenEnablingLock(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testPeerBucketAdoptionPreservesCustomVersioningWhenEnablingLock,
})
}
func testPeerBucketAdoptionPreservesCustomVersioningWhenEnablingLock(_ ObjectLayer, instanceType, bucketName string,
_ http.Handler, _ auth.Credentials, t *testing.T,
) {
versioningXML := []byte(`<VersioningConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Status>Enabled</Status><ExcludeFolders>true</ExcludeFolders><ExcludedPrefixes><Prefix>temporary/</Prefix></ExcludedPrefixes></VersioningConfiguration>`)
if _, err := globalBucketMetadataSys.Update(t.Context(), bucketName, bucketVersioningConfig, versioningXML); err != nil {
t.Fatal(err)
}
before, err := globalBucketMetadataSys.Get(bucketName)
if err != nil {
t.Fatal(err)
}
if err := globalSiteReplicationSys.PeerBucketMakeWithVersioningHandler(t.Context(), bucketName, MakeBucketOptions{
CreatedAt: before.Created,
LockEnabled: true,
}); err != nil {
t.Fatal(err)
}
after, err := globalBucketMetadataSys.Get(bucketName)
if err != nil {
t.Fatal(err)
}
if !bytes.Equal(after.VersioningConfigXML, before.VersioningConfigXML) || !after.VersioningConfigUpdatedAt.Equal(before.VersioningConfigUpdatedAt) {
t.Fatalf("%s: custom versioning changed while enabling Object Lock", instanceType)
}
if !bytes.Equal(after.ObjectLockConfigXML, enabledBucketObjectLockConfig) {
t.Fatalf("%s: Object Lock was not bootstrapped", instanceType)
}
}
func TestPeerBucketAdoptionEnablesSuspendedVersioning(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testPeerBucketAdoptionEnablesSuspendedVersioning,
})
}
func testPeerBucketAdoptionEnablesSuspendedVersioning(_ ObjectLayer, instanceType, bucketName string,
_ http.Handler, _ auth.Credentials, t *testing.T,
) {
suspended := []byte(`<VersioningConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Status>Suspended</Status></VersioningConfiguration>`)
if _, err := globalBucketMetadataSys.Update(t.Context(), bucketName, bucketVersioningConfig, suspended); err != nil {
t.Fatal(err)
}
before, err := globalBucketMetadataSys.Get(bucketName)
if err != nil {
t.Fatal(err)
}
if err := globalSiteReplicationSys.PeerBucketMakeWithVersioningHandler(t.Context(), bucketName, MakeBucketOptions{CreatedAt: before.Created}); err != nil {
t.Fatal(err)
}
after, err := globalBucketMetadataSys.Get(bucketName)
if err != nil {
t.Fatal(err)
}
if after.versioningConfig == nil || !after.versioningConfig.Enabled() {
t.Fatalf("%s: versioning remained disabled: %q", instanceType, after.VersioningConfigXML)
}
if !after.VersioningConfigUpdatedAt.After(before.VersioningConfigUpdatedAt) {
t.Fatalf("%s: versioning update time = %v, want after %v", instanceType, after.VersioningConfigUpdatedAt, before.VersioningConfigUpdatedAt)
}
}
func TestEnablePeerBucketVersioningRepairsInvalidConfig(t *testing.T) {
meta := newBucketMetadata("bucket")
meta.Created = time.Date(2026, time.August, 29, 8, 0, 0, 0, time.UTC)
meta.VersioningConfigXML = []byte(`<VersioningConfiguration>`)
if err := enablePeerBucketVersioning(&meta); err != nil {
t.Fatal(err)
}
if !bytes.Equal(meta.VersioningConfigXML, enabledBucketVersioningConfig) || meta.VersioningConfigUpdatedAt.IsZero() {
t.Fatalf("invalid versioning was not repaired: xml=%q updatedAt=%v", meta.VersioningConfigXML, meta.VersioningConfigUpdatedAt)
}
}
func testPeerBucketAdoptionBootstrapsMissingConfigs(_ ObjectLayer, instanceType, bucketName string,
_ http.Handler, _ auth.Credentials, t *testing.T,
) {
before, err := globalBucketMetadataSys.Get(bucketName)
if err != nil {
t.Fatal(err)
}
if len(before.ObjectLockConfigXML) != 0 || len(before.VersioningConfigXML) != 0 {
t.Fatalf("%s: invalid bootstrap precondition", instanceType)
}
if err := globalSiteReplicationSys.PeerBucketMakeWithVersioningHandler(t.Context(), bucketName, MakeBucketOptions{
CreatedAt: before.Created,
LockEnabled: true,
}); err != nil {
t.Fatalf("%s: adopting existing bucket failed: %v", instanceType, err)
}
after, err := globalBucketMetadataSys.Get(bucketName)
if err != nil {
t.Fatal(err)
}
if !bytes.Equal(after.ObjectLockConfigXML, enabledBucketObjectLockConfig) || !bytes.Equal(after.VersioningConfigXML, enabledBucketVersioningConfig) {
t.Fatalf("%s: missing bootstrap configs: objectLock=%q versioning=%q", instanceType, after.ObjectLockConfigXML, after.VersioningConfigXML)
}
if !after.ObjectLockConfigUpdatedAt.Equal(before.Created) || !after.VersioningConfigUpdatedAt.Equal(before.Created) {
t.Fatalf("%s: bootstrap timestamps = (%v, %v), want %v", instanceType,
after.ObjectLockConfigUpdatedAt, after.VersioningConfigUpdatedAt, before.Created)
}
}
+266
View File
@@ -0,0 +1,266 @@
// Copyright (c) 2015-2026 MinIO, Inc.
//
// This file is part of MinIO 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"
"encoding/base64"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
"github.com/minio/mux"
)
func TestSRBucketObjectLockMetadata(t *testing.T) {
updatedAt := time.Date(2026, time.August, 29, 8, 0, 0, 0, time.UTC)
current := "current"
legacy := "legacy"
event := newSRBucketObjectLockMeta("bucket", &current, updatedAt)
if event.Type != madmin.SRBucketMetaTypeObjectLockConfig || event.Bucket != "bucket" ||
event.ObjectLockConfig == nil || *event.ObjectLockConfig != current || event.Tags != nil || !event.UpdatedAt.Equal(updatedAt) {
t.Fatalf("unexpected Object Lock event: %#v", event)
}
encoded, err := json.Marshal(event)
if err != nil {
t.Fatal(err)
}
var roundTrip madmin.SRBucketMeta
if err := json.Unmarshal(encoded, &roundTrip); err != nil {
t.Fatal(err)
}
if roundTrip.ObjectLockConfig == nil || *roundTrip.ObjectLockConfig != current || roundTrip.Tags != nil {
t.Fatalf("unexpected JSON round trip: %#v", roundTrip)
}
for _, test := range []struct {
name string
item madmin.SRBucketMeta
want *string
}{
{name: "current", item: madmin.SRBucketMeta{ObjectLockConfig: &current}, want: &current},
{name: "legacy", item: madmin.SRBucketMeta{Tags: &legacy}, want: &legacy},
{name: "current wins", item: madmin.SRBucketMeta{ObjectLockConfig: &current, Tags: &legacy}, want: &current},
{name: "missing", item: madmin.SRBucketMeta{}},
} {
t.Run(test.name, func(t *testing.T) {
got := srObjectLockPayload(test.item)
if test.want == nil {
if got != nil {
t.Fatalf("payload = %q, want nil", *got)
}
return
}
if got == nil || *got != *test.want {
t.Fatalf("payload = %v, want %q", got, *test.want)
}
})
}
}
func TestPeerBucketObjectLockMetadataCurrentAndLegacyPayloads(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testPeerBucketObjectLockMetadataCurrentAndLegacyPayloads,
makeBucketOptions: MakeBucketOptions{LockEnabled: true},
})
}
func applySRBucketMetaViaAdmin(t *testing.T, credentials auth.Credentials, item madmin.SRBucketMeta) *httptest.ResponseRecorder {
t.Helper()
body, err := json.Marshal(item)
if err != nil {
t.Fatal(err)
}
adminRouter := mux.NewRouter()
registerAdminRouter(adminRouter, true)
path := adminPathPrefix + adminAPIVersionPrefix + "/site-replication/peer/bucket-meta"
req, err := newTestSignedRequestV4(http.MethodPut, path, int64(len(body)), bytes.NewReader(body),
credentials.AccessKey, credentials.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
adminRouter.ServeHTTP(rec, req)
return rec
}
func testPeerBucketObjectLockMetadataCurrentAndLegacyPayloads(_ ObjectLayer, instanceType, bucketName string,
_ http.Handler, credentials auth.Credentials, t *testing.T,
) {
apply := func(item madmin.SRBucketMeta, wantDays uint64) {
t.Helper()
rec := applySRBucketMetaViaAdmin(t, credentials, item)
if rec.Code != http.StatusOK {
t.Fatalf("%s: admin Object Lock apply returned %d: %s", instanceType, rec.Code, rec.Body.String())
}
config, _, err := globalBucketMetadataSys.GetObjectLockConfig(bucketName)
if err != nil {
t.Fatal(err)
}
if config.Rule == nil || config.Rule.DefaultRetention.Mode != "GOVERNANCE" ||
config.Rule.DefaultRetention.Days == nil || *config.Rule.DefaultRetention.Days != wantDays {
t.Fatalf("%s: persisted Object Lock config = %s, want GOVERNANCE/%d days", instanceType, config, wantDays)
}
}
config30 := base64.StdEncoding.EncodeToString([]byte(`<ObjectLockConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><ObjectLockEnabled>Enabled</ObjectLockEnabled><Rule><DefaultRetention><Mode>GOVERNANCE</Mode><Days>30</Days></DefaultRetention></Rule></ObjectLockConfiguration>`))
apply(newSRBucketObjectLockMeta(bucketName, &config30, UTCNow().Add(time.Hour)), 30)
config45 := base64.StdEncoding.EncodeToString([]byte(`<ObjectLockConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><ObjectLockEnabled>Enabled</ObjectLockEnabled><Rule><DefaultRetention><Mode>GOVERNANCE</Mode><Days>45</Days></DefaultRetention></Rule></ObjectLockConfiguration>`))
apply(madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeObjectLockConfig,
Bucket: bucketName,
Tags: &config45,
UpdatedAt: UTCNow().Add(2 * time.Hour),
}, 45)
}
func TestPeerBucketObjectLockMetadataWithoutLockEnabled(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testPeerBucketObjectLockMetadataWithoutLockEnabled,
})
}
func testPeerBucketObjectLockMetadataWithoutLockEnabled(_ ObjectLayer, instanceType, bucketName string,
_ http.Handler, credentials auth.Credentials, t *testing.T,
) {
config := base64.StdEncoding.EncodeToString([]byte(`<ObjectLockConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><ObjectLockEnabled>Enabled</ObjectLockEnabled><Rule><DefaultRetention><Mode>GOVERNANCE</Mode><Days>30</Days></DefaultRetention></Rule></ObjectLockConfiguration>`))
item := newSRBucketObjectLockMeta(bucketName, &config, UTCNow().Add(time.Hour))
rec := applySRBucketMetaViaAdmin(t, credentials, item)
if rec.Code != http.StatusOK {
t.Fatalf("%s: admin Object Lock apply returned %d: %s", instanceType, rec.Code, rec.Body.String())
}
meta, err := globalBucketMetadataSys.Get(bucketName)
if err != nil {
t.Fatal(err)
}
if meta.objectLockConfig == nil || len(meta.VersioningConfigXML) != 0 {
t.Fatalf("%s: unlocked bucket metadata = objectLock:%v versioning:%q", instanceType, meta.objectLockConfig, meta.VersioningConfigXML)
}
}
func TestHealObjectLockMetadataUsesObjectLockField(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testHealObjectLockMetadataUsesObjectLockField,
})
}
func testHealObjectLockMetadataUsesObjectLockField(obj ObjectLayer, instanceType, bucketName string,
_ http.Handler, credentials auth.Credentials, t *testing.T,
) {
ctx := t.Context()
localID := globalDeploymentID()
remoteID := "remote-object-lock-heal"
updatedAt := UTCNow().Add(time.Hour)
createdAt := updatedAt.Add(-time.Hour)
config := base64.StdEncoding.EncodeToString([]byte(`<ObjectLockConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><ObjectLockEnabled>Enabled</ObjectLockEnabled><Rule><DefaultRetention><Mode>GOVERNANCE</Mode><Days>30</Days></DefaultRetention></Rule></ObjectLockConfiguration>`))
remoteApplies := make(chan madmin.SRBucketMeta, 1)
remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var applied madmin.SRBucketMeta
if err := json.NewDecoder(r.Body).Decode(&applied); err != nil {
t.Errorf("%s: decode remote apply: %v", instanceType, err)
w.WriteHeader(http.StatusBadRequest)
return
}
remoteApplies <- applied
w.WriteHeader(http.StatusOK)
}))
defer remote.Close()
serviceCred, err := auth.CreateCredentials("object-lock-heal-svc", "object-lock-heal-service-secret")
if err != nil {
t.Fatal(err)
}
serviceCred.ParentUser = credentials.AccessKey
if _, err = globalIAMSys.store.AddServiceAccount(ctx, serviceCred); err != nil {
t.Fatal(err)
}
defer globalIAMSys.DeleteServiceAccount(ctx, serviceCred.AccessKey, false)
globalSiteReplicationSys.Lock()
oldEnabled := globalSiteReplicationSys.enabled
oldState := globalSiteReplicationSys.state
globalSiteReplicationSys.enabled = true
globalSiteReplicationSys.state = srState{
Name: "object-lock-heal-test",
ServiceAccountAccessKey: serviceCred.AccessKey,
Peers: map[string]madmin.PeerInfo{
localID: {Name: "local", DeploymentID: localID},
remoteID: {Name: "remote", DeploymentID: remoteID, Endpoint: remote.URL},
},
}
globalSiteReplicationSys.Unlock()
defer func() {
globalSiteReplicationSys.Lock()
globalSiteReplicationSys.enabled = oldEnabled
globalSiteReplicationSys.state = oldState
globalSiteReplicationSys.Unlock()
}()
status := srStatusInfo{
Sites: map[string]madmin.PeerInfo{
localID: {Name: "local", DeploymentID: localID},
remoteID: {Name: "remote", DeploymentID: remoteID, Endpoint: remote.URL},
},
BucketStats: map[string]map[string]srBucketStatsSummary{
bucketName: {
localID: {
SRBucketStatsSummary: madmin.SRBucketStatsSummary{OLockConfigMismatch: true},
meta: srBucketMetaInfo{SRBucketInfo: madmin.SRBucketInfo{
Bucket: bucketName,
CreatedAt: createdAt,
ObjectLockConfig: &config,
ObjectLockConfigUpdatedAt: updatedAt,
}, DeploymentID: localID},
},
remoteID: {
SRBucketStatsSummary: madmin.SRBucketStatsSummary{OLockConfigMismatch: true},
meta: srBucketMetaInfo{SRBucketInfo: madmin.SRBucketInfo{
Bucket: bucketName,
CreatedAt: createdAt,
}, DeploymentID: remoteID},
},
},
},
}
if err := globalSiteReplicationSys.healOLockConfigMetadata(ctx, obj, bucketName, status); err != nil {
t.Fatal(err)
}
select {
case applied := <-remoteApplies:
if applied.Type != madmin.SRBucketMetaTypeObjectLockConfig || applied.Bucket != bucketName ||
applied.ObjectLockConfig == nil || *applied.ObjectLockConfig != config || applied.Tags != nil || !applied.UpdatedAt.Equal(updatedAt) {
t.Fatalf("%s: remote heal apply = %#v", instanceType, applied)
}
case <-time.After(5 * time.Second):
t.Fatalf("%s: remote heal did not dispatch Object Lock metadata", instanceType)
}
}
@@ -0,0 +1,191 @@
// Copyright (c) 2015-2026 MinIO, Inc.
//
// This file is part of MinIO 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 (
"encoding/base64"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"testing"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
)
func TestSiteReplicationStatusAccountsPerSiteAndSurvivesMalformedConfig(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testSiteReplicationStatusAccountsPerSiteAndSurvivesMalformedConfig,
})
}
func testSiteReplicationStatusAccountsPerSiteAndSurvivesMalformedConfig(obj ObjectLayer, instanceType, localBucket string,
_ http.Handler, credentials auth.Credentials, t *testing.T,
) {
ctx := t.Context()
remoteBucket := getRandomBucketName()
if err := obj.MakeBucket(ctx, remoteBucket, MakeBucketOptions{}); err != nil {
t.Fatal(err)
}
remoteBucketMeta, err := loadBucketMetadata(ctx, obj, remoteBucket)
if err != nil {
t.Fatal(err)
}
globalBucketMetadataSys.Set(remoteBucket, remoteBucketMeta)
globalNotificationSys.LoadBucketMetadata(ctx, remoteBucket)
tagXML := []byte(`<Tagging><TagSet><Tag><Key>key</Key><Value>value</Value></Tag></TagSet></Tagging>`)
versioningXML := []byte(`<VersioningConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Status>Enabled</Status></VersioningConfiguration>`)
objectLockXML := []byte(`<ObjectLockConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><ObjectLockEnabled>Enabled</ObjectLockEnabled><Rule><DefaultRetention><Mode>GOVERNANCE</Mode><Days>30</Days></DefaultRetention></Rule></ObjectLockConfiguration>`)
sseXML := []byte(`<ServerSideEncryptionConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ApplyServerSideEncryptionByDefault><SSEAlgorithm>AES256</SSEAlgorithm></ApplyServerSideEncryptionByDefault></Rule></ServerSideEncryptionConfiguration>`)
quotaJSON, err := json.Marshal(madmin.BucketQuota{Type: madmin.HardQuota, Quota: 1024})
if err != nil {
t.Fatal(err)
}
policyJSON := []byte(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:GetObject"],"Resource":["arn:aws:s3:::%s/*"]}]}`, localBucket))
for configFile, data := range map[string][]byte{
bucketTaggingConfig: tagXML,
bucketVersioningConfig: versioningXML,
objectLockConfig: objectLockXML,
bucketSSEConfig: sseXML,
bucketQuotaConfigFile: quotaJSON,
bucketPolicyConfig: policyJSON,
bucketCorsConfig: []byte(testSiteReplicationCORSDoc),
} {
if _, err := globalBucketMetadataSys.Update(ctx, localBucket, configFile, data); err != nil {
t.Fatalf("%s: update %s: %v", instanceType, configFile, err)
}
}
encode := func(data []byte) *string {
encoded := base64.StdEncoding.EncodeToString(data)
return &encoded
}
remotePolicy := []byte(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:GetObject"],"Resource":["arn:aws:s3:::%s/*"]}]}`, remoteBucket))
localMeta, err := globalBucketMetadataSys.Get(localBucket)
if err != nil {
t.Fatal(err)
}
remoteInfo := madmin.SRInfo{
DeploymentID: "remote-status-accounting",
Buckets: map[string]madmin.SRBucketInfo{
localBucket: {
Bucket: localBucket,
CreatedAt: localMeta.Created,
},
remoteBucket: {
Bucket: remoteBucket,
CreatedAt: remoteBucketMeta.Created,
Tags: encode(tagXML),
Versioning: encode(versioningXML),
ObjectLockConfig: encode(objectLockXML),
SSEConfig: encode(sseXML),
QuotaConfig: encode(quotaJSON),
Policy: remotePolicy,
CorsConfig: encode([]byte(testSiteReplicationCORSDoc)),
},
},
}
remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(w).Encode(remoteInfo); err != nil {
t.Errorf("%s: encode remote metadata: %v", instanceType, err)
}
}))
defer remote.Close()
serviceCred, err := auth.CreateCredentials("status-accounting-svc", "status-accounting-service-secret")
if err != nil {
t.Fatal(err)
}
serviceCred.ParentUser = credentials.AccessKey
if _, err = globalIAMSys.store.AddServiceAccount(ctx, serviceCred); err != nil {
t.Fatal(err)
}
defer globalIAMSys.DeleteServiceAccount(ctx, serviceCred.AccessKey, false)
localID := globalDeploymentID()
remoteID := remoteInfo.DeploymentID
globalSiteReplicationSys.Lock()
oldEnabled := globalSiteReplicationSys.enabled
oldState := globalSiteReplicationSys.state
globalSiteReplicationSys.enabled = true
globalSiteReplicationSys.state = srState{
Name: "status-accounting-test",
ServiceAccountAccessKey: serviceCred.AccessKey,
Peers: map[string]madmin.PeerInfo{
localID: {Name: "local", DeploymentID: localID},
remoteID: {Name: "remote", DeploymentID: remoteID, Endpoint: remote.URL},
},
}
globalSiteReplicationSys.Unlock()
defer func() {
globalSiteReplicationSys.Lock()
globalSiteReplicationSys.enabled = oldEnabled
globalSiteReplicationSys.state = oldState
globalSiteReplicationSys.Unlock()
}()
check := func(name string, wantRemoteTags, wantRemoteQuota int) {
t.Helper()
status, err := globalSiteReplicationSys.siteReplicationStatus(ctx, obj, madmin.SRStatusOptions{Buckets: true})
if err != nil {
t.Fatal(err)
}
local := status.StatsSummary[localID]
remote := status.StatsSummary[remoteID]
if local.TotalBucketsCount != 2 || remote.TotalBucketsCount != 2 {
t.Fatalf("%s: bucket totals = local:%d remote:%d", name, local.TotalBucketsCount, remote.TotalBucketsCount)
}
if local.TotalTagsCount != 1 || remote.TotalTagsCount != wantRemoteTags ||
local.TotalLockConfigCount != 1 || remote.TotalLockConfigCount != 1 ||
local.TotalSSEConfigCount != 1 || remote.TotalSSEConfigCount != 1 ||
local.TotalVersioningConfigCount != 1 || remote.TotalVersioningConfigCount != 1 ||
local.TotalBucketPoliciesCount != 1 || remote.TotalBucketPoliciesCount != 1 ||
local.TotalQuotaConfigCount != 1 || remote.TotalQuotaConfigCount != wantRemoteQuota ||
local.TotalCorsConfigCount != 1 || remote.TotalCorsConfigCount != 1 {
t.Fatalf("%s: site totals = local:%+v remote:%+v", name, local, remote)
}
if local.ReplicatedTags != 0 || remote.ReplicatedTags != 0 ||
local.ReplicatedBucketPolicies != 0 || remote.ReplicatedBucketPolicies != 0 ||
local.ReplicatedQuotaConfig != 0 || remote.ReplicatedQuotaConfig != 0 {
t.Fatalf("%s: asymmetric configs counted as replicated: local:%+v remote:%+v", name, local, remote)
}
remoteBucketStatus := status.BucketStats[remoteBucket][remoteID]
if remoteBucketStatus.HasTagsSet != (wantRemoteTags != 0) || remoteBucketStatus.HasQuotaCfgSet != (wantRemoteQuota != 0) {
t.Fatalf("%s: remote bucket presence = tags:%v quota:%v", name, remoteBucketStatus.HasTagsSet, remoteBucketStatus.HasQuotaCfgSet)
}
}
check("valid asymmetric configs", 1, 1)
invalidTags := "not-base64"
info := remoteInfo.Buckets[remoteBucket]
info.Tags = &invalidTags
remoteInfo.Buckets[remoteBucket] = info
check("malformed tags do not drop site", 0, 1)
emptyQuota := base64.StdEncoding.EncodeToString([]byte(`{}`))
info = remoteInfo.Buckets[remoteBucket]
info.QuotaConfig = &emptyQuota
remoteInfo.Buckets[remoteBucket] = info
check("empty quota is absent", 0, 0)
}
+153 -69
View File
@@ -46,6 +46,7 @@ import (
"github.com/minio/minio/internal/bucket/cors"
"github.com/minio/minio/internal/bucket/lifecycle"
sreplication "github.com/minio/minio/internal/bucket/replication"
"github.com/minio/minio/internal/bucket/versioning"
"github.com/minio/minio/internal/logger"
xldap "github.com/minio/pkg/v3/ldap"
"github.com/minio/pkg/v3/policy"
@@ -888,6 +889,32 @@ func (c *SiteReplicationSys) DeleteBucketHook(ctx context.Context, bucket string
return errors.Unwrap(cerr)
}
func enablePeerBucketVersioning(meta *BucketMetadata) error {
if len(meta.VersioningConfigXML) == 0 {
meta.VersioningConfigXML = enabledBucketVersioningConfig
if meta.VersioningConfigUpdatedAt.IsZero() {
meta.VersioningConfigUpdatedAt = meta.Created
}
return nil
}
config, err := versioning.ParseConfig(bytes.NewReader(meta.VersioningConfigXML))
if err != nil {
meta.VersioningConfigXML = enabledBucketVersioningConfig
meta.VersioningConfigUpdatedAt = UTCNow()
return nil
}
if config.Enabled() {
return nil
}
config.Status = versioning.Enabled
meta.VersioningConfigXML, err = xml.Marshal(config)
if err != nil {
return err
}
meta.VersioningConfigUpdatedAt = UTCNow()
return nil
}
// PeerBucketMakeWithVersioningHandler - creates bucket and enables versioning.
func (c *SiteReplicationSys) PeerBucketMakeWithVersioningHandler(ctx context.Context, bucket string, opts MakeBucketOptions) error {
objAPI := newObjectLayerFn()
@@ -916,9 +943,14 @@ func (c *SiteReplicationSys) PeerBucketMakeWithVersioningHandler(ctx context.Con
meta.SetCreatedAt(opts.CreatedAt)
meta.VersioningConfigXML = enabledBucketVersioningConfig
if opts.LockEnabled {
if err := enablePeerBucketVersioning(&meta); err != nil {
return wrapSRErr(err)
}
if opts.LockEnabled && len(meta.ObjectLockConfigXML) == 0 {
meta.ObjectLockConfigXML = enabledBucketObjectLockConfig
if meta.ObjectLockConfigUpdatedAt.IsZero() {
meta.ObjectLockConfigUpdatedAt = meta.Created
}
}
if err := meta.Save(context.Background(), objAPI); err != nil {
@@ -1727,6 +1759,26 @@ func (c *SiteReplicationSys) PeerBucketTaggingHandler(ctx context.Context, bucke
return nil
}
func newSRBucketObjectLockMeta(bucket string, config *string, updatedAt time.Time) madmin.SRBucketMeta {
return madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeObjectLockConfig,
Bucket: bucket,
ObjectLockConfig: config,
UpdatedAt: updatedAt,
}
}
func srObjectLockPayload(item madmin.SRBucketMeta) *string {
if item.ObjectLockConfig != nil {
return item.ObjectLockConfig
}
return item.Tags
}
func (c *SiteReplicationSys) peerBucketObjectLockConfigItem(ctx context.Context, item madmin.SRBucketMeta) error {
return c.PeerBucketObjectLockConfigHandler(ctx, item.Bucket, srObjectLockPayload(item), item.UpdatedAt)
}
// PeerBucketObjectLockConfigHandler - sets object lock on local bucket.
func (c *SiteReplicationSys) PeerBucketObjectLockConfigHandler(ctx context.Context, bucket string, objectLockData *string, updatedAt time.Time) error {
if objectLockData != nil {
@@ -2176,12 +2228,7 @@ func (c *SiteReplicationSys) syncToAllPeers(ctx context.Context, addOpts madmin.
objLockCfgData, tm := meta.ObjectLockConfigXML, meta.ObjectLockConfigUpdatedAt
if len(objLockCfgData) > 0 {
objLockStr := base64.StdEncoding.EncodeToString(objLockCfgData)
err = c.BucketMetaHook(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeObjectLockConfig,
Bucket: bucket,
Tags: &objLockStr,
UpdatedAt: tm,
})
err = c.BucketMetaHook(ctx, newSRBucketObjectLockMeta(bucket, &objLockStr, tm))
if err != nil {
return errSRBucketMetaError(err)
}
@@ -3406,80 +3453,116 @@ func (c *SiteReplicationSys) siteReplicationStatus(ctx context.Context, objAPI O
quotaCfgs := make([]*madmin.BucketQuota, numSites)
sseCfgSet := set.NewStringSet()
versionCfgSet := set.NewStringSet()
var tagCount, olockCfgCount, sseCfgCount, corsCfgCount, versionCfgCount int
validReplCfg := make([]bool, numSites)
validVersionCfg := make([]bool, numSites)
validQuotaCfg := make([]bool, numSites)
validTags := make([]bool, numSites)
validPolicies := make([]bool, numSites)
validObjectLockCfg := make([]bool, numSites)
validSSECfg := make([]bool, numSites)
validCorsCfg := make([]bool, numSites)
var tagCount, olockCfgCount, policyCount, quotaCfgCount, sseCfgCount, corsCfgCount, versionCfgCount int
for i, s := range slc {
logInvalid := func(configType string, err error) {
replLogOnceIf(ctx,
fmt.Errorf("unable to parse %s metadata for bucket %s from site %s: %w", configType, b, s.DeploymentID, err),
"site-replication-status-"+configType+"-"+b+"-"+s.DeploymentID)
}
if s.ReplicationConfig != nil {
cfgBytes, err := base64.StdEncoding.DecodeString(*s.ReplicationConfig)
if err != nil {
continue
if err == nil {
cfg, err := sreplication.ParseConfig(bytes.NewReader(cfgBytes))
if err == nil {
replCfgs[i] = cfg
validReplCfg[i] = true
} else {
logInvalid("replication", err)
}
} else {
logInvalid("replication", err)
}
cfg, err := sreplication.ParseConfig(bytes.NewReader(cfgBytes))
if err != nil {
continue
}
replCfgs[i] = cfg
}
if s.Versioning != nil {
configData, err := base64.StdEncoding.DecodeString(*s.Versioning)
if err != nil {
continue
}
versionCfgCount++
if !versionCfgSet.Contains(string(configData)) {
versionCfgSet.Add(string(configData))
if err == nil {
validVersionCfg[i] = true
versionCfgCount++
if !versionCfgSet.Contains(string(configData)) {
versionCfgSet.Add(string(configData))
}
} else {
logInvalid("versioning", err)
}
}
if s.QuotaConfig != nil {
cfgBytes, err := base64.StdEncoding.DecodeString(*s.QuotaConfig)
if err != nil {
continue
if err == nil {
cfg, err := parseBucketQuota(b, cfgBytes)
if err == nil {
if cfg != nil && *cfg != (madmin.BucketQuota{}) {
quotaCfgs[i] = cfg
validQuotaCfg[i] = true
quotaCfgCount++
}
} else {
logInvalid("quota", err)
}
} else {
logInvalid("quota", err)
}
cfg, err := parseBucketQuota(b, cfgBytes)
if err != nil {
continue
}
quotaCfgs[i] = cfg
}
if s.Tags != nil {
tagBytes, err := base64.StdEncoding.DecodeString(*s.Tags)
if err != nil {
continue
}
tagCount++
if !tagSet.Contains(string(tagBytes)) {
tagSet.Add(string(tagBytes))
if err == nil {
validTags[i] = true
tagCount++
if !tagSet.Contains(string(tagBytes)) {
tagSet.Add(string(tagBytes))
}
} else {
logInvalid("tags", err)
}
}
if len(s.Policy) > 0 {
plcy, err := policy.ParseBucketPolicyConfig(bytes.NewReader(s.Policy), b)
if err != nil {
continue
if err == nil {
policies[i] = plcy
validPolicies[i] = true
policyCount++
} else {
logInvalid("policy", err)
}
policies[i] = plcy
}
if s.ObjectLockConfig != nil {
configData, err := base64.StdEncoding.DecodeString(*s.ObjectLockConfig)
if err != nil {
continue
}
olockCfgCount++
if !olockConfigSet.Contains(string(configData)) {
olockConfigSet.Add(string(configData))
if err == nil {
validObjectLockCfg[i] = true
olockCfgCount++
if !olockConfigSet.Contains(string(configData)) {
olockConfigSet.Add(string(configData))
}
} else {
logInvalid("object-lock", err)
}
}
if s.SSEConfig != nil {
configData, err := base64.StdEncoding.DecodeString(*s.SSEConfig)
if err != nil {
continue
}
sseCfgCount++
if !sseCfgSet.Contains(string(configData)) {
sseCfgSet.Add(string(configData))
if err == nil {
validSSECfg[i] = true
sseCfgCount++
if !sseCfgSet.Contains(string(configData)) {
sseCfgSet.Add(string(configData))
}
} else {
logInvalid("sse", err)
}
}
if s.CorsConfig != nil {
if _, err := decodeCORSReplicationPayload(s.CorsConfig); err == nil {
validCorsCfg[i] = true
corsCfgCount++
} else {
logInvalid("cors", err)
}
}
ss, ok := info.StatsSummary[s.DeploymentID]
@@ -3491,24 +3574,27 @@ func (c *SiteReplicationSys) siteReplicationStatus(ctx context.Context, objAPI O
ss.ReplicatedBuckets++
}
ss.TotalBucketsCount++
if tagCount > 0 {
if validTags[i] {
ss.TotalTagsCount++
}
if olockCfgCount > 0 {
if validObjectLockCfg[i] {
ss.TotalLockConfigCount++
}
if sseCfgCount > 0 {
if validSSECfg[i] {
ss.TotalSSEConfigCount++
}
if s.CorsConfig != nil {
if validCorsCfg[i] {
ss.TotalCorsConfigCount++
}
if versionCfgCount > 0 {
if validVersionCfg[i] {
ss.TotalVersioningConfigCount++
}
if len(policies) > 0 {
if validPolicies[i] {
ss.TotalBucketPoliciesCount++
}
if validQuotaCfg[i] {
ss.TotalQuotaConfigCount++
}
info.StatsSummary[s.DeploymentID] = ss
}
tagMismatch := !isReplicated(tagCount, numSites, tagSet)
@@ -3542,13 +3628,13 @@ func (c *SiteReplicationSys) siteReplicationStatus(ctx context.Context, objAPI O
PolicyMismatch: policyMismatch,
ReplicationCfgMismatch: replCfgMismatch,
QuotaCfgMismatch: quotaCfgMismatch,
HasReplicationCfg: s.ReplicationConfig != nil,
HasTagsSet: s.Tags != nil,
HasOLockConfigSet: s.ObjectLockConfig != nil,
HasPolicySet: s.Policy != nil,
HasReplicationCfg: validReplCfg[i],
HasTagsSet: validTags[i],
HasOLockConfigSet: validObjectLockCfg[i],
HasPolicySet: validPolicies[i],
HasQuotaCfgSet: quotaCfgSet,
HasSSECfgSet: s.SSEConfig != nil,
HasCorsCfgSet: s.CorsConfig != nil,
HasSSECfgSet: validSSECfg[i],
HasCorsCfgSet: validCorsCfg[i],
}
var m srBucketMetaInfo
if len(bucketStats[s.Bucket]) > dIdx {
@@ -3574,12 +3660,15 @@ func (c *SiteReplicationSys) siteReplicationStatus(ctx context.Context, objAPI O
if !corsCfgMismatch && corsCfgCount == numSites {
sum.ReplicatedCorsConfig++
}
if !policyMismatch && len(policies) == numSites {
if !policyMismatch && policyCount == numSites {
sum.ReplicatedBucketPolicies++
}
if !tagMismatch && tagCount == numSites {
sum.ReplicatedTags++
}
if !quotaCfgMismatch && quotaCfgCount == numSites {
sum.ReplicatedQuotaConfig++
}
info.StatsSummary[s.DeploymentID] = sum
}
}
@@ -5322,12 +5411,7 @@ func (c *SiteReplicationSys) healOLockConfigMetadata(ctx context.Context, objAPI
return wrapSRErr(err)
}
peerName := info.Sites[dID].Name
err = admClient.SRPeerReplicateBucketMeta(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeObjectLockConfig,
Bucket: bucket,
Tags: latestObjLockConfig,
UpdatedAt: lastUpdate,
})
err = admClient.SRPeerReplicateBucketMeta(ctx, newSRBucketObjectLockMeta(bucket, latestObjLockConfig, lastUpdate))
if err != nil {
replLogIf(ctx, c.annotatePeerErr(peerName, replicateBucketMetadata,
fmt.Errorf("Unable to heal object lock config metadata for peer %s from peer %s : %w",
+53 -12
View File
@@ -157,7 +157,6 @@ func ChecksumStringToType(alg string) ChecksumType {
case "SHA256":
return ChecksumSHA256
case "CRC64NVME":
// AWS seems to ignore full value, and just assume it.
return ChecksumCRC64NVME
case "":
return ChecksumNone
@@ -192,7 +191,9 @@ func NewChecksumType(alg, objType string) ChecksumType {
}
return ChecksumSHA256
case "CRC64NVME":
// AWS seems to ignore full value, and just assume it.
if objType == xhttp.AmzChecksumTypeComposite {
return ChecksumInvalid
}
return ChecksumCRC64NVME
case "":
if full != 0 {
@@ -657,22 +658,55 @@ func AddChecksumHeader(w http.ResponseWriter, c map[string]string) {
}
}
func isSupportedChecksumHeader(name string) bool {
switch {
case strings.EqualFold(name, xhttp.AmzChecksumAlgo),
strings.EqualFold(name, xhttp.AmzChecksumType),
strings.EqualFold(name, xhttp.AmzChecksumMode):
return true
}
for _, checksumType := range BaseChecksumTypes {
if strings.EqualFold(name, checksumType.Key()) {
return true
}
}
return false
}
func hasUnsupportedChecksumHeader(h http.Header) bool {
for name := range h {
if strings.HasPrefix(strings.ToLower(name), "x-amz-checksum-") && !isSupportedChecksumHeader(name) {
return true
}
}
return false
}
// GetContentChecksum returns content checksum.
// Returns ErrInvalidChecksum if so.
// Returns nil, nil if no checksum.
func GetContentChecksum(h http.Header) (*Checksum, error) {
if hasUnsupportedChecksumHeader(h) {
return nil, ErrInvalidChecksum
}
if trailing := h.Values(xhttp.AmzTrailer); len(trailing) > 0 {
var res *Checksum
for _, header := range trailing {
var duplicates bool
for _, t := range BaseChecksumTypes {
if strings.EqualFold(t.Key(), header) {
duplicates = res != nil
res = NewChecksumWithType(t|ChecksumTrailing, "")
for _, headers := range trailing {
for header := range strings.SplitSeq(headers, ",") {
header = strings.TrimSpace(header)
var duplicates bool
for _, t := range BaseChecksumTypes {
if strings.EqualFold(t.Key(), header) {
duplicates = res != nil
res = NewChecksumWithType(t|ChecksumTrailing, "")
}
}
if strings.HasPrefix(strings.ToLower(header), "x-amz-checksum-") && !isSupportedChecksumHeader(header) {
return nil, ErrInvalidChecksum
}
if duplicates {
return nil, ErrInvalidChecksum
}
}
if duplicates {
return nil, ErrInvalidChecksum
}
}
if res != nil {
@@ -682,7 +716,11 @@ func GetContentChecksum(h http.Header) (*Checksum, error) {
return nil, ErrInvalidChecksum
}
res.Type |= ChecksumFullObject
case xhttp.AmzChecksumTypeComposite, "":
case xhttp.AmzChecksumTypeComposite:
if res.Type.Base().Is(ChecksumCRC64NVME) {
return nil, ErrInvalidChecksum
}
case "":
default:
return nil, ErrInvalidChecksum
}
@@ -748,5 +786,8 @@ func getContentChecksum(h http.Header) (t ChecksumType, s string) {
for _, t := range BaseChecksumTypes {
checkType(t)
}
if t.Base().Is(ChecksumCRC64NVME) && h.Get(xhttp.AmzChecksumType) == xhttp.AmzChecksumTypeComposite {
return ChecksumInvalid, ""
}
return t, s
}
+54
View File
@@ -18,14 +18,58 @@
package hash
import (
"errors"
"net/http"
"net/http/httptest"
"testing"
xhttp "github.com/minio/minio/internal/http"
)
func TestGetContentChecksumRejectsUnsupportedHeaders(t *testing.T) {
unsupported := []string{
"x-amz-checksum-md5",
"x-amz-checksum-sha512",
"x-amz-checksum-xxhash64",
"x-amz-checksum-xxhash3",
"x-amz-checksum-xxhash128",
"x-amz-checksum-future",
}
for _, header := range unsupported {
t.Run("header/"+header, func(t *testing.T) {
h := http.Header{header: {"AA=="}}
if _, err := GetContentChecksum(h); !errors.Is(err, ErrInvalidChecksum) {
t.Fatalf("GetContentChecksum(%s) error = %v, want ErrInvalidChecksum", header, err)
}
})
t.Run("trailer/"+header, func(t *testing.T) {
h := http.Header{xhttp.AmzTrailer: {header}}
if _, err := GetContentChecksum(h); !errors.Is(err, ErrInvalidChecksum) {
t.Fatalf("GetContentChecksum(trailer %s) error = %v, want ErrInvalidChecksum", header, err)
}
})
}
for header, value := range map[string]string{
xhttp.AmzChecksumAlgo: "CRC32",
xhttp.AmzChecksumType: xhttp.AmzChecksumTypeComposite,
xhttp.AmzChecksumMode: "ENABLED",
"x-amz-sdk-checksum-algorithm": "SHA512",
} {
t.Run("control/"+header, func(t *testing.T) {
h := http.Header{header: {value}}
if _, err := GetContentChecksum(h); errors.Is(err, ErrInvalidChecksum) {
t.Fatalf("control header %s was rejected", header)
}
})
}
}
// TestChecksumAddToHeader tests that adding and retrieving a checksum on a header works
func TestChecksumAddToHeader(t *testing.T) {
if got := NewChecksumType("CRC64NVME", xhttp.AmzChecksumTypeComposite); !got.Is(ChecksumInvalid) {
t.Fatalf("CRC64NVME/COMPOSITE = %s, want invalid", got.StringFull())
}
tests := []struct {
name string
checksum ChecksumType
@@ -106,6 +150,16 @@ func TestChecksumAddToHeader(t *testing.T) {
}
}
func TestCRC64NVMECompositeTrailerIsInvalid(t *testing.T) {
h := http.Header{}
h.Set(xhttp.AmzTrailer, ChecksumCRC64NVME.Key())
h.Set(xhttp.AmzChecksumType, xhttp.AmzChecksumTypeComposite)
_, err := GetContentChecksum(h)
if !errors.Is(err, ErrInvalidChecksum) {
t.Fatalf("CRC64NVME/COMPOSITE trailer error = %v, want ErrInvalidChecksum", err)
}
}
// TestChecksumSerializeDeserialize checks AppendTo can be reversed by ChecksumFromBytes
func TestChecksumSerializeDeserialize(t *testing.T) {
myData := []byte("this-is-a-checksum-data-test")