master,filer: fix lock ring poisoning after leader change (#11453)

* cluster: never broadcast an empty lock ring

An empty member list is never a usable ring state, but a delayed
RemoveServer on a former leader can fire after the new leader already
broadcast the recovered ring. That late broadcast carries a newer
wall-clock version, so clients accept the empty ring and permanently
reject the good one.

Skip the broadcast entirely when the member list is empty, keeping the
last non-empty snapshot for reconnecting clients.

* cluster: periodically rebroadcast the lock ring

Ring updates are purely event-driven, so one lost or poisoned update is
permanent until the next membership change — with a single filer that may
never come. Re-arm a per-group timer after every broadcast so the current
leader keeps re-sending the ring; clients reject nothing newer than their
last accepted version, so a re-sent snapshot always heals a stale view.

* filer,s3api: reset the lock ring on master change

Ring versions are per-master monotonic — each master stamps wall-clock
nanoseconds — so a late high-version update accepted from a former leader
makes the new leader's snapshot look stale forever. Detect a leader
change across the reconnect gap (currentMaster is cleared between
attempts, so remember the last served master) and reset the ring to
bootstrap state so the new leader's view always applies.

* cluster: fail lock acquisition when no lock server exists

retryUntilLocked loops forever, so a filer reporting an empty lock ring
wedges every append write indefinitely. Bound only the "no lock server
found" case — ordinary contention is still waited out since the holder
releases eventually. The constructors now return nil on failure: the
filer append path and S3 object writes fail fast, while mounts degrade
to their existing lockless mode.

* cluster: reset only the ring version on master change

Ring versions are per-master monotonic, so a version gate reset is all a
leader change needs. Clearing the whole ring made every filer its own
write owner until the next update and dropped the prior-owner window for
keys the new leader remaps; the last ring now keeps routing until the
new leader's snapshot transitions off it.

* cluster: skip redundant ring installs and defer rebroadcasts

An unchanged member list now only bumps the accepted version instead of
installing a snapshot: periodic rebroadcasts no longer fire the
topology-change callback or restart the prior-owner window. And a
rebroadcast that lands inside a membership stabilization window yields
to the pending timer rather than publishing an intermediate ring.

* cluster,mount: bound lock unavailability, fail ops that cannot lock

Only 'lock already owned' contention retries without bound now; every
other failure — no lock server, or a dead ring member refusing
connections — shares the same unavailability budget, so a ring naming
departed filers can no longer hang a lock forever. Mount open-write,
create, and rename fail with EAGAIN when the required lock cannot be
acquired instead of proceeding without cross-mount serialization.

* cluster: check pending stabilization inside the broadcast critical section

rebroadcast released the mutex between the pending-timer check and
nextBroadcastUpdate, so a membership change arriving in the gap could arm
a stabilization timer while the rebroadcast emitted an intermediate ring.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* mount: acquire path locks before mutating create/rename state

Create took the DLM lock only after the filer create, so a lock failure
returned EAGAIN with an eagerly persisted file left behind. Rename marked
source handles renamed before acquiring locks, so a failed acquisition
left them suppressing old-path flushes for a rename that never happened.
Both now take the locks first; the create's lock is released again if the
entry race loses to another creator and AcquireHandle takes over.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* mount: keep the old-path lock when rename lock migration fails

The migration stopped the handle's lock before acquiring the replacement,
so a nil result left the handle writing with no lock at all. Acquiring the
new-path lock first means failure keeps the existing lock instead of
reporting success with serialization dropped.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* mount: skip new-path rename lock when a handle already holds it

A target file open for write on this mount already carries a lock on
newPath; the lock manager does not grant a second lock to the same
owner, so the rename would wait on itself until the handle closed.
Also avoid locking twice when old and new paths coincide.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* mount: hand the rename's target lock to the migrating handle

The rename holds a lock on newPath for its duration, so the response
migration's fresh acquisition waited on that same lock until the handle
released — under fhLockTable, blocking the handle's own close. Adopt the
rename's lock directly; nested move responses still acquire their own.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* mount: move the replaced target's lock to the renamed handle

When the target path was already locked by an open handle on this
mount, the migrated source handle kept only its stale old-path lock —
the target's close would then release the last lock on the new path
while the renamed handle was still open. Adopt the replaced handle's
lock instead.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* mount: stop the handle lock inside the fh lock on release

ReleaseHandle stopped fh.dlmLock before taking the fhLockTable slot, so
a rename migration holding that slot could still observe and adopt a
lock that was already stopping. Stopping under the fh lock makes the
transfer serialize against the release.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* mount: claim the replaced target's lock for the renamed handle

When the target path is already locked by an open handle on this mount,
adopting it at migration time keeps the renamed path protected after
that handle closes, without waiting on a lock this mount already holds.
If the handle was released mid-migration the claimed lock is stopped,
and a fresh acquire covers the case where it was already gone.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* mount: claim the target handle's lock before the rename runs

Skipping the new-path lock when a handle already holds it let that
handle's close release the lock mid-rename, leaving the path unguarded
until the response migrated it. Take over the lock at check time and
hold it for the rename's duration: the response adopts it for the
migrating handle, or it returns to the target handle / is released on
failure. The target handle lookup also falls back to the entry's stored
inode for a forgotten path mapping.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* mount: read handle locks only under the fh lock during rename

The loose dlmLock reads raced ReleaseHandle, which now mutates the lock
inside the handle lock; check and claim it under the same hold.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
Chris Lu
2026-09-26 11:57:00 +08:00
committed by GitHub
co-authored by Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
parent 4914c14982
commit f31a026b2a
17 changed files with 691 additions and 106 deletions
+68 -12
View File
@@ -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 {
+134
View File
@@ -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
+22
View File
@@ -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()
@@ -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()))
}
+81 -27
View File
@@ -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).
+105
View File
@@ -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)
+4
View File
@@ -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
+3 -3
View File
@@ -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()
}
+30 -10
View File
@@ -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)
+7
View File
@@ -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
+126 -48
View File
@@ -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)
}
+1 -1
View File
@@ -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)
}
+5 -1
View File
@@ -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
+9 -3
View File
@@ -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
}
}
+3
View File
@@ -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
+32 -1
View File
@@ -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
+22
View File
@@ -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) {