mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-21 22:56:55 +00:00
Merge branch 's3-lock-ring-view' into s3-route-by-key
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user