fix(s3): make self-heal pointer persist CAS-bound against concurrent writers (#11627)

* fix(s3): make self-heal pointer persist CAS-bound against concurrent writers

Follow-up to #11618: pointerless reads of a slash key whose regular-path
entry is a physical parent (or a bare-key object) now fall through to
healStaleLatestVersionPointer, which rescans .versions and persists a
repaired pointer. The persist was an unconditional upsert off the pre-scan
snapshot, so a PUT or delete that atomically advanced the pointer on the
owner filer while the heal was rescanning could be rolled back, making
older content or ACLs current again.

Mirror the CAS discipline clearStaleLatestVersionPointer already applies:
re-fetch the live .versions entry, require its pointer fields to still
match the ones the heal observed, and abandon the persist (still returning
the rescanned entry) when a concurrent writer has moved them. Write the
live Extended map so concurrently updated fields are preserved.

* fix(s3): close the check-then-act window in the self-heal pointer persist

The CAS re-fetch added in the previous commit narrows the race but leaves
a gateway-side window: after the live .versions entry is re-read and the
pointer compared, the repair is still written back through an
unconditional RPC, so a PUT or delete committing between the re-fetch and
the persist still ends up rolled back by the stale repair.

Bind the persist to the live image the heal just re-read with an
IF_ENTRY_EQUAL precondition, the same discipline routedSelfCopy applies
to stale self-copies: the filer evaluates the condition under the entry's
path lock and conditional writes route to the owner filer, so a writer
committing inside the window fails the precondition and the winner's
pointer stands. FailedPrecondition and NotFound are authoritative replies
and are not replayed by the failover layer.

The test now also covers a writer committing during the persist, which
reverts the pointer on the previous unconditional write-back.

* s3api: CAS-bind the stale-pointer clear against concurrent writers

The pointer clear re-read the live .versions entry and then wrote it
back unconditionally through mkFile, so a writer committing between the
re-fetch and the persist was rolled back to a cleared pointer. Persist
through the same IF_ENTRY_EQUAL conditional update as the repair path.

* s3api: test the CAS contract on the stale-pointer clear

* s3api: never clear a pointer the clear did not observe as stale

The CAS clear skipped its live-pointer match when the caller's snapshot
carried an empty latest-version id, so a writer promoting a version
between the clear's rescan and its re-fetch had the fresh pointer
CAS-cleared away (expected = the writer's own live entry), briefly
making the just-written version appear absent. With an empty observed
id, reaching the persist at all implies a concurrent promotion (an idle
key short-circuits as already-clear), so make the pointer match
unconditional and abort instead.

Extend TestClearStaleLatestVersionPointerConcurrentWriter with
pointerless-snapshot cases: a post-rescan promotion must survive, and an
idle pointerless key must short-circuit as already-clear. The fake
filer's proto round-trip drops empty Extended maps, so the snapshot is
padded the way real callers do.

---------

Co-authored-by: zhaoyuchen <yc.zhao@yinzon.com>
Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
This commit is contained in:
authored and GitHub committed 2026-10-07 20:04:09 +08:00
1 parent 3288d90b21
commit 576837f1b2
2 files changed
+333 -28

No files matched your search

+119 -28
View File
@@ -26,6 +26,7 @@ import (
"github.com/seaweedfs/seaweedfs/weed/util"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/proto"
)
// ErrDeleteMarker is returned when the latest version is a delete marker (expected condition)
@@ -1850,20 +1851,34 @@ func (s3a *S3ApiServer) clearStaleLatestVersionPointer(bucket, object, bucketDir
// Already cleared by another path.
return true
}
if observedStaleId != "" && string(currentIdBytes) != observedStaleId {
// The live pointer must still be the stale one this clear observed. With an
// empty observed id, any pointer the re-fetch sees was promoted by a
// concurrent writer between the rescan and here — a pointer this clear
// never set out to remove, so leave it alone (the rescan's no-tagged-
// versions result only covered the pre-promotion directory).
if string(currentIdBytes) != observedStaleId {
glog.V(1).Infof("%s: skipping pointer clear for %s/%s, live pointer changed (observed=%s, current=%s)", caller, bucket, object, observedStaleId, string(currentIdBytes))
return false
}
delete(liveEntry.Extended, s3_constants.ExtLatestVersionIdKey)
delete(liveEntry.Extended, s3_constants.ExtLatestVersionFileNameKey)
clearCachedVersionMetadata(liveEntry.Extended)
if mkErr := s3a.mkFile(bucketDir, versionsObjectPath, liveEntry.Chunks, func(updatedEntry *filer_pb.Entry) {
updatedEntry.Extended = liveEntry.Extended
updatedEntry.Attributes = liveEntry.Attributes
updatedEntry.Chunks = liveEntry.Chunks
}); mkErr != nil {
versioningHealWarningf("clear_failed", "bucket=%s key=%s caller=%s err=%v", bucket, object, caller, mkErr)
expected := proto.Clone(liveEntry).(*filer_pb.Entry)
updated := proto.Clone(liveEntry).(*filer_pb.Entry)
delete(updated.Extended, s3_constants.ExtLatestVersionIdKey)
delete(updated.Extended, s3_constants.ExtLatestVersionFileNameKey)
clearCachedVersionMetadata(updated.Extended)
updated.Name = versionsObjectPath
if err := s3a.conditionalUpdateEntry(bucketDir, updated, expected); err != nil {
switch status.Code(err) {
case codes.FailedPrecondition:
// A writer moved the pointer after this heal's snapshot; its
// pointer stands.
glog.V(1).Infof("%s: skipping pointer clear for %s/%s, live entry changed", caller, bucket, object)
case codes.NotFound:
// The .versions entry vanished; nothing left to clear.
return true
default:
versioningHealWarningf("clear_failed", "bucket=%s key=%s caller=%s err=%v", bucket, object, caller, err)
}
return false
}
versioningHealInfof("healed", "bucket=%s key=%s mode=pointer_cleared caller=%s (orphan entries remain in .versions directory)", bucket, object, caller)
@@ -2144,11 +2159,13 @@ func scanLatestVersionEntry(list entryLister, versionsDir string) (latestEntry *
// healStaleLatestVersionPointer is invoked when the .versions directory metadata
// points to a version file that no longer exists. It paginates the directory,
// picks the chronologically newest remaining entry (content version or delete
// marker), updates the directory pointer metadata best-effort, and returns the
// rescanned entry. Downstream handlers detect ExtDeleteMarkerKey on the
// returned entry and render NoSuchKey, so promoting a delete marker preserves
// correct S3 semantics. If no version-tagged entry remains an error is
// returned and the caller surfaces it as not found.
// marker), updates the directory pointer metadata best-effort — atomically
// guarded by an IF_ENTRY_EQUAL precondition, so a pointer a concurrent writer
// advanced in the meantime is left alone — and returns the rescanned entry.
// Downstream handlers detect ExtDeleteMarkerKey on the returned entry and
// render NoSuchKey, so promoting a delete marker preserves correct S3
// semantics. If no version-tagged entry remains an error is returned and the
// caller surfaces it as not found.
func (s3a *S3ApiServer) healStaleLatestVersionPointer(bucket, normalizedObject string, versionsEntry *filer_pb.Entry, stalePointerFile string) (*filer_pb.Entry, error) {
bucketDir := s3a.bucketDir(bucket)
versionsObjectPath := normalizedObject + s3_constants.VersionsFolder
@@ -2194,28 +2211,102 @@ func (s3a *S3ApiServer) healStaleLatestVersionPointer(bucket, normalizedObject s
return nil, fmt.Errorf("%w: no remaining version in %s", filer_pb.ErrNotFound, versionsDir)
}
if versionsEntry.Extended == nil {
versionsEntry.Extended = make(map[string][]byte)
// The persist is CAS-style (mirroring clearStaleLatestVersionPointer) so a
// stale repair cannot supersede a concurrent writer: the rescan takes time,
// during which a PUT or delete may have atomically advanced the pointer on
// the owner filer. Re-fetch the live .versions entry and require its pointer
// fields to still match the ones this heal observed when it decided the
// pointer was unusable. If they moved, abandon the persist — this read
// still returns the scanned entry, and the winner's pointer needs no
// repair. Write the live Extended map so fields a concurrent writer
// updated between the re-fetch and the persist are preserved.
liveEntry, liveErr := s3a.getEntry(bucketDir, versionsObjectPath)
if liveErr != nil {
// The directory is gone (e.g. a concurrent delete's teardown); nothing
// to repair, and this read still returns the scanned entry.
versioningHealInfof("abandoned", "bucket=%s key=%s mode=versions_dir_gone err=%v", bucket, normalizedObject, liveErr)
return latestEntry, nil
}
observedId := string(versionsEntry.Extended[s3_constants.ExtLatestVersionIdKey])
observedFile := string(versionsEntry.Extended[s3_constants.ExtLatestVersionFileNameKey])
liveId := string(liveEntry.Extended[s3_constants.ExtLatestVersionIdKey])
liveFile := string(liveEntry.Extended[s3_constants.ExtLatestVersionFileNameKey])
if liveId != observedId || liveFile != observedFile {
// A concurrent writer promoted the pointer between the heal's snapshot
// and this check; rolling it back to the scanned version would make
// older content or ACLs current again.
versioningHealInfof("abandoned", "bucket=%s key=%s mode=pointer_moved observed_id=%s live_id=%s", bucket, normalizedObject, observedId, liveId)
return latestEntry, nil
}
versionsEntry.Extended[s3_constants.ExtLatestVersionIdKey] = []byte(latestVersionId)
versionsEntry.Extended[s3_constants.ExtLatestVersionFileNameKey] = []byte(latestVersionFileName)
setCachedListMetadata(versionsEntry, latestEntry)
if mkErr := s3a.mkFile(bucketDir, versionsObjectPath, versionsEntry.Chunks, func(updatedEntry *filer_pb.Entry) {
updatedEntry.Extended = versionsEntry.Extended
updatedEntry.Attributes = versionsEntry.Attributes
updatedEntry.Chunks = versionsEntry.Chunks
}); mkErr != nil {
// Clone the live image before any mutation so the precondition compares
// against exactly what the filer stores (nil and empty Extended maps are
// not proto-equal).
expected := proto.Clone(liveEntry).(*filer_pb.Entry)
updated := proto.Clone(liveEntry).(*filer_pb.Entry)
if updated.Extended == nil {
updated.Extended = make(map[string][]byte)
}
updated.Name = versionsObjectPath
updated.Extended[s3_constants.ExtLatestVersionIdKey] = []byte(latestVersionId)
updated.Extended[s3_constants.ExtLatestVersionFileNameKey] = []byte(latestVersionFileName)
setCachedListMetadata(updated, latestEntry)
updateErr := s3a.conditionalUpdateEntry(bucketDir, updated, expected)
switch {
case updateErr == nil:
versioningHealInfof("healed", "bucket=%s key=%s mode=pointer_repaired new_version=%s file=%s delete_marker=%v", bucket, normalizedObject, latestVersionId, latestVersionFileName, isDeleteMarker)
case status.Code(updateErr) == codes.FailedPrecondition, status.Code(updateErr) == codes.NotFound:
// The filer refused the stale repair: the live entry changed between
// the re-fetch and the persist, so the winner's pointer stands.
versioningHealInfof("abandoned", "bucket=%s key=%s mode=entry_changed err=%v", bucket, normalizedObject, updateErr)
default:
// Persisting the repair is best-effort. Surface a warning but still
// return the rescanned entry so the read succeeds; a subsequent write
// on the object will persist a fresh pointer.
versioningHealWarningf("heal_persist_failed", "bucket=%s key=%s err=%v (returning rescanned entry)", bucket, normalizedObject, mkErr)
} else {
versioningHealInfof("healed", "bucket=%s key=%s mode=pointer_repaired new_version=%s file=%s delete_marker=%v", bucket, normalizedObject, latestVersionId, latestVersionFileName, isDeleteMarker)
versioningHealWarningf("heal_persist_failed", "bucket=%s key=%s err=%v (returning rescanned entry)", bucket, normalizedObject, updateErr)
}
return latestEntry, nil
}
// conditionalUpdateEntry persists entry only while the stored image still
// equals expected: the filer evaluates IF_ENTRY_EQUAL under the entry's path
// lock, and conditional writes route to the owner filer, so a writer
// committing between the caller's snapshot and this update fails the
// precondition instead of being rolled back. FailedPrecondition and NotFound
// are authoritative replies, not transport failures — they return without
// retrying or marking the filer down.
func (s3a *S3ApiServer) conditionalUpdateEntry(directory string, entry, expected *filer_pb.Entry) error {
req := &filer_pb.UpdateEntryRequest{
Directory: directory,
Entry: entry,
Condition: &filer_pb.WriteCondition{Clauses: []*filer_pb.WriteCondition_Clause{{
Kind: filer_pb.WriteCondition_IF_ENTRY_EQUAL,
ExpectedEntry: expected,
}}},
}
var conditionErr error
update := func(client filer_pb.SeaweedFilerClient) error {
err := filer_pb.UpdateEntry(context.Background(), client, req)
if wrapped := errors.Unwrap(err); wrapped != nil {
// The helper wraps the RPC error ("UpdateEntry: %w"); unwrap so
// the status code below survives the wrapping.
err = wrapped
}
switch status.Code(err) {
case codes.FailedPrecondition, codes.NotFound:
conditionErr = err
return nil
default:
return err
}
}
if err := s3a.WithFilerClient(false, update); err != nil {
return err
}
return conditionErr
}
// getLatestVersionEntryFromDirectoryEntry creates a logical entry for list operations using cached metadata
// from the .versions directory entry. This achieves SINGLE-SCAN efficiency - no additional getEntry calls needed.
//
@@ -0,0 +1,214 @@
package s3api
import (
"context"
"net/http"
"path"
"testing"
"time"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
"github.com/stretchr/testify/require"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/proto"
)
// anonymousReadHealFiler drives the versioned self-heal path for "folder/": it
// serves a real version file from .versions, runs hooks at the heal's rescan
// and at the persist (simulating a writer committing at either point), and
// evaluates the heal's IF_ENTRY_EQUAL precondition like a filer would, so
// tests can inspect the resulting .versions metadata.
type anonymousReadHealFiler struct {
*anonymousReadFiler
onScan func(*anonymousReadHealFiler)
onWrite func(*anonymousReadHealFiler)
}
// ListEntries serves the version files under .versions for the heal's rescan.
func (f *anonymousReadHealFiler) ListEntries(req *filer_pb.ListEntriesRequest, stream filer_pb.SeaweedFiler_ListEntriesServer) error {
f.mu.Lock()
if f.onScan != nil {
onScan := f.onScan
f.onScan = nil
f.mu.Unlock()
onScan(f)
} else {
f.mu.Unlock()
}
f.mu.Lock()
defer f.mu.Unlock()
for _, entry := range f.entries {
if path.Join(req.Directory, entry.Name) == "/buckets/b/folder/.versions/v_v1" {
if err := stream.Send(&filer_pb.ListEntriesResponse{Entry: proto.Clone(entry).(*filer_pb.Entry)}); err != nil {
return err
}
}
}
return nil
}
// UpdateEntry evaluates the heal's IF_ENTRY_EQUAL precondition the way a
// filer does under the path lock: the stored entry must still equal the live
// image the heal re-read, otherwise the stale repair is rejected.
func (f *anonymousReadHealFiler) UpdateEntry(_ context.Context, req *filer_pb.UpdateEntryRequest) (*filer_pb.UpdateEntryResponse, error) {
f.mu.Lock()
if f.onWrite != nil {
onWrite := f.onWrite
f.onWrite = nil
f.mu.Unlock()
// A writer commits between the heal's re-fetch and the persist.
onWrite(f)
} else {
f.mu.Unlock()
}
f.mu.Lock()
defer f.mu.Unlock()
fullPath := path.Join(req.Directory, req.Entry.Name)
current := f.entries[fullPath]
for _, clause := range req.Condition.Clauses {
if clause.Kind != filer_pb.WriteCondition_IF_ENTRY_EQUAL {
return nil, status.Error(codes.Unimplemented, "unsupported clause")
}
expected := proto.Clone(clause.ExpectedEntry).(*filer_pb.Entry)
actual := proto.Clone(current).(*filer_pb.Entry)
filer_pb.BeforeEntrySerialization(expected.Chunks)
filer_pb.BeforeEntrySerialization(actual.Chunks)
if !proto.Equal(actual, expected) {
return nil, status.Error(codes.FailedPrecondition, "entry changed")
}
}
f.entries[fullPath] = proto.Clone(req.Entry).(*filer_pb.Entry)
return &filer_pb.UpdateEntryResponse{}, nil
}
// TestAnonymousObjectACLHealPointerConcurrentWriter pins the CAS contract of
// the self-heal persist: a pointer that a concurrent writer promoted while the
// heal was rescanning, or between the heal's re-fetch and its persist, must
// survive, instead of being rolled back to the scanned version, which would
// make older content or ACLs current again.
func TestAnonymousObjectACLHealPointerConcurrentWriter(t *testing.T) {
for _, tc := range []struct {
name string
writerDuringScan bool
writerDuringWrite bool
wantPointer string
wantPointerVersion string
}{
{name: "idle key persists the rescanned pointer", wantPointer: "v1", wantPointerVersion: "v_v1"},
{name: "concurrent writer during scan keeps its promotion", writerDuringScan: true, wantPointer: "v2", wantPointerVersion: "v_v2"},
{name: "concurrent writer during persist keeps its promotion", writerDuringWrite: true, wantPointer: "v3", wantPointerVersion: "v_v3"},
} {
t.Run(tc.name, func(t *testing.T) {
f := &anonymousReadHealFiler{anonymousReadFiler: &anonymousReadFiler{entries: make(map[string]*filer_pb.Entry)}}
// The object is "folder/"; the regular path is a physical parent, so a
// pointerless read falls through to the persistent self-heal.
f.entries["/buckets/b/folder"] = &filer_pb.Entry{Name: "folder", IsDirectory: true, Attributes: &filer_pb.FuseAttributes{Mtime: 1700000000}}
f.entries["/buckets/b/folder/child"] = anonymousReadEntry([]byte(`[]`))
f.entries["/buckets/b/folder/.versions"] = &filer_pb.Entry{Name: ".versions", IsDirectory: true, Attributes: &filer_pb.FuseAttributes{Mtime: 1700000000}, Extended: map[string][]byte{s3_constants.ExtLatestVersionIdKey: []byte("")}}
v1 := anonymousReadEntry([]byte(`[]`))
v1.Name = "v_v1"
v1.Extended[s3_constants.ExtVersionIdKey] = []byte("v1")
f.entries["/buckets/b/folder/.versions/v_v1"] = v1
promotePointer := func(versionId, fileName string) func(*anonymousReadHealFiler) {
return func(f *anonymousReadHealFiler) {
// A concurrent PUT commits a newer version.
f.mu.Lock()
defer f.mu.Unlock()
f.entries["/buckets/b/folder/.versions"].Extended[s3_constants.ExtLatestVersionIdKey] = []byte(versionId)
f.entries["/buckets/b/folder/.versions"].Extended[s3_constants.ExtLatestVersionFileNameKey] = []byte(fileName)
}
}
if tc.writerDuringScan {
f.onScan = promotePointer("v2", "v_v2")
}
if tc.writerDuringWrite {
f.onWrite = promotePointer("v3", "v_v3")
}
s3a := newPutTestServer(t, startFakeFiler(t, f))
s3a.iam = &IdentityAccessManagement{isAuthEnabled: true}
s3a.bucketConfigCache = NewBucketConfigCache(time.Minute)
s3a.bucketConfigCache.Set("b", &BucketConfig{Name: "b", Ownership: s3_constants.OwnershipObjectWriter, Versioning: "Enabled"})
s3a.policyEngine = NewBucketPolicyEngine()
s3a.iam.policyEngine = s3a.policyEngine
require.NoError(t, s3a.policyEngine.engine.SetBucketPolicy("b", `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject*","Resource":"arn:aws:s3:::b/folder/"}]}`))
rr := serveAnonymousRead(s3a, http.MethodGet, "folder/", "", nil)
require.Equal(t, http.StatusOK, rr.Code, rr.Body.String())
// This read returns the rescanned entry either way.
require.Equal(t, "hello world", rr.Body.String())
// The persisted pointer must reflect the concurrent writer, not the scan.
pointer := f.entries["/buckets/b/folder/.versions"].Extended[s3_constants.ExtLatestVersionIdKey]
require.Equal(t, tc.wantPointer, string(pointer))
pointerFile := f.entries["/buckets/b/folder/.versions"].Extended[s3_constants.ExtLatestVersionFileNameKey]
require.Equal(t, tc.wantPointerVersion, string(pointerFile))
})
}
}
// TestClearStaleLatestVersionPointerConcurrentWriter pins the same CAS
// contract on the pointer clear: a writer that promotes the pointer between
// the clear's re-fetch and its persist must not be rolled back to a cleared
// pointer. A pointerless caller snapshot must not clear a pointer a writer
// promoted after the clear's rescan either — the only live pointer such a
// clear can observe is one that appeared concurrently, never the stale one
// the clear set out to remove.
func TestClearStaleLatestVersionPointerConcurrentWriter(t *testing.T) {
for _, tc := range []struct {
name string
observedEmpty bool
writerDuringScan bool
writerDuringWrite bool
wantCleared bool
wantPointer string
}{
{name: "idle key clears the stale pointer", wantCleared: true, wantPointer: ""},
{name: "concurrent writer during persist keeps its promotion", writerDuringWrite: true, wantCleared: false, wantPointer: "v2"},
{name: "pointerless snapshot keeps a post-rescan promotion", observedEmpty: true, writerDuringScan: true, wantCleared: false, wantPointer: "v2"},
{name: "pointerless snapshot on an idle key is already clear", observedEmpty: true, wantCleared: true, wantPointer: ""},
} {
t.Run(tc.name, func(t *testing.T) {
f := &anonymousReadHealFiler{anonymousReadFiler: &anonymousReadFiler{entries: make(map[string]*filer_pb.Entry)}}
extended := map[string][]byte{
s3_constants.ExtLatestVersionIdKey: []byte("v1"),
s3_constants.ExtLatestVersionFileNameKey: []byte("v_v1"),
}
if tc.observedEmpty {
extended = map[string][]byte{}
}
f.entries["/buckets/b/folder/.versions"] = &filer_pb.Entry{
Name: ".versions",
IsDirectory: true,
Attributes: &filer_pb.FuseAttributes{Mtime: 1700000000},
Extended: extended,
}
promotePointer := func(f *anonymousReadHealFiler) {
f.mu.Lock()
defer f.mu.Unlock()
f.entries["/buckets/b/folder/.versions"].Extended[s3_constants.ExtLatestVersionIdKey] = []byte("v2")
f.entries["/buckets/b/folder/.versions"].Extended[s3_constants.ExtLatestVersionFileNameKey] = []byte("v_v2")
}
if tc.writerDuringScan {
f.onScan = promotePointer
}
if tc.writerDuringWrite {
f.onWrite = promotePointer
}
s3a := newPutTestServer(t, startFakeFiler(t, f))
versionsEntry := proto.Clone(f.entries["/buckets/b/folder/.versions"]).(*filer_pb.Entry)
if versionsEntry.Extended == nil {
// proto.Clone turns an empty Extended map into nil; real callers
// (updateLatestVersionAfterDeletion) hand the clear a non-nil
// snapshot, so mirror that here.
versionsEntry.Extended = map[string][]byte{}
}
cleared := s3a.clearStaleLatestVersionPointer("b", "folder", "/buckets/b", "folder/.versions", versionsEntry, "test")
require.Equal(t, tc.wantCleared, cleared)
pointer := f.entries["/buckets/b/folder/.versions"].Extended[s3_constants.ExtLatestVersionIdKey]
require.Equal(t, tc.wantPointer, string(pointer))
})
}
}