diff --git a/weed/cluster/lock_client.go b/weed/cluster/lock_client.go index fa4eea934..44529acb0 100644 --- a/weed/cluster/lock_client.go +++ b/weed/cluster/lock_client.go @@ -28,6 +28,13 @@ type LockClient struct { ringMu sync.RWMutex ring *lock_manager.HashRing ringVersion int64 + + // priorRing is the ring before the most recent change, kept for priorWindow so a + // route-by-key caller can consult a just-moved key's previous owner during a + // rebalance. Mirrors the master's lock_manager.LockRing.PriorOwner cooling-off. + priorRing *lock_manager.HashRing + ringChangedAt time.Time + priorWindow time.Duration } func NewLockClient(grpcDialOption grpc.DialOption, seedFiler pb.ServerAddress) *LockClient { @@ -36,6 +43,7 @@ func NewLockClient(grpcDialOption grpc.DialOption, seedFiler pb.ServerAddress) * maxLockDuration: 5 * time.Second, sleepDuration: 2473 * time.Millisecond, seedFiler: seedFiler, + priorWindow: 5 * time.Second, } } @@ -50,10 +58,15 @@ func (lc *LockClient) SetRing(servers []pb.ServerAddress, version int64) { return } lc.ringVersion = version - if lc.ring == nil { - lc.ring = lock_manager.NewHashRing(lock_manager.DefaultVnodeCount) + // Build a fresh ring (not an in-place mutation) so the outgoing ring survives as + // priorRing with its own servers for the cooling-off window. + newRing := lock_manager.NewHashRing(lock_manager.DefaultVnodeCount) + newRing.SetServers(servers) + if lc.ring != nil { + lc.priorRing = lc.ring + lc.ringChangedAt = time.Now() } - lc.ring.SetServers(servers) + lc.ring = newRing } // hostForKey returns the filer that should own key per the current ring view, @@ -82,6 +95,27 @@ func (lc *LockClient) PrimaryForKey(key string) pb.ServerAddress { return lc.ring.GetPrimary(key) } +// PriorOwnerForKey returns key's owner from the prior ring while a rebalance is +// within the cooling-off window and ownership actually moved, else "". Lets a +// route-by-key reader consult a just-moved key's previous owner before the new +// owner's NotFound is final (the new owner may not have replicated the key yet). +func (lc *LockClient) PriorOwnerForKey(key string) pb.ServerAddress { + lc.ringMu.RLock() + defer lc.ringMu.RUnlock() + if lc.ring == nil || lc.priorRing == nil { + return "" + } + if time.Since(lc.ringChangedAt) > lc.priorWindow { + return "" + } + current := lc.ring.GetPrimary(key) + prior := lc.priorRing.GetPrimary(key) + if prior != "" && prior != current { + return prior + } + return "" +} + type LiveLock struct { key string renewToken string diff --git a/weed/cluster/lock_client_test.go b/weed/cluster/lock_client_test.go index 76ba02157..608ddf6da 100644 --- a/weed/cluster/lock_client_test.go +++ b/weed/cluster/lock_client_test.go @@ -1,7 +1,9 @@ package cluster import ( + "fmt" "testing" + "time" "github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager" "github.com/seaweedfs/seaweedfs/weed/pb" @@ -90,3 +92,61 @@ func TestLockClientPrimaryForKey(t *testing.T) { t.Errorf("PrimaryForKey %q disagrees with hostForKey %q", got, lc.hostForKey("k")) } } + +// A moved key reports its previous owner within the cooling-off window; an unmoved +// key reports none. +func TestLockClientPriorOwnerForKey(t *testing.T) { + lc := NewLockClient(nil, "seed:8888") + + setA := []pb.ServerAddress{"filer-a:8888", "filer-b:8888", "filer-c:8888"} + lc.SetRing(setA, 1) + // One ring: nothing to fall back to. + if got := lc.PriorOwnerForKey("any"); got != "" { + t.Fatalf("single ring should have no prior owner, got %q", got) + } + + priorRing := lock_manager.NewHashRing(lock_manager.DefaultVnodeCount) + priorRing.SetServers(setA) + + // Add a server so some keys' ownership moves. + setB := []pb.ServerAddress{"filer-a:8888", "filer-b:8888", "filer-c:8888", "filer-d:8888"} + lc.SetRing(setB, 2) + + var moved, stable string + for i := 0; i < 2000 && (moved == "" || stable == ""); i++ { + key := fmt.Sprintf("key-%d", i) + if lc.PrimaryForKey(key) != priorRing.GetPrimary(key) { + if moved == "" { + moved = key + } + } else if stable == "" { + stable = key + } + } + if moved == "" || stable == "" { + t.Skip("could not find both a moved and a stable key") + } + + if got, want := lc.PriorOwnerForKey(moved), priorRing.GetPrimary(moved); got != want { + t.Fatalf("PriorOwnerForKey(moved)=%q, want %q", got, want) + } + if got := lc.PriorOwnerForKey(stable); got != "" { + t.Fatalf("unmoved key should have no prior owner, got %q", got) + } +} + +// The prior owner is only offered within the cooling-off window. +func TestLockClientPriorOwnerForKeyExpires(t *testing.T) { + lc := NewLockClient(nil, "seed:8888") + lc.priorWindow = 20 * time.Millisecond + + lc.SetRing([]pb.ServerAddress{"filer-a:8888", "filer-b:8888", "filer-c:8888"}, 1) + lc.SetRing([]pb.ServerAddress{"filer-a:8888", "filer-b:8888", "filer-c:8888", "filer-d:8888"}, 2) + + time.Sleep(40 * time.Millisecond) + for i := 0; i < 2000; i++ { + if got := lc.PriorOwnerForKey(fmt.Sprintf("key-%d", i)); got != "" { + t.Fatalf("prior owner should expire after the cooling window, got %q", got) + } + } +} diff --git a/weed/s3api/s3api_handlers.go b/weed/s3api/s3api_handlers.go index 9b57d5807..09fbbe81a 100644 --- a/weed/s3api/s3api_handlers.go +++ b/weed/s3api/s3api_handlers.go @@ -20,7 +20,7 @@ var _ = filer_pb.FilerClient(&S3ApiServer{}) func (s3a *S3ApiServer) WithFilerClient(streamingMode bool, fn func(filer_pb.SeaweedFilerClient) error) error { // Use filerClient for proper connection management and failover if s3a.filerClient != nil { - return s3a.withFilerClientFailover(streamingMode, fn) + return s3a.withFilerClientFailover("", streamingMode, fn) } // Fallback to direct connection if filerClient not initialized @@ -32,69 +32,85 @@ func (s3a *S3ApiServer) WithFilerClient(streamingMode bool, fn func(filer_pb.Sea } -// withFilerClientFailover attempts to execute fn with automatic failover to other filers -func (s3a *S3ApiServer) withFilerClientFailover(streamingMode bool, fn func(filer_pb.SeaweedFilerClient) error) error { - // Get current filer as starting point +// withFilerClientFailover runs fn against preferred (if set) first, then the +// remaining filers (current, then the rest) with healthy ones before unhealthy. +// Failover is for transport errors only: a reached filer's ErrNotFound is +// authoritative (no +// fan-out, no resurrecting a peer's not-yet-replicated tombstone). preferred lets a +// caller route to a key's ring owner for read-after-write; it may be a filer outside +// the static list (the bookkeeping no-ops for untracked addresses). A failover +// updates the current filer; a preferred read does not, as its owner is per-key. +func (s3a *S3ApiServer) withFilerClientFailover(preferred pb.ServerAddress, streamingMode bool, fn func(filer_pb.SeaweedFilerClient) error) error { currentFiler := s3a.filerClient.GetCurrentFiler() - // Try current filer first (fast path) - err := pb.WithGrpcClient(context.Background(), streamingMode, s3a.randomClientId, func(grpcConnection *grpc.ClientConn) error { - client := filer_pb.NewSeaweedFilerClient(grpcConnection) - return fn(client) - }, currentFiler.ToGrpcAddress(), false, s3a.option.GrpcDialOption) - - if err == nil { - s3a.filerClient.RecordFilerSuccess(currentFiler) - return nil - } - - // A reachable filer answering ErrNotFound is authoritative; failover is for - // unreachable/unhealthy filers, not for re-asking about an absent entry. - if errors.Is(err, filer_pb.ErrNotFound) { - return err - } - - s3a.filerClient.RecordFilerFailure(currentFiler) - - // Current filer failed - try all other filers with health-aware selection - filers := s3a.filerClient.GetAllFilers() - var lastErr error = err - - for _, filer := range filers { - if filer == currentFiler { - continue // Already tried this one + candidates := make([]pb.ServerAddress, 0, 2+len(s3a.option.Filers)) + seen := make(map[pb.ServerAddress]bool) + addCandidate := func(filer pb.ServerAddress) { + if filer == "" || seen[filer] { + return } + seen[filer] = true + candidates = append(candidates, filer) + } + addCandidate(preferred) + addCandidate(currentFiler) + for _, filer := range s3a.filerClient.GetAllFilers() { + addCandidate(filer) + } - // Skip filers known to be unhealthy (circuit breaker pattern) - if s3a.filerClient.ShouldSkipUnhealthyFiler(filer) { - glog.V(2).Infof("WithFilerClient: skipping unhealthy filer %s", filer) + // An explicit preferred owner is tried first even if health-flagged: demoting it + // behind a healthy replica would let that replica's authoritative NotFound mask a + // write that has only reached the owner. Health-order the rest, unhealthy ones + // last so a request still progresses when all are flagged. + var healthy, unhealthy []pb.ServerAddress + for _, filer := range candidates { + if filer == preferred { continue } + if s3a.filerClient.ShouldSkipUnhealthyFiler(filer) { + unhealthy = append(unhealthy, filer) + } else { + healthy = append(healthy, filer) + } + } + ordered := make([]pb.ServerAddress, 0, len(candidates)) + if preferred != "" { + ordered = append(ordered, preferred) + } + ordered = append(ordered, healthy...) + ordered = append(ordered, unhealthy...) - err = pb.WithGrpcClient(context.Background(), streamingMode, s3a.randomClientId, func(grpcConnection *grpc.ClientConn) error { - client := filer_pb.NewSeaweedFilerClient(grpcConnection) - return fn(client) + var lastErr error + for _, filer := range ordered { + err := pb.WithGrpcClient(context.Background(), streamingMode, s3a.randomClientId, func(grpcConnection *grpc.ClientConn) error { + return fn(filer_pb.NewSeaweedFilerClient(grpcConnection)) }, filer.ToGrpcAddress(), false, s3a.option.GrpcDialOption) if err == nil { - // Success! Record success and update current filer for future requests s3a.filerClient.RecordFilerSuccess(filer) - s3a.filerClient.SetCurrentFiler(filer) - glog.V(1).Infof("WithFilerClient: failover from %s to %s succeeded", currentFiler, filer) + if filer != currentFiler && filer != preferred { + s3a.filerClient.SetCurrentFiler(filer) + glog.V(1).Infof("WithFilerClient: failover from %s to %s succeeded", currentFiler, filer) + } return nil } - - // Authoritative not-found - stop failing over. if errors.Is(err, filer_pb.ErrNotFound) { return err } s3a.filerClient.RecordFilerFailure(filer) - glog.V(2).Infof("WithFilerClient: failover to %s failed: %v", filer, err) + // A preferred owner is often outside the static filer list, where the health + // tracking above no-ops; flag it so route-by-key reads skip it briefly. + if filer == preferred { + s3a.markOwnerUnreachable(filer) + } + glog.V(2).Infof("WithFilerClient: filer %s failed: %v", filer, err) lastErr = err } - // All filers failed + if lastErr == nil { + lastErr = fmt.Errorf("no filer available") + } return fmt.Errorf("all filers failed, last error: %w", lastErr) } diff --git a/weed/s3api/s3api_object_handlers.go b/weed/s3api/s3api_object_handlers.go index 04d510061..6ce81b249 100644 --- a/weed/s3api/s3api_object_handlers.go +++ b/weed/s3api/s3api_object_handlers.go @@ -2408,8 +2408,7 @@ func (s3a *S3ApiServer) HeadObjectHandler(w http.ResponseWriter, r *http.Request // fetchObjectEntry fetches the filer entry for an object // Returns nil if not found (not an error), or propagates other errors func (s3a *S3ApiServer) fetchObjectEntry(bucket, object string) (*filer_pb.Entry, error) { - objectPath := fmt.Sprintf("%s/%s", s3a.bucketDir(bucket), object) - fetchedEntry, fetchErr := s3a.getEntry("", objectPath) + fetchedEntry, fetchErr := s3a.getObjectEntryRoutedByKey(bucket, object) if fetchErr != nil { if errors.Is(fetchErr, filer_pb.ErrNotFound) { return nil, nil // Not found is not an error for SSE check @@ -2422,8 +2421,7 @@ func (s3a *S3ApiServer) fetchObjectEntry(bucket, object string) (*filer_pb.Entry // fetchObjectEntryRequired fetches the filer entry for an object // Returns an error if the object is not found or any other error occurs func (s3a *S3ApiServer) fetchObjectEntryRequired(bucket, object string) (*filer_pb.Entry, error) { - objectPath := fmt.Sprintf("%s/%s", s3a.bucketDir(bucket), object) - fetchedEntry, fetchErr := s3a.getEntry("", objectPath) + fetchedEntry, fetchErr := s3a.getObjectEntryRoutedByKey(bucket, object) if fetchErr != nil { return nil, fetchErr // Return error for both not-found and other errors } diff --git a/weed/s3api/s3api_object_routed_read.go b/weed/s3api/s3api_object_routed_read.go new file mode 100644 index 000000000..ec5661064 --- /dev/null +++ b/weed/s3api/s3api_object_routed_read.go @@ -0,0 +1,97 @@ +package s3api + +import ( + "context" + "errors" + "time" + + "github.com/seaweedfs/seaweedfs/weed/pb" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/util" +) + +// unreachableOwnerTTL is how long a route-by-key owner that just failed a read is +// skipped before being retried — long enough to spare a dead owner per-request +// dials, short enough to resume owner-first reads soon after it (or the ring) recovers. +const unreachableOwnerTTL = 2 * time.Second + +// getObjectEntryRoutedByKey resolves an object's entry preferring the key's write +// owner (the same route key the write path hashes), so a read sees a just-written +// object without waiting for cross-filer replication. On the owner's ErrNotFound it +// probes the key's prior owner once during a rebalance window; falls back to +// getEntry when no owner is resolvable. +func (s3a *S3ApiServer) getObjectEntryRoutedByKey(bucket, object string) (*filer_pb.Entry, error) { + fullPath := util.NewFullPath(s3a.bucketDir(bucket), object) + + owner := s3a.routableWriteOwner(bucket, object) + if owner == "" || s3a.filerClient == nil { + return filer_pb.GetEntry(context.Background(), s3a, fullPath) + } + + // Skip an owner whose recent read hit a transport error; read local-first until + // it (or the ring) recovers, rather than re-dialing a dead owner every request. + preferred := owner + if s3a.ownerRecentlyUnreachable(owner) { + preferred = "" + } + + dir, name := fullPath.DirAndName() + var entry *filer_pb.Entry + err := s3a.withFilerClientFailover(preferred, false, func(client filer_pb.SeaweedFilerClient) error { + resp, lookupErr := filer_pb.LookupEntry(context.Background(), client, &filer_pb.LookupDirectoryEntryRequest{ + Directory: dir, + Name: name, + }) + if lookupErr != nil { + return lookupErr + } + entry = resp.Entry + return nil + }) + + // A just-moved key may not have replicated to the new owner yet; consult its + // 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 { + return priorEntry, nil + } + } + } + return entry, err +} + +func (s3a *S3ApiServer) priorWriteOwner(bucket, object string) pb.ServerAddress { + if object == "" || s3a.objectWriteLockClient == nil { + return "" + } + return s3a.objectWriteLockClient.PriorOwnerForKey(s3a.objectRouteKey(bucket, object)) +} + +func (s3a *S3ApiServer) markOwnerUnreachable(owner pb.ServerAddress) { + s3a.unreachableOwners.Store(owner, time.Now().Add(unreachableOwnerTTL)) +} + +func (s3a *S3ApiServer) ownerRecentlyUnreachable(owner pb.ServerAddress) bool { + if v, ok := s3a.unreachableOwners.Load(owner); ok { + return time.Now().Before(v.(time.Time)) + } + return false +} + +// lookupEntryOnFiler resolves dir/name against a single filer, without failover. +func (s3a *S3ApiServer) lookupEntryOnFiler(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{ + Directory: dir, + Name: name, + }) + if lookupErr != nil { + return lookupErr + } + entry = resp.Entry + return nil + }) + return entry, err +} diff --git a/weed/s3api/s3api_object_routed_read_test.go b/weed/s3api/s3api_object_routed_read_test.go new file mode 100644 index 000000000..1987eed7c --- /dev/null +++ b/weed/s3api/s3api_object_routed_read_test.go @@ -0,0 +1,24 @@ +package s3api + +import ( + "testing" + + "github.com/seaweedfs/seaweedfs/weed/pb" +) + +// A flagged owner reads back as recently unreachable; an unflagged one does not. +func TestOwnerRecentlyUnreachable(t *testing.T) { + s3a := &S3ApiServer{} + owner := pb.ServerAddress("127.0.0.1:8888.18888") + + if s3a.ownerRecentlyUnreachable(owner) { + t.Fatal("unmarked owner should not be flagged") + } + s3a.markOwnerUnreachable(owner) + if !s3a.ownerRecentlyUnreachable(owner) { + t.Fatal("marked owner should be flagged within the TTL") + } + if s3a.ownerRecentlyUnreachable(pb.ServerAddress("127.0.0.1:9999.19999")) { + t.Fatal("a different owner should not be flagged") + } +} diff --git a/weed/s3api/s3api_server.go b/weed/s3api/s3api_server.go index 2477557bf..88029d70e 100644 --- a/weed/s3api/s3api_server.go +++ b/weed/s3api/s3api_server.go @@ -98,6 +98,12 @@ type S3ApiServer struct { newObjectWriteLock func(bucket, object string) objectWriteLock // objectWriteLockClient resolves a key's owner filer for route-by-key. objectWriteLockClient *cluster.LockClient + // unreachableOwners holds owners (pb.ServerAddress -> expiry time.Time) whose + // last owner-first read hit a transport error, so route-by-key reads briefly + // skip them instead of re-dialing a dead owner every request until the ring + // drops it. Bypasses the gateway's filer health tracking, which no-ops for an + // owner outside the static -filer list. + unreachableOwners sync.Map // Shared ReaderCache used by the S3 GET streaming path. It lives for the // lifetime of the server so that concurrent and repeat reads share a // single in-flight download per chunk, and so that no per-request