From 438f80e885df75c142e3476dcfe90c7f8aa2dc9e Mon Sep 17 00:00:00 2001 From: Joost Jager Date: Wed, 23 Sep 2026 11:15:26 +0200 Subject: [PATCH 1/2] Make PostgreSQL schema setup atomic Create the KV table, record its schema version, and create the listing index in one transaction so failed initialization rolls back schema changes. Keep database creation outside the transaction and preserve the session advisory lock for the store's lifetime. Add a regression test that forces index creation to fail and verifies the schema version update is rolled back. --- src/io/postgres_store/mod.rs | 64 +++++++++++++++++++++++++++++++----- 1 file changed, 56 insertions(+), 8 deletions(-) diff --git a/src/io/postgres_store/mod.rs b/src/io/postgres_store/mod.rs index 4ffb948ff9..ace93f6c83 100644 --- a/src/io/postgres_store/mod.rs +++ b/src/io/postgres_store/mod.rs @@ -443,7 +443,7 @@ impl PostgresStoreInner { Self::create_database_if_not_exists(&config, &tls, logger.as_deref()).await?; - let client = make_config_connection(&config, &tls).await?; + let mut client = make_config_connection(&config, &tls).await?; let lock_id = advisory_lock_id(&db_name, &kv_table_name); let row = client.query_one("SELECT pg_try_advisory_lock($1)", &[&lock_id]).await.map_err( |e| { @@ -462,6 +462,14 @@ impl PostgresStoreInner { )); } + let pool = SmallPool::new(&config, &tls).await?; + let transaction = client.transaction().await.map_err(|e| { + io::Error::new( + io::ErrorKind::Other, + format!("Failed to start PostgreSQL schema setup transaction: {e}"), + ) + })?; + // Create the KV data table if it doesn't exist. `sort_order` uses BIGSERIAL so // the database assigns a fresh, monotonically increasing value on each INSERT and // keeps the previous value untouched on UPSERT-update; the sequence persists across @@ -476,13 +484,13 @@ impl PostgresStoreInner { PRIMARY KEY (primary_namespace, secondary_namespace, key) )" ); - client.execute(sql.as_str(), &[]).await.map_err(|e| { + transaction.execute(sql.as_str(), &[]).await.map_err(|e| { let msg = format!("Failed to create table {kv_table_name}: {e}"); io::Error::new(io::ErrorKind::Other, msg) })?; // Read the schema version from the table comment (analogous to SQLite's PRAGMA user_version). - let row = client + let row = transaction .query_one("SELECT obj_description(to_regclass($1), 'pg_class')", &[&kv_table_name_sql]) .await .map_err(|e| { @@ -513,13 +521,18 @@ impl PostgresStoreInner { if version_res == 0 { // New table, set our SCHEMA_VERSION. let sql = format!("COMMENT ON TABLE {kv_table_name_sql} IS '{SCHEMA_VERSION}'"); - client.execute(sql.as_str(), &[]).await.map_err(|e| { + transaction.execute(sql.as_str(), &[]).await.map_err(|e| { let msg = format!("Failed to set schema version: {e}"); io::Error::new(io::ErrorKind::Other, msg) })?; } else if version_res < SCHEMA_VERSION { - migrations::migrate_schema(&client, &kv_table_name_sql, version_res, SCHEMA_VERSION) - .await?; + migrations::migrate_schema( + transaction.client(), + &kv_table_name_sql, + version_res, + SCHEMA_VERSION, + ) + .await?; } else if version_res > SCHEMA_VERSION { let msg = format!( "Failed to open database: incompatible schema version {version_res}. Expected: {SCHEMA_VERSION}" @@ -532,12 +545,17 @@ impl PostgresStoreInner { let sql = format!( "CREATE INDEX IF NOT EXISTS {index_name_sql} ON {kv_table_name_sql} (primary_namespace, secondary_namespace, sort_order DESC, key ASC)" ); - client.execute(sql.as_str(), &[]).await.map_err(|e| { + transaction.execute(sql.as_str(), &[]).await.map_err(|e| { let msg = format!("Failed to create index on table {kv_table_name}: {e}"); io::Error::new(io::ErrorKind::Other, msg) })?; - let pool = SmallPool::new(&config, &tls).await?; + transaction.commit().await.map_err(|e| { + io::Error::new( + io::ErrorKind::Other, + format!("Failed to commit PostgreSQL schema setup transaction: {e}"), + ) + })?; let write_version_locks = Mutex::new(HashMap::new()); Ok(Self { @@ -1002,6 +1020,36 @@ mod tests { cleanup_store(&store).await; } + #[tokio::test(flavor = "multi_thread")] + async fn test_postgres_store_rolls_back_failed_schema_setup() { + let table_name = "test_pg_schema_setup_rollback"; + let store = create_test_store(table_name).await; + let kv_table = store.inner.kv_table_name_sql.clone(); + let client = make_config_connection(&store.inner.config, &store.inner.tls).await.unwrap(); + drop(store); + + // Force index creation to fail after setup writes the schema version. + client + .batch_execute(&format!( + "COMMENT ON TABLE {kv_table} IS NULL; ALTER TABLE {kv_table} DROP COLUMN sort_order" + )) + .await + .unwrap(); + let err = + PostgresStore::new(test_connection_string(), None, Some(table_name.to_string()), None) + .await + .err() + .expect("schema setup must fail without sort_order"); + assert!(err.to_string().contains("Failed to create index")); + + let row = client + .query_one("SELECT obj_description(to_regclass($1), 'pg_class')", &[&kv_table]) + .await + .unwrap(); + client.execute(&format!("DROP TABLE {kv_table}"), &[]).await.unwrap(); + assert_eq!(row.get::<_, Option<&str>>(0), None); + } + #[tokio::test(flavor = "multi_thread")] async fn read_write_remove_list_persist() { let store = create_test_store("test_rwrl").await; From ce1a0f285205f80af74d1a16f8c0a4111a54aa84 Mon Sep 17 00:00:00 2001 From: Joost Jager Date: Wed, 23 Sep 2026 13:47:49 +0200 Subject: [PATCH 2/2] Replace PostgreSQL advisory locks with renewable leases Replace the temporary session advisory lock with an expiring lease in a companion table. Acquire the lease atomically with schema setup before loading persisted state. Fence each write and delete with a single statement that validates lease ownership and expiry before applying the mutation. Use FOR KEY SHARE and a unique owner_id so mutations and background renewal can run concurrently while takeover waits for the statement to commit. Retry only preparation to avoid replaying a mutation whose commit outcome is unknown. Renew leases only in the background and skip missed ticks. Keep renewal and the store runtime alive through outstanding mutation futures and their spawned tasks. Release by owner ID after the final shared resources reference is dropped, with bounded shutdown that allows Tokio workers to continue driving lease release. Panic on lease rejection and renewal failure or timeout, and document the application panic-abort requirement. Cover setup rollback, renewal and takeover lock compatibility, stale writes and deletes, missing-key removal after expiry, renewal timeouts, shutdown with a blocked pool, and mutation futures outliving the store. Verify immediate lease acquisition after those mutations finish and clean up both KV and lease tables after tests. --- src/builder.rs | 20 +- src/io/postgres_store/mod.rs | 608 +++++++++++++++++++++++++++-------- tests/common/mod.rs | 9 +- 3 files changed, 493 insertions(+), 144 deletions(-) diff --git a/src/builder.rs b/src/builder.rs index 1158044e47..29759e93af 100644 --- a/src/builder.rs +++ b/src/builder.rs @@ -737,11 +737,11 @@ impl NodeBuilder { /// /// # Warning /// - /// Do not point multiple [`Node`] instances at the same database and table. Concurrent access is - /// unsafe and can corrupt node state. You must make sure that only one node accesses each - /// database and table. The store uses a PostgreSQL advisory lock to reduce this risk. This lock - /// is only a temporary safeguard and does not make concurrent access safe. - /// Nodes using a different database or table on the same server may coexist. + /// This acquires an exclusive lease for the selected KV table before reading persisted node + /// state. Nodes may share a database when each node identity uses a distinct `kv_table_name`. + /// Detected lease loss or failed or timed-out background renewals panic. Applications must set + /// `panic = "abort"` in their own Cargo profiles so these panics terminate the process. Recovery + /// requires restarting the process and constructing a fresh node from persisted state. /// /// If `certificate_pem` is `Some`, TLS will be used for database connections and the /// provided PEM-encoded CA certificate will be added to the system's default root @@ -1334,11 +1334,11 @@ impl Builder { /// /// # Warning /// - /// Do not point multiple [`Node`] instances at the same database and table. Concurrent access is - /// unsafe and can corrupt node state. You must make sure that only one node accesses each - /// database and table. The store uses a PostgreSQL advisory lock to reduce this risk. This lock - /// is only a temporary safeguard and does not make concurrent access safe. - /// Nodes using a different database or table on the same server may coexist. + /// This acquires an exclusive lease for the selected KV table before reading persisted node + /// state. Nodes may share a database when each node identity uses a distinct `kv_table_name`. + /// Detected lease loss or failed or timed-out background renewals panic. Applications must set + /// `panic = "abort"` in their own Cargo profiles so these panics terminate the process. Recovery + /// requires restarting the process and constructing a fresh node from persisted state. /// /// If `certificate_pem` is `Some`, TLS will be used for database connections and the /// provided PEM-encoded CA certificate will be added to the system's default root diff --git a/src/io/postgres_store/mod.rs b/src/io/postgres_store/mod.rs index ace93f6c83..241df19cf8 100644 --- a/src/io/postgres_store/mod.rs +++ b/src/io/postgres_store/mod.rs @@ -10,8 +10,8 @@ use std::collections::HashMap; use std::future::Future; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; +use std::time::Duration; -use bitcoin::hashes::{sha256, Hash, HashEngine}; use lightning::io; use lightning::util::persist::{ KVStore, MigratableKVStore, PageToken, PaginatedKVStore, PaginatedListResponse, @@ -24,7 +24,7 @@ use tokio_postgres::{Config, Error as PgError}; use self::pool::{make_config_connection, ClientConnection, PgTlsConnector, SmallPool}; use crate::io::utils::check_namespace_key_validity; -use crate::logger::{log_debug, log_info, LdkLogger, Logger}; +use crate::logger::{log_debug, log_error, log_info, LdkLogger, Logger}; use crate::runtime::StoreRuntime; mod migrations; @@ -45,17 +45,15 @@ const PAGE_SIZE: usize = 50; // Keep this small while still allowing progress if one runtime worker blocks on sync store access. const INTERNAL_RUNTIME_WORKERS: usize = 2; -fn advisory_lock_id(db_name: &str, kv_table_name: &str) -> i64 { - let mut engine = sha256::Hash::engine(); - engine.input(b"ldk-node:postgres-store"); - for component in [db_name, kv_table_name] { - engine.input(&(component.len() as u64).to_be_bytes()); - engine.input(component.as_bytes()); - } +const NODE_LEASE_DURATION: Duration = Duration::from_secs(30); +const NODE_LEASE_RENEWAL_INTERVAL: Duration = Duration::from_secs(10); +const NODE_LEASE_RENEWAL_TIMEOUT: Duration = Duration::from_secs(10); +const NODE_LEASE_RELEASE_TIMEOUT: Duration = Duration::from_secs(5); - let hash = sha256::Hash::from_engine(engine).to_byte_array(); - i64::from_be_bytes(hash[..8].try_into().expect("SHA-256 prefix has the expected length")) -} +const NODE_LEASE_TABLE_SUFFIX: &str = "_node_lease"; +const POSTGRES_IDENTIFIER_MAX_BYTES: usize = 63; +const MAX_KV_TABLE_NAME_BYTES: usize = + POSTGRES_IDENTIFIER_MAX_BYTES - NODE_LEASE_TABLE_SUFFIX.len(); fn sql_identifier(identifier: &str) -> io::Result { if identifier.is_empty() || identifier.contains('\0') { @@ -83,10 +81,24 @@ fn sql_table_identifier(table_name: &str) -> io::Result { Ok(quoted_parts?.join(".")) } -/// Runs a tokio-postgres query and, if the pooled connection dropped mid-flight, reconnects and -/// retries once after asserting that the store's advisory-lock connection is still open. `$store` -/// is the [`PostgresStoreInner`], `$locked` the held client slot guard, `$err_map` an -/// `Fn(PgError) -> io::Error` (called at most once), and `$query` an expression that yields a fresh +fn sql_node_lease_table_identifier(table_name: &str) -> io::Result { + sql_table_identifier(table_name)?; + let table_part = table_name.rsplit_once('.').map_or(table_name, |(_, table)| table); + if table_part.len() > MAX_KV_TABLE_NAME_BYTES { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!( + "PostgreSQL KV table name exceeds the maximum of {MAX_KV_TABLE_NAME_BYTES} bytes: {table_name}" + ), + )); + } + + sql_table_identifier(&format!("{table_name}{NODE_LEASE_TABLE_SUFFIX}")) +} + +/// Runs a standalone tokio-postgres query and, if the pooled connection dropped mid-flight, +/// reconnects and retries once. `$store` is the [`PostgresStoreInner`], `$locked` the held client +/// slot guard, `$err_map` an `FnOnce(PgError) -> io::Error`, and `$query` an expression that yields a fresh /// `Future>` each time it is evaluated. `$query` may be evaluated up to /// twice (once normally, once on retry), so it must be side-effect-free outside of issuing the /// query itself. @@ -95,14 +107,10 @@ macro_rules! query_with_retry { match $query.await { Ok(v) => Ok(v), Err(e) if $locked.is_closed() || e.is_closed() => { - $store.assert_store_lock(); if let Some(logger) = $store.logger.as_ref() { log_debug!(logger, "Reconnecting to PostgreSQL after error: {e}"); } *$locked = make_config_connection(&$store.config, &$store.tls).await?; - // Recheck after awaiting the connection in case the store lock was lost while - // reconnecting. - $store.assert_store_lock(); $query.await.map_err($err_map) }, Err(e) => Err($err_map(e)), @@ -126,6 +134,12 @@ fn handle_runtime_task_result( /// A [`KVStore`] implementation that writes to and reads from a [PostgreSQL] database. /// /// Maintains an internal runtime for the underlying tokio-postgres connection drivers. +/// Each instance exclusively leases its configured KV table and checks the lease within each KV +/// mutation transaction. Only the background renewal task extends the lease. +/// Lease rejection panics in the calling task; a failed or timed-out +/// background renewal panics in the renewal task. Applications must set `panic = "abort"` in their +/// own Cargo profiles so these panics terminate the process. Recovery requires restarting the +/// process and constructing a fresh node from persisted state. /// /// [PostgreSQL]: https://www.postgresql.org pub struct PostgresStore { @@ -135,8 +149,19 @@ pub struct PostgresStore { // operations aren't sensitive to the order of execution. next_write_version: AtomicU64, + // Outstanding mutations retain the lease and its renewal task after the store is dropped. + resources: Arc, +} + +struct PostgresStoreResources { + inner: Arc, + // A store-internal runtime that drives PostgreSQL I/O independently from the node runtime. internal_runtime: Option>, + + lease_renewal_task: tokio::task::JoinHandle<()>, + lease_shutdown_sender: Option>, + lease_shutdown_complete: Mutex>>, } // tokio::sync::Mutex (used for the DB client) contains UnsafeCell which opts out of @@ -159,14 +184,10 @@ impl PostgresStore { /// the default `postgres` database to create it. /// /// The given `kv_table_name` will be used or default to [`DEFAULT_KV_TABLE_NAME`]. + /// A companion lease table is created by appending `_node_lease` to this name. /// - /// # Warning - /// - /// Do not point multiple [`PostgresStore`] instances at the same database and table. Concurrent - /// access is unsafe and can corrupt stored data. You must make sure that only one store accesses - /// each database and table. The store uses a PostgreSQL advisory lock to reduce this risk. This - /// lock is only a temporary safeguard and does not make concurrent access safe. - /// Stores using a different database or table on the same PostgreSQL server may coexist. + /// Construction acquires an exclusive lease for the selected KV table. Returns an error with + /// [`io::ErrorKind::AlreadyExists`] while another store holds an unexpired lease. /// /// If `certificate_pem` is `Some`, TLS will be used for database connections and the /// provided PEM-encoded CA certificate will be added to the system's default root @@ -200,8 +221,74 @@ impl PostgresStore { io::Error::new(io::ErrorKind::Other, format!("PostgreSQL runtime task failed: {}", e)) })??; let inner = Arc::new(inner); - let next_write_version = AtomicU64::new(1); - Ok(Self { inner, next_write_version, internal_runtime: Some(internal_runtime) }) + + let inner_ref = Arc::clone(&inner); + let (lease_shutdown_sender, shutdown_rx) = tokio::sync::oneshot::channel(); + let (completion_sender, lease_shutdown_complete) = std::sync::mpsc::channel(); + let lease_renewal_task = internal_runtime.spawn(async move { + let mut interval = tokio::time::interval(NODE_LEASE_RENEWAL_INTERVAL); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + let lease_duration_secs = NODE_LEASE_DURATION.as_secs() as i64; + let lease_table = &inner_ref.node_lease_table_name_sql; + let update_sql = format!( + "UPDATE {lease_table} + SET expires_at = clock_timestamp() + ($2::bigint * interval '1 second') + WHERE id = 1 AND owner_id = $1 AND expires_at > clock_timestamp()" + ); + let renewal_loop = async { + loop { + // The first tick is immediate; later attempts follow the renewal interval. + interval.tick().await; + let renewal = async { + let mut locked = inner_ref.locked_client().await?; + let err_map = |e| { + io::Error::new( + io::ErrorKind::Other, + format!("Failed to renew node lease: {e}"), + ) + }; + query_with_retry!( + inner_ref, + locked, + err_map, + locked.execute( + &update_sql, + &[&inner_ref.lease_owner_id.as_slice(), &lease_duration_secs] + ) + ) + }; + // Bound both the pool wait and query. + match tokio::time::timeout(NODE_LEASE_RENEWAL_TIMEOUT, renewal).await { + Ok(Ok(1)) => {}, + Ok(Ok(_)) => panic!( + "PostgreSQL node lease was lost; continuing may corrupt node state" + ), + Ok(Err(e)) => panic!("Failed to renew PostgreSQL node lease: {e}"), + Err(_) => panic!("PostgreSQL node lease renewal timed out"), + } + } + }; + tokio::select! { + biased; + _ = shutdown_rx => {}, + _ = renewal_loop => {}, + } + + let result = inner_ref.release_node_lease().await; + let _ = completion_sender.send(result); + }); + + Ok(Self { + inner: Arc::clone(&inner), + next_write_version: AtomicU64::new(1), + resources: Arc::new(PostgresStoreResources { + inner, + internal_runtime: Some(internal_runtime), + lease_renewal_task, + lease_shutdown_sender: Some(lease_shutdown_sender), + lease_shutdown_complete: Mutex::new(lease_shutdown_complete), + }), + }) } fn build_tls_connector(certificate_pem: Option) -> io::Result { @@ -246,12 +333,37 @@ impl PostgresStore { } fn internal_runtime(&self) -> Arc { - Arc::clone(self.internal_runtime.as_ref().expect("PostgreSQL runtime must be available")) + Arc::clone( + self.resources.internal_runtime.as_ref().expect("PostgreSQL runtime must be available"), + ) } } -impl Drop for PostgresStore { +impl Drop for PostgresStoreResources { fn drop(&mut self) { + if let Some(sender) = self.lease_shutdown_sender.take() { + let _ = sender.send(()); + } + // The final mutation may release the last resources reference on a store runtime worker. + // Allow lease shutdown to progress while this worker waits for completion. + let completion = self.lease_shutdown_complete.get_mut().unwrap(); + let wait = || completion.recv_timeout(NODE_LEASE_RELEASE_TIMEOUT); + let result = match tokio::runtime::Handle::try_current() { + Ok(handle) if handle.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread => { + tokio::task::block_in_place(wait) + }, + _ => wait(), + }; + if let Some(logger) = self.inner.logger.as_ref() { + match result { + Ok(Ok(())) => {}, + Ok(Err(e)) => log_error!(logger, "Failed to release PostgreSQL node lease: {e}"), + Err(e) => log_error!(logger, "PostgreSQL lease shutdown did not complete: {e}"), + } + } + // Cancel any remaining work if shutdown failed or timed out. + self.lease_renewal_task.abort(); + if let Some(internal_runtime) = self.internal_runtime.take() { if let Ok(internal_runtime) = Arc::try_unwrap(internal_runtime) { internal_runtime.shutdown_background(); @@ -270,7 +382,6 @@ impl KVStore for PostgresStore { let inner = Arc::clone(&self.inner); let runtime = self.internal_runtime(); async move { - inner.assert_store_lock(); let task = runtime.spawn(async move { inner.read_internal(&primary_namespace, &secondary_namespace, &key).await }); @@ -288,9 +399,10 @@ impl KVStore for PostgresStore { let key = key.to_string(); let inner = Arc::clone(&self.inner); let runtime = self.internal_runtime(); + let resources = Arc::clone(&self.resources); async move { - inner.assert_store_lock(); let task = runtime.spawn(async move { + let _resources = resources; inner .write_internal( inner_lock_ref, @@ -317,9 +429,10 @@ impl KVStore for PostgresStore { let key = key.to_string(); let inner = Arc::clone(&self.inner); let runtime = self.internal_runtime(); + let resources = Arc::clone(&self.resources); async move { - inner.assert_store_lock(); let task = runtime.spawn(async move { + let _resources = resources; inner .remove_internal( inner_lock_ref, @@ -343,7 +456,6 @@ impl KVStore for PostgresStore { let inner = Arc::clone(&self.inner); let runtime = self.internal_runtime(); async move { - inner.assert_store_lock(); let task = runtime.spawn(async move { inner.list_internal(&primary_namespace, &secondary_namespace).await }); @@ -361,7 +473,6 @@ impl PaginatedKVStore for PostgresStore { let inner = Arc::clone(&self.inner); let runtime = self.internal_runtime(); async move { - inner.assert_store_lock(); let task = runtime.spawn(async move { inner .list_paginated_internal(&primary_namespace, &secondary_namespace, page_token) @@ -379,7 +490,6 @@ impl MigratableKVStore for PostgresStore { let inner = Arc::clone(&self.inner); let runtime = self.internal_runtime(); async move { - inner.assert_store_lock(); let task = runtime.spawn(async move { inner.list_all_keys_internal().await }); handle_runtime_task_result(task.await) } @@ -388,11 +498,10 @@ impl MigratableKVStore for PostgresStore { struct PostgresStoreInner { pool: SmallPool, - // PostgreSQL advisory locks are session-scoped, so keep the connection that acquired our lock - // alive for the lifetime of the store. - lock_client: ClientConnection, config: Config, kv_table_name_sql: String, + node_lease_table_name_sql: String, + lease_owner_id: [u8; 32], tls: PgTlsConnector, write_version_locks: Mutex>>>, logger: Option>, @@ -405,6 +514,11 @@ impl PostgresStoreInner { ) -> io::Result { let kv_table_name = kv_table_name.unwrap_or(DEFAULT_KV_TABLE_NAME.to_string()); let kv_table_name_sql = sql_table_identifier(&kv_table_name)?; + let node_lease_table_name_sql = sql_node_lease_table_identifier(&kv_table_name)?; + let mut lease_owner_id = [0u8; 32]; + getrandom::fill(&mut lease_owner_id).map_err(|e| { + io::Error::new(io::ErrorKind::Other, format!("Failed to generate lease owner ID: {e}")) + })?; let mut config: Config = connection_string.parse().map_err(|e: PgError| { let msg = format!("Failed to parse PostgreSQL connection string: {e}"); @@ -444,24 +558,6 @@ impl PostgresStoreInner { Self::create_database_if_not_exists(&config, &tls, logger.as_deref()).await?; let mut client = make_config_connection(&config, &tls).await?; - let lock_id = advisory_lock_id(&db_name, &kv_table_name); - let row = client.query_one("SELECT pg_try_advisory_lock($1)", &[&lock_id]).await.map_err( - |e| { - let msg = format!( - "Failed to acquire PostgreSQL store lock for database {db_name} and table {kv_table_name}: {e}" - ); - io::Error::new(io::ErrorKind::Other, msg) - }, - )?; - if !row.get::<_, bool>(0) { - return Err(io::Error::new( - io::ErrorKind::AlreadyExists, - format!( - "PostgreSQL store for database {db_name} and table {kv_table_name} is already in use" - ), - )); - } - let pool = SmallPool::new(&config, &tls).await?; let transaction = client.transaction().await.map_err(|e| { io::Error::new( @@ -550,6 +646,42 @@ impl PostgresStoreInner { io::Error::new(io::ErrorKind::Other, msg) })?; + // A unique owner ID makes takeover conflict with mutation key-share locks. + let sql = format!( + "CREATE TABLE IF NOT EXISTS {node_lease_table_name_sql} ( + id SMALLINT PRIMARY KEY CHECK (id = 1), + owner_id BYTEA NOT NULL UNIQUE, + expires_at TIMESTAMPTZ NOT NULL + )" + ); + transaction.execute(&sql, &[]).await.map_err(|e| { + io::Error::new(io::ErrorKind::Other, format!("Failed to create node lease table: {e}")) + })?; + + // Acquire the lease after schema setup. + // Schema setup and lease ownership commit or roll back together. + let lease_duration_secs = NODE_LEASE_DURATION.as_secs() as i64; + let acquire_sql = format!( + "INSERT INTO {node_lease_table_name_sql} (id, owner_id, expires_at) + VALUES (1, $1, clock_timestamp() + ($2::bigint * interval '1 second')) + ON CONFLICT (id) DO UPDATE SET + owner_id = EXCLUDED.owner_id, + expires_at = EXCLUDED.expires_at + WHERE {node_lease_table_name_sql}.expires_at <= clock_timestamp()" + ); + let acquired = transaction + .execute(&acquire_sql, &[&lease_owner_id.as_slice(), &lease_duration_secs]) + .await + .map_err(|e| { + io::Error::new(io::ErrorKind::Other, format!("Failed to acquire node lease: {e}")) + })?; + if acquired != 1 { + return Err(io::Error::new( + io::ErrorKind::AlreadyExists, + "PostgreSQL node lease is unavailable", + )); + } + transaction.commit().await.map_err(|e| { io::Error::new( io::ErrorKind::Other, @@ -560,9 +692,10 @@ impl PostgresStoreInner { let write_version_locks = Mutex::new(HashMap::new()); Ok(Self { pool, - lock_client: client, config, kv_table_name_sql, + node_lease_table_name_sql, + lease_owner_id, tls, write_version_locks, logger, @@ -648,18 +781,22 @@ impl PostgresStoreInner { } async fn locked_client(&self) -> io::Result> { - let client = self.pool.get(&self.config, &self.tls, self.logger.as_deref()).await?; - // Recheck after any runtime queue, per-key write lock, and pool wait, immediately before - // the caller accesses PostgreSQL. - self.assert_store_lock(); - Ok(client) + self.pool.get(&self.config, &self.tls, self.logger.as_deref()).await } - fn assert_store_lock(&self) { - assert!( - !self.lock_client.is_closed(), - "PostgreSQL store lock connection closed; continuing may corrupt node state" - ); + async fn release_node_lease(&self) -> io::Result<()> { + let lease_table = &self.node_lease_table_name_sql; + let sql = format!("DELETE FROM {lease_table} WHERE id = 1 AND owner_id = $1"); + let mut locked = self.locked_client().await?; + let err_map = + |e| io::Error::new(io::ErrorKind::Other, format!("Failed to release node lease: {e}")); + query_with_retry!( + self, + locked, + err_map, + locked.execute(&sql, &[&self.lease_owner_id.as_slice()]) + )?; + Ok(()) } fn get_inner_lock_ref(&self, locking_key: String) -> Arc> { @@ -720,11 +857,21 @@ impl PostgresStoreInner { check_namespace_key_validity(primary_namespace, secondary_namespace, Some(key), "write")?; self.execute_locked_write(inner_lock_ref, locking_key, version, async move || { + let kv_table = &self.kv_table_name_sql; + let lease_table = &self.node_lease_table_name_sql; + // PostgreSQL takes the KV table lock before executing the lease check. The key-share lock + // allows concurrent mutations and renewal, but blocks takeover until the statement commits, + // even if the lease expires while the mutation waits for a KV row lock. let sql = format!( - "INSERT INTO {} (primary_namespace, secondary_namespace, key, value) \ - VALUES ($1, $2, $3, $4) \ - ON CONFLICT (primary_namespace, secondary_namespace, key) DO UPDATE SET value = EXCLUDED.value", - self.kv_table_name_sql + "WITH valid_lease AS ( + SELECT 1 FROM {lease_table} + WHERE id = 1 AND owner_id = $1 AND expires_at > clock_timestamp() FOR KEY SHARE + ), mutation AS ( + INSERT INTO {kv_table} (primary_namespace, secondary_namespace, key, value) + SELECT $2, $3, $4, $5 WHERE EXISTS (SELECT 1 FROM valid_lease) + ON CONFLICT (primary_namespace, secondary_namespace, key) DO UPDATE SET value = EXCLUDED.value + ) + SELECT 1 FROM valid_lease" ); let err_map = |e: PgError| { @@ -739,16 +886,26 @@ impl PostgresStoreInner { }; let mut locked = self.locked_client().await?; - query_with_retry!( - self, - locked, - err_map, - locked.execute( - sql.as_str(), - &[&primary_namespace, &secondary_namespace, &key, &buf], + // Retry only preparation, which cannot mutate data. A failed execution may have committed + // before its response was lost, so it must not be retried. + let statement = query_with_retry!(self, locked, &err_map, locked.prepare(&sql))?; + let lease = locked + .query_opt( + &statement, + &[ + &self.lease_owner_id.as_slice(), + &primary_namespace, + &secondary_namespace, + &key, + &buf, + ], ) - ) - .map(|_| ()) + .await + .map_err(err_map)?; + if lease.is_none() { + panic!("PostgreSQL node lease was lost; continuing may corrupt node state"); + } + Ok(()) }) .await } @@ -760,9 +917,18 @@ impl PostgresStoreInner { check_namespace_key_validity(primary_namespace, secondary_namespace, Some(key), "remove")?; self.execute_locked_write(inner_lock_ref, locking_key, version, async move || { + let kv_table = &self.kv_table_name_sql; + let lease_table = &self.node_lease_table_name_sql; + // Return the lease row even when the key does not exist. let sql = format!( - "DELETE FROM {} WHERE primary_namespace=$1 AND secondary_namespace=$2 AND key=$3", - self.kv_table_name_sql + "WITH valid_lease AS ( + SELECT 1 FROM {lease_table} + WHERE id = 1 AND owner_id = $1 AND expires_at > clock_timestamp() FOR KEY SHARE + ), mutation AS ( + DELETE FROM {kv_table} WHERE primary_namespace=$2 AND secondary_namespace=$3 AND key=$4 + AND EXISTS (SELECT 1 FROM valid_lease) + ) + SELECT 1 FROM valid_lease" ); let err_map = |e: PgError| { @@ -777,13 +943,25 @@ impl PostgresStoreInner { }; let mut locked = self.locked_client().await?; - query_with_retry!( - self, - locked, - err_map, - locked.execute(sql.as_str(), &[&primary_namespace, &secondary_namespace, &key]) - ) - .map(|_| ()) + // Retry only preparation, which cannot mutate data. A failed execution may have committed + // before its response was lost, so it must not be retried. + let statement = query_with_retry!(self, locked, &err_map, locked.prepare(&sql))?; + let lease = locked + .query_opt( + &statement, + &[ + &self.lease_owner_id.as_slice(), + &primary_namespace, + &secondary_namespace, + &key, + ], + ) + .await + .map_err(err_map)?; + if lease.is_none() { + panic!("PostgreSQL node lease was lost; continuing may corrupt node state"); + } + Ok(()) }) .await } @@ -978,9 +1156,13 @@ mod tests { } async fn cleanup_store(store: &PostgresStore) { + // Stop renewal before removing its table. + store.resources.lease_renewal_task.abort(); let kv_table = store.inner.kv_table_name_sql.clone(); + let lease_table = &store.inner.node_lease_table_name_sql; let client = store.inner.pool.connections[0].lock().await; - let _ = client.execute(&format!("DROP TABLE IF EXISTS {kv_table}"), &[]).await; + let _ = + client.execute(&format!("DROP TABLE IF EXISTS {kv_table}, {lease_table}"), &[]).await; } #[test] @@ -995,20 +1177,24 @@ mod tests { assert!(sql_identifier("").is_err()); assert!(sql_table_identifier("too.many.parts").is_err()); assert!(sql_table_identifier("schema.").is_err()); - } - - #[test] - fn test_postgres_advisory_lock_id_uses_database_and_table() { - let lock_id = advisory_lock_id("database_a", "table_a"); - assert_eq!(lock_id, advisory_lock_id("database_a", "table_a")); - assert_ne!(lock_id, advisory_lock_id("database_b", "table_a")); - assert_ne!(lock_id, advisory_lock_id("database_a", "table_b")); + assert_eq!(sql_node_lease_table_identifier("tenant-1").unwrap(), "\"tenant-1_node_lease\""); + assert_eq!( + sql_node_lease_table_identifier("tenant.select").unwrap(), + "\"tenant\".\"select_node_lease\"" + ); + assert!(sql_node_lease_table_identifier(&"a".repeat(MAX_KV_TABLE_NAME_BYTES)).is_ok()); + assert!(sql_node_lease_table_identifier(&"a".repeat(MAX_KV_TABLE_NAME_BYTES + 1)).is_err()); } #[tokio::test(flavor = "multi_thread")] - async fn test_postgres_store_advisory_lock() { - let table_name = "test_pg_advisory_lock"; + async fn test_postgres_store_lease() { + let table_name = "test_pg_lease"; let store = create_test_store(table_name).await; + let kv_table = &store.inner.kv_table_name_sql; + let mut client = + make_config_connection(&store.inner.config, &store.inner.tls).await.unwrap(); + // Make the second store attempt a schema change before its lease is rejected. + client.execute(&format!("COMMENT ON TABLE {kv_table} IS NULL"), &[]).await.unwrap(); let err = PostgresStore::new(test_connection_string(), None, Some(table_name.to_string()), None) @@ -1016,15 +1202,169 @@ mod tests { .err() .expect("a second store using the same database and table must fail"); assert_eq!(err.kind(), io::ErrorKind::AlreadyExists); + let row = client + .query_one("SELECT obj_description(to_regclass($1), 'pg_class')", &[kv_table]) + .await + .unwrap(); + assert_eq!(row.get::<_, Option<&str>>(0), None); + + // A write must complete while another transaction holds renewal's non-key update lock. + let transaction = client.transaction().await.unwrap(); + let lease_table = &store.inner.node_lease_table_name_sql; + transaction + .query_one(&format!("SELECT 1 FROM {lease_table} WHERE id = 1 FOR NO KEY UPDATE"), &[]) + .await + .unwrap(); + tokio::time::timeout( + Duration::from_secs(5), + KVStore::write(&store, "test_ns", "", "key", vec![1]), + ) + .await + .expect("renewal's row lock must not block writes") + .unwrap(); + transaction.commit().await.unwrap(); + + // Changing the owner must wait for a mutation's key-share lock to be released. + let transaction = client.transaction().await.unwrap(); + transaction + .query_one(&format!("SELECT 1 FROM {lease_table} WHERE id = 1 FOR KEY SHARE"), &[]) + .await + .unwrap(); + let mut contender = + make_config_connection(&store.inner.config, &store.inner.tls).await.unwrap(); + contender.batch_execute("SET lock_timeout = '1s'").await.unwrap(); + let takeover_sql = format!("UPDATE {lease_table} SET owner_id = $1 WHERE id = 1"); + let new_owner_id = [42u8; 32]; + let err = contender.execute(&takeover_sql, &[&new_owner_id.as_slice()]).await.unwrap_err(); + assert_eq!(err.code(), Some(&tokio_postgres::error::SqlState::LOCK_NOT_AVAILABLE)); + transaction.commit().await.unwrap(); + + let takeover = contender.transaction().await.unwrap(); + assert_eq!(takeover.execute(&takeover_sql, &[&new_owner_id.as_slice()]).await.unwrap(), 1); + takeover.rollback().await.unwrap(); + + // Mutation futures retain the lease until they finish, even after dropping the store. + let write = KVStore::write(&store, "test_ns", "", "key", vec![2]); + let read = KVStore::read(&store, "test_ns", "", "key"); + drop(store); + write.await.unwrap(); + assert_eq!(read.await.unwrap(), vec![2]); + + let store = create_test_store(table_name).await; + let remove = KVStore::remove(&store, "test_ns", "", "key", false); + drop(store); + remove.await.unwrap(); + // The last mutation releases the lease so another store can acquire it immediately. + let store = create_test_store(table_name).await; + assert_eq!( + KVStore::read(&store, "test_ns", "", "key").await.unwrap_err().kind(), + io::ErrorKind::NotFound + ); cleanup_store(&store).await; } + #[tokio::test(flavor = "multi_thread")] + async fn test_background_renewal_and_lease_loss() { + let mut store = create_test_store("test_pg_background_lease_loss").await; + let mut client = store.inner.pool.connections[0].lock().await; + let transaction = client.transaction().await.unwrap(); + let lease_table = &store.inner.node_lease_table_name_sql; + let sql = format!("SELECT expires_at FROM {lease_table} WHERE id = 1 FOR KEY SHARE"); + let mut expires_at = + transaction.query_one(&sql, &[]).await.unwrap().get::<_, std::time::SystemTime>(0); + // Observe two renewals while holding a mutation fence, without KV writes. + tokio::time::timeout( + NODE_LEASE_RENEWAL_INTERVAL * 2 + NODE_LEASE_RENEWAL_TIMEOUT + Duration::from_secs(5), + async { + for _ in 0..2 { + loop { + let renewed_until = transaction + .query_one(&sql, &[]) + .await + .unwrap() + .get::<_, std::time::SystemTime>(0); + if renewed_until > expires_at { + expires_at = renewed_until; + break; + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + } + }, + ) + .await + .expect("background renewal must extend the lease while a mutation fence is held"); + transaction.commit().await.unwrap(); + drop(client); + + expire_lease(&store).await; + let error = tokio::time::timeout( + NODE_LEASE_RENEWAL_INTERVAL + NODE_LEASE_RENEWAL_TIMEOUT + Duration::from_secs(5), + &mut Arc::get_mut(&mut store.resources).unwrap().lease_renewal_task, + ) + .await + .expect("background renewal must detect lease loss without another store operation") + .expect_err("background renewal must panic on lease loss"); + let panic = error.into_panic(); + assert!(panic.downcast_ref::<&str>().unwrap().contains("PostgreSQL node lease was lost")); + cleanup_store(&store).await; + } + + #[tokio::test(flavor = "multi_thread")] + async fn test_background_renewal_panics_on_timeout() { + let mut store = create_test_store("test_pg_background_lease_timeout").await; + // Exhaust the pool to verify the renewal timeout includes waiting for a connection. + let first = store.inner.pool.connections[0].lock().await; + let second = store.inner.pool.connections[1].lock().await; + let error = tokio::time::timeout( + NODE_LEASE_RENEWAL_INTERVAL + NODE_LEASE_RENEWAL_TIMEOUT + Duration::from_secs(5), + &mut Arc::get_mut(&mut store.resources).unwrap().lease_renewal_task, + ) + .await + .expect("background renewal must time out while waiting for the pool") + .expect_err("background renewal must panic on timeout"); + let panic = error.into_panic(); + assert!(panic + .downcast_ref::<&str>() + .unwrap() + .contains("PostgreSQL node lease renewal timed out")); + drop(first); + drop(second); + // The panicked renewal task cannot release its lease during shutdown. + store.inner.release_node_lease().await.unwrap(); + cleanup_store(&store).await; + } + + #[tokio::test(flavor = "multi_thread")] + async fn test_drop_with_blocked_pool() { + let store = create_test_store("test_pg_drop_blocked_pool").await; + let inner = Arc::clone(&store.inner); + let client = make_config_connection(&inner.config, &inner.tls).await.unwrap(); + let _first = inner.pool.connections[0].lock().await; + let _second = inner.pool.connections[1].lock().await; + let start = std::time::Instant::now(); + drop(store); + assert!(start.elapsed() < NODE_LEASE_RELEASE_TIMEOUT + Duration::from_secs(2)); + + client + .execute( + &format!( + "DROP TABLE {}, {}", + inner.kv_table_name_sql, inner.node_lease_table_name_sql + ), + &[], + ) + .await + .unwrap(); + } + #[tokio::test(flavor = "multi_thread")] async fn test_postgres_store_rolls_back_failed_schema_setup() { let table_name = "test_pg_schema_setup_rollback"; let store = create_test_store(table_name).await; let kv_table = store.inner.kv_table_name_sql.clone(); + let lease_table = store.inner.node_lease_table_name_sql.clone(); let client = make_config_connection(&store.inner.config, &store.inner.tls).await.unwrap(); drop(store); @@ -1046,7 +1386,7 @@ mod tests { .query_one("SELECT obj_description(to_regclass($1), 'pg_class')", &[&kv_table]) .await .unwrap(); - client.execute(&format!("DROP TABLE {kv_table}"), &[]).await.unwrap(); + client.execute(&format!("DROP TABLE {kv_table}, {lease_table}"), &[]).await.unwrap(); assert_eq!(row.get::<_, Option<&str>>(0), None); } @@ -1090,20 +1430,26 @@ mod tests { } async fn kill_connection(store: &PostgresStore) { - // Terminate every backend in the pool so the next op deterministically - // hits a closed connection regardless of which slot `get` selects. + // Terminate each pooled backend to exercise reconnection. Background renewal may + // reconnect a slot before the next store operation. for mutex in &store.inner.pool.connections { let client = mutex.lock().await; let _ = client.execute("SELECT pg_terminate_backend(pg_backend_pid())", &[]).await; } } - async fn kill_lock_connection(store: &PostgresStore) { - let client = &store.inner.lock_client; - let _ = client.execute("SELECT pg_terminate_backend(pg_backend_pid())", &[]).await; - while !client.is_closed() { - tokio::task::yield_now().await; - } + async fn expire_lease(store: &PostgresStore) { + let client = store.inner.pool.connections[0].lock().await; + let lease_table = &store.inner.node_lease_table_name_sql; + client + .execute( + &format!( + "UPDATE {lease_table} SET expires_at = clock_timestamp() - interval '1 second' WHERE id = 1" + ), + &[], + ) + .await + .unwrap(); } #[tokio::test(flavor = "multi_thread")] @@ -1131,29 +1477,16 @@ mod tests { } #[tokio::test(flavor = "multi_thread")] - #[should_panic( - expected = "PostgreSQL store lock connection closed; continuing may corrupt node state" - )] - async fn test_postgres_store_panics_when_lock_connection_closes() { - let table_name = "test_pg_lock_connection_closed"; - let store = create_test_store(table_name).await; - - kill_lock_connection(&store).await; - cleanup_store(&store).await; - KVStore::write(&store, "test_ns", "test_sub", "key", vec![1u8]).await.unwrap(); - } - - #[tokio::test(flavor = "multi_thread")] - async fn test_queued_write_rechecks_closed_lock_connection() { - let table_name = "test_pg_queued_write_lock_connection_closed"; + async fn test_queued_write_rechecks_lease() { + let table_name = "test_pg_queued_write_lease"; let store = create_test_store(table_name).await; + KVStore::remove(&store, "test_ns", "test_sub", "missing", false).await.unwrap(); let locking_key = store.build_locking_key("test_ns", "test_sub", "key"); let inner_lock_ref = store.inner.get_inner_lock_ref(locking_key); let inner_lock = inner_lock_ref.lock().await; let mut write = Box::pin(KVStore::write(&store, "test_ns", "test_sub", "key", vec![1u8])); - // Poll the public method once so its initial check passes and its internal write is queued - // on the per-key lock before closing the store lock connection. + // Start the write while holding the per-key lock so it cannot proceed before lease expiry. std::future::poll_fn(|cx| match write.as_mut().poll(cx) { std::task::Poll::Pending => std::task::Poll::Ready(()), std::task::Poll::Ready(result) => { @@ -1162,7 +1495,10 @@ mod tests { }) .await; - kill_lock_connection(&store).await; + expire_lease(&store).await; + // Even a delete that affects no rows must reject an expired lease before takeover. + let remove = tokio::spawn(KVStore::remove(&store, "test_ns", "test_sub", "missing", false)); + assert!(remove.await.unwrap_err().is_panic()); let second_store = create_test_store(table_name).await; drop(inner_lock); @@ -1171,9 +1507,21 @@ mod tests { assert!(err.is_panic()); let err = KVStore::read(&second_store, "test_ns", "test_sub", "key") .await - .expect_err("the queued write must not access PostgreSQL"); + .expect_err("the queued write must not persist data"); assert_eq!(err.kind(), io::ErrorKind::NotFound); + KVStore::write(&second_store, "test_ns", "test_sub", "key", vec![2]).await.unwrap(); + let remove = tokio::spawn(KVStore::remove(&store, "test_ns", "test_sub", "key", false)); + assert!(remove.await.unwrap_err().is_panic()); + // Exercise release even if the old renewal task has already panicked. Neither release + // nor dropping the old owner may remove the replacement's lease. + store.inner.release_node_lease().await.unwrap(); + drop(store); + assert_eq!( + KVStore::read(&second_store, "test_ns", "test_sub", "key").await.unwrap(), + vec![2] + ); + KVStore::write(&second_store, "test_ns", "test_sub", "key", vec![3]).await.unwrap(); cleanup_store(&second_store).await; } diff --git a/tests/common/mod.rs b/tests/common/mod.rs index 5a409b8104..d214ebfbb6 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -2198,17 +2198,18 @@ pub(crate) fn test_connection_string() -> String { .unwrap_or_else(|_| "host=localhost user=postgres password=postgres".to_string()) } -/// Drops the given table from the `ldk_db` database, ignoring the case where the database doesn't -/// exist yet. Used to ensure a clean slate before and after Postgres-backed tests. +// Best-effort cleanup of the KV and lease tables in `ldk_db` between tests. #[cfg(feature = "storage-postgres")] pub(crate) async fn drop_table(table_name: &str) { let connection_string = format!("{} dbname=ldk_db", test_connection_string()); let Ok((client, connection)) = tokio_postgres::connect(&connection_string, tokio_postgres::NoTls).await else { - // Database doesn't exist yet — nothing to drop. + // Ignore connection failures, including when the database does not exist yet. return; }; tokio::spawn(connection); - let _ = client.execute(&format!("DROP TABLE IF EXISTS {table_name}"), &[]).await; + let _ = client + .execute(&format!("DROP TABLE IF EXISTS {table_name}, {table_name}_node_lease"), &[]) + .await; }