Postgres engine: named locks — try_lock (validated entry, duration guard, opportunistic expiry-delete + insert-or-reacquire + holder read-back over the locks table), PgLockHandle (full-window renew, consuming owner-scoped release with both boolean arms), same-owner re-acquire matched to the SQLite arm, silent-lapse posture pinned, second open re-acquires after expiry (task pg-engine-locks)

This commit is contained in:
glm-5.3-flash committed 2026-10-09 09:17:39 +00:00
1 parent 3dd83791ff
commit c3591c2d44
6 files changed
+1001 -25

No files matched your search

+1
View File
@@ -33,6 +33,7 @@
//! scheduler/outbox); each task replaces the prior stubs.
mod forwarder;
mod lock;
mod notify;
mod opts;
mod queue;
+236
View File
@@ -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<Self>)` — 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<bool> {
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<String> = 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<AtomicBool>,
name: String,
owner: String,
ttl: i64,
) -> BoxedFuture<'static, Result<Option<Box<dyn Lock>>>> {
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<dyn Lock>
}))
})
}
/// 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<AtomicBool>,
}
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<bool>> {
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<Self>) -> BoxedFuture<'static, Result<bool>> {
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)
})
}
}
+17 -10
View File
@@ -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<Option<Box<dyn alkstore::Lock>>>> {
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;
+644
View File
@@ -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<String> {
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<tokio_postgres::Client> {
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<Fixture> {
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<i64> {
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;
}
+12 -5
View File
@@ -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",