From 62476e553ea9acd79a3e2e8414028545ba4a97d8 Mon Sep 17 00:00:00 2001 From: Ettore Di Giacinto Date: Tue, 1 Sep 2026 09:06:42 +0000 Subject: [PATCH] fix(cluster): keep a session close out of the per-node claim gate The gate is justified by being held for one claim round trip, and Attach held it across the close of the session it superseded. Closing a yamux session closes the underlying conn and then waits for both its send and recv loops to exit, and the send loop can be inside a write bounded only by ConnectionWriteTimeout, so that is a wait on other goroutines. It must not stand between a worker re-dialling this node and its claim. The gate is now released after the store and before the close, which also makes Attach match reclaimOne, where it has always been released explicitly on every path. This is safe because a superseded session is no longer reachable from the map by the time it is closed: the next re-dial replaces an entry that already names the new session. Pin the re-claim half of the gate too. A worker that re-dials between a re-claim's commit and its record leaves the row carrying the re-dial's epoch while the entry carries the re-claim's, so the attachment holding the socket releases an epoch the row does not have and the row outlives it, with nothing to sweep it while this replica is alive. Only Attach's half of the serialisation was asserted; keying the two apart left every spec green. Also take the test hook's action under the lock that guards whether it has fired. It was written from the spec's goroutine and read from whichever goroutine issued the statement, which is a race in the harness that pins the serialisation specs. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto --- core/services/cluster/tunnel.go | 15 +++++- core/services/cluster/tunnel_test.go | 79 +++++++++++++++++++++++++--- 2 files changed, 86 insertions(+), 8 deletions(-) diff --git a/core/services/cluster/tunnel.go b/core/services/cluster/tunnel.go index 891bf77b8..36f7c91d3 100644 --- a/core/services/cluster/tunnel.go +++ b/core/services/cluster/tunnel.go @@ -170,10 +170,10 @@ func (t *TunnelRegistry) Attach(ctx context.Context, nodeID string, sess *yamux. if err := t.enterClaim(ctx, nodeID); err != nil { return 0, fmt.Errorf("attaching tunnel for node %q: %w", nodeID, err) } - defer t.leaveClaim(nodeID) epoch, err := t.reg.Claim(ctx, nodeID, t.selfID) if err != nil { + t.leaveClaim(nodeID) return 0, err } @@ -181,7 +181,20 @@ func (t *TunnelRegistry) Attach(ctx context.Context, nodeID string, sess *yamux. previous := t.tunnels[nodeID] t.tunnels[nodeID] = &heldTunnel{sess: sess, token: epoch, claim: epoch} t.mu.Unlock() + t.leaveClaim(nodeID) + // Closed after the gate is released, not under it. The gate is justified by + // being held for one claim round trip, and closing a session is not that: + // yamux closes the underlying conn and then waits for both its send and + // recv loops to exit (go-yamux/v5@v5.1.0/session.go:330-332), and the send + // loop can be inside a write bounded only by ConnectionWriteTimeout. That + // is a wait on other goroutines, and it must not stand between a worker + // re-dialling this node and its claim. + // + // Releasing first is safe because the superseded session is no longer + // reachable from the map: whoever re-dials next replaces an entry that + // already names the new session, and this close can only ever affect the + // one it just displaced. if previous != nil && previous.sess != sess { xlog.Debug("worker re-dialled this replica, dropping its previous tunnel", "node", nodeID) _ = previous.sess.Close() diff --git a/core/services/cluster/tunnel_test.go b/core/services/cluster/tunnel_test.go index ac7ff4ad8..67993130e 100644 --- a/core/services/cluster/tunnel_test.go +++ b/core/services/cluster/tunnel_test.go @@ -43,16 +43,26 @@ func newClaimHook(match func(sql string) bool) *claimHook { return &claimHook{Interface: gormlogger.Default.LogMode(gormlogger.Silent), match: match} } +// setAction installs what the hook runs, under the same lock that guards fired. +// The action is written from the spec's goroutine and read from whichever +// goroutine issues the statement, which need not be the same one. +func (h *claimHook) setAction(action func(sql string)) { + h.mu.Lock() + defer h.mu.Unlock() + h.action = action +} + func (h *claimHook) Trace(ctx context.Context, begin time.Time, fc func() (string, int64), err error) { sql, rows := fc() h.mu.Lock() - fire := !h.fired && h.action != nil && h.match(sql) + action := h.action + fire := !h.fired && action != nil && h.match(sql) if fire { h.fired = true } h.mu.Unlock() if fire { - h.action(sql) + action(sql) } h.Interface.Trace(ctx, begin, func() (string, int64) { return sql, rows }, err) } @@ -334,7 +344,7 @@ var _ = Describe("The worker tunnel registry", func() { secondSession, _ := workerTunnel() secondStarted := make(chan struct{}) secondEpochs := make(chan int64, 1) - hook.action = func(string) { + hook.setAction(func(string) { // Launched from inside the first claim, so the second Attach is // provably reaching for the same node while the first is between // its claim and its store. Starting it before the call would leave @@ -349,7 +359,7 @@ var _ = Describe("The worker tunnel registry", func() { <-secondStarted Consistently(secondEpochs, serializationProbe, 10*time.Millisecond).ShouldNot(Receive(), "a second Attach for this node claimed AND recorded its epoch while the first was between its own claim and store") - } + }) firstSession, _ := workerTunnel() firstEpoch, err := hooked.Attach(ctx, "w1", firstSession) @@ -575,7 +585,7 @@ var _ = Describe("Re-claiming tunnels after this replica's rows were reaped", fu epoch int64 } redialled := make(chan redial, 1) - hook.action = func(sql string) { + hook.setAction(func(sql string) { node := "w2" if claimedNode(sql, "w1", "w2") == "w2" { node = "w1" @@ -584,7 +594,7 @@ var _ = Describe("Re-claiming tunnels after this replica's rows were reaped", fu epoch, err := hooked.Attach(ctx, node, frontend) Expect(err).ToNot(HaveOccurred()) redialled <- redial{node: node, epoch: epoch} - } + }) count, err := hooked.Reclaim(ctx) Expect(err).ToNot(HaveOccurred()) @@ -604,6 +614,61 @@ var _ = Describe("Re-claiming tunnels after this replica's rows were reaped", fu "the re-claim was recorded on nothing, so the attachment that holds the socket cannot release its row") }) + It("serialises a re-claim against an Attach for the same node", func() { + // The re-claim takes the same gate Attach does, and for the same + // reason. Unserialised, a worker re-dialling between the re-claim's + // commit and its record leaves the row carrying the re-dial's epoch + // while the entry carries the re-claim's, so the attachment holding the + // socket releases an epoch the row does not have. Nothing sweeps the + // row that is left, because this replica is alive: Owner keeps naming + // it as the owner of a tunnel it does not hold. + hook := newClaimHook(isClaimOf("w1")) + hooked := cluster.NewTunnelRegistry( + cluster.NewRegistry(db.Session(&gorm.Session{Logger: hook})), "me") + + first, _ := workerTunnel() + attachEpoch, err := hooked.Attach(ctx, "w1", first) + Expect(err).ToNot(HaveOccurred()) + + redialSession, _ := workerTunnel() + redialStarted := make(chan struct{}) + redialEpochs := make(chan int64, 1) + hook.setAction(func(string) { + // Launched from inside the re-claim's own claim, so the re-dial is + // provably reaching for this node while the re-claim is between its + // commit and its record. + go func() { + defer GinkgoRecover() + close(redialStarted) + epoch, err := hooked.Attach(ctx, "w1", redialSession) + Expect(err).ToNot(HaveOccurred()) + redialEpochs <- epoch + }() + <-redialStarted + Consistently(redialEpochs, serializationProbe, 10*time.Millisecond).ShouldNot(Receive(), + "a worker re-dialled and recorded its claim while a re-claim for the same node was between its own claim and record") + }) + + count, err := hooked.Reclaim(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(count).To(Equal(1)) + + var redialEpoch int64 + Eventually(redialEpochs, "10s").Should(Receive(&redialEpoch)) + Expect(redialEpoch).ToNot(Equal(attachEpoch)) + + // The re-dial claimed after the re-claim, so the row carries its epoch + // and its attachment is the one that has to be able to release it. The + // superseded token must still be a no-op. + hooked.Detach("w1", attachEpoch) + Expect(hooked.Held()).To(ConsistOf("w1")) + hooked.Detach("w1", redialEpoch) + Expect(hooked.Held()).To(BeEmpty()) + _, _, err = reg.OwnerRow(ctx, "w1") + Expect(err).To(MatchError(cluster.ErrNoConnection), + "the attachment that re-dialled could not release its row, so the row outlived the socket") + }) + It("releases a re-claim whose attachment detached while the claim was in flight", func() { // Detach is not gated against a re-claim, so this interleave is real: // the claim commits, then the socket dies and Detach releases the epoch @@ -619,7 +684,7 @@ var _ = Describe("Re-claiming tunnels after this replica's rows were reaped", fu epoch, err := hooked.Attach(ctx, "w1", frontend) Expect(err).ToNot(HaveOccurred()) - hook.action = func(string) { hooked.Detach("w1", epoch) } + hook.setAction(func(string) { hooked.Detach("w1", epoch) }) count, err := hooked.Reclaim(ctx) Expect(err).ToNot(HaveOccurred())