Postgres engine integration: StoreFactory (fresh schema per open, owned idempotent CASCADE teardown, isolation/idempotence pinned), the engine's backlog column (all ten rows green against the harness server — exemplar verified, the three ADR-023 rows verified not rewritten, the five engine-scoped rows, the new pg-arm PayloadTooLarge row discharging the SQLite task's deferred adoption), test-observation accessors cfg(test)-gated with the unused PgStore::new cut, stale stub-era doc text removed, lib docs stating the finished-engine posture (task pg-engine-integration)
This commit is contained in:
1 parent
77619c5e93
commit
fb37da617d
7 files changed
+509
-59
No files matched your search
@@ -27,10 +27,13 @@
|
||||
//! default `alkstore`; co-tenancy with consumer tables is the
|
||||
//! deployment posture).
|
||||
//!
|
||||
//! Wave 4 build order: schema bootstrap (this module — the dependency
|
||||
//! root), then `open`/`PgOpts` + pool + listener wiring, the tx seam,
|
||||
//! the LISTEN forwarder, the mechanisms (queues, streams, locks,
|
||||
//! scheduler/outbox); each task replaces the prior stubs.
|
||||
//! The engine's surface re-exports as the consumer entry points:
|
||||
//! [`open`] + [`PgOpts`] construct the store, [`PgStore`] is the
|
||||
//! concrete type behind the [`alkstore::Store`] trait (the contract
|
||||
//! surface), [`PgTxHandle`] the transaction seam, and the schema
|
||||
//! module's bootstrap the DDL ground. The verification backlog's pg
|
||||
//! column (the ADR-022 contract suite rows this engine owns) runs in
|
||||
//! this crate's `contract_suite` test target.
|
||||
|
||||
mod forwarder;
|
||||
mod lock;
|
||||
|
||||
@@ -21,14 +21,6 @@
|
||||
//! — the server session ends and the listener connection is released)
|
||||
//! and the pool is closed; [`Drop`](PgStore::drop) delegates to
|
||||
//! [`PgStore::close`]. Post-close trait ops fail closed (`Database`).
|
||||
//!
|
||||
//! No `spawn_blocking` seam on this engine (the crate docs' posture);
|
||||
//! 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), `try_lock` (the locks task's) — the wake
|
||||
//! contract's pg arm rides the forwarder.
|
||||
|
||||
use alkstore::Store;
|
||||
use std::sync::Arc;
|
||||
@@ -54,10 +46,12 @@ pub struct PgStore {
|
||||
pool: deadpool_postgres::Pool,
|
||||
forwarder: Arc<Forwarder>,
|
||||
schema: String,
|
||||
/// Read by `listener_application_name()` below (the notify-listen
|
||||
/// task's kill-targeting surface; tests exercise it now): gated in
|
||||
/// the lib build until then.
|
||||
#[cfg_attr(not(test), allow(dead_code))]
|
||||
/// Test-observation surface (the kill-targetability probes assert
|
||||
/// against the exact assigned name) — honestly `#[cfg(test)]`
|
||||
///-gated; the name itself rides `listener_cfg` into the forwarder
|
||||
/// (the reconnect re-issues carry it — that is the correctness
|
||||
/// carrier, not this field).
|
||||
#[cfg(test)]
|
||||
listener_application_name: String,
|
||||
closed: Arc<AtomicBool>,
|
||||
}
|
||||
@@ -82,23 +76,27 @@ impl PgStore {
|
||||
}
|
||||
|
||||
/// The deadpool pool (queries/claims/non-transactional work) —
|
||||
/// the mechanism tasks' access point (tests exercise it now):
|
||||
/// gated in the lib build until the mechanism tasks wire it.
|
||||
#[cfg_attr(not(test), allow(dead_code))]
|
||||
/// test-observation surface for the engine's test modules
|
||||
/// (`#[cfg(test)]`-gated; the lib build reaches the pool through
|
||||
/// the trait wiring and `engine_ctx` directly).
|
||||
#[cfg(test)]
|
||||
pub(crate) fn pool(&self) -> &deadpool_postgres::Pool {
|
||||
&self.pool
|
||||
}
|
||||
|
||||
/// The forwarder handle (the wake substrate) — the mechanism
|
||||
/// tasks' access point.
|
||||
#[allow(dead_code)]
|
||||
/// The forwarder handle (the wake substrate) —
|
||||
/// test-observation surface for the engine's test modules
|
||||
/// (`#[cfg(test)]`-gated; the lib build reaches the forwarder
|
||||
/// through the trait wiring directly).
|
||||
#[cfg(test)]
|
||||
pub(crate) fn forwarder(&self) -> &Arc<Forwarder> {
|
||||
&self.forwarder
|
||||
}
|
||||
|
||||
/// The engine-owned schema name — the mechanism tasks'
|
||||
/// schema-qualified SQL composition point (`quote_identifier`).
|
||||
#[allow(dead_code)]
|
||||
/// The engine-owned schema name — test-observation surface for the
|
||||
/// engine's test modules (`#[cfg(test)]`-gated; the lib build
|
||||
/// composes schema-qualified SQL from the field directly).
|
||||
#[cfg(test)]
|
||||
pub(crate) fn schema(&self) -> &str {
|
||||
&self.schema
|
||||
}
|
||||
@@ -198,36 +196,22 @@ pub(crate) async fn open_store(config: &str, opts: PgOpts) -> alkstore::Result<P
|
||||
let (client, connection) = listener_cfg.connect(NoTls).await.map_err(pg_error)?;
|
||||
let forwarder = Forwarder::spawn(client, connection, listener_cfg);
|
||||
|
||||
Ok(PgStore::new(
|
||||
Ok(PgStore {
|
||||
pool,
|
||||
forwarder,
|
||||
forwarder: Arc::new(forwarder),
|
||||
schema,
|
||||
#[cfg(test)]
|
||||
listener_application_name,
|
||||
))
|
||||
closed: Arc::new(AtomicBool::new(false)),
|
||||
})
|
||||
}
|
||||
|
||||
impl PgStore {
|
||||
fn new(
|
||||
pool: deadpool_postgres::Pool,
|
||||
forwarder: Forwarder,
|
||||
schema: String,
|
||||
listener_application_name: String,
|
||||
) -> PgStore {
|
||||
PgStore {
|
||||
pool,
|
||||
forwarder: Arc::new(forwarder),
|
||||
schema,
|
||||
listener_application_name,
|
||||
closed: Arc::new(AtomicBool::new(false)),
|
||||
}
|
||||
}
|
||||
|
||||
/// This store's listener `application_name` — the per-instance
|
||||
/// `{prefix}-{pid}-{seq}` unique name (kill-targetable,
|
||||
/// `pg_stat_activity`-assertable; deployment.md's ops note; the
|
||||
/// notify-listen task's backend-kill tests target it — tests
|
||||
/// exercise it now). Gated in the lib build until then.
|
||||
#[cfg_attr(not(test), allow(dead_code))]
|
||||
/// `pg_stat_activity`-assertable; deployment.md's ops note).
|
||||
/// Test-observation surface, `#[cfg(test)]`-gated.
|
||||
#[cfg(test)]
|
||||
pub(crate) fn listener_application_name(&self) -> &str {
|
||||
&self.listener_application_name
|
||||
}
|
||||
|
||||
@@ -616,7 +616,7 @@ async fn open_fails_database_on_unparseable_and_unreachable() {
|
||||
/// tests' scope); `with_tx` surfaces the wired `begin_tx` naturally;
|
||||
/// no panics.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn store_trait_methods_are_wiring_stubs() {
|
||||
async fn store_trait_methods_are_wired() {
|
||||
let Some(dsn) = harness_dsn() else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
|
||||
@@ -0,0 +1,296 @@
|
||||
//! The suite-driving test target: the Postgres engine's backlog column
|
||||
//! (ADR-022's option (a) — the rows live in `alkstore-contract-suite`,
|
||||
//! this target is the executor that runs them against the pg factory;
|
||||
//! wave 5 runs the cross-engine equivalence rows on top).
|
||||
//!
|
||||
//! The factory: a **fresh schema per `open`** (the SQLite factory's
|
||||
//! fresh-temp-file precedent; the POC's shared-server
|
||||
//! parallel-interference caveat is answered by exactly this isolation —
|
||||
//! each property's rows live in their own schema, so parallel property
|
||||
//! runs never observe one another), teardown `DROP SCHEMA IF EXISTS …
|
||||
//! CASCADE` over exactly the schemas this factory instance minted,
|
||||
//! tracked per-instance (an idempotent, owned teardown — never the
|
||||
//! shared server's other schemas). The drop is safe against an
|
||||
//! abandoned store: `CASCADE` removes the dependents server-side while
|
||||
//! the store's pooled connections fail closed on their next use (the
|
||||
//! documented post-close posture). Close-arm rows drop the store handle
|
||||
//! (the dispose carrier) — teardown under a live store is not a close
|
||||
//! mechanism.
|
||||
//!
|
||||
//! Server-less environments skip cleanly: the rows are gated behind a
|
||||
//! reachability probe so the workspace gates stay green; wave
|
||||
//! acceptance runs them against the harness server (dockerized
|
||||
//! `postgres:16-alpine` on :15432, connection settings riding the
|
||||
//! environment — the schema task's convention).
|
||||
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
use alkstore::{BoxedFuture, Error, Result, Store};
|
||||
use alkstore_contract_suite::StoreFactory;
|
||||
use alkstore_postgres::{PgOpts, open, quote_identifier};
|
||||
|
||||
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}"
|
||||
))
|
||||
}
|
||||
|
||||
async fn admin_connect(dsn: &str) -> Result<tokio_postgres::Client> {
|
||||
let (client, connection) = tokio_postgres::connect(dsn, tokio_postgres::NoTls)
|
||||
.await
|
||||
.map_err(Error::database)?;
|
||||
tokio::spawn(async move {
|
||||
let _ = connection.await;
|
||||
});
|
||||
Ok(client)
|
||||
}
|
||||
|
||||
/// The Postgres suite factory: a fresh engine-owned schema per `open`
|
||||
/// (unique per call — the isolation guarantee), tracked on the
|
||||
/// instance, tearing down with `DROP SCHEMA IF EXISTS … CASCADE` over
|
||||
/// exactly the schemas it minted. `Send + Sync` — instance state is
|
||||
/// the DSN, a counter, and the schema registry. Default `PgOpts` (pool
|
||||
/// size, `synchronous_commit` ship config) ride; only the schema name
|
||||
/// is factory policy.
|
||||
struct PgFactory {
|
||||
dsn: String,
|
||||
schema_seq: AtomicU64,
|
||||
schemas: std::sync::Mutex<Vec<String>>,
|
||||
}
|
||||
|
||||
impl PgFactory {
|
||||
fn new(tag: &str) -> Self {
|
||||
let dsn = harness_dsn().expect("the harness server is configured for the suite");
|
||||
PgFactory {
|
||||
dsn,
|
||||
// Seeded per (tag, pid, time): two factory instances never
|
||||
// share a schema name, racing opens included.
|
||||
schema_seq: AtomicU64::new(
|
||||
(std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|d| d.subsec_nanos() as u64)
|
||||
.unwrap_or(0))
|
||||
^ (tag.len() as u64)
|
||||
^ ((std::process::id() as u64) << 8),
|
||||
),
|
||||
schemas: std::sync::Mutex::new(Vec::new()),
|
||||
}
|
||||
}
|
||||
|
||||
async fn fresh_schema(&self) -> String {
|
||||
let name = format!(
|
||||
"alkstore_suite_{}_{}_{}",
|
||||
std::process::id(),
|
||||
std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|d| d.as_nanos())
|
||||
.unwrap_or(0),
|
||||
self.schema_seq.fetch_add(1, Ordering::SeqCst),
|
||||
);
|
||||
self.schemas
|
||||
.lock()
|
||||
.expect("factory schema registry poisoned")
|
||||
.push(name.clone());
|
||||
name
|
||||
}
|
||||
}
|
||||
|
||||
impl StoreFactory for PgFactory {
|
||||
fn open(&self) -> BoxedFuture<'_, Result<Box<dyn Store>>> {
|
||||
Box::pin(async move {
|
||||
let schema = self.fresh_schema().await;
|
||||
let opts = PgOpts {
|
||||
schema,
|
||||
..PgOpts::default()
|
||||
};
|
||||
// Unreachable-server: the typed `Database` error surfaces
|
||||
// (tests gate on reachability first so server-less runs
|
||||
// skip).
|
||||
open(&self.dsn, opts).await
|
||||
})
|
||||
}
|
||||
|
||||
fn teardown(&self) -> BoxedFuture<'_, Result<()>> {
|
||||
let dsn = self.dsn.clone();
|
||||
let schemas = self
|
||||
.schemas
|
||||
.lock()
|
||||
.expect("factory schema registry poisoned")
|
||||
.clone();
|
||||
Box::pin(async move {
|
||||
let client = admin_connect(&dsn).await?;
|
||||
for schema in &schemas {
|
||||
let sql = format!("DROP SCHEMA IF EXISTS {} CASCADE", quote_identifier(schema));
|
||||
client.batch_execute(&sql).await.map_err(Error::database)?;
|
||||
}
|
||||
Ok(())
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/// The server reachability probe: server-less environments skip the
|
||||
/// rows (the gates stay green); the harness server's presence runs
|
||||
/// them. One cheap admin connection per test; `teardown` reconnects
|
||||
/// independently (no shared admin client lifetime).
|
||||
async fn harness_ready() -> bool {
|
||||
let Some(dsn) = harness_dsn() else {
|
||||
eprintln!("skip: no harness server configured");
|
||||
return false;
|
||||
};
|
||||
match tokio::time::timeout(std::time::Duration::from_secs(3), admin_connect(&dsn)).await {
|
||||
Ok(Ok(_)) => true,
|
||||
Ok(Err(e)) => {
|
||||
eprintln!("skip: harness server unreachable ({e})");
|
||||
false
|
||||
}
|
||||
Err(_) => {
|
||||
eprintln!("skip: harness server unreachable (timeout)");
|
||||
false
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
mod rows {
|
||||
use alkstore_contract_suite::properties::{
|
||||
drop_rollback_leaves_no_ghosts, duration_refusal_on_non_positive_ttl,
|
||||
enqueue_opts_resolution, extent_clamp_semantics, in_tx_reads_see_own_writes,
|
||||
name_validation_rejects_empty_and_reserved, payload_round_trip_stores_exact_encoding,
|
||||
payload_too_large_produced_on_pg, receiver_close_and_save_arms,
|
||||
};
|
||||
|
||||
use super::PgFactory;
|
||||
|
||||
fn factory(tag: &str) -> PgFactory {
|
||||
PgFactory::new(tag)
|
||||
}
|
||||
|
||||
macro_rules! harness_row {
|
||||
($fn_name:ident, $row_fn:path, $tag:literal) => {
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn $fn_name() {
|
||||
if !super::harness_ready().await {
|
||||
return;
|
||||
}
|
||||
$row_fn(&factory($tag)).await;
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
harness_row!(
|
||||
row_name_validation_rejects_empty_and_reserved,
|
||||
name_validation_rejects_empty_and_reserved,
|
||||
"row-name-validation"
|
||||
);
|
||||
harness_row!(
|
||||
row_extent_clamp_semantics,
|
||||
extent_clamp_semantics,
|
||||
"row-extent-clamp"
|
||||
);
|
||||
harness_row!(
|
||||
row_duration_refusal_on_non_positive_ttl,
|
||||
duration_refusal_on_non_positive_ttl,
|
||||
"row-duration-refusal"
|
||||
);
|
||||
harness_row!(
|
||||
row_payload_round_trip_stores_exact_encoding,
|
||||
payload_round_trip_stores_exact_encoding,
|
||||
"row-payload-round-trip"
|
||||
);
|
||||
harness_row!(
|
||||
row_payload_too_large_produced_on_pg,
|
||||
payload_too_large_produced_on_pg,
|
||||
"row-too-large"
|
||||
);
|
||||
harness_row!(
|
||||
row_drop_rollback_leaves_no_ghosts,
|
||||
drop_rollback_leaves_no_ghosts,
|
||||
"row-drop-rollback"
|
||||
);
|
||||
harness_row!(
|
||||
row_in_tx_reads_see_own_writes,
|
||||
in_tx_reads_see_own_writes,
|
||||
"row-in-tx-ryow"
|
||||
);
|
||||
harness_row!(
|
||||
row_enqueue_opts_resolution,
|
||||
enqueue_opts_resolution,
|
||||
"row-opts-resolution"
|
||||
);
|
||||
harness_row!(
|
||||
row_receiver_close_and_save_arms,
|
||||
receiver_close_and_save_arms,
|
||||
"row-receiver-arms"
|
||||
);
|
||||
}
|
||||
|
||||
mod factory_shape {
|
||||
use alkstore_contract_suite::StoreFactory;
|
||||
|
||||
use super::PgFactory;
|
||||
|
||||
/// ADR-022's factory contract: every open yield is isolated (no two
|
||||
/// opens share rows — proved by a row one store lands being
|
||||
/// invisible to the other) and teardown is idempotent per instance;
|
||||
/// the failing-assertion early-exit (the panic path) leaves teardown
|
||||
/// safe against the abandoned store (`CASCADE`).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn opens_are_isolated_and_teardown_is_idempotent() {
|
||||
if !super::harness_ready().await {
|
||||
return;
|
||||
}
|
||||
let factory = PgFactory::new("factory-shape");
|
||||
|
||||
let a = factory.open().await.unwrap();
|
||||
let b = factory.open().await.unwrap();
|
||||
|
||||
// Isolation: a row `a` publishes lands under `a`'s schema alone
|
||||
// — `b`'s read of the same stream name sees nothing (each
|
||||
// store owns its schema's tables; parallel property runs never
|
||||
// observe one another).
|
||||
let sa = a.stream("iso").await.unwrap();
|
||||
let event_offset = sa.publish(serde_json::json!({"iso": true})).await.unwrap();
|
||||
assert!(event_offset > 0);
|
||||
let sb = b.stream("iso").await.unwrap();
|
||||
let seen = sb.read_since(0, 100).await.unwrap();
|
||||
assert!(
|
||||
seen.is_empty(),
|
||||
"two opens must never share rows — b saw {seen:?}"
|
||||
);
|
||||
drop(sa);
|
||||
drop(sb);
|
||||
drop(a);
|
||||
drop(b);
|
||||
|
||||
let schemas: Vec<String> = {
|
||||
let guard = factory.schemas.lock().expect("schema registry poisoned");
|
||||
guard.clone()
|
||||
};
|
||||
factory.teardown().await.unwrap();
|
||||
factory.teardown().await.unwrap();
|
||||
|
||||
// Idempotent teardown: both opens' schemas are gone server-side.
|
||||
let client = super::admin_connect(&factory.dsn).await.unwrap();
|
||||
for schema in schemas {
|
||||
let count: i64 = client
|
||||
.query_one(
|
||||
"SELECT count(*) FROM pg_namespace WHERE nspname = $1",
|
||||
&[&schema],
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.get(0);
|
||||
assert_eq!(count, 0, "teardown dropped the schema {schema}");
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user