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
15 changes: 14 additions & 1 deletion 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 Down Expand Up @@ -721,7 +724,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 @@ -1217,6 +1226,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 Down Expand Up @@ -2335,6 +2347,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();
}
}
6 changes: 3 additions & 3 deletions src/io/postgres_store/migrations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,15 +6,15 @@
// 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, 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?;
// migrate_v1_to_v2(transaction, kv_table_name).await?;
// from_version = 2;
// }
Ok(())
Expand Down
Loading
Loading