Skip to content

Replace PostgreSQL advisory locking with fenced leases - #1111

Open
joostjager wants to merge 3 commits into
lightningdevkit:mainfrom
joostjager:postgres-leases
Open

joostjager wants to merge 3 commits into
lightningdevkit:mainfrom
joostjager:postgres-leases

Conversation

@joostjager

@joostjager joostjager commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

Replace the temporary PostgreSQL advisory lock introduced in #1012 with a renewable lease in a companion table. The advisory-lock mechanism has not been released, so no migration is needed.

Acquire the lease atomically with schema setup, before loading persisted state. Validate ownership and expiry and renew the lease within each KV mutation transaction, so stale writes and deletes cannot commit. This replaces the separate lock checks before and after operations.

Renew idle leases in the background and release them by owner ID on drop. Lease loss panics in the calling task for mutations or in the background renewal task; background renewal errors and timeouts also panic. No recovery API or node lifecycle changes are introduced.

Tests cover setup rollback, contention, idle renewal, renewal timeouts, stale mutations, and release after takeover. Also verified contention, takeover after a crash, graceful restart, and stale-owner failure using two real LDK Server instances sharing a database.

@ldk-reviews-bot

ldk-reviews-bot commented Sep 21, 2026 •

Copy link
Copy Markdown

👋 Thanks for assigning @benthecarman as a reviewer!
I'll wait for their review and will help manage the review process.
Once they submit their review, I'll check if a second reviewer would be helpful.

@tnull tnull left a comment •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This needs a rebase for CI to run. Let me know when fully ready for review.

@tnull tnull added this to the 0.8 milestone Sep 22, 2026
@joostjager

Copy link
Copy Markdown
Contributor Author

If we indeed add this on the 0.8 milestone, I'll remove the migration because the temp. locking solution won't be in a release.

@tnull

tnull commented Sep 23, 2026

Copy link
Copy Markdown
Collaborator

If we indeed add this on the 0.8 milestone, I'll remove the migration because the temp. locking solution won't be in a release.

Do it.

Create the KV table, record its schema version, and create the listing
index in one transaction so failed initialization rolls back schema
changes. Keep database creation outside the transaction and preserve the
session advisory lock for the store's lifetime.

Add a regression test that forces index creation to fail and verifies
the schema version update is rolled back.
Move pooled connection acquisition and query retry handling for writes
and removals into execute_mutation. Preserve advisory lock checks,
per-key write ordering, SQL, and error mapping so lease enforcement can
be added to the shared helper separately.
Replace the temporary session advisory lock with an expiring lease in a
companion table. Acquire it after schema setup and validate and renew it
after each KV mutation in the same transaction, so stale writes and
deletes roll back before commit.

Renew idle leases in the background and release them by owner ID on
drop. Explicitly panic on lease loss in the calling task or background
renewal task, and on background failure or timeout, without adding
recovery APIs.

Cover rejected setup rollback, repeated idle renewal, renewal timeouts,
stale writes and deletes, and releasing an old owner after takeover.
@joostjager
joostjager marked this pull request as ready for review September 23, 2026 12:21
@joostjager
joostjager removed the request for review from TheBlueMatt September 23, 2026 12:23
@joostjager

Copy link
Copy Markdown
Contributor Author

I’ve removed the migration. This now directly replaces advisory locks with the lease table.

I’m particularly happy to see all the scattered checks before and after operations to confirm we still hold the lock disappear. PostgreSQL now checks lease validity within the transaction.

I had a good back-and-forth with AI to minimize the diff. I think it’s looking pretty good now.

@joostjager
joostjager requested a review from tnull September 23, 2026 12:24

@tnull tnull left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, yet to do a very detailed review. Also tagging @benthecarman as a secondary reviewer as he did the first approach.

.execute(&update_sql, &[&self.lease_owner_id.as_slice(), &lease_duration_secs])
.await?;
if updated != 1 {
panic!("PostgreSQL node lease was lost; continuing may corrupt node state");

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This probably should be std::process:abort as panicking the tokio task won't abort the whole process, but just have the task return a JoinError in the end.

@joostjager joostjager Sep 24, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The previous advisory-lock implementation used assert! rather than abort, and LDK Server sets panic = "abort" for both dev and release builds. But definitely seems safer to use std::process:abort, will change.

Note that not panicking wouldn't be a data consistency issue, because the consistency is guarded in each transaction.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Although, it seems in other places in ldk-node, it's not an explicit process abort?

@tnull
tnull requested a review from benthecarman September 24, 2026 12:41
);
}

