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,