From e59fb854ee2fb769819976540e4fc79b60c3e17f Mon Sep 17 00:00:00 2001 From: Ettore Di Giacinto Date: Sat, 26 Sep 2026 22:31:39 +0000 Subject: [PATCH] fix(failover): free the leader lock soon after the leader's host dies A leader whose host died without closing its connection kept the advisory lock for about two hours of OS keepalive defaults, and no other frontend could probe. The lock session now sets short TCP keepalives and tcp_user_timeout, so the server drops it within about 30 seconds. Shutdown now closes the lock for good, so a tick that runs after it cannot take the lock back. Assisted-by: Claude:claude-opus-5-5 Signed-off-by: Ettore Di Giacinto --- core/application/failover_distributed.go | 4 +- core/services/advisorylock/held_lock.go | 118 ++++++++++++++----- core/services/advisorylock/held_lock_test.go | 78 ++++++++++++ docs/content/features/model-failover.md | 6 +- 4 files changed, 175 insertions(+), 31 deletions(-) diff --git a/core/application/failover_distributed.go b/core/application/failover_distributed.go index 37eb94edc..87adbdeb5 100644 --- a/core/application/failover_distributed.go +++ b/core/application/failover_distributed.go @@ -87,6 +87,8 @@ func (a *Application) stopFailoverDistributed() { } } if a.failoverLock != nil { - a.failoverLock.Release() + // Close, not Release: Run may still be ticking and would take the + // lock straight back. + a.failoverLock.Close() } } diff --git a/core/services/advisorylock/held_lock.go b/core/services/advisorylock/held_lock.go index 92e8acbe8..46f2db6f5 100644 --- a/core/services/advisorylock/held_lock.go +++ b/core/services/advisorylock/held_lock.go @@ -16,22 +16,41 @@ import ( // database cannot stall the caller's loop. const heldLockCheckTimeout = 5 * time.Second +// heldLockSessionSettings make the server notice a lock holder whose host +// died without closing the connection (crash, power loss, partition) within +// about 30 seconds, instead of after the OS keepalive default of over two +// hours, during which nobody else could take the lock. +var heldLockSessionSettings = []string{ + "SET tcp_keepalives_idle = 10", + "SET tcp_keepalives_interval = 5", + "SET tcp_keepalives_count = 3", +} + +// heldLockUserTimeout bounds how long unacknowledged data (for example a +// keepalive reply to a dead host) may stay in flight. PostgreSQL 12 added it, +// so it is applied separately and an older server's error is ignored. +const heldLockUserTimeout = "SET tcp_user_timeout = 30000" + // HeldLock is an advisory lock that stays taken across calls, for leader // election where leadership must be sticky. TryWithLockCtx releases the lock // when fn returns; with a short fn, every contender wins it in turn and // leadership flips on each tick. // // On PostgreSQL the lock belongs to one session, so HeldLock keeps a -// dedicated connection out of the pool while it holds the lock. If that -// session dies, the server drops the lock and another instance can take it. -// On other dialects it holds the package's in-process lock for the key. +// dedicated connection out of the pool. The same session is reused for +// later attempts while the lock is taken elsewhere, so a follower does not +// open a connection per try. If the holding session dies, the server drops +// the lock and another instance can take it. On other dialects it holds the +// package's in-process lock for the key. type HeldLock struct { db *gorm.DB key int64 - mu sync.Mutex - conn *sql.Conn // PostgreSQL: the session that holds the lock - local bool // other dialects: this lock holds the in-process slot + mu sync.Mutex + conn *sql.Conn // PostgreSQL: the dedicated session, holding the lock or not + held bool // PostgreSQL: conn holds the lock + local bool // other dialects: this lock holds the in-process slot + closed bool // Close was called; the lock is never taken again } // NewHeldLock returns a lock for key on db. It takes nothing until TryAcquire. @@ -44,7 +63,10 @@ func NewHeldLock(db *gorm.DB, key int64) *HeldLock { func (l *HeldLock) TryAcquire(ctx context.Context) (bool, error) { l.mu.Lock() defer l.mu.Unlock() - if l.conn != nil || l.local { + if l.closed { + return false, nil + } + if l.held || l.local { return true, nil } @@ -58,28 +80,48 @@ func (l *HeldLock) TryAcquire(ctx context.Context) (bool, error) { } } - sqlDB, err := l.db.DB() - if err != nil { - return false, fmt.Errorf("get sql.DB: %w", err) - } - conn, err := sqlDB.Conn(ctx) - if err != nil { - return false, fmt.Errorf("advisory lock conn: %w", err) + if l.conn == nil { + conn, err := l.openSession(ctx) + if err != nil { + return false, err + } + l.conn = conn } var acquired bool - if err := conn.QueryRowContext(ctx, "SELECT pg_try_advisory_lock($1)", l.key).Scan(&acquired); err != nil { + if err := l.conn.QueryRowContext(ctx, "SELECT pg_try_advisory_lock($1)", l.key).Scan(&acquired); err != nil { // The lock may have been granted before the error (a cancelled // context, say); discarding the session is the only way to be sure // it is not left held. - discardConn(conn) + discardConn(l.conn) + l.conn = nil return false, fmt.Errorf("pg_try_advisory_lock: %w", err) } - if !acquired { - _ = conn.Close() // holds nothing, so it may go back to the pool - return false, nil + l.held = acquired + return acquired, nil +} + +// openSession takes a connection out of the pool for this lock and applies +// the keepalive settings to it. The session never goes back to the pool, so +// other pool users do not inherit those settings. +func (l *HeldLock) openSession(ctx context.Context) (*sql.Conn, error) { + sqlDB, err := l.db.DB() + if err != nil { + return nil, fmt.Errorf("get sql.DB: %w", err) } - l.conn = conn - return true, nil + conn, err := sqlDB.Conn(ctx) + if err != nil { + return nil, fmt.Errorf("advisory lock conn: %w", err) + } + for _, stmt := range heldLockSessionSettings { + if _, err := conn.ExecContext(ctx, stmt); err != nil { + discardConn(conn) + return nil, fmt.Errorf("advisory lock session %q: %w", stmt, err) + } + } + if _, err := conn.ExecContext(ctx, heldLockUserTimeout); err != nil { + xlog.Debug("advisory lock session: tcp_user_timeout not supported, relying on keepalives", "error", err) + } + return conn, nil } // Held reports whether this HeldLock believes it holds the lock. It does not @@ -87,7 +129,7 @@ func (l *HeldLock) TryAcquire(ctx context.Context) (bool, error) { func (l *HeldLock) Held() bool { l.mu.Lock() defer l.mu.Unlock() - return l.conn != nil || l.local + return l.held || l.local } // Verify reports whether the lock is still held. On PostgreSQL it checks @@ -99,7 +141,7 @@ func (l *HeldLock) Verify(ctx context.Context) bool { if l.local { return true } - if l.conn == nil { + if !l.held { return false } pctx, cancel := context.WithTimeout(ctx, heldLockCheckTimeout) @@ -110,15 +152,32 @@ func (l *HeldLock) Verify(ctx context.Context) bool { // alive; discarding closes it, so the server releases the lock. discardConn(l.conn) l.conn = nil + l.held = false return false } return true } -// Release gives up the lock. It is safe to call when the lock is not held. +// Release gives up the lock; a later TryAcquire may take it again. It is safe +// to call when the lock is not held. func (l *HeldLock) Release() { l.mu.Lock() defer l.mu.Unlock() + l.releaseLocked() +} + +// Close gives up the lock for good: TryAcquire returns false afterwards. Use +// it on shutdown, where a loop still running could otherwise take the lock +// back right after Release (and, on the in-process fallback, keep it for the +// rest of the process). +func (l *HeldLock) Close() { + l.mu.Lock() + defer l.mu.Unlock() + l.closed = true + l.releaseLocked() +} + +func (l *HeldLock) releaseLocked() { if l.local { <-localLockChan(l.key) l.local = false @@ -127,15 +186,18 @@ func (l *HeldLock) Release() { if l.conn == nil { return } - ctx, cancel := context.WithTimeout(context.Background(), heldLockCheckTimeout) - defer cancel() - if _, err := l.conn.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", l.key); err != nil { - xlog.Warn("advisory lock unlock failed, closing its session instead", "key", l.key, "error", err) + if l.held { + ctx, cancel := context.WithTimeout(context.Background(), heldLockCheckTimeout) + defer cancel() + if _, err := l.conn.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", l.key); err != nil { + xlog.Warn("advisory lock unlock failed, closing its session instead", "key", l.key, "error", err) + } } // Never hand the session back to the pool: if the unlock failed it still // holds the lock, and a later pool user would silently inherit it. discardConn(l.conn) l.conn = nil + l.held = false } // discardConn closes conn's underlying session instead of returning it to diff --git a/core/services/advisorylock/held_lock_test.go b/core/services/advisorylock/held_lock_test.go index fa8322d10..17f2c17b9 100644 --- a/core/services/advisorylock/held_lock_test.go +++ b/core/services/advisorylock/held_lock_test.go @@ -11,6 +11,29 @@ import ( "gorm.io/gorm" ) +// expectCloseIsFinal checks that Close frees the lock for a rival and that the +// closed HeldLock never takes it back, as a still-running loop would try to. +func expectCloseIsFinal(db *gorm.DB, key int64) { + ctx := context.Background() + l, rival := NewHeldLock(db, key), NewHeldLock(db, key) + DeferCleanup(rival.Release) + + ok, err := l.TryAcquire(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(ok).To(BeTrue()) + l.Close() + Expect(l.Held()).To(BeFalse()) + + ok, err = l.TryAcquire(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(ok).To(BeFalse(), "a closed lock is never taken again") + + ok, err = rival.TryAcquire(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(ok).To(BeTrue(), "Close frees the lock for others") + l.Close() // idempotent +} + // expectStickyHandover checks the contract both backends share: the first // holder keeps the lock across repeated checks while a rival is denied, and // the rival gets it once the holder releases. @@ -50,6 +73,12 @@ var _ = Describe("HeldLock (SQLite fallback)", Label("sqlite"), func() { expectStickyHandover(db, 12101) }) + It("never takes the lock again after Close", func() { + db, err := gorm.Open(sqlite.Open("file::memory:?cache=shared"), &gorm.Config{}) + Expect(err).ToNot(HaveOccurred()) + expectCloseIsFinal(db, 12103) + }) + It("is idempotent: acquiring twice and releasing twice is safe", func() { db, err := gorm.Open(sqlite.Open("file::memory:?cache=shared"), &gorm.Config{}) Expect(err).ToNot(HaveOccurred()) @@ -80,6 +109,55 @@ var _ = Describe("HeldLock (PostgreSQL)", func() { expectStickyHandover(db, 12201) }) + It("never takes the lock again after Close", func() { + expectCloseIsFinal(db, 12203) + }) + + It("sets short TCP keepalives on its session so a dead host's lock expires", func() { + // Killing a host without closing its socket is not practical in a + // test; check the settings that make the server notice one. + ctx := context.Background() + l := NewHeldLock(db, 12204) + DeferCleanup(l.Release) + ok, err := l.TryAcquire(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(ok).To(BeTrue()) + + show := func(name string) string { + var v string + Expect(l.conn.QueryRowContext(ctx, "SHOW "+name).Scan(&v)).To(Succeed()) + return v + } + Expect(show("tcp_keepalives_idle")).To(Equal("10")) + Expect(show("tcp_keepalives_interval")).To(Equal("5")) + Expect(show("tcp_keepalives_count")).To(Equal("3")) + Expect(show("tcp_user_timeout")).To(Equal("30000")) // milliseconds + }) + + It("reuses one session while the lock is taken elsewhere", func() { + ctx := context.Background() + holder, follower := NewHeldLock(db, 12205), NewHeldLock(db, 12205) + DeferCleanup(holder.Release) + DeferCleanup(follower.Release) + ok, err := holder.TryAcquire(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(ok).To(BeTrue()) + + pid := func() int { + var p int + Expect(follower.conn.QueryRowContext(ctx, "SELECT pg_backend_pid()").Scan(&p)).To(Succeed()) + return p + } + ok, err = follower.TryAcquire(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(ok).To(BeFalse()) + before := pid() + ok, err = follower.TryAcquire(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(ok).To(BeFalse()) + Expect(pid()).To(Equal(before), "a follower must not open a connection per attempt") + }) + It("drops leadership when the holding session dies, freeing the lock", func() { ctx := context.Background() first, second := NewHeldLock(db, 12202), NewHeldLock(db, 12202) diff --git a/docs/content/features/model-failover.md b/docs/content/features/model-failover.md index f4f42027b..32e91af96 100644 --- a/docs/content/features/model-failover.md +++ b/docs/content/features/model-failover.md @@ -187,8 +187,10 @@ share one failover state: all frontends converge on the same target for a chain. - One frontend, the probe leader, runs the health checks, decides fail-over and fail-back, and loads warm targets. The leader holds a PostgreSQL - advisory lock and keeps it until it shuts down or its database connection - fails. Then another frontend takes the lock and becomes the leader. + advisory lock and keeps it until it stops or its database connection fails. + Then another frontend takes the lock and becomes the leader: immediately + when the leader shuts down or its process exits, and within about 30 seconds + when the leader's host or network fails. - Warm targets stay loaded on the workers. The router and the replica reconciler treat them like pinned models and do not evict them. - A frontend that starts late gets the current state within 10 seconds,