async fn execute_mutation<F: FnOnce(PgError) -> io::Error>(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hmm, not the biggest fan of extracting helpers that are less than 5-7 lines of code, as they simply tend to increase clutter and make it harder to see what's going on. Is there a particular benefit from this?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

In the later commit, it is expanded and no longer 5-7 lines.

Comment thread src/builder.rs
/// This acquires an exclusive lease for the selected KV table before reading persisted node
/// state. Nodes may share a database when each node identity uses a distinct `kv_table_name`.
/// Mutations panic on detected lease loss. Failed or timed-out renewals panic in the background
/// renewal task. Node recovery is not handled automatically.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good to mention it's not handled automatically, but should we at least give the user half a sentence of guidance what this effectively means, i.e., what they are supposed to do if it panics?

Comment thread tests/common/mod.rs

/// Drops the given table from the `ldk_db` database, ignoring the case where the database doesn't
/// exist yet. Used to ensure a clean slate before and after Postgres-backed tests.
/// Also drops the companion lease table.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do we want/need these docs in test-only code to begin with?

Comment thread src/io/postgres_store/mod.rs
Ok(Self { inner, next_write_version, internal_runtime: Some(internal_runtime) })

let inner_ref = Arc::clone(&inner);
let lease_renewal_task = internal_runtime.spawn(async move {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Seems we need to make sure we hard panic here, too:

Codex:

  1. High — Renewal failure leaves the node running. The renewal task (/home/tnull/workspace/ldk-node-pr-1111-20260928-cPRuYZ/src/io/postgres_store/mod.rs:216) panics on failure, but nothing observes its handle during normal operation. Renewal stops while the node continues running; the
    process wide abort discussed earlier is still absent.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The downside of hard panic is that we can't cleanly unit test this behavior anymore, and need to spawn a child process to let it panic. I tend towards just documenting to set panic=abort. ldk-server already has that.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes we can, see for example

async fn catch_future_unwind<F: Future>(future: F) -> std::thread::Result<F::Output> {
let mut future = std::pin::pin!(future);
std::future::poll_fn(|cx| {
match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| future.as_mut().poll(cx))) {
Ok(std::task::Poll::Ready(output)) => std::task::Poll::Ready(Ok(output)),
Ok(std::task::Poll::Pending) => std::task::Poll::Pending,
Err(panic) => std::task::Poll::Ready(Err(panic)),
}
})
.await
}
async fn assert_invalid_write_fails<K: KVStore + RefUnwindSafe>(
kv_store: &K, primary_namespace: &str, secondary_namespace: &str, key: &str, data: Vec<u8>,
) {
let res = std::panic::catch_unwind(|| {
KVStore::write(kv_store, primary_namespace, secondary_namespace, key, data)
});
if let Ok(fut) = res {
if let Ok(write_res) = catch_future_unwind(fut).await {
assert!(write_res.is_err());
}
}
}

let mut interval = tokio::time::interval(NODE_LEASE_RENEWAL_INTERVAL);
loop {
// The first tick is immediate; later attempts follow the renewal interval.
interval.tick().await;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Default behavior is Burst - should we set something like:

 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);

let inner_ref = Arc::clone(&inner);
let lease_renewal_task = internal_runtime.spawn(async move {
let mut interval = tokio::time::interval(NODE_LEASE_RENEWAL_INTERVAL);
loop {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should this have a graceful shutdown path like:

 let mut interval = tokio::time::interval(NODE_LEASE_RENEWAL_INTERVAL);
 loop {
-    interval.tick().await;
+    tokio::select! {
+        biased;
+        _ = &mut shutdown_rx => break,
+        _ = interval.tick() => {}
+    }
+
     let renewal = async {
         let mut locked = inner_ref.locked_client().await?;
         // Existing renewal body...
     };
-    tokio::time::timeout(NODE_LEASE_RENEWAL_TIMEOUT, renewal)
-        .await 
-        .expect("PostgreSQL node lease renewal timed out")
-        .expect("Failed to renew PostgreSQL node lease");
+    tokio::select! {
+        biased;
+        _ = &mut shutdown_rx => break,
+        result = tokio::time::timeout(NODE_LEASE_RENEWAL_TIMEOUT, renewal) => {
+            result
+                .expect("PostgreSQL node lease renewal timed out")
+                .expect("Failed to renew PostgreSQL node lease");
+        }
+    }
 }
+
+let release_result =
+    tokio::time::timeout(NODE_LEASE_RELEASE_TIMEOUT, inner_ref.release_node_lease()).await;
+// Send release_result to a synchronous completion channel observed by Drop.


let runtime_handle = internal_runtime.handle().clone();
let inner = Arc::clone(&self.inner);
let _ = std::thread::spawn(move || {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think instead of spawning yet another thread here we'll rather want to first send a stop signal and handle graceful shutdown as one arm in the lease task loop, something like:

 impl Drop for PostgresStore {
     fn drop(&mut self) {
-        if let Some(internal_runtime) = self.internal_runtime.as_ref() {
-            let renewal_task = self.lease_renewal_task.take();
-            if let Some(task) = renewal_task.as_ref() {
-                task.abort();
-            }
-
-            let runtime_handle = internal_runtime.handle().clone();
-            let inner = Arc::clone(&self.inner);
-            let _ = std::thread::spawn(move || {
-                runtime_handle.block_on(async move {
-                    if let Some(task) = renewal_task {
-                        let _ = task.await;
-                    }
-
-                    let _ = tokio::time::timeout(
-                        NODE_LEASE_RELEASE_TIMEOUT,
-                        inner.release_node_lease(),
-                    )
-                    .await;
-                });
-            })
-            .join();
+        if let Some(sender) = self.lease_shutdown_sender.take() {
+            let _ = sender.send(());
+        }
+        if let Some(receiver) = self.lease_shutdown_complete.take() {
+            let wait = NODE_LEASE_RELEASE_TIMEOUT + Duration::from_secs(1);
+            let result = receiver.recv_timeout(wait);
+            if let Some(logger) = self.inner.logger.as_ref() {
+                match result {
+                    Ok(Ok(())) => {},
+                    Ok(Err(e)) => log_error!(logger, "Failed to release PostgreSQL node lease: {e}"),
+                    Err(e) => log_error!(logger, "PostgreSQL lease shutdown did not complete: {e}"),
+                }
+            }
+        }
+        if let Some(task) = self.lease_renewal_task.take() {
+            task.abort(); // Only matters if the completion wait failed or timed out.
         }

         if let Some(internal_runtime) = self.internal_runtime.take() {
             if let Ok(internal_runtime) = Arc::try_unwrap(internal_runtime) {
                 internal_runtime.shutdown_background();
             }
         }
     }
 }

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants