diff --git a/weed/s3api/s3api_object_versioning.go b/weed/s3api/s3api_object_versioning.go index bffc531ef..8ef46b7cc 100644 --- a/weed/s3api/s3api_object_versioning.go +++ b/weed/s3api/s3api_object_versioning.go @@ -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. // diff --git a/weed/s3api/s3api_object_versioning_heal_test.go b/weed/s3api/s3api_object_versioning_heal_test.go new file mode 100644 index 000000000..cb8b4b690 --- /dev/null +++ b/weed/s3api/s3api_object_versioning_heal_test.go @@ -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)) + }) + } +}