Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 22 additions & 15 deletions src/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -683,6 +683,9 @@ impl NodeBuilder {
/// Builds a [`Node`] instance with a [PostgreSQL] backend and according to the options
/// previously configured.
///
/// 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`.
///
/// Connects to the PostgreSQL database at the given `connection_string`, e.g.,
/// `"postgres://user:password@localhost/ldk_db"`.
///
Expand All @@ -696,13 +699,10 @@ impl NodeBuilder {
/// The given `kv_table_name` will be used or default to
/// [`DEFAULT_KV_TABLE_NAME`](io::postgres_store::DEFAULT_KV_TABLE_NAME).
///
/// # 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.
/// Opening a schema-v1 store upgrades it to the lease-aware schema v2. Stop all processes using
/// the v1 store before upgrading. For the first v2 open, use the same resolved database name and
/// byte-for-byte same `kv_table_name` spelling, including schema qualification, so its transition
/// lock matches v1. Older releases cannot reopen a v2 store, so downgrading is unsupported.
///
/// 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
Expand All @@ -729,7 +729,13 @@ impl NodeBuilder {
log_error!(logger, "Failed to set up Postgres store: {e}");
BuildError::KVStoreSetupFailed
})?;
self.build_with_store_runtime_and_logger(node_entropy, kv_store, runtime, logger)
let node_lease = kv_store.node_lease();
let mut node =
self.build_with_store_runtime_and_logger(node_entropy, kv_store, runtime, logger)?;
if !node.install_node_lease(node_lease) {
return Err(BuildError::KVStoreSetupFailed);
}
Ok(node)
}

