mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-19 04:50:54 +02:00
s3api: delete orphaned chunks only when the entry is confirmed absent (#11389)
* s3api: test for chunks deleted under an entry the filer committed Issue #11387: the filer can report a create failure after inserting the entry (e.g. a parent-directory creation failing post-insert). The error arrives in the response rather than as a transport status, so it maps to a definitive error and putToFiler deletes the chunks of the live entry. * s3api: confirmCreateLanded also reports a confirmed-absent entry The verification a failed create runs can answer both directions: the entry matching the uploaded chunks proves the write landed, and an authoritative not-found proves the uploaded chunks are orphaned. Return both outcomes so the cleanup path can gate on the fact rather than the error class. An empty upload can never prove a landing, so a zero-chunk entry match no longer upgrades the outcome. * s3api: delete orphaned chunks only when the entry is confirmed absent A failed create no longer skips verification based on the error class: the filer can fail after inserting the entry (issue #11387) and a partially-applied routed transaction can leave it behind too, both surfacing as definitive errors. Every failed create now resolves the entry's fate, and the uploaded chunks are deleted only when the entry is confirmed absent; anything unverifiable keeps them for vacuum. * s3api: confirm absence on every filer the create could have committed on A lock-path create fails over across filers, so the entry can live on a replica the routed owner has not caught up to; one not-found does not prove absence. The confirmation now queries the owner, the prior owner, and the failover set, declaring absent only when none of them has the entry. * s3api: bound the reconciliation lookups confirmCreateLanded runs The lookups ran on context.Background() under the object write lock, so a connected filer that never replies could stall the write path. One timeout now covers the whole enumeration; an expired budget fails the remaining lookups as uncertain, which keeps the chunks.
This commit is contained in:
@@ -22,6 +22,7 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/operation"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/s3_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
@@ -995,15 +996,15 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader
|
||||
|
||||
// If the entry was never created, the uploaded chunks are orphaned and must be deleted.
|
||||
if !entryCreated {
|
||||
// A transport failure is ambiguous: the filer may have committed the
|
||||
// entry anyway (issue #11366), so a retryable error never deletes the
|
||||
// uploaded chunks — it is only upgraded to success when the write owner
|
||||
// proves the entry landed with these chunks.
|
||||
ambiguous := createErr != nil && filerErrorToS3Error(createErr) == s3err.ErrServiceUnavailable
|
||||
if ambiguous && len(chunkResult.FileChunks) > 0 && s3a.confirmCreateLanded(filePath, bucket, object, entry, chunkResult.FileChunks, finalize) {
|
||||
// A failed create does not prove the entry is absent: a lost response
|
||||
// can hide a commit (issue #11366) and the filer can fail after
|
||||
// inserting the entry (issue #11387), so the entry's presence — not
|
||||
// the error class — decides the chunks' fate.
|
||||
landed, absent := s3a.confirmCreateLanded(filePath, bucket, object, entry, chunkResult.FileChunks, finalize)
|
||||
if landed {
|
||||
createCode = s3err.ErrNone
|
||||
}
|
||||
if createCode != s3err.ErrNone && !ambiguous {
|
||||
if createCode != s3err.ErrNone && absent {
|
||||
orphaned := chunkResult.FileChunks
|
||||
if manifestChunks, _ := filer.SeparateManifestChunks(entry.GetChunks()); len(manifestChunks) > 0 {
|
||||
orphaned = append(manifestChunks, orphaned...)
|
||||
@@ -1044,28 +1045,51 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader
|
||||
return etag, s3err.ErrNone, responseMetadata
|
||||
}
|
||||
|
||||
// confirmCreateLanded checks whether a create that failed ambiguously still
|
||||
// committed: the stored entry's resolved chunks must be exactly the uploaded
|
||||
// ones. On a match the finalization the error skipped runs under the object
|
||||
// write lock, and true reports the write as successful.
|
||||
func (s3a *S3ApiServer) confirmCreateLanded(filePath, bucket, object string, entry *filer_pb.Entry, uploaded []*filer_pb.FileChunk, finalize *putFinalize) bool {
|
||||
// createLookupTimeout bounds the entry lookups confirmCreateLanded runs under
|
||||
// the object write lock, so a hung filer cannot stall the write path.
|
||||
const createLookupTimeout = 10 * time.Second
|
||||
|
||||
// confirmCreateLanded resolves a create whose outcome is uncertain: a stored
|
||||
// entry resolving to the uploaded chunks confirms the write landed — the
|
||||
// finalization the error skipped then runs under the object write lock, and
|
||||
// landed reports success — while absent requires every filer the create could
|
||||
// have committed on to lack the entry, the only outcome where the uploaded
|
||||
// chunks are orphaned.
|
||||
func (s3a *S3ApiServer) confirmCreateLanded(filePath, bucket, object string, entry *filer_pb.Entry, uploaded []*filer_pb.FileChunk, finalize *putFinalize) (landed, absent bool) {
|
||||
dir, name := path.Dir(filePath), path.Base(filePath)
|
||||
owner := s3a.routableWriteOwner(bucket, object)
|
||||
confirmed := false
|
||||
lookupCtx, cancel := context.WithTimeout(context.Background(), createLookupTimeout)
|
||||
defer cancel()
|
||||
// Verify, finalize, and roll back inside one critical section: a concurrent
|
||||
// write to the same key must not slip in between them.
|
||||
s3a.withObjectWriteLock(bucket, object, nil, func() s3err.ErrorCode {
|
||||
existing, lookupErr := s3a.lookupEntryPreferringOwner(owner, dir, name)
|
||||
if lookupErr != nil || existing == nil {
|
||||
var existing *filer_pb.Entry
|
||||
uncertain, queried := false, false
|
||||
for _, target := range s3a.createTargetFilers(owner, bucket, object) {
|
||||
e, lookupErr := s3a.lookupEntryOnFiler(lookupCtx, target, dir, name)
|
||||
queried = true
|
||||
if e != nil {
|
||||
existing = e
|
||||
break
|
||||
}
|
||||
if lookupErr != nil && !errors.Is(lookupErr, filer_pb.ErrNotFound) {
|
||||
uncertain = true
|
||||
}
|
||||
}
|
||||
if existing == nil {
|
||||
absent = queried && !uncertain
|
||||
return s3err.ErrNone
|
||||
}
|
||||
resolved, _, resolveErr := filer.ResolveChunkManifest(context.Background(), s3a.createLookupFileIdFunction(), existing.GetChunks(), 0, math.MaxInt64, s3a.filerClient)
|
||||
if len(uploaded) == 0 {
|
||||
return s3err.ErrNone
|
||||
}
|
||||
resolved, _, resolveErr := filer.ResolveChunkManifest(lookupCtx, s3a.createLookupFileIdFunction(), existing.GetChunks(), 0, math.MaxInt64, s3a.filerClient)
|
||||
if resolveErr != nil || !sameFileChunks(resolved, uploaded) {
|
||||
return s3err.ErrNone
|
||||
}
|
||||
glog.Warningf("putToFiler: create entry for %s failed but the entry exists, treating the write as successful", filePath)
|
||||
if finalize == nil || finalize.afterCreate == nil {
|
||||
confirmed = true
|
||||
landed = true
|
||||
return s3err.ErrNone
|
||||
}
|
||||
if code := finalize.afterCreate(entry); code != s3err.ErrNone {
|
||||
@@ -1075,10 +1099,36 @@ func (s3a *S3ApiServer) confirmCreateLanded(filePath, bucket, object string, ent
|
||||
}
|
||||
return s3err.ErrNone
|
||||
}
|
||||
confirmed = true
|
||||
landed = true
|
||||
return s3err.ErrNone
|
||||
})
|
||||
return confirmed
|
||||
return landed, absent
|
||||
}
|
||||
|
||||
// createTargetFilers lists the filers a failed create could have committed on:
|
||||
// the routed owner and the prior one mid-rebalance first, then the failover
|
||||
// set the lock path dials. Deduped, empty addresses skipped.
|
||||
func (s3a *S3ApiServer) createTargetFilers(owner pb.ServerAddress, bucket, object string) []pb.ServerAddress {
|
||||
var filers []pb.ServerAddress
|
||||
seen := map[pb.ServerAddress]bool{}
|
||||
add := func(f pb.ServerAddress) {
|
||||
if f != "" && !seen[f] {
|
||||
seen[f] = true
|
||||
filers = append(filers, f)
|
||||
}
|
||||
}
|
||||
add(owner)
|
||||
add(s3a.priorWriteOwner(bucket, object))
|
||||
if s3a.filerClient != nil {
|
||||
add(s3a.filerClient.GetCurrentFiler())
|
||||
for _, f := range s3a.filerClient.GetAllFilers() {
|
||||
add(f)
|
||||
}
|
||||
}
|
||||
for _, f := range s3a.option.Filers {
|
||||
add(f)
|
||||
}
|
||||
return filers
|
||||
}
|
||||
|
||||
// sameFileChunks reports whether two chunk lists reference the same needles,
|
||||
|
||||
@@ -88,6 +88,7 @@ type ambiguousPutFiler struct {
|
||||
entries map[string]*filer_pb.Entry
|
||||
apply bool
|
||||
createErr error
|
||||
respError string
|
||||
lookupErr error
|
||||
lookupFailKey string
|
||||
nextKey uint64
|
||||
@@ -119,7 +120,7 @@ func (f *ambiguousPutFiler) CreateEntry(_ context.Context, req *filer_pb.CreateE
|
||||
if f.createErr != nil {
|
||||
return nil, f.createErr
|
||||
}
|
||||
return &filer_pb.CreateEntryResponse{}, nil
|
||||
return &filer_pb.CreateEntryResponse{Error: f.respError}, nil
|
||||
}
|
||||
|
||||
func (f *ambiguousPutFiler) LookupDirectoryEntry(_ context.Context, req *filer_pb.LookupDirectoryEntryRequest) (*filer_pb.LookupDirectoryEntryResponse, error) {
|
||||
@@ -148,16 +149,16 @@ func (f *ambiguousPutFiler) LookupVolume(_ context.Context, req *filer_pb.Lookup
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func newPutTestServer(t *testing.T, filerAddr pb.ServerAddress) *S3ApiServer {
|
||||
func newPutTestServer(t *testing.T, filerAddrs ...pb.ServerAddress) *S3ApiServer {
|
||||
t.Helper()
|
||||
dialOption := grpc.WithTransportCredentials(insecure.NewCredentials())
|
||||
return &S3ApiServer{
|
||||
option: &S3ApiServerOption{
|
||||
Filers: []pb.ServerAddress{filerAddr},
|
||||
Filers: filerAddrs,
|
||||
GrpcDialOption: dialOption,
|
||||
BucketsPath: "/buckets",
|
||||
},
|
||||
filerClient: wdclient.NewFilerClient([]pb.ServerAddress{filerAddr}, dialOption, ""),
|
||||
filerClient: wdclient.NewFilerClient(filerAddrs, dialOption, ""),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -193,6 +194,64 @@ func TestPutToFilerAmbiguousCreateKeepsChunks(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// Issue 11387: the filer can report a create failure after inserting the entry
|
||||
// (e.g. a parent-directory creation failing post-insert). The failure arrives
|
||||
// in the response rather than as a transport status, so it maps to a
|
||||
// definitive error — but the entry exists and deleting its chunks would
|
||||
// tombstone live needles.
|
||||
func TestPutToFilerPostCommitErrorKeepsChunks(t *testing.T) {
|
||||
volume := startFakeVolumeServer(t)
|
||||
filerImpl := &ambiguousPutFiler{
|
||||
volume: volume,
|
||||
entries: map[string]*filer_pb.Entry{},
|
||||
apply: true,
|
||||
respError: "create parent directories of /buckets/b: i/o timeout",
|
||||
}
|
||||
s3a := newPutTestServer(t, startFakeFiler(t, filerImpl))
|
||||
|
||||
etag, code := putTestObject(t, s3a)
|
||||
if code != s3err.ErrNone {
|
||||
t.Fatalf("putToFiler returned %v, want success once the entry is confirmed on the filer", code)
|
||||
}
|
||||
if etag == "" {
|
||||
t.Fatal("expected an etag")
|
||||
}
|
||||
if deleted := volume.deleted(); len(deleted) != 0 {
|
||||
t.Fatalf("chunks under a live entry were deleted: %v", deleted)
|
||||
}
|
||||
}
|
||||
|
||||
// Issue 11387, multi-filer: a create that fails over mid-flight can commit on
|
||||
// a filer the confirmation does not ask first. A not-found from one replica
|
||||
// does not authorize deleting chunks an entry on another filer references.
|
||||
func TestPutToFilerPostCommitErrorOnFailoverFilerKeepsChunks(t *testing.T) {
|
||||
volume := startFakeVolumeServer(t)
|
||||
filerA := &ambiguousPutFiler{
|
||||
volume: volume,
|
||||
entries: map[string]*filer_pb.Entry{},
|
||||
apply: false,
|
||||
createErr: status.Error(codes.Unavailable, "connect: connection refused"),
|
||||
}
|
||||
filerB := &ambiguousPutFiler{
|
||||
volume: volume,
|
||||
entries: map[string]*filer_pb.Entry{},
|
||||
apply: true,
|
||||
respError: "create parent directories of /buckets/b: i/o timeout",
|
||||
}
|
||||
s3a := newPutTestServer(t, startFakeFiler(t, filerA), startFakeFiler(t, filerB))
|
||||
|
||||
etag, code := putTestObject(t, s3a)
|
||||
if code != s3err.ErrNone {
|
||||
t.Fatalf("putToFiler returned %v, want success once the entry is confirmed on the filer", code)
|
||||
}
|
||||
if etag == "" {
|
||||
t.Fatal("expected an etag")
|
||||
}
|
||||
if deleted := volume.deleted(); len(deleted) != 0 {
|
||||
t.Fatalf("chunks under a live entry were deleted: %v", deleted)
|
||||
}
|
||||
}
|
||||
|
||||
// A create the filer definitively refused still cleans up the uploaded chunks.
|
||||
func TestPutToFilerConfirmedFailureDeletesOrphans(t *testing.T) {
|
||||
volume := startFakeVolumeServer(t)
|
||||
|
||||
@@ -54,7 +54,7 @@ func (s3a *S3ApiServer) getObjectEntryRoutedByKey(bucket, object string) (*filer
|
||||
// prior owner once while the ring change is within the cooling-off window.
|
||||
if errors.Is(err, filer_pb.ErrNotFound) {
|
||||
if prior := s3a.priorWriteOwner(bucket, object); prior != "" && prior != owner {
|
||||
if priorEntry, priorErr := s3a.lookupEntryOnFiler(prior, dir, name); priorErr == nil {
|
||||
if priorEntry, priorErr := s3a.lookupEntryOnFiler(context.Background(), prior, dir, name); priorErr == nil {
|
||||
return priorEntry, nil
|
||||
}
|
||||
}
|
||||
@@ -87,14 +87,14 @@ func (s3a *S3ApiServer) lookupEntryPreferringOwner(owner pb.ServerAddress, dir,
|
||||
if owner == "" {
|
||||
return s3a.getEntry(dir, name)
|
||||
}
|
||||
return s3a.lookupEntryOnFiler(owner, dir, name)
|
||||
return s3a.lookupEntryOnFiler(context.Background(), owner, dir, name)
|
||||
}
|
||||
|
||||
// lookupEntryOnFiler resolves dir/name against a single filer, without failover.
|
||||
func (s3a *S3ApiServer) lookupEntryOnFiler(filer pb.ServerAddress, dir, name string) (*filer_pb.Entry, error) {
|
||||
func (s3a *S3ApiServer) lookupEntryOnFiler(ctx context.Context, filer pb.ServerAddress, dir, name string) (*filer_pb.Entry, error) {
|
||||
var entry *filer_pb.Entry
|
||||
err := pb.WithFilerClient(false, 0, filer, s3a.option.GrpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
|
||||
resp, lookupErr := filer_pb.LookupEntry(context.Background(), client, &filer_pb.LookupDirectoryEntryRequest{
|
||||
resp, lookupErr := filer_pb.LookupEntry(ctx, client, &filer_pb.LookupDirectoryEntryRequest{
|
||||
Directory: dir,
|
||||
Name: name,
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user