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 4ffb948ff9..71a0502f7d 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, @@ -20,11 +20,12 @@ use lightning_types::string::PrintableString; use native_tls::TlsConnector; use postgres_native_tls::MakeTlsConnector; use tokio_postgres::config::SslMode; -use tokio_postgres::{Config, Error as PgError}; +use tokio_postgres::types::ToSql; +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; -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 +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') { @@ -83,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. @@ -95,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)), @@ -126,6 +135,11 @@ 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 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 { @@ -137,6 +151,10 @@ pub struct PostgresStore { // A store-internal runtime that drives PostgreSQL I/O independently from the node runtime. internal_runtime: Option>, + + lease_renewal_task: Option>, + lease_shutdown_sender: Option>, + lease_shutdown_complete: Mutex>>, } // tokio::sync::Mutex (used for the DB client) contains UnsafeCell which opts out of @@ -159,14 +177,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 +214,58 @@ 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, mut 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); + loop { + // The first tick is immediate; later attempts follow the renewal interval. + tokio::select! { + biased; + _ = &mut shutdown_rx => break, + _ = interval.tick() => {}, + } + 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, but allow shutdown to cancel either. + let result = tokio::select! { + biased; + _ = &mut shutdown_rx => break, + result = tokio::time::timeout(NODE_LEASE_RENEWAL_TIMEOUT, renewal) => result, + }; + result + .expect("PostgreSQL node lease renewal timed out") + .expect("Failed to renew PostgreSQL node lease"); + } + + let result = inner_ref.release_node_lease().await; + let _ = completion_sender.send(result); + }); + + Ok(Self { + inner, + next_write_version: AtomicU64::new(1), + internal_runtime: Some(internal_runtime), + lease_renewal_task: Some(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 { @@ -252,6 +316,28 @@ impl PostgresStore { impl Drop for PostgresStore { fn drop(&mut self) { + if let Some(sender) = self.lease_shutdown_sender.take() { + let _ = sender.send(()); + } + // The store runtime drives shutdown independently. Bound the entire wait here and cancel + // the task below if release does not finish in time. + let result = self + .lease_shutdown_complete + .get_mut() + .unwrap() + .recv_timeout(NODE_LEASE_RELEASE_TIMEOUT); + 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}"), + } + } + if let Some(task) = self.lease_renewal_task.take() { + // Cancel any remaining work if shutdown failed or timed out. + 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 +356,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 }); @@ -289,7 +374,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( @@ -318,7 +402,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( @@ -343,7 +426,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 +443,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 +460,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 +468,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 +484,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}"); @@ -443,24 +527,14 @@ impl PostgresStoreInner { Self::create_database_if_not_exists(&config, &tls, logger.as_deref()).await?; - let 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 mut client = make_config_connection(&config, &tls).await?; + 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 @@ -476,13 +550,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 +587,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,19 +611,64 @@ 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?; + 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, + format!("Failed to commit PostgreSQL schema setup transaction: {e}"), + ) + })?; + + // 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, @@ -630,18 +754,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 + } + + 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 assert_store_lock(&self) { - assert!( - !self.lock_client.is_closed(), - "PostgreSQL store lock connection closed; continuing may corrupt node state" + 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?; + // 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> { @@ -720,17 +906,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 } @@ -758,14 +939,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 } @@ -977,20 +1156,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) @@ -998,10 +1180,141 @@ 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); + // 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 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; @@ -1042,20 +1355,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")] @@ -1083,29 +1402,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) => { @@ -1114,7 +1419,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); @@ -1123,9 +1428,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; }