/// Builds a [`Node`] instance with a [`FilesystemStoreV2`] backend and according to the options
Expand Down Expand Up @@ -1225,6 +1231,9 @@ impl ArcedNodeBuilder {
/// Builds a [`Node`] instance with a [PostgreSQL] backend and according to the options
/// previously configured.
///
/// 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`.
///
/// Connects to the PostgreSQL database at the given `connection_string`, e.g.,
/// `"postgres://user:password@localhost/ldk_db"`.
///
Expand All @@ -1238,13 +1247,10 @@ impl ArcedNodeBuilder {
/// The given `kv_table_name` will be used or default to
/// [`DEFAULT_KV_TABLE_NAME`](io::postgres_store::DEFAULT_KV_TABLE_NAME).
///
/// # 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.
/// Opening a schema-v1 store upgrades it to the lease-aware schema v2. Stop all processes using
/// the v1 store before upgrading. For the first v2 open, use the same resolved database name and
/// byte-for-byte same `kv_table_name` spelling, including schema qualification, so its transition
/// lock matches v1. Older releases cannot reopen a v2 store, so downgrading is unsupported.
///
/// 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
Expand Down Expand Up @@ -2377,6 +2383,7 @@ fn build_with_store_internal(
payment_store,
lnurl_auth,
is_running,
node_lease: None,
node_metrics,
om_mailbox,
async_payments_role,
Expand Down
2 changes: 2 additions & 0 deletions src/io/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@

//! Objects and traits for data persistence.

#[cfg_attr(not(feature = "postgres"), allow(dead_code))]
pub(crate) mod node_lease;
#[cfg(feature = "postgres")]
pub mod postgres_store;
pub mod sqlite_store;
Expand Down
173 changes: 173 additions & 0 deletions src/io/node_lease.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,173 @@
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// http://opensource.org/licenses/MIT>, at your option. You may not use this file except in
// accordance with one or both of these licenses.

use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};

use lightning::io;

pub(crate) const NODE_LEASE_DURATION: Duration = Duration::from_secs(30);
// Fail closed before the database lease expires, leaving time for process termination.
pub(crate) const NODE_LEASE_RENEWAL_DEADLINE: Duration = Duration::from_secs(20);
pub(crate) const NODE_LEASE_RENEWAL_INTERVAL: Duration = Duration::from_secs(10);
pub(crate) const NODE_LEASE_RETRY_INTERVAL: Duration = Duration::from_secs(1);
pub(crate) const NODE_LEASE_RELEASE_TIMEOUT: Duration = Duration::from_secs(5);

type LeaseLossHandler = Box<dyn FnOnce() + Send>;

pub(crate) struct NodeLease {
owner_id: [u8; 32],
lease_lost: AtomicBool,
last_confirmed_renewal: Mutex<Instant>,
loss_sender: tokio::sync::watch::Sender<bool>,
loss_handler: Mutex<Option<LeaseLossHandler>>,
}

impl NodeLease {
pub(crate) fn new() -> io::Result<Arc<Self>> {
let mut owner_id = [0u8; 32];
getrandom::fill(&mut owner_id).map_err(|e| {
io::Error::new(io::ErrorKind::Other, format!("Failed to generate lease owner ID: {e}"))
})?;
let (loss_sender, _) = tokio::sync::watch::channel(false);
Ok(Arc::new(Self {
owner_id,
lease_lost: AtomicBool::new(false),
last_confirmed_renewal: Mutex::new(Instant::now()),
loss_sender,
loss_handler: Mutex::new(None),
}))
}

pub(crate) fn owner_id(&self) -> &[u8; 32] {
&self.owner_id
}

pub(crate) fn is_lost(&self) -> bool {
self.lease_lost.load(Ordering::Acquire)
}

pub(crate) fn record_renewal_started_at(&self, renewal_started_at: Instant) {
if !self.is_lost() {
let mut last_confirmed_renewal = self.last_confirmed_renewal.lock().expect("lock");
*last_confirmed_renewal = (*last_confirmed_renewal).max(renewal_started_at);
}
}

pub(crate) fn renewal_deadline_elapsed(&self) -> bool {
self.last_confirmed_renewal.lock().expect("lock").elapsed() >= NODE_LEASE_RENEWAL_DEADLINE
}

pub(crate) async fn wait_for_renewal_deadline(&self) {
loop {
let last_confirmed_renewal = *self.last_confirmed_renewal.lock().expect("lock");
let deadline = last_confirmed_renewal + NODE_LEASE_RENEWAL_DEADLINE;
tokio::time::sleep_until(tokio::time::Instant::from_std(deadline)).await;
if self.renewal_deadline_elapsed() {
return;
}
}
}

pub(crate) fn ensure_operation_active(&self) -> io::Result<()> {
if self.is_lost() || self.renewal_deadline_elapsed() {
self.mark_lost();
Err(lease_lost_error())
} else {
Ok(())
}
}

pub(crate) fn map_operation_error(&self, error: io::Error) -> io::Error {
// Preserve transient database errors until they outlive the local safety margin.
self.ensure_operation_active().err().unwrap_or(error)
}

pub(crate) fn mark_lost(&self) {
if self
.lease_lost
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return;
}

// Run any installed containment handler before publishing lease loss.
if let Some(handler) = self.loss_handler.lock().expect("lock").take() {
handler();
}
self.loss_sender.send_replace(true);
}

pub(crate) fn set_loss_handler(&self, handler: LeaseLossHandler) {
let mut locked_handler = self.loss_handler.lock().expect("lock");
if self.is_lost() {
drop(locked_handler);
handler();
} else {
*locked_handler = Some(handler);
}
}

pub(crate) async fn wait_for_loss(self: Arc<Self>) {
let mut receiver = self.loss_sender.subscribe();
let _ = receiver.wait_for(|lost| *lost).await;
}
}

pub(crate) fn lease_lost_error() -> io::Error {
io::Error::new(io::ErrorKind::PermissionDenied, "PostgreSQL node lease was lost")
}

#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicBool, Ordering};

use super::*;

#[test]
fn expired_operation_marks_loss_before_returning_error() {
let lease = NodeLease::new().unwrap();
let handler_ran = Arc::new(AtomicBool::new(false));
let handler_ran_ref = Arc::clone(&handler_ran);
lease.set_loss_handler(Box::new(move || {
handler_ran_ref.store(true, Ordering::Release);
}));
*lease.last_confirmed_renewal.lock().unwrap() =
Instant::now() - NODE_LEASE_RENEWAL_DEADLINE;

let error = lease.map_operation_error(io::Error::from(io::ErrorKind::Other));

assert_eq!(error.kind(), io::ErrorKind::PermissionDenied);
assert!(lease.is_lost());
assert!(handler_ran.load(Ordering::Acquire));
}

#[test]
fn confirmed_renewal_uses_attempt_time_and_does_not_regress() {
let lease = NodeLease::new().unwrap();
let renewal_started_at = Instant::now() - Duration::from_secs(1);
*lease.last_confirmed_renewal.lock().unwrap() = renewal_started_at - Duration::from_secs(1);

lease.record_renewal_started_at(renewal_started_at);
lease.record_renewal_started_at(renewal_started_at - Duration::from_secs(1));

assert_eq!(*lease.last_confirmed_renewal.lock().unwrap(), renewal_started_at);
}

#[tokio::test]
async fn expired_renewal_deadline_completes_immediately() {
let lease = NodeLease::new().unwrap();
*lease.last_confirmed_renewal.lock().unwrap() =
Instant::now() - NODE_LEASE_RENEWAL_DEADLINE;

tokio::time::timeout(Duration::from_secs(1), lease.wait_for_renewal_deadline())
.await
.unwrap();
}
}
33 changes: 26 additions & 7 deletions src/io/postgres_store/migrations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,16 +6,35 @@
// accordance with one or both of these licenses.

use lightning::io;
use tokio_postgres::Client;
use tokio_postgres::Transaction;

pub(super) async fn migrate_schema(
_client: &Client, _kv_table_name: &str, from_version: u16, to_version: u16,
transaction: &Transaction<'_>, kv_table_name: &str, mut from_version: u16, to_version: u16,
) -> io::Result<()> {
assert!(from_version < to_version);
// Future migrations go here, e.g.:
// if from_version == 1 && to_version >= 2 {
// migrate_v1_to_v2(client, kv_table_name).await?;
// from_version = 2;
// }
if from_version == 1 && to_version >= 2 {
migrate_v1_to_v2(transaction, kv_table_name).await?;
from_version = 2;
}

if from_version != to_version {
return Err(io::Error::new(
io::ErrorKind::Other,
format!("No PostgreSQL schema migration from version {from_version} to {to_version}"),
));
}
Ok(())
}

async fn migrate_v1_to_v2(transaction: &Transaction<'_>, kv_table_name: &str) -> io::Result<()> {
// Schema v2 marks the transition from the legacy session advisory lock to fenced node leases.
// Older releases reject this version instead of reopening the store without lease fencing.
let sql = format!("COMMENT ON TABLE {kv_table_name} IS '2'");
transaction.execute(&sql, &[]).await.map_err(|e| {
io::Error::new(
io::ErrorKind::Other,
format!("Failed to set PostgreSQL schema version 2: {e}"),
)
})?;
Ok(())
}
Loading
Loading