diff --git a/weed/cluster/lock_client.go b/weed/cluster/lock_client.go index d84f51bef..18ffc61bf 100644 --- a/weed/cluster/lock_client.go +++ b/weed/cluster/lock_client.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "strings" + "sync" "sync/atomic" "time" @@ -19,6 +20,14 @@ type LockClient struct { maxLockDuration time.Duration sleepDuration time.Duration seedFiler pb.ServerAddress + + // ring is an optional client-side view of the filer lock hash ring. When + // populated, a new lock starts at the key's primary filer instead of the + // seed filer, avoiding the seed->primary forward hop. A stale view stays + // correct: the filer forwards to the real primary as a fallback. + ringMu sync.RWMutex + ring *lock_manager.HashRing + ringVersion int64 } func NewLockClient(grpcDialOption grpc.DialOption, seedFiler pb.ServerAddress) *LockClient { @@ -30,6 +39,37 @@ func NewLockClient(grpcDialOption grpc.DialOption, seedFiler pb.ServerAddress) * } } +// SetRing updates the client-side view of the filer lock ring. It mirrors the +// master's LockRingUpdate broadcasts so the client computes the same primary +// the filers do. Updates older than the current version are ignored to tolerate +// reordered broadcasts; version 0 is always accepted (bootstrap). +func (lc *LockClient) SetRing(servers []pb.ServerAddress, version int64) { + lc.ringMu.Lock() + defer lc.ringMu.Unlock() + if version != 0 && version < lc.ringVersion { + return + } + lc.ringVersion = version + if lc.ring == nil { + lc.ring = lock_manager.NewHashRing(lock_manager.DefaultVnodeCount) + } + lc.ring.SetServers(servers) +} + +// hostForKey returns the filer that should own key per the current ring view, +// falling back to the seed filer when no view has been received yet. +func (lc *LockClient) hostForKey(key string) pb.ServerAddress { + lc.ringMu.RLock() + defer lc.ringMu.RUnlock() + if lc.ring == nil { + return lc.seedFiler + } + if primary := lc.ring.GetPrimary(key); primary != "" { + return primary + } + return lc.seedFiler +} + type LiveLock struct { key string renewToken string @@ -51,7 +91,7 @@ type LiveLock struct { func (lc *LockClient) NewShortLivedLock(key string, owner string) (lock *LiveLock) { lock = &LiveLock{ key: key, - hostFiler: lc.seedFiler, + hostFiler: lc.hostForKey(key), cancelCh: make(chan struct{}), expireAtNs: time.Now().Add(5 * time.Second).UnixNano(), grpcDialOption: lc.grpcDialOption, @@ -72,7 +112,7 @@ func (lc *LockClient) NewBlockingLongLivedLock(key, owner string, lockTTL time.D } lock := &LiveLock{ key: key, - hostFiler: lc.seedFiler, + hostFiler: lc.hostForKey(key), cancelCh: make(chan struct{}), expireAtNs: time.Now().Add(lockTTL).UnixNano(), grpcDialOption: lc.grpcDialOption, @@ -110,7 +150,7 @@ func (lc *LockClient) NewBlockingLongLivedLock(key, owner string, lockTTL time.D func (lc *LockClient) StartLongLivedLock(key string, owner string, onLockOwnerChange func(newLockOwner string), lockTTL time.Duration) (lock *LiveLock) { lock = &LiveLock{ key: key, - hostFiler: lc.seedFiler, + hostFiler: lc.hostForKey(key), cancelCh: make(chan struct{}), expireAtNs: time.Now().Add(lockTTL).UnixNano(), grpcDialOption: lc.grpcDialOption, diff --git a/weed/cluster/lock_client_test.go b/weed/cluster/lock_client_test.go new file mode 100644 index 000000000..4b2d6904c --- /dev/null +++ b/weed/cluster/lock_client_test.go @@ -0,0 +1,73 @@ +package cluster + +import ( + "testing" + + "github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager" + "github.com/seaweedfs/seaweedfs/weed/pb" +) + +// The gateway must resolve a lock key to the same primary the filers do, +// otherwise it dials the wrong filer and the lock still gets forwarded. Both +// sides use the same HashRing over the same server set, so for every key the +// client's hostForKey must equal the filer ring's GetPrimary. +func TestLockClientHostMatchesFilerRing(t *testing.T) { + servers := []pb.ServerAddress{ + "filer-a:8888", "filer-b:8888", "filer-c:8888", "filer-d:8888", + } + + filerRing := lock_manager.NewHashRing(lock_manager.DefaultVnodeCount) + filerRing.SetServers(servers) + + lc := NewLockClient(nil, "seed:8888") + lc.SetRing(servers, 1) + + for _, key := range []string{ + "s3.object.write:/buckets/b/obj-0", + "s3.object.write:/buckets/b/obj-1", + "s3.object.write:/buckets/b/obj-2", + "s3.object.write:/buckets/gosbench-0/w0obj-kilo-0877", + "some/other/key", + } { + if got, want := lc.hostForKey(key), filerRing.GetPrimary(key); got != want { + t.Errorf("key %q: client host %q != filer primary %q", key, got, want) + } + } +} + +// Without a ring view, the client falls back to the seed filer (which the filer +// forwards from), preserving the pre-optimization behavior. +func TestLockClientHostFallsBackToSeed(t *testing.T) { + lc := NewLockClient(nil, "seed:8888") + if got := lc.hostForKey("any-key"); got != "seed:8888" { + t.Errorf("expected seed fallback, got %q", got) + } + + // An empty ring (no members yet) also falls back to the seed. + lc.SetRing(nil, 1) + if got := lc.hostForKey("any-key"); got != "seed:8888" { + t.Errorf("expected seed fallback on empty ring, got %q", got) + } +} + +// A stale (older-version) update must not regress a newer ring view, while +// version 0 always applies as a bootstrap. +func TestLockClientSetRingVersionGuard(t *testing.T) { + lc := NewLockClient(nil, "seed:8888") + + newer := []pb.ServerAddress{"filer-a:8888", "filer-b:8888"} + lc.SetRing(newer, 10) + primaryAt10 := lc.hostForKey("k") + + // Older version is ignored. + lc.SetRing([]pb.ServerAddress{"filer-z:8888"}, 5) + if got := lc.hostForKey("k"); got != primaryAt10 { + t.Errorf("stale update applied: host changed to %q", got) + } + + // version 0 is always accepted. + lc.SetRing([]pb.ServerAddress{"filer-z:8888"}, 0) + if got := lc.hostForKey("k"); got != "filer-z:8888" { + t.Errorf("bootstrap update not applied, got %q", got) + } +} diff --git a/weed/s3api/s3api_server.go b/weed/s3api/s3api_server.go index b3e6d9d25..607cb4dc6 100644 --- a/weed/s3api/s3api_server.go +++ b/weed/s3api/s3api_server.go @@ -25,6 +25,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/iam/policy" "github.com/seaweedfs/seaweedfs/weed/iam/sts" "github.com/seaweedfs/seaweedfs/weed/pb" + "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" "github.com/seaweedfs/seaweedfs/weed/pb/s3_lifecycle_pb" "github.com/seaweedfs/seaweedfs/weed/pb/s3_pb" "github.com/seaweedfs/seaweedfs/weed/s3api/policy_engine" @@ -152,6 +153,7 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl // Uses the battle-tested vidMap with filer-based lookups // Supports multiple filer addresses with automatic failover for high availability var filerClient *wdclient.FilerClient + var masterClient *wdclient.MasterClient if len(option.Masters) > 0 { // Enable filer discovery via master masterMap := make(map[string]pb.ServerAddress) @@ -162,7 +164,7 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl if clientHost == "0.0.0.0" || clientHost == "" { clientHost = util.DetectedHostAddress() } - masterClient := wdclient.NewMasterClient(option.GrpcDialOption, option.FilerGroup, cluster.S3Type, pb.ServerAddress(util.JoinHostPort(clientHost, option.GrpcPort)), "", "", *pb.NewServiceDiscoveryFromMap(masterMap)) + masterClient = wdclient.NewMasterClient(option.GrpcDialOption, option.FilerGroup, cluster.S3Type, pb.ServerAddress(util.JoinHostPort(clientHost, option.GrpcPort)), "", "", *pb.NewServiceDiscoveryFromMap(masterMap)) // Start the master client connection loop - required for GetMaster() to work go masterClient.KeepConnectedToMaster(context.Background()) @@ -264,6 +266,18 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl if len(option.Filers) > 0 { objectWriteLockClient := cluster.NewLockClient(option.GrpcDialOption, option.Filers[0]) + // Mirror the master's lock-ring view so each object lock dials the key's + // primary filer directly instead of forwarding through the seed filer. + // The masterClient already filters updates to this server's filer group. + if masterClient != nil { + masterClient.SetOnLockRingUpdateFn(func(update *master_pb.LockRingUpdate) { + servers := make([]pb.ServerAddress, 0, len(update.Servers)) + for _, s := range update.Servers { + servers = append(servers, pb.ServerAddress(s)) + } + objectWriteLockClient.SetRing(servers, update.Version) + }) + } s3ApiServer.newObjectWriteLock = func(bucket, object string) objectWriteLock { lockKey := fmt.Sprintf("s3.object.write:%s", s3ApiServer.toFilerPath(bucket, object)) owner := fmt.Sprintf("s3api-%d", s3ApiServer.randomClientId) diff --git a/weed/server/master_grpc_server.go b/weed/server/master_grpc_server.go index fd7709d67..aef686959 100644 --- a/weed/server/master_grpc_server.go +++ b/weed/server/master_grpc_server.go @@ -393,7 +393,13 @@ func (ms *MasterServer) KeepConnected(stream master_pb.Seaweed_KeepConnectedServ } func (ms *MasterServer) initialLockRingUpdate(clientType string, filerGroup string) *master_pb.KeepConnectedResponse { - if clientType != cluster.FilerType || ms.LockRingManager == nil { + if ms.LockRingManager == nil { + return nil + } + // Filers are ring members; S3 gateways are lock clients that need the same + // view to dial a key's primary directly. Both get the initial snapshot; + // later membership changes already broadcast to every connected client. + if clientType != cluster.FilerType && clientType != cluster.S3Type { return nil }