From 438f80e885df75c142e3476dcfe90c7f8aa2dc9e Mon Sep 17 00:00:00 2001 From: Joost Jager Date: Wed, 23 Sep 2026 11:15:26 +0200 Subject: [PATCH 1/3] 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 4ffb948ff..ace93f6c8 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 effe76e3f56102365c427619ddef3f2976e9e1a2 Mon Sep 17 00:00:00 2001 From: Joost Jager Date: Wed, 23 Sep 2026 11:16:31 +0200 Subject: [PATCH 2/3] Extract shared PostgreSQL mutation execution Move pooled connection acquisition and query retry handling for writes and removals into execute_mutation. Preserve advisory lock checks, per-key write ordering, SQL, and error mapping so lease enforcement can be added to the shared helper separately. --- src/io/postgres_store/mod.rs | 31 ++++++++++++++++--------------- 1 file changed, 16 insertions(+), 15 deletions(-) diff --git a/src/io/postgres_store/mod.rs b/src/io/postgres_store/mod.rs index ace93f6c8..5e7dc697d 100644 --- a/src/io/postgres_store/mod.rs +++ b/src/io/postgres_store/mod.rs @@ -20,6 +20,7 @@ use lightning_types::string::PrintableString; use native_tls::TlsConnector; use postgres_native_tls::MakeTlsConnector; use tokio_postgres::config::SslMode; +use tokio_postgres::types::ToSql; use tokio_postgres::{Config, Error as PgError}; use self::pool::{make_config_connection, ClientConnection, PgTlsConnector, SmallPool}; @@ -662,6 +663,13 @@ impl PostgresStoreInner { ); } + async fn execute_mutation io::Error>( + &self, sql: &str, params: &[&(dyn ToSql + Sync)], err_map: F, + ) -> io::Result<()> { + let mut locked = self.locked_client().await?; + query_with_retry!(self, locked, err_map, locked.execute(sql, params)).map(|_| ()) + } + fn get_inner_lock_ref(&self, locking_key: String) -> Arc> { let mut outer_lock = self.write_version_locks.lock().unwrap(); Arc::clone(&outer_lock.entry(locking_key).or_default()) @@ -738,17 +746,12 @@ impl PostgresStoreInner { io::Error::new(io::ErrorKind::Other, msg) }; - let mut locked = self.locked_client().await?; - query_with_retry!( - self, - locked, + self.execute_mutation( + sql.as_str(), + &[&primary_namespace, &secondary_namespace, &key, &buf], err_map, - locked.execute( - sql.as_str(), - &[&primary_namespace, &secondary_namespace, &key, &buf], - ) ) - .map(|_| ()) + .await }) .await } @@ -776,14 +779,12 @@ impl PostgresStoreInner { io::Error::new(io::ErrorKind::Other, msg) }; - let mut locked = self.locked_client().await?; - query_with_retry!( - self, - locked, + self.execute_mutation( + sql.as_str(), + &[&primary_namespace, &secondary_namespace, &key], err_map, - locked.execute(sql.as_str(), &[&primary_namespace, &secondary_namespace, &key]) ) - .map(|_| ()) + .await }) .await } From 6532d1b2108cab05bd15c2337312e21c1cff5929 Mon Sep 17 00:00:00 2001 From: Joost Jager Date: Wed, 23 Sep 2026 13:47:49 +0200 Subject: [PATCH 3/3] Replace PostgreSQL advisory locks with renewable leases Replace the temporary session advisory lock with an expiring lease in a companion table. Acquire it after schema setup and validate and renew it after each KV mutation in the same transaction, so stale writes and deletes roll back before commit. Renew idle leases in the background and release them by owner ID on drop. Explicitly panic on lease loss in the calling task or background renewal task, and on background failure or timeout, without adding recovery APIs. Cover rejected setup rollback, repeated idle renewal, renewal timeouts, stale writes and deletes, and releasing an old owner after takeover. --- src/builder.rs | 18 +- src/io/postgres_store/mod.rs | 432 ++++++++++++++++++++++++++--------- tests/common/mod.rs | 7 +- 3 files changed, 341 insertions(+), 116 deletions(-) diff --git a/src/builder.rs b/src/builder.rs index 1158044e4..8c7695a36 100644 --- a/src/builder.rs +++ b/src/builder.rs @@ -737,11 +737,10 @@ 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`. + /// Mutations panic on detected lease loss. Failed or timed-out renewals panic in the background + /// renewal task. Node recovery is not handled automatically. /// /// 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 +1333,10 @@ 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`. + /// Mutations panic on detected lease loss. Failed or timed-out renewals panic in the background + /// renewal task. Node recovery is not handled automatically. /// /// 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 5e7dc697d..7aa13b6cb 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, @@ -21,7 +21,7 @@ use native_tls::TlsConnector; use postgres_native_tls::MakeTlsConnector; use tokio_postgres::config::SslMode; use tokio_postgres::types::ToSql; -use tokio_postgres::{Config, Error as PgError}; +use tokio_postgres::{Config, Error as PgError, GenericClient}; use self::pool::{make_config_connection, ClientConnection, PgTlsConnector, SmallPool}; use crate::io::utils::check_namespace_key_validity; @@ -46,17 +46,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') { @@ -84,10 +82,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. @@ -96,14 +108,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)), @@ -127,6 +135,10 @@ 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. Lease rejection during a mutation panics in the calling task. A failed or +/// timed-out background renewal panics in the renewal task. This does not provide node shutdown +/// or recovery. /// /// [PostgreSQL]: https://www.postgresql.org pub struct PostgresStore { @@ -138,6 +150,8 @@ pub struct PostgresStore { // A store-internal runtime that drives PostgreSQL I/O independently from the node runtime. internal_runtime: Option>, + + lease_renewal_task: Option>, } // tokio::sync::Mutex (used for the DB client) contains UnsafeCell which opts out of @@ -160,14 +174,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 @@ -201,8 +211,42 @@ 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_renewal_task = internal_runtime.spawn(async move { + let mut interval = tokio::time::interval(NODE_LEASE_RENEWAL_INTERVAL); + 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, + inner_ref.renew_node_lease(&**locked) + ) + }; + // Bound both the pool wait and query so an unreachable database also causes a panic. + tokio::time::timeout(NODE_LEASE_RENEWAL_TIMEOUT, renewal) + .await + .expect("PostgreSQL node lease renewal timed out") + .expect("Failed to renew PostgreSQL node lease"); + } + }); + + Ok(Self { + inner, + next_write_version: AtomicU64::new(1), + internal_runtime: Some(internal_runtime), + lease_renewal_task: Some(lease_renewal_task), + }) } fn build_tls_connector(certificate_pem: Option) -> io::Result { @@ -253,6 +297,30 @@ impl PostgresStore { impl Drop for PostgresStore { fn drop(&mut self) { + if let Some(internal_runtime) = self.internal_runtime.as_ref() { + let renewal_task = self.lease_renewal_task.take(); + if let Some(task) = renewal_task.as_ref() { + task.abort(); + } + + let runtime_handle = internal_runtime.handle().clone(); + let inner = Arc::clone(&self.inner); + let _ = std::thread::spawn(move || { + runtime_handle.block_on(async move { + if let Some(task) = renewal_task { + let _ = task.await; + } + + let _ = tokio::time::timeout( + NODE_LEASE_RELEASE_TIMEOUT, + inner.release_node_lease(), + ) + .await; + }); + }) + .join(); + } + if let Some(internal_runtime) = self.internal_runtime.take() { if let Ok(internal_runtime) = Arc::try_unwrap(internal_runtime) { internal_runtime.shutdown_background(); @@ -271,7 +339,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 }); @@ -290,7 +357,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 .write_internal( @@ -319,7 +385,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 .remove_internal( @@ -344,7 +409,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 }); @@ -362,7 +426,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) @@ -380,7 +443,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) } @@ -389,11 +451,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>, @@ -406,6 +467,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}"); @@ -445,24 +511,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( @@ -551,6 +599,42 @@ impl PostgresStoreInner { io::Error::new(io::ErrorKind::Other, msg) })?; + 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, + 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, matching mutation lock order. + // 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() + RETURNING id" + ); + let row = transaction + .query_opt(&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 row.is_none() { + 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, @@ -558,12 +642,16 @@ impl PostgresStoreInner { ) })?; + // Drop the setup client; the pool has its own connections. + drop(client); + 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, @@ -649,25 +737,80 @@ 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(()) + } + + async fn renew_node_lease(&self, client: &impl GenericClient) -> Result<(), PgError> { + let lease_duration_secs = NODE_LEASE_DURATION.as_secs() as i64; + let lease_table = &self.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 updated = client + .execute(&update_sql, &[&self.lease_owner_id.as_slice(), &lease_duration_secs]) + .await?; + if updated != 1 { + panic!("PostgreSQL node lease was lost; continuing may corrupt node state"); + } + Ok(()) } async fn execute_mutation io::Error>( &self, sql: &str, params: &[&(dyn ToSql + Sync)], err_map: F, ) -> io::Result<()> { let mut locked = self.locked_client().await?; - query_with_retry!(self, locked, err_map, locked.execute(sql, params)).map(|_| ()) + // Retry only BEGIN. After it succeeds, the lease check and mutation must use the same + // transaction. + let transaction_result = locked.transaction().await; + let reconnect = transaction_result.as_ref().is_err_and(PgError::is_closed); + let transaction_result = if reconnect { + if let (Some(logger), Err(e)) = (self.logger.as_ref(), &transaction_result) { + log_debug!(logger, "Reconnecting to PostgreSQL after error: {e}"); + } + drop(transaction_result); + *locked = make_config_connection(&self.config, &self.tls).await?; + locked.transaction().await + } else { + transaction_result + }; + let transaction = transaction_result.map_err(|e| { + io::Error::new( + io::ErrorKind::Other, + format!("Failed to start fenced mutation transaction: {e}"), + ) + })?; + transaction.execute(sql, params).await.map_err(err_map)?; + // Renew after the mutation to give the lease a later expiry. Rejection rolls back the + // mutation; success holds the lease row lock until commit. + self.renew_node_lease(&transaction).await.map_err(|e| { + io::Error::new(io::ErrorKind::Other, format!("Failed to renew node lease: {e}")) + })?; + // A connection failure can hide a successful COMMIT if its response is lost. Return the + // error without retrying, since the mutation may already be committed. + transaction.commit().await.map_err(|e| { + io::Error::new( + io::ErrorKind::Other, + format!("Failed to commit fenced mutation transaction: {e}"), + ) + })?; + Ok(()) } fn get_inner_lock_ref(&self, locking_key: String) -> Arc> { @@ -996,20 +1139,23 @@ 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 client = store.inner.pool.connections[0].lock().await; + // 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) @@ -1017,10 +1163,86 @@ 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); + drop(client); 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 client = store.inner.pool.connections[0].lock().await; + let lease_table = &store.inner.node_lease_table_name_sql; + let sql = format!("SELECT expires_at FROM {lease_table} WHERE id = 1"); + let mut expires_at = + client.query_one(&sql, &[]).await.unwrap().get::<_, std::time::SystemTime>(0); + // Observe two extensions without KV writes, so a single renewal at startup is insufficient. + 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 = client + .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 the store is idle"); + drop(client); + + expire_lease(&store).await; + let renewal_task = store.lease_renewal_task.take().unwrap(); + let error = tokio::time::timeout( + NODE_LEASE_RENEWAL_INTERVAL + NODE_LEASE_RENEWAL_TIMEOUT + Duration::from_secs(5), + 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; + let renewal_task = store.lease_renewal_task.take().unwrap(); + // 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), + 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::() + .unwrap() + .contains("PostgreSQL node lease renewal timed out")); + drop(first); + drop(second); + 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"; @@ -1091,20 +1313,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")] @@ -1132,29 +1360,15 @@ 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; 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) => { @@ -1163,7 +1377,7 @@ mod tests { }) .await; - kill_lock_connection(&store).await; + expire_lease(&store).await; let second_store = create_test_store(table_name).await; drop(inner_lock); @@ -1172,9 +1386,19 @@ 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()); + // Dropping the old owner must not release the replacement's lease. + 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 5a409b810..4bc5963ba 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -2200,15 +2200,18 @@ pub(crate) fn test_connection_string() -> 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. +/// Also drops the companion lease table. #[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; }