diff --git a/weed/cluster/lock_client.go b/weed/cluster/lock_client.go index bc5a8ac01..9d14f383f 100644 --- a/weed/cluster/lock_client.go +++ b/weed/cluster/lock_client.go @@ -3,6 +3,7 @@ package cluster import ( "context" "fmt" + "slices" "strings" "sync" "sync/atomic" @@ -27,6 +28,7 @@ type LockClient struct { // correct: the filer forwards to the real primary as a fallback. ringMu sync.RWMutex ring *lock_manager.HashRing + ringServers []pb.ServerAddress ringVersion int64 // priorRing is the ring before the most recent change, kept for priorWindow so a @@ -35,15 +37,24 @@ type LockClient struct { priorRing *lock_manager.HashRing ringChangedAt time.Time priorWindow time.Duration + + // noLockServerRetryPeriod bounds retries when every filer reports "no + // lock server found": an empty lock ring is a systemic fault that waiting + // on a lock holder cannot resolve, unlike ordinary contention. The bound + // is long enough to ride out a master leader change (ring reset + fresh + // snapshot) but short enough that a write fails instead of hanging + // forever. + noLockServerRetryPeriod time.Duration } func NewLockClient(grpcDialOption grpc.DialOption, seedFiler pb.ServerAddress) *LockClient { return &LockClient{ - grpcDialOption: grpcDialOption, - maxLockDuration: 5 * time.Second, - sleepDuration: 2473 * time.Millisecond, - seedFiler: seedFiler, - priorWindow: 5 * time.Second, + grpcDialOption: grpcDialOption, + maxLockDuration: 5 * time.Second, + sleepDuration: 2473 * time.Millisecond, + seedFiler: seedFiler, + priorWindow: 5 * time.Second, + noLockServerRetryPeriod: 15 * time.Second, } } @@ -58,10 +69,16 @@ func (lc *LockClient) SetRing(servers []pb.ServerAddress, version int64) { return } lc.ringVersion = version + sorted := slices.Clone(servers) + slices.Sort(sorted) + if slices.Equal(sorted, lc.ringServers) { + return + } + lc.ringServers = sorted // 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) + newRing.SetServers(sorted) if lc.ring != nil { lc.priorRing = lc.ring lc.ringChangedAt = time.Now() @@ -69,6 +86,18 @@ func (lc *LockClient) SetRing(servers []pb.ServerAddress, version int64) { lc.ring = newRing } +// ResetRing clears only the version gate so the first update from a +// different master applies unconditionally: ring versions are per-master +// monotonic and a high version accepted from a former leader must not +// reject the new leader's snapshot. The last ring keeps routing during the +// gap rather than falling back to the seed filer, and the arriving ring +// becomes the prior ring for the cooling-off window. +func (lc *LockClient) ResetRing() { + lc.ringMu.Lock() + defer lc.ringMu.Unlock() + lc.ringVersion = 0 +} + // 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 { @@ -133,7 +162,9 @@ type LiveLock struct { consecutiveFailures int // Track connection failures to trigger fallback } -// NewShortLivedLock creates a lock with a 5-second duration +// NewShortLivedLock creates a lock with a 5-second duration. +// It returns nil when the lock cannot be acquired because no lock server +// exists; ordinary contention is still waited out. func (lc *LockClient) NewShortLivedLock(key string, owner string) (lock *LiveLock) { lock = &LiveLock{ key: key, @@ -144,7 +175,10 @@ func (lc *LockClient) NewShortLivedLock(key string, owner string) (lock *LiveLoc self: owner, lc: lc, } - lock.retryUntilLocked(5 * time.Second) + if err := lock.retryUntilLocked(5 * time.Second); err != nil { + glog.Warningf("create lock %s: %v", key, err) + return nil + } return } @@ -167,7 +201,10 @@ func (lc *LockClient) NewBlockingLongLivedLock(key, owner string, lockTTL time.D lockTTL: lockTTL, } // Block until acquired - lock.retryUntilLocked(lockTTL) + if err := lock.retryUntilLocked(lockTTL); err != nil { + glog.Warningf("create lock %s: %v", key, err) + return nil + } // Start renewal goroutine using a ticker for interruptible sleep lock.renewalDone = make(chan struct{}) go func() { @@ -266,12 +303,31 @@ func (lc *LockClient) StartLongLivedLock(key string, owner string, onLockOwnerCh // several seconds): when a holder on another mount releases the lock, the // waiter must pick it up promptly, otherwise cross-mount write handoff stalls // long enough to time out clients. -func (lock *LiveLock) retryUntilLocked(lockDuration time.Duration) { +func (lock *LiveLock) retryUntilLocked(lockDuration time.Duration) error { + var unavailableSince time.Time for lock.renewToken == "" { - if err := lock.AttemptToLock(lockDuration); err != nil { - glog.V(1).Infof("create lock %s: %v", lock.key, err) + err := lock.AttemptToLock(lockDuration) + if err == nil { + unavailableSince = time.Time{} + continue + } + glog.V(1).Infof("create lock %s: %v", lock.key, err) + if strings.Contains(err.Error(), "lock already owned") { + // Ordinary contention: a reachable server holds the lock, so + // waiting is the point and has no bound. + unavailableSince = time.Time{} + continue + } + // Anything else — "no lock server found", a dead ring member refusing + // connections — is a systemic fault waiting cannot fix; give up once + // it persists past the retry period. + if unavailableSince.IsZero() { + unavailableSince = time.Now() + } else if time.Since(unavailableSince) > lock.lc.noLockServerRetryPeriod { + return err } } + return nil } func (lock *LiveLock) AttemptToLock(lockDuration time.Duration) error { diff --git a/weed/cluster/lock_client_test.go b/weed/cluster/lock_client_test.go index c53e8c6cf..7a6a89950 100644 --- a/weed/cluster/lock_client_test.go +++ b/weed/cluster/lock_client_test.go @@ -1,13 +1,18 @@ package cluster import ( + "context" "fmt" + "net" "sync/atomic" "testing" "time" "github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager" "github.com/seaweedfs/seaweedfs/weed/pb" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" ) // The gateway must resolve a lock key to the same primary the filers do, @@ -152,6 +157,135 @@ func TestLockClientPriorOwnerForKeyExpires(t *testing.T) { } } +// A master change clears the version gate so the new leader's (lower-versioned) +// snapshot applies — versions are only comparable within one master's stream. +func TestLockClientResetRing(t *testing.T) { + lc := NewLockClient(nil, "seed:8888") + + lc.SetRing([]pb.ServerAddress{"filer-a:8888", "filer-b:8888"}, 100) + lc.ResetRing() + + // The last ring keeps routing during the gap; only version acceptance + // is reset so the new leader's (lower-versioned) snapshot applies. + if got := lc.hostForKey("k"); got == "seed:8888" { + t.Fatal("expected the previous ring to keep routing after reset") + } + lc.SetRing([]pb.ServerAddress{"filer-z:8888"}, 50) + if got := lc.hostForKey("k"); got != "filer-z:8888" { + t.Fatalf("lower version from new master not applied, got %q", got) + } +} + +type noLockServerFiler struct { + filer_pb.UnimplementedSeaweedFilerServer +} + +func (s *noLockServerFiler) DistributedLock(ctx context.Context, req *filer_pb.LockRequest) (*filer_pb.LockResponse, error) { + return &filer_pb.LockResponse{Error: lock_manager.NoLockServerError.Error()}, nil +} + +// When every filer reports an empty lock ring, lock acquisition must fail +// after a bounded period instead of hanging the write forever. +func TestNewShortLivedLockFailsFastOnNoLockServer(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + grpcServer := grpc.NewServer() + filer_pb.RegisterSeaweedFilerServer(grpcServer, &noLockServerFiler{}) + go grpcServer.Serve(listener) + defer grpcServer.Stop() + + dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) + host, port, err := net.SplitHostPort(listener.Addr().String()) + if err != nil { + t.Fatalf("split host port: %v", err) + } + // "host:httpPort.grpcPort" dials the fake filer's port directly. + lc := NewLockClient(dialOption, pb.ServerAddress(fmt.Sprintf("%s:0.%s", host, port))) + lc.noLockServerRetryPeriod = 200 * time.Millisecond + + start := time.Now() + lock := lc.NewShortLivedLock("test-key", "test-owner") + elapsed := time.Since(start) + + if lock != nil { + t.Fatal("expected nil lock when no lock server exists") + } + if elapsed > 10*time.Second { + t.Fatalf("lock acquisition took %v, expected fail-fast", elapsed) + } +} + +type contendedLockFiler struct { + filer_pb.UnimplementedSeaweedFilerServer +} + +func (s *contendedLockFiler) DistributedLock(ctx context.Context, req *filer_pb.LockRequest) (*filer_pb.LockResponse, error) { + return &filer_pb.LockResponse{Error: "lock already owned by someone"}, nil +} + +// A lock held by another owner is ordinary contention: acquisition waits it +// out rather than failing on the unavailability bound. +func TestNewShortLivedLockWaitsOutContention(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + grpcServer := grpc.NewServer() + filer_pb.RegisterSeaweedFilerServer(grpcServer, &contendedLockFiler{}) + go grpcServer.Serve(listener) + defer grpcServer.Stop() + + dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) + host, port, err := net.SplitHostPort(listener.Addr().String()) + if err != nil { + t.Fatalf("split host port: %v", err) + } + lc := NewLockClient(dialOption, pb.ServerAddress(fmt.Sprintf("%s:0.%s", host, port))) + lc.noLockServerRetryPeriod = 200 * time.Millisecond + + done := make(chan *LiveLock, 1) + go func() { + done <- lc.NewShortLivedLock("test-key", "test-owner") + }() + select { + case <-done: + t.Fatal("ordinary lock contention must not hit the unavailability bound") + case <-time.After(3 * lc.noLockServerRetryPeriod): + } +} + +// A ring member that refuses connections is unavailability, not contention: +// acquisition fails on the same bound as "no lock server found". +func TestNewShortLivedLockFailsFastOnUnreachableFiler(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + addr := listener.Addr().String() + listener.Close() + + dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) + host, port, err := net.SplitHostPort(addr) + if err != nil { + t.Fatalf("split host port: %v", err) + } + lc := NewLockClient(dialOption, pb.ServerAddress(fmt.Sprintf("%s:0.%s", host, port))) + lc.noLockServerRetryPeriod = 200 * time.Millisecond + + start := time.Now() + lock := lc.NewShortLivedLock("test-key", "test-owner") + elapsed := time.Since(start) + + if lock != nil { + t.Fatal("expected nil lock when the ring member is unreachable") + } + if elapsed > 10*time.Second { + t.Fatalf("lock acquisition took %v, expected fail-fast", elapsed) + } +} + // LiveLock.generation is a fencing token written and read with 64-bit atomic // operations. On 32-bit platforms (GOARCH=386 and GOARCH=arm) a 64-bit atomic // op requires an 8-byte-aligned address, which Go only guarantees for the diff --git a/weed/cluster/lock_manager/lock_ring.go b/weed/cluster/lock_manager/lock_ring.go index d5079f605..8b87e6243 100644 --- a/weed/cluster/lock_manager/lock_ring.go +++ b/weed/cluster/lock_manager/lock_ring.go @@ -1,6 +1,7 @@ package lock_manager import ( + "slices" "sort" "sync" "time" @@ -54,6 +55,14 @@ func (r *LockRing) SetSnapshot(servers []pb.ServerAddress, version int64) bool { r.Unlock() return false } + // An unchanged member list is only a version refresh: installing it as a + // new snapshot would run the topology-change callback and restart the + // prior-owner window on every periodic rebroadcast. + if len(r.snapshots) > 0 && slices.Equal(servers, r.snapshots[0].servers) { + r.version = version + r.Unlock() + return true + } r.version = version // Update the ring while holding the lock so version and ring state // are always consistent — prevents a concurrent SetSnapshot from @@ -74,6 +83,19 @@ func (r *LockRing) SetSnapshot(servers []pb.ServerAddress, version int64) bool { return true } +// Reset clears only the version gate so the first update from a different +// master always applies: ring versions are per-master monotonic, and a high +// version accepted from a former leader must not reject the new leader's +// view. The ring itself stays installed — writes keep routing to the last +// known owner during the gap instead of every filer treating itself as the +// owner, and the arriving snapshot transitions off it with the usual +// prior-owner window. +func (r *LockRing) Reset() { + r.Lock() + defer r.Unlock() + r.version = 0 +} + // Version returns the current ring version. func (r *LockRing) Version() int64 { r.RLock() diff --git a/weed/cluster/lock_manager/lock_ring_test.go b/weed/cluster/lock_manager/lock_ring_test.go index 39c65c9da..61da67232 100644 --- a/weed/cluster/lock_manager/lock_ring_test.go +++ b/weed/cluster/lock_manager/lock_ring_test.go @@ -96,3 +96,42 @@ func TestLockRing_VersionRejectsStale(t *testing.T) { assert.True(t, ok) assert.Equal(t, 1, len(r.GetSnapshot())) } + +func TestLockRing_SetSnapshotUnchangedOnlyBumpsVersion(t *testing.T) { + r := NewLockRing(100 * time.Millisecond) + callbacks := 0 + r.SetTakeSnapshotCallback(func(snapshot []pb.ServerAddress) { callbacks++ }) + + assert.True(t, r.SetSnapshot([]pb.ServerAddress{"a:1", "b:2"}, 100)) + assert.Equal(t, 1, callbacks) + + // A periodic rebroadcast with the same members refreshes the version + // without a new snapshot or another topology-change callback. + assert.True(t, r.SetSnapshot([]pb.ServerAddress{"b:2", "a:1"}, 200)) + assert.Equal(t, int64(200), r.Version()) + assert.Equal(t, 1, r.GetSnapshotCount()) + assert.Equal(t, 1, callbacks, "unchanged ring must not fire the topology callback") + + assert.True(t, r.SetSnapshot([]pb.ServerAddress{"a:1", "b:2", "c:3"}, 300)) + assert.Equal(t, 2, callbacks) +} + +func TestLockRing_Reset(t *testing.T) { + r := NewLockRing(100 * time.Millisecond) + + // A high version accepted from a former leader must not reject the new + // leader's view once the client has moved masters. + ok := r.SetSnapshot([]pb.ServerAddress{"a:1", "b:2"}, 100) + assert.True(t, ok) + + r.Reset() + assert.Equal(t, int64(0), r.Version()) + // The operational ring survives the reset: writes keep routing to the + // last known owner until the new leader's snapshot arrives. + assert.Equal(t, 2, len(r.GetSnapshot())) + assert.NotEqual(t, "", string(r.GetPrimary("key"))) + + ok = r.SetSnapshot([]pb.ServerAddress{"c:1"}, 50) + assert.True(t, ok, "lower version from a different master must apply after reset") + assert.Equal(t, 1, len(r.GetSnapshot())) +} diff --git a/weed/cluster/lock_ring_manager.go b/weed/cluster/lock_ring_manager.go index 486e8b07a..036db3c4f 100644 --- a/weed/cluster/lock_ring_manager.go +++ b/weed/cluster/lock_ring_manager.go @@ -11,29 +11,38 @@ import ( const LockRingStabilizationInterval = 1 * time.Second +// LockRingRebroadcastInterval is how often the ring is re-sent even when +// membership has not changed. Broadcasts are otherwise purely event-driven, +// so a single lost or poisoned update would be permanent without this. +const LockRingRebroadcastInterval = 30 * time.Second + // LockRingManager tracks filer membership for the distributed lock ring. // It batches rapid topology changes (e.g., node drop + join) with a // stabilization timer, then broadcasts the complete member list atomically // so filers receive a single consistent ring update instead of multiple // intermediate states. type LockRingManager struct { - mu sync.Mutex - members map[FilerGroupName]map[pb.ServerAddress]struct{} - version map[FilerGroupName]int64 - lastBroadcast map[FilerGroupName]*master_pb.LockRingUpdate - pendingTimer map[FilerGroupName]*time.Timer - broadcastFn func(resp *master_pb.KeepConnectedResponse) - stabilizeDelay time.Duration + mu sync.Mutex + members map[FilerGroupName]map[pb.ServerAddress]struct{} + version map[FilerGroupName]int64 + lastBroadcast map[FilerGroupName]*master_pb.LockRingUpdate + pendingTimer map[FilerGroupName]*time.Timer + rebroadcastTimer map[FilerGroupName]*time.Timer + broadcastFn func(resp *master_pb.KeepConnectedResponse) + stabilizeDelay time.Duration + rebroadcastInterval time.Duration } func NewLockRingManager(broadcastFn func(resp *master_pb.KeepConnectedResponse)) *LockRingManager { return &LockRingManager{ - members: make(map[FilerGroupName]map[pb.ServerAddress]struct{}), - version: make(map[FilerGroupName]int64), - lastBroadcast: make(map[FilerGroupName]*master_pb.LockRingUpdate), - pendingTimer: make(map[FilerGroupName]*time.Timer), - broadcastFn: broadcastFn, - stabilizeDelay: LockRingStabilizationInterval, + members: make(map[FilerGroupName]map[pb.ServerAddress]struct{}), + version: make(map[FilerGroupName]int64), + lastBroadcast: make(map[FilerGroupName]*master_pb.LockRingUpdate), + pendingTimer: make(map[FilerGroupName]*time.Timer), + rebroadcastTimer: make(map[FilerGroupName]*time.Timer), + broadcastFn: broadcastFn, + stabilizeDelay: LockRingStabilizationInterval, + rebroadcastInterval: LockRingRebroadcastInterval, } } @@ -116,16 +125,64 @@ func (lrm *LockRingManager) scheduleBroadcast(filerGroup FilerGroupName) { func (lrm *LockRingManager) doBroadcast(filerGroup FilerGroupName) { lrm.mu.Lock() + delete(lrm.pendingTimer, filerGroup) + lrm.mu.Unlock() + lrm.emit(filerGroup) +} + +// rebroadcast re-sends the current ring unless a membership broadcast is +// still stabilizing — emitting mid-window would publish an intermediate +// topology that the pending timer immediately replaces. The check and the +// update must sit in one critical section or a membership change can slip +// a pending timer in between. +func (lrm *LockRingManager) rebroadcast(filerGroup FilerGroupName) { + lrm.mu.Lock() + var update *master_pb.LockRingUpdate + if _, pending := lrm.pendingTimer[filerGroup]; !pending { + update = lrm.nextBroadcastUpdate(filerGroup) + } + lrm.mu.Unlock() + lrm.sendUpdate(filerGroup, update) +} + +func (lrm *LockRingManager) emit(filerGroup FilerGroupName) { + lrm.mu.Lock() + update := lrm.nextBroadcastUpdate(filerGroup) + lrm.mu.Unlock() + lrm.sendUpdate(filerGroup, update) +} + +func (lrm *LockRingManager) sendUpdate(filerGroup FilerGroupName, update *master_pb.LockRingUpdate) { + if update == nil { + return + } + glog.V(0).Infof("LockRing: broadcasting ring update for group %q version %d: %v", filerGroup, update.Version, update.Servers) + if lrm.broadcastFn != nil { + lrm.broadcastFn(&master_pb.KeepConnectedResponse{ + LockRingUpdate: update, + }) + } +} + +// nextBroadcastUpdate stamps the current members into an update and re-arms +// the periodic rebroadcast. It returns nil for an empty member list: an empty +// lock ring is never usable, so a late "last member removed" event from a +// former leader must not propagate and wedge every lock client. The last +// non-empty broadcast stays in lastBroadcast for reconnecting clients. +// Caller must hold lrm.mu. +func (lrm *LockRingManager) nextBroadcastUpdate(filerGroup FilerGroupName) *master_pb.LockRingUpdate { + members := lrm.members[filerGroup] + if len(members) == 0 { + return nil + } // Use wall-clock nanoseconds so the version survives master restarts // without persistence — a restarted master produces a version greater // than any pre-restart value (assuming clocks don't jump backward). version := time.Now().UnixNano() lrm.version[filerGroup] = version - servers := make([]string, 0) - if members, ok := lrm.members[filerGroup]; ok { - for addr := range members { - servers = append(servers, string(addr)) - } + servers := make([]string, 0, len(members)) + for addr := range members { + servers = append(servers, string(addr)) } update := &master_pb.LockRingUpdate{ FilerGroup: string(filerGroup), @@ -133,16 +190,13 @@ func (lrm *LockRingManager) doBroadcast(filerGroup FilerGroupName) { Version: version, } lrm.lastBroadcast[filerGroup] = update - delete(lrm.pendingTimer, filerGroup) - lrm.mu.Unlock() - - glog.V(0).Infof("LockRing: broadcasting ring update for group %q version %d: %v", filerGroup, version, servers) - - if lrm.broadcastFn != nil { - lrm.broadcastFn(&master_pb.KeepConnectedResponse{ - LockRingUpdate: update, - }) + if timer, ok := lrm.rebroadcastTimer[filerGroup]; ok { + timer.Stop() } + lrm.rebroadcastTimer[filerGroup] = time.AfterFunc(lrm.rebroadcastInterval, func() { + lrm.rebroadcast(filerGroup) + }) + return update } // FlushPending fires any pending timer immediately (for testing or shutdown). diff --git a/weed/cluster/lock_ring_manager_test.go b/weed/cluster/lock_ring_manager_test.go index 981d13d9c..91fcdb98e 100644 --- a/weed/cluster/lock_ring_manager_test.go +++ b/weed/cluster/lock_ring_manager_test.go @@ -213,6 +213,111 @@ func TestLockRingManager_NoBroadcastWithoutFn(t *testing.T) { time.Sleep(50 * time.Millisecond) // should not panic } +func TestLockRingManager_EmptyRingNotBroadcast(t *testing.T) { + var mu sync.Mutex + var broadcasts []*master_pb.LockRingUpdate + + lrm := NewLockRingManager(func(resp *master_pb.KeepConnectedResponse) { + mu.Lock() + if resp.LockRingUpdate != nil { + broadcasts = append(broadcasts, resp.LockRingUpdate) + } + mu.Unlock() + }) + lrm.stabilizeDelay = 50 * time.Millisecond + + group := FilerGroupName("default") + + lrm.AddServer(group, "filer1:8888") + lrm.FlushPending(group) + + mu.Lock() + broadcasts = nil + mu.Unlock() + + // Removing the last member must not propagate an empty ring. + lrm.RemoveServer(group, "filer1:8888") + time.Sleep(100 * time.Millisecond) + + mu.Lock() + assert.Equal(t, 0, len(broadcasts), "empty ring must not be broadcast") + mu.Unlock() + + // The last non-empty snapshot is still served to reconnecting clients. + update := lrm.GetLastUpdate(group) + require.NotNil(t, update) + assert.Equal(t, []string{"filer1:8888"}, update.Servers) +} + +func TestLockRingManager_PeriodicRebroadcast(t *testing.T) { + var mu sync.Mutex + var broadcasts []*master_pb.LockRingUpdate + + lrm := NewLockRingManager(func(resp *master_pb.KeepConnectedResponse) { + mu.Lock() + if resp.LockRingUpdate != nil { + broadcasts = append(broadcasts, resp.LockRingUpdate) + } + mu.Unlock() + }) + lrm.stabilizeDelay = 20 * time.Millisecond + lrm.rebroadcastInterval = 60 * time.Millisecond + + group := FilerGroupName("default") + lrm.AddServer(group, "filer1:8888") + + // Without any further membership change, the ring keeps being re-sent so + // a lost or poisoned update cannot be permanent. + time.Sleep(200 * time.Millisecond) + + mu.Lock() + require.GreaterOrEqual(t, len(broadcasts), 2, "ring should rebroadcast periodically") + for i := 1; i < len(broadcasts); i++ { + assert.Greater(t, broadcasts[i].Version, broadcasts[i-1].Version) + } + mu.Unlock() +} + +func TestLockRingManager_RebroadcastDefersToPendingStabilization(t *testing.T) { + var mu sync.Mutex + var broadcasts []*master_pb.LockRingUpdate + + lrm := NewLockRingManager(func(resp *master_pb.KeepConnectedResponse) { + mu.Lock() + if resp.LockRingUpdate != nil { + broadcasts = append(broadcasts, resp.LockRingUpdate) + } + mu.Unlock() + }) + lrm.stabilizeDelay = 100 * time.Millisecond + lrm.rebroadcastInterval = 30 * time.Millisecond + + group := FilerGroupName("default") + lrm.AddServer(group, "filer1:8888") + lrm.FlushPending(group) + + mu.Lock() + require.Len(t, broadcasts, 1) + mu.Unlock() + + // A membership change just before the periodic tick: the rebroadcast must + // not publish the unsettled ring ahead of the stabilization timer. + lrm.RemoveServer(group, "filer1:8888") + lrm.AddServer(group, "filer2:8888") + time.Sleep(2 * lrm.rebroadcastInterval) + + mu.Lock() + assert.Len(t, broadcasts, 1, "rebroadcast during stabilization should be deferred") + mu.Unlock() + + time.Sleep(2 * lrm.stabilizeDelay) + + mu.Lock() + require.GreaterOrEqual(t, len(broadcasts), 2) + assert.Equal(t, []string{"filer2:8888"}, broadcasts[1].Servers) + mu.Unlock() +} + func TestLockRingManager_GetLastUpdateReturnsBroadcastState(t *testing.T) { lrm := NewLockRingManager(nil) diff --git a/weed/filer/filer.go b/weed/filer/filer.go index 0eab4ef7e..88943d4c0 100644 --- a/weed/filer/filer.go +++ b/weed/filer/filer.go @@ -181,6 +181,10 @@ func (f *Filer) AggregateFromPeers(self pb.ServerAddress, existingNodes []*maste glog.V(0).Infof("LockRing: applying master ring update v%d: %v", update.Version, servers) f.Dlm.LockRing.SetSnapshot(servers, update.Version) }) + f.MasterClient.SetOnMasterChangeFn(func(previous, current pb.ServerAddress) { + glog.V(0).Infof("LockRing: master changed %s -> %s, resetting ring", previous, current) + f.Dlm.LockRing.Reset() + }) // Subscribe to the local filer first: its events reach the aggregated // buffer only through this subscription, and the peer watermarks must diff --git a/weed/mount/filehandle.go b/weed/mount/filehandle.go index 8514b505e..d40d7c035 100644 --- a/weed/mount/filehandle.go +++ b/weed/mount/filehandle.go @@ -211,6 +211,9 @@ func (fh *FileHandle) AddChunks(chunks []*filer_pb.FileChunk) { } func (fh *FileHandle) ReleaseHandle() { + fhActiveLock := fh.wfs.fhLockTable.AcquireLock("ReleaseHandle", fh.fh, util.ExclusiveLock) + defer fh.wfs.fhLockTable.ReleaseLock(fh.fh, fhActiveLock) + // Release distributed lock before cleaning up, so other mounts can // proceed as soon as this handle is done flushing. if fh.dlmLock != nil { @@ -219,9 +222,6 @@ func (fh *FileHandle) ReleaseHandle() { glog.V(1).Infof("DLM lock released for inode %d", fh.inode) } - fhActiveLock := fh.wfs.fhLockTable.AcquireLock("ReleaseHandle", fh.fh, util.ExclusiveLock) - defer fh.wfs.fhLockTable.ReleaseLock(fh.fh, fhActiveLock) - if fh.entryChunkGroup != nil { _ = fh.entryChunkGroup.Close() } diff --git a/weed/mount/weedfs_file_mkrm.go b/weed/mount/weedfs_file_mkrm.go index 7db012365..0e7e02db1 100644 --- a/weed/mount/weedfs_file_mkrm.go +++ b/weed/mount/weedfs_file_mkrm.go @@ -9,6 +9,7 @@ import ( "github.com/seaweedfs/go-fuse/v2/fuse" "google.golang.org/protobuf/proto" + "github.com/seaweedfs/seaweedfs/weed/cluster" "github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager" "github.com/seaweedfs/seaweedfs/weed/filer" "github.com/seaweedfs/seaweedfs/weed/glog" @@ -72,10 +73,26 @@ func (wfs *WFS) Create(cancel <-chan struct{}, in *fuse.CreateIn, name string, o return code } + // Acquire the DLM lock before the filer create: a lock failure after an + // eager create would return EAGAIN while the file stays persisted. + var dlmLock *cluster.LiveLock + if wfs.lockClient != nil { + owner := fmt.Sprintf("mount-%d", wfs.signature) + dlmLock = wfs.lockClient.NewBlockingLongLivedLock( + string(entryFullPath), owner, lock_manager.LiveLockTTL, + ) + if dlmLock == nil { + return fuse.Status(syscall.EAGAIN) + } + } + inode, newEntry, code = wfs.createRegularFile(dirFullPath, name, in.Mode, in.Uid, in.Gid, 0, !wfs.option.EagerFilerCreate, true) if code == fuse.Status(syscall.EEXIST) && in.Flags&syscall.O_EXCL == 0 { // Race: another process created the file between our check and create. - // Reopen the winner's entry. + // Reopen the winner's entry; AcquireHandle takes its own lock. + if dlmLock != nil { + dlmLock.Stop() + } newEntry, _, code = wfs.maybeLoadEntry(entryFullPath) if code != fuse.OK { return code @@ -100,6 +117,9 @@ func (wfs *WFS) Create(cancel <-chan struct{}, in *fuse.CreateIn, name string, o out.OpenFlags = 0 return fuse.OK } else if code != fuse.OK { + if dlmLock != nil { + dlmLock.Stop() + } return code } else { inode = wfs.inodeToPath.Lookup(entryFullPath, newEntry.Attributes.Crtime, false, false, inode, true) @@ -122,15 +142,15 @@ func (wfs *WFS) Create(cancel <-chan struct{}, in *fuse.CreateIn, name string, o // persisted the entry, so its handle starts clean. fileHandle.dirtyMetadata = !wfs.option.EagerFilerCreate - // Acquire DLM lock for new file creation (Create bypasses AcquireHandle - // so we must acquire the lock here). Always lock on Create since file - // creation is inherently a write operation. - if wfs.lockClient != nil && fileHandle.dlmLock == nil { - owner := fmt.Sprintf("mount-%d", wfs.signature) - fileHandle.dlmLock = wfs.lockClient.NewBlockingLongLivedLock( - string(entryFullPath), owner, lock_manager.LiveLockTTL, - ) - glog.V(1).Infof("DLM lock acquired for new file %s", entryFullPath) + // Create bypasses AcquireHandle, so attach the lock acquired above. + // A surviving handle may already hold one; ours is redundant then. + if dlmLock != nil { + if fileHandle.dlmLock == nil { + fileHandle.dlmLock = dlmLock + glog.V(1).Infof("DLM lock acquired for new file %s", entryFullPath) + } else { + dlmLock.Stop() + } } out.Fh = uint64(fileHandle.fh) diff --git a/weed/mount/weedfs_filehandle.go b/weed/mount/weedfs_filehandle.go index 7009d982d..d7f16b76e 100644 --- a/weed/mount/weedfs_filehandle.go +++ b/weed/mount/weedfs_filehandle.go @@ -2,6 +2,7 @@ package mount import ( "fmt" + "syscall" "github.com/seaweedfs/go-fuse/v2/fuse" "github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager" @@ -77,6 +78,12 @@ func (wfs *WFS) AcquireHandle(inode uint64, flags, uid, gid uint32) (fileHandle fileHandle.dlmLock = wfs.lockClient.NewBlockingLongLivedLock( string(path), owner, lock_manager.LiveLockTTL, ) + if fileHandle.dlmLock == nil { + // No lock server is reachable: proceeding would silently drop + // cross-mount write serialization, so fail the open instead. + wfs.fhMap.ReleaseByHandle(fileHandle.fh) + return nil, fuse.Status(syscall.EAGAIN) + } glog.V(1).Infof("DLM lock acquired for %s", path) } return fileHandle, fuse.OK diff --git a/weed/mount/weedfs_rename.go b/weed/mount/weedfs_rename.go index e4fd86682..6a000f3a2 100644 --- a/weed/mount/weedfs_rename.go +++ b/weed/mount/weedfs_rename.go @@ -9,6 +9,7 @@ import ( "syscall" "github.com/seaweedfs/go-fuse/v2/fuse" + "github.com/seaweedfs/seaweedfs/weed/cluster" "github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager" "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" @@ -16,10 +17,10 @@ import ( ) // doRename tries the streaming mux first, falling back to unary on transport errors. -func (wfs *WFS) doRename(ctx context.Context, request *filer_pb.StreamRenameEntryRequest, oldPath, newPath util.FullPath) error { +func (wfs *WFS) doRename(ctx context.Context, request *filer_pb.StreamRenameEntryRequest, oldPath, newPath util.FullPath, newPathLock **cluster.LiveLock) error { if wfs.streamMutate != nil && wfs.streamMutate.IsAvailable() { err := wfs.streamMutate.Rename(ctx, request, func(resp *filer_pb.StreamRenameEntryResponse) error { - return wfs.handleRenameResponse(ctx, resp) + return wfs.handleRenameResponse(ctx, resp, newPath, newPathLock) }) if err == nil || !errors.Is(err, ErrStreamTransport) { return err // success or application error @@ -39,7 +40,7 @@ func (wfs *WFS) doRename(ctx context.Context, request *filer_pb.StreamRenameEntr } return fmt.Errorf("dir Rename %s => %s receive: %v", oldPath, newPath, recvErr) } - if err := wfs.handleRenameResponse(ctx, resp); err != nil { + if err := wfs.handleRenameResponse(ctx, resp, newPath, newPathLock); err != nil { return err } } @@ -240,6 +241,102 @@ func (wfs *WFS) Rename(cancel <-chan struct{}, in *fuse.RenameIn, oldName string glog.V(4).Infof("dir Rename %s => %s", oldPath, newPath) + // Acquire DLM locks on both old and new paths to prevent another mount + // from opening either path for writing during the rename. Lock in + // sorted order to prevent deadlocks when two mounts rename in opposite + // directions (A→B vs B→A). + // + // Skip the old-path lock if this mount already holds it via an open + // file handle (otherwise we'd deadlock trying to re-acquire our own lock). + // Acquiring before the handle marks below keeps a lock failure from + // leaving source handles flagged for a rename that never happened. + var heldLocks []*cluster.LiveLock + var newPathLock *cluster.LiveLock + if wfs.lockClient != nil { + // A handle open on the target holds its lock; claiming it up front + // keeps the path locked through the rename (its close would otherwise + // release it mid-flight) without waiting on our own ownership. + var newPathLockFh *FileHandle + defer func() { + for _, l := range heldLocks { + l.Stop() + } + if newPathLock != nil { + // Never adopted by a migrated handle — return it to the + // handle it came from if that is still open. + if newPathLockFh != nil { + if cur, ok := wfs.fhMap.FindFileHandle(newPathLockFh.inode); ok && cur == newPathLockFh { + lk := wfs.fhLockTable.AcquireLock("renameDLM", cur.fh, util.ExclusiveLock) + if cur.dlmLock == nil { + cur.dlmLock = newPathLock + newPathLock = nil + } + wfs.fhLockTable.ReleaseLock(cur.fh, lk) + } + } + if newPathLock != nil { + newPathLock.Stop() + } + } + }() + owner := fmt.Sprintf("mount-%d", wfs.signature) + + // Check if the source file handle already holds a DLM lock on oldPath + oldPathAlreadyLocked := false + sourceInode, sourceMapped := wfs.inodeToPath.GetInode(oldPath) + if !sourceMapped && oldEntry != nil && oldEntry.Attributes != nil { + sourceInode = oldEntry.Attributes.Inode + } + if sourceInode != 0 { + if fh, ok := wfs.fhMap.FindFileHandle(sourceInode); ok { + lk := wfs.fhLockTable.AcquireLock("renameDLM", fh.fh, util.ExclusiveLock) + oldPathAlreadyLocked = fh.dlmLock != nil + wfs.fhLockTable.ReleaseLock(fh.fh, lk) + } + } + targetInode, targetMapped := wfs.inodeToPath.GetInode(newPath) + if !targetMapped && newEntry != nil && newEntry.Attributes != nil { + targetInode = newEntry.Attributes.Inode + } + if targetInode != 0 && targetInode != sourceInode { + if targetFh, ok := wfs.fhMap.FindFileHandle(targetInode); ok { + targetFhLock := wfs.fhLockTable.AcquireLock("renameDLM", targetFh.fh, util.ExclusiveLock) + if targetFh.dlmLock != nil { + newPathLock = targetFh.dlmLock + targetFh.dlmLock = nil + newPathLockFh = targetFh + } + wfs.fhLockTable.ReleaseLock(targetFh.fh, targetFhLock) + } + } + + // Determine which paths need new DLM locks + pathsToLock := []string{} + if newPathLock == nil { + pathsToLock = append(pathsToLock, string(newPath)) + } + if !oldPathAlreadyLocked && string(oldPath) != string(newPath) { + pathsToLock = append(pathsToLock, string(oldPath)) + } + // Sort for consistent lock ordering + if len(pathsToLock) == 2 && pathsToLock[0] > pathsToLock[1] { + pathsToLock[0], pathsToLock[1] = pathsToLock[1], pathsToLock[0] + } + + for _, p := range pathsToLock { + dlmLock := wfs.lockClient.NewBlockingLongLivedLock(p, owner, lock_manager.LiveLockTTL) + if dlmLock == nil { + return fuse.Status(syscall.EAGAIN) + } + if p == string(newPath) { + newPathLock = dlmLock + } else { + heldLocks = append(heldLocks, dlmLock) + } + } + glog.V(1).Infof("DLM locks acquired for rename %s => %s (oldPathAlreadyLocked=%v)", oldPath, newPath, oldPathAlreadyLocked) + } + // Ensure the source file's metadata exists on the filer before renaming. // Two cases can leave the entry only in the local cache: // 1. deferFilerCreate=true — file handle still open, dirtyMetadata set. @@ -300,41 +397,6 @@ func (wfs *WFS) Rename(cancel <-chan struct{}, in *fuse.RenameIn, oldName string } } - // Acquire DLM locks on both old and new paths to prevent another mount - // from opening either path for writing during the rename. Lock in - // sorted order to prevent deadlocks when two mounts rename in opposite - // directions (A→B vs B→A). - // - // Skip the old-path lock if this mount already holds it via an open - // file handle (otherwise we'd deadlock trying to re-acquire our own lock). - if wfs.lockClient != nil { - owner := fmt.Sprintf("mount-%d", wfs.signature) - - // Check if the source file handle already holds a DLM lock on oldPath - oldPathAlreadyLocked := false - if sourceInode, found := wfs.inodeToPath.GetInode(oldPath); found { - if fh, ok := wfs.fhMap.FindFileHandle(sourceInode); ok && fh.dlmLock != nil { - oldPathAlreadyLocked = true - } - } - - // Determine which paths need new DLM locks - pathsToLock := []string{string(newPath)} - if !oldPathAlreadyLocked { - pathsToLock = append(pathsToLock, string(oldPath)) - } - // Sort for consistent lock ordering - if len(pathsToLock) == 2 && pathsToLock[0] > pathsToLock[1] { - pathsToLock[0], pathsToLock[1] = pathsToLock[1], pathsToLock[0] - } - - for _, p := range pathsToLock { - dlmLock := wfs.lockClient.NewBlockingLongLivedLock(p, owner, lock_manager.LiveLockTTL) - defer dlmLock.Stop() - } - glog.V(1).Infof("DLM locks acquired for rename %s => %s (oldPathAlreadyLocked=%v)", oldPath, newPath, oldPathAlreadyLocked) - } - // update remote filer request := &filer_pb.StreamRenameEntryRequest{ OldDirectory: string(oldDir), @@ -345,7 +407,7 @@ func (wfs *WFS) Rename(cancel <-chan struct{}, in *fuse.RenameIn, oldName string } ctx := context.Background() - err := wfs.doRename(ctx, request, oldPath, newPath) + err := wfs.doRename(ctx, request, oldPath, newPath, &newPathLock) if err != nil { glog.V(0).Infof("Rename %s => %s: %v", oldPath, newPath, err) // Map error strings to FUSE status codes. String matching is used @@ -377,7 +439,7 @@ func (wfs *WFS) Rename(cancel <-chan struct{}, in *fuse.RenameIn, oldName string } -func (wfs *WFS) handleRenameResponse(ctx context.Context, resp *filer_pb.StreamRenameEntryResponse) error { +func (wfs *WFS) handleRenameResponse(ctx context.Context, resp *filer_pb.StreamRenameEntryResponse, renameNewPath util.FullPath, renameNewPathLock **cluster.LiveLock) error { // comes from filer StreamRenameEntry, can only be create or delete entry glog.V(4).Infof("dir Rename %+v", resp.EventNotification) @@ -411,16 +473,32 @@ func (wfs *WFS) handleRenameResponse(ctx context.Context, resp *filer_pb.StreamR // Migrate the DLM lock from old path to new path so the // lock key matches the current file location. Hold the // fhLockTable to prevent ReleaseHandle from concurrently - // stopping the lock during migration. - if wfs.lockClient != nil { + // stopping the lock during migration. Acquire before + // releasing: a failed migration keeps the old lock rather + // than leaving the handle unlocked. + if wfs.lockClient != nil && oldPath != newPath { fhActiveLock := wfs.fhLockTable.AcquireLock("renameDLM", fh.fh, util.ExclusiveLock) if fh.dlmLock != nil { - owner := fmt.Sprintf("mount-%d", wfs.signature) - fh.dlmLock.Stop() - fh.dlmLock = wfs.lockClient.NewBlockingLongLivedLock( - string(newPath), owner, lock_manager.LiveLockTTL, - ) - glog.V(1).Infof("DLM lock migrated from %s to %s", oldPath, newPath) + var newLock *cluster.LiveLock + if newPath == renameNewPath && renameNewPathLock != nil && *renameNewPathLock != nil { + // The rename already holds a lock on the target; + // adopt it — re-acquiring our own lock would block + // until this handle is released. + newLock = *renameNewPathLock + *renameNewPathLock = nil + } else { + owner := fmt.Sprintf("mount-%d", wfs.signature) + newLock = wfs.lockClient.NewBlockingLongLivedLock( + string(newPath), owner, lock_manager.LiveLockTTL, + ) + } + if newLock != nil { + fh.dlmLock.Stop() + fh.dlmLock = newLock + glog.V(1).Infof("DLM lock migrated from %s to %s", oldPath, newPath) + } else { + glog.Warningf("DLM lock migration to %s failed; keeping lock on %s", newPath, oldPath) + } } wfs.fhLockTable.ReleaseLock(fh.fh, fhActiveLock) } diff --git a/weed/mount/weedfs_rename_test.go b/weed/mount/weedfs_rename_test.go index 406a093c3..6eb47ca8a 100644 --- a/weed/mount/weedfs_rename_test.go +++ b/weed/mount/weedfs_rename_test.go @@ -69,7 +69,7 @@ func TestHandleRenameResponseLeavesUncachedTargetOutOfCache(t *testing.T) { }, } - if err := wfs.handleRenameResponse(context.Background(), resp); err != nil { + if err := wfs.handleRenameResponse(context.Background(), resp, targetPath, nil); err != nil { t.Fatalf("handle rename response: %v", err) } diff --git a/weed/s3api/s3api_object_handlers_put.go b/weed/s3api/s3api_object_handlers_put.go index 30a9cfcf9..9883187d8 100644 --- a/weed/s3api/s3api_object_handlers_put.go +++ b/weed/s3api/s3api_object_handlers_put.go @@ -385,7 +385,11 @@ func (s3a *S3ApiServer) withObjectWriteLock(bucket, object string, preconditionF return fn() } - lock := s3a.newObjectWriteLock(bucket, object) + lock, err := s3a.newObjectWriteLock(bucket, object) + if err != nil { + glog.Warningf("withObjectWriteLock: %v", err) + return s3err.ErrServiceUnavailable + } if lock == nil { if errCode := runPrecondition(); errCode != s3err.ErrNone { return errCode diff --git a/weed/s3api/s3api_server.go b/weed/s3api/s3api_server.go index 6f060636b..b9c6205ee 100644 --- a/weed/s3api/s3api_server.go +++ b/weed/s3api/s3api_server.go @@ -108,7 +108,7 @@ type S3ApiServer struct { icebergCredentialRole string icebergCredentialDuration int64 cipher bool // encrypt data on volume servers - newObjectWriteLock func(bucket, object string) objectWriteLock + newObjectWriteLock func(bucket, object string) (objectWriteLock, error) // objectWriteLockClient resolves a key's owner filer for route-by-key. objectWriteLockClient *cluster.LockClient // unreachableOwners holds owners (pb.ServerAddress -> expiry time.Time) whose @@ -225,6 +225,9 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl } objectWriteLockClient.SetRing(servers, update.Version) }) + masterClient.SetOnMasterChangeFn(func(previous, current pb.ServerAddress) { + objectWriteLockClient.ResetRing() + }) } // Start the master client connection loop - required for GetMaster() to work go masterClient.KeepConnectedToMaster(context.Background()) @@ -345,16 +348,19 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl objectWriteLockClient = cluster.NewLockClient(option.GrpcDialOption, option.Filers[0]) } s3ApiServer.objectWriteLockClient = objectWriteLockClient - s3ApiServer.newObjectWriteLock = func(bucket, object string) objectWriteLock { + s3ApiServer.newObjectWriteLock = func(bucket, object string) (objectWriteLock, error) { lockKey := objectWriteRouteKeyPrefix + s3ApiServer.toFilerPath(bucket, object) owner := fmt.Sprintf("s3api-%d", s3ApiServer.randomClientId) lock := objectWriteLockClient.NewShortLivedLock(lockKey, owner) + if lock == nil { + return nil, fmt.Errorf("objectWriteLock: failed to acquire lock for %s", lockKey) + } if err := lock.AttemptToLock(objectWriteLockTTL); err != nil { // The initial acquisition already succeeded with the default short TTL. // Renewal to a longer TTL is opportunistic to cover slower metadata paths. glog.Warningf("objectWriteLock: failed to extend lock TTL for %s: %v", lockKey, err) } - return lock + return lock, nil } } diff --git a/weed/server/filer_grpc_server.go b/weed/server/filer_grpc_server.go index e99c9abe2..c70b02bc8 100644 --- a/weed/server/filer_grpc_server.go +++ b/weed/server/filer_grpc_server.go @@ -816,6 +816,9 @@ func (fs *FilerServer) AppendToEntry(ctx context.Context, req *filer_pb.AppendTo lockClient := cluster.NewLockClient(fs.grpcDialOption, fs.option.Host) lock := lockClient.NewShortLivedLock(string(fullpath), string(fs.option.Host)) + if lock == nil { + return nil, fmt.Errorf("failed to acquire lock for %s", fullpath) + } defer lock.StopShortLivedLock() // The cluster lock serializes appenders across filers; the path lock makes diff --git a/weed/wdclient/masterclient.go b/weed/wdclient/masterclient.go index beced07cf..6d638b6d3 100644 --- a/weed/wdclient/masterclient.go +++ b/weed/wdclient/masterclient.go @@ -150,6 +150,7 @@ type MasterClient struct { clientHost pb.ServerAddress rack string currentMaster pb.ServerAddress + lastServedMaster pb.ServerAddress currentMasterLock sync.RWMutex masters pb.ServerDiscovery grpcDialOption grpc.DialOption @@ -158,6 +159,8 @@ type MasterClient struct { OnPeerUpdateLock sync.RWMutex OnLockRingUpdate func(update *master_pb.LockRingUpdate) OnLockRingUpdateLock sync.RWMutex + OnMasterChange func(previous, current pb.ServerAddress) + OnMasterChangeLock sync.RWMutex } func NewMasterClient(grpcDialOption grpc.DialOption, filerGroup string, clientType string, clientHost pb.ServerAddress, clientDataCenter string, rack string, masters pb.ServerDiscovery) *MasterClient { @@ -192,6 +195,12 @@ func (mc *MasterClient) SetOnLockRingUpdateFn(fn func(update *master_pb.LockRing mc.OnLockRingUpdateLock.Unlock() } +func (mc *MasterClient) SetOnMasterChangeFn(fn func(previous, current pb.ServerAddress)) { + mc.OnMasterChangeLock.Lock() + mc.OnMasterChange = fn + mc.OnMasterChangeLock.Unlock() +} + func (mc *MasterClient) tryAllMasters(ctx context.Context) { var nextHintedLeader pb.ServerAddress failedMasters := make(map[pb.ServerAddress]struct{}) @@ -280,7 +289,14 @@ func (mc *MasterClient) tryConnectToMaster(ctx context.Context, master pb.Server // Still need to reset cache to ensure we don't use stale data from previous master mc.resetVidMap() } - mc.setCurrentMaster(master) + if previous := mc.markServingMaster(master); previous != "" { + mc.OnMasterChangeLock.RLock() + if mc.OnMasterChange != nil { + glog.V(0).Infof("%s.%s masterClient master changed %s -> %s", mc.FilerGroup, mc.clientType, previous, master) + mc.OnMasterChange(previous, master) + } + mc.OnMasterChangeLock.RUnlock() + } for { resp, err := stream.Recv() @@ -500,6 +516,21 @@ func (mc *MasterClient) setCurrentMaster(master pb.ServerAddress) { mc.currentMasterLock.Unlock() } +// markServingMaster records the master now serving this client and returns +// the previously served one when it differs. Unlike currentMaster, +// lastServedMaster survives the disconnected gap between reconnect attempts, +// so a leader change is still detected. +func (mc *MasterClient) markServingMaster(master pb.ServerAddress) (previous pb.ServerAddress) { + mc.currentMasterLock.Lock() + defer mc.currentMasterLock.Unlock() + if mc.lastServedMaster != "" && !mc.lastServedMaster.Equals(master) { + previous = mc.lastServedMaster + } + mc.currentMaster = master + mc.lastServedMaster = master + return +} + // GetMaster returns the current master address, blocking until connected. // // IMPORTANT: This method blocks until KeepConnectedToMaster successfully establishes diff --git a/weed/wdclient/masterclient_test.go b/weed/wdclient/masterclient_test.go index c6450ce64..dc10f50e9 100644 --- a/weed/wdclient/masterclient_test.go +++ b/weed/wdclient/masterclient_test.go @@ -155,6 +155,28 @@ func TestWithClientStopsWaitingOnDeadline(t *testing.T) { } } +// TestMarkServingMaster verifies leader-change detection survives the +// disconnected gap where currentMaster is cleared between reconnect attempts. +func TestMarkServingMaster(t *testing.T) { + mc := NewMasterClient(grpc.EmptyDialOption{}, "test-group", "test-client", "", "", "", pb.ServerDiscovery{}) + + if prev := mc.markServingMaster("master1:9333"); prev != "" { + t.Fatalf("first connect should not report a change, got %q", prev) + } + if prev := mc.markServingMaster("master1:9333"); prev != "" { + t.Fatalf("reconnect to same master should not report a change, got %q", prev) + } + + // Simulate the gap between stream attempts. + mc.setCurrentMaster("") + if prev := mc.markServingMaster("master2:9333"); prev != "master1:9333" { + t.Fatalf("leader change across disconnect not detected, got %q", prev) + } + if got := mc.getCurrentMaster(); got != "master2:9333" { + t.Fatalf("current master = %q", got) + } +} + // TestWithClientStopsBackoffOnCancel verifies that a cancellation arriving while // the retry is backing off cuts the backoff short rather than sleeping it out. func TestWithClientStopsBackoffOnCancel(t *testing.T) {