s3: dial the object lock's primary filer directly

The S3 object write lock builds a fresh short-lived lock per write, each
starting at the seed filer. When the seed isn't the key's hash-ring primary
the filer forwards the request to the primary, and in multi-cluster setups
that forward crosses clusters on every write.

Give the lock client a view of the filer lock ring, fed by the master's
LockRingUpdate broadcasts the gateway already receives, so it dials the
primary directly. The view tracks filer membership by version; a stale view
stays correct because the filer still forwards as a fallback.

Also send the initial ring snapshot to S3 clients, not just filers.
This commit is contained in:
Chris Lu
2026-05-22 21:36:49 -07:00
parent d1665750e1
commit 71c33cc1e2
4 changed files with 138 additions and 5 deletions
+43 -3
View File
@@ -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,
+73
View File
@@ -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)
}
}
+15 -1
View File
@@ -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)
+7 -1
View File
@@ -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
}