mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-22 07:06:51 +00:00
s3: subscribe to lock-ring updates before starting the master loop
The master delivers the initial LockRingUpdate once, on connect. Registering the callback after KeepConnectedToMaster started left a window where that first update could arrive before the handler was set and be dropped, delaying the ring view until the next membership change. Build the lock client and register the callback in the masters block before launching the loop; the filers block reuses that client (or creates a plain one when no masters are configured).
This commit is contained in:
+20
-12
@@ -154,6 +154,7 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl
|
||||
// Supports multiple filer addresses with automatic failover for high availability
|
||||
var filerClient *wdclient.FilerClient
|
||||
var masterClient *wdclient.MasterClient
|
||||
var objectWriteLockClient *cluster.LockClient
|
||||
if len(option.Masters) > 0 {
|
||||
// Enable filer discovery via master
|
||||
masterMap := make(map[string]pb.ServerAddress)
|
||||
@@ -165,6 +166,21 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl
|
||||
clientHost = util.DetectedHostAddress()
|
||||
}
|
||||
masterClient = wdclient.NewMasterClient(option.GrpcDialOption, option.FilerGroup, cluster.S3Type, pb.ServerAddress(util.JoinHostPort(clientHost, option.GrpcPort)), "", "", *pb.NewServiceDiscoveryFromMap(masterMap))
|
||||
// Build the object-write lock client and subscribe to the master's
|
||||
// lock-ring updates BEFORE starting the master loop, so the initial
|
||||
// LockRingUpdate sent on connect isn't dropped (the master only delivers
|
||||
// it once per connect). The masterClient already filters updates to this
|
||||
// server's filer group.
|
||||
if len(option.Filers) > 0 {
|
||||
objectWriteLockClient = cluster.NewLockClient(option.GrpcDialOption, option.Filers[0])
|
||||
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)
|
||||
})
|
||||
}
|
||||
// Start the master client connection loop - required for GetMaster() to work
|
||||
go masterClient.KeepConnectedToMaster(context.Background())
|
||||
|
||||
@@ -265,18 +281,10 @@ 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)
|
||||
})
|
||||
// Reuse the lock client built in the masters block (already subscribed to
|
||||
// ring updates); create a plain one when no masters are configured.
|
||||
if objectWriteLockClient == nil {
|
||||
objectWriteLockClient = cluster.NewLockClient(option.GrpcDialOption, option.Filers[0])
|
||||
}
|
||||
s3ApiServer.newObjectWriteLock = func(bucket, object string) objectWriteLock {
|
||||
lockKey := fmt.Sprintf("s3.object.write:%s", s3ApiServer.toFilerPath(bucket, object))
|
||||
|
||||
Reference in New Issue
Block a user