diff --git a/alkstore-postgres/src/lib.rs b/alkstore-postgres/src/lib.rs index e750e94..5d0dc3c 100644 --- a/alkstore-postgres/src/lib.rs +++ b/alkstore-postgres/src/lib.rs @@ -33,6 +33,7 @@ //! scheduler/outbox); each task replaces the prior stubs. mod forwarder; +mod lock; mod notify; mod opts; mod queue; diff --git a/alkstore-postgres/src/lock.rs b/alkstore-postgres/src/lock.rs new file mode 100644 index 0000000..af98e6f --- /dev/null +++ b/alkstore-postgres/src/lock.rs @@ -0,0 +1,236 @@ +//! The named-lock mechanism's Postgres arm (`pg-engine-locks`) — the +//! substrate's lock-op semantics re-derived over the schema task's +//! locks table (the POC's lock probe pinned acquire/release/renew with +//! TTL and expiry re-acquisition directly): +//! +//! - `try_lock(name, owner, ttl)` — shared-namespace + non-empty +//! owner validation, the **duration guard (ADR-023 §2)** on `ttl <= +//! 0` at the entry (opaque `Err(Database)`, detail in the source +//! chain — an exclusion mechanism must not silently hand back a +//! lease the caller believes it holds), the closed-store check, +//! then the acquire op on a pool connection: opportunistic- +//! expiry-delete + insert-or-reacquire + read-back. `Ok(None)` = +//! held elsewhere (no-work is a value); `Ok(Some(..))` = the boxed +//! [`Lock`] handle. +//! - `Lock::renew(ttl)` — the same duration guard; the new **full** +//! TTL window from now (not additive); `false` = lost it (expired +//! and re-acquired elsewhere — exclusion lapses silently at TTL +//! expiry, no revocation event, ADR-008 §7). +//! - `Lock::release(self: Box)` — consuming (the RAII shape, +//! ADR-019 §1); deletes only the owner's row (`false` = already +//! expired/re-acquired — a no-op, not an error). No handle-local +//! "held" flag caching: the row-scoped delete is the truth-teller. +//! +//! No `lock_tx` exists — deliberate (core-contract.md's named-locks +//! section): acquisition is a separate auto-commit op; the TTL +//! discipline governs. There is no general expiry sweeper either: +//! expired rows hang until the next acquire's opportunistic delete +//! deletes them (the inherited substrate posture, documented) — the +//! `locks_expiry_idx` index keeps that delete cheap. + +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; + +use deadpool_postgres::Pool; + +use alkstore::{BoxedFuture, Lock, Result, validate_local_name, validate_shared_name}; + +use crate::resolution::now_unix; +use crate::schema::{QualifiedTable, tables}; +use crate::seam::{database_error, pg_error, pool_error}; + +/// The duration guard's message (ADR-023 §2's engine detail — rides +/// the source chain, `Database`'s Display is opaque). +fn duration_error(ttl: i64) -> alkstore::Error { + database_error(format!( + "ttl must be a positive duration (seconds), got {ttl}" + )) +} + +/// Qualified table reference for this engine-owned schema. +fn table(schema: &str, name: &str) -> String { + QualifiedTable { schema, name }.to_string() +} + +/// The acquire op (see the module docs): the opportunistic +/// expired-row delete for this name, then the insert-or-reacquire +/// (the PK on `name` prevents dual acquisition — `ON CONFLICT DO +/// NOTHING` keeps the incumbent row exactly as-is, so a same-owner +/// re-acquire keeps the original row with its TTL *not* refreshed, +/// the substrate's `INSERT OR IGNORE` shape), then the read-back — +/// the holder column decides, not the insert's row count (the D-4 +/// swallow's lesson: never trust a write result for ownership). +/// +/// One statement sequence on one pool connection's implicit +/// transaction (auto-committed as a unit on return-to-pool; the +/// expiry-delete + insert + read-back race window is closed by the +/// PK constraint — a concurrent grant of the same name either +/// lands first or reads as foreign). +async fn lock_acquire( + client: &tokio_postgres::Client, + schema: &str, + name: &str, + owner: &str, + ttl: i64, +) -> Result { + let now = now_unix(); + client + .execute( + &format!( + "DELETE FROM {locks} WHERE name = $1 AND expires_at <= $2", + locks = table(schema, tables::LOCKS) + ), + &[&name, &now], + ) + .await + .map_err(pg_error)?; + client + .execute( + &format!( + "INSERT INTO {locks} (name, owner, expires_at) + VALUES ($1, $2, $3::bigint + $4::bigint) + ON CONFLICT (name) DO NOTHING", + locks = table(schema, tables::LOCKS) + ), + &[&name, &owner, &now, &ttl], + ) + .await + .map_err(pg_error)?; + let holder: Option = client + .query_opt( + &format!( + "SELECT owner FROM {locks} WHERE name = $1", + locks = table(schema, tables::LOCKS) + ), + &[&name], + ) + .await + .map_err(pg_error)? + .map(|row| row.get(0)); + Ok(holder.as_deref() == Some(owner)) +} + +/// The acquire path (see the module docs): validation, duration +/// guard, closed-store check, then the acquire op on a pool +/// connection. `Ok(None)` is the held-by-someone-else arm. +pub(crate) fn try_lock( + pool: Pool, + schema: String, + closed: Arc, + name: String, + owner: String, + ttl: i64, +) -> BoxedFuture<'static, Result>>> { + if let Err(e) = validate_shared_name(&name) { + return Box::pin(async move { Err(e) }); + } + if let Err(e) = validate_local_name(&owner) { + return Box::pin(async move { Err(e) }); + } + if ttl <= 0 { + return Box::pin(async move { Err(duration_error(ttl)) }); + } + if closed.load(Ordering::Acquire) { + return Box::pin(async { Err(database_error("the store is closed")) }); + } + Box::pin(async move { + let client = pool.get().await.map_err(pool_error)?; + let granted = lock_acquire(&client, &schema, &name, &owner, ttl).await?; + Ok(granted.then(|| { + Box::new(PgLockHandle { + name, + owner, + pool, + schema, + closed, + }) as Box + })) + }) +} + +/// The held named-lock handle (see the module docs). [`Lock::release`] +/// consumes it; `renew` is the repeatable `&self` case. +pub struct PgLockHandle { + name: String, + owner: String, + pool: Pool, + schema: String, + closed: Arc, +} + +impl std::fmt::Debug for PgLockHandle { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("PgLockHandle") + .field("name", &self.name) + .finish_non_exhaustive() + } +} + +impl Lock for PgLockHandle { + fn name(&self) -> &str { + &self.name + } + + fn renew<'a>(&'a self, ttl: i64) -> BoxedFuture<'a, Result> { + if ttl <= 0 { + return Box::pin(async move { Err(duration_error(ttl)) }); + } + if self.closed.load(Ordering::Acquire) { + return Box::pin(async { Err(database_error("the store is closed")) }); + } + let name = self.name.clone(); + let owner = self.owner.clone(); + let pool = self.pool.clone(); + let schema = self.schema.clone(); + Box::pin(async move { + let now = now_unix(); + let client = pool.get().await.map_err(pool_error)?; + // The owner-scoped update: 0 rows = expired-and-gone or + // re-acquired by another owner — `false`, the lost-it arm. + let n = client + .execute( + &format!( + "UPDATE {locks} SET expires_at = $3::bigint + $4::bigint + WHERE name = $1 AND owner = $2", + locks = table(&schema, tables::LOCKS) + ), + &[&name, &owner, &now, &ttl], + ) + .await + .map_err(pg_error)?; + Ok(n > 0) + }) + } + + fn release(self: Box) -> BoxedFuture<'static, Result> { + let PgLockHandle { + name, + owner, + pool, + schema, + closed, + } = *self; + Box::pin(async move { + if closed.load(Ordering::Acquire) { + return Err(database_error("the store is closed")); + } + let client = pool.get().await.map_err(pool_error)?; + // The owner-scoped delete — row-scoped truth-teller: `true` + // = the row was there and is gone by this owner's hand + // (including the stale row past its TTL, which no sweeper + // has touched); `false` = row-less or foreign (the + // no-op arm). + let n = client + .execute( + &format!( + "DELETE FROM {locks} WHERE name = $1 AND owner = $2", + locks = table(&schema, tables::LOCKS) + ), + &[&name, &owner], + ) + .await + .map_err(pg_error)?; + Ok(n > 0) + }) + } +} diff --git a/alkstore-postgres/src/store.rs b/alkstore-postgres/src/store.rs index f7fc30b..a56705d 100644 --- a/alkstore-postgres/src/store.rs +++ b/alkstore-postgres/src/store.rs @@ -26,9 +26,9 @@ //! the trait stubs below are the wave-3 posture — they return //! `Err(Database("… wiring lands with the … task"))` until the //! mechanism tasks replace them. `notify`/`listen` are wired (the -//! notify-listen task's), `stream` (the streams task's), `queue` -//! (the queues task's) — the wake contract's pg arm rides the -//! forwarder. +//! notify-listen task's), `stream` (the streams task's), `queue` (the +//! queues task's), `try_lock` (the locks task's) — the wake +//! contract's pg arm rides the forwarder. use alkstore::Store; use std::sync::Arc; @@ -314,14 +314,18 @@ impl Store for PgStore { fn try_lock<'a>( &'a self, - _name: &str, - _owner: &str, - _ttl: i64, + name: &str, + owner: &str, + ttl: i64, ) -> alkstore::BoxedFuture<'a, alkstore::Result>>> { - if let Err(e) = self.closed_check() { - return Box::pin(async move { Err(e) }); - } - Box::pin(async { Err(stub("lock wiring lands with the locks task")) }) + crate::lock::try_lock( + self.pool.clone(), + self.schema.clone(), + self.closed.clone(), + name.to_string(), + owner.to_string(), + ttl, + ) } fn schedule<'a>( @@ -370,6 +374,9 @@ mod open_tests; #[cfg(test)] mod queue_tests; +#[cfg(test)] +mod lock_tests; + #[cfg(test)] mod notify_tests; diff --git a/alkstore-postgres/src/store/lock_tests.rs b/alkstore-postgres/src/store/lock_tests.rs new file mode 100644 index 0000000..dc1ea5c --- /dev/null +++ b/alkstore-postgres/src/store/lock_tests.rs @@ -0,0 +1,644 @@ +//! The named-lock mechanism's acceptance tests (`pg-engine-locks`): +//! exclusive grant over the harness server with the loser's `Ok(None)` +//! arm, release re-enabling acquisition, the duration guards (ADR-023 +//! §2 — `ttl <= 0` rejected opaque on `try_lock` and `renew`, the +//! engine guard's detail in the source chain), `renew`'s +//! full-window-from-now semantics and lapsed-holder refusal, the +//! inherited same-owner re-acquire shape (granted, TTL not refreshed — +//! matched to the SQLite arm deliberately for cross-engine +//! equivalence), `release`'s owner-scoped delete + both boolean arms +//! (the stale row deletes `true`; after a foreign re-acquire `false`), +//! the silent TTL lapse with lapse-then-reacquire (ADR-008 §7's +//! guarantee row), the engine-external lapse probe (a second store +//! instance on the same schema acquires after the TTL lapses — the +//! verification backlog's pg re-acquisition pin at engine level), the +//! entry-point validation posture, and the closed-store fail-closed +//! shape on all lock ops. ADR-019 §1–§2; ADR-023 §2. +//! +//! Harness convention (the schema/open tasks'): connection settings +//! ride the environment (`ALKSTORE_PG_HOST/PORT/USER/PASSWORD/DB`), +//! never hardcoded; tests without a reachable server skip cleanly so +//! the workspace gates stay green server-less. Isolation is a fresh +//! unique schema per test. The short-window lapse steps backdate the +//! lock row's expiry directly through the raw admin connection (the +//! queue tests' honest cross-connection probe shape) instead of +//! sleeping out multi-second windows. + +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +use alkstore::{Error, Store}; + +use crate::store::{PgStore, open_store}; + +const ENV_HOST: &str = "ALKSTORE_PG_HOST"; +const ENV_PORT: &str = "ALKSTORE_PG_PORT"; +const ENV_USER: &str = "ALKSTORE_PG_USER"; +const ENV_PASSWORD: &str = "ALKSTORE_PG_PASSWORD"; +const ENV_DB: &str = "ALKSTORE_PG_DB"; + +fn harness_dsn() -> Option { + let host = std::env::var(ENV_HOST).ok()?; + let port: u16 = std::env::var(ENV_PORT).ok()?.parse().ok()?; + let user = std::env::var(ENV_USER).ok()?; + let password = std::env::var(ENV_PASSWORD).ok()?; + let db = std::env::var(ENV_DB).unwrap_or_else(|_| "postgres".to_string()); + Some(format!( + "host={host} port={port} user={user} password={password} dbname={db}" + )) +} + +fn instance_namer(tag: &str) -> impl Fn() -> String + use<'_> { + let counter = AtomicU64::new(0); + move || { + format!( + "{tag}_{}_{}_{}", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos(), + counter.fetch_add(1, Ordering::SeqCst), + ) + } +} + +/// A raw admin client for cross-connection probes (row backdating) +/// — separate from the store's pool. +async fn harness_client() -> Option { + let (client, connection) = tokio_postgres::connect(&harness_dsn()?, tokio_postgres::NoTls) + .await + .ok()?; + tokio::spawn(async move { + let _ = connection.await; + }); + Some(client) +} + +async fn drop_schema(client: &tokio_postgres::Client, schema: &str) { + if schema == crate::schema::DEFAULT_SCHEMA { + return; + } + let sql = format!( + "DROP SCHEMA IF EXISTS {} CASCADE", + crate::schema::quote_identifier(schema) + ); + let _ = client.batch_execute(&sql).await; +} + +fn test_opts(schema: &str) -> crate::opts::PgOpts { + crate::opts::PgOpts { + schema: schema.to_string(), + ..crate::opts::PgOpts::default() + } +} + +/// Open a store on the harness server with a fresh unique schema — +/// `None` = no harness server (the skip posture). The returned store +/// closes on `Fixture` drop; each test tears its schema down explicitly +/// at the end (the queue tests' convention). +struct Fixture { + store: PgStore, + admin: tokio_postgres::Client, + schema: String, +} + +impl Fixture { + async fn open(namer: &impl Fn() -> String) -> Option { + let dsn = harness_dsn()?; + let schema = namer(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let admin = harness_client().await?; + Some(Fixture { + store, + admin, + schema, + }) + } +} + +impl Drop for Fixture { + fn drop(&mut self) { + self.store.close(); + } +} + +fn now_unix() -> i64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_secs() as i64 +} + +/// Read the lock row's `expires_at` through the admin connection (the +/// cross-connection probe — the expiry column is what the TTL tests +/// assert and backdate). +async fn lock_expires_at(admin: &tokio_postgres::Client, schema: &str, name: &str) -> Option { + admin + .query_opt( + &format!( + "SELECT expires_at FROM {} WHERE name = $1", + crate::schema::QualifiedTable { + schema, + name: crate::schema::tables::LOCKS, + } + ), + &[&name], + ) + .await + .unwrap() + .map(|row| row.get(0)) +} + +/// Backdate a lock row's expiry past the wall clock (the honest +/// cross-connection probe for lapse windows — a fresh acquire would +/// carry a future expiry, so the admin UPDATE backdates it: +/// schema-qualified, typed binds). +async fn backdate_expiry(admin: &tokio_postgres::Client, schema: &str, name: &str) { + admin + .execute( + &format!( + "UPDATE {} SET expires_at = $2 WHERE name = $1", + crate::schema::QualifiedTable { + schema, + name: crate::schema::tables::LOCKS, + } + ), + &[&name, &(now_unix() - 1)], + ) + .await + .unwrap(); +} + +/// Acceptance: exclusive grant — the second acquirer gets `Ok(None)` +/// (no-work is a value, not an error), a different lock name is +/// independent, and release makes it acquirable again. The same-owner +/// re-acquire is the inherited shape: granted (`Some`), TTL **not** +/// refreshed — the SQLite arm's pinned disposition, matched here for +/// cross-engine equivalence (`ON CONFLICT DO NOTHING` keeps the +/// incumbent row exactly as the substrate's `INSERT OR IGNORE` does); +/// `renew` is the refresh path. +#[tokio::test(flavor = "multi_thread")] +async fn try_lock_grants_exclusively_and_release_reopens() { + let Some(fx) = Fixture::open(&instance_namer("lk")).await else { + eprintln!("skip: no harness server"); + return; + }; + + let a = fx.store.try_lock("leader", "a", 300).await.unwrap(); + let lock = a.expect("first acquirer holds it"); + assert_eq!(lock.name(), "leader"); + let grant_expires = lock_expires_at(&fx.admin, &fx.schema, "leader") + .await + .expect("the grant stamped the row"); + + assert!( + fx.store + .try_lock("leader", "b", 300) + .await + .unwrap() + .is_none(), + "someone holds it: Ok(None), never an error" + ); + + let reacquired = fx + .store + .try_lock("leader", "a", 300) + .await + .unwrap() + .expect("same-owner re-acquire is granted"); + let reacquire_expires = lock_expires_at(&fx.admin, &fx.schema, "leader") + .await + .expect("the row survives"); + assert_eq!( + reacquire_expires, grant_expires, + "re-acquire must not refresh the TTL (the SQLite arm's inherited \ + shape, matched): the row is the original insert" + ); + assert!( + fx.store + .try_lock("leader", "b", 300) + .await + .unwrap() + .is_none(), + "still exclusive after the no-refresh re-acquire" + ); + + assert!( + fx.store + .try_lock("other-lock", "a", 300) + .await + .unwrap() + .is_some(), + "a different lock name is independent" + ); + + assert!(lock.release().await.unwrap(), "release reports it was held"); + assert!( + fx.store + .try_lock("leader", "b", 300) + .await + .unwrap() + .is_some(), + "release makes it acquirable again" + ); + drop(reacquired); + + drop_schema(&fx.admin, &fx.schema).await; +} + +/// Acceptance: the duration guards (ADR-023 §2) — `try_lock(ttl <= 0)` +/// and `renew(ttl <= 0)` are opaque `Err` (`Database`), never a +/// granted lock, never a panic; the engine guard's detail +/// ("ttl must be a positive duration") rides the source chain +/// (`#[source]`, ADR-008 §5). The guard fires before any round trip: +/// no lock row exists after a rejected acquire. +#[tokio::test(flavor = "multi_thread")] +async fn duration_guards_reject_non_positive_ttl_opaquely() { + let Some(fx) = Fixture::open(&instance_namer("lk")).await else { + eprintln!("skip: no harness server"); + return; + }; + + for ttl in [0, -1] { + let err = match fx.store.try_lock("guarded", "a", ttl).await { + Ok(granted) => panic!( + "try_lock(ttl = {ttl}) must be a rejection, granted: {}", + granted.is_some() + ), + Err(e) => e, + }; + assert!( + matches!(err, Error::Database(_)), + "opaque Database shape, got: {err:?}" + ); + assert!( + !matches!(err, Error::InvalidName { .. } | Error::ReservedName { .. }), + "the duration rule must not be misreported as a name error: {err:?}" + ); + let detail = std::error::Error::source(&err) + .map(|s| s.to_string()) + .unwrap_or_default(); + assert!( + detail.contains("ttl must be a positive duration (seconds)"), + "the guard's detail rides the source chain, got: {detail:?}" + ); + } + assert!( + lock_expires_at(&fx.admin, &fx.schema, "guarded") + .await + .is_none(), + "the guard fires before the acquire: no lock row was created" + ); + + let lock = fx + .store + .try_lock("renew-guard", "a", 300) + .await + .unwrap() + .unwrap(); + for ttl in [0, -1] { + let err = match lock.renew(ttl).await { + Ok(landed) => panic!("renew({ttl}) must be a rejection, landed: {landed}"), + Err(e) => e, + }; + assert!(matches!(err, Error::Database(_)), "opaque, got: {err:?}"); + let detail = std::error::Error::source(&err) + .map(|s| s.to_string()) + .unwrap_or_default(); + assert!( + detail.contains("ttl must be a positive duration (seconds)"), + "the renew guard's detail rides the source chain too, got: {detail:?}" + ); + } + assert!(lock.renew(600).await.unwrap(), "a positive ttl renews"); + + let expires = lock_expires_at(&fx.admin, &fx.schema, "renew-guard") + .await + .unwrap(); + assert!(expires > now_unix(), "the renew stamped a future expiry"); + + drop_schema(&fx.admin, &fx.schema).await; +} + +/// Acceptance: `renew` is a new full TTL window from now (not +/// additive); a lapsed holder's renew is `false` (expired and +/// re-acquired elsewhere — exclusion lapses silently, ADR-008 §7). +#[tokio::test(flavor = "multi_thread")] +async fn renew_is_full_window_and_refuses_lapsed_holders() { + let Some(fx) = Fixture::open(&instance_namer("lk")).await else { + eprintln!("skip: no harness server"); + return; + }; + + let lock = fx + .store + .try_lock("leader", "a", 300) + .await + .unwrap() + .unwrap(); + assert!(lock.renew(600).await.unwrap()); + + let before = now_unix(); + let expires = lock_expires_at(&fx.admin, &fx.schema, "leader") + .await + .expect("row present"); + assert!( + ((before + 600 - 3)..=(before + 600 + 3)).contains(&expires), + "renew sets a full window from now (not additive): \ + {expires} vs ~{before}+600" + ); + + // Lapse + re-acquire, the lapse-then-reacquire shape: backdate the + // expiry (the honest cross-connection probe at second precision), + // re-acquire as a foreign owner, then the stale holder's renew is + // false (lost it). + backdate_expiry(&fx.admin, &fx.schema, "leader").await; + let _other = fx + .store + .try_lock("leader", "b", 1) + .await + .unwrap() + .expect("re-acquired after expiry"); + assert!(!lock.renew(300).await.unwrap(), "lost it"); + assert!( + !lock.release().await.unwrap(), + "the stale holder's release is the no-op arm after the foreign re-acquire" + ); + + drop_schema(&fx.admin, &fx.schema).await; +} + +/// Acceptance: `release` deletes only the owner's row (the +/// `(name, owner)` predicate — a foreign owner string deletes +/// nothing), and the post-lapse release arms (the SQLite task's pinned +/// shape, matched): immediately after a lapse (the stale row still +/// present — no general expiry sweeper touched it) the stale holder's +/// release deletes its own stale row (`true`) and the lock is +/// immediately acquirable; only a *removed* or *foreign-held* row +/// releases as `false` (the no-op arm). +#[tokio::test(flavor = "multi_thread")] +async fn release_owner_scoped_with_both_post_lapse_arms() { + let Some(fx) = Fixture::open(&instance_namer("lk")).await else { + eprintln!("skip: no harness server"); + return; + }; + + let lock = fx + .store + .try_lock("leader", "a", 300) + .await + .unwrap() + .unwrap(); + + // The row survives a foreign owner's release attempt (the + // predicate is (name, owner), probed raw). + let deleted = fx + .admin + .execute( + &format!( + "DELETE FROM {} WHERE name = $1 AND owner = $2", + crate::schema::QualifiedTable { + schema: &fx.schema, + name: crate::schema::tables::LOCKS, + } + ), + &[&"leader", &"nobody"], + ) + .await + .unwrap(); + assert_eq!(deleted, 0, "a foreign owner string must not delete the row"); + + // The stale-row arm: lapse (backdate), the row still present (no + // sweeper), the stale holder's release deletes its own stale row + // → `true`; the lock is immediately acquirable again. + backdate_expiry(&fx.admin, &fx.schema, "leader").await; + assert!( + lock_expires_at(&fx.admin, &fx.schema, "leader") + .await + .is_some(), + "the stale row persists (no general expiry sweeper — the \ + next acquire's opportunistic delete is the mechanism)" + ); + assert!(lock.release().await.unwrap(), "the stale row deletes"); + assert!( + fx.store + .try_lock("leader", "b", 300) + .await + .unwrap() + .is_some(), + "acquirable right after the stale holder's release" + ); + + // The row-less arm: releasing when nothing is there is `false`, a + // no-op, not an error. + assert!(fx.store.try_lock("l2", "a", 300).await.unwrap().is_some()); + // (the owner `a`'s second handle is the fresh grant; release it.) + let fresh = fx.store.try_lock("l2", "b", 300).await.unwrap(); + assert!(fresh.is_none(), "held by a — exclusive"); + let held = fx.store.try_lock("l2", "a", 300).await.unwrap().unwrap(); + assert!(held.release().await.unwrap(), "held: release deletes"); + let after = fx.store.try_lock("l2", "a", 300).await.unwrap().unwrap(); + assert!(after.release().await.unwrap(), "re-granted, released once"); + // Double release: the first consumed the row — the second is the + // no-op arm (`false`) — released on the fresh grant above; probe + // the consumed-row shape via a foreign-backed delete instead: + let deleted = fx + .admin + .execute( + &format!( + "DELETE FROM {} WHERE name = $1 AND owner = $2", + crate::schema::QualifiedTable { + schema: &fx.schema, + name: crate::schema::tables::LOCKS, + } + ), + &[&"l2", &"a"], + ) + .await + .unwrap(); + assert_eq!(deleted, 0, "the row is gone: nothing left to release"); + + drop_schema(&fx.admin, &fx.schema).await; +} + +/// Acceptance: silent lapse (ADR-008 §7's guarantee row) — no +/// revocation event; the next acquire's opportunistic expired-row +/// delete is the lapse mechanism, pinned lapse-then-reacquire: after +/// expiry another owner acquires while the stale holder's ops +/// disown (`renew`/`release` `false`, never an error). +#[tokio::test(flavor = "multi_thread")] +async fn ttl_lapses_silently_and_another_owner_acquires() { + let Some(fx) = Fixture::open(&instance_namer("lk")).await else { + eprintln!("skip: no harness server"); + return; + }; + + let lock = fx + .store + .try_lock("leader", "a", 300) + .await + .unwrap() + .unwrap(); + backdate_expiry(&fx.admin, &fx.schema, "leader").await; + + // Lapse is silent: the stale holder's ops report through values, + // never an error or a revocation event. The foreign acquire here + // also exercises the opportunistic expired-row delete (the + // expired row goes first, then b's insert lands). + assert!( + fx.store + .try_lock("leader", "b", 300) + .await + .unwrap() + .is_some(), + "after expiry another owner acquires (the expired-row delete is the mechanism)" + ); + assert!( + !lock.renew(300).await.unwrap(), + "the stale holder's renew is false: expired and re-acquired elsewhere" + ); + assert!( + !lock.release().await.unwrap(), + "the stale holder's release is the no-op arm after the foreign re-acquire" + ); + assert!( + fx.store + .try_lock("leader", "c", 300) + .await + .unwrap() + .is_none(), + "b's fresh acquisition still excludes others" + ); + + drop_schema(&fx.admin, &fx.schema).await; +} + +/// Acceptance: the lapse is visible across engine instances — a second +/// `open` on the same schema (the engine-external probe) acquires +/// after the TTL lapses (the verification backlog's +/// re-acquire-after-expiry pin on Postgres at engine level; the +/// contract-suite row remains for the suite itself). The polling loop +/// is bounded (the harness convention — no tight-race assertion). +#[tokio::test(flavor = "multi_thread")] +async fn re_acquired_after_expiry_from_a_second_engine_instance() { + let Some(fx) = Fixture::open(&instance_namer("lk")).await else { + eprintln!("skip: no harness server"); + return; + }; + let dsn = harness_dsn().expect("checked by Fixture::open"); + + let lock = fx.store.try_lock("leader", "a", 1).await.unwrap().unwrap(); + backdate_expiry(&fx.admin, &fx.schema, "leader").await; + + let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + let other = loop { + assert!( + tokio::time::Instant::now() < deadline, + "the lock never lapsed for the second instance" + ); + let other = open_store(&dsn, test_opts(&fx.schema)).await.unwrap(); + match other.try_lock("leader", "b", 300).await.unwrap() { + Some(_) => break other, + None => other.close(), + } + }; + assert!(!lock.renew(300).await.unwrap(), "lapsed: renew refuses"); + assert!( + !lock.release().await.unwrap(), + "lapsed + foreign re-acquire: the no-op arm" + ); + other.close(); + + drop_schema(&fx.admin, &fx.schema).await; +} + +/// Acceptance: the entry-point validation posture — the lock name's +/// shared-namespace rules (empty → `InvalidName`; the reserved +/// `__alkstore_` prefix → `ReservedName`) and the owner's non-empty +/// rule (→ `InvalidName`), all at the entry (the residue-free rejects +/// shape — nothing landed on the server for a rejected acquire; the +/// row probe confirms the name space stays clean). +#[tokio::test(flavor = "multi_thread")] +async fn entry_points_validate_names() { + let Some(fx) = Fixture::open(&instance_namer("lk")).await else { + eprintln!("skip: no harness server"); + return; + }; + + let err = match fx.store.try_lock("", "a", 300).await { + Ok(_) => panic!("empty lock name must be rejected"), + Err(e) => e, + }; + assert!(matches!(err, Error::InvalidName { .. }), "got: {err:?}"); + let err = match fx.store.try_lock("__alkstore_leader", "a", 300).await { + Ok(_) => panic!("reserved lock name must be rejected"), + Err(e) => e, + }; + assert!(matches!(err, Error::ReservedName { .. }), "got: {err:?}"); + let err = match fx.store.try_lock("leader", "", 300).await { + Ok(_) => panic!("empty owner must be rejected"), + Err(e) => e, + }; + assert!(matches!(err, Error::InvalidName { .. }), "got: {err:?}"); + let err = match fx.store.try_lock("leader", " ", 300).await { + Ok(_) => panic!("whitespace-only owner must be rejected"), + Err(e) => e, + }; + assert!( + matches!(err, Error::InvalidName { .. }), + "whitespace-only owner is rejected too, got: {err:?}" + ); + // Validation rejects leave no residue: the acquire after them + // proceeds cleanly. + assert!( + fx.store + .try_lock("leader", "a", 300) + .await + .unwrap() + .is_some(), + "validation rejects never poisoned the name space" + ); + + drop_schema(&fx.admin, &fx.schema).await; +} + +/// Acceptance: the closed-store fail-closed posture — after `close()` +/// (or drop), `try_lock` and the held handle's `renew`/`release` fail +/// with opaque `Database`, never `Ok` values (the engine-wide shape). +#[tokio::test(flavor = "multi_thread")] +async fn closed_store_fails_closed_on_all_lock_ops() { + let Some(fx) = Fixture::open(&instance_namer("lk")).await else { + eprintln!("skip: no harness server"); + return; + }; + + let lock = fx + .store + .try_lock("leader", "a", 300) + .await + .unwrap() + .unwrap(); + fx.store.close(); + + let err = match fx.store.try_lock("leader", "b", 300).await { + Ok(granted) => panic!( + "post-close try_lock must fail closed, granted: {}", + granted.is_some() + ), + Err(e) => e, + }; + assert!(matches!(err, Error::Database(_)), "got: {err:?}"); + let err = match lock.renew(300).await { + Ok(_) => panic!("post-close renew must fail closed"), + Err(e) => e, + }; + assert!(matches!(err, Error::Database(_)), "got: {err:?}"); + let err = match lock.release().await { + Ok(released) => panic!("post-close release must fail closed, released: {released}"), + Err(e) => e, + }; + assert!(matches!(err, Error::Database(_)), "got: {err:?}"); + + drop_schema(&fx.admin, &fx.schema).await; +} diff --git a/alkstore-postgres/src/store/open_tests.rs b/alkstore-postgres/src/store/open_tests.rs index 3a3c2f8..a678241 100644 --- a/alkstore-postgres/src/store/open_tests.rs +++ b/alkstore-postgres/src/store/open_tests.rs @@ -620,9 +620,9 @@ async fn open_fails_database_on_unparseable_and_unreachable() { } } -/// Acceptance: the trait surface — `begin_tx`, `notify`, and -/// `listen` are wired (the seam + notify-listen tasks); the remaining -/// methods are the wave-3 stub posture — every unwired `Store` method +/// Acceptance: the trait surface — `begin_tx`, `notify`, `listen`, +/// `stream`, `queue`, and `try_lock` are wired; the remaining methods +/// are the wave-3 stub posture — every unwired `Store` method /// returns `Err(Database(… wiring lands with the … task))`; `with_tx` /// surfaces the wired `begin_tx` naturally; no panics. #[tokio::test(flavor = "multi_thread")] @@ -675,10 +675,17 @@ async fn store_trait_methods_are_wiring_stubs() { .enqueue(serde_json::json!({}), alkstore::EnqueueOpts::default()) .await .unwrap(); + // The lock constructor is wired (the locks task's): a valid + // acquire grants — full behavior the lock tests' scope. + let lock = store + .try_lock("stub_l_probe", "stub_owner", 60) + .await + .unwrap() + .expect("the wired try_lock grants"); + assert_eq!(lock.name(), "stub_l_probe"); + assert!(lock.release().await.unwrap(), "the wired release deletes"); let err = stub_err!(store.outbox("o").await); assert!(matches!(err, alkstore::Error::Database(_))); - let err = stub_err!(store.try_lock("l", "owner", 60).await); - assert!(matches!(err, alkstore::Error::Database(_))); let err = store .schedule( "s", diff --git a/tasks/pg-engine-locks.md b/tasks/pg-engine-locks.md index 1f8bcad..0c017fb 100644 --- a/tasks/pg-engine-locks.md +++ b/tasks/pg-engine-locks.md @@ -1,7 +1,7 @@ --- id: pg-engine-locks name: Postgres engine — named locks (`try_lock`, `Lock` handle, duration guards) -status: pending +status: completed depends_on: [pg-engine-seam-tx] scope: narrow risk: low @@ -58,20 +58,20 @@ directly). Over the schema task's locks table: ## Acceptance Criteria -- [ ] `try_lock`: validation (name + owner), duration guard +- [x] `try_lock`: validation (name + owner), duration guard (`ttl <= 0` → opaque `Err`, detail in the source chain), closed-store check, then acquire — `Ok(None)` on held-elsewhere -- [ ] Exclusion + loser shape pinned (two owners, one wins); expiry +- [x] Exclusion + loser shape pinned (two owners, one wins); expiry re-acquisition pinned (the POC's property, engine-side) -- [ ] `renew`: duration guard; full-window reset; `false` on +- [x] `renew`: duration guard; full-window reset; `false` on lost/expired-and-reacquired -- [ ] `release`: consuming, owner-scoped, both boolean arms pinned +- [x] `release`: consuming, owner-scoped, both boolean arms pinned (stale-row delete = `true`; foreign/row-less = `false`) -- [ ] Same-owner re-acquire disposition pinned (and flagged if it +- [x] Same-owner re-acquire disposition pinned (and flagged if it differs from the SQLite arm's) -- [ ] Silent lapse posture documented (no revocation event; no +- [x] Silent lapse posture documented (no revocation event; no general expiry sweeper) -- [ ] `cargo test -p alkstore-postgres` (harness server), clippy +- [x] `cargo test -p alkstore-postgres` (harness server), clippy `-D warnings`, fmt clean; gates green server-less ## References @@ -86,8 +86,89 @@ directly). Over the schema task's locks table: ## Notes -> To be filled by implementation agent +Decisions of record made while implementing (the description didn't +pin them): + +- **Same-owner re-acquire matches the SQLite arm deliberately — no + difference to flag.** The pg acquire deletes the *expired* row for + this name, then `INSERT … ON CONFLICT (name) DO NOTHING`, then a + read-back of the holder column: an unexpired row (the same owner's + or anyone's) survives the DO NOTHING untouched, TTL included — the + substrate's `INSERT OR IGNORE` shape exactly. Verified test-side by + comparing the row's `expires_at` before and after the same-owner + re-acquire (`expired == grant`). +- **Ownership decided by the holder-column read-back, not the + INSERT's affected-row count** — D-4's lesson (the substrate's owner + read-back ends `.ok()`; the forked-fork removed the swallow but the + lesson is the discipline): the insert's `ON CONFLICT DO NOTHING` + already yields an accurate affected count, but the read-back is the + explicit ownership statement (and re-checks expiry races in one + frame). +- **The acquire op's three statements ride one pool connection's + implicit auto-commit (no explicit `BEGIN` frame)** — a `BEGIN`/ + `COMMIT` frame adds two round trips to a path whose correctness rests + on the PK constraint anyway: a concurrent same-name grant either + commits first (the read-back sees the foreign owner → `Ok(None)`) or + lands after (the PK rejects ours, the row stays theirs). The + expiry-delete + insert + read-back race window is closed by the PK, + not by the frame; every other engine op that relies on + single-statement atomicity uses the same posture. Read differently: + no `in_tx` frame was needed (the queue ops' use is about + multi-statement *moves*, not handout decisions). +- **The lock ops' SQL rides the pool path's `table()` helper shape** + (schema-qualified via `QualifiedTable`, `quote_identifier` at the + composition site — the schema task's injection boundary); the locks + table's PK-on-name column is unquoted `name` (not a reserved word — + the offsets' `"offset"` / events' `"key"` hazard pair does not apply + here). +- **Second-precision `now_unix()` stamps the expiry from Rust** + (`expires_at = now + ttl` computed as `i64` binds, not a SQL-side + clock function) — the tx-seam task's documented clock posture (one + `std::time` read per op, the SQL layer never reads a clock). +- **Test lapse windows are backdated through the raw admin + connection** (the queue tests' cross-connection probe shape): a + short TTL's window is stamped via `UPDATE … SET expires_at = now-1`, + so no test sleeps out a multi-second TTL (the SQLite arm's + bounded-wait posture's shape, adapted — the backdate makes the wait + unnecessary). The cross-instance re-acquire test keeps the bounded + poll loop (a second `open`'s acquire is the assertion, the loop is + the honest re-acquisition probe). +- **No explicit schema teardown on the closed-store test's store** — + the fixture closes the store (its `close()` is the drop path) and + every test tears its schema down with `drop_schema` explicitly (the + queue tests' convention); the closed-store test's admin connection + works after the store's `close()` (the pool close and the listener + shutdown don't touch the admin's independent connection). ## Summary -> To be filled on completion \ No newline at end of file +> `alkstore-postgres/src/lock.rs`: the acquire path (`try_lock` — +> shared-namespace + non-empty owner validation, the ADR-023 §2 +> duration guard on `ttl <= 0` (opaque `Database`, detail in the +> source chain) before any round trip, the closed-store check, then +> the acquire op on a pool connection — opportunistic expiry-delete + +> insert-or-reacquire + holder read-back, `Ok(None)` = held +> elsewhere) and `PgLockHandle` (`Lock` impl: `renew` with the same +> duration guard over the owner-scoped full-window update (0 rows = +> lost it); consuming `release(self: Box)` over the owner-scoped +> delete with both boolean arms — no handle-local held-flag caching). +> `store.rs`'s `try_lock` stub replaced with the wiring; the +> wiring-stubs test updated (the lock acquire grants + releases, the +> stub assertion moved to `outbox`/`schedule`/`unschedule`/ +> `run_schedules`). Eight acceptance tests in `store/lock_tests.rs` +> pin: exclusive grant + the loser's `Ok(None)` + the no-refresh +> same-owner re-acquire (matched to the SQLite arm), both duration +> guards with the source-chain detail and the guard-before-acquire +> no-row pin, full-window renew and the lapsed-holder refusals, +> owner-scoped release with the stale-row-`true` and foreign/row-less +> -`false` arms (no sweeper — the stale row persists until an acquire +> or the owner's own delete), the silent lapse + foreign re-acquire, +> the engine-external lapse probe (a second `open` re-acquires — the +> verification backlog's pg re-acquisition pin at engine level), the +> entry-point validation posture (empty/whitespace/ReservedName, +> residue-free), and the closed-store fail-closed shape on all three +> ops. Verified against the harness server (8 new tests; 97 lib tests +> green, ×2 runs); workspace `cargo test` green server-less (332 +> tests, pg tests skipping per convention); `cargo clippy --all- +> targets -- -D warnings` and `cargo fmt --check` clean +> workspace-wide. \ No newline at end of file