Compare commits

...

7 Commits

Author SHA1 Message Date
Feng Ruohang 022722a7a7 chore: align R4 contribution notices and record final review
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 00:08:09 +08:00
Feng Ruohang 03027727d1 fix(replication): preserve SSE-KMS tag timestamps
Carry the already-parsed source tagging timestamp through the KMS options
constructor so replica COPY can apply newer tag updates on explicitly or
automatically encrypted destinations.

Cover all option encryption modes and signed COPY persistence for newer,
stale, duplicate and timestamp-less updates, including bucket defaults.
Preserve the real Opus 5.0/max plan review, consensus and local validation.

Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 00:03:52 +08:00
Feng Ruohang 9ebe81c1b3 Merge pull request #192 from pgsty/codex/iam-revision-tombstones
fix(iam): retain revocation versions through replay and recovery
2026-09-15 23:14:22 +08:00
Feng Ruohang e5f5c9e7f6 Merge pull request #191 from pgsty/codex/iam-peer-delete-reload
fix(iam): reload committed state on peer deletion notifications
2026-09-15 23:10:38 +08:00
Feng Ruohang 7b4cacc392 chore: align IAM file notices with contribution policy
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-15 22:59:04 +08:00
Feng Ruohang a0dd7dae9b fix(iam): resume site healing after leadership changes
Keep one healing loop per process across replication configuration reloads. Reacquire leadership after a lease is canceled and allow shutdown while waiting, so temporary quorum loss cannot permanently stop revocation propagation.

Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-15 22:59:04 +08:00
Feng Ruohang 709d50a916 fix(iam): persist revocations across site replay and recovery
Retain source-ordered tombstones and parent grant boundaries across both IAM backends, cache reloads, and deliberate identity recreation. Reconcile deletions through a versioned, bounded replication protocol with restart-aware acknowledgements.

