mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-11 16:57:45 +02:00
s3api: verify under the write lock before re-committing an ambiguous routed PUT (#11649)
* s3api: verify under the write lock before re-committing an ambiguous routed PUT A routed PUT whose response was lost after the owner committed, or whose transaction returned a store error, fell straight into the lock path and re-sent the same entry. A concurrent PUT that had superseded the commit in the meantime had already deleted this entry's chunks as old, so the re-commit restored metadata pointing at dead needles and the object read back 404 permanently. The lock path now resolves the routed attempt's outcome inside the write lock first: a stored entry with the uploaded chunks means the route landed; a different stored entry or an unresolved lookup refuses the re-commit with ServiceUnavailable so the client retries with a fresh upload; only a proven-absent entry falls through to the normal create. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * s3api: tighten ambiguous routed PUT recovery - a match found only after an unanswered filer is not authoritative; the unreachable filer may hold a newer entry, so refuse instead of finalizing on it (both recovery paths) - a failed post-recovery finalization now rolls the recovered entry back, matching the create path's undo - a proven-absent entry is only re-committed after probing that the uploaded chunks' needles are still alive; a committed-and-deleted PUT would otherwise write back an entry pointing at reclaimed needles - the .versions latest pointer no longer flips back to an older version when a newer one committed while the write was uncertain Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * s3api: stamp late versions noncurrent when a newer latest pointer wins The keep-newer early return in updateLatestVersionInDirectory left the version just stored without ExtNoncurrentSinceNsKey, so the lifecycle engine could never age it out. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * s3api: only a failure past mutation 0 makes a routed PUT ambiguous A deterministic refusal at the PUT mutation (e.g. "existing entry is a directory") applied nothing, but marking it ambiguous sent the lock-path fallback through conservative recovery, which found the directory entry and refused with 503 instead of the correct 409. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * s3: keep all transaction errors ambiguous; route directory conflicts through the lock path * s3: verify uploaded chunks before re-committing over a directory A stored directory does not prove the routed PUT never committed: the route can land, a delete can remove the entry and its chunks, and a nested write can recreate the directory before recovery takes the lock. Re-committing then stores an entry pointing at dead needles. Probe the uploaded chunks first and refuse with ServiceUnavailable when they can no longer be verified. --------- Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
1 parent
e337443176
commit
5d80ac39b8
3 files changed
+382
-17
No files matched your search
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"path"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
@@ -21,6 +22,7 @@ type fakeLookupFiler struct {
|
||||
filer_pb.UnimplementedSeaweedFilerServer
|
||||
entry *filer_pb.Entry
|
||||
lookupErr error
|
||||
deleted []string
|
||||
}
|
||||
|
||||
func (f *fakeLookupFiler) LookupDirectoryEntry(ctx context.Context, req *filer_pb.LookupDirectoryEntryRequest) (*filer_pb.LookupDirectoryEntryResponse, error) {
|
||||
@@ -30,6 +32,12 @@ func (f *fakeLookupFiler) LookupDirectoryEntry(ctx context.Context, req *filer_p
|
||||
return &filer_pb.LookupDirectoryEntryResponse{Entry: f.entry}, nil
|
||||
}
|
||||
|
||||
func (f *fakeLookupFiler) DeleteEntry(ctx context.Context, req *filer_pb.DeleteEntryRequest) (*filer_pb.DeleteEntryResponse, error) {
|
||||
f.deleted = append(f.deleted, path.Join(req.Directory, req.Name))
|
||||
f.entry = nil
|
||||
return &filer_pb.DeleteEntryResponse{}, nil
|
||||
}
|
||||
|
||||
func newHeadBucketTestServer(t *testing.T, impl filer_pb.SeaweedFilerServer) *S3ApiServer {
|
||||
t.Helper()
|
||||
filers := []pb.ServerAddress{startFakeFiler(t, impl)}
|
||||
|
||||
@@ -0,0 +1,212 @@
|
||||
package s3api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math"
|
||||
"path"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
"github.com/seaweedfs/seaweedfs/weed/wdclient"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
)
|
||||
|
||||
func ambiguousRouteServer(t *testing.T, addrs ...pb.ServerAddress) *S3ApiServer {
|
||||
t.Helper()
|
||||
dialOption := grpc.WithTransportCredentials(insecure.NewCredentials())
|
||||
return &S3ApiServer{
|
||||
option: &S3ApiServerOption{GrpcDialOption: dialOption, Filers: addrs},
|
||||
filerClient: wdclient.NewFilerClient(addrs, dialOption, ""),
|
||||
}
|
||||
}
|
||||
|
||||
func TestCreateAfterAmbiguousRoute(t *testing.T) {
|
||||
uploaded := []*filer_pb.FileChunk{{FileId: "5,01637037d6", Size: 55}}
|
||||
filePath := "/buckets/b/o"
|
||||
chunksAlive := func(context.Context, []*filer_pb.FileChunk) bool { return true }
|
||||
chunksDead := func(context.Context, []*filer_pb.FileChunk) bool { return false }
|
||||
newS3a := func(t *testing.T, f *fakeLookupFiler) (*S3ApiServer, pb.ServerAddress) {
|
||||
addr := startFakeFiler(t, f)
|
||||
return ambiguousRouteServer(t, addr), addr
|
||||
}
|
||||
|
||||
t.Run("stored entry is this PUT's — route landed", func(t *testing.T) {
|
||||
f := &fakeLookupFiler{entry: &filer_pb.Entry{Name: "o", Chunks: []*filer_pb.FileChunk{{FileId: "5,01637037d6", Size: 55}}}}
|
||||
s3a, owner := newS3a(t, f)
|
||||
entryCreated := false
|
||||
ran := false
|
||||
code := s3a.createAfterAmbiguousRoute(filePath, "b", "o", owner, &filer_pb.Entry{Name: "o"}, uploaded, nil, &entryCreated, chunksAlive, func() s3err.ErrorCode { ran = true; return s3err.ErrNone })
|
||||
if code != s3err.ErrNone || !entryCreated {
|
||||
t.Fatalf("code=%v entryCreated=%v", code, entryCreated)
|
||||
}
|
||||
if ran {
|
||||
t.Fatal("re-committed an entry the route already stored")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("stored entry is a newer write — refuse stale re-commit", func(t *testing.T) {
|
||||
f := &fakeLookupFiler{entry: &filer_pb.Entry{Name: "o", Chunks: []*filer_pb.FileChunk{{FileId: "5,01637037d7", Size: 55}}}}
|
||||
s3a, owner := newS3a(t, f)
|
||||
entryCreated := false
|
||||
ran := false
|
||||
code := s3a.createAfterAmbiguousRoute(filePath, "b", "o", owner, &filer_pb.Entry{Name: "o"}, uploaded, nil, &entryCreated, chunksAlive, func() s3err.ErrorCode { ran = true; return s3err.ErrNone })
|
||||
if code != s3err.ErrServiceUnavailable {
|
||||
t.Fatalf("code = %v, want ServiceUnavailable", code)
|
||||
}
|
||||
if ran || entryCreated {
|
||||
t.Fatal("stale entry was committed over a newer write")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("stored entry is a directory, chunks alive — lock path answers", func(t *testing.T) {
|
||||
f := &fakeLookupFiler{entry: &filer_pb.Entry{Name: "o", IsDirectory: true}}
|
||||
s3a, owner := newS3a(t, f)
|
||||
entryCreated := false
|
||||
ran := false
|
||||
code := s3a.createAfterAmbiguousRoute(filePath, "b", "o", owner, &filer_pb.Entry{Name: "o"}, uploaded, nil, &entryCreated, chunksAlive, func() s3err.ErrorCode { ran = true; return s3err.ErrExistingObjectIsDirectory })
|
||||
if code != s3err.ErrExistingObjectIsDirectory || !ran {
|
||||
t.Fatalf("code=%v ran=%v — a directory conflict with live chunks must reach createUnderLock, which maps it", code, ran)
|
||||
}
|
||||
})
|
||||
|
||||
// The routed PUT may have committed before a delete removed the entry and
|
||||
// its chunks and a nested write recreated the name as a directory; dead
|
||||
// chunks mean re-committing would store an entry pointing at them.
|
||||
t.Run("stored entry is a directory, chunks unverifiable — refuse", func(t *testing.T) {
|
||||
f := &fakeLookupFiler{entry: &filer_pb.Entry{Name: "o", IsDirectory: true}}
|
||||
s3a, owner := newS3a(t, f)
|
||||
entryCreated := false
|
||||
ran := false
|
||||
code := s3a.createAfterAmbiguousRoute(filePath, "b", "o", owner, &filer_pb.Entry{Name: "o"}, uploaded, nil, &entryCreated, chunksDead, func() s3err.ErrorCode { ran = true; return s3err.ErrNone })
|
||||
if code != s3err.ErrServiceUnavailable || ran {
|
||||
t.Fatalf("code=%v ran=%v — re-committed over a directory with possibly-deleted chunks", code, ran)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("match behind an unanswered filer — refuse", func(t *testing.T) {
|
||||
deadAddr := startFakeFiler(t, &fakeLookupFiler{lookupErr: context.DeadlineExceeded})
|
||||
matchAddr := startFakeFiler(t, &fakeLookupFiler{entry: &filer_pb.Entry{Name: "o", Chunks: []*filer_pb.FileChunk{{FileId: "5,01637037d6", Size: 55}}}})
|
||||
s3a := ambiguousRouteServer(t, deadAddr, matchAddr)
|
||||
entryCreated := false
|
||||
ran := false
|
||||
code := s3a.createAfterAmbiguousRoute(filePath, "b", "o", deadAddr, &filer_pb.Entry{Name: "o"}, uploaded, nil, &entryCreated, chunksAlive, func() s3err.ErrorCode { ran = true; return s3err.ErrNone })
|
||||
if code != s3err.ErrServiceUnavailable || entryCreated || ran {
|
||||
t.Fatalf("code=%v entryCreated=%v ran=%v — a match behind an unresolved lookup must not finalize", code, entryCreated, ran)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("entry proven absent, chunks alive — normal create proceeds", func(t *testing.T) {
|
||||
f := &fakeLookupFiler{}
|
||||
s3a, owner := newS3a(t, f)
|
||||
entryCreated := false
|
||||
ran := false
|
||||
code := s3a.createAfterAmbiguousRoute(filePath, "b", "o", owner, &filer_pb.Entry{Name: "o"}, uploaded, nil, &entryCreated, chunksAlive, func() s3err.ErrorCode { ran = true; return s3err.ErrNone })
|
||||
if code != s3err.ErrNone || !ran {
|
||||
t.Fatalf("code=%v ran=%v — proven-absent entry did not reach createUnderLock", code, ran)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("entry absent but chunks unverifiable — refuse", func(t *testing.T) {
|
||||
f := &fakeLookupFiler{}
|
||||
s3a, owner := newS3a(t, f)
|
||||
entryCreated := false
|
||||
ran := false
|
||||
code := s3a.createAfterAmbiguousRoute(filePath, "b", "o", owner, &filer_pb.Entry{Name: "o"}, uploaded, nil, &entryCreated, chunksDead, func() s3err.ErrorCode { ran = true; return s3err.ErrNone })
|
||||
if code != s3err.ErrServiceUnavailable || ran {
|
||||
t.Fatalf("code=%v ran=%v — re-committed an entry over possibly-deleted chunks", code, ran)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("lookup cannot resolve state — refuse", func(t *testing.T) {
|
||||
f := &fakeLookupFiler{lookupErr: context.DeadlineExceeded}
|
||||
s3a, owner := newS3a(t, f)
|
||||
entryCreated := false
|
||||
ran := false
|
||||
code := s3a.createAfterAmbiguousRoute(filePath, "b", "o", owner, &filer_pb.Entry{Name: "o"}, uploaded, nil, &entryCreated, chunksAlive, func() s3err.ErrorCode { ran = true; return s3err.ErrNone })
|
||||
if code != s3err.ErrServiceUnavailable || ran {
|
||||
t.Fatalf("code=%v ran=%v — uncertain state still committed", code, ran)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("finalize fails — recovered entry rolls back", func(t *testing.T) {
|
||||
f := &fakeLookupFiler{entry: &filer_pb.Entry{Name: "o", Chunks: []*filer_pb.FileChunk{{FileId: "5,01637037d6", Size: 55}}}}
|
||||
s3a, owner := newS3a(t, f)
|
||||
entryCreated := false
|
||||
ran := false
|
||||
finalize := &putFinalize{afterCreate: func(*filer_pb.Entry) s3err.ErrorCode { return s3err.ErrInternalError }}
|
||||
code := s3a.createAfterAmbiguousRoute(filePath, "b", "o", owner, &filer_pb.Entry{Name: "o"}, uploaded, finalize, &entryCreated, chunksAlive, func() s3err.ErrorCode { ran = true; return s3err.ErrNone })
|
||||
if code != s3err.ErrInternalError || ran {
|
||||
t.Fatalf("code=%v ran=%v", code, ran)
|
||||
}
|
||||
if len(f.deleted) != 1 || f.deleted[0] != filePath {
|
||||
t.Fatalf("recovered entry was not rolled back: deleted=%v", f.deleted)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
type fakeVersionedDirFiler struct {
|
||||
filer_pb.UnimplementedSeaweedFilerServer
|
||||
entries map[string]*filer_pb.Entry
|
||||
updates map[string]*filer_pb.Entry
|
||||
}
|
||||
|
||||
func (f *fakeVersionedDirFiler) LookupDirectoryEntry(ctx context.Context, req *filer_pb.LookupDirectoryEntryRequest) (*filer_pb.LookupDirectoryEntryResponse, error) {
|
||||
if e, ok := f.entries[path.Join(req.Directory, req.Name)]; ok {
|
||||
return &filer_pb.LookupDirectoryEntryResponse{Entry: e}, nil
|
||||
}
|
||||
return nil, filer_pb.ErrNotFound
|
||||
}
|
||||
|
||||
func (f *fakeVersionedDirFiler) UpdateEntry(ctx context.Context, req *filer_pb.UpdateEntryRequest) (*filer_pb.UpdateEntryResponse, error) {
|
||||
f.updates[path.Join(req.Directory, req.Entry.Name)] = req.Entry
|
||||
return &filer_pb.UpdateEntryResponse{}, nil
|
||||
}
|
||||
|
||||
// A late versioned write keeps the newer latest pointer, but the version it
|
||||
// just stored is born noncurrent — without a NoncurrentSinceNs stamp the
|
||||
// lifecycle engine never ages it out.
|
||||
func TestUpdateLatestVersionInDirectoryBornNoncurrent(t *testing.T) {
|
||||
now := time.Now().UnixNano()
|
||||
newerId := fmt.Sprintf("%016x%s", math.MaxInt64-now, "1111111111111111")
|
||||
olderId := fmt.Sprintf("%016x%s", math.MaxInt64-(now-int64(time.Hour)), "2222222222222222")
|
||||
|
||||
bucketDir := "/buckets/b"
|
||||
versionsPath := path.Join(bucketDir, "o"+s3_constants.VersionsFolder)
|
||||
f := &fakeVersionedDirFiler{
|
||||
entries: map[string]*filer_pb.Entry{
|
||||
versionsPath: {
|
||||
Name: "o" + s3_constants.VersionsFolder,
|
||||
IsDirectory: true,
|
||||
Extended: map[string][]byte{
|
||||
s3_constants.ExtLatestVersionIdKey: []byte(newerId),
|
||||
s3_constants.ExtLatestVersionFileNameKey: []byte("newer.v"),
|
||||
},
|
||||
},
|
||||
path.Join(versionsPath, "late.v"): {Name: "late.v"},
|
||||
},
|
||||
updates: map[string]*filer_pb.Entry{},
|
||||
}
|
||||
s3a := ambiguousRouteServer(t, startFakeFiler(t, f))
|
||||
s3a.option.BucketsPath = "/buckets"
|
||||
|
||||
err := s3a.updateLatestVersionInDirectory("b", "o", olderId, "late.v", &filer_pb.Entry{})
|
||||
if err != nil {
|
||||
t.Fatalf("updateLatestVersionInDirectory: %v", err)
|
||||
}
|
||||
stamped := f.updates[path.Join(versionsPath, "late.v")]
|
||||
if stamped == nil {
|
||||
t.Fatal("late version never got its noncurrent stamp")
|
||||
}
|
||||
if stamped.Extended[s3_constants.ExtNoncurrentSinceNsKey] == nil {
|
||||
t.Fatal("stamp did not set ExtNoncurrentSinceNsKey")
|
||||
}
|
||||
if _, pointerTouched := f.updates[versionsPath]; pointerTouched {
|
||||
t.Fatal("the newer latest pointer must not move")
|
||||
}
|
||||
}
|
||||
@@ -31,6 +31,7 @@ import (
|
||||
weed_server "github.com/seaweedfs/seaweedfs/weed/server"
|
||||
stats_collect "github.com/seaweedfs/seaweedfs/weed/stats"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/constants"
|
||||
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
@@ -1007,7 +1008,9 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader
|
||||
// it); conditional/object-lock/non-reducible cases fall back to the lock.
|
||||
var createCode s3err.ErrorCode
|
||||
routed := false
|
||||
if owner := s3a.routableWriteOwner(bucket, object); owner != "" {
|
||||
routedAmbiguous := false
|
||||
var owner pb.ServerAddress
|
||||
if owner = s3a.routableWriteOwner(bucket, object); owner != "" {
|
||||
if cond, ok := routeWriteCondition(r, uniqueWritePath); ok {
|
||||
// Routed mutations ride in the PUT's transaction (committing atomically),
|
||||
// so lockKey is the object path they carry, not the version file path.
|
||||
@@ -1018,11 +1021,16 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader
|
||||
resp, err := s3a.routedPut(owner, s3a.objectRouteKey(bucket, object), lockKey, filePath, entry, cond, "", finalizeMutations)
|
||||
switch {
|
||||
case err != nil:
|
||||
routedAmbiguous = true
|
||||
glog.Warningf("putToFiler: routed PUT to %s failed for %s, falling back to lock: %v", owner, filePath, err)
|
||||
case resp.ErrorCode == filer_pb.FilerError_PRECONDITION_FAILED:
|
||||
createCode, routed = s3err.ErrPreconditionFailed, true
|
||||
case resp.Error != "":
|
||||
// Non-precondition mutation error: fall back so the lock path maps it.
|
||||
// A mutation can store the entry before a later step of the
|
||||
// same mutation fails, so even an error from the first mutation
|
||||
// does not prove nothing committed — every transaction error is
|
||||
// ambiguous until the recovery lookup decides.
|
||||
routedAmbiguous = true
|
||||
glog.Warningf("putToFiler: routed PUT to %s returned %q for %s, falling back to lock", owner, resp.Error, filePath)
|
||||
default:
|
||||
entryCreated, routed, createCode = true, true, s3err.ErrNone
|
||||
@@ -1035,7 +1043,13 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader
|
||||
}
|
||||
}
|
||||
if !routed {
|
||||
createCode = s3a.withObjectWriteLock(bucket, object, preconditionFn, createUnderLock)
|
||||
createFn := createUnderLock
|
||||
if routedAmbiguous && len(chunkResult.FileChunks) > 0 {
|
||||
createFn = func() s3err.ErrorCode {
|
||||
return s3a.createAfterAmbiguousRoute(filePath, bucket, object, owner, entry, chunkResult.FileChunks, finalize, &entryCreated, s3a.uploadedChunksExist, createUnderLock)
|
||||
}
|
||||
}
|
||||
createCode = s3a.withObjectWriteLock(bucket, object, preconditionFn, createFn)
|
||||
}
|
||||
if createCode != s3err.ErrNone {
|
||||
if createErr != nil {
|
||||
@@ -1097,6 +1111,131 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader
|
||||
// the object write lock, so a hung filer cannot stall the write path.
|
||||
const createLookupTimeout = 10 * time.Second
|
||||
|
||||
// createAfterAmbiguousRoute decides, inside the object write lock, what an
|
||||
// ambiguous routed PUT left behind before re-committing its entry. A stored
|
||||
// entry resolving to the uploaded chunks means the route landed and the create
|
||||
// is already done. A different stored entry means a concurrent write
|
||||
// superseded the commit, and its old-chunk cleanup may have deleted this
|
||||
// entry's chunks — committing would write back a stale entry pointing at dead
|
||||
// needles, so the request fails for the client to retry with a fresh upload.
|
||||
// Unresolvable lookups fail the same way; only a proven-absent entry or a
|
||||
// directory conflict with verifiably-live chunks falls through to the normal
|
||||
// create.
|
||||
func (s3a *S3ApiServer) createAfterAmbiguousRoute(filePath, bucket, object string, owner pb.ServerAddress, entry *filer_pb.Entry, uploaded []*filer_pb.FileChunk, finalize *putFinalize, entryCreated *bool, chunksExist func(ctx context.Context, chunks []*filer_pb.FileChunk) bool, createUnderLock func() s3err.ErrorCode) s3err.ErrorCode {
|
||||
dir, name := path.Dir(filePath), path.Base(filePath)
|
||||
lookupCtx, cancel := context.WithTimeout(context.Background(), createLookupTimeout)
|
||||
defer cancel()
|
||||
existing, uncertain := s3a.lookupCommittedEntry(lookupCtx, s3a.createTargetFilers(owner, bucket, object), dir, name)
|
||||
if existing != nil {
|
||||
if uncertain {
|
||||
// A filer earlier in the list could not answer and may hold a
|
||||
// newer entry; this match is not authoritative.
|
||||
glog.Warningf("putToFiler: ambiguous routed PUT for %s matches an entry but a filer could not be checked; not recovering", filePath)
|
||||
return s3err.ErrServiceUnavailable
|
||||
}
|
||||
if existing.IsDirectory {
|
||||
// The directory can postdate this PUT's commit: the route may
|
||||
// have stored the file entry, a delete removed it, and a nested
|
||||
// write recreated the directory before recovery took the lock.
|
||||
// Re-committing is only safe while the uploaded chunks verifiably
|
||||
// survive — that delete may already have reclaimed the needles,
|
||||
// and a create retrying over the directory would store an entry
|
||||
// pointing at dead chunks.
|
||||
if !chunksExist(lookupCtx, uploaded) {
|
||||
glog.Warningf("putToFiler: ambiguous routed PUT for %s resolved to a directory and its chunks cannot be verified; not re-applying entry", filePath)
|
||||
return s3err.ErrServiceUnavailable
|
||||
}
|
||||
return createUnderLock()
|
||||
}
|
||||
resolved, _, resolveErr := filer.ResolveChunkManifest(lookupCtx, s3a.createLookupFileIdFunction(), existing.GetChunks(), 0, math.MaxInt64, s3a.filerClient)
|
||||
if resolveErr != nil || !sameFileChunks(resolved, uploaded) {
|
||||
glog.Warningf("putToFiler: ambiguous routed PUT for %s was superseded by another write; not committing stale entry", filePath)
|
||||
return s3err.ErrServiceUnavailable
|
||||
}
|
||||
*entryCreated = true
|
||||
if finalize != nil && finalize.afterCreate != nil {
|
||||
if code := finalize.afterCreate(entry); code != s3err.ErrNone {
|
||||
// Same undo the create path applies when post-create
|
||||
// finalization fails.
|
||||
if rbErr := s3a.rmObject(context.Background(), dir, name, true, false); rbErr != nil {
|
||||
glog.Errorf("putToFiler: failed to rollback recovered entry for %s: %v", filePath, rbErr)
|
||||
}
|
||||
return code
|
||||
}
|
||||
}
|
||||
return s3err.ErrNone
|
||||
}
|
||||
if uncertain {
|
||||
glog.Warningf("putToFiler: cannot confirm what ambiguous routed PUT for %s committed; not re-applying entry", filePath)
|
||||
return s3err.ErrServiceUnavailable
|
||||
}
|
||||
if !chunksExist(lookupCtx, uploaded) {
|
||||
glog.Warningf("putToFiler: ambiguous routed PUT for %s left no entry and its chunks cannot be verified; not re-applying entry", filePath)
|
||||
return s3err.ErrServiceUnavailable
|
||||
}
|
||||
return createUnderLock()
|
||||
}
|
||||
|
||||
// lookupCommittedEntry returns the entry any of the candidate filers stores
|
||||
// for dir/name. uncertain reports that a filer earlier in the list could not
|
||||
// answer, so a found entry is not authoritative — the unreachable filer may
|
||||
// hold a newer one.
|
||||
func (s3a *S3ApiServer) lookupCommittedEntry(ctx context.Context, targets []pb.ServerAddress, dir, name string) (existing *filer_pb.Entry, uncertain bool) {
|
||||
for _, target := range targets {
|
||||
e, lookupErr := s3a.lookupEntryOnFiler(ctx, target, dir, name)
|
||||
if e != nil {
|
||||
return e, uncertain
|
||||
}
|
||||
if lookupErr != nil && !errors.Is(lookupErr, filer_pb.ErrNotFound) {
|
||||
uncertain = true
|
||||
}
|
||||
}
|
||||
return nil, uncertain
|
||||
}
|
||||
|
||||
// uploadedChunksExist probes every uploaded chunk's needle on the volumes the
|
||||
// master reports for it. A proven-absent entry is only safe to re-commit when
|
||||
// its chunks survive: if the routed PUT committed and a delete then removed
|
||||
// the entry, its chunk cleanup may have reclaimed the needles too.
|
||||
func (s3a *S3ApiServer) uploadedChunksExist(ctx context.Context, uploaded []*filer_pb.FileChunk) bool {
|
||||
lookupFileId := s3a.createLookupFileIdFunction()
|
||||
for _, chunk := range uploaded {
|
||||
fileId := chunk.GetFileIdString()
|
||||
urls, err := lookupFileId(ctx, fileId)
|
||||
if err != nil || len(urls) == 0 {
|
||||
return false
|
||||
}
|
||||
jwt := filer.ChunkReadJwt(urls, fileId)
|
||||
found := false
|
||||
for _, fileUrl := range urls {
|
||||
if chunkNeedleExists(ctx, fileUrl, jwt) {
|
||||
found = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func chunkNeedleExists(ctx context.Context, fileUrl, jwt string) bool {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodHead, fileUrl, nil)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
if jwt != "" {
|
||||
req.Header.Set("Authorization", security.BearerPrefix+string(jwt))
|
||||
}
|
||||
resp, err := util_http.Do(req)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
util_http.CloseResponse(resp)
|
||||
return resp.StatusCode == http.StatusOK
|
||||
}
|
||||
|
||||
// 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
|
||||
@@ -1111,21 +1250,15 @@ func (s3a *S3ApiServer) confirmCreateLanded(filePath, bucket, object string, ent
|
||||
// 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 {
|
||||
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
|
||||
}
|
||||
}
|
||||
targets := s3a.createTargetFilers(owner, bucket, object)
|
||||
existing, uncertain := s3a.lookupCommittedEntry(lookupCtx, targets, dir, name)
|
||||
if existing == nil {
|
||||
absent = queried && !uncertain
|
||||
absent = len(targets) > 0 && !uncertain
|
||||
return s3err.ErrNone
|
||||
}
|
||||
if uncertain {
|
||||
// The match may sit on a replica behind the filer that could not
|
||||
// answer; it is not authoritative enough to finalize on.
|
||||
return s3err.ErrNone
|
||||
}
|
||||
if len(uploaded) == 0 {
|
||||
@@ -1937,6 +2070,18 @@ func (s3a *S3ApiServer) updateLatestVersionInDirectory(bucket, object, versionId
|
||||
// detected by filename equality and skip the stamp.
|
||||
prevLatestFileName := string(versionsEntry.Extended[s3_constants.ExtLatestVersionFileNameKey])
|
||||
|
||||
// A newer version may have committed while this write's outcome was
|
||||
// uncertain; a pointer already naming one newer than this version must
|
||||
// keep latest where it belongs instead of pointing back.
|
||||
prevLatestVersionId := string(versionsEntry.Extended[s3_constants.ExtLatestVersionIdKey])
|
||||
if prevLatestVersionId != "" && versionId != "null" && compareVersionIds(prevLatestVersionId, versionId) < 0 {
|
||||
// This version is born noncurrent: stamp it so the lifecycle engine
|
||||
// can compute NoncurrentDays, but leave latest where it belongs.
|
||||
s3a.markVersionNoncurrent(bucketDir, versionsObjectPath, versionFileName, time.Now().UnixNano())
|
||||
glog.V(2).Infof("updateLatestVersionInDirectory: %s/%s already points at newer version %s; keeping it", bucket, object, prevLatestVersionId)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Stamp the demoted entry BEFORE updating the .versions/ directory
|
||||
// pointer. The pointer-flip emits a meta-log event that the
|
||||
// lifecycle router consumes; that router then looks up the demoted
|
||||
|
||||
Reference in new issue
Block a user