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 numberDiff line numberDiff 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 DownExpand 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 DownExpand 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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 numberDiff line numberDiff 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 DownExpand 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 DownExpand 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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 numberDiff line numberDiff 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 DownExpand 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 DownExpand 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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 numberDiff line numberDiff 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 DownExpand 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 DownExpand 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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 numberDiff line numberDiff 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 DownExpand 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 DownExpand 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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 numberDiff line numberDiff 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 DownExpand 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 DownExpand 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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 numberDiff line numberDiff 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 DownExpand 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 DownExpand 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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 numberDiff line numberDiff 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 DownExpand 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 DownExpand 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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 numberDiff line numberDiff 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