Cover inherited group grants, STS retention, same-key service recreation, absolute expiration, and failures after the durable commit. Document coordinated upgrades and the remaining consistency boundaries.

Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-15 22:59:04 +08:00
42 changed files with 5881 additions and 413 deletions
@@ -829,6 +829,7 @@
"/site-replication/peer/bucket-ops",
"/site-replication/peer/edit",
"/site-replication/peer/iam-item",
"/site-replication/peer/iam-revisions",
"/site-replication/peer/idp-settings",
"/site-replication/peer/join",
"/site-replication/peer/remove",
@@ -877,6 +878,7 @@
"/v2/metrics/cluster",
"/v2/metrics/node",
"/v2/metrics/resource",
"/v3/site-replication/peer/iam-revisions",
"/var/vcap/bosh",
"/verifybinary",
"/version",
+1
View File
@@ -388,6 +388,7 @@ func registerAdminRouter(router *mux.Router, enableConfigOps bool) {
adminRouter.Methods(http.MethodPut).Path(adminVersion + "/site-replication/peer/join").HandlerFunc(adminMiddleware(adminAPI.SRPeerJoin))
adminRouter.Methods(http.MethodPut).Path(adminVersion+"/site-replication/peer/bucket-ops").HandlerFunc(adminMiddleware(adminAPI.SRPeerBucketOps)).Queries("bucket", "{bucket:.*}").Queries("operation", "{operation:.*}")
adminRouter.Methods(http.MethodPut).Path(adminVersion + "/site-replication/peer/iam-item").HandlerFunc(adminMiddleware(adminAPI.SRPeerReplicateIAMItem))
adminRouter.Methods(http.MethodGet, http.MethodPut).Path(adminVersion + "/site-replication/peer/iam-revisions").HandlerFunc(adminMiddleware(adminAPI.SRPeerIAMRevisions))
adminRouter.Methods(http.MethodPut).Path(adminVersion + "/site-replication/peer/bucket-meta").HandlerFunc(adminMiddleware(adminAPI.SRPeerReplicateBucketItem))
adminRouter.Methods(http.MethodGet).Path(adminVersion + "/site-replication/peer/idp-settings").HandlerFunc(adminMiddleware(adminAPI.SRPeerGetIDPSettings))
adminRouter.Methods(http.MethodPut).Path(adminVersion + "/site-replication/edit").HandlerFunc(adminMiddleware(adminAPI.SiteReplicationEdit))
+111
View File
@@ -0,0 +1,111 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"errors"
"testing"
"time"
"github.com/minio/minio/internal/auth"
)
func TestIAMCredentialRetention(t *testing.T) {
for _, backend := range []string{"object", "etcd"} {
t.Run(backend, func(t *testing.T) {
ctx, sys, _ := prepareIAMRevisionFixture(t, backend)
secret, err := getTokenSigningKey()
mustIAM(t, err)
parent := "external-idp-parent"
credential := func(exp time.Time) auth.Credentials {
cred, err := auth.GetNewCredentialsWithMetadata(map[string]any{"exp": exp.Unix(), parentClaim: parent}, secret)
mustIAM(t, err)
cred.ParentUser = parent
return cred
}
// Disablement of an external identity must include cached STS,
// which are kept separately from regular and service accounts.
cred := credential(UTCNow().Add(time.Hour))
_, err = sys.SetTempUser(ctx, cred.AccessKey, cred, "")
mustIAM(t, err)
mustIAM(t, sys.store.DeleteUsers(ctx, []string{parent}))
r, err := loadIAMRevision(ctx, sys.store, getUserIdentityPath(cred.AccessKey, stsUser))
mustIAM(t, err)
if !r.Deleted || !r.ExpiresAt.Equal(cred.Expiration.Add(globalMaxSkewTime)) || r.Credentials.SessionToken != "" || r.Credentials.SecretKey != "" {
t.Fatal("early STS revocation lost its retention boundary or retained a secret")
}
if _, ok := sys.store.GetUser(cred.AccessKey); ok {
t.Fatal("external disablement left the STS cache live")
}
_, err = sys.SetTempUser(withIAMReplicationTime(ctx, UTCNow().Add(time.Minute)), cred.AccessKey, cred, "")
if !errors.Is(err, errIAMStaleUpdate) {
t.Fatalf("same revoked token was reissued by replay: %v", err)
}
var mp MappedPolicy
err = sys.store.loadIAMConfig(ctx, &mp, getMappedPolicyPath(cred.AccessKey, stsUser, false))
if !errors.Is(err, errConfigNotFound) {
t.Fatalf("random STS key produced a permanent mapping: %v", err)
}
// Seed genuinely expired immutable tokens, as an ordinary startup
// loader sees them. Natural expiry leaves no permanent tombstone.
expired := credential(UTCNow().Add(-time.Hour))
path := getUserIdentityPath(expired.AccessKey, stsUser)
mustIAM(t, sys.store.saveIAMConfig(ctx, &UserIdentity{Version: 1, Credentials: expired, UpdatedAt: UTCNow().Add(-2 * time.Hour)}, path))
_ = sys.store.loadUser(ctx, expired.AccessKey, stsUser, make(map[string]UserIdentity))
var u UserIdentity
if err := sys.store.loadIAMConfig(ctx, &u, path); !errors.Is(err, errConfigNotFound) {
t.Fatalf("natural expiration retained a random key: %v", err)
}
// A retained early-revocation record is collectable only after the
// immutable token's expiration plus the skew allowance.
tomb := UserIdentity{Version: 1, Deleted: true, UpdatedAt: UTCNow().Add(-2 * time.Hour), ExpiresAt: expired.Expiration.Add(globalMaxSkewTime)}
mustIAM(t, sys.store.saveIAMConfig(ctx, &tomb, path))
_ = sys.store.loadUser(ctx, expired.AccessKey, stsUser, make(map[string]UserIdentity))
if err := sys.store.loadIAMConfig(ctx, &u, path); !errors.Is(err, errConfigNotFound) {
t.Fatalf("expired STS revocation not collected: %v", err)
}
if _, ok := sys.store.revisionIndex().snapshot()[path]; ok {
t.Fatal("expired STS retained an index entry")
}
})
}
}
func TestIAMPolicyDeletionRemainsExplicit(t *testing.T) {
for _, backend := range []string{"object", "etcd"} {
t.Run(backend, func(t *testing.T) {
ctx, sys, _ := prepareIAMRevisionFixture(t, backend)
mustIAM(t, sys.DeletePolicy(ctx, "misspelled-policy", true))
r, err := loadIAMRevision(ctx, sys.store, getPolicyDocPath("misspelled-policy"))
mustIAM(t, err)
if r.Deleted {
t.Fatal("local nonexistent policy created a tombstone")
}
p, err := sys.store.GetPolicy("readwrite")
mustIAM(t, err)
if err := sys.DeletePolicy(ctx, "readwrite", true); err == nil {
t.Fatal("local pristine builtin policy became deletable")
}
_, err = sys.SetPolicy(ctx, "readwrite", p)
mustIAM(t, err)
mustIAM(t, sys.DeletePolicy(ctx, "readwrite", true))
mustIAM(t, sys.store.LoadIAMCache(ctx, false))
if _, err := sys.store.GetPolicy("readwrite"); !errors.Is(err, errNoSuchPolicy) {
t.Fatalf("reload restored an explicitly deleted override: %v", err)
}
_, err = sys.SetPolicy(ctx, "readwrite", p)
mustIAM(t, err)
if _, err := sys.store.GetPolicy("readwrite"); err != nil {
t.Fatal("explicit policy recreation failed", err)
}
mustIAM(t, globalSiteReplicationSys.PeerAddPolicyHandler(ctx, "remote-unknown-policy", nil, UTCNow()))
r, err = loadIAMRevision(ctx, sys.store, getPolicyDocPath("remote-unknown-policy"))
mustIAM(t, err)
if !r.Deleted {
t.Fatal("replicated unknown deletion lost its version")
}
})
}
}
+65 -71
View File
@@ -26,7 +26,6 @@ import (
"sync"
"time"
jsoniter "github.com/json-iterator/go"
"github.com/minio/minio-go/v7/pkg/set"
"github.com/minio/minio/internal/config"
"github.com/minio/minio/internal/kms"
@@ -62,6 +61,7 @@ type IAMEtcdStore struct {
sync.RWMutex
*iamCache
index iamRevisionIndex
usersSysType UsersSysType
@@ -69,13 +69,17 @@ type IAMEtcdStore struct {
}
func newIAMEtcdStore(client *etcd.Client, usersSysType UsersSysType) *IAMEtcdStore {
return &IAMEtcdStore{
store := &IAMEtcdStore{
iamCache: newIamCache(),
client: client,
usersSysType: usersSysType,
}
store.revisions = &store.index
return store
}
func (ies *IAMEtcdStore) revisionIndex() *iamRevisionIndex { return &ies.index }
func (ies *IAMEtcdStore) rlock() *iamCache {
ies.RLock()
return ies.iamCache
@@ -103,6 +107,7 @@ func (ies *IAMEtcdStore) saveIAMConfig(ctx context.Context, item any, itemPath s
if err != nil {
return err
}
plain := data
if GlobalKMS != nil {
data, err = config.EncryptBytes(GlobalKMS, data, kms.Context{
minioMetaBucket: path.Join(minioMetaBucket, itemPath),
@@ -111,24 +116,28 @@ func (ies *IAMEtcdStore) saveIAMConfig(ctx context.Context, item any, itemPath s
return err
}
}
return saveKeyEtcd(ctx, ies.client, itemPath, data, opts...)
if err := saveKeyEtcd(ctx, ies.client, itemPath, data, opts...); err != nil {
return err
}
ies.index.observe(itemPath, plain)
return nil
}
func getIAMConfig(item any, data []byte, itemPath string) error {
data, err := decryptData(data, itemPath)
func (ies *IAMEtcdStore) decodeIAMConfig(item any, data []byte, path string) error {
data, err := decryptData(data, path)
if err != nil {
return err
}
json := jsoniter.ConfigCompatibleWithStandardLibrary
ies.index.observe(path, data)
return json.Unmarshal(data, item)
}
func (ies *IAMEtcdStore) loadIAMConfig(ctx context.Context, item any, path string) error {
data, err := readKeyEtcd(ctx, ies.client, path)
data, err := ies.loadIAMConfigBytes(ctx, path)
if err != nil {
return err
}
return getIAMConfig(item, data, path)
return json.Unmarshal(data, item)
}
func (ies *IAMEtcdStore) loadIAMConfigBytes(ctx context.Context, path string) ([]byte, error) {
@@ -136,11 +145,19 @@ func (ies *IAMEtcdStore) loadIAMConfigBytes(ctx context.Context, path string) ([
if err != nil {
return nil, err
}
return decryptData(data, path)
data, err = decryptData(data, path)
if err == nil {
ies.index.observe(path, data)
}
return data, err
}
func (ies *IAMEtcdStore) deleteIAMConfig(ctx context.Context, path string) error {
return deleteKeyEtcd(ctx, ies.client, path)
if err := deleteKeyEtcd(ctx, ies.client, path); err != nil {
return err
}
ies.index.forget(path)
return nil
}
func (ies *IAMEtcdStore) loadPolicyDocWithRetry(ctx context.Context, policy string, m map[string]PolicyDoc, _ int) error {
@@ -162,6 +179,9 @@ func (ies *IAMEtcdStore) loadPolicyDoc(ctx context.Context, policy string, m map
return err
}
if p.Deleted {
return errNoSuchPolicy
}
m[policy] = p
return nil
}
@@ -181,7 +201,11 @@ func (ies *IAMEtcdStore) getPolicyDocKV(ctx context.Context, kvs *mvccpb.KeyValu
return err
}
ies.index.observe(string(kvs.Key), data)
policy := extractPathPrefixAndSuffix(string(kvs.Key), iamConfigPoliciesPrefix, path.Base(string(kvs.Key)))
if p.Deleted {
return errNoSuchPolicy
}
m[policy] = p
return nil
}
@@ -207,7 +231,7 @@ func (ies *IAMEtcdStore) loadPolicyDocs(ctx context.Context, m map[string]Policy
func (ies *IAMEtcdStore) getUserKV(ctx context.Context, userkv *mvccpb.KeyValue, userType IAMUserType, m map[string]UserIdentity, basePrefix string) error {
var u UserIdentity
err := getIAMConfig(&u, userkv.Value, string(userkv.Key))
err := ies.decodeIAMConfig(&u, userkv.Value, string(userkv.Key))
if err != nil {
if err == errConfigNotFound {
return errNoSuchUser
@@ -219,10 +243,14 @@ func (ies *IAMEtcdStore) getUserKV(ctx context.Context, userkv *mvccpb.KeyValue,
}
func (ies *IAMEtcdStore) addUser(ctx context.Context, user string, userType IAMUserType, u UserIdentity, m map[string]UserIdentity) error {
if u.Deleted {
if userType == stsUser && !u.ExpiresAt.IsZero() && UTCNow().After(u.ExpiresAt) {
bestEffortIAMExpiration(ctx, ies, getUserIdentityPath(user, userType))
}
return errNoSuchUser
}
if u.Credentials.IsExpired() {
// Delete expired identity.
deleteKeyEtcd(ctx, ies.client, getUserIdentityPath(user, userType))
deleteKeyEtcd(ctx, ies.client, getMappedPolicyPath(user, userType, false))
bestEffortIAMExpiration(ctx, ies, getUserIdentityPath(user, userType))
return nil
}
if u.Credentials.AccessKey == "" {
@@ -231,16 +259,17 @@ func (ies *IAMEtcdStore) addUser(ctx context.Context, user string, userType IAMU
if u.Credentials.SessionToken != "" {
jwtClaims, err := extractJWTClaims(u)
if err != nil {
if u.Credentials.IsTemp() {
// We should delete such that the client can re-request
// for the expiring credentials.
deleteKeyEtcd(ctx, ies.client, getUserIdentityPath(user, userType))
deleteKeyEtcd(ctx, ies.client, getMappedPolicyPath(user, userType, false))
}
// A temporarily unavailable signing key is not proof of expiration.
return nil
}
u.Credentials.Claims = jwtClaims.Map()
}
if err := checkIAMParentRevision(ctx, ies, u.Credentials); err != nil {
if errors.Is(err, errIAMStaleUpdate) {
return errNoSuchUser
}
return err
}
if u.Credentials.Description == "" {
u.Credentials.Description = u.Credentials.Comment
}
@@ -258,6 +287,9 @@ func (ies *IAMEtcdStore) loadSecretKey(ctx context.Context, user string, userTyp
}
return "", err
}
if u.Deleted {
return "", errNoSuchUser
}
return u.Credentials.SecretKey, nil
}
@@ -274,6 +306,7 @@ func (ies *IAMEtcdStore) loadUser(ctx context.Context, user string, userType IAM
}
func (ies *IAMEtcdStore) loadUsers(ctx context.Context, userType IAMUserType, m map[string]UserIdentity) error {
ctx = withIAMExpirationCleanup(ctx)
var basePrefix string
switch userType {
case svcUser:
@@ -312,6 +345,9 @@ func (ies *IAMEtcdStore) loadGroup(ctx context.Context, group string, m map[stri
}
return err
}
if gi.Deleted {
return errNoSuchGroup
}
m[group] = gi
return nil
}
@@ -349,13 +385,16 @@ func (ies *IAMEtcdStore) loadMappedPolicy(ctx context.Context, name string, user
}
return err
}
if !ies.index.mappingAllowed(getMappedPolicyPath(name, userType, isGroup), p) {
return errNoSuchPolicy
}
m.Store(name, p)
return nil
}
func getMappedPolicy(kv *mvccpb.KeyValue, m *xsync.MapOf[string, MappedPolicy], basePrefix string) error {
func (ies *IAMEtcdStore) getMappedPolicy(kv *mvccpb.KeyValue, m *xsync.MapOf[string, MappedPolicy], basePrefix string) error {
var p MappedPolicy
err := getIAMConfig(&p, kv.Value, string(kv.Key))
err := ies.decodeIAMConfig(&p, kv.Value, string(kv.Key))
if err != nil {
if err == errConfigNotFound {
return errNoSuchPolicy
@@ -363,6 +402,9 @@ func getMappedPolicy(kv *mvccpb.KeyValue, m *xsync.MapOf[string, MappedPolicy],
return err
}
name := extractPathPrefixAndSuffix(string(kv.Key), basePrefix, ".json")
if !ies.index.mappingAllowed(string(kv.Key), p) {
return errNoSuchPolicy
}
m.Store(name, p)
return nil
}
@@ -392,61 +434,13 @@ func (ies *IAMEtcdStore) loadMappedPolicies(ctx context.Context, userType IAMUse
// Parse all policies mapping to create the proper data model
for _, kv := range r.Kvs {
if err = getMappedPolicy(kv, m, basePrefix); err != nil && !errors.Is(err, errNoSuchPolicy) {
if err = ies.getMappedPolicy(kv, m, basePrefix); err != nil && !errors.Is(err, errNoSuchPolicy) {
return err
}
}
return nil
}
func (ies *IAMEtcdStore) savePolicyDoc(ctx context.Context, policyName string, p PolicyDoc) error {
return ies.saveIAMConfig(ctx, &p, getPolicyDocPath(policyName))
}
func (ies *IAMEtcdStore) saveMappedPolicy(ctx context.Context, name string, userType IAMUserType, isGroup bool, mp MappedPolicy, opts ...options) error {
return ies.saveIAMConfig(ctx, mp, getMappedPolicyPath(name, userType, isGroup), opts...)
}
func (ies *IAMEtcdStore) saveUserIdentity(ctx context.Context, name string, userType IAMUserType, u UserIdentity, opts ...options) error {
return ies.saveIAMConfig(ctx, u, getUserIdentityPath(name, userType), opts...)
}
func (ies *IAMEtcdStore) saveGroupInfo(ctx context.Context, name string, gi GroupInfo) error {
return ies.saveIAMConfig(ctx, gi, getGroupInfoPath(name))
}
func (ies *IAMEtcdStore) deletePolicyDoc(ctx context.Context, name string) error {
err := ies.deleteIAMConfig(ctx, getPolicyDocPath(name))
if err == errConfigNotFound {
err = errNoSuchPolicy
}
return err
}
func (ies *IAMEtcdStore) deleteMappedPolicy(ctx context.Context, name string, userType IAMUserType, isGroup bool) error {
err := ies.deleteIAMConfig(ctx, getMappedPolicyPath(name, userType, isGroup))
if err == errConfigNotFound {
err = errNoSuchPolicy
}
return err
}
func (ies *IAMEtcdStore) deleteUserIdentity(ctx context.Context, name string, userType IAMUserType) error {
err := ies.deleteIAMConfig(ctx, getUserIdentityPath(name, userType))
if err == errConfigNotFound {
err = errNoSuchUser
}
return err
}
func (ies *IAMEtcdStore) deleteGroupInfo(ctx context.Context, name string) error {
err := ies.deleteIAMConfig(ctx, getGroupInfoPath(name))
if err == errConfigNotFound {
err = errNoSuchGroup
}
return err
}
func (ies *IAMEtcdStore) watch(ctx context.Context, keyPath string) <-chan iamWatchEvent {
ch := make(chan iamWatchEvent)
+164
View File
@@ -0,0 +1,164 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"maps"
"slices"
"time"
"github.com/minio/minio-go/v7/pkg/set"
)
type (
iamGroupGrantsKey struct{}
iamGroupMutationKey struct{}
iamGroupMutation struct {
Members []string
Remove bool
StatusOnly bool
}
)
// Merge the intended mutation with the record read under the distributed
// revision lock, not the older cache used to prepare the request.
func mergeIAMGroupMutation(ctx context.Context, previous GroupInfo, next *GroupInfo) {
op, ok := ctx.Value(iamGroupMutationKey{}).(iamGroupMutation)
if !ok || previous.Deleted || previous.Version == 0 {
return
}
members := set.CreateStringSet(previous.Members...)
grants := maps.Clone(previous.MemberGrants)
if grants == nil {
grants = make(map[string]time.Time)
}
switch {
case op.StatusOnly:
// Only the status changes.
case op.Remove:
for _, member := range op.Members {
members.Remove(member)
delete(grants, member)
}
next.Status = previous.Status
default:
requested := set.CreateStringSet(next.Members...)
for _, member := range op.Members {
if !requested.Contains(member) {
continue
}
at := next.MemberGrants[member]
if at.Before(grants[member]) {
continue
}
members.Add(member)
grants[member] = at
}
next.Status = previous.Status
}
next.Members, next.MemberGrants = members.ToSlice(), grants
slices.Sort(next.Members)
}
// A non-nil map is supplied by the versioned peer envelope, including for
// snapshots. Missing times are unknown, never the snapshot's newer timestamp.
func withIAMGroupGrants(ctx context.Context, grants map[string]time.Time) context.Context {
return context.WithValue(ctx, iamGroupGrantsKey{}, grants)
}
func (c *iamCache) effectiveGroupMembers(gi GroupInfo) []string {
var members []string
for _, member := range gi.Members {
if c.groupMemberAllowed(member, gi.MemberGrants[member], gi.RevokedBefore) {
members = append(members, member)
}
}
return members
}
func (c *iamCache) effectiveUserGroups(user string) []string {
var groups []string
for group := range c.iamUserGroupMemberships[user] {
gi, ok := c.iamGroupsMap[group]
r := c.revisions.get(getGroupInfoPath(group))
if r.RevokedBefore.After(gi.RevokedBefore) {
gi.RevokedBefore = r.RevokedBefore
}
if ok && !r.Deleted && c.groupMemberAllowed(user, gi.MemberGrants[user], gi.RevokedBefore) {
groups = append(groups, group)
}
}
return groups
}
func (c *iamCache) addGroupMembers(ctx context.Context, gi GroupInfo, members []string) (GroupInfo, error) {
grants, versioned := ctx.Value(iamGroupGrantsKey{}).(map[string]time.Time)
if boundary, ok := ctx.Value(iamRecordBoundaryKey{}).(time.Time); ok && boundary.After(gi.RevokedBefore) {
gi.RevokedBefore = boundary
}
origin, replicated := iamReplicationTime(ctx)
gi.Members = slices.Clone(gi.Members)
gi.MemberGrants = maps.Clone(gi.MemberGrants)
if gi.MemberGrants == nil {
gi.MemberGrants = make(map[string]time.Time)
}
current := set.CreateStringSet(gi.Members...)
gi.UpdatedAt = UTCNow()
if replicated {
gi.UpdatedAt = origin
}
for _, member := range members {
at := gi.UpdatedAt
r := c.userRevocation(member)
if replicated {
switch {
case versioned:
at = grants[member]
if at.After(origin) {
return gi, errInvalidArgument
}
case !r.RevokedBefore.IsZero() || r.Deleted || !gi.RevokedBefore.IsZero():
// Legacy snapshots cannot prove a post-revocation grant.
continue
case current.Contains(member):
continue
}
if !c.groupMemberAllowed(member, at, gi.RevokedBefore) {
continue
}
} else {
if current.Contains(member) && c.groupMemberAllowed(member, gi.MemberGrants[member], gi.RevokedBefore) {
continue // Editing the group is not reissuing every grant.
}
if !at.After(gi.RevokedBefore) {
at = gi.RevokedBefore.Add(time.Nanosecond)
}
if !at.After(r.RevokedBefore) {
at = r.RevokedBefore.Add(time.Nanosecond)
}
if !at.After(gi.MemberGrants[member]) {
at = gi.MemberGrants[member].Add(time.Nanosecond)
}
if at.After(gi.UpdatedAt) {
gi.UpdatedAt = at
}
}
u, ok := c.iamUsersMap[member]
if !ok {
return gi, errNoSuchUser
}
if u.Credentials.IsTemp() || u.Credentials.IsServiceAccount() {
return gi, errIAMActionNotAllowed
}
if previous := gi.MemberGrants[member]; previous.After(at) {
continue
}
current.Add(member)
gi.MemberGrants[member] = at
}
gi.Members = current.ToSlice()
slices.Sort(gi.Members)
return gi, nil
}
+75
View File
@@ -0,0 +1,75 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"testing"
"time"
)
func TestIAMHealingResumesAfterLeadershipLoss(t *testing.T) {
previous := globalLeaderLock
locks := make(chan LockContext)
globalLeaderLock = &sharedLock{lockContext: locks}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
c := &SiteReplicationSys{}
go func() { c.startHealRoutine(ctx, nil); close(done) }()
t.Cleanup(func() {
cancel()
select {
case <-done:
case locks <- LockContext{ctx: ctx}:
}
<-done
globalLeaderLock = previous
})
first, loseFirst := context.WithCancel(ctx)
defer loseFirst()
select {
case locks <- LockContext{ctx: first}:
case <-time.After(time.Second):
t.Fatal("healer did not acquire its first leader context")
}
loseFirst() // A transient quorum loss cancels the distributed lease.
select {
case <-done:
t.Fatal("healer permanently exited after temporary leadership loss")
case locks <- LockContext{ctx: ctx}:
case <-time.After(time.Second):
t.Fatal("healer did not wait for reacquired leadership")
}
// Shutdown must also interrupt the wait for leadership after lease loss.
cancel()
select {
case <-done:
case <-time.After(time.Second):
t.Fatal("healer did not stop with its owning context")
}
}
func TestIAMHealingLeadershipWaitCancels(t *testing.T) {
previous := globalLeaderLock
locks := make(chan LockContext)
globalLeaderLock = &sharedLock{lockContext: locks}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
c := &SiteReplicationSys{}
go func() { c.startHealRoutine(ctx, nil); close(done) }()
cancel()
t.Cleanup(func() {
select {
case <-done:
case locks <- LockContext{ctx: ctx}:
}
<-done
globalLeaderLock = previous
})
select {
case <-done:
case <-time.After(time.Second):
t.Fatal("healer ignored shutdown while waiting for leadership")
}
}
+60
View File
@@ -0,0 +1,60 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"time"
"github.com/minio/madmin-go/v3"
)
// Measures steady-state index traversal, sorting and the capability request.
// The network peer acknowledges real batches but performs no disk I/O; this
// benchmark deliberately does not claim durable catch-up throughput.
func BenchmarkIAMRevisionConvergedHealing(b *testing.B) {
for _, n := range []int{1000, 10000} {
b.Run(fmt.Sprint(n), func(b *testing.B) {
ctx, sys, _ := prepareIAMRevisionFixture(b)
_, err := sys.CreateUser(ctx, "benchmark-sync", madmin.AddOrUpdateUserReq{SecretKey: "valid-sync-password", Status: madmin.AccountEnabled})
mustIAM(b, err)
for i := range n {
at := UTCNow().Add(time.Duration(i) * time.Nanosecond)
data, err := json.Marshal(iamRevision{Deleted: true, UpdatedAt: at, RevokedBefore: at})
mustIAM(b, err)
sys.store.revisionIndex().observe(getUserIdentityPath(fmt.Sprintf("deleted-%06d", i), regUser), data)
}
var puts atomic.Int64
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/minio/health/live" {
w.WriteHeader(http.StatusOK)
return
}
if r.Method == http.MethodPut {
puts.Add(1)
}
_ = json.NewEncoder(w).Encode(iamRevisionResponse{iamRevisionStatus: iamRevisionStatus{Version: iamRevisionProtocol, Node: "node-1", Instance: "benchmark-peer", Digest: "constant"}})
}))
defer server.Close()
c := &SiteReplicationSys{enabled: true, state: srState{ServiceAccountAccessKey: "benchmark-sync", Peers: map[string]madmin.PeerInfo{globalDeploymentID(): {DeploymentID: globalDeploymentID(), Name: "local"}, "remote": {DeploymentID: "remote", Name: "remote", Endpoint: server.URL}}}}
mustIAM(b, c.healIAMDeletions(ctx))
before := puts.Load()
b.ReportAllocs()
b.ResetTimer()
for b.Loop() {
mustIAM(b, c.healIAMDeletions(ctx))
}
b.StopTimer()
b.ReportMetric(float64(puts.Load()-before)/float64(b.N), "PUT/op")
if puts.Load() != before {
b.Fatal("steady-state healing replayed acknowledged records")
}
})
}
}
+56 -61
View File
@@ -45,6 +45,7 @@ type IAMObjectStore struct {
sync.RWMutex
*iamCache
index iamRevisionIndex
usersSysType UsersSysType
@@ -52,13 +53,17 @@ type IAMObjectStore struct {
}
func newIAMObjectStore(objAPI ObjectLayer, usersSysType UsersSysType) *IAMObjectStore {
return &IAMObjectStore{
store := &IAMObjectStore{
iamCache: newIamCache(),
objAPI: objAPI,
usersSysType: usersSysType,
}
store.revisions = &store.index
return store
}
func (iamOS *IAMObjectStore) revisionIndex() *iamRevisionIndex { return &iamOS.index }
func (iamOS *IAMObjectStore) rlock() *iamCache {
iamOS.RLock()
return iamOS.iamCache
@@ -87,6 +92,7 @@ func (iamOS *IAMObjectStore) saveIAMConfig(ctx context.Context, item any, objPat
if err != nil {
return err
}
plain := data
if GlobalKMS != nil {
data, err = config.EncryptBytes(GlobalKMS, data, kms.Context{
minioMetaBucket: path.Join(minioMetaBucket, objPath),
@@ -95,7 +101,11 @@ func (iamOS *IAMObjectStore) saveIAMConfig(ctx context.Context, item any, objPat
return err
}
}
return saveConfig(ctx, iamOS.objAPI, objPath, data)
if err := saveConfig(ctx, iamOS.objAPI, objPath, data); err != nil {
return err
}
iamOS.index.observe(objPath, plain)
return nil
}
func decryptData(data []byte, objPath string) ([]byte, error) {
@@ -133,6 +143,7 @@ func (iamOS *IAMObjectStore) loadIAMConfigBytesWithMetadata(ctx context.Context,
if err != nil {
return nil, meta, err
}
iamOS.index.observe(objPath, data)
return data, meta, nil
}
@@ -146,7 +157,11 @@ func (iamOS *IAMObjectStore) loadIAMConfig(ctx context.Context, item any, objPat
}
func (iamOS *IAMObjectStore) deleteIAMConfig(ctx context.Context, path string) error {
return deleteConfig(ctx, iamOS.objAPI, path)
if err := deleteConfig(ctx, iamOS.objAPI, path); err != nil {
return err
}
iamOS.index.forget(path)
return nil
}
func (iamOS *IAMObjectStore) loadPolicyDocWithRetry(ctx context.Context, policy string, m map[string]PolicyDoc, retries int) error {
@@ -171,6 +186,10 @@ func (iamOS *IAMObjectStore) loadPolicyDocWithRetry(ctx context.Context, policy
return err
}
if p.Deleted {
return errNoSuchPolicy
}
if p.Version == 0 {
// This means that policy was in the old version (without any
// timestamp info). We fetch the mod time of the file and save
@@ -200,6 +219,10 @@ func (iamOS *IAMObjectStore) loadPolicy(ctx context.Context, policy string) (Pol
return p, err
}
if p.Deleted {
return PolicyDoc{}, errNoSuchPolicy
}
if p.Version == 0 {
// This means that policy was in the old version (without any
// timestamp info). We fetch the mod time of the file and save
@@ -245,6 +268,9 @@ func (iamOS *IAMObjectStore) loadSecretKey(ctx context.Context, user string, use
}
return "", err
}
if u.Deleted {
return "", errNoSuchUser
}
return u.Credentials.SecretKey, nil
}
@@ -258,10 +284,15 @@ func (iamOS *IAMObjectStore) loadUserIdentity(ctx context.Context, user string,
return u, err
}
if u.Deleted {
if userType == stsUser && !u.ExpiresAt.IsZero() && UTCNow().After(u.ExpiresAt) {
bestEffortIAMExpiration(ctx, iamOS, getUserIdentityPath(user, userType))
}
return UserIdentity{}, errNoSuchUser
}
if u.Credentials.IsExpired() {
// Delete expired identity - ignoring errors here.
iamOS.deleteIAMConfig(ctx, getUserIdentityPath(user, userType))
iamOS.deleteIAMConfig(ctx, getMappedPolicyPath(user, userType, false))
bestEffortIAMExpiration(ctx, iamOS, getUserIdentityPath(user, userType))
return u, errNoSuchUser
}
@@ -272,16 +303,18 @@ func (iamOS *IAMObjectStore) loadUserIdentity(ctx context.Context, user string,
if u.Credentials.SessionToken != "" {
jwtClaims, err := extractJWTClaims(u)
if err != nil {
if u.Credentials.IsTemp() {
// We should delete such that the client can re-request
// for the expiring credentials.
iamOS.deleteIAMConfig(ctx, getUserIdentityPath(user, userType))
iamOS.deleteIAMConfig(ctx, getMappedPolicyPath(user, userType, false))
}
return u, errNoSuchUser
// During startup the site signing key may not be available yet.
// Reject this load without deleting a credential that has not expired.
return UserIdentity{}, errNoSuchUser
}
u.Credentials.Claims = jwtClaims.Map()
}
if err := checkIAMParentRevision(ctx, iamOS, u.Credentials); err != nil {
if errors.Is(err, errIAMStaleUpdate) {
return UserIdentity{}, errNoSuchUser
}
return UserIdentity{}, err
}
if u.Credentials.Description == "" {
u.Credentials.Description = u.Credentials.Comment
@@ -320,6 +353,7 @@ func (iamOS *IAMObjectStore) loadUser(ctx context.Context, user string, userType
}
func (iamOS *IAMObjectStore) loadUsers(ctx context.Context, userType IAMUserType, m map[string]UserIdentity) error {
ctx = withIAMExpirationCleanup(ctx)
var basePrefix string
switch userType {
case svcUser:
@@ -354,6 +388,9 @@ func (iamOS *IAMObjectStore) loadGroup(ctx context.Context, group string, m map[
}
return err
}
if g.Deleted {
return errNoSuchGroup
}
m[group] = g
return nil
}
@@ -391,6 +428,9 @@ func (iamOS *IAMObjectStore) loadMappedPolicyWithRetry(ctx context.Context, name
goto retry
}
if !iamOS.index.mappingAllowed(getMappedPolicyPath(name, userType, isGroup), p) {
return errNoSuchPolicy
}
m.Store(name, p)
return nil
}
@@ -405,6 +445,9 @@ func (iamOS *IAMObjectStore) loadMappedPolicyInternal(ctx context.Context, name
}
return p, err
}
if !iamOS.index.mappingAllowed(getMappedPolicyPath(name, userType, isGroup), p) {
return MappedPolicy{}, errNoSuchPolicy
}
return p, nil
}
@@ -824,54 +867,6 @@ func (iamOS *IAMObjectStore) loadAllFromObjStore(ctx context.Context, cache *iam
return nil
}
func (iamOS *IAMObjectStore) savePolicyDoc(ctx context.Context, policyName string, p PolicyDoc) error {
return iamOS.saveIAMConfig(ctx, &p, getPolicyDocPath(policyName))
}
func (iamOS *IAMObjectStore) saveMappedPolicy(ctx context.Context, name string, userType IAMUserType, isGroup bool, mp MappedPolicy, opts ...options) error {
return iamOS.saveIAMConfig(ctx, mp, getMappedPolicyPath(name, userType, isGroup), opts...)
}
func (iamOS *IAMObjectStore) saveUserIdentity(ctx context.Context, name string, userType IAMUserType, u UserIdentity, opts ...options) error {
return iamOS.saveIAMConfig(ctx, u, getUserIdentityPath(name, userType), opts...)
}
func (iamOS *IAMObjectStore) saveGroupInfo(ctx context.Context, name string, gi GroupInfo) error {
return iamOS.saveIAMConfig(ctx, gi, getGroupInfoPath(name))
}
func (iamOS *IAMObjectStore) deletePolicyDoc(ctx context.Context, name string) error {
err := iamOS.deleteIAMConfig(ctx, getPolicyDocPath(name))
if err == errConfigNotFound {
err = errNoSuchPolicy
}
return err
}
func (iamOS *IAMObjectStore) deleteMappedPolicy(ctx context.Context, name string, userType IAMUserType, isGroup bool) error {
err := iamOS.deleteIAMConfig(ctx, getMappedPolicyPath(name, userType, isGroup))
if err == errConfigNotFound {
err = errNoSuchPolicy
}
return err
}
func (iamOS *IAMObjectStore) deleteUserIdentity(ctx context.Context, name string, userType IAMUserType) error {
err := iamOS.deleteIAMConfig(ctx, getUserIdentityPath(name, userType))
if err == errConfigNotFound {
err = errNoSuchUser
}
return err
}
func (iamOS *IAMObjectStore) deleteGroupInfo(ctx context.Context, name string) error {
err := iamOS.deleteIAMConfig(ctx, getGroupInfoPath(name))
if err == errConfigNotFound {
err = errNoSuchGroup
}
return err
}
// Lists objects in the minioMetaBucket at the given path prefix. All returned
// items have the pathPrefix removed from their names.
func listIAMConfigItems(ctx context.Context, objAPI ObjectLayer, pathPrefix string) <-chan itemOrErr[string] {
+131
View File
@@ -0,0 +1,131 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"fmt"
"os"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
)
// Uses APIs shared with the pre-revision tree so the same benchmark can be
// overlaid on that tree for a comparable local baseline.
func prepareIAMPerformanceFixture(b *testing.B) (context.Context, *IAMSys) {
b.Helper()
resetTestGlobals()
ctx, cancel := context.WithCancel(context.Background())
disks, err := getRandomDisks(1)
if err != nil {
b.Fatal(err)
}
obj, _, err := initObjectLayer(ctx, mustGetPoolEndpoints(0, disks...))
if err != nil {
b.Fatal(err)
}
initAllSubsystems(ctx)
globalIAMSys.initStore(obj, nil)
if err := globalIAMSys.Load(ctx, true); err != nil {
b.Fatal(err)
}
b.Cleanup(func() { cancel(); obj.Shutdown(context.Background()); os.RemoveAll(disks[0]); resetTestGlobals() })
return ctx, globalIAMSys
}
func BenchmarkIAMCachedCredential(b *testing.B) {
for _, kind := range []string{"user", "service", "sts"} {
b.Run(kind, func(b *testing.B) {
ctx, sys := prepareIAMPerformanceFixture(b)
const parent = "benchmark-parent"
_, err := sys.CreateUser(ctx, parent, madmin.AddOrUpdateUserReq{SecretKey: "benchmark-user-password", Status: madmin.AccountEnabled})
if err != nil {
b.Fatal(err)
}
key := parent
if kind == "service" {
c, _, err := sys.NewServiceAccount(ctx, parent, nil, newServiceAccountOpts{accessKey: "benchmark-service", secretKey: "benchmark-service-password"})
if err != nil {
b.Fatal(err)
}
key = c.AccessKey
}
if kind == "sts" {
secret, err := getTokenSigningKey()
if err != nil {
b.Fatal(err)
}
c, err := auth.GetNewCredentialsWithMetadata(map[string]any{"exp": UTCNow().Add(time.Hour).Unix(), parentClaim: parent}, secret)
if err != nil {
b.Fatal(err)
}
c.ParentUser = parent
if _, err := sys.SetTempUser(ctx, c.AccessKey, c, ""); err != nil {
b.Fatal(err)
}
key = c.AccessKey
}
b.ReportAllocs()
b.ResetTimer()
for b.Loop() {
if _, ok := sys.store.GetUser(key); !ok {
b.Fatal("credential missing")
}
}
})
}
}
func BenchmarkIAMSetTempUser(b *testing.B) {
ctx, sys := prepareIAMPerformanceFixture(b)
const parent = "benchmark-sts-parent"
_, err := sys.CreateUser(ctx, parent, madmin.AddOrUpdateUserReq{SecretKey: "benchmark-user-password", Status: madmin.AccountEnabled})
if err != nil {
b.Fatal(err)
}
secret, err := getTokenSigningKey()
if err != nil {
b.Fatal(err)
}
cred, err := auth.GetNewCredentialsWithMetadata(map[string]any{"exp": UTCNow().Add(time.Hour).Unix(), parentClaim: parent}, secret)
if err != nil {
b.Fatal(err)
}
cred.ParentUser = parent
b.ReportAllocs()
b.ResetTimer()
for b.Loop() {
if _, err := sys.SetTempUser(ctx, cred.AccessKey, cred, "readwrite"); err != nil {
b.Fatal(err)
}
}
}
// Run with -benchtime=1x. Preparation is outside the timer; each measured load
// sees a fresh set of expired reusable service-account records.
func BenchmarkIAMColdLoadExpiredServices(b *testing.B) {
for _, count := range []int{100, 1000} {
b.Run(fmt.Sprint(count), func(b *testing.B) {
ctx, sys := prepareIAMPerformanceFixture(b)
b.ReportAllocs()
for i := 0; i < b.N; i++ {
b.StopTimer()
for j := 0; j < count; j++ {
key := fmt.Sprintf("expired-benchmark-%d-%d", i, j)
u := UserIdentity{Version: 1, UpdatedAt: UTCNow().Add(-2 * time.Hour), Credentials: auth.Credentials{AccessKey: key, SecretKey: "expired-benchmark-password", ParentUser: "absent-idp-parent", Expiration: UTCNow().Add(-time.Hour), Status: auth.AccountOn}}
if err := sys.store.saveIAMConfig(ctx, &u, getUserIdentityPath(key, svcUser)); err != nil {
b.Fatal(err)
}
}
b.StartTimer()
if err := sys.store.LoadIAMCache(ctx, true); err != nil {
b.Fatal(err)
}
}
})
}
}
+168
View File
@@ -0,0 +1,168 @@
package cmd
import (
"context"
"errors"
"os"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/pgsty/silo-pkg/v3/policy"
)
func TestReviewIAMRevokedUserReplay(t *testing.T) {
resetTestGlobals()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
obj, disk, err := prepareFS(ctx)
if err != nil {
t.Fatal(err)
}
defer os.RemoveAll(disk)
defer obj.Shutdown(ctx)
defer resetTestGlobals()
user := "review-revoked-user"
req := madmin.AddOrUpdateUserReq{SecretKey: "review-valid-password", Status: madmin.AccountEnabled}
created, err := globalIAMSys.CreateUser(ctx, user, req)
if err != nil {
t.Fatal(err)
}
policyAt, err := globalIAMSys.PolicyDBSet(ctx, user, "readwrite", regUser, false)
if err != nil {
t.Fatal(err)
}
args := policy.Args{AccountName: user, Action: policy.GetObjectAction, BucketName: "review-bucket", ObjectName: "review-object"}
if !globalIAMSys.IsAllowed(args) {
t.Fatal("seed must allow object read")
}
if err := globalIAMSys.DeleteUser(ctx, user, false); err != nil {
t.Fatal(err)
}
if err := globalIAMSys.store.LoadIAMCache(ctx, false); err != nil {
t.Fatal(err)
}
if globalIAMSys.IsAllowed(args) {
t.Fatal("deletion did not remove initial permission")
}
if err := globalSiteReplicationSys.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, UserReq: &req}, created); err != nil {
t.Fatal(err)
}
if _, err := globalIAMSys.GetUserInfo(ctx, user); !errors.Is(err, errNoSuchUser) {
t.Errorf("revoked user restored by an older replicated create, GetUserInfo error = %v", err)
}
if err := globalSiteReplicationSys.PeerPolicyMappingHandler(ctx, &madmin.SRPolicyMapping{UserOrGroup: user, UserType: int(regUser), Policy: "readwrite"}, policyAt); err != nil {
t.Fatal(err)
}
if globalIAMSys.IsAllowed(args) {
t.Error("older replicated identity and policy events restored revoked S3 read permission")
}
}
func TestReviewIAMSourceTimestampOrder(t *testing.T) {
resetTestGlobals()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
obj, disk, err := prepareFS(ctx)
if err != nil {
t.Fatal(err)
}
defer os.RemoveAll(disk)
defer obj.Shutdown(ctx)
defer resetTestGlobals()
user := "review-ordered-user"
req := madmin.AddOrUpdateUserReq{SecretKey: "review-valid-password", Status: madmin.AccountEnabled}
origin := UTCNow().Add(-time.Hour)
if err := globalSiteReplicationSys.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, UserReq: &req}, origin); err != nil {
t.Fatal(err)
}
if err := globalSiteReplicationSys.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, IsDeleteReq: true}, origin.Add(time.Minute)); err != nil {
t.Fatal(err)
}
if _, err := globalIAMSys.GetUserInfo(ctx, user); !errors.Is(err, errNoSuchUser) {
t.Fatalf("newer source deletion skipped after delayed creation, GetUserInfo error = %v", err)
}
}
// A user's old group grant must not return after deletion and deliberate recreation.
func TestR3CandidateOldGroupReplayAfterRecreation(t *testing.T) {
resetTestGlobals()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
obj, disk, err := prepareFS(ctx)
if err != nil {
t.Fatal(err)
}
defer os.RemoveAll(disk)
defer obj.Shutdown(ctx)
defer resetTestGlobals()
must := func(err error) {
t.Helper()
if err != nil {
t.Fatal(err)
}
}
user, group := "r3-group-member", "r3-granting-group"
origin := UTCNow().Add(-time.Hour)
req := madmin.AddOrUpdateUserReq{SecretKey: "valid-r3-user-password", Status: madmin.AccountEnabled}
peer := &globalSiteReplicationSys
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, UserReq: &req}, origin))
add := &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: group, Members: []string{user}}}
must(peer.PeerGroupInfoChangeHandler(ctx, add, origin.Add(time.Minute)))
must(peer.PeerPolicyMappingHandler(ctx, &madmin.SRPolicyMapping{UserOrGroup: group, IsGroup: true, UserType: int(regUser), Policy: "readwrite"}, origin.Add(time.Minute)))
args := policy.Args{AccountName: user, Action: policy.GetObjectAction, BucketName: "r3-bucket", ObjectName: "probe"}
if !globalIAMSys.IsAllowed(args) {
t.Fatal("fixture must grant through group")
}
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, IsDeleteReq: true}, origin.Add(2*time.Minute)))
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, UserReq: &req}, origin.Add(3*time.Minute)))
must(globalIAMSys.store.LoadIAMCache(ctx, false))
if globalIAMSys.IsAllowed(args) {
t.Fatal("recreation must start without deleted group membership")
}
must(peer.PeerGroupInfoChangeHandler(ctx, add, origin.Add(time.Minute)))
must(globalIAMSys.store.LoadIAMCache(ctx, false))
if globalIAMSys.IsAllowed(args) {
t.Fatal("old group event restored the deleted user's read grant after recreation and durable reload")
}
}
// The user delete is also a revocation of its earlier group memberships.
func TestR3CandidateLateDeleteRetainsOldGroupGrant(t *testing.T) {
resetTestGlobals()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
obj, disk, err := prepareFS(ctx)
if err != nil {
t.Fatal(err)
}
defer os.RemoveAll(disk)
defer obj.Shutdown(ctx)
defer resetTestGlobals()
must := func(err error) {
t.Helper()
if err != nil {
t.Fatal(err)
}
}
user, group := "r3-late-member", "r3-late-group"
origin := UTCNow().Add(-time.Hour)
req := madmin.AddOrUpdateUserReq{SecretKey: "valid-r3-user-password", Status: madmin.AccountEnabled}
peer := &globalSiteReplicationSys
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, UserReq: &req}, origin))
must(peer.PeerGroupInfoChangeHandler(ctx, &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: group, Members: []string{user}}}, origin.Add(time.Minute)))
must(peer.PeerPolicyMappingHandler(ctx, &madmin.SRPolicyMapping{UserOrGroup: group, IsGroup: true, UserType: int(regUser), Policy: "readwrite"}, origin.Add(time.Minute)))
args := policy.Args{AccountName: user, Action: policy.GetObjectAction, BucketName: "r3-bucket", ObjectName: "probe"}
if !globalIAMSys.IsAllowed(args) {
t.Fatal("fixture must grant through group")
}
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, UserReq: &req}, origin.Add(3*time.Minute)))
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, IsDeleteReq: true}, origin.Add(2*time.Minute)))
must(globalIAMSys.store.LoadIAMCache(ctx, false))
if _, ok := globalIAMSys.GetUser(ctx, user); !ok {
t.Fatal("newer identity must survive")
}
if globalIAMSys.IsAllowed(args) {
t.Fatal("late user deletion retained the older group grant on the recreated identity")
}
}
+321
View File
@@ -0,0 +1,321 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"sort"
"strings"
"sync/atomic"
"time"
"github.com/minio/madmin-go/v3"
xhttp "github.com/minio/minio/internal/http"
"github.com/pgsty/silo-pkg/v3/policy"
)
const (
iamRevisionProtocol = 1
iamRevisionPeerPath = "/v3/site-replication/peer/iam-revisions"
iamUserBoundaryType = "silo-user-revocation"
iamGroupBoundaryType = "silo-group-revocation"
maxIAMRevisionBatch = 128
)
var iamRevisionInstance = mustGetUUID()
type iamUserBoundary struct {
User string `json:"user"`
Before time.Time `json:"before"`
}
type iamGroupBoundary struct {
Group string `json:"group"`
Before time.Time `json:"before"`
}
// The server owns this additive protocol, without changing the client SDK or
// overloading a policy/document field. Old servers reject the dedicated route
// before applying any change that would lose revocation or member metadata.
type iamReplicationItem struct {
madmin.SRIAMItem
GroupGrants map[string]time.Time `json:"groupGrants,omitempty"`
GroupSnapshot bool `json:"groupSnapshot,omitempty"`
UserRevocation *iamUserBoundary `json:"userRevocation,omitempty"`
GroupRevocation *iamGroupBoundary `json:"groupRevocation,omitempty"`
RevokedBefore time.Time `json:"revokedBefore,omitempty"`
}
type iamRevisionBatch struct {
Version int `json:"version"`
Items []iamReplicationItem `json:"items"`
}
type iamRevisionStatus struct {
Version int `json:"version"`
Node string `json:"node"`
Instance string `json:"instance"`
Digest string `json:"digest"`
}
type iamRevisionResponse struct {
iamRevisionStatus
Errors []string `json:"errors,omitempty"`
}
type iamRevisionBatchError struct{ failures []string }
func (e *iamRevisionBatchError) Error() string {
return "IAM revision batch: " + strings.Join(e.failures, "; ")
}
type iamRevisionProgress struct {
Instances map[string]string
Acknowledged map[string]string
}
type iamRevisionMetrics struct {
healFailures atomic.Uint64
healLastSuccess atomic.Int64
healDurationMillis atomic.Int64
}
func iamRevisionDigest(items map[string]iamRevision) string {
paths := make([]string, 0, len(items))
for path := range items {
paths = append(paths, path)
}
sort.Strings(paths)
h := sha256.New()
for _, path := range paths {
r := items[path]
fmt.Fprintf(h, "%q %s %t %s\n", path, r.timestamp().UTC().Format(time.RFC3339Nano), r.Deleted, r.RevokedBefore.UTC().Format(time.RFC3339Nano))
}
return hex.EncodeToString(h.Sum(nil))
}
func (store *IAMStoreSys) iamRevisionStatus() iamRevisionStatus {
node := globalLocalNodeName
if node == "" {
node = "local"
}
return iamRevisionStatus{Version: iamRevisionProtocol, Node: node, Instance: iamRevisionInstance, Digest: store.revisionIndex().digest()}
}
func executeIAMRevisionRequest(ctx context.Context, client *madmin.AdminClient, method string, batch *iamRevisionBatch) (status iamRevisionStatus, err error) {
var content []byte
if batch != nil {
content, err = json.Marshal(batch)
if err != nil {
return status, err
}
}
resp, err := client.ExecuteMethod(ctx, method, madmin.RequestData{RelPath: iamRevisionPeerPath, QueryValues: url.Values{"api-version": {madmin.SiteReplAPIVersion}}, Content: content})
if resp != nil {
defer xhttp.DrainBody(resp.Body)
}
if err != nil {
return status, err
}
if resp.StatusCode != http.StatusOK {
var remote madmin.ErrorResponse
if json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(&remote) == nil && remote.Code != "" {
return status, remote
}
return status, fmt.Errorf("IAM revision protocol requires upgraded peers: %s", resp.Status)
}
var response iamRevisionResponse
if err = json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(&response); err != nil {
return status, err
}
status = response.iamRevisionStatus
if status.Version != iamRevisionProtocol || status.Node == "" || status.Instance == "" || status.Digest == "" {
return status, errors.New("peer did not acknowledge the IAM revision protocol")
}
if len(response.Errors) != 0 {
return status, &iamRevisionBatchError{failures: response.Errors}
}
return status, nil
}
type (
iamRecordBoundaryKey struct{}
iamGroupSnapshotKey struct{}
)
func (c *SiteReplicationSys) replicationItem(ctx context.Context, item madmin.SRIAMItem) (iamReplicationItem, error) {
out := iamReplicationItem{SRIAMItem: item}
if item.Type == madmin.SRIAMItemSvcAcc && item.SvcAccChange != nil {
var key string
if item.SvcAccChange.Create != nil {
key = item.SvcAccChange.Create.AccessKey
} else if item.SvcAccChange.Update != nil {
key = item.SvcAccChange.Update.AccessKey
}
if key != "" {
r, err := loadIAMRevision(ctx, globalIAMSys.store, getUserIdentityPath(key, svcUser))
if err != nil {
return out, err
}
if r.Deleted || r.timestamp().After(item.UpdatedAt) {
return out, errIAMStaleUpdate
}
out.RevokedBefore = r.RevokedBefore
}
}
if item.Type == madmin.SRIAMItemGroupInfo && item.GroupInfo != nil && !item.GroupInfo.UpdateReq.IsRemove {
out.GroupSnapshot = true
var gi GroupInfo
if err := globalIAMSys.store.loadIAMConfig(ctx, &gi, getGroupInfoPath(item.GroupInfo.UpdateReq.Group)); err != nil {
return out, err
}
// The matching persisted snapshot carries member grant times. If a
// later write won before sending, propagate that whole newer state.
if gi.Deleted {
return out, errIAMStaleUpdate
}
out.UpdatedAt = gi.UpdatedAt
out.RevokedBefore = gi.RevokedBefore
out.GroupInfo = &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: item.GroupInfo.UpdateReq.Group, Status: madmin.GroupStatus(gi.Status)}}
cache := globalIAMSys.store.rlock()
out.GroupInfo.UpdateReq.Members = cache.effectiveGroupMembers(gi)
globalIAMSys.store.runlock()
out.GroupGrants = make(map[string]time.Time, len(out.GroupInfo.UpdateReq.Members))
for _, member := range out.GroupInfo.UpdateReq.Members {
out.GroupGrants[member] = gi.MemberGrants[member]
}
}
if item.Type == madmin.SRIAMItemGroupInfo && item.GroupInfo != nil && item.GroupInfo.UpdateReq.IsRemove && len(item.GroupInfo.UpdateReq.Members) == 0 {
r, err := loadIAMRevision(ctx, globalIAMSys.store, getGroupInfoPath(item.GroupInfo.UpdateReq.Group))
if err != nil {
return out, err
}
if !r.Deleted && !r.RevokedBefore.IsZero() {
out.Type, out.GroupInfo = iamGroupBoundaryType, nil
out.GroupRevocation = &iamGroupBoundary{Group: item.GroupInfo.UpdateReq.Group, Before: r.RevokedBefore}
out.UpdatedAt = r.RevokedBefore
}
}
if item.Type == madmin.SRIAMItemIAMUser && item.IAMUser != nil {
r, err := loadIAMRevision(ctx, globalIAMSys.store, getUserIdentityPath(item.IAMUser.AccessKey, regUser))
if err != nil {
return out, err
}
if item.IAMUser.IsDeleteReq && !r.Deleted && !r.RevokedBefore.IsZero() {
out.Type = iamUserBoundaryType
out.IAMUser = nil
out.UserRevocation = &iamUserBoundary{User: item.IAMUser.AccessKey, Before: r.RevokedBefore}
out.UpdatedAt = r.RevokedBefore
} else if !item.IAMUser.IsDeleteReq {
if r.Deleted || r.timestamp().After(item.UpdatedAt) {
return out, errIAMStaleUpdate
}
out.RevokedBefore = r.RevokedBefore
}
}
return out, nil
}
func applyIAMReplicationItem(ctx context.Context, item iamReplicationItem) error {
if item.GroupInfo != nil {
if item.GroupSnapshot {
ctx = context.WithValue(ctx, iamGroupSnapshotKey{}, true)
}
// A nil map also explicitly denotes unknown legacy grants. Do not
// turn an unrelated group edit into a new grant after a revocation.
ctx = withIAMGroupGrants(ctx, item.GroupGrants)
}
if !item.RevokedBefore.IsZero() {
if item.RevokedBefore.After(item.UpdatedAt) {
return errSRInvalidRequest(errInvalidArgument)
}
ctx = context.WithValue(ctx, iamRecordBoundaryKey{}, item.RevokedBefore)
}
switch item.Type {
case iamUserBoundaryType:
if item.UserRevocation == nil || item.UserRevocation.User == "" || item.UserRevocation.Before.IsZero() {
return errSRInvalidRequest(errInvalidArgument)
}
return iamReplicationError(globalIAMSys.DeleteUser(withIAMReplicationTime(ctx, item.UserRevocation.Before), item.UserRevocation.User, true))
case iamGroupBoundaryType:
if item.GroupRevocation == nil || item.GroupRevocation.Group == "" || item.GroupRevocation.Before.IsZero() {
return errSRInvalidRequest(errInvalidArgument)
}
_, err := globalIAMSys.RemoveUsersFromGroup(withIAMReplicationTime(ctx, item.GroupRevocation.Before), item.GroupRevocation.Group, nil)
return iamReplicationError(err)
case madmin.SRIAMItemPolicy:
if len(item.Policy) == 0 {
return globalSiteReplicationSys.PeerAddPolicyHandler(ctx, item.Name, nil, item.UpdatedAt)
}
p, err := policy.ParseConfig(bytes.NewReader(item.Policy))
if err != nil {
return err
}
if p.IsEmpty() {
p = nil
}
return globalSiteReplicationSys.PeerAddPolicyHandler(ctx, item.Name, p, item.UpdatedAt)
case madmin.SRIAMItemSvcAcc:
return globalSiteReplicationSys.PeerSvcAccChangeHandler(ctx, item.SvcAccChange, item.UpdatedAt)
case madmin.SRIAMItemPolicyMapping:
return globalSiteReplicationSys.PeerPolicyMappingHandler(ctx, item.PolicyMapping, item.UpdatedAt)
case madmin.SRIAMItemSTSAcc:
return globalSiteReplicationSys.PeerSTSAccHandler(ctx, item.STSCredential, item.UpdatedAt)
case madmin.SRIAMItemIAMUser:
return globalSiteReplicationSys.PeerIAMUserChangeHandler(ctx, item.IAMUser, item.UpdatedAt)
case madmin.SRIAMItemGroupInfo:
return globalSiteReplicationSys.PeerGroupInfoChangeHandler(ctx, item.GroupInfo, item.UpdatedAt)
default:
return errSRInvalidRequest(errInvalidArgument)
}
}
func (a adminAPIHandlers) SRPeerIAMRevisions(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
if obj, _ := validateAdminReq(ctx, w, r, policy.SiteReplicationOperationAction); obj == nil {
return
}
var failures []string
if r.Method == http.MethodPut {
var batch iamRevisionBatch
if err := parseJSONBody(ctx, r.Body, &batch, ""); err != nil {
writeErrorResponseJSON(ctx, w, toAdminAPIErr(ctx, err), r.URL)
return
}
if batch.Version != iamRevisionProtocol || len(batch.Items) == 0 || len(batch.Items) > maxIAMRevisionBatch {
writeErrorResponseJSON(ctx, w, toAdminAPIErr(ctx, errSRInvalidRequest(errInvalidArgument)), r.URL)
return
}
for i, item := range batch.Items {
if err := applyIAMReplicationItem(ctx, item); err != nil {
failures = append(failures, fmt.Sprintf("item %d (%s): %v", i, item.Type, err))
}
}
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(iamRevisionResponse{iamRevisionStatus: globalIAMSys.store.iamRevisionStatus(), Errors: failures})
}
// A site endpoint can balance requests across nodes sharing durable IAM state.
// Switching between known node incarnations preserves ACKs; a new incarnation
// conservatively invalidates them so restoring an old backend cannot inherit
// acknowledgements from before the restore.
func (p *iamRevisionProgress) observePeer(status iamRevisionStatus) {
if p.Instances == nil {
p.Instances = make(map[string]string)
}
if p.Instances[status.Node] != status.Instance || p.Acknowledged == nil {
p.Instances[status.Node] = status.Instance
p.Acknowledged = make(map[string]string)
}
}
+161
View File
@@ -0,0 +1,161 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"sync"
"sync/atomic"
"testing"
"time"
"github.com/minio/madmin-go/v3"
)
func TestIAMRevisionProtocolDoesNotFallBackToLegacy(t *testing.T) {
var requests atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
requests.Add(1)
if r.URL.Path != "/minio/admin/v3/site-replication/peer/iam-revisions" {
t.Errorf("unsafe fallback path: %s", r.URL.Path)
}
w.WriteHeader(http.StatusNotFound)
_, _ = w.Write([]byte(`{"Code":"NotImplemented","Message":"old server"}`))
}))
defer server.Close()
client, err := madmin.New(strings.TrimPrefix(server.URL, "http://"), "test-access", "valid-test-secret", false)
mustIAM(t, err)
_, err = executeIAMRevisionRequest(context.Background(), client, http.MethodPut, &iamRevisionBatch{Version: iamRevisionProtocol, Items: []iamReplicationItem{{SRIAMItem: madmin.SRIAMItem{Type: iamUserBoundaryType}, UserRevocation: &iamUserBoundary{User: "recreated", Before: UTCNow()}}}})
if err == nil || requests.Load() != 1 {
t.Fatalf("old peer must reject without fallback, err=%v requests=%d", err, requests.Load())
}
}
type iamNoHealingScanStore struct{ IAMStorageAPI }
func (s *iamNoHealingScanStore) listIAMConfigPaths(context.Context) ([]string, error) {
panic("healing must use the loaded revision index")
}
func TestIAMRevisionHealingAcknowledgements(t *testing.T) {
for _, balanced := range []bool{false, true} {
t.Run(fmt.Sprintf("load_balanced_%t", balanced), func(t *testing.T) { testIAMRevisionHealingAcknowledgements(t, balanced) })
}
}
func testIAMRevisionHealingAcknowledgements(t *testing.T, balanced bool) {
ctx, sys, _ := prepareIAMRevisionFixture(t)
_, err := sys.CreateUser(ctx, "ack-sync", madmin.AddOrUpdateUserReq{SecretKey: "valid-sync-password", Status: madmin.AccountEnabled})
mustIAM(t, err)
for i := range maxIAMRevisionBatch*2 + 1 {
at := UTCNow().Add(time.Duration(i) * time.Nanosecond)
mustIAM(t, sys.store.saveIAMConfig(ctx, &UserIdentity{Version: 1, Deleted: true, UpdatedAt: at, RevokedBefore: at}, getUserIdentityPath(fmt.Sprintf("ack-%04d", i), regUser)))
}
sys.store.IAMStorageAPI = &iamNoHealingScanStore{IAMStorageAPI: sys.store.IAMStorageAPI}
var mu sync.Mutex
var puts, gets int
var applied int
instance := "boot-1"
failSecondBatch := true
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/minio/health/live" {
w.WriteHeader(http.StatusOK)
return
}
mu.Lock()
defer mu.Unlock()
if r.URL.Path != "/minio/admin/v3/site-replication/peer/iam-revisions" {
t.Errorf("unexpected request: %s", r.URL.Path)
w.WriteHeader(404)
return
}
var failures []string
if r.Method == http.MethodGet {
gets++
} else {
puts++
var batch iamRevisionBatch
if err := json.NewDecoder(r.Body).Decode(&batch); err != nil {
t.Error(err)
w.WriteHeader(400)
return
}
if len(batch.Items) > maxIAMRevisionBatch {
t.Error("batch exceeds limit")
}
if failSecondBatch && puts == 2 {
failures = []string{"injected item error"}
} else {
applied += len(batch.Items)
}
}
node := "node-1"
if balanced {
node = fmt.Sprintf("node-%d", (gets+puts)%2+1)
}
_ = json.NewEncoder(w).Encode(iamRevisionResponse{iamRevisionStatus: iamRevisionStatus{Version: iamRevisionProtocol, Node: node, Instance: node + instance, Digest: fmt.Sprintf("%d", applied)}, Errors: failures})
}))
defer server.Close()
c := &SiteReplicationSys{enabled: true, state: srState{ServiceAccountAccessKey: "ack-sync", Peers: map[string]madmin.PeerInfo{globalDeploymentID(): {DeploymentID: globalDeploymentID(), Name: "local"}, "remote": {DeploymentID: "remote", Name: "remote", Endpoint: server.URL}}}}
if err := c.healIAMDeletions(ctx); err == nil {
t.Fatal("item failure was hidden")
}
mu.Lock()
if puts != 3 || applied != maxIAMRevisionBatch+1 {
t.Errorf("failed middle batch blocked later revocations: puts=%d applied=%d", puts, applied)
}
failSecondBatch = false
applied++ // Unrelated remote mutation changes its digest.
mu.Unlock()
at := UTCNow()
mustIAM(t, sys.store.saveIAMConfig(ctx, &UserIdentity{Version: 1, Deleted: true, UpdatedAt: at, RevokedBefore: at}, getUserIdentityPath("ack-new-local", regUser)))
mustIAM(t, c.healIAMDeletions(ctx))
mu.Lock()
if puts != 5 {
t.Errorf("did not resume at unacknowledged batch: puts=%d", puts)
}
mu.Unlock()
mustIAM(t, c.healIAMDeletions(ctx))
mu.Lock()
if puts != 5 || gets != 3 {
t.Errorf("converged records were replayed: puts=%d gets=%d", puts, gets)
}
instance = "boot-2"
mu.Unlock()
mustIAM(t, c.healIAMDeletions(ctx))
mu.Lock()
defer mu.Unlock()
if puts != 8 {
t.Fatalf("peer restart reused an old acknowledgement: puts=%d", puts)
}
}
func TestIAMRevisionIndexRebuildsFromStorage(t *testing.T) {
ctx, sys, obj := prepareIAMRevisionFixture(t)
const user = "index-parent"
req := madmin.AddOrUpdateUserReq{SecretKey: "valid-parent-password", Status: madmin.AccountEnabled}
_, err := sys.CreateUser(ctx, user, req)
mustIAM(t, err)
mustIAM(t, sys.DeleteUser(ctx, user, false))
before := sys.store.revisionIndex().snapshot()
store := &IAMStoreSys{IAMStorageAPI: newIAMObjectStore(obj, MinIOUsersSysType)}
mustIAM(t, store.LoadIAMCache(ctx, true))
if iamRevisionDigest(before) != iamRevisionDigest(store.revisionIndex().snapshot()) {
t.Fatal("ordinary IAM loading did not restore the deletion index")
}
_, err = store.AddUser(ctx, user, req)
mustIAM(t, err)
r := store.revisionIndex().get(getUserIdentityPath(user, regUser))
if r.Deleted || r.RevokedBefore.IsZero() {
t.Fatal("recreation discarded the retained boundary")
}
if r.Credentials.SecretKey != "" || r.Credentials.SessionToken != "" {
t.Fatal("index retained credentials")
}
}
+393
View File
@@ -0,0 +1,393 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"errors"
"fmt"
"os"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/grid"
xnet "github.com/pgsty/silo-pkg/v3/net"
"github.com/pgsty/silo-pkg/v3/policy"
etcd "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/namespace"
)
func prepareIAMRevisionFixture(t testing.TB, backend ...string) (context.Context, *IAMSys, ObjectLayer) {
t.Helper()
resetTestGlobals()
ctx, cancel := context.WithCancel(context.Background())
disks, err := getRandomDisks(1)
mustIAM(t, err)
obj, _, err := initObjectLayer(ctx, mustGetPoolEndpoints(0, disks...))
mustIAM(t, err)
initAllSubsystems(ctx)
// Deliberately omit the periodic refresh goroutine. Fault injection can
// replace this fixture's storage interface without racing initialization.
var client *etcd.Client
if len(backend) != 0 && backend[0] == "etcd" {
endpoint := os.Getenv("SILO_TEST_IAM_REVOCATION_ETCD")
if endpoint == "" {
cancel()
obj.Shutdown(context.Background())
os.RemoveAll(disks[0])
t.Skip("set SILO_TEST_IAM_REVOCATION_ETCD to a disposable etcd endpoint")
}
client, err = etcd.New(etcd.Config{Endpoints: strings.Split(endpoint, ","), DialTimeout: 5 * time.Second})
mustIAM(t, err)
prefix := fmt.Sprintf("/silo-boundary-test/%d/", time.Now().UnixNano())
client.KV = namespace.NewKV(client.KV, prefix)
client.Watcher = namespace.NewWatcher(client.Watcher, prefix)
t.Cleanup(func() { client.Delete(context.Background(), "", etcd.WithPrefix()); client.Close() })
}
globalIAMSys.initStore(obj, client)
mustIAM(t, globalIAMSys.Load(ctx, true))
t.Cleanup(func() { cancel(); obj.Shutdown(context.Background()); os.RemoveAll(disks[0]); resetTestGlobals() })
return ctx, globalIAMSys, obj
}
func mustIAM(t testing.TB, err error) {
t.Helper()
if err != nil {
t.Fatal(err)
}
}
var errIAMInjectedWrite = errors.New("injected IAM persistence failure")
type iamFailingCleanupStore struct {
IAMStorageAPI
parentPath string
beforeCommit bool
}
func (s *iamFailingCleanupStore) saveIAMConfig(ctx context.Context, item any, path string, opts ...options) error {
if s.beforeCommit || path != s.parentPath {
return errIAMInjectedWrite
}
return s.IAMStorageAPI.saveIAMConfig(ctx, item, path, opts...)
}
func TestIAMRevocationCommitBoundary(t *testing.T) {
for _, before := range []bool{true, false} {
name := "after_identity_commit"
if before {
name = "before_identity_commit"
}
t.Run(name, func(t *testing.T) {
ctx, sys, obj := prepareIAMRevisionFixture(t)
const user = "commit-boundary-user"
origin := UTCNow().Add(-time.Hour)
req := madmin.AddOrUpdateUserReq{SecretKey: "valid-test-password", Status: madmin.AccountEnabled}
_, err := sys.CreateUser(withIAMReplicationTime(ctx, origin), user, req)
mustIAM(t, err)
_, err = sys.PolicyDBSet(withIAMReplicationTime(ctx, origin.Add(time.Minute)), user, "readwrite", regUser, false)
mustIAM(t, err)
_, err = sys.AddUsersToGroup(withIAMReplicationTime(ctx, origin.Add(time.Minute)), "commit-group", []string{user})
mustIAM(t, err)
_, err = sys.PolicyDBSet(ctx, "commit-group", "readwrite", regUser, true)
mustIAM(t, err)
child, _, err := sys.NewServiceAccount(withIAMReplicationTime(ctx, origin), user, nil, newServiceAccountOpts{accessKey: "commit-child", secretKey: "valid-child-password"})
mustIAM(t, err)
args := policy.Args{AccountName: user, Action: policy.GetObjectAction, BucketName: "bucket", ObjectName: "object"}
if !sys.IsAllowed(args) {
t.Fatal("fixture has no grant")
}
siblingStore := &IAMStoreSys{IAMStorageAPI: newIAMObjectStore(obj, MinIOUsersSysType)}
mustIAM(t, siblingStore.LoadIAMCache(ctx, true))
sibling := &IAMSys{store: siblingStore, usersSysType: MinIOUsersSysType}
tg, err := grid.SetupTestGrid(2)
mustIAM(t, err)
defer tg.Cleanup()
var notifications atomic.Int32
mustIAM(t, deleteUserRPC.Register(tg.Managers[1], func(r *grid.MSS) (grid.NoPayload, *grid.RemoteErr) {
notifications.Add(1)
if err := sibling.LoadUserAfterDelete(ctx, r.Get(peerRESTUser)); err != nil {
return grid.NoPayload{}, grid.NewRemoteErr(err)
}
return grid.NoPayload{}, nil
}))
host, err := xnet.ParseHost(strings.TrimPrefix(tg.Hosts[1], "http://"))
mustIAM(t, err)
globalNotificationSys = &NotificationSys{peerClients: []*peerRESTClient{{host: host, gridConn: func() *grid.Connection { return tg.Managers[0].Connection(tg.Hosts[1]) }}}}
original := sys.store.IAMStorageAPI
sys.store.IAMStorageAPI = &iamFailingCleanupStore{IAMStorageAPI: original, parentPath: getUserIdentityPath(user, regUser), beforeCommit: before}
boundary := origin.Add(2 * time.Minute)
err = sys.DeleteUser(withIAMReplicationTime(ctx, boundary), user, true)
if !errors.Is(err, errIAMInjectedWrite) {
t.Fatalf("expected write failure, got %v", err)
}
sys.store.IAMStorageAPI = original
r, err := loadIAMRevision(ctx, original, getUserIdentityPath(user, regUser))
mustIAM(t, err)
if before {
if r.Deleted || !sys.IsAllowed(args) || !sibling.IsAllowed(args) || notifications.Load() != 0 {
t.Fatal("failure before commit changed the identity or grant")
}
return
}
if !r.Deleted || !r.RevokedBefore.Equal(boundary) {
t.Fatal("cleanup failure lost durable revocation")
}
if sys.IsAllowed(args) || sibling.IsAllowed(args) || notifications.Load() != 1 {
t.Fatal("cleanup failure retained old permission")
}
// Subsequent fixture writes need no additional RPC handlers.
globalNotificationSys = &NotificationSys{}
// Recreate after the partial cleanup. The old mapping, group member
// and child still exist in storage; none may authorize this identity.
_, err = sys.CreateUser(withIAMReplicationTime(ctx, origin.Add(3*time.Minute)), user, req)
mustIAM(t, err)
reloaded := &IAMStoreSys{IAMStorageAPI: newIAMObjectStore(obj, MinIOUsersSysType)}
mustIAM(t, reloaded.LoadIAMCache(ctx, true))
fresh := &IAMSys{store: reloaded, usersSysType: MinIOUsersSysType}
if fresh.IsAllowed(args) {
t.Fatal("cold reload restored partially cleaned-up grants")
}
if _, ok := reloaded.GetUser(child.AccessKey); ok {
t.Fatal("cold reload restored the old child")
}
gd, err := reloaded.GetGroupDescription("commit-group")
mustIAM(t, err)
if len(gd.Members) != 0 {
t.Fatalf("listing exposed a revoked group relation: %v", gd.Members)
}
_, err = sys.AddUsersToGroup(ctx, "commit-group", []string{user})
mustIAM(t, err)
if !sys.IsAllowed(args) {
t.Fatal("explicit new group grant was not accepted")
}
})
}
}
func TestIAMGroupGrantVersionsSurviveSnapshotsAndRecreation(t *testing.T) {
ctx, sys, _ := prepareIAMRevisionFixture(t)
origin := UTCNow().Add(-time.Hour)
req := madmin.AddOrUpdateUserReq{SecretKey: "valid-test-password", Status: madmin.AccountEnabled}
for _, user := range []string{"grant-alice", "grant-bob"} {
_, err := sys.CreateUser(withIAMReplicationTime(ctx, origin), user, req)
mustIAM(t, err)
}
grant := origin.Add(time.Minute)
_, err := sys.AddUsersToGroup(withIAMReplicationTime(ctx, grant), "grant-group", []string{"grant-alice"})
mustIAM(t, err)
_, err = sys.PolicyDBSet(ctx, "grant-group", "readwrite", regUser, true)
mustIAM(t, err)
boundary := origin.Add(2 * time.Minute)
mustIAM(t, sys.DeleteUser(withIAMReplicationTime(ctx, boundary), "grant-alice", false))
_, err = sys.CreateUser(withIAMReplicationTime(ctx, origin.Add(3*time.Minute)), "grant-alice", req)
mustIAM(t, err)
_, err = sys.AddUsersToGroup(withIAMReplicationTime(ctx, origin.Add(4*time.Minute)), "grant-group", []string{"grant-bob"})
mustIAM(t, err)
_, err = sys.SetGroupStatus(withIAMReplicationTime(ctx, origin.Add(5*time.Minute)), "grant-group", true)
mustIAM(t, err)
var gi GroupInfo
mustIAM(t, sys.store.loadIAMConfig(ctx, &gi, getGroupInfoPath("grant-group")))
if !gi.MemberGrants["grant-alice"].Equal(grant) {
t.Fatal("unrelated group edits refreshed an old grant")
}
args := policy.Args{AccountName: "grant-alice", Action: policy.GetObjectAction, BucketName: "bucket", ObjectName: "object"}
for _, stale := range []time.Time{grant, boundary, {}} {
item := iamReplicationItem{SRIAMItem: madmin.SRIAMItem{Type: madmin.SRIAMItemGroupInfo, UpdatedAt: origin.Add(6 * time.Minute), GroupInfo: &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: "grant-group", Members: []string{"grant-alice", "grant-bob"}}}}, GroupSnapshot: true, GroupGrants: map[string]time.Time{"grant-alice": stale, "grant-bob": origin.Add(4 * time.Minute)}}
mustIAM(t, applyIAMReplicationItem(ctx, item))
mustIAM(t, sys.store.LoadIAMCache(ctx, false))
if sys.IsAllowed(args) {
t.Fatalf("snapshot restored revoked grant %s", stale)
}
gd, err := sys.GetGroupDescription("grant-group")
mustIAM(t, err)
if len(gd.Members) != 1 || gd.Members[0] != "grant-bob" {
t.Fatalf("inconsistent effective members: %v", gd.Members)
}
}
// Only an explicit post-revocation grant restores access.
freshAt, err := sys.AddUsersToGroup(ctx, "grant-group", []string{"grant-alice"})
mustIAM(t, err)
if !sys.IsAllowed(args) {
t.Fatal("explicit regrant rejected")
}
mustIAM(t, sys.store.LoadIAMCache(ctx, false))
mustIAM(t, sys.store.loadIAMConfig(ctx, &gi, getGroupInfoPath("grant-group")))
if !gi.MemberGrants["grant-alice"].Equal(freshAt) {
t.Fatal("new grant version was not persisted")
}
if !gi.MemberGrants["grant-bob"].Equal(origin.Add(4 * time.Minute)) {
t.Fatal("regranting Alice changed Bob's grant")
}
}
func TestIAMGroupRevocationCommitAndRecreation(t *testing.T) {
for _, backend := range []string{"object", "etcd"} {
t.Run(backend, func(t *testing.T) { testIAMGroupRevocationCommitAndRecreation(t, backend) })
}
}
func testIAMGroupRevocationCommitAndRecreation(t *testing.T, backend string) {
ctx, sys, obj := prepareIAMRevisionFixture(t, backend)
origin := UTCNow().Add(-time.Hour)
user, group := "group-boundary-user", "group-boundary"
_, err := sys.CreateUser(withIAMReplicationTime(ctx, origin), user, madmin.AddOrUpdateUserReq{SecretKey: "valid-user-password", Status: madmin.AccountEnabled})
mustIAM(t, err)
grant, boundary := origin.Add(time.Minute), origin.Add(2*time.Minute)
_, err = sys.AddUsersToGroup(withIAMReplicationTime(ctx, grant), group, []string{user})
mustIAM(t, err)
// A newer mapping must not veto the authoritative group deletion.
_, err = sys.PolicyDBSet(withIAMReplicationTime(ctx, origin.Add(3*time.Minute)), group, "readwrite", regUser, true)
mustIAM(t, err)
_, err = sys.RemoveUsersFromGroup(withIAMReplicationTime(ctx, boundary), group, nil)
mustIAM(t, err)
r, err := loadIAMRevision(ctx, sys.store, getGroupInfoPath(group))
mustIAM(t, err)
if !r.Deleted || !r.RevokedBefore.Equal(boundary) {
t.Fatal("newer mapping swallowed group deletion")
}
_, err = sys.AddUsersToGroup(withIAMReplicationTime(ctx, origin.Add(4*time.Minute)), group, nil)
mustIAM(t, err)
for _, at := range []time.Time{grant, boundary, {}} {
item := iamReplicationItem{SRIAMItem: madmin.SRIAMItem{Type: madmin.SRIAMItemGroupInfo, UpdatedAt: origin.Add(5 * time.Minute), GroupInfo: &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: group, Members: []string{user}}}}, GroupSnapshot: true, GroupGrants: map[string]time.Time{user: at}}
mustIAM(t, applyIAMReplicationItem(ctx, item))
gd, err := sys.GetGroupDescription(group)
mustIAM(t, err)
if len(gd.Members) != 0 {
t.Fatalf("group recreation restored grant %s", at)
}
}
_, err = sys.AddUsersToGroup(ctx, group, []string{user})
mustIAM(t, err)
args := policy.Args{AccountName: user, Action: policy.GetObjectAction, BucketName: "bucket", ObjectName: "object"}
if !sys.IsAllowed(args) {
t.Fatal("explicit group regrant was rejected")
}
// The newer live snapshot may arrive before an older group deletion.
lateBoundary := origin.Add(6 * time.Minute)
_, err = sys.RemoveUsersFromGroup(withIAMReplicationTime(ctx, lateBoundary), group, nil)
mustIAM(t, err)
r, err = loadIAMRevision(ctx, sys.store, getGroupInfoPath(group))
mustIAM(t, err)
if r.Deleted || !r.RevokedBefore.Equal(lateBoundary) {
t.Fatal("late deletion lost the live group's revocation boundary")
}
// The old mapping is now revoked; a new explicit mapping restores access.
if sys.IsAllowed(args) {
t.Fatal("late group boundary retained an old mapping")
}
_, err = sys.PolicyDBSet(ctx, group, "readwrite", regUser, true)
mustIAM(t, err)
store := &IAMStoreSys{IAMStorageAPI: newIAMObjectStore(obj, MinIOUsersSysType)}
if es, ok := sys.store.IAMStorageAPI.(*IAMEtcdStore); ok {
store.IAMStorageAPI = newIAMEtcdStore(es.client, MinIOUsersSysType)
}
mustIAM(t, store.LoadIAMCache(ctx, true))
fresh := &IAMSys{store: store, usersSysType: MinIOUsersSysType}
if !fresh.IsAllowed(args) {
t.Fatal("reload lost explicit grants after a retained group boundary")
}
item, err := globalSiteReplicationSys.replicationItem(ctx, madmin.SRIAMItem{Type: madmin.SRIAMItemGroupInfo, GroupInfo: &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: group}}, UpdatedAt: r.timestamp()})
mustIAM(t, err)
if !item.RevokedBefore.Equal(lateBoundary) || !item.GroupGrants[user].After(lateBoundary) {
t.Fatal("group snapshot lost revision metadata")
}
}
// A committed revision is observable before all cached dependents have been
// cleaned up. Every authorization read must apply that boundary in this window.
func TestIAMCachedMappingHonorsCommittedRevision(t *testing.T) {
ctx, sys, _ := prepareIAMRevisionFixture(t)
origin := UTCNow().Add(-time.Hour)
parent := "cached-external-parent"
_, err := sys.PolicyDBSet(withIAMReplicationTime(ctx, origin), parent, "readwrite", stsUser, false)
mustIAM(t, err)
policies, err := sys.PolicyDBGet(parent)
mustIAM(t, err)
if len(policies) == 0 {
t.Fatal("fixture has no STS-parent mapping")
}
mustIAM(t, sys.store.saveIAMConfig(ctx, &MappedPolicy{Version: 1, Deleted: true, UpdatedAt: origin.Add(time.Minute)}, getMappedPolicyPath(parent, stsUser, false)))
policies, err = sys.PolicyDBGet(parent)
mustIAM(t, err)
if len(policies) != 0 {
t.Fatal("cached STS mapping ignored its own namespace tombstone")
}
user, group := "cached-group-user", "cached-group"
_, err = sys.CreateUser(withIAMReplicationTime(ctx, origin), user, madmin.AddOrUpdateUserReq{SecretKey: "valid-user-password", Status: madmin.AccountEnabled})
mustIAM(t, err)
grant := origin.Add(5 * time.Minute)
_, err = sys.AddUsersToGroup(withIAMReplicationTime(ctx, grant), group, []string{user})
mustIAM(t, err)
_, err = sys.PolicyDBSet(withIAMReplicationTime(ctx, origin), group, "readwrite", regUser, true)
mustIAM(t, err)
args := policy.Args{AccountName: user, Action: policy.GetObjectAction, BucketName: "bucket", ObjectName: "object"}
if !sys.IsAllowed(args) {
t.Fatal("fixture has no group grant")
}
// A late deletion preserves the newer member grant but revokes the older
// policy mapping. Simulate the interval before mapping cleanup completes.
gi := GroupInfo{Version: 1, Status: statusEnabled, Members: []string{user}, MemberGrants: map[string]time.Time{user: grant}, UpdatedAt: grant, RevokedBefore: origin.Add(2 * time.Minute)}
mustIAM(t, sys.store.saveIAMConfig(ctx, &gi, getGroupInfoPath(group)))
if sys.IsAllowed(args) {
t.Fatal("cached group mapping ignored the committed group boundary")
}
gd, err := sys.GetGroupDescription(group)
mustIAM(t, err)
if gd.Policy != "" {
t.Fatal("group listing exposed a revoked mapping")
}
}
type (
iamExpiryLockFailure struct {
ObjectLayer
path string
}
iamFailedExpiryLock struct{ RWLocker }
)
func (o *iamExpiryLockFailure) NewNSLock(bucket string, objects ...string) RWLocker {
lock := o.ObjectLayer.NewNSLock(bucket, objects...)
if bucket == minioMetaBucket && len(objects) == 1 && objects[0] == o.path+".revision-lock" {
return &iamFailedExpiryLock{RWLocker: lock}
}
return lock
}
func (l *iamFailedExpiryLock) GetLock(context.Context, *dynamicTimeout) (LockContext, error) {
return LockContext{}, errIAMInjectedWrite
}
func TestIAMExpiredCredentialCleanupDoesNotBlockLoading(t *testing.T) {
ctx, sys, obj := prepareIAMRevisionFixture(t)
_, err := sys.CreateUser(ctx, "healthy-user", madmin.AddOrUpdateUserReq{SecretKey: "healthy-user-password", Status: madmin.AccountEnabled})
mustIAM(t, err)
_, err = sys.PolicyDBSet(ctx, "healthy-user", "readwrite", regUser, false)
mustIAM(t, err)
c, _, err := sys.NewServiceAccount(ctx, "healthy-user", nil, newServiceAccountOpts{accessKey: "expired-service", secretKey: "expired-service-password"})
mustIAM(t, err)
c.Expiration = UTCNow().Add(-time.Hour)
path := getUserIdentityPath(c.AccessKey, svcUser)
mustIAM(t, sys.store.saveIAMConfig(ctx, &UserIdentity{Version: 1, Credentials: c, UpdatedAt: UTCNow()}, path))
// A cold loader sees the existing version but cannot acquire the cleanup
// write lock. Healthy users must still load; the expired one stays denied.
fresh := &IAMStoreSys{IAMStorageAPI: newIAMObjectStore(&iamExpiryLockFailure{ObjectLayer: obj, path: path}, MinIOUsersSysType)}
mustIAM(t, fresh.LoadIAMCache(ctx, true))
if _, ok := fresh.GetUser("healthy-user"); !ok {
t.Fatal("cleanup failure prevented healthy IAM state from loading")
}
if _, ok := fresh.GetUser(c.AccessKey); ok {
t.Fatal("cleanup failure admitted an expired service account")
}
r, err := loadIAMRevision(ctx, fresh, path)
mustIAM(t, err)
if r.Deleted || !r.Credentials.IsExpired() {
t.Fatal("failed cleanup lost the existing expired revision")
}
}
+220
View File
@@ -0,0 +1,220 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"encoding/json"
"fmt"
"maps"
"strings"
"sync"
"time"
"github.com/minio/minio/internal/auth"
)
// This index is rebuilt by the existing IAM loaders and updated by successful
// storage operations. It avoids a second full IAM walk during every heal pass.
// It is an optimization of the durable records, never a reason to delete them.
// The index contains no secrets or grants.
type iamParentRevision struct {
deleted bool
before time.Time
}
type iamRevisionIndex struct {
mu sync.RWMutex
items map[string]iamRevision
parents map[string]iamParentRevision
floors map[string]time.Time
generation uint64
}
func (idx *iamRevisionIndex) observe(path string, data []byte) {
if !strings.HasPrefix(path, iamConfigPrefix+"/") {
return
}
var r iamRevision
if json.Unmarshal(data, &r) != nil {
return // The caller reports malformed data using its normal decoder.
}
r.Credentials = auth.Credentials{ParentUser: r.Credentials.ParentUser, Expiration: r.Credentials.Expiration}
idx.mu.Lock()
defer idx.mu.Unlock()
if strings.HasPrefix(path, iamConfigUsersPrefix) {
// Keep a compact name-keyed view for the authentication hot path;
// constructing a config path on every S3 request allocates needlessly.
defer func() {
name := strings.TrimSuffix(strings.TrimPrefix(path, iamConfigUsersPrefix), "/"+iamIdentityFile)
if current, ok := idx.items[path]; ok {
if idx.parents == nil {
idx.parents = make(map[string]iamParentRevision)
}
idx.parents[name] = iamParentRevision{deleted: current.Deleted, before: current.RevokedBefore}
} else {
delete(idx.parents, name)
}
}()
}
if floor, ok := idx.floors[path]; ok && r.timestamp().Before(floor) {
return
}
if previous, ok := idx.items[path]; ok {
// A concurrent read that began before a write must not roll it back.
if previous.timestamp().After(r.timestamp()) || (previous.Deleted && !r.Deleted && !r.timestamp().After(previous.timestamp())) {
return
}
if previous.RevokedBefore.After(r.RevokedBefore) {
r.RevokedBefore = previous.RevokedBefore
}
if previous.timestamp().Equal(r.timestamp()) && previous.Deleted == r.Deleted && previous.RevokedBefore.Equal(r.RevokedBefore) {
return
}
}
if r.Deleted && !r.ExpiresAt.IsZero() && UTCNow().After(r.ExpiresAt) {
if _, tracked := idx.items[path]; tracked {
delete(idx.items, path)
idx.generation++
}
delete(idx.floors, path)
return
}
if !r.Deleted && r.RevokedBefore.IsZero() {
_, tracked := idx.items[path]
_, hasFloor := idx.floors[path]
if tracked || hasFloor {
if idx.floors == nil {
idx.floors = make(map[string]time.Time)
}
idx.floors[path] = r.timestamp()
}
if tracked {
delete(idx.items, path)
idx.generation++
}
return
}
if idx.items == nil {
idx.items = make(map[string]iamRevision)
}
idx.items[path] = r
delete(idx.floors, path)
idx.generation++
}
func (idx *iamRevisionIndex) get(path string) iamRevision {
if idx == nil {
return iamRevision{}
}
idx.mu.RLock()
defer idx.mu.RUnlock()
return idx.items[path]
}
func (idx *iamRevisionIndex) snapshot() map[string]iamRevision {
idx.mu.Lock()
defer idx.mu.Unlock()
for path, r := range idx.items {
if r.Deleted && !r.ExpiresAt.IsZero() && UTCNow().After(r.ExpiresAt) {
delete(idx.items, path)
delete(idx.floors, path)
idx.generation++
}
}
return maps.Clone(idx.items)
}
func (idx *iamRevisionIndex) count() int {
idx.mu.RLock()
defer idx.mu.RUnlock()
return len(idx.items)
}
// A process-local generation plus the protocol's instance ID is sufficient
// for acknowledgements. Avoid hashing the entire index on every IAM write.
func (idx *iamRevisionIndex) digest() string {
idx.mu.RLock()
defer idx.mu.RUnlock()
return fmt.Sprintf("%x:%x", idx.generation, len(idx.items))
}
func (idx *iamRevisionIndex) forget(path string) {
idx.mu.Lock()
if _, ok := idx.items[path]; ok {
delete(idx.items, path)
idx.generation++
}
delete(idx.floors, path)
if strings.HasPrefix(path, iamConfigUsersPrefix) {
delete(idx.parents, strings.TrimSuffix(strings.TrimPrefix(path, iamConfigUsersPrefix), "/"+iamIdentityFile))
}
idx.mu.Unlock()
}
func (c *iamCache) userRevocation(user string) iamRevision {
r := c.revisions.parentRevision(user)
if u, ok := c.iamUsersMap[user]; ok && u.RevokedBefore.After(r.RevokedBefore) {
r.RevokedBefore = u.RevokedBefore
}
return r
}
func (c *iamCache) groupMemberAllowed(member string, grantedAt, groupBoundary time.Time) bool {
r := c.userRevocation(member)
return !r.Deleted && (r.RevokedBefore.IsZero() || grantedAt.After(r.RevokedBefore)) && (groupBoundary.IsZero() || grantedAt.After(groupBoundary))
}
func iamMappingParentPath(path string) string {
kind, name, ok := strings.Cut(strings.TrimPrefix(path, iamConfigPolicyDBPrefix), "/")
if !ok {
return ""
}
name = strings.TrimSuffix(name, ".json")
switch kind {
case "users", "sts-users":
return getUserIdentityPath(name, regUser)
case "service-accounts":
return getUserIdentityPath(name, svcUser)
case "groups":
return getGroupInfoPath(name)
}
return ""
}
func (idx *iamRevisionIndex) mappingAllowed(path string, mp MappedPolicy) bool {
if mp.Deleted || idx.get(path).Deleted {
return false
}
r := idx.get(iamMappingParentPath(path))
return !r.Deleted && (r.RevokedBefore.IsZero() || mp.UpdatedAt.After(r.RevokedBefore))
}
// Apply the persisted commit boundary even before dependent cache cleanup has
// completed. The map namespace is part of the authorization record's identity.
func (c *iamCache) cachedMappedPolicy(name string, userType IAMUserType, isGroup bool) (MappedPolicy, bool) {
var mp MappedPolicy
var ok bool
switch {
case isGroup:
mp, ok = c.iamGroupPolicyMap.Load(name)
case userType == stsUser:
mp, ok = c.iamSTSPolicyMap.Load(name)
default:
mp, ok = c.iamUserPolicyMap.Load(name)
}
if !ok || !c.revisions.mappingAllowed(getMappedPolicyPath(name, userType, isGroup), mp) {
return MappedPolicy{}, false
}
return mp, true
}
func (idx *iamRevisionIndex) parentRevision(user string) iamRevision {
if idx == nil {
return iamRevision{}
}
idx.mu.RLock()
p := idx.parents[user]
idx.mu.RUnlock()
return iamRevision{Deleted: p.deleted, RevokedBefore: p.before}
}
+305
View File
@@ -0,0 +1,305 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"crypto/sha256"
"fmt"
"os"
"strings"
"sync"
"testing"
"time"
"github.com/minio/madmin-go/v3"
etcd "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/concurrency"
"go.etcd.io/etcd/client/v3/namespace"
)
type iamRevisionLockObserver struct {
ObjectLayer
path string
waiting chan struct{}
once sync.Once
}
func (o *iamRevisionLockObserver) NewNSLock(bucket string, objects ...string) RWLocker {
lock := o.ObjectLayer.NewNSLock(bucket, objects...)
if bucket == minioMetaBucket && len(objects) == 1 && objects[0] == o.path {
return &iamRevisionObservedLock{RWLocker: lock, observe: func() { o.once.Do(func() { close(o.waiting) }) }}
}
return lock
}
type iamRevisionObservedLock struct {
RWLocker
observe func()
}
func (l *iamRevisionObservedLock) GetLock(ctx context.Context, timeout *dynamicTimeout) (LockContext, error) {
l.observe()
return l.RWLocker.GetLock(ctx, timeout)
}
type iamRevisionWatchObserver struct {
etcd.Watcher
waiting chan struct{}
once sync.Once
}
func (w *iamRevisionWatchObserver) Watch(ctx context.Context, key string, opts ...etcd.OpOption) etcd.WatchChan {
w.once.Do(func() { close(w.waiting) })
return w.Watcher.Watch(ctx, key, opts...)
}
// Simulate an unavailable cleanup RPC. Mutex.Lock calls Delete after its wait
// is canceled; that RPC must inherit a deadline too, not Client.Ctx() forever.
type iamRevisionCleanupBlocker struct {
etcd.KV
release chan struct{}
}
func (b *iamRevisionCleanupBlocker) Delete(ctx context.Context, key string, opts ...etcd.OpOption) (*etcd.DeleteResponse, error) {
if strings.Contains(key, "/iam-revision-locks/") {
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-b.release:
}
}
return b.KV.Delete(ctx, key, opts...)
}
type iamRevisionReadBlocker struct {
IAMStorageAPI
path string
after int
waiting chan struct{}
}
func (b *iamRevisionReadBlocker) loadIAMConfig(ctx context.Context, item any, path string) error {
if path == b.path {
b.after--
if b.after == 0 {
close(b.waiting)
<-ctx.Done()
return ctx.Err()
}
}
return b.IAMStorageAPI.loadIAMConfig(ctx, item, path)
}
func TestIAMRevisionReadDoesNotBlockAuthentication(t *testing.T) {
for _, stage := range []struct {
name string
offset time.Duration
}{{"deletion", time.Minute}, {"retained_revocation", -time.Minute}} {
t.Run(stage.name, func(t *testing.T) {
resetTestGlobals()
t.Cleanup(resetTestGlobals)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
disks, err := getRandomDisks(1)
if err != nil {
t.Fatal(err)
}
obj, _, err := initObjectLayer(ctx, mustGetPoolEndpoints(0, disks...))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
obj.Shutdown(context.Background())
os.RemoveAll(disks[0])
})
store := &IAMStoreSys{IAMStorageAPI: newIAMObjectStore(obj, MinIOUsersSysType)}
const user = "read-blocked-parent"
created, err := store.AddUser(ctx, user, madmin.AddOrUpdateUserReq{SecretKey: "original-password", Status: madmin.AccountEnabled})
if err != nil {
t.Fatal(err)
}
blocked := &iamRevisionReadBlocker{IAMStorageAPI: store.IAMStorageAPI, path: getUserIdentityPath(user, regUser), after: 1, waiting: make(chan struct{})}
store.IAMStorageAPI = blocked
done := make(chan error, 1)
go func() {
done <- store.DeleteUser(withIAMReplicationTime(ctx, created.Add(stage.offset)), user, regUser)
}()
defer func() { cancel(); <-done }()
select {
case <-blocked.waiting:
case <-time.After(5 * time.Second):
t.Fatal("revision read was not attempted")
}
read := make(chan bool, 1)
go func() {
u, ok := store.GetUser(user)
read <- ok && u.Credentials.SecretKey == "original-password"
}()
select {
case ok := <-read:
if !ok {
t.Fatal("pending revision read changed the cached identity")
}
case <-time.After(time.Second):
t.Fatal("revision read blocked cached authentication")
}
})
}
}
func TestIAMRevisionLockContention(t *testing.T) {
for _, backend := range []string{"object", "etcd"} {
t.Run(backend, func(t *testing.T) {
endpoint := os.Getenv("SILO_TEST_IAM_REVOCATION_ETCD")
if backend == "etcd" && endpoint == "" {
t.Skip("set SILO_TEST_IAM_REVOCATION_ETCD to a disposable etcd endpoint")
}
for _, outcome := range []string{"release", "cancel", "default_timeout"} {
t.Run(outcome, func(t *testing.T) {
resetTestGlobals()
t.Cleanup(resetTestGlobals)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
oldTimeout := defaultContextTimeout
defaultContextTimeout = 2 * time.Second
t.Cleanup(func() { defaultContextTimeout = oldTimeout })
must := func(err error) {
t.Helper()
if err != nil {
t.Fatal(err)
}
}
const user = "contended-user"
path := getUserIdentityPath(user, regUser)
waiting := make(chan struct{})
var store *IAMStoreSys
var hold func() func()
unblockCleanup := func() {}
if backend == "object" {
disks, err := getRandomDisks(1)
must(err)
obj, _, err := initObjectLayer(ctx, mustGetPoolEndpoints(0, disks...))
must(err)
t.Cleanup(func() {
obj.Shutdown(context.Background())
os.RemoveAll(disks[0])
})
observed := &iamRevisionLockObserver{ObjectLayer: obj, path: path + ".revision-lock", waiting: waiting}
store = &IAMStoreSys{IAMStorageAPI: newIAMObjectStore(obj, MinIOUsersSysType)}
hold = func() func() {
lock := obj.NewNSLock(minioMetaBucket, observed.path)
lc, err := lock.GetLock(ctx, newDynamicTimeout(time.Second, time.Second))
must(err)
store.IAMStorageAPI.(*IAMObjectStore).objAPI = observed
return func() { lock.Unlock(lc) }
}
} else {
client, err := etcd.New(etcd.Config{Endpoints: strings.Split(endpoint, ","), DialTimeout: time.Second})
must(err)
t.Cleanup(func() { client.Close() })
prefix := fmt.Sprintf("/silo-lock-test/%d/", time.Now().UnixNano())
client.KV = namespace.NewKV(client.KV, prefix)
client.Watcher = namespace.NewWatcher(client.Watcher, prefix)
store = &IAMStoreSys{IAMStorageAPI: newIAMEtcdStore(client, MinIOUsersSysType)}
hold = func() func() {
session, err := concurrency.NewSession(client, concurrency.WithContext(ctx))
must(err)
lock := concurrency.NewMutex(session, fmt.Sprintf("%s/iam-revision-locks/%x", minioConfigPrefix, sha256.Sum256([]byte(path))))
must(lock.Lock(ctx))
client.Watcher = &iamRevisionWatchObserver{Watcher: client.Watcher, waiting: waiting}
blocker := &iamRevisionCleanupBlocker{KV: client.KV, release: make(chan struct{})}
client.KV = blocker
unblockCleanup = sync.OnceFunc(func() { close(blocker.release) })
t.Cleanup(unblockCleanup)
return func() { session.Close() }
}
}
request := func(secret string) madmin.AddOrUpdateUserReq {
return madmin.AddOrUpdateUserReq{SecretKey: secret, Status: madmin.AccountEnabled}
}
_, err := store.AddUser(ctx, user, request("original-password"))
must(err)
release := sync.OnceFunc(hold())
t.Cleanup(release)
writeCtx, cancelWrite := context.WithCancel(ctx)
defer cancelWrite()
first, second := make(chan error, 1), make(chan error, 1)
var writers sync.WaitGroup
t.Cleanup(func() {
cancelWrite()
unblockCleanup()
release()
writers.Wait()
})
writers.Go(func() {
_, err := store.AddUser(writeCtx, user, request("first-password"))
first <- err
})
select {
case <-waiting:
case <-time.After(5 * time.Second):
t.Fatal("writer did not attempt the held revision lock")
}
// A second writer must queue without taking the cache's RWMutex:
// Go's writer preference would otherwise block every new reader.
writers.Go(func() {
_, err := store.AddUser(ctx, user, request("second-password"))
second <- err
})
select {
case err := <-second:
t.Fatalf("second writer bypassed the first: %v", err)
case <-time.After(50 * time.Millisecond):
}
read := make(chan UserIdentity, 1)
go func() {
u, _ := store.GetUser(user)
read <- u
}()
select {
case u := <-read:
if u.Credentials.SecretKey != "original-password" {
t.Fatal("pending write changed the cached credential")
}
case <-time.After(time.Second):
t.Fatal("distributed lock contention blocked cached authentication")
}
switch outcome {
case "release":
release()
case "cancel":
cancelWrite()
}
select {
case err := <-first:
if outcome == "release" {
must(err)
} else if err == nil {
t.Fatal("canceled or timed-out write succeeded")
}
case <-time.After(5 * time.Second):
t.Fatal("lock wait or cancellation cleanup exceeded its deadline")
}
release()
select {
case err := <-second:
must(err)
case <-time.After(5 * time.Second):
t.Fatal("queued writer did not recover after the first completed")
}
cached, ok := store.GetUser(user)
if !ok || cached.Credentials.SecretKey != "second-password" {
t.Fatal("cached write order was lost")
}
var persisted UserIdentity
must(store.loadIAMConfig(ctx, &persisted, path))
if persisted.Credentials.SecretKey != cached.Credentials.SecretKey || !persisted.UpdatedAt.Equal(cached.UpdatedAt) {
t.Fatal("persistent and cached revisions differ")
}
})
}
})
}
}
+749
View File
@@ -0,0 +1,749 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"crypto/sha256"
"errors"
"fmt"
"net/http"
"sort"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
etcd "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/concurrency"
)
var errIAMStaleUpdate = errors.New("IAM update predates a stored revision or revocation")
// The parent is still live; callers must not broadcast a user deletion when
// only its revocation boundary was retained.
var errIAMRevocationRetained = errors.New("IAM revocation recorded without deleting the record")
// A revocation advances the boundary even when a newer identity already
// exists. Keep this operation distinct from replacing/deleting that identity.
type iamUserRevocation struct {
UserIdentity
retained bool
}
type iamGroupRevocation struct {
GroupInfo
retained bool
requireEmpty bool
}
// Natural expiration is distinct from revoking a live credential. An expired
// immutable STS token can be removed; a reusable service-account key retains
// its revision so an older non-expiring credential cannot return.
type iamExpireIdentity struct{}
// The authoritative revocation is durable even if dependent cleanup fails.
// Callers must publish it to sibling caches before returning the error.
type iamCommittedCleanupError struct {
err error
retained bool
}
func (e *iamCommittedCleanupError) Error() string {
return "IAM revocation committed; cleanup failed: " + e.err.Error()
}
func (e *iamCommittedCleanupError) Unwrap() error { return e.err }
type iamReplicationTimeKey struct{}
func withIAMReplicationTime(ctx context.Context, at time.Time) context.Context {
return context.WithValue(ctx, iamReplicationTimeKey{}, at)
}
func iamReplicationTime(ctx context.Context) (time.Time, bool) {
at, ok := ctx.Value(iamReplicationTimeKey{}).(time.Time)
return at, ok
}
func iamReplicationError(err error) error {
if errors.Is(err, errIAMStaleUpdate) {
// Retrying an obsolete event cannot change the result.
return nil
}
return wrapSRErr(err)
}
// Deletions occupy the original IAM config path. They contain no secret or
// grant and are hidden by the normal loaders, but remain available to heal
// and to timestamp comparisons after a restart. Do not age them out: a peer
// can be offline indefinitely.
type iamRevision struct {
UpdatedAt time.Time `json:"updatedAt"`
UpdateDate time.Time `json:"UpdateDate"`
Deleted bool `json:"deleted"`
RevokedBefore time.Time `json:"revokedBefore"`
ExpiresAt time.Time `json:"expiresAt,omitempty"`
Credentials auth.Credentials `json:"credentials"`
}
func (r iamRevision) timestamp() time.Time {
if r.UpdateDate.After(r.UpdatedAt) {
return r.UpdateDate
}
return r.UpdatedAt
}
func loadIAMRevision(ctx context.Context, store IAMStorageAPI, path string) (iamRevision, error) {
var r iamRevision
err := store.loadIAMConfig(ctx, &r, path)
if errors.Is(err, errConfigNotFound) {
err = nil
}
return r, err
}
func (store *IAMStoreSys) checkIAMRevision(ctx context.Context, path string, deleting bool) error {
at, replicated := iamReplicationTime(ctx)
if !replicated {
return nil
}
return store.withIAMStorage(ctx, func(ctx context.Context) error {
r, err := loadIAMRevision(ctx, store.IAMStorageAPI, path)
if err != nil {
return err
}
if r.timestamp().After(at) || (r.Deleted && !deleting && !at.After(r.timestamp())) {
return errIAMStaleUpdate
}
return nil
})
}
// This signed claim records the parent's revocation boundary at issuance.
// Unlike UpdatedAt, it cannot advance when an offline site edits an old child.
// It travels in the existing service-account Claims and STS SessionToken fields.
const iamParentRevocationClaim = "siloParentRevocation"
func setIAMParentRevocationClaim(ctx context.Context, store IAMStorageAPI, parent string, claims map[string]any) error {
delete(claims, iamParentRevocationClaim)
if parent == "" || parent == globalActiveCred.AccessKey {
return nil
}
r, err := loadIAMRevision(ctx, store, getUserIdentityPath(parent, regUser))
if err != nil {
return err
}
if r.Deleted {
return errIAMStaleUpdate
}
if !r.RevokedBefore.IsZero() {
claims[iamParentRevocationClaim] = r.RevokedBefore.Format(time.RFC3339Nano)
}
return nil
}
func iamCredentialSurvivesRevocation(cred auth.Credentials, at time.Time) bool {
if at.IsZero() {
return true
}
s, _ := cred.Claims[iamParentRevocationClaim].(string)
issuedAfter, err := time.Parse(time.RFC3339Nano, s)
return err == nil && !issuedAfter.Before(at)
}
// Parent revocations delete old children even if an offline peer has edited
// them later. Preserve children that prove issuance after this revocation.
func iamChildDeletionContext(ctx context.Context, child UserIdentity) (context.Context, bool) {
if at, replicated := iamReplicationTime(ctx); replicated {
if !at.IsZero() && iamCredentialSurvivesRevocation(child.Credentials, at) {
return ctx, false
}
if child.UpdatedAt.After(at) {
ctx = withIAMReplicationTime(ctx, child.UpdatedAt)
}
}
return ctx, true
}
// A delayed service account or STS event must not outlive deletion of its
// built-in parent. The caller must populate Claims from the verified token.
func checkIAMParentRevision(ctx context.Context, store IAMStorageAPI, cred auth.Credentials) error {
parent := cred.ParentUser
if parent == "" || parent == globalActiveCred.AccessKey {
return nil
}
r, err := loadIAMRevision(ctx, store, getUserIdentityPath(parent, regUser))
if err != nil {
return err
}
if r.Deleted || !iamCredentialSurvivesRevocation(cred, r.RevokedBefore) {
return errIAMStaleUpdate
}
return nil
}
// Called with the IAM writer mutex and cache lock held. Persistence only
// touches the caller's record, not the cache. Keep writers serialized while
// allowing cached authentication reads throughout storage and lock waits.
func (store *IAMStoreSys) withIAMStorage(ctx context.Context, fn func(context.Context) error) error {
store.IAMStorageAPI.unlock()
defer store.IAMStorageAPI.lock()
ctx, cancel := context.WithTimeout(ctx, defaultContextTimeout)
defer cancel()
return fn(ctx)
}
func (store *IAMStoreSys) saveIAMRevision(ctx context.Context, path string, item any, opts ...options) error {
return store.withIAMStorage(ctx, func(ctx context.Context) error {
return saveIAMRevision(ctx, store.IAMStorageAPI, path, item, opts...)
})
}
func (store *IAMStoreSys) checkIAMParentRevision(ctx context.Context, cred auth.Credentials) error {
return store.withIAMStorage(ctx, func(ctx context.Context) error {
return checkIAMParentRevision(ctx, store.IAMStorageAPI, cred)
})
}
// Update the caller's record with the persisted revision before it is cached.
func saveIAMRevision(ctx context.Context, store IAMStorageAPI, path string, item any, opts ...options) error {
ctx, cancel := context.WithTimeout(ctx, defaultContextTimeout)
defer cancel()
// Serialize compare-and-write across nodes, as well as goroutines. Use a
// separate lock name so saving the config does not reacquire this lock.
switch s := store.(type) {
case *IAMObjectStore:
lock := s.objAPI.NewNSLock(minioMetaBucket, path+".revision-lock")
lc, err := lock.GetLock(ctx, globalOperationTimeout)
if err != nil {
return err
}
defer lock.Unlock(lc)
ctx = lc.Context()
case *IAMEtcdStore:
// Mutex.Lock also uses Client.Ctx() for cleanup after cancellation.
// Borrow the existing services with the operation's bounded context;
// never close this facade, which does not own those services.
client := etcd.NewCtxClient(ctx, etcd.WithZapLogger(s.client.GetLogger()))
client.KV, client.Lease, client.Watcher = s.client.KV, s.client.Lease, s.client.Watcher
session, err := concurrency.NewSession(client, concurrency.WithContext(ctx))
if err != nil {
return err
}
defer func() {
session.Orphan()
// A canceled operation must still release its lease when etcd is
// reachable. If it is unavailable, stop waiting and let it expire.
cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), defaultContextTimeout)
defer cancel()
_, _ = s.client.Revoke(cleanupCtx, session.Lease())
}()
lock := concurrency.NewMutex(session, fmt.Sprintf("%s/iam-revision-locks/%x", minioConfigPrefix, sha256.Sum256([]byte(path))))
if err = lock.Lock(ctx); err != nil {
return err
}
// Revoking the session lease releases the lock, including on cancellation.
}
previous, err := loadIAMRevision(ctx, store, path)
if err != nil {
return err
}
if _, expiring := item.(*iamExpireIdentity); expiring {
sts := strings.HasPrefix(path, iamConfigSTSPrefix)
if previous.Deleted {
if sts && !previous.ExpiresAt.IsZero() && UTCNow().After(previous.ExpiresAt) {
return expireIAMSTSConfig(ctx, store, path)
}
return nil
}
if previous.timestamp().IsZero() || !previous.Credentials.IsExpired() {
return nil
}
if sts {
return expireIAMSTSConfig(ctx, store, path)
}
item = &UserIdentity{Version: 1, Deleted: true}
ctx = withIAMReplicationTime(ctx, previous.timestamp())
}
var revocation *iamUserRevocation
if op, ok := item.(*iamUserRevocation); ok {
revocation = op
op.UserIdentity = UserIdentity{Version: 1, Deleted: true}
if origin, replicated := iamReplicationTime(ctx); replicated && previous.timestamp().After(origin) {
if previous.Deleted || !origin.After(previous.RevokedBefore) {
return errIAMStaleUpdate
}
op.retained = true
op.UserIdentity = UserIdentity{Version: 1, Credentials: previous.Credentials, UpdatedAt: previous.timestamp(), RevokedBefore: origin}
ctx = withIAMReplicationTime(ctx, previous.timestamp())
}
item = &op.UserIdentity
}
var groupRevocation *iamGroupRevocation
if op, ok := item.(*iamGroupRevocation); ok {
groupRevocation = op
var group GroupInfo
if err := store.loadIAMConfig(ctx, &group, path); err != nil && !errors.Is(err, errConfigNotFound) {
return err
}
if op.requireEmpty && !group.Deleted {
for _, member := range group.Members {
r := store.revisionIndex().get(getUserIdentityPath(member, regUser))
at := group.MemberGrants[member]
if !r.Deleted && (r.RevokedBefore.IsZero() || at.After(r.RevokedBefore)) && (group.RevokedBefore.IsZero() || at.After(group.RevokedBefore)) {
return errGroupNotEmpty
}
}
}
op.GroupInfo = GroupInfo{Version: 1, Deleted: true}
if origin, replicated := iamReplicationTime(ctx); replicated && previous.timestamp().After(origin) {
if previous.Deleted || !origin.After(previous.RevokedBefore) {
return errIAMStaleUpdate
}
op.retained = true
op.GroupInfo = group
op.RevokedBefore = origin
ctx = withIAMReplicationTime(ctx, previous.timestamp())
}
item = &op.GroupInfo
}
var at *time.Time
var deleted bool
switch v := item.(type) {
case *UserIdentity:
at, deleted = &v.UpdatedAt, v.Deleted
if boundary, ok := ctx.Value(iamRecordBoundaryKey{}).(time.Time); ok && boundary.After(v.RevokedBefore) {
v.RevokedBefore = boundary
}
if previous.RevokedBefore.After(v.RevokedBefore) {
v.RevokedBefore = previous.RevokedBefore
}
case *GroupInfo:
at, deleted = &v.UpdatedAt, v.Deleted
if boundary, ok := ctx.Value(iamRecordBoundaryKey{}).(time.Time); ok && boundary.After(v.RevokedBefore) {
v.RevokedBefore = boundary
}
if !deleted {
var group GroupInfo
if err := store.loadIAMConfig(ctx, &group, path); err != nil && !errors.Is(err, errConfigNotFound) {
return err
}
mergeIAMGroupMutation(ctx, group, v)
}
if previous.RevokedBefore.After(v.RevokedBefore) {
v.RevokedBefore = previous.RevokedBefore
}
case *MappedPolicy:
at, deleted = &v.UpdatedAt, v.Deleted
case *PolicyDoc:
at, deleted = &v.UpdateDate, v.Deleted
default:
return errInvalidArgument
}
if strings.HasPrefix(path, iamConfigSTSPrefix) && previous.Deleted && !deleted {
// STS access keys identify immutable tokens, not reusable user names.
return errIAMStaleUpdate
}
if origin, replicated := iamReplicationTime(ctx); replicated {
*at = origin
if previous.timestamp().After(origin) || (previous.Deleted && !deleted && !origin.After(previous.timestamp())) {
return errIAMStaleUpdate
}
if strings.HasPrefix(path, iamConfigServiceAccountsPrefix) && !deleted && previous.Credentials.AccessKey != "" && previous.timestamp().Equal(origin) {
// Duplicate service snapshots are acknowledgements, not new creates
// or edits. Reload the winner without writing, so even a stale
// sibling cache is refreshed by the retry before acknowledging it.
return store.loadIAMConfig(ctx, item, path)
}
if previous.Deleted && deleted && !origin.After(previous.timestamp()) {
// An already-applied tombstone needs no further persistent write.
if v, ok := item.(*UserIdentity); ok {
v.RevokedBefore = previous.RevokedBefore
}
if v, ok := item.(*GroupInfo); ok {
v.RevokedBefore = previous.RevokedBefore
}
return nil
}
} else {
if previous.Deleted && deleted {
// A peer notification without an originating revision must not
// advance a tombstone past a subsequent deliberate recreation.
*at = previous.timestamp()
if v, ok := item.(*UserIdentity); ok {
v.RevokedBefore = previous.RevokedBefore
}
if v, ok := item.(*GroupInfo); ok {
v.RevokedBefore = previous.RevokedBefore
}
return nil
}
if at.IsZero() {
*at = UTCNow()
}
if !at.After(previous.timestamp()) {
*at = previous.timestamp().Add(time.Nanosecond)
}
}
if v, ok := item.(*UserIdentity); ok {
if deleted {
// Retain only the parent name for root-account exclusion during heal.
v.Credentials = auth.Credentials{ParentUser: previous.Credentials.ParentUser}
v.RevokedBefore = *at
if strings.HasPrefix(path, iamConfigSTSPrefix) && !previous.Credentials.Expiration.IsZero() && !previous.Credentials.Expiration.Equal(timeSentinel) {
// The signed STS token cannot authorize beyond this time, even
// if an offline site replays it with a newer event timestamp.
v.ExpiresAt = previous.Credentials.Expiration.Add(globalMaxSkewTime)
opts = []options{{ttl: max(1, int64(time.Until(v.ExpiresAt).Seconds())+1)}}
}
} else {
if v.Credentials.SessionToken != "" && v.Credentials.Claims == nil {
claims, err := extractJWTClaims(*v)
if err != nil {
return err
}
v.Credentials.Claims = claims.Map()
}
if err = checkIAMParentRevision(ctx, store, v.Credentials); err != nil {
return err
}
}
}
if v, ok := item.(*GroupInfo); ok && deleted {
v.RevokedBefore = *at
v.Members, v.MemberGrants = nil, nil
}
if _, ok := item.(*MappedPolicy); ok && !deleted {
if parentPath := iamMappingParentPath(path); parentPath != "" {
parent, err := loadIAMRevision(ctx, store, parentPath)
if err != nil {
return err
}
if parent.Deleted {
return errIAMStaleUpdate
}
if !parent.RevokedBefore.IsZero() && !at.After(parent.RevokedBefore) {
if _, replicated := iamReplicationTime(ctx); replicated {
return errIAMStaleUpdate
}
*at = parent.RevokedBefore.Add(time.Nanosecond)
}
}
}
if err := store.saveIAMConfig(ctx, item, path, opts...); err != nil {
return err
}
if revocation != nil && revocation.retained {
return errIAMRevocationRetained
}
if groupRevocation != nil && groupRevocation.retained {
return errIAMRevocationRetained
}
return nil
}
func (iamOS *IAMObjectStore) listIAMConfigPaths(ctx context.Context) ([]string, error) {
ctx, cancel := context.WithCancel(ctx)
defer cancel()
var paths []string
for item := range listIAMConfigItems(ctx, iamOS.objAPI, iamConfigPrefix+"/") {
if item.Err != nil {
return nil, item.Err
}
paths = append(paths, iamConfigPrefix+"/"+item.Item)
}
return paths, nil
}
func (ies *IAMEtcdStore) listIAMConfigPaths(ctx context.Context) ([]string, error) {
ctx, cancel := context.WithTimeout(ctx, defaultContextTimeout)
defer cancel()
r, err := ies.client.Get(ctx, iamConfigPrefix+"/", etcd.WithPrefix(), etcd.WithKeysOnly())
if err != nil {
return nil, err
}
paths := make([]string, 0, len(r.Kvs))
for _, kv := range r.Kvs {
paths = append(paths, string(kv.Key))
}
return paths, nil
}
func iamDeletionItem(path string, r iamRevision) (item madmin.SRIAMItem, ok bool) {
if (strings.HasPrefix(path, iamConfigUsersPrefix) || strings.HasPrefix(path, iamConfigGroupsPrefix)) && !r.RevokedBefore.IsZero() {
// Recreating a parent does not cancel its older revocation of derived
// credentials. Replay this boundary even after the parent is live again.
r.Deleted = true
r.UpdatedAt, r.UpdateDate = r.RevokedBefore, time.Time{}
}
if !r.Deleted {
return item, false
}
item.UpdatedAt = r.timestamp()
switch {
case strings.HasPrefix(path, iamConfigUsersPrefix):
name := strings.TrimSuffix(strings.TrimPrefix(path, iamConfigUsersPrefix), "/"+iamIdentityFile)
item.Type = madmin.SRIAMItemIAMUser
item.IAMUser = &madmin.SRIAMUser{AccessKey: name, IsDeleteReq: true}
case strings.HasPrefix(path, iamConfigServiceAccountsPrefix):
name := strings.TrimSuffix(strings.TrimPrefix(path, iamConfigServiceAccountsPrefix), "/"+iamIdentityFile)
if name == siteReplicatorSvcAcc || r.Credentials.ParentUser == globalActiveCred.AccessKey {
return item, false
}
item.Type = madmin.SRIAMItemSvcAcc
item.SvcAccChange = &madmin.SRSvcAccChange{Delete: &madmin.SRSvcAccDelete{AccessKey: name}}
case strings.HasPrefix(path, iamConfigGroupsPrefix):
name := strings.TrimSuffix(strings.TrimPrefix(path, iamConfigGroupsPrefix), "/"+iamGroupMembersFile)
item.Type = madmin.SRIAMItemGroupInfo
item.GroupInfo = &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: name, IsRemove: true}}
case strings.HasPrefix(path, iamConfigPoliciesPrefix):
item.Type = madmin.SRIAMItemPolicy
item.Name = strings.TrimSuffix(strings.TrimPrefix(path, iamConfigPoliciesPrefix), "/"+iamPolicyFile)
case strings.HasPrefix(path, iamConfigPolicyDBPrefix):
prefix, name, found := strings.Cut(strings.TrimPrefix(path, iamConfigPolicyDBPrefix), "/")
if !found {
return item, false
}
typ := regUser
switch prefix {
case "sts-users":
typ = stsUser
case "service-accounts":
typ = svcUser
}
item.Type = madmin.SRIAMItemPolicyMapping
item.PolicyMapping = &madmin.SRPolicyMapping{UserOrGroup: strings.TrimSuffix(name, ".json"), UserType: int(typ), IsGroup: prefix == "groups"}
default:
// Expired STS credentials are not replayed. Parent revocations and
// their retained timestamp reject delayed copies of derived tokens.
return item, false
}
return item, true
}
func iamDeletionPath(item madmin.SRIAMItem) string {
switch item.Type {
case madmin.SRIAMItemIAMUser:
if item.IAMUser != nil && item.IAMUser.IsDeleteReq {
return getUserIdentityPath(item.IAMUser.AccessKey, regUser)
}
case madmin.SRIAMItemSvcAcc:
if item.SvcAccChange != nil && item.SvcAccChange.Delete != nil {
return getUserIdentityPath(item.SvcAccChange.Delete.AccessKey, svcUser)
}
case madmin.SRIAMItemGroupInfo:
if item.GroupInfo != nil && item.GroupInfo.UpdateReq.IsRemove && len(item.GroupInfo.UpdateReq.Members) == 0 {
return getGroupInfoPath(item.GroupInfo.UpdateReq.Group)
}
case madmin.SRIAMItemPolicy:
if len(item.Policy) == 0 {
return getPolicyDocPath(item.Name)
}
case madmin.SRIAMItemPolicyMapping:
if p := item.PolicyMapping; p != nil && p.Policy == "" {
return getMappedPolicyPath(p.UserOrGroup, IAMUserType(p.UserType), p.IsGroup)
}
}
return ""
}
func (c *SiteReplicationSys) healIAMDeletions(ctx context.Context) (err error) {
started := time.Now()
defer func() {
c.iamRevisionMetrics.healDurationMillis.Store(time.Since(started).Milliseconds())
if err != nil {
c.iamRevisionMetrics.healFailures.Add(1)
} else {
c.iamRevisionMetrics.healLastSuccess.Store(time.Now().Unix())
}
}()
c.iamHealMu.Lock()
defer c.iamHealMu.Unlock()
c.RLock()
defer c.RUnlock()
if !c.enabled {
return nil
}
snapshot := globalIAMSys.store.revisionIndex().snapshot()
paths := make([]string, 0, len(snapshot))
for path := range snapshot {
paths = append(paths, path)
}
sort.Strings(paths)
byType := make(map[string][]iamReplicationItem)
for _, path := range paths {
r := snapshot[path]
item, ok := iamDeletionItem(path, r)
if !ok {
continue
}
out := iamReplicationItem{SRIAMItem: item}
if item.Type == madmin.SRIAMItemIAMUser && !r.Deleted {
out.Type, out.IAMUser = iamUserBoundaryType, nil
out.UserRevocation = &iamUserBoundary{User: item.IAMUser.AccessKey, Before: r.RevokedBefore}
}
if item.Type == madmin.SRIAMItemGroupInfo && !r.Deleted {
out.Type, out.GroupInfo = iamGroupBoundaryType, nil
out.GroupRevocation = &iamGroupBoundary{Group: item.GroupInfo.UpdateReq.Group, Before: r.RevokedBefore}
}
byType[item.Type] = append(byType[item.Type], out)
}
var items []iamReplicationItem
for _, typ := range []string{madmin.SRIAMItemPolicyMapping, madmin.SRIAMItemIAMUser, madmin.SRIAMItemSvcAcc, madmin.SRIAMItemGroupInfo, madmin.SRIAMItemPolicy} {
items = append(items, byType[typ]...)
}
if len(items) == 0 {
return nil
}
if c.iamRevisionProgress == nil {
c.iamRevisionProgress = make(map[string]iamRevisionProgress)
}
for id := range c.iamRevisionProgress {
if _, present := c.state.Peers[id]; !present {
delete(c.iamRevisionProgress, id)
}
}
var progressMu sync.Mutex
cerr := c.concDo(nil, func(id string, p madmin.PeerInfo) error {
// Bound each pass, but retain acknowledgements independently of the
// pass deadline or unrelated changes at either site.
peerCtx, cancel := context.WithTimeout(ctx, defaultContextTimeout)
defer cancel()
client, err := c.getAdminClient(peerCtx, id)
if err != nil {
return err
}
remote, err := executeIAMRevisionRequest(peerCtx, client, http.MethodGet, nil)
if err != nil {
return err
}
progressMu.Lock()
progress := c.iamRevisionProgress[id]
progressMu.Unlock()
progress.observePeer(remote)
defer func() {
progressMu.Lock()
c.iamRevisionProgress[id] = progress
progressMu.Unlock()
}()
var pending []iamReplicationItem
for _, item := range items {
path, version := iamReplicationMarker(item)
if progress.Acknowledged[path] != version {
pending = append(pending, item)
}
}
// Acknowledgements are only a replay optimization, never GC proof.
for path := range progress.Acknowledged {
if _, retained := snapshot[path]; !retained {
delete(progress.Acknowledged, path)
}
}
var failures []error
for next := 0; next < len(pending); {
end := min(next+maxIAMRevisionBatch, len(pending))
batch := pending[next:end]
remote, err = executeIAMRevisionRequest(peerCtx, client, http.MethodPut, &iamRevisionBatch{Version: iamRevisionProtocol, Items: batch})
if err != nil {
var batchErr *iamRevisionBatchError
if !errors.As(err, &batchErr) {
return errors.Join(append(failures, err)...)
}
failures = append(failures, err)
} else {
progress.observePeer(remote)
for _, item := range batch {
path, version := iamReplicationMarker(item)
progress.Acknowledged[path] = version
}
}
next = end
}
return errors.Join(failures...)
}, "IAM revision convergence")
return errors.Unwrap(cerr)
}
func iamReplicationMarker(item iamReplicationItem) (path, version string) {
path = iamDeletionPath(item.SRIAMItem)
if item.UserRevocation != nil {
path = getUserIdentityPath(item.UserRevocation.User, regUser)
}
if item.GroupRevocation != nil {
path = getGroupInfoPath(item.GroupRevocation.Group)
}
return path, item.Type + ":" + item.UpdatedAt.UTC().Format(time.RFC3339Nano)
}
func (store *IAMStoreSys) savePolicyDoc(ctx context.Context, policyName string, p *PolicyDoc) error {
return store.saveIAMRevision(ctx, getPolicyDocPath(policyName), p)
}
func (store *IAMStoreSys) saveMappedPolicy(ctx context.Context, name string, userType IAMUserType, isGroup bool, mp *MappedPolicy, opts ...options) error {
return store.saveIAMRevision(ctx, getMappedPolicyPath(name, userType, isGroup), mp, opts...)
}
func (store *IAMStoreSys) saveUserIdentity(ctx context.Context, name string, userType IAMUserType, u *UserIdentity, opts ...options) error {
return store.saveIAMRevision(ctx, getUserIdentityPath(name, userType), u, opts...)
}
func (store *IAMStoreSys) saveGroupInfo(ctx context.Context, name string, gi *GroupInfo) error {
return store.saveIAMRevision(ctx, getGroupInfoPath(name), gi)
}
func (store *IAMStoreSys) deletePolicyDoc(ctx context.Context, name string) error {
return store.saveIAMRevision(ctx, getPolicyDocPath(name), &PolicyDoc{Version: 1, Deleted: true})
}
func (store *IAMStoreSys) deleteMappedPolicy(ctx context.Context, name string, userType IAMUserType, isGroup bool) error {
return store.saveIAMRevision(ctx, getMappedPolicyPath(name, userType, isGroup), &MappedPolicy{Version: 1, Deleted: true})
}
func (store *IAMStoreSys) deleteUserIdentity(ctx context.Context, name string, userType IAMUserType) error {
return store.saveIAMRevision(ctx, getUserIdentityPath(name, userType), &UserIdentity{Version: 1, Deleted: true})
}
// Called under the identity's distributed revision lock, after verifying that
// its immutable STS token (or early-revocation retention) has expired. Only the
// old token-key mapping is removed; the reusable parent mapping is unaffected.
func expireIAMSTSConfig(ctx context.Context, store IAMStorageAPI, path string) error {
key := strings.TrimSuffix(strings.TrimPrefix(path, iamConfigSTSPrefix), "/"+iamIdentityFile)
if err := store.deleteIAMConfig(ctx, getMappedPolicyPath(key, stsUser, false)); err != nil && !errors.Is(err, errConfigNotFound) {
return err
}
return store.deleteIAMConfig(ctx, path)
}
type (
iamExpirationCleanupKey struct{}
iamExpirationCleanupState struct{ failed atomic.Bool }
)
func withIAMExpirationCleanup(ctx context.Context) context.Context {
if _, ok := ctx.Value(iamExpirationCleanupKey{}).(*iamExpirationCleanupState); ok {
return ctx
}
return context.WithValue(ctx, iamExpirationCleanupKey{}, &iamExpirationCleanupState{})
}
func bestEffortIAMExpiration(ctx context.Context, store IAMStorageAPI, path string) {
state, _ := ctx.Value(iamExpirationCleanupKey{}).(*iamExpirationCleanupState)
if ctx.Err() != nil || (state != nil && state.failed.Load()) {
return
}
ctx, cancel := context.WithTimeout(ctx, time.Second)
defer cancel()
// Failure leaves the expired record and its existing version intact.
// Stop optional reclamation for this load, while still loading healthy
// users. Healthy cleanup has no per-scan quota that could build a backlog.
if err := saveIAMRevision(ctx, store, path, &iamExpireIdentity{}); err != nil {
if state != nil {
state.failed.Store(true)
}
iamLogIf(ctx, err)
}
}
+330
View File
@@ -0,0 +1,330 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"encoding/json"
"errors"
"os"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/grid"
xnet "github.com/pgsty/silo-pkg/v3/net"
)
// Count physical saves: comparing timestamps alone would miss identical
// tombstones being rewritten on every heal pass.
type iamRevisionWriteCounter struct {
IAMStorageAPI
data []byte
writes int
}
func (s *iamRevisionWriteCounter) loadIAMConfig(_ context.Context, item any, _ string) error {
return json.Unmarshal(s.data, item)
}
func (s *iamRevisionWriteCounter) saveIAMConfig(_ context.Context, item any, _ string, _ ...options) error {
data, err := json.Marshal(item)
if err == nil {
s.data = data
s.writes++
}
return err
}
func TestIAMRevocationTombstoneReplayIsIdempotent(t *testing.T) {
at := time.Date(2026, 9, 14, 12, 0, 0, 0, time.UTC)
for _, record := range []struct {
name string
new func(bool) any
}{
{"user", func(deleted bool) any { return &UserIdentity{Version: 1, Deleted: deleted} }},
{"group", func(deleted bool) any { return &GroupInfo{Version: 1, Deleted: deleted} }},
{"policy", func(deleted bool) any { return &PolicyDoc{Version: 1, Deleted: deleted} }},
{"mapping", func(deleted bool) any { return &MappedPolicy{Version: 1, Deleted: deleted} }},
} {
t.Run(record.name, func(t *testing.T) {
data, err := json.Marshal(iamRevision{Deleted: true, UpdatedAt: at, RevokedBefore: at})
if err != nil {
t.Fatal(err)
}
store := &iamRevisionWriteCounter{data: data}
ctx := context.Background()
for range 3 {
// Site heal carries the original timestamp. Sibling notifications
// have no timestamp; both must leave an applied deletion untouched.
for _, replay := range []context.Context{withIAMReplicationTime(ctx, at), ctx} {
if err := saveIAMRevision(replay, store, record.name, record.new(true)); err != nil {
t.Fatal(err)
}
}
}
if store.writes != 0 || string(store.data) != string(data) {
t.Fatalf("replayed tombstone changed storage: writes=%d, record=%s", store.writes, store.data)
}
for _, deleted := range []bool{false, true} {
err := saveIAMRevision(withIAMReplicationTime(ctx, at.Add(-time.Second)), store, record.name, record.new(deleted))
if !errors.Is(err, errIAMStaleUpdate) {
t.Fatalf("older event accepted, deleted=%t: %v", deleted, err)
}
}
if err := saveIAMRevision(withIAMReplicationTime(ctx, at), store, record.name, record.new(false)); !errors.Is(err, errIAMStaleUpdate) {
t.Fatalf("equal-time recreation accepted: %v", err)
}
newer := at.Add(time.Minute)
if err := saveIAMRevision(withIAMReplicationTime(ctx, newer), store, record.name, record.new(true)); err != nil {
t.Fatal(err)
}
r, err := loadIAMRevision(ctx, store, record.name)
if err != nil || store.writes != 1 || !r.timestamp().Equal(newer) || !r.Deleted {
t.Fatalf("newer deletion did not advance storage: writes=%d, revision=%+v, error=%v", store.writes, r, err)
}
if err := saveIAMRevision(withIAMReplicationTime(ctx, newer.Add(time.Minute)), store, record.name, record.new(false)); err != nil {
t.Fatalf("newer recreation rejected: %v", err)
}
})
}
}
// Run with both object storage and etcd through TestIAMRevocation*Lifecycle.
// The object-store case uses the real peer RPC and deletion handler, so a
// spurious notification actually destroys the parent instead of only counting it.
func testIAMRevocationReplayAfterRecreation(ctx context.Context, t *testing.T, sys *IAMSys) {
t.Helper()
peer := &globalSiteReplicationSys
must := func(err error) {
t.Helper()
if err != nil {
t.Fatal(err)
}
}
user := "heal-recreated-parent"
req := madmin.AddOrUpdateUserReq{SecretKey: "valid-test-password", Status: madmin.AccountEnabled}
origin := UTCNow().Add(-time.Hour).Truncate(time.Millisecond)
deleted, recreated := origin.Add(time.Minute), origin.Add(3*time.Minute)
create := func(at time.Time) {
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, UserReq: &req}, at))
}
revoke := func(at time.Time) {
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, IsDeleteReq: true}, at))
}
create(origin)
revoke(deleted)
create(recreated)
child, _, err := sys.NewServiceAccount(ctx, user, nil, newServiceAccountOpts{
accessKey: "heal-recreated-child", secretKey: "valid-service-password",
})
must(err)
tg, err := grid.SetupTestGrid(2)
must(err)
t.Cleanup(tg.Cleanup)
var deletes atomic.Int32
server := &peerRESTServer{}
must(deleteUserRPC.Register(tg.Managers[1], func(req *grid.MSS) (grid.NoPayload, *grid.RemoteErr) {
deletes.Add(1)
return server.DeleteUserHandler(req)
}))
// Future user updates still use the normal peer reload notification.
must(loadUserRPC.Register(tg.Managers[1], server.LoadUserHandler))
host, err := xnet.ParseHost(strings.TrimPrefix(tg.Hosts[1], "http://"))
must(err)
previousNotifications := globalNotificationSys
globalNotificationSys = &NotificationSys{peerClients: []*peerRESTClient{{
host: host,
gridConn: func() *grid.Connection {
return tg.Managers[0].Connection(tg.Hosts[1])
},
}}}
t.Cleanup(func() { globalNotificationSys = previousNotifications })
assertLive := func(key string) {
t.Helper()
if _, ok := sys.GetUser(ctx, key); !ok {
t.Fatalf("live credential %s lost during deletion replay", key)
}
}
assertNoDelete := func() {
t.Helper()
if n := deletes.Load(); n != 0 {
t.Fatalf("retained revocation sent %d destructive sibling notifications", n)
}
}
for range 3 {
r, err := loadIAMRevision(ctx, sys.store, getUserIdentityPath(user, regUser))
must(err)
item, ok := iamDeletionItem(getUserIdentityPath(user, regUser), r)
if !ok || item.IAMUser == nil || !item.UpdatedAt.Equal(deleted) {
t.Fatal("recreated user lost its durable revocation replay")
}
must(peer.PeerIAMUserChangeHandler(ctx, item.IAMUser, item.UpdatedAt))
must(sys.store.LoadIAMCache(ctx, false))
assertLive(user)
assertLive(child.AccessKey)
assertNoDelete()
}
// A divergent site sends a previously unseen revocation between our old
// boundary and recreation. Retain it and revoke old children, but never
// turn it into an unversioned delete of the recreated parent.
delayed := deleted.Add(time.Minute)
revoke(delayed)
assertNoDelete()
must(sys.store.LoadIAMCache(ctx, false))
assertLive(user)
if _, ok := sys.GetUser(ctx, child.AccessKey); ok {
t.Fatal("child from before the delayed revocation remains usable")
}
r, err := loadIAMRevision(ctx, sys.store, getUserIdentityPath(user, regUser))
must(err)
if r.Deleted || !r.RevokedBefore.Equal(delayed) || !r.timestamp().Equal(recreated) {
t.Fatalf("retained revocation damaged the recreated identity: %+v", r)
}
// A genuinely newer deletion must still reach siblings and remove the
// parent plus credentials issued under its latest revocation boundary.
fresh, _, err := sys.NewServiceAccount(ctx, user, nil, newServiceAccountOpts{
accessKey: "heal-fresh-child", secretKey: "valid-service-password",
})
must(err)
latest := recreated.Add(time.Minute)
revoke(latest)
wantDeletes := int32(1)
if sys.HasWatcher() {
wantDeletes = 0
}
if n := deletes.Load(); n != wantDeletes {
t.Fatalf("new deletion notifications=%d, want %d", n, wantDeletes)
}
for _, key := range []string{user, fresh.AccessKey} {
if _, ok := sys.GetUser(ctx, key); ok {
t.Fatalf("newer deletion left credential %s usable", key)
}
}
// Exercise the actual sibling handler again against the already persisted
// tombstone. Its context has no revision; it must not re-stamp the record.
_, remoteErr := server.DeleteUserHandler(grid.NewMSSWith(map[string]string{peerRESTUser: user}))
if remoteErr != nil {
t.Fatal(remoteErr)
}
r, err = loadIAMRevision(ctx, sys.store, getUserIdentityPath(user, regUser))
must(err)
if !r.Deleted || !r.timestamp().Equal(latest) {
t.Fatalf("sibling re-stamped the tombstone: got %s, want %s", r.timestamp(), latest)
}
create(latest.Add(time.Minute))
_, remoteErr = server.DeleteUserHandler(grid.NewMSSWith(map[string]string{peerRESTUser: user}))
if remoteErr != nil {
t.Fatal(remoteErr)
}
assertLive(user)
revoke(latest)
must(sys.store.LoadIAMCache(ctx, false))
assertLive(user)
}
// Counts what a retained revocation actually sends to sibling nodes.
func TestIAMRevocationRetainedReloadsSibling(t *testing.T) {
resetTestGlobals()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
disks, err := getRandomDisks(1)
if err != nil {
t.Fatal(err)
}
obj, _, err := initObjectLayer(ctx, mustGetPoolEndpoints(0, disks...))
if err != nil {
t.Fatal(err)
}
initAllSubsystems(ctx)
globalIAMSys.Init(ctx, obj, nil, 2*time.Second)
defer os.RemoveAll(disks[0])
defer obj.Shutdown(ctx)
defer resetTestGlobals()
sys, peer := globalIAMSys, &globalSiteReplicationSys
must := func(err error) {
t.Helper()
if err != nil {
t.Fatal(err)
}
}
user := "retained-parent"
req := madmin.AddOrUpdateUserReq{SecretKey: "valid-test-password", Status: madmin.AccountEnabled}
origin := UTCNow().Add(-time.Hour).Truncate(time.Millisecond)
deleted, recreated := origin.Add(time.Minute), origin.Add(3*time.Minute)
create := func(at time.Time) {
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, UserReq: &req}, at))
}
revoke := func(at time.Time) {
must(peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, IsDeleteReq: true}, at))
}
create(origin)
revoke(deleted)
create(recreated)
child, _, err := sys.NewServiceAccount(ctx, user, nil, newServiceAccountOpts{
accessKey: "retained-child", secretKey: "valid-service-password",
})
must(err)
// A sibling shares persistent state but has an independent IAM cache.
sibling := &IAMStoreSys{IAMStorageAPI: newIAMObjectStore(obj, sys.usersSysType)}
must(sibling.LoadIAMCache(ctx, false))
if _, ok := sibling.GetUser(child.AccessKey); !ok {
t.Fatal("sibling fixture did not load child")
}
tg, err := grid.SetupTestGrid(2)
must(err)
t.Cleanup(tg.Cleanup)
var deletes, loads atomic.Int32
server := &peerRESTServer{}
must(deleteUserRPC.Register(tg.Managers[1], func(r *grid.MSS) (grid.NoPayload, *grid.RemoteErr) {
deletes.Add(1)
return server.DeleteUserHandler(r)
}))
must(loadUserRPC.Register(tg.Managers[1], func(r *grid.MSS) (grid.NoPayload, *grid.RemoteErr) {
loads.Add(1)
// LoadUserHandler delegates to this same cache reload method.
if err := sibling.UserNotificationHandler(ctx, r.Get(peerRESTUser), regUser); err != nil {
return grid.NoPayload{}, grid.NewRemoteErr(err)
}
return grid.NoPayload{}, nil
}))
host, err := xnet.ParseHost(strings.TrimPrefix(tg.Hosts[1], "http://"))
must(err)
prev := globalNotificationSys
globalNotificationSys = &NotificationSys{peerClients: []*peerRESTClient{{
host: host,
gridConn: func() *grid.Connection { return tg.Managers[0].Connection(tg.Hosts[1]) },
}}}
t.Cleanup(func() { globalNotificationSys = prev })
delayed := deleted.Add(time.Minute)
revoke(delayed)
t.Logf("sibling notifications after a retained revocation: destructive=%d reload=%d", deletes.Load(), loads.Load())
if deletes.Load() != 0 {
t.Errorf("destructive sibling delete sent: %d", deletes.Load())
}
if loads.Load() == 0 {
t.Errorf("retained revocation did not notify the sibling")
}
if _, ok := sibling.GetUser(child.AccessKey); ok {
t.Error("sibling still resolves revoked child")
}
if _, ok := sibling.GetUser(user); !ok {
t.Error("sibling lost live parent")
}
if _, ok := sys.store.GetUser(user); !ok {
t.Error("live parent lost")
}
if _, ok := sys.store.GetUser(child.AccessKey); ok {
t.Error("revoked child still resolves on the receiving node")
}
}
+535
View File
@@ -0,0 +1,535 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"net/http/httptest"
"os"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
etcd "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/namespace"
)
// Exercise the persisted IAM store and the same peer handler used by site heal.
// A delete must survive a cache reload and an older create arriving afterwards.
func TestIAMRevocationRejectsOfflineUser(t *testing.T) {
resetTestGlobals()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
obj, disk, err := prepareFS(ctx)
if err != nil {
t.Fatal(err)
}
defer os.RemoveAll(disk)
defer obj.Shutdown(ctx)
defer resetTestGlobals()
user := "offline-revoked-user"
req := madmin.AddOrUpdateUserReq{SecretKey: "test-password-valid", Status: madmin.AccountEnabled}
created, err := globalIAMSys.CreateUser(ctx, user, req)
if err != nil {
t.Fatal(err)
}
if err = globalIAMSys.DeleteUser(ctx, user, false); err != nil {
t.Fatal(err)
}
if err = globalIAMSys.store.LoadIAMCache(ctx, false); err != nil {
t.Fatal(err)
}
if err = globalSiteReplicationSys.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, UserReq: &req}, created); err != nil {
t.Fatal(err)
}
if _, err = globalIAMSys.GetUserInfo(ctx, user); !errors.Is(err, errNoSuchUser) {
t.Fatalf("revoked user restored by old peer event: %v", err)
}
}
func TestIAMRevocationHealingContinuesAfterPeerRejectsDelete(t *testing.T) {
resetTestGlobals()
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
obj, disk, err := prepareFS(ctx)
if err != nil {
t.Fatal(err)
}
defer os.RemoveAll(disk)
defer obj.Shutdown(ctx)
defer resetTestGlobals()
req := madmin.AddOrUpdateUserReq{SecretKey: "valid-test-password", Status: madmin.AccountEnabled}
if _, err := globalIAMSys.CreateUser(ctx, "heal-sync", req); err != nil {
t.Fatal(err)
}
if _, err := globalIAMSys.CreateUser(ctx, "heal-deleted", req); err != nil {
t.Fatal(err)
}
if err := globalIAMSys.DeleteUser(ctx, "heal-deleted", false); err != nil {
t.Fatal(err)
}
p, err := globalIAMSys.store.GetPolicy("readwrite")
if err != nil {
t.Fatal(err)
}
if _, err := globalIAMSys.SetPolicy(ctx, "heal-new-policy", p); err != nil {
t.Fatal(err)
}
var liveUpdates atomic.Int32
peer := func(id string, rejectDelete bool) *httptest.Server {
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
switch {
case strings.HasSuffix(r.URL.Path, "/metainfo"):
_ = json.NewEncoder(w).Encode(madmin.SRInfo{DeploymentID: id})
case r.URL.Path == "/minio/admin/v3/site-replication/peer/iam-revisions":
if r.Method == http.MethodGet {
_ = json.NewEncoder(w).Encode(iamRevisionStatus{Version: iamRevisionProtocol, Node: "node-1", Instance: id, Digest: "fixture"})
return
}
var batch iamRevisionBatch
if err := json.NewDecoder(r.Body).Decode(&batch); err != nil {
t.Error(err)
w.WriteHeader(http.StatusBadRequest)
return
}
for _, item := range batch.Items {
if rejectDelete && iamDeletionPath(item.SRIAMItem) != "" {
w.WriteHeader(http.StatusForbidden)
_, _ = w.Write([]byte(`{"Code":"AccessDenied","Message":"delete rejected"}`))
return
}
if id == "healthy" && item.Type == madmin.SRIAMItemPolicy && item.Name == "heal-new-policy" && len(item.Policy) > 0 {
liveUpdates.Add(1)
}
}
_ = json.NewEncoder(w).Encode(iamRevisionStatus{Version: iamRevisionProtocol, Node: "node-1", Instance: id, Digest: "fixture"})
default:
t.Errorf("unexpected peer request %s", r.URL.Path)
w.WriteHeader(http.StatusNotFound)
}
}))
}
healthy, rejected := peer("healthy", false), peer("rejected", true)
defer healthy.Close()
defer rejected.Close()
c := &SiteReplicationSys{enabled: true, state: srState{
ServiceAccountAccessKey: "heal-sync",
Peers: map[string]madmin.PeerInfo{
globalDeploymentID(): {Name: "local", DeploymentID: globalDeploymentID()},
"healthy": {Name: "healthy", DeploymentID: "healthy", Endpoint: healthy.URL},
"rejected": {Name: "rejected", DeploymentID: "rejected", Endpoint: rejected.URL},
},
}}
if err := c.healIAMSystem(ctx, obj); err == nil {
t.Fatal("deletion failure was not reported")
}
if liveUpdates.Load() == 0 {
t.Fatal("one peer rejecting a deletion blocked unrelated live IAM healing to a healthy peer")
}
}
func TestIAMRevocationLifecycle(t *testing.T) {
testIAMRevocationLifecycle(t, nil)
}
func TestIAMRevocationEtcdLifecycle(t *testing.T) {
endpoint := os.Getenv("SILO_TEST_IAM_REVOCATION_ETCD")
if endpoint == "" {
t.Skip("set SILO_TEST_IAM_REVOCATION_ETCD to a disposable etcd endpoint")
}
connection, err := etcd.New(etcd.Config{Endpoints: strings.Split(endpoint, ","), DialTimeout: 5 * time.Second})
if err != nil {
t.Fatal(err)
}
defer connection.Close()
// The facade borrows the connection's services. Close the owning client,
// not namespace.Watcher while IAM's canceled watch loop is winding down.
ctx, cancel := context.WithCancel(connection.Ctx())
defer cancel()
client := etcd.NewCtxClient(ctx, etcd.WithZapLogger(connection.GetLogger()))
prefix := fmt.Sprintf("/silo-revocation-test/%d/", time.Now().UnixNano())
client.KV = namespace.NewKV(connection.KV, prefix)
client.Watcher = namespace.NewWatcher(connection.Watcher, prefix)
client.Lease = connection.Lease
testIAMRevocationLifecycle(t, client)
}
func testIAMRevocationLifecycle(t *testing.T, client *etcd.Client) {
resetTestGlobals()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
disks, err := getRandomDisks(1)
if err != nil {
t.Fatal(err)
}
disk := disks[0]
obj, _, err := initObjectLayer(ctx, mustGetPoolEndpoints(0, disks...))
if err == nil {
initAllSubsystems(ctx)
globalIAMSys.Init(ctx, obj, client, 2*time.Second)
}
if err != nil {
t.Fatal(err)
}
defer os.RemoveAll(disk)
defer obj.Shutdown(ctx)
defer resetTestGlobals()
sys, peer := globalIAMSys, &globalSiteReplicationSys
must := func(t *testing.T, err error) {
t.Helper()
if err != nil {
t.Fatal(err)
}
}
reload := func(t *testing.T) { t.Helper(); must(t, sys.store.LoadIAMCache(ctx, false)) }
req := madmin.AddOrUpdateUserReq{SecretKey: "valid-test-password", Status: madmin.AccountEnabled}
origin := UTCNow().Add(-time.Hour).Truncate(time.Millisecond)
createUser := func(t *testing.T, name string) {
t.Helper()
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: name, UserReq: &req}, origin))
}
assertAbsent := func(t *testing.T, name string) {
t.Helper()
if _, ok := sys.GetUser(ctx, name); ok {
t.Fatalf("revoked credential %s is usable", name)
}
}
t.Run("replay after recreation", func(t *testing.T) {
testIAMRevocationReplayAfterRecreation(ctx, t, sys)
})
t.Run("origin timestamp and recreation", func(t *testing.T) {
name := "revocation-recreate"
createUser(t, name)
ui, ok := sys.store.GetUser(name)
if !ok || !ui.UpdatedAt.Equal(origin) {
t.Fatalf("origin time changed: %v", ui.UpdatedAt)
}
must(t, sys.DeleteUser(ctx, name, false))
reload(t)
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: name, UserReq: &req}, time.Time{}))
assertAbsent(t, name)
newTime := UTCNow().Add(time.Minute)
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: name, UserReq: &req}, newTime))
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: name, IsDeleteReq: true}, origin.Add(time.Second)))
reload(t)
ui, ok = sys.store.GetUser(name)
if !ok || !ui.UpdatedAt.Equal(newTime) || ui.RevokedBefore.IsZero() {
t.Fatalf("newer recreation lost, or deletion boundary missing: present=%v", ok)
}
})
t.Run("groups policies and mappings", func(t *testing.T) {
user, group, name := "revocation-member", "revocation-group", "revocation-policy"
createUser(t, user)
p, err := sys.store.GetPolicy("readwrite")
must(t, err)
must(t, peer.PeerAddPolicyHandler(ctx, name, &p, origin))
add := &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: group, Members: []string{user}}}
must(t, peer.PeerGroupInfoChangeHandler(ctx, add, origin))
for _, isGroup := range []bool{false, true} {
entity := user
if isGroup {
entity = group
}
mp := &madmin.SRPolicyMapping{UserOrGroup: entity, Policy: name, UserType: int(regUser), IsGroup: isGroup}
must(t, peer.PeerPolicyMappingHandler(ctx, mp, origin))
_, err = sys.PolicyDBSet(ctx, entity, "", regUser, isGroup)
must(t, err)
must(t, peer.PeerPolicyMappingHandler(ctx, mp, origin))
if _, ok := sys.store.GetMappedPolicy(entity, isGroup); ok {
t.Fatal("old grant restored")
}
}
// This receiver never saw the member-removal event preceding deletion.
must(t, peer.PeerGroupInfoChangeHandler(ctx, &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: group, IsRemove: true}}, UTCNow()))
must(t, sys.DeletePolicy(ctx, name, true))
reload(t)
must(t, peer.PeerGroupInfoChangeHandler(ctx, add, origin))
must(t, peer.PeerAddPolicyHandler(ctx, name, &p, origin))
if _, err = sys.GetGroupDescription(group); !errors.Is(err, errNoSuchGroup) {
t.Fatalf("group restored: %v", err)
}
if _, err = sys.store.GetPolicyDoc(name); !errors.Is(err, errNoSuchPolicy) {
t.Fatalf("policy restored: %v", err)
}
paths, err := sys.store.listIAMConfigPaths(ctx)
must(t, err)
found := make(map[string]bool)
for _, path := range paths {
r, err := loadIAMRevision(ctx, sys.store, path)
must(t, err)
if item, ok := iamDeletionItem(path, r); ok {
found[iamDeletionPath(item)] = true
if item.UpdatedAt.IsZero() {
t.Fatal("undated delete replay")
}
}
}
for _, path := range []string{getGroupInfoPath(group), getPolicyDocPath(name), getMappedPolicyPath(user, regUser, false), getMappedPolicyPath(group, regUser, true)} {
if !found[path] {
t.Errorf("deletion missing from heal: %s", path)
}
}
})
t.Run("parent revokes service accounts and STS", func(t *testing.T) {
parent := "revocation-parent"
createUser(t, parent)
svc, svcAt, err := sys.NewServiceAccount(withIAMReplicationTime(ctx, origin), parent, nil, newServiceAccountOpts{accessKey: "revocation-service", secretKey: "valid-service-password"})
must(t, err)
secret, err := getTokenSigningKey()
must(t, err)
sts, err := auth.GetNewCredentialsWithMetadata(map[string]any{"exp": UTCNow().Add(time.Hour).Unix(), parentClaim: parent}, secret)
must(t, err)
sts.ParentUser = parent
_, err = sys.SetTempUser(withIAMReplicationTime(ctx, origin), sts.AccessKey, sts, "readwrite")
must(t, err)
must(t, sys.DeleteUser(ctx, parent, false))
reload(t)
assertAbsent(t, parent)
assertAbsent(t, svc.AccessKey)
assertAbsent(t, sts.AccessKey)
// Recreate the parent, then deliver old child events from the offline site.
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, UserReq: &req}, UTCNow()))
must(t, peer.PeerSvcAccChangeHandler(ctx, &madmin.SRSvcAccChange{Create: &madmin.SRSvcAccCreate{Parent: parent, AccessKey: svc.AccessKey, SecretKey: svc.SecretKey}}, svcAt))
must(t, peer.PeerSTSAccHandler(ctx, &madmin.SRSTSCredential{AccessKey: sts.AccessKey, SecretKey: sts.SecretKey, ParentUser: parent, SessionToken: sts.SessionToken, ParentPolicyMapping: "readwrite"}, origin))
reload(t)
assertAbsent(t, svc.AccessKey)
assertAbsent(t, sts.AccessKey)
// A freshly issued credential is still supported after deliberate recreation.
_, _, err = sys.NewServiceAccount(ctx, parent, nil, newServiceAccountOpts{accessKey: "new-service", secretKey: "valid-service-password"})
must(t, err)
if _, ok := sys.GetUser(ctx, "new-service"); !ok {
t.Fatal("fresh service account rejected")
}
})
t.Run("delete before first create", func(t *testing.T) {
name := "revocation-unseen"
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: name, IsDeleteReq: true}, UTCNow()))
createUser(t, name)
assertAbsent(t, name)
svc := "unseen-service"
must(t, peer.PeerSvcAccChangeHandler(ctx, &madmin.SRSvcAccChange{Delete: &madmin.SRSvcAccDelete{AccessKey: svc}}, UTCNow()))
must(t, peer.PeerSvcAccChangeHandler(ctx, &madmin.SRSvcAccChange{Create: &madmin.SRSvcAccCreate{Parent: "revocation-recreate", AccessKey: svc, SecretKey: "valid-service-password"}}, origin))
assertAbsent(t, svc)
})
t.Run("recreation arrives before revocation", func(t *testing.T) {
parent := "reordered-parent"
createUser(t, parent)
child, _, err := sys.NewServiceAccount(withIAMReplicationTime(ctx, origin), parent, nil, newServiceAccountOpts{accessKey: "reordered-child", secretKey: "valid-service-password"})
must(t, err)
newTime, deleteTime := UTCNow(), origin.Add(time.Minute)
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, UserReq: &req}, newTime))
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, IsDeleteReq: true}, deleteTime))
assertAbsent(t, child.AccessKey)
reload(t)
assertAbsent(t, child.AccessKey)
u, ok := sys.GetUser(ctx, parent)
if !ok || !u.UpdatedAt.Equal(newTime) || !u.RevokedBefore.Equal(deleteTime) {
t.Fatal("reordered revocation damaged the new parent or lost its boundary")
}
r, err := loadIAMRevision(ctx, sys.store, getUserIdentityPath(parent, regUser))
must(t, err)
item, ok := iamDeletionItem(getUserIdentityPath(parent, regUser), r)
if !ok || !item.UpdatedAt.Equal(deleteTime) {
t.Fatal("recreation erased deletion replay")
}
})
t.Run("user cleanup does not supersede group deletion", func(t *testing.T) {
user, group := "cascade-user", "cascade-group"
createUser(t, user)
must(t, peer.PeerGroupInfoChangeHandler(ctx, &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: group, Members: []string{user}}}, origin))
// On the origin site the group was removed before the user, but the
// recovering receiver processes those independent events in reverse.
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: user, IsDeleteReq: true}, origin.Add(2*time.Minute)))
must(t, peer.PeerGroupInfoChangeHandler(ctx, &madmin.SRGroupInfo{UpdateReq: madmin.GroupAddRemove{Group: group, IsRemove: true}}, origin.Add(time.Minute)))
reload(t)
if _, err := sys.GetGroupDescription(group); !errors.Is(err, errNoSuchGroup) {
t.Fatalf("deleted group survived reordered cleanup: %v", err)
}
groups, err := sys.ListGroups(ctx)
must(t, err)
for _, name := range groups {
if name == group {
t.Fatal("deleted group listed")
}
}
})
t.Run("parent revocation covers later updates to existing children", func(t *testing.T) {
parent, key := "late-update-parent", "late-update-child"
createUser(t, parent)
_, _, err := sys.NewServiceAccount(withIAMReplicationTime(ctx, origin), parent, nil, newServiceAccountOpts{accessKey: key, secretKey: "valid-service-password"})
must(t, err)
// This site missed the deletion and subsequently edited an old child.
_, err = sys.UpdateServiceAccount(withIAMReplicationTime(ctx, origin.Add(2*time.Minute)), key, updateServiceAccountOpts{description: "edited while the peer was offline"})
must(t, err)
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, IsDeleteReq: true}, origin.Add(time.Minute)))
assertAbsent(t, key)
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, UserReq: &req}, origin.Add(3*time.Minute)))
reload(t)
assertAbsent(t, key)
})
t.Run("old generation cannot return with a newer event timestamp", func(t *testing.T) {
parent := "generation-parent"
createUser(t, parent)
deleteTime := origin.Add(time.Minute)
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, IsDeleteReq: true}, deleteTime))
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, UserReq: &req}, origin.Add(2*time.Minute)))
// Another offline site issued this child under the original parent,
// after this site's delete/recreate. Wall-clock ordering cannot identify it.
late := origin.Add(3 * time.Minute)
must(t, peer.PeerSvcAccChangeHandler(ctx, &madmin.SRSvcAccChange{Create: &madmin.SRSvcAccCreate{Parent: parent, AccessKey: "old-gen-service", SecretKey: "valid-service-password"}}, late))
secret, err := getTokenSigningKey()
must(t, err)
sts, err := auth.GetNewCredentialsWithMetadata(map[string]any{"exp": UTCNow().Add(time.Hour).Unix(), parentClaim: parent}, secret)
must(t, err)
must(t, peer.PeerSTSAccHandler(ctx, &madmin.SRSTSCredential{AccessKey: sts.AccessKey, SecretKey: sts.SecretKey, ParentUser: parent, SessionToken: sts.SessionToken}, late))
reload(t)
assertAbsent(t, "old-gen-service")
assertAbsent(t, sts.AccessKey)
// A local issuer knows the new boundary and signs it into both kinds
// of child. Untrusted inherited claims cannot select that boundary.
child, _, err := sys.NewServiceAccount(ctx, parent, nil, newServiceAccountOpts{accessKey: "new-gen-service", secretKey: "valid-service-password", claims: map[string]any{iamParentRevocationClaim: "forged"}})
must(t, err)
newClaims := map[string]any{"exp": UTCNow().Add(time.Hour).Unix(), parentClaim: parent}
must(t, setIAMParentRevocationClaim(ctx, sys.store, parent, newClaims))
fresh, err := auth.GetNewCredentialsWithMetadata(newClaims, secret)
must(t, err)
fresh.ParentUser = parent
_, err = sys.SetTempUser(ctx, fresh.AccessKey, fresh, "")
must(t, err)
reload(t)
// A periodic reload retains the STS cache. Explicitly clear it to
// exercise the cold credential load performed after process restart.
cache := sys.store.lock()
cache.iamSTSAccountsMap = make(map[string]UserIdentity)
sys.store.unlock()
for _, key := range []string{child.AccessKey, fresh.AccessKey} {
u, ok := sys.GetUser(ctx, key)
if !ok || !iamCredentialSurvivesRevocation(u.Credentials, deleteTime) {
t.Fatalf("new-generation credential %s rejected", key)
}
}
})
t.Run("late revocation preserves proven new-generation children", func(t *testing.T) {
for _, recreateFirst := range []bool{false, true} {
parent := fmt.Sprintf("gen-parent-%t", recreateFirst)
createUser(t, parent)
deleteTime, createTime := origin.Add(time.Minute), origin.Add(2*time.Minute)
if recreateFirst {
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, UserReq: &req}, createTime))
}
key := fmt.Sprintf("gen-child-%t", recreateFirst)
must(t, peer.PeerSvcAccChangeHandler(ctx, &madmin.SRSvcAccChange{Create: &madmin.SRSvcAccCreate{Parent: parent, AccessKey: key, SecretKey: "valid-service-password", Claims: map[string]any{iamParentRevocationClaim: deleteTime.Format(time.RFC3339Nano)}}}, createTime.Add(time.Second)))
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, IsDeleteReq: true}, deleteTime))
if !recreateFirst {
assertAbsent(t, key)
must(t, peer.PeerIAMUserChangeHandler(ctx, &madmin.SRIAMUser{AccessKey: parent, UserReq: &req}, createTime))
}
reload(t)
if _, ok := sys.GetUser(ctx, key); !ok {
t.Fatal("late revocation deleted a child issued by the recreated parent")
}
}
})
t.Run("cold loading preserves site-signed STS", func(t *testing.T) {
parent := "cold-sts-parent"
createUser(t, parent)
secret := "site-signing-key-valid"
_, _, err := sys.NewServiceAccount(ctx, globalActiveCred.AccessKey, nil, newServiceAccountOpts{
accessKey: siteReplicatorSvcAcc, secretKey: secret, allowSiteReplicatorAccount: true,
})
must(t, err)
setReplication := func(enabled bool) {
globalSiteReplicationSys.Lock()
globalSiteReplicationSys.enabled = enabled
globalSiteReplicationSys.Unlock()
globalSiteReplicatorCred.Set("")
}
setReplication(true)
defer setReplication(false)
cred, err := auth.GetNewCredentialsWithMetadata(map[string]any{"exp": UTCNow().Add(time.Hour).Unix(), parentClaim: parent}, secret)
must(t, err)
cred.ParentUser = parent
_, err = sys.SetTempUser(ctx, cred.AccessKey, cred, "")
must(t, err)
// IAM can load before the site replication manager during startup.
// A signing key that is not available yet must not delete live tokens.
setReplication(false)
for range 3 {
unverified := make(map[string]UserIdentity)
_ = sys.store.loadUser(ctx, cred.AccessKey, stsUser, unverified)
if _, ok := unverified[cred.AccessKey]; ok {
t.Fatal("accepted STS before the signing key became available")
}
}
r, err := loadIAMRevision(ctx, sys.store, getUserIdentityPath(cred.AccessKey, stsUser))
must(t, err)
if r.Credentials.SessionToken == "" {
t.Fatal("cold IAM load physically deleted a non-expired site-signed STS credential")
}
setReplication(true)
loaded := make(map[string]UserIdentity)
must(t, sys.store.loadUser(ctx, cred.AccessKey, stsUser, loaded))
if _, ok := loaded[cred.AccessKey]; !ok {
t.Fatal("STS credential did not recover when the signing key became available")
}
})
t.Run("unverifiable STS stay denied and expired STS are removed", func(t *testing.T) {
parent := "invalid-sts-parent"
createUser(t, parent)
// Keep the etcd watcher from cleaning half of the fixture before the
// second record is seeded; this subtest exercises the loader directly.
sys.store.lock()
defer sys.store.unlock()
for _, expired := range []bool{false, true} {
cred, err := auth.GetNewCredentialsWithMetadata(map[string]any{"exp": UTCNow().Add(time.Hour).Unix(), parentClaim: parent}, "unavailable-test-signing-key")
must(t, err)
cred.ParentUser = parent
if expired {
cred.Expiration = UTCNow().Add(-time.Minute)
}
identityPath := getUserIdentityPath(cred.AccessKey, stsUser)
mappingPath := getMappedPolicyPath(cred.AccessKey, stsUser, false)
// Seed disk directly to exercise loading, including existing records
// whose key is unknown. The write API should not accept such tokens.
must(t, sys.store.saveIAMConfig(ctx, &UserIdentity{Version: 1, Credentials: cred, UpdatedAt: UTCNow()}, identityPath))
must(t, sys.store.saveIAMConfig(ctx, &MappedPolicy{Version: 1, Policies: "readwrite"}, mappingPath))
loaded := make(map[string]UserIdentity)
_ = sys.store.loadUser(ctx, cred.AccessKey, stsUser, loaded)
if _, ok := loaded[cred.AccessKey]; ok {
t.Fatalf("invalid STS accepted, expired=%t", expired)
}
for _, path := range []string{identityPath, mappingPath} {
var record map[string]any
err := sys.store.loadIAMConfig(ctx, &record, path)
if expired {
if !errors.Is(err, errConfigNotFound) {
t.Fatalf("expired STS data not cleaned up at %s: %v", path, err)
}
} else {
must(t, err)
}
}
}
})
}
+239
View File
@@ -0,0 +1,239 @@
// Copyright (c) 2026 PGSTY
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"encoding/json"
"errors"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
)
// The receiver missed a deletion and still has the previous service key.
// A later full snapshot must replace it, including when the owner changed.
func TestIAMServiceAccountRecreation(t *testing.T) {
for _, backend := range []string{"object", "etcd"} {
t.Run(backend, func(t *testing.T) {
ctx, sys, _ := prepareIAMRevisionFixture(t, backend)
origin := UTCNow().Add(-time.Hour)
for _, parent := range []string{"old-owner", "new-owner"} {
_, err := sys.CreateUser(withIAMReplicationTime(ctx, origin), parent, madmin.AddOrUpdateUserReq{SecretKey: "valid-owner-password", Status: madmin.AccountEnabled})
mustIAM(t, err)
}
const key = "reusable-service"
old := &madmin.SRSvcAccChange{Create: &madmin.SRSvcAccCreate{Parent: "old-owner", AccessKey: key, SecretKey: "old-service-password"}}
mustIAM(t, globalSiteReplicationSys.PeerSvcAccChangeHandler(ctx, old, origin))
oldIdentity, _ := sys.store.GetUser(key)
_, err := sys.PolicyDBSet(withIAMReplicationTime(ctx, origin), key, "readwrite", svcUser, false)
mustIAM(t, err)
boundary, newer := origin.Add(time.Minute), origin.Add(2*time.Minute)
fresh := iamReplicationItem{SRIAMItem: madmin.SRIAMItem{Type: madmin.SRIAMItemSvcAcc, UpdatedAt: newer, SvcAccChange: &madmin.SRSvcAccChange{Create: &madmin.SRSvcAccCreate{Parent: "new-owner", AccessKey: key, SecretKey: "new-service-password", Status: auth.AccountOff}}}, RevokedBefore: boundary}
mustIAM(t, applyIAMReplicationItem(ctx, fresh))
// A sibling may have missed the notification of the committed
// replacement. An equal-version retry must refresh that cache too.
staleCache := sys.store.lock()
staleCache.iamUsersMap[key] = oldIdentity
sys.store.unlock()
mustIAM(t, applyIAMReplicationItem(ctx, fresh)) // duplicate delivery is acknowledged
if current, _ := sys.store.GetUser(key); current.Credentials.SecretKey != fresh.SvcAccChange.Create.SecretKey {
t.Fatal("duplicate snapshot acknowledged without refreshing the stale cache")
}
mustIAM(t, globalSiteReplicationSys.PeerSvcAccChangeHandler(ctx, old, origin))
mustIAM(t, globalSiteReplicationSys.PeerSvcAccChangeHandler(ctx, &madmin.SRSvcAccChange{Delete: &madmin.SRSvcAccDelete{AccessKey: key}}, boundary))
mustIAM(t, sys.store.LoadIAMCache(ctx, false))
u, ok := sys.store.GetUser(key)
if !ok || u.Credentials.SecretKey != fresh.SvcAccChange.Create.SecretKey || u.Credentials.ParentUser != "new-owner" || u.Credentials.Status != auth.AccountOff || !u.UpdatedAt.Equal(newer) || !u.RevokedBefore.Equal(boundary) {
t.Fatal("recreation did not retain the new identity, disabled status, source version and revocation")
}
if _, ok := sys.GetUser(ctx, key); ok {
t.Fatal("replicated disabled service can authenticate")
}
cache := sys.store.rlock()
_, mapped := cache.cachedMappedPolicy(key, svcUser, false)
sys.store.runlock()
if mapped {
t.Fatal("recreated service inherited an older mapping")
}
_, err = sys.PolicyDBSet(withIAMReplicationTime(ctx, origin), key, "readwrite", svcUser, false)
if !errors.Is(err, errIAMStaleUpdate) {
t.Fatalf("old service mapping replay was accepted: %v", err)
}
_, _, err = sys.NewServiceAccount(ctx, "new-owner", nil, newServiceAccountOpts{accessKey: key, secretKey: "local-service-password"})
if !errors.Is(err, errIAMServiceAccountNotAllowed) {
t.Fatalf("local duplicate creation must remain rejected: %v", err)
}
// Outbound snapshots must carry the retained service boundary too.
out, err := globalSiteReplicationSys.replicationItem(ctx, fresh.SRIAMItem)
mustIAM(t, err)
if !out.RevokedBefore.Equal(boundary) {
t.Fatal("outbound service snapshot lost its revocation")
}
})
}
}
func TestIAMServiceAccountReplicationRejectsOtherCredentialKinds(t *testing.T) {
ctx, sys, _ := prepareIAMRevisionFixture(t)
_, err := sys.CreateUser(ctx, "builtin-collision", madmin.AddOrUpdateUserReq{SecretKey: "valid-user-password", Status: madmin.AccountEnabled})
mustIAM(t, err)
secret, err := getTokenSigningKey()
mustIAM(t, err)
token, err := auth.GetNewCredentialsWithMetadata(map[string]any{"exp": UTCNow().Add(time.Hour).Unix(), parentClaim: "builtin-collision"}, secret)
mustIAM(t, err)
token.ParentUser = "builtin-collision"
_, err = sys.SetTempUser(ctx, token.AccessKey, token, "")
mustIAM(t, err)
for _, key := range []string{"builtin-collision", token.AccessKey} {
_, _, err := sys.NewServiceAccount(withIAMReplicationTime(ctx, UTCNow().Add(time.Minute)), "another-owner", nil, newServiceAccountOpts{accessKey: key, secretKey: "valid-service-password"})
if !errors.Is(err, errIAMServiceAccountNotAllowed) {
t.Fatalf("service replication replaced another credential kind: %v", err)
}
}
}
// SR configuration can be temporarily unreadable even though a service token
// is signed with its own valid secret. Do not acknowledge a failed cache load.
func TestIAMServiceAccountRetryReportsClaimLoadFailure(t *testing.T) {
ctx, sys, _ := prepareIAMRevisionFixture(t)
globalSiteReplicatorCred.RLock()
previousSigningKey := globalSiteReplicatorCred.secretKey
globalSiteReplicatorCred.RUnlock()
globalSiteReplicatorCred.Set("")
t.Cleanup(func() { globalSiteReplicatorCred.Set(previousSigningKey) })
_, err := sys.CreateUser(ctx, "retry-owner", madmin.AddOrUpdateUserReq{SecretKey: "valid-owner-password", Status: madmin.AccountEnabled})
mustIAM(t, err)
opts := newServiceAccountOpts{accessKey: "retry-service", secretKey: "valid-service-password"}
_, _, err = sys.NewServiceAccount(ctx, "retry-owner", nil, opts)
mustIAM(t, err)
old, _ := sys.store.GetUser(opts.accessKey)
opts.secretKey = "replacement-service-password"
at, err := sys.UpdateServiceAccount(ctx, opts.accessKey, updateServiceAccountOpts{secretKey: opts.secretKey})
mustIAM(t, err)
cache := sys.store.lock()
cache.iamUsersMap[opts.accessKey] = old // Missed sibling notification.
sys.store.unlock()
globalSiteReplicationSys.Lock()
globalSiteReplicationSys.enabled = true // No site-replicator credential is installed.
globalSiteReplicationSys.Unlock()
_, _, err = sys.NewServiceAccount(withIAMReplicationTime(ctx, at), "retry-owner", nil, opts)
if err == nil {
t.Fatal("acknowledged service retry despite failed claims loading")
}
if _, ok := sys.store.GetUser(opts.accessKey); ok {
t.Fatal("failed cache refresh retained the superseded service secret")
}
globalSiteReplicationSys.Lock()
globalSiteReplicationSys.enabled = false
globalSiteReplicationSys.Unlock()
_, _, err = sys.NewServiceAccount(withIAMReplicationTime(ctx, at), "retry-owner", nil, opts)
mustIAM(t, err)
}
// A delayed snapshot still has its original absolute expiration. Reapplying
// the local minimum issuance lifetime would leave the old unexpired key alive.
func TestIAMServiceAccountReplicationPreservesExpiration(t *testing.T) {
for _, backend := range []string{"object", "etcd"} {
for _, action := range []string{"create", "update"} {
for _, expired := range []bool{false, true} {
name := backend + "/" + action + "/near_expiry"
if expired {
name = backend + "/" + action + "/expired"
}
t.Run(name, func(t *testing.T) {
ctx, sys, _ := prepareIAMRevisionFixture(t, backend)
origin := UTCNow().Add(-time.Hour)
_, err := sys.CreateUser(ctx, "expiry-owner", madmin.AddOrUpdateUserReq{SecretKey: "valid-owner-password", Status: madmin.AccountEnabled})
mustIAM(t, err)
old := &madmin.SRSvcAccChange{Create: &madmin.SRSvcAccCreate{Parent: "expiry-owner", AccessKey: "expiry-service", SecretKey: "old-service-password"}}
mustIAM(t, globalSiteReplicationSys.PeerSvcAccChangeHandler(ctx, old, origin))
expires := UTCNow().Add(time.Minute)
if expired {
expires = UTCNow().Add(-time.Minute)
}
change := &madmin.SRSvcAccChange{Create: &madmin.SRSvcAccCreate{Parent: "expiry-owner", AccessKey: "expiry-service", SecretKey: "new-service-password", Expiration: &expires}}
if action == "update" {
change = &madmin.SRSvcAccChange{Update: &madmin.SRSvcAccUpdate{AccessKey: "expiry-service", SecretKey: "new-service-password", Expiration: &expires}}
}
mustIAM(t, globalSiteReplicationSys.PeerSvcAccChangeHandler(ctx, change, origin.Add(2*time.Minute)))
u, ok := sys.store.GetUser("expiry-service")
if !ok || u.Credentials.SecretKey != "new-service-password" || !u.Credentials.Expiration.Equal(expires) {
t.Fatal("delayed snapshot lost its new secret or absolute expiration")
}
_, allowed := sys.GetUser(ctx, "expiry-service")
if allowed == expired {
t.Fatal("credential validity disagrees with its absolute expiration")
}
mustIAM(t, sys.store.LoadIAMCache(ctx, false))
mustIAM(t, globalSiteReplicationSys.PeerSvcAccChangeHandler(ctx, old, origin))
r, err := loadIAMRevision(ctx, sys.store, getUserIdentityPath("expiry-service", svcUser))
mustIAM(t, err)
if r.Credentials.SecretKey == old.Create.SecretKey || (expired && !r.Deleted) {
t.Fatal("old non-expiring credential returned after reload")
}
_, _, err = sys.NewServiceAccount(ctx, "expiry-owner", nil, newServiceAccountOpts{accessKey: "local-expiry", secretKey: "valid-service-password", expiration: &expires})
if !errors.Is(err, errInvalidSvcAcctExpiration) {
t.Fatalf("local issuance lifetime check changed: %v", err)
}
})
}
}
}
}
// Status-only summaries intentionally omit secrets. Different revisions must
// still trigger live healing, and disabled identities must be eligible sources.
func TestIAMServiceAccountHealingNewerSnapshot(t *testing.T) {
ctx, sys, obj := prepareIAMRevisionFixture(t)
origin := UTCNow().Add(-time.Hour)
_, err := sys.CreateUser(ctx, "heal-svc-owner", madmin.AddOrUpdateUserReq{SecretKey: "valid-owner-password", Status: madmin.AccountEnabled})
mustIAM(t, err)
_, _, err = sys.NewServiceAccount(withIAMReplicationTime(ctx, origin), "heal-svc-owner", nil, newServiceAccountOpts{accessKey: "heal-service", secretKey: "valid-service-password"})
mustIAM(t, err)
at, err := sys.UpdateServiceAccount(ctx, "heal-service", updateServiceAccountOpts{status: auth.AccountOff})
mustIAM(t, err)
var sent atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/minio/health/live" {
return
}
if r.URL.Path != "/minio/admin/v3/site-replication/peer/iam-revisions" || r.Method != http.MethodPut {
t.Errorf("unexpected request %s %s", r.Method, r.URL.Path)
w.WriteHeader(http.StatusNotFound)
return
}
var batch iamRevisionBatch
if err := json.NewDecoder(r.Body).Decode(&batch); err != nil {
t.Error(err)
w.WriteHeader(http.StatusBadRequest)
return
}
for _, item := range batch.Items {
if item.SvcAccChange != nil && item.SvcAccChange.Create != nil && item.SvcAccChange.Create.AccessKey == "heal-service" && item.SvcAccChange.Create.Status == auth.AccountOff && item.UpdatedAt.Equal(at) {
sent.Add(1)
}
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(iamRevisionResponse{iamRevisionStatus: iamRevisionStatus{Version: iamRevisionProtocol, Node: "node", Instance: "boot", Digest: "fixture"}})
}))
defer server.Close()
peers := map[string]madmin.PeerInfo{globalDeploymentID(): {Name: "local", DeploymentID: globalDeploymentID()}, "remote": {Name: "remote", DeploymentID: "remote", Endpoint: server.URL}}
c := &SiteReplicationSys{enabled: true, state: srState{ServiceAccountAccessKey: "heal-svc-owner", Peers: peers}}
local := madmin.UserInfo{Status: madmin.AccountStatus(auth.AccountOff), UpdatedAt: at}
remote := local
remote.UpdatedAt = origin
if isUserInfoReplicated(2, 2, []madmin.UserInfo{local, remote}) {
t.Fatal("status-only summaries concealed different service revisions")
}
info := srStatusInfo{Sites: peers, UserStats: map[string]map[string]srUserStatsSummary{"heal-service": {globalDeploymentID(): {userInfo: srUserInfo{UserInfo: local}}, "remote": {SRUserStatsSummary: madmin.SRUserStatsSummary{UserInfoMismatch: true}, userInfo: srUserInfo{UserInfo: remote}}}}}
mustIAM(t, c.healUsers(ctx, obj, "heal-service", info))
if sent.Load() == 0 {
t.Fatal("disabled newer service snapshot was not healed")
}
}
+337 -221
View File
File diff suppressed because it is too large Load Diff
+56 -21
View File
@@ -605,10 +605,22 @@ func (sys *IAMSys) DeletePolicy(ctx context.Context, policyName string, notifyPe
return errServerNotInitialized
}
for _, v := range policy.DefaultPolicies {
if v.Name == policyName {
if err := checkConfig(ctx, globalObjectAPI, getPolicyDocPath(policyName)); err != nil && err == errConfigNotFound {
return fmt.Errorf("inbuilt policy `%s` not allowed to be deleted", policyName)
if _, replicated := iamReplicationTime(ctx); !replicated && notifyPeers {
for _, v := range policy.DefaultPolicies {
if v.Name == policyName {
var err error
if objectStore, ok := sys.store.IAMStorageAPI.(*IAMObjectStore); ok {
err = checkConfig(ctx, objectStore.objAPI, getPolicyDocPath(policyName))
} else {
var r iamRevision
err = sys.store.loadIAMConfig(ctx, &r, getPolicyDocPath(policyName))
}
if errors.Is(err, errConfigNotFound) {
return fmt.Errorf("inbuilt policy `%s` not allowed to be deleted", policyName)
}
if err != nil {
return err
}
}
}
}
@@ -714,20 +726,30 @@ func (sys *IAMSys) DeleteUser(ctx context.Context, accessKey string, notifyPeers
return errServerNotInitialized
}
if err := sys.store.DeleteUser(ctx, accessKey, regUser); err != nil {
err := sys.store.DeleteUser(ctx, accessKey, regUser)
var cleanupErr *iamCommittedCleanupError
retained := errors.Is(err, errIAMRevocationRetained)
if errors.As(err, &cleanupErr) {
retained = cleanupErr.retained
} else if err != nil && !retained {
return err
}
// Notify all other MinIO peers to delete user.
// Publish the committed state even when dependent cleanup must be retried.
if notifyPeers && !sys.HasWatcher() {
for _, nerr := range globalNotificationSys.DeleteUser(ctx, accessKey) {
if nerr.Err != nil {
logger.GetReqInfo(ctx).SetTags("peerAddress", nerr.Host.String())
iamLogIf(ctx, nerr.Err)
if retained {
sys.notifyForUser(ctx, accessKey, false)
} else {
for _, nerr := range globalNotificationSys.DeleteUser(ctx, accessKey) {
if nerr.Err != nil {
logger.GetReqInfo(ctx).SetTags("peerAddress", nerr.Host.String())
iamLogIf(ctx, nerr.Err)
}
}
}
}
if cleanupErr != nil {
return cleanupErr
}
return nil
}
@@ -1061,6 +1083,7 @@ type newServiceAccountOpts struct {
sessionPolicy *policy.Policy
accessKey string
secretKey string
status string // Used by replication snapshots; local creates default to enabled.
name, description string
expiration *time.Time
allowSiteReplicatorAccount bool // allow creating internal service account for site-replication.
@@ -1125,6 +1148,11 @@ func (sys *IAMSys) NewServiceAccount(ctx context.Context, parentUser string, gro
m[k] = v
}
}
if _, replicated := iamReplicationTime(ctx); !replicated {
if err := setIAMParentRevocationClaim(ctx, sys.store, parentUser, m); err != nil {
return auth.Credentials{}, time.Time{}, err
}
}
var accessKey, secretKey string
var err error
@@ -1143,12 +1171,19 @@ func (sys *IAMSys) NewServiceAccount(ctx context.Context, parentUser string, gro
cred.ParentUser = parentUser
cred.Groups = groups
cred.Status = string(auth.AccountOn)
switch opts.status {
case "", auth.AccountOn, string(madmin.AccountEnabled):
case auth.AccountOff, string(madmin.AccountDisabled):
cred.Status = auth.AccountOff
default:
return auth.Credentials{}, time.Time{}, errInvalidArgument
}
cred.Name = opts.name
cred.Description = opts.description
if opts.expiration != nil {
expirationInUTC := opts.expiration.UTC()
if err := validateSvcExpirationInUTC(expirationInUTC); err != nil {
if err := validateSvcExpirationInUTC(ctx, expirationInUTC); err != nil {
return auth.Credentials{}, time.Time{}, err
}
cred.Expiration = expirationInUTC
@@ -1379,7 +1414,7 @@ func (sys *IAMSys) DeleteServiceAccount(ctx context.Context, accessKey string, n
}
sa, ok := sys.store.GetUser(accessKey)
if !ok || !sa.Credentials.IsServiceAccount() {
if _, replicated := iamReplicationTime(ctx); (!ok || !sa.Credentials.IsServiceAccount()) && !replicated {
return nil
}
@@ -1483,8 +1518,8 @@ func (sys *IAMSys) purgeExpiredCredentialsForExternalSSO(ctx context.Context) {
}
}
// We ignore any errors
_ = sys.store.DeleteUsers(ctx, expiredUsers)
// Keep failed revocations visible so the next purge can retry.
iamLogIf(ctx, sys.store.DeleteUsers(ctx, expiredUsers))
}
// purgeExpiredCredentialsForLDAP - validates if local credentials are still
@@ -1512,8 +1547,8 @@ func (sys *IAMSys) purgeExpiredCredentialsForLDAP(ctx context.Context) {
return
}
// We ignore any errors
_ = sys.store.DeleteUsers(ctx, expiredUsers)
// Keep failed revocations visible so the next purge can retry.
iamLogIf(ctx, sys.store.DeleteUsers(ctx, expiredUsers))
}
// updateGroupMembershipsForLDAP - updates the list of groups associated with the credential.
@@ -1934,12 +1969,12 @@ func (sys *IAMSys) RemoveUsersFromGroup(ctx context.Context, group string, membe
}
updatedAt, err = sys.store.RemoveUsersFromGroup(ctx, group, members)
if err != nil {
var cleanupErr *iamCommittedCleanupError
if err != nil && !errors.As(err, &cleanupErr) {
return updatedAt, err
}
sys.notifyForGroup(ctx, group)
return updatedAt, nil
return updatedAt, err
}
// SetGroupStatus - enable/disabled a group
+14
View File
@@ -34,6 +34,10 @@ const (
sinceLastSyncMillis = "since_last_sync_millis"
syncFailures = "sync_failures"
syncSuccesses = "sync_successes"
revocationRecords = "revocation_records"
revocationHealFailures = "revocation_heal_failures"
revocationHealDurationMillis = "revocation_heal_duration_millis"
revocationHealLastSuccess = "revocation_heal_last_success_timestamp_seconds"
)
var (
@@ -47,10 +51,20 @@ var (
sinceLastSyncMillisMD = NewCounterMD(sinceLastSyncMillis, "Time (in milliseconds) since last successful IAM data sync.")
syncFailuresMD = NewCounterMD(syncFailures, "Number of failed IAM data syncs since server start.")
syncSuccessesMD = NewCounterMD(syncSuccesses, "Number of successful IAM data syncs since server start.")
revocationRecordsMD = NewGaugeMD(revocationRecords, "Retained IAM deletion records and revocation boundaries in this node's index.")
revocationHealFailuresMD = NewCounterMD(revocationHealFailures, "Failed IAM revocation convergence passes since server start.")
revocationHealDurationMillisMD = NewGaugeMD(revocationHealDurationMillis, "Duration of the last IAM revocation convergence pass in milliseconds.")
revocationHealLastSuccessMD = NewGaugeMD(revocationHealLastSuccess, "Unix timestamp of the last successful IAM revocation convergence pass.")
)
// loadClusterIAMMetrics - `MetricsLoaderFn` for cluster IAM metrics.
func loadClusterIAMMetrics(_ context.Context, m MetricValues, _ *metricsCache) error {
if globalIAMSys.Initialized() {
m.Set(revocationRecords, float64(globalIAMSys.store.revisionIndex().count()))
}
m.Set(revocationHealFailures, float64(globalSiteReplicationSys.iamRevisionMetrics.healFailures.Load()))
m.Set(revocationHealDurationMillis, float64(globalSiteReplicationSys.iamRevisionMetrics.healDurationMillis.Load()))
m.Set(revocationHealLastSuccess, float64(globalSiteReplicationSys.iamRevisionMetrics.healLastSuccess.Load()))
m.Set(lastSyncDurationMillis, float64(atomic.LoadUint64(&globalIAMSys.LastRefreshDurationMilliseconds)))
pluginAuthNMetrics := globalAuthNPlugin.Metrics()
m.Set(pluginAuthnServiceFailedRequestsMinute, float64(pluginAuthNMetrics.FailedRequests))
+4
View File
@@ -323,6 +323,10 @@ func newMetricGroups(r *prometheus.Registry) *metricsV3Collection {
sinceLastSyncMillisMD,
syncFailuresMD,
syncSuccessesMD,
revocationRecordsMD,
revocationHealFailuresMD,
revocationHealDurationMillisMD,
revocationHealLastSuccessMD,
},
loadClusterIAMMetrics,
)
+113
View File
@@ -0,0 +1,113 @@
// 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 (
"encoding/base64"
"net/http"
"reflect"
"strconv"
"strings"
"testing"
"time"
xhttp "github.com/minio/minio/internal/http"
)
func TestPutOptsFromHeadersReplicationTimestamps(t *testing.T) {
stamp := time.Date(2026, 9, 15, 1, 2, 3, 123456789, time.UTC)
context := base64.StdEncoding.EncodeToString([]byte(`{"purpose":"tag-replication"}`))
for _, encryption := range []struct {
name string
headers map[string]string
}{
{name: "none"},
{name: "SSE-S3", headers: map[string]string{xhttp.AmzServerSideEncryption: xhttp.AmzEncryptionAES}},
{name: "SSE-KMS", headers: map[string]string{xhttp.AmzServerSideEncryption: xhttp.AmzEncryptionKMS}},
{name: "SSE-KMS-context", headers: map[string]string{
xhttp.AmzServerSideEncryption: xhttp.AmzEncryptionKMS, xhttp.AmzServerSideEncryptionKmsID: "tag-replication-key",
xhttp.AmzServerSideEncryptionKmsContext: context,
}},
{name: "SSE-C", headers: ssecKeyHeaders([]byte("01234567890123456789012345678901"), false)},
} {
t.Run(encryption.name, func(t *testing.T) {
for _, trusted := range []bool{false, true} {
t.Run("trusted="+strconv.FormatBool(trusted), func(t *testing.T) {
for _, tagging := range []struct {
name, header string
want time.Time
invalid bool
}{
{name: "absent"},
{name: "nanoseconds", header: stamp.Format(time.RFC3339Nano), want: stamp},
{name: "offset-whitespace", header: " " + stamp.In(time.FixedZone("UTC+8", 8*60*60)).Format(time.RFC3339Nano) + " ", want: stamp},
{name: "invalid", header: "not-a-timestamp", invalid: true},
} {
t.Run(tagging.name, func(t *testing.T) {
for _, metadata := range []map[string]string{nil, {"x-amz-meta-test": "kept"}} {
hdr := make(http.Header)
wantEncryption := make(http.Header)
for key, value := range encryption.headers {
hdr.Set(key, value)
wantEncryption.Set(key, value)
}
hdr.Set(xhttp.MinIOSourceTaggingTimestamp, tagging.header)
hdr.Set(xhttp.MinIOSourceMTime, stamp.Add(-time.Hour).Format(time.RFC3339Nano))
hdr.Set(xhttp.MinIOSourceObjectRetentionTimestamp, stamp.Add(-time.Minute).Format(time.RFC3339Nano))
hdr.Set(xhttp.MinIOSourceObjectLegalHoldTimestamp, stamp.Add(-time.Second).Format(time.RFC3339Nano))
hdr.Set(xhttp.MinIOSourceETag, "source-etag")
opts, err := putOptsFromHeaders(t.Context(), hdr, metadata, trusted)
if trusted && tagging.invalid {
if err == nil || !strings.Contains(err.Error(), xhttp.MinIOSourceTaggingTimestamp) {
t.Fatalf("malformed trusted timestamp: got %v", err)
}
continue
}
if err != nil {
t.Fatal(err)
}
wantTag, wantMTime, wantRetention, wantLegalhold, wantETag := time.Time{}, time.Time{}, time.Time{}, time.Time{}, ""
if trusted {
wantTag, wantMTime = tagging.want, stamp.Add(-time.Hour)
wantRetention, wantLegalhold, wantETag = stamp.Add(-time.Minute), stamp.Add(-time.Second), "source-etag"
}
if !opts.ReplicationSourceTaggingTimestamp.Equal(wantTag) {
t.Errorf("tag timestamp=%s, want %s", opts.ReplicationSourceTaggingTimestamp, wantTag)
}
if !opts.MTime.Equal(wantMTime) || !opts.ReplicationSourceRetentionTimestamp.Equal(wantRetention) ||
!opts.ReplicationSourceLegalholdTimestamp.Equal(wantLegalhold) || opts.PreserveETag != wantETag || opts.ReplicationRequest != trusted {
t.Error("other source fields did not preserve the replication trust boundary")
}
if opts.UserDefined == nil || (metadata != nil && !reflect.DeepEqual(opts.UserDefined, metadata)) {
t.Errorf("metadata=%v, want nonnil map preserving %v", opts.UserDefined, metadata)
}
gotEncryption := make(http.Header)
if opts.ServerSideEncryption != nil {
opts.ServerSideEncryption.Marshal(gotEncryption)
}
if !reflect.DeepEqual(gotEncryption, wantEncryption) {
t.Errorf("SSE headers=%v, want %v", gotEncryption, wantEncryption)
}
}
})
}
})
}
})
}
}
+3 -3
View File
@@ -452,11 +452,11 @@ func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[strin
MTime: mtime,
PreserveETag: etag,
ReplicationRequest: trustedReplication,
// The Object Lock timestamps order replicated retention and legal
// hold updates. Dropping them here would leave every update on an
// SSE-KMS destination unordered.
// These timestamps order replicated retention, legal hold and tagging
// updates on an SSE-KMS destination.
ReplicationSourceLegalholdTimestamp: lholdtimestmp,
ReplicationSourceRetentionTimestamp: retaintimestmp,
ReplicationSourceTaggingTimestamp: taggingtimestmp,
}
return op, nil
}
+143
View File
@@ -0,0 +1,143 @@
// 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"
"net/http"
"net/http/httptest"
"testing"
"time"
"github.com/minio/minio/internal/auth"
"github.com/minio/minio/internal/crypto"
xhttp "github.com/minio/minio/internal/http"
"github.com/minio/minio/internal/kms"
)
// TestAPICopyObjectReplicaTaggingTimestampUnderKMS covers signed replica COPY
// requests through encryption, metadata replacement and disk persistence. Both
// the single-disk and 16-disk fixtures are single-pool backends.
func TestAPICopyObjectReplicaTaggingTimestampUnderKMS(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: testAPICopyObjectReplicaTaggingTimestampUnderKMS})
}
func testAPICopyObjectReplicaTaggingTimestampUnderKMS(obj ObjectLayer, instance, bucket string, router http.Handler, creds auth.Credentials, t *testing.T) {
// Ignore the host free-space percentage while retaining real disk I/O.
for _, pool := range obj.(*erasureServerPools).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 }
defer func() { set.getDisks = original }()
}
}
oldKMS, oldAuto := GlobalKMS, globalAutoEncryption
GlobalKMS = kms.NewStub("replica-tags-key")
globalAutoEncryption = false
defer func() { GlobalKMS, globalAutoEncryption = oldKMS, oldAuto }()
if _, err := globalBucketMetadataSys.Update(t.Context(), bucket, bucketVersioningConfig, enabledBucketVersioningConfig); err != nil {
t.Fatal(err)
}
const tsKey = ReservedMetadataPrefixLower + TaggingTimestamp
stamp := time.Date(2026, 9, 15, 1, 0, 0, 123456789, time.UTC)
for _, mode := range []string{"none", "explicit-sse-s3", "explicit-kms", "auto-kms", "bucket-kms"} {
t.Run(instance+"/"+mode, func(t *testing.T) {
globalAutoEncryption = mode == "auto-kms"
if mode == "bucket-kms" {
sseXML := []byte(`<ServerSideEncryptionConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ApplyServerSideEncryptionByDefault><SSEAlgorithm>aws:kms</SSEAlgorithm><KMSMasterKeyID>replica-tags-key</KMSMasterKeyID></ApplyServerSideEncryptionByDefault></Rule></ServerSideEncryptionConfiguration>`)
if _, err := globalBucketMetadataSys.Update(t.Context(), bucket, bucketSSEConfig, sseXML); err != nil {
t.Fatal(err)
}
}
const data = "encrypted replica copy remains readable"
oi, err := obj.PutObject(t.Context(), bucket, mode, mustGetPutObjReader(t, bytes.NewReader([]byte(data)), int64(len(data)), "", ""), ObjectOptions{
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "key=old", tsKey: stamp.Format(time.RFC3339Nano)},
})
if err != nil {
t.Fatal(err)
}
for _, event := range []struct {
name, tags, wantTags string
delta, wantDelta time.Duration
missingTimestamp bool
}{
{"newer", "key=new", "key=new", 2, 2, false},
{"stale", "key=stale", "key=new", 1, 2, false},
{"duplicate", "key=new", "key=new", 2, 2, false},
{"newer-again", "key=latest", "key=latest", 3, 3, false},
{"missing-timestamp", "key=unordered", "key=latest", 0, 3, true},
} {
headers := map[string]string{
xhttp.AmzCopySource: "/" + bucket + "/" + mode + "?versionId=" + oi.VersionID,
xhttp.AmzMetadataDirective: "REPLACE", xhttp.AmzTagDirective: "REPLACE",
xhttp.AmzObjectTagging: event.tags, xhttp.MinIOSourceReplicationRequest: "true",
xhttp.AmzBucketReplicationStatus: "REPLICA", xhttp.MinIOSourceTaggingTimestamp: stamp.Add(event.delta).Format(time.RFC3339Nano),
xhttp.MinIOSourceMTime: oi.ModTime.Format(time.RFC3339Nano), xhttp.MinIOSourceETag: oi.ETag,
}
if event.missingTimestamp {
delete(headers, xhttp.MinIOSourceTaggingTimestamp)
}
if mode == "explicit-sse-s3" {
headers[xhttp.AmzServerSideEncryption] = xhttp.AmzEncryptionAES
}
if mode == "explicit-kms" {
headers[xhttp.AmzServerSideEncryption] = "aws:kms"
headers[xhttp.AmzServerSideEncryptionKmsID] = "replica-tags-key"
}
req, err := newTestSignedRequestV4(http.MethodPut, "/"+bucket+"/"+mode+"?versionId="+oi.VersionID, 0, nil, creds.AccessKey, creds.SecretKey, headers)
if err != nil {
t.Fatal(err)
}
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if w.Code != http.StatusOK {
t.Fatalf("%s: COPY %d %s", event.name, w.Code, w.Body.String())
}
got, err := obj.GetObjectInfo(t.Context(), bucket, mode, ObjectOptions{VersionID: oi.VersionID})
if err != nil {
t.Fatal(err)
}
t.Logf("%s: tags=%q timestamp=%q kms=%v", event.name, got.UserTags, got.UserDefined[tsKey], crypto.S3KMS.IsEncrypted(got.UserDefined))
if got.UserTags != event.wantTags || got.UserDefined[tsKey] != stamp.Add(event.wantDelta).Format(time.RFC3339Nano) {
t.Errorf("%s: incorrect persisted tags/timestamp", event.name)
}
wantKMS := mode == "explicit-kms" || mode == "auto-kms" || mode == "bucket-kms"
if crypto.S3KMS.IsEncrypted(got.UserDefined) != wantKMS || crypto.S3.IsEncrypted(got.UserDefined) != (mode == "explicit-sse-s3") {
t.Errorf("%s: unexpected destination encryption", event.name)
}
if got.VersionID != oi.VersionID {
t.Errorf("%s: version=%q, want %q", event.name, got.VersionID, oi.VersionID)
}
req, err = newTestSignedRequestV4(http.MethodGet, "/"+bucket+"/"+mode+"?versionId="+oi.VersionID, 0, nil, creds.AccessKey, creds.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
w = httptest.NewRecorder()
router.ServeHTTP(w, req)
if w.Code != http.StatusOK || w.Body.String() != data {
t.Fatalf("%s: GET %d %q", event.name, w.Code, w.Body.String())
}
}
})
}
}
+6 -2
View File
@@ -204,7 +204,9 @@ func (s *peerRESTServer) DeleteServiceAccountHandler(mss *grid.MSS) (np grid.NoP
return np, grid.NewRemoteErr(errors.New("service account name is missing"))
}
if err := globalIAMSys.DeleteServiceAccount(context.Background(), accessKey, false); err != nil {
ctx, cancel := context.WithTimeout(GlobalContext, defaultContextTimeout)
defer cancel()
if err := globalIAMSys.LoadServiceAccount(ctx, accessKey); err != nil {
return np, grid.NewRemoteErr(err)
}
@@ -274,7 +276,9 @@ func (s *peerRESTServer) LoadUserHandler(mss *grid.MSS) (np grid.NoPayload, nerr
userType = stsUser
}
if err = globalIAMSys.LoadUser(context.Background(), objAPI, accessKey, userType); err != nil {
ctx, cancel := context.WithTimeout(GlobalContext, defaultContextTimeout)
defer cancel()
if err = globalIAMSys.LoadUser(ctx, objAPI, accessKey, userType); err != nil {
return np, grid.NewRemoteErr(err)
}
@@ -274,6 +274,10 @@ func TestBucketMetadataInitialSyncPhysicalCreated(t *testing.T) {
}
events = append(events, event)
}
if r.URL.Path == "/minio/admin/v3/site-replication/peer/iam-revisions" {
_ = json.NewEncoder(w).Encode(iamRevisionResponse{iamRevisionStatus: iamRevisionStatus{Version: iamRevisionProtocol, Node: "initial-peer", Instance: "initial-boot", Digest: "ack"}})
return
}
w.WriteHeader(http.StatusOK)
}))
defer peer.Close()
+87 -34
View File
@@ -28,6 +28,7 @@ import (
"fmt"
"maps"
"math/rand"
"net/http"
"net/url"
"reflect"
"runtime"
@@ -209,7 +210,11 @@ type SiteReplicationSys struct {
// In-memory and persisted multi-site replication state.
state srState
iamMetaCache srIAMCache
iamMetaCache srIAMCache
healOnce sync.Once // Configuration reloads must not spawn more healing loops.
iamHealMu sync.Mutex
iamRevisionProgress map[string]iamRevisionProgress
iamRevisionMetrics iamRevisionMetrics
}
type srState srStateV1
@@ -234,7 +239,7 @@ type srStateData struct {
// Init - initialize the site replication manager.
func (c *SiteReplicationSys) Init(ctx context.Context, objAPI ObjectLayer) error {
go c.startHealRoutine(ctx, objAPI)
c.healOnce.Do(func() { go c.startHealRoutine(ctx, objAPI) })
r := rand.New(rand.NewSource(time.Now().UnixNano()))
for {
err := c.loadFromDisk(ctx, objAPI)
@@ -1257,13 +1262,36 @@ func (c *SiteReplicationSys) IAMChangeHook(ctx context.Context, item madmin.SRIA
return nil
}
if path := iamDeletionPath(item); path != "" {
r, err := loadIAMRevision(ctx, globalIAMSys.store, path)
if err != nil {
return err
}
if !r.Deleted && ((item.Type != madmin.SRIAMItemIAMUser && item.Type != madmin.SRIAMItemGroupInfo) || r.RevokedBefore.IsZero()) {
// A concurrent recreation has already superseded this delete.
return nil
}
item.UpdatedAt = r.timestamp()
if !r.Deleted {
item.UpdatedAt = r.RevokedBefore
}
}
versioned, err := c.replicationItem(ctx, item)
if errors.Is(err, errIAMStaleUpdate) {
return nil
}
if err != nil {
return err
}
cerr := c.concDo(nil, func(d string, p madmin.PeerInfo) error {
admClient, err := c.getAdminClient(ctx, d)
if err != nil {
return wrapSRErr(err)
}
return c.annotatePeerErr(p.Name, replicateIAMItem, admClient.SRPeerReplicateIAMItem(ctx, item))
_, err = executeIAMRevisionRequest(ctx, admClient, http.MethodPut, &iamRevisionBatch{Version: iamRevisionProtocol, Items: []iamReplicationItem{versioned}})
return c.annotatePeerErr(p.Name, replicateIAMItem, err)
},
replicateIAMItem,
)
@@ -1273,6 +1301,7 @@ func (c *SiteReplicationSys) IAMChangeHook(ctx context.Context, item madmin.SRIA
// PeerAddPolicyHandler - copies IAM policy to local. A nil policy argument,
// causes the named policy to be deleted.
func (c *SiteReplicationSys) PeerAddPolicyHandler(ctx context.Context, policyName string, p *policy.Policy, updatedAt time.Time) error {
ctx = withIAMReplicationTime(ctx, updatedAt)
var err error
// skip overwrite of local update if peer sent stale info
if !updatedAt.IsZero() {
@@ -1286,18 +1315,19 @@ func (c *SiteReplicationSys) PeerAddPolicyHandler(ctx context.Context, policyNam
_, err = globalIAMSys.SetPolicy(ctx, policyName, *p)
}
if err != nil {
return wrapSRErr(err)
return iamReplicationError(err)
}
return nil
}
// PeerIAMUserChangeHandler - copies IAM user to local.
func (c *SiteReplicationSys) PeerIAMUserChangeHandler(ctx context.Context, change *madmin.SRIAMUser, updatedAt time.Time) error {
ctx = withIAMReplicationTime(ctx, updatedAt)
if change == nil {
return errSRInvalidRequest(errInvalidArgument)
}
// skip overwrite of local update if peer sent stale info
if !updatedAt.IsZero() {
if !change.IsDeleteReq && !updatedAt.IsZero() {
if ui, err := globalIAMSys.GetUserInfo(ctx, change.AccessKey); err == nil && ui.UpdatedAt.After(updatedAt) {
return nil
}
@@ -1326,13 +1356,14 @@ func (c *SiteReplicationSys) PeerIAMUserChangeHandler(ctx context.Context, chang
}
}
if err != nil {
return wrapSRErr(err)
return iamReplicationError(err)
}
return nil
}
// PeerGroupInfoChangeHandler - copies group changes to local.
func (c *SiteReplicationSys) PeerGroupInfoChangeHandler(ctx context.Context, change *madmin.SRGroupInfo, updatedAt time.Time) error {
ctx = withIAMReplicationTime(ctx, updatedAt)
if change == nil {
return errSRInvalidRequest(errInvalidArgument)
}
@@ -1340,7 +1371,7 @@ func (c *SiteReplicationSys) PeerGroupInfoChangeHandler(ctx context.Context, cha
var err error
// skip overwrite of local update if peer sent stale info
if !updatedAt.IsZero() {
if !updatedAt.IsZero() && (!updReq.IsRemove || len(updReq.Members) != 0) {
if gd, err := globalIAMSys.GetGroupDescription(updReq.Group); err == nil && gd.UpdatedAt.After(updatedAt) {
return nil
}
@@ -1349,7 +1380,8 @@ func (c *SiteReplicationSys) PeerGroupInfoChangeHandler(ctx context.Context, cha
if updReq.IsRemove {
_, err = globalIAMSys.RemoveUsersFromGroup(ctx, updReq.Group, updReq.Members)
} else {
if updReq.Status != "" && len(updReq.Members) == 0 {
snapshot, _ := ctx.Value(iamGroupSnapshotKey{}).(bool)
if !snapshot && updReq.Status != "" && len(updReq.Members) == 0 {
_, err = globalIAMSys.SetGroupStatus(ctx, updReq.Group, updReq.Status == madmin.GroupEnabled)
} else {
if globalIAMSys.LDAPConfig.Enabled() {
@@ -1365,13 +1397,14 @@ func (c *SiteReplicationSys) PeerGroupInfoChangeHandler(ctx context.Context, cha
}
}
if err != nil && !errors.Is(err, errNoSuchGroup) {
return wrapSRErr(err)
return iamReplicationError(err)
}
return nil
}
// PeerSvcAccChangeHandler - copies service-account change to local.
func (c *SiteReplicationSys) PeerSvcAccChangeHandler(ctx context.Context, change *madmin.SRSvcAccChange, updatedAt time.Time) error {
ctx = withIAMReplicationTime(ctx, updatedAt)
if change == nil {
return errSRInvalidRequest(errInvalidArgument)
}
@@ -1382,7 +1415,7 @@ func (c *SiteReplicationSys) PeerSvcAccChangeHandler(ctx context.Context, change
if len(change.Create.SessionPolicy) > 0 {
sp, err = policy.ParseConfig(bytes.NewReader(change.Create.SessionPolicy))
if err != nil {
return wrapSRErr(err)
return iamReplicationError(err)
}
}
// skip overwrite of local update if peer sent stale info
@@ -1394,6 +1427,7 @@ func (c *SiteReplicationSys) PeerSvcAccChangeHandler(ctx context.Context, change
opts := newServiceAccountOpts{
accessKey: change.Create.AccessKey,
secretKey: change.Create.SecretKey,
status: change.Create.Status,
sessionPolicy: sp,
claims: change.Create.Claims,
name: change.Create.Name,
@@ -1402,7 +1436,7 @@ func (c *SiteReplicationSys) PeerSvcAccChangeHandler(ctx context.Context, change
}
_, _, err = globalIAMSys.NewServiceAccount(ctx, change.Create.Parent, change.Create.Groups, opts)
if err != nil {
return wrapSRErr(err)
return iamReplicationError(err)
}
case change.Update != nil:
@@ -1411,7 +1445,7 @@ func (c *SiteReplicationSys) PeerSvcAccChangeHandler(ctx context.Context, change
if len(change.Update.SessionPolicy) > 0 {
sp, err = policy.ParseConfig(bytes.NewReader(change.Update.SessionPolicy))
if err != nil {
return wrapSRErr(err)
return iamReplicationError(err)
}
}
// skip overwrite of local update if peer sent stale info
@@ -1431,7 +1465,7 @@ func (c *SiteReplicationSys) PeerSvcAccChangeHandler(ctx context.Context, change
_, err = globalIAMSys.UpdateServiceAccount(ctx, change.Update.AccessKey, opts)
if err != nil {
return wrapSRErr(err)
return iamReplicationError(err)
}
case change.Delete != nil:
@@ -1442,7 +1476,7 @@ func (c *SiteReplicationSys) PeerSvcAccChangeHandler(ctx context.Context, change
}
}
if err := globalIAMSys.DeleteServiceAccount(ctx, change.Delete.AccessKey, true); err != nil {
return wrapSRErr(err)
return iamReplicationError(err)
}
}
@@ -1451,24 +1485,17 @@ func (c *SiteReplicationSys) PeerSvcAccChangeHandler(ctx context.Context, change
// PeerPolicyMappingHandler - copies policy mapping to local.
func (c *SiteReplicationSys) PeerPolicyMappingHandler(ctx context.Context, mapping *madmin.SRPolicyMapping, updatedAt time.Time) error {
ctx = withIAMReplicationTime(ctx, updatedAt)
if mapping == nil {
return errSRInvalidRequest(errInvalidArgument)
}
// skip overwrite of local update if peer sent stale info
if !updatedAt.IsZero() {
mp, ok := globalIAMSys.store.GetMappedPolicy(mapping.Policy, mapping.IsGroup)
if ok && mp.UpdatedAt.After(updatedAt) {
return nil
}
}
// When LDAP is enabled, we verify that the user or group exists in LDAP and
// use the normalized form of the entityName (which will be an LDAP DN).
userType := IAMUserType(mapping.UserType)
isGroup := mapping.IsGroup
entityName := mapping.UserOrGroup
if globalIAMSys.GetUsersSysType() == LDAPUsersSysType && userType == stsUser {
if mapping.Policy != "" && globalIAMSys.GetUsersSysType() == LDAPUsersSysType && userType == stsUser {
// Validate that the user or group exists in LDAP and use the normalized
// form of the entityName (which will be an LDAP DN).
var err error
@@ -1491,19 +1518,20 @@ func (c *SiteReplicationSys) PeerPolicyMappingHandler(ctx context.Context, mappi
entityName = foundUserDN.NormDN
}
if err != nil {
return wrapSRErr(err)
return iamReplicationError(err)
}
}
_, err := globalIAMSys.PolicyDBSet(ctx, entityName, mapping.Policy, userType, isGroup)
if err != nil {
return wrapSRErr(err)
return iamReplicationError(err)
}
return nil
}
// PeerSTSAccHandler - replicates STS credential locally.
func (c *SiteReplicationSys) PeerSTSAccHandler(ctx context.Context, stsCred *madmin.SRSTSCredential, updatedAt time.Time) error {
ctx = withIAMReplicationTime(ctx, updatedAt)
if stsCred == nil {
return errSRInvalidRequest(errInvalidArgument)
}
@@ -1555,7 +1583,7 @@ func (c *SiteReplicationSys) PeerSTSAccHandler(ctx context.Context, stsCred *mad
// Set these credentials to IAM.
if _, err := globalIAMSys.SetTempUser(ctx, cred.AccessKey, cred, stsCred.ParentPolicyMapping); err != nil {
return fmt.Errorf("unable to save STS credential and/or parent policy mapping: %w", err)
return iamReplicationError(fmt.Errorf("unable to save STS credential and/or parent policy mapping: %w", err))
}
return nil
@@ -4193,7 +4221,8 @@ func (c *SiteReplicationSys) SiteReplicationMetaInfo(ctx context.Context, objAPI
}
info.UserInfoMap[k] = madmin.UserInfo{
Status: madmin.AccountStatus(v.Credentials.Status),
Status: madmin.AccountStatus(v.Credentials.Status),
UpdatedAt: v.UpdatedAt,
}
}
}
@@ -4549,9 +4578,25 @@ func (c *SiteReplicationSys) PeerStateEditReq(ctx context.Context, arg madmin.SR
const siteHealTimeInterval = 30 * time.Second
func (c *SiteReplicationSys) startHealRoutine(ctx context.Context, objAPI ObjectLayer) {
ctx, cancel := globalLeaderLock.GetLock(ctx)
defer cancel()
for ctx.Err() == nil {
var leadership LockContext
select {
case <-ctx.Done():
return
case leadership = <-globalLeaderLock.lockContext:
}
if leadership.Context().Err() != nil {
continue
}
leaderCtx, cancel := mergeContext(leadership.Context(), ctx)
c.healWithLeadership(leaderCtx, objAPI)
cancel()
// Quorum loss cancels a leadership lease, not the subsystem. Wait
// for a new lease so revocations can still reach offline sites.
}
}
func (c *SiteReplicationSys) healWithLeadership(ctx context.Context, objAPI ObjectLayer) {
healTimer := time.NewTimer(siteHealTimeInterval)
defer healTimer.Stop()
@@ -5170,13 +5215,16 @@ func (c *SiteReplicationSys) healBucketReplicationConfig(ctx context.Context, ob
}
func (c *SiteReplicationSys) healIAMSystem(ctx context.Context, objAPI ObjectLayer) error {
// A peer rejecting a deletion must not stop unrelated live updates from
// reaching healthy peers. Retain and report the error for the next retry.
deletionErr := c.healIAMDeletions(ctx)
info, err := c.siteReplicationStatus(ctx, objAPI, madmin.SRStatusOptions{
Users: true,
Policies: true,
Groups: true,
})
if err != nil {
return err
return errors.Join(deletionErr, err)
}
for policy := range info.PolicyStats {
c.healPolicies(ctx, objAPI, policy, info)
@@ -5194,7 +5242,7 @@ func (c *SiteReplicationSys) healIAMSystem(ctx context.Context, objAPI ObjectLay
c.healGroupPolicies(ctx, objAPI, group, info)
}
return nil
return deletionErr
}
// heal iam policies present on this site to peers, provided current cluster has the most recent update.
@@ -5429,8 +5477,10 @@ func (c *SiteReplicationSys) healUsers(ctx context.Context, objAPI ObjectLayer,
peerName := info.Sites[dID].Name
u, ok := globalIAMSys.GetUser(ctx, user)
if !ok {
// Disabled identities are valid replication sources. CheckKey returns
// their stored record even though authentication is denied.
u, _, err := globalIAMSys.CheckKey(ctx, user)
if err != nil || u.Credentials.AccessKey == "" || u.Credentials.IsExpired() {
continue
}
creds := u.Credentials
@@ -5629,7 +5679,10 @@ func isGroupDescEqual(g1, g2 madmin.GroupDesc) bool {
}
func isUserInfoEqual(u1, u2 madmin.UserInfo) bool {
if u1.PolicyName != u2.PolicyName ||
// Full-site summaries omit secrets and claims. Equal status alone cannot
// distinguish a recreated identity or an edited service-account policy.
if !u1.UpdatedAt.Equal(u2.UpdatedAt) ||
u1.PolicyName != u2.PolicyName ||
u1.Status != u2.Status ||
u1.SecretKey != u2.SecretKey {
return false
+4
View File
@@ -624,6 +624,10 @@ func (sts *stsAPIHandlers) AssumeRole(w http.ResponseWriter, r *http.Request) {
claims[expClaim] = UTCNow().Add(duration).Unix()
claims[parentClaim] = user.AccessKey
if err := setIAMParentRevocationClaim(ctx, globalIAMSys.store, user.AccessKey, claims); err != nil {
writeSTSErrorResponse(ctx, w, ErrSTSInternalError, err)
return
}
tokenRevokeType := r.Form.Get(stsRevokeTokenType)
if tokenRevokeType != "" {
+33
View File
@@ -0,0 +1,33 @@
# R4 plan consensus and review disposition
## Agreed version
- Plan: [plan v1](plan-v1.md), SHA-256 `ad539f2071155de6955b583991684ed33c4bfe2e29660005840cdc97d7e1a754`. The frozen file remains unchanged.
- Baseline: `9ebe81c1b3611f9cc73e676b5b741c2be62c467a`.
- Actual reviewer: Claude Code 2.1.270, every assistant model in the review stream is `claude-opus-5`; explicit `--effort max`.
- Opus: **GO_WITH_NONBLOCKING_NOTES**, zero blockers; explicitly agrees that this exact plan can enter local implementation. [Unedited returned review](opus-v1-review.md), [machine-readable provenance](opus-v1.metadata.json).
- Codex: agrees that adding the already-parsed timestamp to the KMS literal fixes R4, and accepts the nonblocking dispositions below. **No blocking disagreement remains on plan v1.** No production source edits were made before this record was saved.
- The agreement permits the planned local implementation and tests; it is not implementation acceptance, a merge decision or production release approval.
## Item-by-item disposition
| Opus ID | Disposition |
|---|---|
| R4-01 | Accepted citation correction here, leaving the agreed hash frozen: `ReplicaLockReconcile` is at baseline `object-handlers.go:1847`; encryption merge is at `:1903`. |
| R4-02 | Accepted scope clarification: ErasureSD and Erasure16 are both single-pool local backends. KMS rewrites use PutObject under-lock reconciliation. Multi-pool and multi-site validation are optional and deferred to the wider integration gate. Test comments and the final report will identify this boundary. |
| R4-03 | Accepted wording clarification: source encryption alone does not request destination encryption. Source-only SSE-C copy headers do not prevent destination bucket/default auto-KMS from selecting KMS. The three destination trigger categories stay unchanged. |
| R4-04 | Accepted intent. The regression matrix uses identical expected mtime, ETag, trust and all three source timestamps across all encryption modes, giving field-by-field equivalence without constructing expected values through the production function. The temporary expanded baseline matrix fails only trusted valid KMS tag timestamps. |
| R4-05 | Accepted optional test within the existing scope: a signed KMS COPY with nonempty tags and no source tag timestamp must preserve the stored value/time. This adds evidence, not production behavior. |
| R4-06 | Registered as a separate unverified-impact finding: KMS construction also omits `ProxyHeaderSet`, `ProxyRequest`, `Speedtest` relative to `getDefaultOpts`. No R4 fix or correctness claim for those flags. Send the observation to the parent for separate triage; do not assign it to R5. |
| R4-07 | Accepted. Assertions target final disk state; the REPLACE handler rebuilds metadata, while final stored-tag rejection occurs under the storage write lock. HTTP 200 alone is not acceptance. |
| R4-08 | Resolved provenance uncertainty by Codex: SHA-256 recomputed before/after review, baseline identity and current GitHub main/PR query captured in `baseline-identity.txt`. History was inspected locally with `git blame` / `git show`. Opus's read-only tools did not independently recompute the hash or check GitHub; those facts remain attributed to the local commands. |
## Raw evidence
Directory: `/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/`.
- `review-prompt-v1.md`, `opus-review-v1.jsonl`, `opus-review-v1.stderr.log`, `opus-review-v1.exit`.
- `baseline-identity.txt`, `r4_repro_test.go`, `overlay.json`, `baseline-repro.log`.
- `options_repro_test.go`, `options-overlay.json`, `baseline-options.log`.
The stream includes an attempted Write to Claude's own plan file. Its tool was disabled; the reviewer returned the full result in text and did not edit production source. The successful result and actual assistant models are checked separately from rate-limit status and auxiliary-model usage.
@@ -0,0 +1,39 @@
已完成独立复核:读了三份改动文件在 HEAD 的实际内容、`putOptsFromHeaders` 全函数与两条返回路径、`CopyObjectHandler` 的默认加密与标签排序段、`reconcileStoredObjectTags` 的三处生产调用点,以及全部指定日志与证据文件。
# R4 实现复核结论
**Verdict: GO_WITH_NONBLOCKING_NOTES(0 阻断项)**
- 复核 HEAD:`dbcf8dec589deb5d91e17d295cb70997635f5b55`
- 代码/测试 diff SHA-256(按提供值记录):`c8cd6648f8ecea835ec74a038cdeaa82acaa3f36250395f97ead3260dc2fc0a5`
- 本会话无 shell,未重算该哈希;改为逐行比对 diff 与工作树三份文件,内容完全一致(`cmd/object-api-options.go`、`cmd/object-api-options-replication_test.go`、`cmd/object-copy-replication-tagging_test.go`)。
## 核验到的事实
- 生产改动确为一个字段 + 相邻注释:`cmd/object-api-options.go:459` 的 `ReplicationSourceTaggingTimestamp: taggingtimestmp`,变量来自 `:419-425` 已解析值,与非 KMS 路径 `:473` 对齐。未动解析、信任判定、KMS key/context、返回结构。
- 影响面封闭:全仓该字段唯一消费点是 `cmd/object-handlers.go:1820`(COPY 标签排序)。PUT/POST/multipart 虽同经 `putOptsFromReq`,但无消费者,故不可能回归——与 R4/R5 切分一致。
- 三条 KMS 触发路径真实可达:`object-handlers.go:1428-1433` 在 `copyDstOpts`(`:1454`)之前套用目的端默认;`bucket-sse-config.go:139-151` 在 `nil 配置 + AutoEncrypt` 与桶默认 KMS 两种情况下都写入 `aws:kms`,因此 explicit / auto / bucket 三种模式均进入 KMS 分支。
- 回归证明成立:`baseline-final.log` 用 `-overlay` 换回未修复 constructor,失败面精确为「trusted × 有效标签时间戳 × SSE-KMS / SSE-KMS-context」和 6 个 KMS COPY 子测试(`tags="key=old"`、`kms=true`、HTTP 200),`none`/`SSE-S3`/`SSE-C`/非 trusted 全通过。修复后 `focused.log:194-204` 全 PASS。
- 测试确实覆盖被要求的维度:信任边界(trusted=false 时 mtime/ETag/三时间戳全归零)、错误路径(trusted + 畸形值必须报错且错误串含头名)、SSE 序列化回环(KMS keyID/context 原样还原)、磁盘终态(每事件 `obj.GetObjectInfo` 读真实盘)、版本一致性、签名 GET 明文可读。全局 `GlobalKMS`/`globalAutoEncryption`/`set.getDisks` 均 defer 还原。
- `race` exit 0、`vet` 空输出、`golangci-lint` 0 issues,均记录了与 HEAD 一致的三文件哈希。
- 未发现 `verification.md` / `verification.json` / `consensus.md` 中与日志矛盾的陈述。(评审者版本/模型/effort 这类 provenance 声明不在我可验证范围,未作背书。)
## 发现清单
| ID | 内容 | 阻断 |
|---|---|---|
| IMPL-01 | 单字段修复正确且充分,位置、变量、注释与 `:473` 语义一致 | 否(确认项) |
| IMPL-02 | 基线失败/修复通过的判别力成立,对照组不误报 | 否(确认项) |
| IMPL-03 | KMS 字面量相对 `getDefaultOpts` 仍缺 `ProxyHeaderSet`/`ProxyRequest`/`Speedtest`(`object-api-options.go:40-44` vs `:449-460`)。R4 范围外,已登记为 R4-06 | 否,不设为新合并门槛 |
| IMPL-04 | `metadata-directive: REPLACE` 下 `getCpObjMetadataFromHeader`(`:1143-1156`)返回全新 map,故 `:1818` 的 `lastTaggingTimestamp` 为空、`:1822` 解析失败使 handler 侧比较恒「incoming 胜」;真正的 stale 拒绝发生在写锁内的 `reconcileStoredObjectTags`(`erasure-object.go:1312-1315`)。测试终态断言仍正确,文档 R4-07 已明示此分工 | 否(R5 上下文) |
| IMPL-05 | 测试卫生:`bucket-kms` 模式写入的 `bucketSSEConfig` 未还原,仅因它是最后一个 mode、且 `ExecObjectLayerAPITest` 每后端重建对象层并 `resetTestGlobals()` 才安全;后续若在其后追加 mode 会继承默认 KMS | 否 |
| IMPL-06 | `object-api-options-replication_test.go:35` 局部变量名 `context` 遮蔽标准包名(本文件未导入该包),纯观感 | 否 |
| IMPL-07 | `focused.log` exit 1 的唯一失败是既有 `TestAPICopyObjectReplicaRetentionRemovalUnderBucketKMS`(`replication-trust_test.go:1284`,"Storage reached its minimum free drive threshold"),属本机磁盘余量环境问题,非本次引入;容量 overlay 是测试专用、未提交。新增 COPY 测试自带 `tagTestCapacityDisk` 包装,不受该阈值影响 | 否 |
**没有发现阻断性正确性问题。** 生产语义、存储格式、API 与既有排序规则均未改变,无任何既有测试断言旧(缺陷)行为。
## 合并适配性
`dbcf8dec5` 直接位于实时 main `9ebe81c1b` 之上,可快进合并。按仓库 CI 通过为前提,本实现适合合入 main。
*(我未运行任何测试,也未查询 GitHub;以上仅基于源码阅读与所提供日志。)*
@@ -0,0 +1,33 @@
{
"baseline": "9ebe81c1b3611f9cc73e676b5b741c2be62c467a",
"reviewed_head": "dbcf8dec589deb5d91e17d295cb70997635f5b55",
"requested_model": "claude-opus-5",
"requested_effort": "max",
"cli_version": "2.1.270",
"diff_sha256": "c8cd6648f8ecea835ec74a038cdeaa82acaa3f36250395f97ead3260dc2fc0a5",
"started_at": "2026-09-15T15:59:41.004963+00:00",
"status": "completed",
"command": "/opt/homebrew/bin/claude --print --model claude-opus-5 --effort max --safe-mode --permission-mode plan --tools Read,Grep,Glob --strict-mcp-config --no-session-persistence --add-dir /Users/vonng/tmp/silo-r4-evidence-20260915-a9cb --output-format stream-json --verbose",
"completed_at": "2026-09-15T16:03:52.331640+00:00",
"assistant_models": [
"claude-opus-5"
],
"observed_model": "claude-opus-5",
"verdict": "GO_WITH_NONBLOCKING_NOTES",
"blocking_findings": 0,
"session_id": "c09bc4fb-85f1-4af8-a1de-c453398e4a20",
"duration_ms": 161469,
"subtype": "success",
"is_error": false,
"used_tools": {
"Read": 20,
"Glob": 3,
"Grep": 14,
"ExitPlanMode": 1
},
"raw_stream": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/merge-review-1/opus.jsonl",
"stream_sha256": "875f617454b27e15ab44d9b777de89643ee352e0ccb948faad570fc546261acc",
"review_sha256": "962ff88d2ffcb75cd692ec17017411d624de80dc120f8dedc675e0fe25009335",
"prompt_sha256": "e1341c721161229431b943ba18d89b740e94470803c099b9ae3d597fd50544a4",
"review_extraction": "The substantive review is an earlier assistant text block; result.result only repeats CLI plan-mode merge limitations. Full raw stream and all assistant text are retained."
}
@@ -0,0 +1,84 @@
{
"original_reviewed_head": "dbcf8dec589deb5d91e17d295cb70997635f5b55",
"dco_signed_equivalent_head": "03027727d1d1b97d8beb83ac55569ea9a83dab23",
"notice_equivalence": {
"cmd/object-api-options-replication_test.go": {
"before_sha256": "1ea2a060987e32a4c76fce96ee974df475944c2d6ab482a4893e33daf7bca849",
"after_sha256": "c21fc8889a079085d9a882499a1cbe868278a3517580651f3bed1102e2a6aef8",
"package_body_sha256": "096f143c0b0a068581f9bb892f35ded0d65b6b60ab711f043236d27fbf51ca33",
"body_unchanged": true
},
"cmd/object-copy-replication-tagging_test.go": {
"before_sha256": "5437a77e68736b4ce69de9c777675251fef24b0352dfe30bd8a836fc7ee810e3",
"after_sha256": "73f066ed7258d430f078ecc90e551ece878bd3d4672bc762094d434ff8fec23d",
"package_body_sha256": "6b5173db2ded2d54055073c3259be208a4d7c8eac0367687082877f1fd3bef15",
"body_unchanged": true
}
},
"checks": {
"verifiers": {
"command": [
"make",
"verifiers",
"GOLANGCI=/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/merge-review-1/golangci-serial"
],
"exit_code": 0,
"started_at": "2026-09-15T16:04:55.536605+00:00",
"finished_at": "2026-09-15T16:07:03.252820+00:00",
"cwd": "/Users/vonng/.codex/worktrees/a9cb/silo",
"env_override": {
"GOMAXPROCS": "2",
"GOFLAGS": "-p=2"
},
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "c21fc8889a079085d9a882499a1cbe868278a3517580651f3bed1102e2a6aef8",
"cmd/object-copy-replication-tagging_test.go": "73f066ed7258d430f078ecc90e551ece878bd3d4672bc762094d434ff8fec23d"
},
"log_sha256": "e42a5bb55f5c1ebfcf02cebebf6d82cf1ec5a2d74590cdf838deba16dd80bfdf"
},
"build": {
"command": [
"make",
"build"
],
"exit_code": 0,
"started_at": "2026-09-15T16:07:03.253715+00:00",
"finished_at": "2026-09-15T16:07:35.594062+00:00",
"cwd": "/Users/vonng/.codex/worktrees/a9cb/silo",
"env_override": {
"GOMAXPROCS": "2",
"GOFLAGS": "-p=2"
},
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "c21fc8889a079085d9a882499a1cbe868278a3517580651f3bed1102e2a6aef8",
"cmd/object-copy-replication-tagging_test.go": "73f066ed7258d430f078ecc90e551ece878bd3d4672bc762094d434ff8fec23d"
},
"log_sha256": "6ba9b545236be964861749c72e7609edf12b8f470df30d1ede8fd62f497e629b"
},
"binary-version": {
"command": [
"./silo",
"--version"
],
"exit_code": 0,
"started_at": "2026-09-15T16:07:35.594918+00:00",
"finished_at": "2026-09-15T16:07:37.616706+00:00",
"cwd": "/Users/vonng/.codex/worktrees/a9cb/silo",
"env_override": {
"GOMAXPROCS": "2",
"GOFLAGS": "-p=2"
},
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "c21fc8889a079085d9a882499a1cbe868278a3517580651f3bed1102e2a6aef8",
"cmd/object-copy-replication-tagging_test.go": "73f066ed7258d430f078ecc90e551ece878bd3d4672bc762094d434ff8fec23d"
},
"log_sha256": "36317d06b691593fe0d74f88d053a24485500c15fc2001e857f2fc6fa5ba752a"
}
},
"binary_version": "silo version DEVELOPMENT.2026-09-15T16-03-52Z (commit-id=03027727d1d1b97d8beb83ac55569ea9a83dab23)\nRuntime: go1.27.1 darwin/arm64\nLicense: GNU AGPLv3 - https://www.gnu.org/licenses/agpl-3.0.html\nCopyright: 2015-2025 MinIO, Inc.\nModifications: Copyright 2025-2026 PGSTY\nSource compatibility: based on MinIO technology\n",
"raw_evidence_directory": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/merge-review-1",
"all_function_and_test_bodies_identical_to_opus_reviewed_version": true
}
@@ -0,0 +1,36 @@
# R4 合并前复核
用户已明确追加授权:使用 Opus 5 max 核实最终实现,确认无误后合并 main。本轮授权取代此前只交付本地补丁的范围限制。
## 真实实现评审
- 独立新调用:Claude Code 2.1.270,`--model claude-opus-5 --effort max`。
- 复核代码提交:`dbcf8dec589deb5d91e17d295cb70997635f5b55`;当时实时 main 与 fetch 结果均为 `9ebe81c1b3611f9cc73e676b5b741c2be62c467a`。
- 实际 assistant 模型只有 `claude-opus-5`。结论 **GO_WITH_NONBLOCKING_NOTES,0 阻断项**,明确表示仓库 CI 通过后适合合入 main。
- [原始实现评审正文](implementation-review.md)、[实际模型与输出哈希](implementation-review.metadata.json) 已保存。
- 原始流、全部 assistant 正文与最终 result 位于 `/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/merge-review-1/`。实质评审出现在较早的 assistant 消息;最终 result 只重复 Claude 只读会话不能自行合并的工具限制,不是对修复结论的撤回。本任务由 Codex 按用户明确授权完成合并。
## 意见处置
| 条目 | 处置 |
|---|---|
| IMPL-01 / IMPL-02 | 确认单字段修复和基线失败/修复通过的测试判别力,无需追加修改。 |
| IMPL-03 | Proxy/Speedtest 选项遗漏已交父任务单独核验,维持范围外,不纳入 R4 合并。 |
| IMPL-04 | REPLACE 请求的旧标签拒绝由写锁内对账完成,测试断言真实落盘状态,已有文档准确说明。 |
| IMPL-05 | 当前测试固定以 bucket-kms 为最后一种模式,且每后端重新初始化;现有执行顺序安全。后续增添模式需同步隔离桶默认配置,本次保持已评审测试逻辑。 |
| IMPL-06 | 局部变量 context 命名建议为可选观感项,不改动已评审逻辑。 |
| IMPL-07 | 既有锁测试的磁盘余量限制及仅测试容量 overlay 已如实记录;新测试和 race 不使用生产代码 overlay。 |
## 提交规范调整
按 `CONTRIBUTING.md` 补齐提交作者对应的 DCO sign-off,并将两个新原创测试文件的文件头改为 `Copyright (c) 2026 Feng Ruohang`,保留 AGPL-3.0-or-later。原有生产文件的继承声明保持原样。
生产函数和测试的 `package cmd` 之后内容与 Opus 审查版本逐字节相同。`merge-review-1/notice-equivalence.json` 记录了旧/新文件哈希及不变的代码正文哈希。原 `verification.json` 保留当时原始验证记录,不覆盖历史哈希;本轮 PR 的 CI 对最终提交重新验证。
## 合并门槛
`make verifiers` 已通过:全仓 lint 为 0 issues,生成文件检查通过,rebrand 兼容性清单未变化,交付/运行时标识检查和 entrypoint 参数兼容性测试通过。首次执行曾遇到其他任务持有 golangci-lint 进程锁;使用工具自带 `--allow-serial-runners` 串行等待后完成全部检查。可选 typos 工具未安装,由仓库 Makefile 按既有规则跳过。
`make build` 通过,已生成本地 `silo` 并成功执行 `./silo --version`。最终三个源文件哈希与本轮校验记录一致,详情见 [本轮验证清单](merge-verification.json)。
接下来由 PR CI 验证最终候选,并在合并前再次核对 main 和精确 PR head。CI 与合并事实以 GitHub PR 状态和本机原始合并证据为准,评审意见不等同于合并或发布。
+46
View File
@@ -0,0 +1,46 @@
## R4 独立评审(Opus 5.0,只读复核)
**计划**:plan v1 — `docs/investigations/r4/plan-v1.md`
**哈希(按任务给定)**:`ad539f2071155de6955b583991684ed33c4bfe2e29660005840cdc97d7e1a754`
**Baseline**:`9ebe81c1b3611f9cc73e676b5b741c2be62c467a`
### 裁定:GO_WITH_NONBLOCKING_NOTES
无阻断项。**我同意这份 exact plan(上述哈希)可以进入本地实现。** 下列 R4-01…R4-08 均为非阻断;其中 R4-02/04/05 的测试建议是**可选**的,不自动并入必做范围。
(说明:本会话 Write 工具被禁用,评审仅以正文返回,未写入任何文件,也未改动任何源码。)
### 我实际核验到的关键事实(支撑"单字段补丁正确且充分")
1. **缺陷确认**:`cmd/object-api-options.go:449-460` 的 KMS 字面量带了 MTime/PreserveETag/ReplicationRequest + 两个 Object Lock 时间戳,独缺 tagging;默认路径 `:473` 有。补丁片段中的变量名 `taggingtimestmp` 与 `:419` 完全一致,可直接编译;gofmt 对齐由更长的两个 Lock 键决定,不会扰动他行。
2. **影响面封闭**:全仓 `ReplicationSourceTaggingTimestamp` 只在 `cmd/object-handlers.go:1820` 被读取(定义于 `object-api-interface.go:99`)。因此该字段对 PUT/分段路径天然无效果——既印证 R4/R5 的切分合理,也说明补丁不可能回归其他路径。
3. **充分性的关键点(我重点查证的风险)**:`encMetadata` 只有在 SSE-C 轮换分支 `object-handlers.go:1648-1659` 才批量快照全部保留键,而该分支与 KMS options 分支互斥(目的端是 SSE-C 时 `crypto.S3KMS.IsRequested` 为假)。故 `:1903` 的 `maps.Copy(srcInfo.UserDefined, encMetadata)` **不会**覆盖 KMS COPY 新写入的 tags/时间戳 —— 单字段补丁在 R4 边界内充分。
4. **三个触发点准确**:`bucket-sse-config.go:135-153`(显式请求优先 → nil 配置 + AutoEncrypt → KMS → bucket 默认 KMS 写 header+keyID;默认 AES 走 AES 分支),配合 `object-handlers.go:1428-1433` 仅在非联邦时套用目的端默认。
5. **REPLACE 副本路径准确**:`reconcileStoredObjectTags`(`erasure-server-pool-consistency.go:232-243`)语义即"存量有效时间戳胜过缺失/更旧/相等的 incoming,并连同 tag 值一起还原"。KMS 目的端因 `isTargetEncrypted` 使 `metadataOnly=false`,实际落到 `erasure-server-pool.go:1499-1513`(`ReplicaLockReconcile` 经 `:1509` 透传)→ `erasure-object.go:1276-1316`,在 `cloneMSS`(:1324) 之前于写锁内完成对账;纯元数据路径走 `erasure-object.go:136-139`。计划同时引用 `:136` 与 `:1509`,判断正确。
6. **证据可信**:`baseline-repro.log` 中 options 用例非 KMS 保留 `...123456789Z`、KMS 返回零值;COPY 用例 6/6(ErasureSD + Erasure16 × explicit/auto/bucket KMS)失败,且均为 200、`kms=true`、明文 GET 通过、tags 停在 `key=old`。即"请求成功、加密正常,但复制标签被静默丢弃",与计划表述一致,未夸大。
7. **修复后推演**:newer/stale/duplicate/newer-again 在 handler(:1817-1833) 与写锁对账的双重排序下分别得到 new/new/new/latest,与测试期望吻合;旧发送端不带 `X-Minio-Source-Tagging-Timestamp` 时仍为零值 → 行为不变,兼容性主张成立。
### 问题清单
| ID | 阻断 | 内容与建议 |
|---|---|---|
| **R4-01** | 否 | 行号漂移:计划写的 `1851/1910`,实际是 `object-handlers.go:1847`(`ReplicaLockReconcile`)与 `:1903`(encMetadata merge)。建议更正引用。 |
| **R4-02** | 否(建议可选) | `ExecObjectLayerAPITest` 两种后端均为**单 pool**(`test-utils_test.go:216` `mustGetPoolEndpoints(0, ...)`),故 `erasure-server-pool.go:1443` 多池分支未被覆盖;且 KMS 目的端命中的是 PutObject 重写对账而非 `CopyObject:136`。建议在计划或测试注释中点明"单盘/16 盘均为单池";补多池覆盖**可选**,不必进必做范围。 |
| **R4-03** | 否 | 措辞:`crypto.Requested`(`internal/crypto/sse.go:74`)只检查**目的端** SSE 头,因此仅带 SSE-C *copy-source* 头的请求在 KMS 默认桶/自动加密下仍会进入 KMS 分支(归入触发点 2/3,枚举仍完整)。建议澄清 "source encryption alone…" 一句。 |
| **R4-04** | 否(**可选**) | 建议在 options 矩阵里加一条 KMS 分支 vs 默认分支的**逐字段等价断言**(MTime/PreserveETag/ReplicationRequest/三个复制时间戳)。这是阻止第三次复发最廉价的护栏(2021 漏、2026 补了两个 Lock 时间戳仍漏此项)。计划第 1 条已基本覆盖,此为结构化建议。 |
| **R4-05** | 否(**可选**) | 建议加一例"KMS 目的端 + 有 tags 但无 tagging 时间戳头 → 存量不变",把兼容性主张钉在 handler 层而不仅在 options 层。 |
| **R4-06** | 否(范围外,仅登记) | 同一 KMS 字面量相对 `getDefaultOpts`(`object-api-options.go:40-44`) 还遗漏 `ProxyHeaderSet/ProxyRequest/Speedtest`;`opts.Speedtest` 在 `erasure-object.go:1625` 被读取,全局自动加密下 speedtest PUT 会丢该标志。**不要在 R4 修**,且当前也不在 R5 声明范围内,建议单列条目登记。 |
| **R4-07** | 否 | REPLACE 时 `getCpObjMetadataFromHeader:1143-1156` 会重建 map,`lastTaggingTimestamp` 为空 → handler 对 stale 事件**恒接受**,真正的拒绝来自写锁内对账。因此回归测试必须断言**最终落盘状态**(现有复现已如此),不要改为断言 handler 层行为。 |
| **R4-08** | 否(不确定性) | 本会话无 shell,无法独立复算计划 SHA-256、验证 `c4373ef290 / b2dca43fda / cfefc049c` 历史归属与 PR #184/#187。可由 `shasum -a 256 docs/investigations/r4/plan-v1.md` 与 `git log -L` 输出消解;均不影响补丁正确性。`cfefc049c` 的 KMS context 编码修复实体(`:436-444` 的 `sdkContext`)仍在,补丁不触碰。 |
### 对计划各主张的逐项裁定
- 单字段补丁**正确且对本 bounded issue 充分**:同意(依据 2/3/7)。
- 三个目的端 KMS 触发点**描述准确**:同意(R4-03 仅措辞澄清)。
- 当前 REPLACE 副本路径**描述准确**:同意(R4-01/02 属引用精度)。
- 回归矩阵与存量状态说明**充分**:同意;存量部分"不自动回填、丢失源时间不可重建、并列/更旧事件不保证修复"的表述与 `reconcileStoredObjectTags` 实际语义一致。
- 信任边界、错误行为、加密 key/context、tie 语义、兼容性:补丁均未触碰,维持不变。
### 交付
本轮为**计划共识**,非实现验收。我未作任何源码或文件修改;R4 可按 plan v1 在本地实施,实施后的差异与测试证据需另行验收。
@@ -0,0 +1,36 @@
{
"baseline": "9ebe81c1b3611f9cc73e676b5b741c2be62c467a",
"plan": "docs/investigations/r4/plan-v1.md",
"plan_sha256": "ad539f2071155de6955b583991684ed33c4bfe2e29660005840cdc97d7e1a754",
"requested_model": "claude-opus-5",
"requested_effort": "max",
"cli_version": "2.1.270",
"started_at": "2026-09-15T15:45:46.438630+00:00",
"status": "completed",
"raw_output": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/opus-review-v1.jsonl",
"raw_stderr": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/opus-review-v1.stderr.log",
"command": "/opt/homebrew/bin/claude --print --model claude-opus-5 --effort max --safe-mode --permission-mode plan --tools Read,Grep,Glob --strict-mcp-config --no-session-persistence --add-dir /Users/vonng/tmp/silo-r4-evidence-20260915-a9cb --output-format stream-json --verbose",
"completed_at": "2026-09-15T15:49:48.150062+00:00",
"assistant_models": [
"claude-opus-5"
],
"observed_model": "claude-opus-5",
"verdict": "GO_WITH_NONBLOCKING_NOTES",
"blocking_findings": 0,
"result_subtype": "success",
"is_error": false,
"session_id": "63e14a68-8565-41fa-9746-3e405fb63e9f",
"duration_ms": 161784,
"num_turns": 35,
"used_tools": {
"Read": 18,
"Glob": 4,
"Grep": 11,
"Write": 1
},
"stream_sha256": "eb0918d8a6185b180dddcfc664a96682f05502ecf3b686b08a0547f09879d57d",
"review_sha256": "e1dc12dd99326ae432623ff8de201813e6e84e7ed16a5556c21f9c514d663676",
"prompt_sha256": "07c225beff1523e056c154b3a387cf1ae345def4b4d0173065b882140ac5abdf",
"tool_scope_note": "Read/Grep/Glob allowed. Claude attempted Write to its own plan; the tool was disabled and no file was written. git diff before consensus showed no production source changes.",
"auxiliary_model_note": "assistant_models records actual reviewing assistant messages. Auxiliary usage is distinct. --effort max is explicit in the command, not inferred from model usage."
}
+65
View File
@@ -0,0 +1,65 @@
# R4 plan v1: preserve the replicated tag timestamp for SSE-KMS
## Baseline and ownership
- Baseline: `9ebe81c1b3611f9cc73e676b5b741c2be62c467a`, verified against GitHub main on 2026-09-15.
- Branch: `codex/r4-kms-tag-timestamp`; worktree: `/Users/vonng/.codex/worktrees/a9cb/silo`.
- Live open PRs at inspection: #184 and #187, neither owns this options change.
- The worktree lacks the ignored `AGENTS.md`; the parent explicitly confirms `/Users/vonng/pgsty/silo/AGENTS.md` applies. Maintain the PGSTY product graph and inexpensive compatibility.
- R4 owns only the missing field in `cmd/object-api-options.go` and its regression tests. R5 owns DELETE/empty tags, PUT/multipart receiving, sender propagation and full receiver ordering. R4 will supply a standalone source patch to R5; neither task edits the other's worktree.
## Proven defect and actual trigger
`putOptsFromHeaders` parses the trusted source tag timestamp before selecting encryption. The SSE-KMS branch constructs and returns another `ObjectOptions` carrying mtime, ETag, replication trust and both Object Lock timestamps, but omits `ReplicationSourceTaggingTimestamp`. The normal path retains it. The parser accepts and preserves RFC3339 fractional seconds even though its layout is `time.RFC3339`; the reproduction uses nanoseconds.
`CopyObjectHandler` applies local destination encryption configuration before `copyDstOpts` → `putOptsFromReq` → `putOpts` → `putOptsFromHeaders`. The omission is reached by:
1. Explicit destination SSE-KMS request headers (with or without a key ID/context).
2. A destination bucket with default SSE-KMS, when the request has no explicit SSE choice.
3. Global automatic encryption with no bucket SSE override and no explicit SSE choice.
Explicit AES256/SSE-C takes its existing branch; source encryption alone does not select the destination KMS branch. Remote federation skips local destination defaults. The relevant trigger is trusted metadata entering the destination KMS branch, not every SSE-KMS object or every tag operation.
At `CopyObjectHandler`'s tag decision, a zero source timestamp skips the tag update. Current under-lock reconciliation can preserve the stored tag/timestamp when metadata REPLACE reconstructs the map with no timestamp. In the observed same-version replica COPY, the request succeeds, destination encryption is valid, and the old tags/timestamp remain. A missing field in the options layer is not itself proof of a content-read failure.
PUT and multipart consumers' independent failure to persist a parsed tag timestamp remain R5's responsibility. R4 does not claim to fix all tag replication by correcting this constructor.
## Source and reproduction evidence
- `cmd/object-api-options.go`: trusted parsing at 383–426; KMS construction at 433–460; normal assignments at 469–475.
- `cmd/object-handlers.go`: destination default encryption at 1425–1435; `copyDstOpts` at 1454; tag timestamp consumption at 1807–1834; replica reconciliation enabled at 1851; encryption metadata merge at 1910.
- `internal/bucket/encryption/bucket-sse-config.go:135`: explicit request wins, absent config + auto encryption selects KMS, otherwise configured bucket algorithm/key ID applies.
- `cmd/erasure-server-pool-consistency.go:232`: stored valid timestamp wins over absent, older or equal incoming timestamp; writes preserve the stored tag value alongside its timestamp.
- `cmd/erasure-object.go:136` and `cmd/erasure-server-pool.go:1509`: same-version replica COPY reaches existing under-lock tag reconciliation, including object-data rewrites.
- History: the omission exists in `c4373ef290` (2021-09-18); `b2dca43fda` (2026-09-05) added the two Object Lock timestamps but not the tag timestamp. `cfefc049c` fixed KMS context encoding independently and must remain intact.
- Fresh temporary reproduction: `/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/r4_repro_test.go` and `baseline-repro.log` (overlay; no production edits).
- Command: `GOMAXPROCS=2 go test -p 2 -overlay /Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/overlay.json ./cmd -run '^TestReviewR4' -count=1 -timeout 5m -v`.
- Result: expected failure. Unencrypted and AES256 options preserve `2026-09-15T01:00:00.123456789Z`; KMS returns zero. Signed metadata REPLACE COPY on ErasureSD and Erasure (16 disks), across explicit/default/automatic KMS, returns 200 but retains `key=old` and the old timestamp for newer events. Actual encrypted metadata and plaintext GET roundtrips pass. Test deltas are 1–3 nanoseconds.
- These are in-process signed HTTP router and real local disk tests. `kms.NewStub` replaces the remote key service; the normal server encryption/decryption code still runs. Existing `tagTestCapacityDisk` avoids the host's free-space percentage threshold; it delegates all object data/metadata I/O to real test disks.
## Proposed production change
Add exactly this field to the existing KMS `ObjectOptions` literal:
```go
ReplicationSourceTaggingTimestamp: taggingtimestmp,
```
Update the neighboring explanatory comment to include tagging alongside retention/legal hold. Do not refactor the common return paths, change parsing/fallback/equal-timestamp semantics, change encryption context encoding, modify trust decisions, add SDK dependencies, or change storage/wire format. Those changes are unnecessary to restore the missing existing contract.
## Required validation after consensus
1. Add an options regression matrix covering unencrypted, SSE-S3, SSE-KMS with no context, SSE-KMS with a context, and SSE-C. Validate trusted/untrusted requests, missing/valid/malformed tag timestamps, nanosecond and timezone/whitespace handling, all three replication timestamps, mtime/ETag/trust, nonnil metadata, and unchanged SSE header serialization (including KMS key/context).
2. Promote the temporary COPY reproduction into a named, isolated regression test. Use actual signed same-version metadata COPY with REPLACE metadata and tagging directives, on single-disk and 16-disk backends. For explicit, bucket-default and automatic SSE-KMS, check newer update, older delivery, duplicate replay, and a second newer update. Verify stored tags, exact timestamp, version ID, encryption kind and plaintext GET after each operation. Include an unencrypted/SSE-S3 control if the fixture can do so without expanding implementation scope.
3. Fail the final regression tests against unmodified baseline using an overlay. Then run them on the fixed source, alongside existing replication-trust/options and bucket-KMS Object Lock tests. Check `gofmt`, `git diff --check`, and `go vet ./cmd`.
4. Run the new focused tests under `-race`. Use `GOMAXPROCS=2` and `-p 2` while sibling tasks share the host. A one-field pure option fix does not justify concurrent full-repository suites in all five tasks; full Linux CI and multi-site validation remain separate delivery gates.
5. If a test exposes a separate handler/storage defect, report evidence and coordinate with R5. Do not broaden R4's production patch to make unrelated tests pass.
## Compatibility, existing state, effort and delivery
- Public API, header names, stored key names, KMS context/key handling and supported dependencies remain unchanged. Untrusted source headers stay ignored; malformed trusted timestamps continue to fail; absent timestamp remains zero. Existing non-KMS behavior remains unchanged.
- No automatic rewrite/backfill. Lost source tag times cannot be reconstructed from the receiver alone. Upgrading permits subsequent properly timestamped events to be consumed. Review source-of-truth and target state before any targeted resync; full historical convergence also depends on R5. Repeated events subject to existing timestamp/tie semantics are not a universal repair guarantee.
- The source fix can land independently; complete deletion/empty-tag and mixed-encryption convergence needs R5 plus its integration evidence.
- Expected effort: approximately 0.5–1 engineer-day including reproduction, review and local validation; key-service deployment, multi-site failures and existing-state remediation are separate.
- After actual Opus 5.0/max agreement on this exact plan hash, implement locally without another user permission prompt. Preserve raw review, assistant model identity, request effort, baseline and plan hash, issue-by-issue disposition and explicit consensus before source edits.
- Deliver a reviewable local diff, tests and evidence. No main merge, remote publication/release, deployment or existing-state rewrite is authorized by this plan.
+18
View File
@@ -0,0 +1,18 @@
# R4 research log
## Verified baseline
2026-09-15: local clean HEAD and GitHub main both `9ebe81c1b3611f9cc73e676b5b741c2be62c467a`. Branch created as `codex/r4-kms-tag-timestamp`. Live GitHub open PRs #184 (`6addf9eb916b5a4b837480cf534cd1efa5407d3c`) and #187 (`b8f2fdde41dff3dc3b8db669c1d42d30ca5c1d3d`) concern other tasks. Claude Code reports `2.1.270`; Go reports `go1.27.1 darwin/arm64`.
## Coordination
- Parent task: `01a0a5ab-ee43-7911-bddd-1aca6f8afcc8`.
- R5: `01a0a5b9-602d-7470-9882-4817cf5fdcd1`, `/Users/vonng/.codex/worktrees/77ad/silo`.
- Parent and R5 acknowledged the ownership boundary: R4 options constructor and nonempty KMS COPY tests; R5 producer/receiver ordering and empty values. R5 will consume R4's minimal patch for combined KMS acceptance.
- Initial conservative expectation separated metadata COPY from REPLACE ordering. Inspection of current `ReplicaLockReconcile` and `reconcileStoredObjectTags` shows that stored timestamps are also reconciled under the write lock for REPLACE. The temporary reproduction therefore uses REPLACE directly; the fix must demonstrate the actual sender-shaped path without changing the handler.
## Baseline reproduction
Temporary overlay test source and output are in `/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/`. The options and HTTP/disk reproductions fail for the expected missing timestamp. Every KMS HTTP request completed with 200; newer tags remained old; encrypted object metadata and subsequent ordinary plaintext GET succeeded. The result is narrower than claiming all KMS replication fails, and stronger than merely comparing options.
The temporary source is not a production implementation. See [plan v1](plan-v1.md) for exact scope and required acceptance. The plan is frozen by SHA-256 before invoking real Opus.
+358
View File
@@ -0,0 +1,358 @@
{
"baseline": "9ebe81c1b3611f9cc73e676b5b741c2be62c467a",
"branch": "codex/r4-kms-tag-timestamp",
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "1ea2a060987e32a4c76fce96ee974df475944c2d6ab482a4893e33daf7bca849",
"cmd/object-copy-replication-tagging_test.go": "5437a77e68736b4ce69de9c777675251fef24b0352dfe30bd8a836fc7ee810e3"
},
"source_files_match_all_test_runs": true,
"checks": [
{
"name": "baseline-final",
"command": [
"go",
"test",
"-p",
"2",
"-overlay",
"/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/baseline-final-overlay.json",
"./cmd",
"-run",
"^Test(PutOptsFromHeadersReplicationTimestamps|APICopyObjectReplicaTaggingTimestampUnderKMS)$",
"-count=1",
"-timeout=5m",
"-v"
],
"cwd": "/Users/vonng/.codex/worktrees/a9cb/silo",
"env_override": {
"GOMAXPROCS": "2"
},
"started_at": "2026-09-15T15:50:49.430860+00:00",
"finished_at": "2026-09-15T15:51:20.637630+00:00",
"exit_code": 1,
"log": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/baseline-final.log",
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "1ea2a060987e32a4c76fce96ee974df475944c2d6ab482a4893e33daf7bca849",
"cmd/object-copy-replication-tagging_test.go": "5437a77e68736b4ce69de9c777675251fef24b0352dfe30bd8a836fc7ee810e3"
},
"log_sha256": "f849c2082235213764e3db7a314d52af75c4859102478a1e8f837afcaa8f1ea8",
"overlay_sha256": "f4cbc16e4ffccaf63191de2e8476162876055796adb4c38cda2d8c569a3bbabc",
"overlay_sources": {
"/Users/vonng/.codex/worktrees/a9cb/silo/cmd/object-api-options.go": {
"path": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/object-api-options.baseline.go",
"sha256": "16a560d0990ae929393f682f22b32ecd2e7f4d9390484b03b54e176fcd00cff5"
}
},
"assessment": "Expected baseline regression failure; KMS timestamp loss. Non-KMS controls pass."
},
{
"name": "focused",
"command": [
"go",
"test",
"-p",
"2",
"./cmd",
"-run",
"^Test(PutOptsFromHeadersReplicationTimestamps|APICopyObjectReplicaTaggingTimestampUnderKMS|ReplicationTrustControlsInternalOptionsAndEvents|GetAndValidateAttributesOpts.*|APICopyObjectReplicaRetentionRemovalUnderBucketKMS)$",
"-count=1",
"-timeout=5m",
"-v"
],
"cwd": "/Users/vonng/.codex/worktrees/a9cb/silo",
"env_override": {
"GOMAXPROCS": "2"
},
"started_at": "2026-09-15T15:51:20.638633+00:00",
"finished_at": "2026-09-15T15:51:48.382212+00:00",
"exit_code": 1,
"log": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/focused.log",
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "1ea2a060987e32a4c76fce96ee974df475944c2d6ab482a4893e33daf7bca849",
"cmd/object-copy-replication-tagging_test.go": "5437a77e68736b4ce69de9c777675251fef24b0352dfe30bd8a836fc7ee810e3"
},
"log_sha256": "9a64cbc46e9fd6186b1d3031034b720857851ce70b1466b1da5da6d1a52b76ba",
"assessment": "New tests and options/trust pass; pre-existing KMS lock fixture blocked by host disk free-space percentage."
},
{
"name": "focused-capacity-adapted",
"command": [
"go",
"test",
"-p",
"2",
"-overlay",
"/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/capacity-overlay.json",
"./cmd",
"-run",
"^Test(PutOptsFromHeadersReplicationTimestamps|APICopyObjectReplicaTaggingTimestampUnderKMS|ReplicationTrustControlsInternalOptionsAndEvents|GetAndValidateAttributesOpts.*|APICopyObjectReplicaRetentionRemovalUnderBucketKMS)$",
"-count=1",
"-timeout=5m",
"-v"
],
"cwd": "/Users/vonng/.codex/worktrees/a9cb/silo",
"env_override": {
"GOMAXPROCS": "2"
},
"started_at": "2026-09-15T15:52:20.113475+00:00",
"finished_at": "2026-09-15T15:52:50.842773+00:00",
"exit_code": 0,
"log": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/focused-capacity-adapted.log",
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "1ea2a060987e32a4c76fce96ee974df475944c2d6ab482a4893e33daf7bca849",
"cmd/object-copy-replication-tagging_test.go": "5437a77e68736b4ce69de9c777675251fef24b0352dfe30bd8a836fc7ee810e3"
},
"log_sha256": "13f7aa54564a433bef4dddcf2c5ad1fe46fa03869527fb902a255ce4c3263bc9",
"overlay_sha256": "76a7e6fb364bdaaa3469b1dc79f9ea318059f287a4888f322639920c76dfe53a",
"overlay_sources": {
"/Users/vonng/.codex/worktrees/a9cb/silo/cmd/replication-trust_test.go": {
"path": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/replication-trust-capacity_test.go",
"sha256": "c81b526ae51983881bbf464199e6b90f074fa695300fa8ff005e427e4d3c8208"
}
},
"assessment": "PASS"
},
{
"name": "race",
"command": [
"go",
"test",
"-p",
"2",
"-race",
"./cmd",
"-run",
"^Test(PutOptsFromHeadersReplicationTimestamps|APICopyObjectReplicaTaggingTimestampUnderKMS)$",
"-count=1",
"-timeout=5m",
"-v"
],
"cwd": "/Users/vonng/.codex/worktrees/a9cb/silo",
"env_override": {
"GOMAXPROCS": "2"
},
"started_at": "2026-09-15T15:52:50.843898+00:00",
"finished_at": "2026-09-15T15:53:43.789325+00:00",
"exit_code": 0,
"log": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/race.log",
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "1ea2a060987e32a4c76fce96ee974df475944c2d6ab482a4893e33daf7bca849",
"cmd/object-copy-replication-tagging_test.go": "5437a77e68736b4ce69de9c777675251fef24b0352dfe30bd8a836fc7ee810e3"
},
"log_sha256": "ae66a8c9e569a5c4b57ae56e75afc85c06a8b76f1567187519da0726605a4de4",
"assessment": "PASS"
},
{
"name": "vet",
"command": [
"go",
"vet",
"-p",
"2",
"./cmd"
],
"cwd": "/Users/vonng/.codex/worktrees/a9cb/silo",
"env_override": {
"GOMAXPROCS": "2"
},
"started_at": "2026-09-15T15:53:43.790197+00:00",
"finished_at": "2026-09-15T15:53:51.129773+00:00",
"exit_code": 0,
"log": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/vet.log",
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "1ea2a060987e32a4c76fce96ee974df475944c2d6ab482a4893e33daf7bca849",
"cmd/object-copy-replication-tagging_test.go": "5437a77e68736b4ce69de9c777675251fef24b0352dfe30bd8a836fc7ee810e3"
},
"log_sha256": "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855",
"assessment": "PASS"
},
{
"name": "lint",
"command": [
"/Users/vonng/pgsty/silo/.bin/golangci/v2.13.1/golangci-lint",
"run",
"--build-tags",
"kqueue",
"--timeout=10m",
"--config",
"./.golangci.yml",
"./cmd/..."
],
"cwd": "/Users/vonng/.codex/worktrees/a9cb/silo",
"env_override": {
"GOMAXPROCS": "2"
},
"started_at": "2026-09-15T15:53:51.130456+00:00",
"finished_at": "2026-09-15T15:55:39.114303+00:00",
"exit_code": 0,
"log": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/lint.log",
"source_sha256": {
"cmd/object-api-options.go": "25e2e9484fafd94d1b2c857b94373758e481893ad93fb1a063edf7746277accc",
"cmd/object-api-options-replication_test.go": "1ea2a060987e32a4c76fce96ee974df475944c2d6ab482a4893e33daf7bca849",
"cmd/object-copy-replication-tagging_test.go": "5437a77e68736b4ce69de9c777675251fef24b0352dfe30bd8a836fc7ee810e3"
},
"log_sha256": "e92606b0bf483111dff0a120c315ea165821348f31365020e2468a0059095c47",
"assessment": "PASS"
}
],
"format_checks": [
{
"command": [
"gofmt",
"-l",
"cmd/object-api-options.go",
"cmd/object-api-options-replication_test.go",
"cmd/object-copy-replication-tagging_test.go"
],
"exit_code": 0,
"output": ""
},
{
"command": [
"git",
"diff",
"--check"
],
"exit_code": 0,
"output": ""
}
],
"evidence_directory": "/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb",
"evidence_files": {
"baseline-final-overlay.json": {
"size": 177,
"sha256": "f4cbc16e4ffccaf63191de2e8476162876055796adb4c38cda2d8c569a3bbabc"
},
"baseline-final.json": {
"size": 1584,
"sha256": "5bf7be51e5a0eb41a40dc5fc5d3a1aba5df7733ad5dcb001f8d870a01c4233ba"
},
"baseline-final.log": {
"size": 22493,
"sha256": "f849c2082235213764e3db7a314d52af75c4859102478a1e8f837afcaa8f1ea8"
},
"baseline-identity.txt": {
"size": 1676,
"sha256": "9b21841e19a0cbb8ded18c2597488a527a27bedc65109d05d4ff598103073b68"
},
"baseline-options.log": {
"size": 11395,
"sha256": "d692f0a4bc9c58e2ac0087afa356ddf48f86e1040ea8138d68ffc4d0992bf3cb"
},
"baseline-repro.log": {
"size": 5693,
"sha256": "aa89b76f4723c6a3ce224faa7796403628d978a8707544bd97848b8887de2113"
},
"capacity-fixture.diff": {
"size": 870,
"sha256": "8d01e0b0068441f37ecee37125b81424d1f30d7c4fb37d435ea0cfe2e4617e5e"
},
"capacity-overlay.json": {
"size": 185,
"sha256": "76a7e6fb364bdaaa3469b1dc79f9ea318059f287a4888f322639920c76dfe53a"
},
"final_copy_repro_test.go": {
"size": 5942,
"sha256": "8949e07d96d2949a79f5a9e83c7a7c0473733d77b407e9c51da477d4ab74f1a8"
},
"focused-capacity-adapted.json": {
"size": 1661,
"sha256": "1a599b41caabfc5eb44db8d89c8b7e4f4f84f5f036d008369155fd337f65bdf9"
},
"focused-capacity-adapted.log": {
"size": 20384,
"sha256": "13f7aa54564a433bef4dddcf2c5ad1fe46fa03869527fb902a255ce4c3263bc9"
},
"focused.json": {
"size": 1253,
"sha256": "d9bb4979ea8aeaabb809cdc6e400a8673530bc83abf3dc2b2a06853a8523d0d9"
},
"focused.log": {
"size": 20475,
"sha256": "9a64cbc46e9fd6186b1d3031034b720857851ce70b1466b1da5da6d1a52b76ba"
},
"format-checks.json": {
"size": 352,
"sha256": "7569260900a799d5efdfb39db1f575ab1dadbbb04ace222e036968e66b6b59e7"
},
"lint.json": {
"size": 991,
"sha256": "a7144713b069f470a94b1ebe6fca6683a4b866a891a2756e15f28c280666ca14"
},
"lint.log": {
"size": 10,
"sha256": "e92606b0bf483111dff0a120c315ea165821348f31365020e2468a0059095c47"
},
"object-api-options.baseline.go": {
"size": 16653,
"sha256": "16a560d0990ae929393f682f22b32ecd2e7f4d9390484b03b54e176fcd00cff5"
},
"options-overlay.json": {
"size": 185,
"sha256": "5506b9c3b998b32f01c45af3cf01605eae9e4fb262c9ff3a6b0040abe719d4d9"
},
"options_repro_test.go": {
"size": 4192,
"sha256": "1a57a47bdd370042fa0f0d2d90efe447abedee9b9ef48a938d4bed631d83ec0b"
},
"opus-review-v1.exit": {
"size": 2,
"sha256": "9a271f2a916b0b6ee6cecb2426f0b3206ef074578be55d9bc94f6f3fe3ab86aa"
},
"opus-review-v1.jsonl": {
"size": 410466,
"sha256": "eb0918d8a6185b180dddcfc664a96682f05502ecf3b686b08a0547f09879d57d"
},
"opus-review-v1.stderr.log": {
"size": 0,
"sha256": "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
},
"overlay.json": {
"size": 158,
"sha256": "95f7c3f7e206fe36731e6d7e4a48f90c7403f07ae8155c125c6b86d2f1c2d487"
},
"r4-kms-tag-timestamp.patch": {
"size": 874,
"sha256": "2d4806d986bbd94ba4bc3951f3aeee48401ee1921c28ded0988fa09ca76ca26f"
},
"r4_repro_test.go": {
"size": 5508,
"sha256": "9bcefb6da2416b577b58085485cad60f677c2265e9dfa84d02e465e1b203766b"
},
"race.json": {
"size": 1026,
"sha256": "50b73e4acbc2426f3dcfadde78d0f0a86f10702345d2939a30204600bc750a13"
},
"race.log": {
"size": 18566,
"sha256": "ae66a8c9e569a5c4b57ae56e75afc85c06a8b76f1567187519da0726605a4de4"
},
"replication-trust-capacity_test.go": {
"size": 61271,
"sha256": "c81b526ae51983881bbf464199e6b90f074fa695300fa8ff005e427e4d3c8208"
},
"review-prompt-v1.md": {
"size": 2770,
"sha256": "07c225beff1523e056c154b3a387cf1ae345def4b4d0173065b882140ac5abdf"
},
"run-checks.py": {
"size": 2266,
"sha256": "ddecea5220bc9c286df18c0e9eca101f3d8ab731c307cd4acef937f3ecd11e65"
},
"vet.json": {
"size": 853,
"sha256": "e0964444bc91640ed6cf78229050bad5b94af4210bdf44b11d1a848f1ee930a6"
},
"vet.log": {
"size": 0,
"sha256": "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
}
},
"scope": "Darwin arm64; real signed HTTP + disk I/O; KMS service stub; single-pool single/16-disk fixtures; no remote CI, multi-site, release or deployment."
}
+52
View File
@@ -0,0 +1,52 @@
# R4 修复与本地验收
这是 2026-09-15 的本地验收快照。用户后续授权的实现级复核、提交规范调整与合并流程见 [合并前复核](merge-verification.md);以下原始测试记录及哈希保留当时状态。
## 结果
在 `putOptsFromHeaders` 的 SSE-KMS 选项构造中补齐 `ReplicationSourceTaggingTimestamp`。目的端使用显式 SSE-KMS、桶默认 KMS 或自动加密时,可信复制 COPY 现在能消费来源标签时间戳,并在现有存储锁内完成排序。
生产修改只有一个字段和相邻注释。API、存储格式、KMS key/context、信任判断和既有排序规则保持兼容。R5 的删除/空标签及 PUT/multipart 时间戳传播独立交付。
## 方案与 Opus 共识
- 基线:`9ebe81c1b3611f9cc73e676b5b741c2be62c467a`,已重新查询 GitHub main。
- 分支:`codex/r4-kms-tag-timestamp`。
- [冻结方案 v1](plan-v1.md):SHA-256 `ad539f2071155de6955b583991684ed33c4bfe2e29660005840cdc97d7e1a754`。
- 真实评审为本机 Claude Code 2.1.270,实际 assistant 模型 `claude-opus-5`,显式 `--effort max`。结论 **GO_WITH_NONBLOCKING_NOTES,0 个阻断项**。
- [逐条意见处置与双方共识](consensus.md)、[原始返回评审正文](opus-v1-review.md)、[模型与哈希记录](opus-v1.metadata.json) 已保存。先保存共识,再修改生产源码。
## 变更与测试
| 文件 | 内容 |
|---|---|
| `cmd/object-api-options.go` | 在 KMS 字面量中保留已解析的来源标签时间戳。 |
| `cmd/object-api-options-replication_test.go` | 无加密、SSE-S3、SSE-KMS、带 key/context 的 KMS、SSE-C;可信/非可信;缺失、有效、无效标签时间;纳秒、时区与空格;mtime/ETag/三个时间戳、metadata 与 SSE 序列化。 |
| `cmd/object-copy-replication-tagging_test.go` | 两种单池后端 × 五种目的端加密模式 × 五个有序事件,共 50 次签名 COPY 和 50 次普通 GET。每步检查最终标签、精确时间戳、对象版本、加密类型及明文内容。 |
COPY 使用 `metadata=REPLACE`、`tagging=REPLACE` 和可信复制身份。事件为较新更新、乱序旧更新、重复事件、再次更新,以及不带来源标签时间戳的请求。更新间隔仅 1–3 纳秒,防止时间精度退化被秒级测试掩盖。无加密与 AES256 是对照;KMS 覆盖显式、桶默认和自动加密入口。
## 验证状态
| 检查 | 结果 | 证据文件 |
|---|---|---|
| 最终测试 + 未修复基线 constructor overlay | 预期失败;只有可信 KMS 有效标签时间戳及 KMS COPY 更新失败,对照通过 | `baseline-final.log/json` |
| 修复后最终新增测试与既有 trust/options 测试 | 通过;未使用生产源码 overlay | `focused.log` |
| 既有 KMS Object Lock 回归 | 首次受宿主机磁盘余量阈值阻挡;仅适配测试容量报告后,与上述定向测试一起通过 | `focused-capacity-adapted.log/json`、`capacity-fixture.diff` |
| 新增测试 `-race` | 通过 | `race.log/json` |
| `go vet -p 2 ./cmd` | 通过 | `vet.log/json` |
| 仓库配置的 golangci-lint,范围 `./cmd/...`、`kqueue` build tag | 通过,0 issues | `lint.log/json` |
| gofmt、git diff --check | 通过 | `format-checks.json` |
原始日志和每条命令的运行记录位于 `/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/`。每份检查 JSON 都记录命令、退出码、时间和三个源码/测试文件的 SHA-256;最终交付已逐一确认文件哈希一致。[验证清单](verification.json) 另记录 overlay 的实际替换文件哈希,避免混淆基线与修复版执行代码。
测试使用真实签名 HTTP 路由、实际本地对象数据/元数据读写、服务器加解密代码;远程密钥服务由 `kms.NewStub` 代替。ErasureSD 与 16 盘 Erasure 均为单池。容量适配只使用已有 `tagTestCapacityDisk`,避免本机磁盘使用比例触发防写阈值,所有对象 I/O 仍由真实测试磁盘承担;未调整生产容量保护。
## 交付与剩余边界
- 本地实现和要求的定向验证均已完成,将源码、回归、研究、共识和验收记录作为一个本地提交交付。
- R4 的独立生产补丁已提供给 R5:`/Users/vonng/tmp/silo-r4-evidence-20260915-a9cb/r4-kms-tag-timestamp.patch`,SHA-256 `2d4806d986bbd94ba4bc3951f3aeee48401ee1921c28ded0988fa09ca76ca26f`。
- Opus 共识为方案级共识;本地测试结论来自实际运行,不把它记作 Opus 执行了测试。
- 多池/多站点故障恢复、外部 KMS 服务、完整 Linux CI、主干合并、远端发布和部署尚未执行。
- 没有改写存量。丢失的来源时间戳不能仅从接收端推导;后续重放/重同步须核对来源权威性及 R5 的全链路处理,不保证旧事件重放可以修复全部历史状态。
- 另登记 KMS 字面量缺少 Proxy/Speedtest 标志的范围外观察,已交父任务单独核验,本次未扩大修复。
+224
View File
@@ -0,0 +1,224 @@
# Durable IAM revocations in site replication
SILO retains the version of a deleted IAM record so an offline site cannot
restore an older identity or grant when it reconnects. This covers built-in
users, service accounts, groups, policy documents, and policy mappings in their
actual user/STS-parent/group namespaces. It also revokes the deleted built-in
user's older service accounts, STS credentials and group grants across deliberate
same-name recreation.
## Ordering and persistence
Deletion records occupy the original IAM configuration paths. They contain the
originating timestamp, `Deleted`, and, where required, `RevokedBefore`; identity
deletion records contain no secret key or session token. Normal IAM listings and
authorization hide deleted records. Object storage and etcd both serialize each
path's version comparison and write with a distributed lock. The source timestamp
is persisted without replacing it with the receiving node's clock.
Older events cannot overwrite a newer revision. At an identical timestamp, a
deletion wins over a live record. A local deliberate recreation receives a
version newer than the stored deletion. An older user/group deletion arriving
after recreation retains its revocation boundary while preserving the newer
live record. Receiving an already-applied tombstone does not rewrite or advance
it. These rules also apply after cold loading persistent IAM state.
The user or group revision is the commit point of deletion. Cleanup of mappings
and children follows that commit; cleanup failure cannot undo it. The API returns
the cleanup error and still notifies sibling nodes to reload the committed
state. Such an error does **not** mean the identity is still active. Retrying the
operation is safe. When the API returns a committed-cleanup error, the admin
handler does not send its immediate cross-site hook; cross-site propagation
relies on the normal deletion-healing retry. Sibling deletion notifications reload current shared state,
so a delayed notification cannot delete a subsequently recreated identity.
Etcd siblings also receive persistent changes through watches. Failed
notifications/watches remain eventual propagation, not a distributed
instantaneous revocation transaction.
Keep node clocks synchronized and monitor offsets. Ordering uses source wall
clock timestamps, with monotonic advancement for local writes to the same path.
It does not establish causal order between concurrent writers at different sites
or resolve conflicting live updates with identical timestamps deterministically.
## Users, child credentials and groups
A recreated user retains `RevokedBefore`. New service accounts and built-in STS
issuances carry a signed `siloParentRevocation` claim identifying the parent
boundary known at issuance. Editing or replaying an old child does not update
this claim. Old children remain invalid even when their own update timestamp is
newer than the parent deletion; children issued for the recreated parent remain
valid. Claims are read from verified tokens on credential load/write.
A newer replicated service-account snapshot can replace an older service with
the same access key, including a deliberate change of owner or secret. Local
duplicate creates remain rejected. Snapshots preserve disabled status and the
service's own revocation boundary, so earlier mappings cannot attach to the new
service. Equal-version retries reload the committed identity without rewriting
it. Periodic live healing compares source versions even when the public status
summary is unchanged, and includes disabled identities as healing sources.
Service snapshots also preserve their absolute expiration. The receiving site
does not reapply the local minimum lifetime for a newly issued credential. A
newer already-expired snapshot still supersedes the old key, is denied by
authentication, and is collected into a durable service tombstone by normal
loading. Cache/claims loading failures are returned for retry, not acknowledged;
after a committed replacement the stale cached secret is evicted immediately.
Collisions with an existing built-in or cached STS identity report an error and
require an explicit administrative resolution; replication cannot change its
credential kind. Concurrent conflicting service updates with exactly the same
timestamp can retain different winners at different sites; the status summary
does not resolve that case.
Each group member has its own `MemberGrants` timestamp. Changing another member
or the group's enabled status does not reissue everyone else's grants. Effective
membership requires the grant to be newer than both the user's and the group's
retained boundaries. Listings and policy evaluation use the same effective
membership. Peer snapshots preserve grant times, including unknown legacy grant
times; they cannot treat a recent snapshot time as a fresh grant to a revoked
identity. A new explicit administrative group grant can restore access.
This does not implement a general conflict-resolution protocol for all group
membership edits. In particular, the inherited live-group snapshot merge adds
members and does not reconcile a missed ordinary member removal. Removing a
member from a live group during a site outage is a separate known limitation;
do not infer that this change resolves it. User/group deletion boundaries and
same-name recreation are covered here.
## Retention and expiration
Permanent identities, groups, policy documents and mappings have no automatic
tombstone TTL. A disconnected peer or an old backup may return arbitrarily late.
Successful replay acknowledgements are an optimization, **not** permission to
garbage-collect this history.
Natural expiration of an immutable STS token physically removes its token-key
record and any legacy token-key mapping, without generating a permanent
tombstone. An early STS revocation is retained until that token's expiration plus
the existing clock-skew allowance. Replaying the same revoked token with a later
event timestamp cannot recreate it. Etcd uses an expiration lease; object storage
collects expired STS tombstones during its existing credential loading/purge.
A record with unknown expiration is retained conservatively. Cleanup writes are
best effort and use a short lock budget. If one fails, that load stops optional
reclamation, reports the error and still loads healthy users; expired credentials
stay denied and retain their existing durable version. The next load retries.
Healthy cleanup has no fixed record quota. The reusable STS
parent policy mapping is not assigned the token's TTL by deletion cleanup.
External-IDP disablement is an early revocation, not natural token expiration,
and its cached STS and service accounts are included in cleanup. Expiring service
accounts retain a durable revision because their access keys are reusable and
an older version might have no expiration.
Direct per-token `RevokeTokens` delivery between sites is not a new guarantee of
this change. The guarantee for built-in parent deletion follows from the durable
parent boundary, including children not currently present in the deleting node's
cache.
## Healing, failures and operational cost
Each process starts one healing loop. Losing its distributed leadership lease
pauses work until leadership is reacquired; it does not permanently terminate
healing. Configuration reloads do not create extra loops. The 30-second interval
starts after leadership is acquired and after each completed pass. Initial lock
retries and endpoint recovery can add further delay; it is not a convergence SLA.
The normal IAM loaders maintain an in-memory index of deletion records and
retained boundaries, without secrets. Healing uses this index; it does not add a
second full walk of `config/iam/` every cycle. Existing full IAM loading still
scans persistent records, including tombstones, at startup and on refresh.
Each peer receives batches of at most 128 records. The sender remembers which
path/version each peer acknowledged. Unrelated new changes at either site do not
reset that progress. A failed batch remains pending while later independent
batches can progress; a lost response may cause safe idempotent replay. A pass
has a bounded duration, and its successful acknowledgements survive that timeout.
The protocol reports each node name and process instance. Switching between
known node instances behind a load balancer preserves acknowledgements. A new
node instance conservatively invalidates prior acknowledgements once; repeated
switches among those known instances do not reset progress. Restore persistent state only with the affected processes
stopped, so a restore cannot reuse an old process acknowledgement.
Steady-state healing still traverses/sorts the retained in-memory set and checks
the peer's protocol status. It suppresses repeated deletion PUTs once acknowledged.
Memory use scales with retained paths and peers; startup storage reads scale with
history. This release does not provide general history compaction.
After upgrading, older live records whose receiving sites originally assigned
different timestamps can require an initial reconciliation wave. Allow for its
storage writes and sibling notifications when planning the maintenance window.
The cluster IAM metrics include `revocation_records`,
`revocation_heal_failures`, `revocation_heal_duration_millis`, and
`revocation_heal_last_success_timestamp_seconds`. Errors are also logged. A
nominal 30-second scheduler interval is not a convergence deadline: outages,
large backlogs, lock contention and failed requests can require more passes.
Cached credential lookup checks the in-memory parent revision index and performs
no additional storage read. STS issuance and cold credential loading still consult
the persistent parent revision. Revision I/O and distributed lock waits release
the IAM cache lock while a separate local writer mutex preserves write order.
These operations have bounded contexts, including etcd lock and lease cleanup.
## Protocol and supported upgrade
The server-owned versioned route is
`/minio/admin/v3/site-replication/peer/iam-revisions`. It carries source versions,
member grant times and distinct user/group revocation items without changing the
admin client SDK or S3 API. Older servers reject this route. The sender reports
the failure and does not fall back to a route that would discard the metadata.
The existing legacy IAM route remains readable for best-effort compatibility;
this does not confer the new guarantees on an older peer.
All participating servers must be upgraded for the guarantee in this document.
Mixed old/new nodes sharing an IAM backend and rolling downgrade are unsupported:
older binaries do not interpret tombstones or signed parent boundaries correctly.
Use a maintenance window for coordinated upgrade:
1. Pause IAM changes and isolate any offline site or backup whose state is unknown.
2. Back up each site's complete IAM storage and required encryption material. A
live IAM admin export omits deletion history and is not an adequate backup.
3. Stop all nodes sharing each site's IAM backend, replace their binaries, and
restart them on the upgraded version. Complete this for every participating
site before relying on the new revocation semantics.
4. Check IAM loading, site-replication errors and revocation convergence. Verify
representative old credentials are denied and deliberately reissued ones work.
5. Resolve pre-upgrade revocations explicitly. Absence cannot reconstruct an
already-lost deletion version: remove surviving old records on the sites that
still have them, and rebuild stale offline peers from approved state before
admitting them. Do not reconnect an unknown old snapshot just to discover its
deleted credentials.
Credentials issued by an older server for a recreated parent lack the required
signed boundary and must be reissued by an upgraded server. Parents without any
retained revocation history preserve existing credential behavior.
A pristine built-in policy remains protected from local deletion. If an
administrator explicitly overrides that policy and later deletes the override,
the durable deletion now suppresses automatic recreation of the built-in policy
on reload. This prevents reload from undoing the deletion. Restore the policy by
an explicit policy-create operation if desired. Local deletion of a nonexistent
policy remains idempotent and does not create a new tombstone; replicated
unknown deletions retain their version.
For rollback, stop and isolate the affected sites and assess changes since the
backup before restoring compatible state. Restoring an older backup can itself
lose later revocations and requires reconciliation/rekeying before access is
reopened. Do not delete tombstones online or convert only live IAM records to
make an older binary start. Server, client, Console, package and deployment
acceptance remain separate delivery gates of the maintained PGSTY stack.
## Regression and performance checks
Focused coverage is in `iam-revocation_test.go`, `iam-revision_test.go`,
`iam-revision-lock_test.go`, `iam-revision-boundary_test.go`,
`iam-replication-protocol_test.go`, `iam-credential-retention_test.go`,
`iam-peer-reload_test.go`, and `iam-replay_test.go`. Set
`SILO_TEST_IAM_REVOCATION_ETCD` to a disposable etcd endpoint to include backend
lifecycle/locking/boundary tests; they use isolated key namespaces.
`BenchmarkIAMCachedCredential` and `BenchmarkIAMSetTempUser` can be run against
the pre-change source for a comparable local baseline.
`BenchmarkIAMRevisionConvergedHealing` covers 1,000 and 10,000 retained records;
it measures steady-state index/network work and asserts zero repeated PUTs. Its
fake peer does not measure durable catch-up throughput.
`BenchmarkIAMColdLoadExpiredServices` measures loading and cleanup with 100 or
1,000 expired reusable credentials; run it with `-benchtime=1x`. Use actual multi-site
signed S3/STS tests and deployment-specific latency/scale measurements in
addition to these component tests.