Replace PostgreSQL advisory locking with fenced leases - #1111
joostjager wants to merge 5 commits into
Conversation
|
👋 Thanks for assigning @benthecarman as a reviewer! |
1c20a0c to
a964bfc
Compare
|
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.
a964bfc to
8f19873
Compare
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.
8f19873 to
6532d1b
Compare
|
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. |
tnull
left a comment
There was a problem hiding this comment.
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"); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Although, it seems in other places in ldk-node, it's not an explicit process abort?
| ); | ||
| } | ||
|
|
||
| async fn execute_mutation<F: FnOnce(PgError) -> io::Error>( |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
In the later commit, it is expanded and no longer 5-7 lines.
| /// 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. |
There was a problem hiding this comment.
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?
|
|
||
| /// 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. |
There was a problem hiding this comment.
Do we want/need these docs in test-only code to begin with?
There was a problem hiding this comment.
Reduced and made comment more generic
| 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 { |
There was a problem hiding this comment.
Seems we need to make sure we hard panic here, too:
Codex:
- 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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Yes we can, see for example
Lines 172 to 195 in 96c1225
There was a problem hiding this comment.
But can this catch a process abort? The docs say:
This function only catches unwinding panics, not those that abort the process.
There was a problem hiding this comment.
I left it a panic! because of the unit tests, and also consistency with panics elsewhere in ldk-node. But open to different trade offs.
There was a problem hiding this comment.
Hmm, so it should def. be consistent between the different cases, and it seems we want to enforce process abortion consistently. You're indeed correct about the catch_unwind behavior, forgot about that, but that still doesn't change the core issue, IMO. So either we just leave this behavior untested or could do something where we switch to a catchable panic! under cfg(test) while doing a 'real' process exit in production?
There was a problem hiding this comment.
I see the benefit, but I’d prefer to keep normal panics for this PR. LDK Server already uses panic = "abort", and we document that requirement for direct LDK Node users. Without it, lease checks still prevent stale writes, though the process may remain partially running.
I’d rather address library-enforced process termination consistently in a separate change than introduce test-specific behavior just for lease loss.
| 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; |
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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.There was a problem hiding this comment.
Refactored the shutdown path
|
|
||
| let runtime_handle = internal_runtime.handle().clone(); | ||
| let inner = Arc::clone(&self.inner); | ||
| let _ = std::thread::spawn(move || { |
There was a problem hiding this comment.
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();
}
}
}
}Skip missed renewal ticks and let the renewal task release its lease when shutdown is requested. Bound the entire Drop wait with the existing release timeout, cancel unfinished work, and log shutdown failures without spawning an extra thread. Document the application panic-abort requirement and restart-based recovery. Cover shutdown with a blocked pool, verify that an old owner cannot release its replacement's lease, and adapt timeout-test cleanup to the renewal task owning lease release.
|
I plan to follow up with an LDK Server PR adding a small retry loop when the lease is already held. A deployment supervisor could also handle this by restarting the server after failed acquisition, but keeping the standby running and waiting for the lease inside LDK Server seems cleaner. |
Hmm, I think conceptually stuff like that (which easily will get bigger/more complicated) should really live in LDK Node. IMO, we'll want to keep LDK Server as close to a simple API wrapper as possible. |
cb2e065 to
0447a63
Compare
Yes, we can also put a small retry loop in LDK Node. The trade-off is that the synchronous build call would remain blocked while waiting for the lease. Keeping the retry outside that call makes it easier for the application to expose its standby status and control cancellation. We don’t currently have a way to cancel an ongoing build. A broader node lifecycle API could change that, but my immediate focus is failover for server deployments, which seems achievable with a small change today. I’m open to either location, provided we can keep that scope narrow. |
|
Regardless of where we put the retry loop, I don’t think it needs to block this PR. The priority here is replacing the temporary advisory locks with fenced leases. We can settle the retry behavior and its location in a follow-up. |
| transaction.execute(sql, params).await.map_err(err_map)?; | ||
| // Renew after the mutation to give the lease a later expiry. Rejection rolls back the | ||
| // mutation; success holds the lease row lock until commit. | ||
| self.renew_node_lease(&transaction).await.map_err(|e| { |
There was a problem hiding this comment.
This UPDATE holds the lease row lock until COMMIT, so every write across the store is serialized, even with the 10-connection pool. Could we use SELECT 1 FROM <lease> WHERE id = 1 AND owner_id = $1 AND expires_at > clock_timestamp() FOR SHARE here instead? Takeover and renewal still conflict with the share lock, so it fences the same way, and writes to different keys can run in parallel again. Keep it after the mutation so the lock order matches setup. The trade-off is that writes stop extending the lease, so only the background renewal keeps it alive.
There was a problem hiding this comment.
Sounds like a good idea: less writing and less locking, leaving renewal to the background task. I've pushed a fixup for that.
I also inlined renew_node_lease now that it only had one caller, which allowed some further simplification of the renewal and shutdown code.
There was a problem hiding this comment.
One catch with claude's FOR SHARE suggestion: PostgreSQL lets new share lockers through even while an UPDATE is waiting, so continuous writes can starve the renewal past its 10s timeout. (The new test aborts the renewal task for this reason.) Adding UNIQUE (owner_id) to the lease table and using FOR KEY SHARE here fixes it. The renewal only touches expires_at, so it no longer conflicts, but a takeover changes owner_id, which is now a key column, so it still blocks until our commit. I checked this in psql and with an updated test_postgres_store_lease.
One gotcha: CREATE TABLE IF NOT EXISTS won't add UNIQUE to an existing lease table, and without it FOR KEY SHARE doesn't block takeover. So drop any *_node_lease tables created by earlier revisions of this PR. cleanup_store should drop the lease table too, since leftover tables from earlier runs hit exactly this.
Validate each KV mutation with SELECT FOR SHARE after the mutation, allowing concurrent lease checks while fencing takeover until commit. Let only the background task extend the lease and cover shared-lock compatibility in the existing lease test. Inline renewal and construct its SQL once, use one shutdown select, and store the task handle directly. Handle lease loss, renewal errors, and timeouts with explicit panics. Check the affected-row count during acquisition and let the setup connection drop at function exit.
da2c7d1 to
49188c8
Compare
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.