Postgres engine: open constructor, PgOpts, pool + listener wiring — forwarder skeleton with the POC-pinned pitfalls structurally excluded, seam error mappings, wave-3 stub surface (task pg-engine-open-opts)

This commit is contained in:
glm-5.3-flash committed 2026-10-09 06:00:31 +00:00
1 parent 2f1353bd41
commit c6a7eeaa45
9 files changed
+1815 -13

No files matched your search

Generated
+1
View File
@@ -28,6 +28,7 @@ dependencies = [
"alkstore",
"alkstore-contract-suite",
"deadpool-postgres",
"serde_json",
"thiserror",
"tokio",
"tokio-postgres",
+3 -2
View File
@@ -7,10 +7,11 @@ repository.workspace = true
[dependencies]
alkstore = { version = "0.1", path = "../alkstore" }
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
tokio-postgres = "0.7"
deadpool-postgres = "0.14"
serde_json = "1"
thiserror = "2"
tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync", "time"] }
tokio-postgres = "0.7"
[dev-dependencies]
alkstore-contract-suite = { path = "../alkstore-contract-suite" }
+463
View File
@@ -0,0 +1,463 @@
//! The LISTEN forwarder (ADR-004): the wake substrate every later
//! mechanism task rides. This module carries the **skeleton** — the
//! dedicated non-pooled connection, the `poll_message` loop fanning
//! out into a bounded broadcast channel, the dynamic channel set the
//! notify-listen task's `listen` adds channels to — with the full
//! behavior (reconnect honesty contract, receiver bridging) the
//! notify-listen task's.
//!
//! # The two POC-pinned deadlock pitfalls, structurally excluded
//!
//! 1. **Query-vs-poll starvation deadlock** (POC failure mode 1): the
//! `poll_message` loop must be *running before the first client
//! query* on the listener connection. Every connection generation
//! is structured the same way: the poll task is spawned first and
//! owns the `Connection`; LISTEN/UNLISTEN queries ride the `Client`
//! only afterwards — and even mid-connect the request cannot
//! complete before the poll task drives the connection's response.
//! The command select in the loop keeps the poll task forever
//! co-resident with client queries on this connection (the
//! exclusion is held permanently, not just at spawn).
//! 2. **Client-drop closes the session** (POC failure mode 2): the
//! listener's `Client` is owned by the forwarder task and kept
//! alive for the connection generation's lifetime — nobody else
//! holds it, and it is only dropped at shutdown or once the session
//! is already dead server-side (reconnect path).
//!
//! # Connection budget
//!
//! The listener connection is outside the pool (deadpool#360 — pooled
//! connections register LISTEN server-side but can never deliver): a
//! per-process budget line `max_size + 1` per LISTEN-ing process
//! (deployment.md). The listener sets `application_name` for
//! kill-targetability (the deployment.md ops note, POC-carried).
//!
//! # Fanout
//!
//! `poll_message`'s notifications fan into one bounded broadcast
//! channel (capacity 1024 per the POC's shape); lag is **surfaced,
//! not silent** — a `Lagged(n)` receiver resumes streaming (the POC's
//! bridge posture). The wake content (`Wake { channel }` only) and
//! per-receiver bridging are the notify-listen task's; the skeleton
//! broadcasts raw notification channel names.
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use tokio::sync::{broadcast, mpsc, watch};
use tokio_postgres::NoTls;
use crate::seam::database_error;
/// The listener connection's concrete type for the `NoTls` posture
/// (the engine's TLS story rides the consumer's `Config` sslmode for
/// the pooled path; the dedicated listener connects with `NoTls` —
/// TLS hardening is the deployment's ingress posture, the notify-
/// listen task revisits if a task or ADR demands listener TLS).
pub(crate) type ListenerConnection = tokio_postgres::Connection<
tokio_postgres::Socket,
<tokio_postgres::NoTls as tokio_postgres::tls::MakeTlsConnect<tokio_postgres::Socket>>::Stream,
>;
/// The forwarder's reconnect backoff: exponential 50 ms → 2 s cap
/// (the POC's policy; ADR-004 — the notify-listen task owns the full
/// reconnect behavior on this skeleton's loop).
pub(crate) const RECONNECT_BASE_MS: u64 = 50;
pub(crate) const RECONNECT_MAX_MS: u64 = 2000;
/// The bounded broadcast capacity (the POC's 1024 shape; lag surfaced,
/// not silent).
const FANOUT_CAPACITY: usize = 1024;
/// Quote a channel name as a SQL identifier (the POC's `quote_ident`
/// shape — Postgres doubles embedded `"`). Channel names are arbitrary
/// legal names (colons etc.); quoting is what makes them work.
pub(crate) fn quote_ident(name: &str) -> String {
format!("\"{}\"", name.replace('"', "\"\""))
}
/// The raw notification the poll loop fans out — the skeleton's
/// broadcast content; the notify-listen task maps it to the opaque
/// [`alkstore::Wake`].
///
/// `channel` is read by the notify-listen task's receiver bridging
/// (and this wave's fanout-proof tests); ungated in tests, gated in
/// the lib build until that task wires it (the wave-3 intermediate
/// posture — clippy `-D warnings` rejects dead code).
#[derive(Debug, Clone)]
pub(crate) struct RawNotification {
#[cfg_attr(not(test), allow(dead_code))]
pub channel: String,
}
/// The dynamic channel set: the registry the forwarder re-issues
/// LISTENs from after every reconnect. The notify-listen task's
/// `listen` (registration) and receiver-drop (unregistration, the
/// refcounted last-subscriber rule) ride [`Forwarder`]'s register /
/// unregister.
#[derive(Debug, Default)]
pub(crate) struct ChannelSet {
channels: std::sync::Mutex<std::collections::BTreeSet<String>>,
}
impl ChannelSet {
// `insert`/`remove` are the notify-listen task's registry entry
// points (this wave's tests exercise them); gated in the lib
// build until that task wires them.
#[cfg_attr(not(test), allow(dead_code))]
fn insert(&self, channel: &str) -> bool {
self.channels
.lock()
.expect("channel registry mutex poisoned")
.insert(channel.to_string())
}
#[cfg_attr(not(test), allow(dead_code))]
fn remove(&self, channel: &str) -> bool {
self.channels
.lock()
.expect("channel registry mutex poisoned")
.remove(channel)
}
fn snapshot(&self) -> Vec<String> {
self.channels
.lock()
.expect("channel registry mutex poisoned")
.iter()
.cloned()
.collect()
}
}
/// The forwarder handle the store holds: the subscription point for
/// the fanout, the dynamic channel registry + LISTEN/UNLISTEN issuer,
/// and the shutdown switch.
///
/// Shutdown (`Forwarder::shutdown`, the store's `close`/`Drop`): the
/// loop exits, dropping the listener `Client` — the server session
/// ends and the listener connection is gone (the teardown's
/// "listener connection dropped" arm).
#[derive(Debug)]
pub(crate) struct Forwarder {
/// Read by the notify-listen task's receiver bridging (this
/// wave's tests exercise it); gated in the lib build until then.
#[cfg_attr(not(test), allow(dead_code))]
fanout: broadcast::Sender<RawNotification>,
/// Read by the forwarder's reconnect re-issue within the spawned
/// loop's own context (`Arc` move); the field itself is read by
/// the register/unregister surfaces below.
#[cfg_attr(not(test), allow(dead_code))]
channels: Arc<ChannelSet>,
/// LISTEN/UNLISTEN commands to the live connection (the loop
/// issues them; queued commands replay harmlessly across
/// reconnects alongside the registry re-issue). Read by the
/// register/unregister surfaces below.
#[cfg_attr(not(test), allow(dead_code))]
commands_tx: mpsc::UnboundedSender<ListenCommand>,
/// Whether a listener connection is currently live — the
/// register-path's honest transient state (registration during a
/// gap fails with `Database`; the registry entry's recovery rides
/// the reconnect re-issue). Read by `register` below.
#[cfg_attr(not(test), allow(dead_code))]
connected: Arc<AtomicBool>,
shutdown: watch::Sender<bool>,
}
/// Commands to the live listener connection (built by
/// `register`/`unregister` below — gated in the lib build until the
/// notify-listen task wires them).
#[derive(Debug)]
enum ListenCommand {
#[cfg_attr(not(test), allow(dead_code))]
Listen(String),
#[cfg_attr(not(test), allow(dead_code))]
Unlisten(String),
}
impl Forwarder {
/// Take over an established listener connection (the `open`
/// constructor connects so `open` fails `Database` synchronously
/// on an unreachable server) and spawn the forwarder loop.
pub(crate) fn spawn(
client: tokio_postgres::Client,
connection: ListenerConnection,
reconnect_config: tokio_postgres::Config,
) -> Forwarder {
let (fanout, _) = broadcast::channel(FANOUT_CAPACITY);
let channels = Arc::new(ChannelSet::default());
let (commands_tx, commands_rx) = mpsc::unbounded_channel();
let (shutdown_tx, shutdown_rx) = watch::channel(false);
// The connection handed in was just established — the liveness
// flag starts true and the loop clears it at each generation's
// death (the register path's honest transient state).
let connected = Arc::new(AtomicBool::new(true));
// The loop owns both the Client and the Connection — the
// pitfall-2 exclusion (kept alive for the connection's
// lifetime).
tokio::spawn(forwarder_loop(LoopCtx {
client,
connection: Some(connection),
reconnect_config,
fanout: fanout.clone(),
channels: channels.clone(),
commands: commands_rx,
connected: connected.clone(),
shutdown: shutdown_rx,
}));
Forwarder {
fanout,
channels,
commands_tx,
connected,
shutdown: shutdown_tx,
}
}
/// Subscribe to the raw fanout (the notify-listen task bridges
/// broadcast receivers into `WakeReceiver`s). Gated in the lib
/// build until that task wires it (tests exercise it now).
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn subscribe(&self) -> broadcast::Receiver<RawNotification> {
self.fanout.subscribe()
}
/// Register a channel AND issue `LISTEN` on the live connection.
/// The registry write happens first: if the connection dies
/// between the write and the listendelivery, the reconnect path
/// re-issues from the snapshot — a registration is never lost
/// (the ordering is the registration-recovery discipline).
///
/// Fails `Database` only when no listener connection is currently
/// live (mid-reconnect) — transient; a retry (or the notify-listen
/// task's refined posture) recovers it.
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn register(&self, channel: &str) -> Result<(), alkstore::Error> {
self.channels.insert(channel);
if !self.connected.load(Ordering::Acquire) {
return Err(database_error(
"the listener connection is not currently up (mid-reconnect); \
the channel is registered and re-issued on reconnect",
));
}
let _ = self
.commands_tx
.send(ListenCommand::Listen(channel.to_string()));
Ok(())
}
/// Unregister a channel AND issue `UNLISTEN` (the notify-listen
/// task's last-subscriber rule). Registry write first (ordering
/// symmetric with [`Forwarder::register`]'s recovery logic); the
/// UNLISTEN is best-effort — a lost one re-issues nothing
/// (UNLISTEN is not recovery-relevant: the snapshot won't carry
/// the removed channel).
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn unregister(&self, channel: &str) {
self.channels.remove(channel);
let _ = self
.commands_tx
.send(ListenCommand::Unlisten(channel.to_string()));
}
/// Shutdown: flip the switch — the loop exits, dropping the
/// listener `Client` and its session (the listener connection is
/// gone). Idempotent.
pub(crate) fn shutdown(&self) {
let _ = self.shutdown.send(true);
}
}
/// The forwarder loop (the POC's ~90-line shape, restructured for the
/// dynamic channel set + command issuer + shutdown).
///
/// Per connection generation: spawn the poll task (it owns the
/// `Connection` and fans out `poll_message`'s notifications) → issue
/// LISTENs from the live channel-set snapshot → serve commands and
/// wait for the connection's death → reconnect with exponential
/// backoff (50 ms → 2 s cap). The loop consumes the handed-in
/// connection on its first generation and reconnects into the slot
/// afterwards.
struct LoopCtx {
client: tokio_postgres::Client,
connection: Option<ListenerConnection>,
reconnect_config: tokio_postgres::Config,
fanout: broadcast::Sender<RawNotification>,
channels: Arc<ChannelSet>,
commands: mpsc::UnboundedReceiver<ListenCommand>,
connected: Arc<AtomicBool>,
shutdown: watch::Receiver<bool>,
}
async fn forwarder_loop(
LoopCtx {
mut client,
mut connection,
reconnect_config,
fanout,
channels,
mut commands,
connected,
mut shutdown,
}: LoopCtx,
) {
let mut backoff_ms = RECONNECT_BASE_MS;
let mut first_listen = true;
loop {
// The dedicated poll loop, spawned BEFORE any client query can
// ride this connection (pitfall 1's structural exclusion —
// held every generation, and the command select below keeps
// the poll task co-resident with client queries forever). The
// poll task owns the Connection; the loop drives queries
// through the Client.
let Some(poll_connection) = connection.take() else {
return;
};
let mut poll_task = tokio::spawn(poll_loop(poll_connection, fanout.clone()));
// Issue LISTENs from the dynamic channel set's snapshot (on
// first connect it carries pre-registered channels; on
// reconnect it carries everything `listen` ever registered).
// With no channels yet, a `SELECT 1` liveness probe stands in
// — `connected` must reflect a verified-alive session, not an
// assumed one.
let listen_sql = channels
.snapshot()
.iter()
.map(|ch| format!("LISTEN {};\n", quote_ident(ch)))
.collect::<String>();
let listen_ok = if listen_sql.is_empty() {
client.simple_query("SELECT 1").await.is_ok()
} else {
client.batch_execute(&listen_sql).await.is_ok()
};
if listen_ok {
backoff_ms = RECONNECT_BASE_MS;
connected.store(true, Ordering::Release);
if !first_listen {
// The synthetic reconnect-wake (ADR-008 §4's reserved
// channel; ADR-006's recovery semantics) broadcast to
// every subscriber's receiver — the no-replay hole's
// honest communication. The wake-content mapping is
// the notify-listen task's; the skeleton broadcasts
// the reserved channel's name so the fanout carries it.
let _ = fanout.send(RawNotification {
channel: alkstore::RESERVED_LISTENER_RECONNECTED.to_string(),
});
}
first_listen = false;
} else {
// The connection died before/during the LISTEN — fall
// through to the reconnect path with the poll task's
// completion below.
connected.store(false, Ordering::Release);
}
// Drive the connection until it dies (or shutdown). Client
// queries (LISTEN/UNLISTEN commands) run while the poll task
// polls the same session — the co-residency that excludes
// pitfall 1's starvation permanently.
if listen_ok {
loop {
tokio::select! {
biased;
// Shutdown: exit the whole forwarder (dropping the
// client closes the listener session).
_ = shutdown.changed() => {
poll_task.abort();
return;
}
cmd = commands.recv() => {
match cmd {
Some(ListenCommand::Listen(ch)) => {
let _ = client
.batch_execute(&format!("LISTEN {};", quote_ident(&ch)))
.await;
}
Some(ListenCommand::Unlisten(ch)) => {
let _ = client
.batch_execute(&format!("UNLISTEN {};", quote_ident(&ch)))
.await;
}
// The command sender is gone (the store —
// and its forwarder handle — went away
// without a shutdown flip): serve no more,
// exit.
None => {
poll_task.abort();
return;
}
}
}
// The poll task ended = the connection died (a
// fatal poll error or server close). The session
// is already gone server-side; reconnect below.
_ = &mut poll_task => {
break;
}
}
}
} else {
// A failed probe/LISTEN means the connection died: wait
// the poll task out (it owns the connection and is its
// reaper).
let _ = poll_task.await;
}
// Reconnect: the Client is dead with the session (dropping it
// here closes nothing server-side — pitfall 2 was satisfied
// for the connection's whole lifetime). Back off exponentially
// (50 ms → 2 s cap) and reconnect.
connected.store(false, Ordering::Release);
tokio::time::sleep(Duration::from_millis(backoff_ms)).await;
backoff_ms = backoff_ms.saturating_mul(2).min(RECONNECT_MAX_MS);
if *shutdown.borrow() {
return;
}
match reconnect_config.connect(NoTls).await {
Ok((new_client, new_connection)) => {
client = new_client;
connection = Some(new_connection);
}
Err(_e) => {
// Still unreachable — back off again and retry. (The
// skeleton carries the loop shape; the notify-listen
// task owns the reconnect behavior's full contract.)
continue;
}
}
}
}
/// The poll loop over one connection generation: `poll_message`
/// forever, fanning notifications into the bounded broadcast. Returns
/// when the connection dies (a fatal poll error or server close) —
/// the forwarder loop reconnects.
async fn poll_loop(mut connection: ListenerConnection, fanout: broadcast::Sender<RawNotification>) {
loop {
match std::future::poll_fn(|cx| connection.poll_message(cx)).await {
Some(Ok(tokio_postgres::AsyncMessage::Notification(n))) => {
// Bounded fanout — lag is surfaced, not silent: the
// POC's posture is a `Lagged(n)` receiver keeps
// streaming (the receiver bridge's `Lagged` arm logs /
// recovers; capacity pressure must not kill the
// fanout). The send-error arm is "no receivers
// subscribed" — the normal idle posture.
let _ = fanout.send(RawNotification {
channel: n.channel().to_string(),
});
}
Some(Ok(_)) => {}
Some(Err(_)) | None => return,
}
}
}
+20 -1
View File
@@ -8,20 +8,39 @@
//! state; nothing assumes a shared host — this engine's deployment
//! boundary is its crate identity (compile-time engine discovery:
//! pick this crate, get this posture). No runtime capability surface.
//! The per-process connection budget is `max_size` pooled connections
//! **+ 1 non-pooled listener connection per LISTEN-ing process** (the
//! forwarder; deadpool#360 — pooled connections cannot deliver;
//! deployment.md's budget line) — consumers size their server for the
//! sum.
//!
//! **Natively async** (ADR-004/ADR-007): no blocking bridge, no
//! `spawn_blocking` seam — the driver's client is `Send + Sync` and
//! the tx handle holds the pooled object directly (POC #2
//! compile-probe verified).
//! compile-probe verified). Every engine call is a straight `.await`.
//!
//! Durability knobs, pool sizing, and the engine-owned schema name are
//! engine-configuration concerns on [`PgOpts`] (ADR-008 §6) — never
//! contract surface: `synchronous_commit` per-session (ship config on,
//! ADR-004/deployment.md's measured trade), `max_size` the deadpool
//! bound, `schema` the one engine-owned PostgreSQL schema (ADR-010 §8,
//! 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.
mod forwarder;
mod opts;
mod schema;
mod seam;
mod store;
pub use opts::{DEFAULT_MAX_SIZE, PgOpts};
pub use schema::{
BootstrapError, BootstrapResult, DEFAULT_SCHEMA, QualifiedTable, bootstrap, quote_identifier,
tables,
};
pub use store::{PgStore, open};
+72
View File
@@ -0,0 +1,72 @@
//! The pg engine's option struct (ADR-008 §6 — constructors and their
//! option structs live in the engine crates; engine-configuration
//! concerns never touch the contract surface).
//!
//! `#[non_exhaustive]` deliberately does not apply (ADR-017 §3): opts
//! structs are consumer-*constructed*, and the attribute on them would
//! push every construction through a builder. New fields may be added
//! with `Default` fallbacks, per the core opts modules' posture.
use crate::DEFAULT_SCHEMA;
/// The default deadpool pool size (the deployment matrix's budget
/// line: the listener connection adds +1 outside the pool per
/// LISTEN-ing process — deployment.md's connection budgets).
pub const DEFAULT_MAX_SIZE: usize = 8;
/// The listener connection's `application_name` prefix — the base of
/// the per-instance name `{prefix}-{pid}-{seq}` the store assigns
/// (kill-targetable and diagnosable server-side; deployment.md's ops
/// note, the POC's listener-setting carried; the unique suffix makes
/// parallel instances distinguishable in `pg_stat_activity` without
/// losing the ops prefix).
pub(crate) const LISTENER_APPLICATION_NAME: &str = "alkstore-pg-listener";
/// Construction options for the pg engine's
/// [`open`](crate::open) constructor.
///
/// Durability knobs, pool sizing, and the engine-owned schema name are
/// engine-configuration concerns — they never appear on the contract
/// surface (ADR-008 §6); deployment.md carries the documented tuning
/// facts (the `synchronous_commit` trade, the `max_size + 1` listener
/// budget line).
///
/// The connection settings themselves do not ride this struct — they
/// are `open`'s first argument (the tokio-postgres
/// [`Config`](::tokio_postgres::Config)-parseable form, ADR-008 §6's
/// `open(url, PgOpts)` pin); the engine never invents defaults for
/// them. `open` fails with [`Error::Database`] if the config is
/// unparseable or the server unreachable.
#[derive(Debug, Clone)]
pub struct PgOpts {
/// The engine-owned PostgreSQL schema (ADR-010 §8). Default
/// [`DEFAULT_SCHEMA`] (`"alkstore"`). All engine tables live in
/// this one schema; consumer tables co-tenant the instance
/// untouched. The name is a quoted identifier boundary
/// (`quote_identifier`) — reserved-word and hostile names are
/// safe.
pub schema: String,
/// Upper bound on pooled connections (deadpool's `max_size`) —
/// queries, claims, and all non-transactional work ride these.
/// Default [`DEFAULT_MAX_SIZE`] = 8 (the deployment matrix's
/// budget line; the listener adds +1 per LISTEN-ing process
/// *outside* the pool — sizing rule `max_size + 1`).
pub max_size: usize,
/// The per-session `synchronous_commit` durability knob, wired as
/// a connect-options `-c synchronous_commit=…` on every pooled
/// connection (POC-verified per-session SET mechanics). Default
/// `true` — ship config (p50 2.40 ms seam); `false` trades
/// max-tail (40.9 ms) for slightly better p50 — the measured,
/// honest trade (ADR-004, deployment.md).
pub synchronous_commit: bool,
}
impl Default for PgOpts {
fn default() -> Self {
Self {
schema: DEFAULT_SCHEMA.to_string(),
max_size: DEFAULT_MAX_SIZE,
synchronous_commit: true,
}
}
}
+43
View File
@@ -0,0 +1,43 @@
//! The error-mapping posture (the seam task's foundation): driver and
//! pool errors map into the contract taxonomy's opaque
//! [`Error::Database`](alkstore::Error) fallback with the source chain
//! preserved (ADR-008 §5). No engine-specific variants are ever
//! minted; no panics.
//!
//! There is **no `spawn_blocking` seam on this engine** (ADR-004, the
//! crate docs' posture): tokio-postgres's client is natively
//! `Send + Sync` (POC #2 compile-probe) — every engine call is a
//! straight `.await`. These mappings are the two the later mechanism
//! tasks reuse ([`pg_error`], [`pool_error`], or one via
//! [`database_error`]).
use alkstore::Error;
use deadpool_postgres::PoolError;
/// Map a `tokio_postgres::Error` into the contract taxonomy, source
/// chain preserved (ADR-008 §5's opaque fallback — no engine-specific
/// variants minted for its contents).
pub(crate) fn pg_error(e: tokio_postgres::Error) -> Error {
Error::database(e)
}
/// Map a deadpool pool error (`Pool::get`'s return) into the contract
/// taxonomy, source chain preserved (a backend error — the
/// `tokio_postgres::Error` the manager created/recycled with — rides
/// `#[source]`).
pub(crate) fn pool_error(e: PoolError) -> Error {
Error::database(e)
}
/// Map a deadpool pool-build error (`PoolBuilder::build`'s return)
/// into the contract taxonomy, source chain preserved.
pub(crate) fn build_error(e: deadpool_postgres::BuildError) -> Error {
Error::database(e)
}
/// Map a string failure (engine-composed context around an op) into
/// the contract taxonomy (the SQLite twin's
/// `database_error(msg)` shape — `io::Error` source).
pub(crate) fn database_error(message: impl Into<String>) -> Error {
Error::database(std::io::Error::other(message.into()))
}
+342
View File
@@ -0,0 +1,342 @@
//! The pg engine's store: the `open` constructor, the connection
//! architecture (engine-postgres.md's "Connection architecture"), and
//! the [`Store`] trait wiring.
//!
//! Connection architecture: a deadpool pool (`RecyclingMethod::Fast` —
//! no `DISCARD ALL` recycling; the POC's zero-error posture) for
//! queries, claims, and all non-transactional work; one dedicated,
//! **non-pooled** listener connection per process carrying the LISTEN
//! forwarder (deadpool#360 — pooled connections cannot deliver; the
//! per-process budget line is `max_size + 1`, deployment.md); and the
//! engine-owned schema bootstrap run at open.
//!
//! Boot order: pool build → pool-checkout bootstrap (idempotent DDL;
//! open fails `Database` if any step fails) → forwarder spawn (the
//! wake substrate the mechanism tasks ride). Every failure at any step
//! is a typed `Error::Database` with the source chain preserved
//! (ADR-008 §5) — see [`crate::seam`], the mappings the mechanism
//! tasks reuse.
//!
//! Drop/close: the forwarder is shut down (the listener `Client` drops
//! — 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.
use alkstore::Store;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use deadpool_postgres::{Manager, ManagerConfig, RecyclingMethod};
use tokio_postgres::NoTls;
use crate::forwarder::Forwarder;
use crate::opts::{LISTENER_APPLICATION_NAME, PgOpts};
use crate::schema::DEFAULT_SCHEMA;
use crate::seam::{build_error, database_error, pg_error, pool_error};
/// The per-instance listener `application_name` suffix counter (the
/// `{prefix}-{pid}-{seq}` unique name; the ops prefix is shared).
static LISTENER_SEQ: AtomicU64 = AtomicU64::new(0);
/// The pg engine's store handle. Consumers hold
/// [`alkstore::Store`] — the trait is the contract surface; the type
/// exists so the engine can own its machinery and so tests can reach
/// the close path directly.
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))]
listener_application_name: String,
closed: Arc<AtomicBool>,
}
impl std::fmt::Debug for PgStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PgStore")
.field("schema", &self.schema)
.finish_non_exhaustive()
}
}
impl PgStore {
/// Explicit close: shuts the forwarder down (the listener
/// connection's session ends — its `Client` drops inside the
/// forwarder task) and closes the pool. Idempotent. [`Drop`]
/// delegates here. Post-close trait ops fail closed.
pub fn close(&self) {
self.closed.store(true, Ordering::SeqCst);
self.forwarder.shutdown();
self.pool.close();
}
/// 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))]
pub(crate) fn pool(&self) -> &deadpool_postgres::Pool {
&self.pool
}
/// The forwarder handle (the wake substrate) — the mechanism
/// tasks' access point.
#[allow(dead_code)]
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)]
pub(crate) fn schema(&self) -> &str {
&self.schema
}
fn closed_check(&self) -> Result<(), alkstore::Error> {
if self.closed.load(Ordering::Acquire) {
return Err(database_error("the store is closed"));
}
Ok(())
}
}
impl Drop for PgStore {
fn drop(&mut self) {
self.closed.store(true, Ordering::SeqCst);
self.forwarder.shutdown();
self.pool.close();
}
}
/// Open a Postgres-backed store from a connection string (the
/// tokio-postgres `Config`-parseable form) with engine options.
///
/// The connection string/config is the consumer's server address and
/// credentials — the engine never invents defaults for them: `open`
/// fails [`alkstore::Error::Database`] if the config is unparseable or
/// the server unreachable. The full connection architecture boots: the
/// deadpool pool (`RecyclingMethod::Fast`), the engine-owned schema
/// bootstrap (idempotent; run on a pool connection), and the LISTEN
/// forwarder on its dedicated non-pooled connection (the wake
/// substrate; `max_size + 1` per-process budget line).
pub async fn open(config: &str, opts: PgOpts) -> alkstore::Result<Box<dyn Store>> {
Ok(Box::new(open_store(config, opts).await?))
}
pub(crate) async fn open_store(config: &str, opts: PgOpts) -> alkstore::Result<PgStore> {
let cfg: tokio_postgres::Config = config.parse().map_err(pg_error)?;
// Per-session durability knob: the POC's connect-options SET
// mechanics (`-c synchronous_commit=…` on every pooled
// connection; POC-verified — pool.rs's `set_sync_commit_session`
// probe shape).
let mut pool_cfg = cfg.clone();
pool_cfg.options(format!(
"-c synchronous_commit={}",
if opts.synchronous_commit { "on" } else { "off" }
));
let mgr_cfg = ManagerConfig {
recycling_method: RecyclingMethod::Fast,
};
let manager = Manager::from_config(pool_cfg, NoTls, mgr_cfg);
let pool = deadpool_postgres::Pool::builder(manager)
.max_size(opts.max_size)
.build()
.map_err(build_error)?;
// Bootstrap the engine-owned schema on a pool connection
// (idempotent DDL; the listener connection is NOT used — it is
// the wake substrate, not DDL ground). Open fails `Database` if
// the bootstrap fails.
let schema = if opts.schema.is_empty() {
DEFAULT_SCHEMA.to_string()
} else {
opts.schema.clone()
};
{
let conn = pool.get().await.map_err(pool_error)?;
crate::schema::bootstrap(&conn, &schema)
.await
.map_err(|e| pg_error(e.0))?;
}
// Spawn the forwarder on its dedicated non-pooled connection. The
// listener sets application_name for kill-targetability
// (deployment.md's ops note). Connect happens here (not inside the
// task) so an unreachable server fails `open` synchronously with
// `Database`.
let listener_application_name = format!(
"{}-{}-{}",
LISTENER_APPLICATION_NAME,
std::process::id(),
LISTENER_SEQ.fetch_add(1, Ordering::SeqCst)
);
let mut listener_cfg = cfg.clone();
listener_cfg.application_name(listener_application_name.clone());
let (client, connection) = listener_cfg.connect(NoTls).await.map_err(pg_error)?;
let forwarder = Forwarder::spawn(client, connection, listener_cfg);
Ok(PgStore::new(
pool,
forwarder,
schema,
listener_application_name,
))
}
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))]
pub(crate) fn listener_application_name(&self) -> &str {
&self.listener_application_name
}
}
fn stub(wiring: &'static str) -> alkstore::Error {
database_error(wiring)
}
impl Store for PgStore {
fn begin_tx(
&self,
) -> alkstore::BoxedFuture<'_, alkstore::Result<Box<dyn alkstore::TxHandle + Send>>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("begin_tx wiring lands with the seam task")) })
}
fn notify<'a>(
&'a self,
_channel: &str,
_payload: serde_json::Value,
) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("notify wiring lands with the notify-listen task")) })
}
fn listen<'a>(
&'a self,
_channel: &str,
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::WakeReceiver>>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("listen wiring lands with the notify-listen task")) })
}
fn stream<'a>(
&'a self,
_name: &str,
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::StreamHandle>>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("stream wiring lands with the streams task")) })
}
fn queue<'a>(
&'a self,
_name: &str,
_opts: alkstore::QueueOpts,
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::Queue>>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("queue wiring lands with the queues task")) })
}
fn outbox<'a>(
&'a self,
_name: &str,
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::Outbox>>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("outbox wiring lands with the scheduler-outbox task")) })
}
fn try_lock<'a>(
&'a self,
_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")) })
}
fn schedule<'a>(
&'a self,
_name: &str,
_spec: &str,
_queue: &str,
_payload: serde_json::Value,
_opts: alkstore::ScheduleOpts,
) -> alkstore::BoxedFuture<'a, alkstore::Result<alkstore::Schedule>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("schedule wiring lands with the scheduler-outbox task")) })
}
fn unschedule<'a>(&'a self, _name: &str) -> alkstore::BoxedFuture<'a, alkstore::Result<bool>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async {
Err(stub(
"unschedule wiring lands with the scheduler-outbox task",
))
})
}
fn run_schedules<'a>(
&'a self,
_stop: alkstore::StopToken,
) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async {
Err(stub(
"scheduler wiring lands with the scheduler-outbox task",
))
})
}
}
#[cfg(test)]
mod open_tests;
+738
View File
@@ -0,0 +1,738 @@
//! The `open`/`PgOpts`/forwarder-skeleton acceptance tests
//! (`pg-engine-open-opts`): the connection architecture boots against
//! the harness server, opts flow, the forwarder's minimal fanout proof
//! (a LISTEN issued through it delivers), and teardown.
//!
//! Harness convention (the schema task's): 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 POC's shared-server parallel-interference
//! caveat answered by schema-per-test isolation).
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use crate::forwarder::RawNotification;
use crate::opts::{DEFAULT_MAX_SIZE, PgOpts};
use crate::schema::DEFAULT_SCHEMA;
use crate::store::open_store;
use alkstore::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";
/// The unique per-instance listener `application_name` prefix (the
/// `{prefix}-{pid}-{seq}` scheme the store assigns; the pg_stat_activity
/// probes assert against this store's own instance name only).
const LISTENER_APP_PREFIX: &str = "alkstore-pg-listener";
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 unreachable_dsn() -> String {
"host=127.0.0.1 port=1 user=x password=x dbname=x connect_timeout=2".to_string()
}
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 teardown probes (server-level `SHOW`s,
/// `pg_stat_activity` accounting) — 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 == DEFAULT_SCHEMA {
// Never drop the shared default schema in tests — tests that
// exercise the default name only assert its shape.
return;
}
let sql = format!(
"DROP SCHEMA IF EXISTS {} CASCADE",
crate::schema::quote_identifier(schema)
);
let _ = client.batch_execute(&sql).await;
}
/// The test-owned default PgOpts with a fresh schema name.
fn test_opts(schema: &str) -> PgOpts {
PgOpts {
schema: schema.to_string(),
..PgOpts::default()
}
}
/// Count this store's listener connections server-side — by the
/// store's exact assigned `{prefix}-{pid}-{seq}` application_name
/// (parallel tests' listeners never interfere).
async fn listener_sessions(admin: &tokio_postgres::Client, app_name: &str) -> i64 {
admin
.query_one(
"SELECT count(*) FROM pg_stat_activity
WHERE application_name = $1 AND pid <> pg_backend_pid()",
&[&app_name],
)
.await
.unwrap()
.get(0)
}
/// The forwarder's fanout subscription surface for the minimal-delivery
/// proof (bridging into `WakeReceiver`s is the notify-listen task's;
/// the skeleton's proof fans the raw notification channel names out).
async fn wait_for_channel(
rx: &mut tokio::sync::broadcast::Receiver<RawNotification>,
want: &str,
timeout: Duration,
) -> bool {
let deadline = tokio::time::Instant::now() + timeout;
loop {
let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
if remaining.is_zero() {
return false;
}
match tokio::time::timeout(remaining, rx.recv()).await {
Ok(Ok(n)) if n.channel == want => return true,
Ok(Ok(_)) => continue,
Ok(Err(tokio::sync::broadcast::error::RecvError::Lagged(_))) => continue,
Ok(Err(tokio::sync::broadcast::error::RecvError::Closed)) => return false,
Err(_) => return false,
}
}
}
/// Extract the `Err` arm from a stub call whose `Ok` type is not
/// `Debug` (the boxed-handle returns: `begin_tx`, `listen`, `stream`,
/// `queue`, `outbox`, `try_lock`).
macro_rules! stub_err {
($call:expr) => {
match $call {
Err(e) => e,
Ok(_) => panic!(concat!(stringify!($call), " must be a stub")),
}
};
}
#[tokio::test(flavor = "multi_thread")]
async fn open_boots_pool_bootstrap_and_forwarder() {
let Some(dsn) = harness_dsn() else {
eprintln!("skip: no harness server");
return;
};
let schema = instance_namer("boot")();
let store = open_store(&dsn, test_opts(&schema)).await.unwrap();
// Pool up.
let conn = store.pool().get().await.unwrap();
let one: i32 = conn.query_one("SELECT 1", &[]).await.unwrap().get(0);
assert_eq!(one, 1);
drop(conn);
// Bootstrap ran: the 6-table family exists under the configured
// schema.
let admin = harness_client().await.unwrap();
let tables: i64 = admin
.query_one(
"SELECT count(*) FROM pg_tables WHERE schemaname = $1",
&[&schema],
)
.await
.unwrap()
.get(0);
assert_eq!(tables, 6, "the engine-owned table family exists");
// Forwarder skeleton up: register a channel, the LISTEN rides the
// dedicated connection, and a NOTIFY from an unrelated client
// delivers through the fanout (the minimal fanout proof — the
// notify-listen task owns the receiver bridging).
let channel = format!("forwarder_probe_{}", instance_namer("ch")());
store.forwarder().register(&channel).unwrap();
let mut rx = store.forwarder().subscribe();
let admin2 = harness_client().await.unwrap();
let notify_sql = format!(
"NOTIFY {}, 'ping'",
crate::schema::quote_identifier(&channel)
);
let mut delivered = false;
for _ in 0..20 {
let _ = admin2.batch_execute(&notify_sql).await;
if wait_for_channel(&mut rx, &channel, Duration::from_millis(150)).await {
delivered = true;
break;
}
}
assert!(
delivered,
"a LISTEN issued through the forwarder must deliver through its fanout"
);
store.close();
drop_schema(&admin, &schema).await;
}
/// The two POC-pinned deadlock pitfalls, structurally excluded and
/// pinned here:
/// 1. A client query on the listener connection issued immediately
/// after spawn completes (the poll loop was already running — no
/// query-vs-poll starvation deadlock).
/// 2. Notifications deliver well after open (the listener `Client`
/// stays alive for the connection's lifetime — no session death
/// from an unseen drop).
#[tokio::test(flavor = "multi_thread")]
async fn forwarder_poll_loop_survives_command_churn() {
let Some(dsn) = harness_dsn() else {
eprintln!("skip: no harness server");
return;
};
let schema = instance_namer("churn")();
let store = open_store(&dsn, test_opts(&schema)).await.unwrap();
let channel = format!("churn_probe_{}", instance_namer("ch")());
let mut rx = store.forwarder().subscribe();
// Pitfall 1: LISTEN issued immediately after open — if the poll
// loop were not already running this would deadlock.
store.forwarder().register(&channel).unwrap();
let admin = harness_client().await.unwrap();
for i in 0..5 {
admin
.batch_execute(&format!(
"NOTIFY {}, 'burst-{i}'",
crate::schema::quote_identifier(&channel)
))
.await
.unwrap();
}
// Pitfall 2: delivery works long after open (session alive).
assert!(
wait_for_channel(&mut rx, &channel, Duration::from_secs(3)).await,
"notifications deliver long after open (listener session alive — pitfall 2 excluded)"
);
store.forwarder().unregister(&channel);
store.close();
drop_schema(&admin, &schema).await;
}
/// The skeleton's reconnect path: killing the listener's backend
/// server-side terminates its session; the loop reconnects (its
/// exponential backoff, 50 ms → 2 s cap), re-LISTENs from the registry,
/// and a post-reconnect notification delivers. The synthetic
/// reconnect-wake's fanout arm (the notify-listen task's full
/// receiver-bridging behavior) is pinned by the fanout carrying the
/// reserved channel's name after the kill.
#[tokio::test(flavor = "multi_thread")]
async fn forwarder_reconnects_after_backend_kill_and_delivers() {
let Some(dsn) = harness_dsn() else {
eprintln!("skip: no harness server");
return;
};
let schema = instance_namer("reconn")();
let store = open_store(&dsn, test_opts(&schema)).await.unwrap();
let channel = format!("reconn_probe_{}", instance_namer("ch")());
let mut rx = store.forwarder().subscribe();
store.forwarder().register(&channel).unwrap();
let admin = harness_client().await.unwrap();
// Deliver once pre-kill (the subscription path is verified).
admin
.batch_execute(&format!(
"NOTIFY {}, 'pre-kill'",
crate::schema::quote_identifier(&channel)
))
.await
.unwrap();
assert!(
wait_for_channel(&mut rx, &channel, Duration::from_secs(3)).await,
"pre-kill delivery must work first"
);
// Kill the listener's backend (the POC's pg_terminate_backend
// shape, targeted at this store's own unique application_name).
let killed = admin
.query_one(
"SELECT pg_terminate_backend(pid) FROM pg_stat_activity
WHERE application_name = $1 AND pid <> pg_backend_pid()",
&[&store.listener_application_name()],
)
.await
.unwrap()
.get::<_, bool>(0);
assert!(killed, "the listener backend was found and killed");
// The loop reconnects (backoff ≤ 2 s cap) and re-LISTENs from the
// registry: the reserved reconnect-wake arrives on the fanout (the
// skeleton's proof of the re-LISTEN + synthetic-wake arm; the
// full wake-contract behavior is the notify-listen task's), and a
// post-reconnect notification on the registered channel delivers.
let got_wake = wait_for_channel(
&mut rx,
alkstore::RESERVED_LISTENER_RECONNECTED,
Duration::from_secs(10),
)
.await;
let post_sql = format!(
"NOTIFY {}, 'post-reconnect'",
crate::schema::quote_identifier(&channel)
);
let mut delivered = false;
for _ in 0..20 {
let _ = admin.batch_execute(&post_sql).await;
if wait_for_channel(&mut rx, &channel, Duration::from_millis(150)).await {
delivered = true;
break;
}
}
assert!(
got_wake && delivered,
"after a backend kill the forwarder must reconnect (reconnect-wake \
{got_wake}) and deliver post-reconnect notifications ({delivered})"
);
store.close();
drop_schema(&admin, &schema).await;
}
/// Acceptance: pool bounds — `max_size` caps concurrent checkouts (the
/// +1 listener sits outside the pool).
#[tokio::test(flavor = "multi_thread")]
async fn pool_size_is_bounded_by_opts() {
let Some(dsn) = harness_dsn() else {
eprintln!("skip: no harness server");
return;
};
let schema = instance_namer("pool")();
let store = open_store(
&dsn,
PgOpts {
schema: schema.to_string(),
max_size: 2,
..PgOpts::default()
},
)
.await
.unwrap();
assert_eq!(store.pool().status().max_size, 2);
let c1 = store.pool().get().await.unwrap();
let c2 = store.pool().get().await.unwrap();
let pool = store.pool().clone();
let blocked = tokio::spawn(async move { pool.get().await });
tokio::time::sleep(Duration::from_millis(200)).await;
assert!(
!blocked.is_finished(),
"checkout above max_size must block (pool bound enforced)"
);
drop(c2);
let c3 = blocked.await.unwrap().unwrap();
drop(c3);
drop(c1);
store.close();
let admin = harness_client().await.unwrap();
drop_schema(&admin, &schema).await;
}
/// The documented defaults: pool size (the deployment budget line),
/// schema (ADR-010 §8), `synchronous_commit` (ship config on) — and
/// the opts-struct exemption (ADR-017 §3): struct-literal construction
/// with `..Default::default()` compiles.
#[test]
fn defaults_are_the_documented_constants() {
assert_eq!(DEFAULT_MAX_SIZE, 8, "the documented default pool size");
assert_eq!(DEFAULT_SCHEMA, "alkstore", "ADR-010 §8's default schema");
let opts = PgOpts::default();
assert_eq!(opts.max_size, DEFAULT_MAX_SIZE);
assert_eq!(opts.schema, "alkstore");
assert!(
opts.synchronous_commit,
"ship config: synchronous_commit on"
);
let constructed = PgOpts {
max_size: 2,
..PgOpts::default()
};
assert_eq!(constructed.max_size, 2);
assert_eq!(constructed.schema, "alkstore");
}
/// Acceptance: opts flow — the schema name lands in the created
/// schema's name.
#[tokio::test(flavor = "multi_thread")]
async fn schema_name_flows_from_opts() {
let Some(dsn) = harness_dsn() else {
eprintln!("skip: no harness server");
return;
};
let schema = instance_namer("schname")();
let store = open_store(&dsn, test_opts(&schema)).await.unwrap();
assert_eq!(store.schema(), schema);
let admin = harness_client().await.unwrap();
let tables: i64 = admin
.query_one(
"SELECT count(*) FROM pg_tables WHERE schemaname = $1",
&[&schema],
)
.await
.unwrap()
.get(0);
assert_eq!(tables, 6, "the schema name rides PgOpts into the DDL");
store.close();
drop_schema(&admin, &schema).await;
}
/// Acceptance: `synchronous_commit` observable via `SHOW` — the knob
/// rides pool connections as a connect-options `-c` SET (the POC's
/// verified per-session mechanics).
#[tokio::test(flavor = "multi_thread")]
async fn synchronous_commit_knob_is_observable_via_show() {
let Some(dsn) = harness_dsn() else {
eprintln!("skip: no harness server");
return;
};
for (sync_commit, expected) in [(true, "on"), (false, "off")] {
let schema = instance_namer("sync")();
let store = open_store(
&dsn,
PgOpts {
schema: schema.to_string(),
synchronous_commit: sync_commit,
..PgOpts::default()
},
)
.await
.unwrap();
let conn = store.pool().get().await.unwrap();
let row = conn
.query_one("SHOW synchronous_commit", &[])
.await
.unwrap();
let value: String = row.get(0);
assert_eq!(
value, expected,
"the knob's session value must match PgOpts (pool connections)"
);
drop(conn);
store.close();
let admin = harness_client().await.unwrap();
drop_schema(&admin, &schema).await;
}
}
/// Acceptance: the listener connection is (a) outside the pool and (b)
/// exactly one per store — the `max_size + 1` budget line's server-side
/// accounting shape (the POC's pgdiag-st-4 posture, per-instance
/// application_name).
#[tokio::test(flavor = "multi_thread")]
async fn listener_connection_sits_outside_the_pool() {
let Some(dsn) = harness_dsn() else {
eprintln!("skip: no harness server");
return;
};
let schema = instance_namer("outside")();
let store = open_store(
&dsn,
PgOpts {
schema: schema.to_string(),
max_size: 2,
..PgOpts::default()
},
)
.await
.unwrap();
let admin = harness_client().await.unwrap();
let app_name = store.listener_application_name().to_string();
assert!(
app_name.starts_with(LISTENER_APP_PREFIX),
"the listener application_name carries the ops prefix: {app_name}"
);
assert_eq!(
listener_sessions(&admin, &app_name).await,
1,
"exactly one dedicated listener connection per store (outside the pool)"
);
// Hold both pool slots; the listener connection is unaffected (it
// is not pooled).
let _c1 = store.pool().get().await.unwrap();
let _c2 = store.pool().get().await.unwrap();
assert_eq!(
listener_sessions(&admin, &app_name).await,
1,
"pool saturation does not touch the listener"
);
store.close();
drop_schema(&admin, &schema).await;
}
/// Acceptance: close/Drop teardown clean — the listener's server
/// session terminates, the pool closes (checkouts fail), and post-close
/// trait ops fail closed (`Database`, not panic). `Drop` runs the same
/// teardown.
#[tokio::test(flavor = "multi_thread")]
async fn close_and_drop_teardown_cleanly() {
let Some(dsn) = harness_dsn() else {
eprintln!("skip: no harness server");
return;
};
// Explicit close: the listener's server session terminates.
let schema = instance_namer("close")();
let store = open_store(&dsn, test_opts(&schema)).await.unwrap();
let admin = harness_client().await.unwrap();
let app_name = store.listener_application_name().to_string();
assert_eq!(
listener_sessions(&admin, &app_name).await,
1,
"the listener session is up pre-close"
);
store.close();
tokio::time::sleep(Duration::from_millis(300)).await;
let after = listener_sessions(&admin, &app_name).await;
assert_eq!(after, 0, "close terminated this store's listener session");
// Post-close ops fail closed (Database, not panic).
let err = store
.notify("__alkstore_post_close", serde_json::json!({}))
.await
.unwrap_err();
assert!(
matches!(err, alkstore::Error::Database(_)),
"post-close notify fails closed, got: {err:?}"
);
let err = stub_err!(store.begin_tx().await);
assert!(
matches!(err, alkstore::Error::Database(_)),
"post-close begin_tx fails closed, got: {err:?}"
);
// Pool closed: checkouts fail.
let err = store.pool().get().await.unwrap_err();
assert!(
err.to_string().contains("closed"),
"close closes the pool, got: {err}"
);
drop_schema(&admin, &schema).await;
// Drop path: the same teardown reached through drop.
let schema = instance_namer("dropped")();
let store = open_store(&dsn, test_opts(&schema)).await.unwrap();
let app_name = store.listener_application_name().to_string();
drop(store);
tokio::time::sleep(Duration::from_millis(300)).await;
let admin = harness_client().await.unwrap();
assert_eq!(
listener_sessions(&admin, &app_name).await,
0,
"drop terminated the listener session"
);
drop_schema(&admin, &schema).await;
}
/// Acceptance: unreachable server — `open` fails synchronously with
/// the typed `Database` error (source chain preserved), for an
/// unparseable config and a refused connection alike.
#[tokio::test]
async fn open_fails_database_on_unparseable_and_unreachable() {
// Unparseable config string.
let err = open_store(
"definitely not a valid connection string!!!",
PgOpts::default(),
)
.await
.unwrap_err();
match err {
alkstore::Error::Database(source) => {
assert!(
!source.to_string().is_empty(),
"the source chain carries the parse detail"
);
}
other => panic!("unparseable config must be Database, got {other:?}"),
}
// Unreachable server (refused connection, bounded by connect
// timeout) — only when the harness itself is reachable (a
// server-less environment cannot distinguish "unreachable" from
// "no server" meaningfully; the skip posture keeps gates green).
if harness_dsn().is_some() {
let err = open_store(&unreachable_dsn(), PgOpts::default())
.await
.unwrap_err();
match err {
alkstore::Error::Database(source) => {
assert!(
!source.to_string().is_empty(),
"the source chain carries the connect detail"
);
}
other => panic!("unreachable server must be Database, got {other:?}"),
}
} else {
eprintln!("skip: no harness server (unreachable-server arm)");
}
}
/// Acceptance: the trait surface is the wave-3 stub posture — every
/// `Store` method returns `Err(Database(… wiring lands with the …
/// task))`; `with_tx` surfaces the `begin_tx` stub naturally; no
/// panics.
#[tokio::test(flavor = "multi_thread")]
async fn store_trait_methods_are_wiring_stubs() {
let Some(dsn) = harness_dsn() else {
eprintln!("skip: no harness server");
return;
};
let schema = instance_namer("stubs")();
let store = open_store(&dsn, test_opts(&schema)).await.unwrap();
let err = match store.begin_tx().await {
Err(e) => e,
Ok(_) => panic!("begin_tx must be a stub"),
};
assert!(
matches!(err, alkstore::Error::Database(_)),
"begin_tx stub is a Database error"
);
let err = store
.with_tx(Box::new(|_tx| Box::pin(async { Ok(()) })))
.await
.unwrap_err();
assert!(
matches!(err, alkstore::Error::Database(_)),
"with_tx surfaces the begin_tx stub naturally"
);
let err = store
.notify("ch", serde_json::json!("x"))
.await
.unwrap_err();
assert!(matches!(err, alkstore::Error::Database(_)));
let err = match store.listen("ch").await {
Err(e) => e,
Ok(_) => panic!("listen must be a stub"),
};
assert!(matches!(err, alkstore::Error::Database(_)));
let err = match store.stream("s").await {
Err(e) => e,
Ok(_) => panic!("stream must be a stub"),
};
assert!(matches!(err, alkstore::Error::Database(_)));
let err = stub_err!(store.queue("q", alkstore::QueueOpts::default()).await);
assert!(matches!(err, alkstore::Error::Database(_)));
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",
"@every 10s",
"q",
serde_json::json!({}),
alkstore::ScheduleOpts::default(),
)
.await
.unwrap_err();
assert!(matches!(err, alkstore::Error::Database(_)));
let err = store.unschedule("s").await.unwrap_err();
assert!(matches!(err, alkstore::Error::Database(_)));
let err = store
.run_schedules(alkstore::StopToken::new())
.await
.unwrap_err();
assert!(matches!(err, alkstore::Error::Database(_)));
// The stub messages name their landing task (the wave-3 posture's
// "… wiring lands with the … task" shape) — in the source chain
// (`Database`'s Display is opaque; ADR-008 §5's detail-carriage).
let err = stub_err!(store.begin_tx().await);
match err {
alkstore::Error::Database(source) => {
let msg = source.to_string();
assert!(
msg.contains("wiring lands with"),
"stub errors name their landing task, got: {msg}"
);
}
other => panic!("begin_tx stub must be Database, got {other:?}"),
}
drop(store);
let admin = harness_client().await.unwrap();
drop_schema(&admin, &schema).await;
}
/// The error-mapping helpers exist and are the documented reuse point:
/// a mapped error is `Error::Database` with the source chain preserved
/// (the mapping trio the mechanism tasks reuse; the io-Error string
/// helper carries its context through the chain).
#[test]
fn seam_mappings_preserve_the_source_chain() {
use crate::seam::{database_error, pg_error, pool_error};
let parse_err = "!!!".parse::<tokio_postgres::Config>().unwrap_err();
let mapped = pg_error(parse_err);
assert!(matches!(mapped, alkstore::Error::Database(_)));
let source = match &mapped {
alkstore::Error::Database(s) => s.to_string(),
_ => unreachable!(),
};
assert!(
!source.is_empty(),
"the tokio_postgres::Error detail rides the chain"
);
let mapped = database_error("engine context around the failure");
let source = match &mapped {
alkstore::Error::Database(s) => s.to_string(),
_ => unreachable!(),
};
assert!(
source.contains("engine context"),
"the string helper's context rides the chain, got: {source}"
);
let mapped = pool_error(deadpool_postgres::PoolError::Closed);
assert!(matches!(mapped, alkstore::Error::Database(_)));
}
+133 -10
View File
@@ -1,7 +1,7 @@
---
id: pg-engine-open-opts
name: Postgres engine — `open` constructor, `PgOpts`, pool + listener wiring
status: pending
status: completed
depends_on: [pg-engine-schema]
scope: moderate
risk: medium
@@ -77,20 +77,20 @@ name; `synchronous_commit` observable via `SHOW`).
## Acceptance Criteria
- [ ] `PgOpts` with connection config, `schema` (default
- [x] `PgOpts` with connection config, `schema` (default
`alkstore`), `max_size`, `synchronous_commit` (default on);
documented; not `#[non_exhaustive]`
- [ ] `open` boots pool + bootstrap + forwarder skeleton; failure at
- [x] `open` boots pool + bootstrap + forwarder skeleton; failure at
any step is a typed `Database` error (source chain preserved)
- [ ] Forwarder skeleton: dedicated non-pooled connection, poll loop
- [x] Forwarder skeleton: dedicated non-pooled connection, poll loop
before first query, bounded broadcast fanout, dynamic channel
set; the two POC deadlock pitfalls structurally excluded
- [ ] `close()`/`Drop` teardown clean (listener connection dropped,
- [x] `close()`/`Drop` teardown clean (listener connection dropped,
pool closed); post-close ops fail closed
- [ ] Error-mapping helpers exist and are the documented reuse point
- [ ] Crate docs carry the `# Posture` statements (multi-host,
- [x] Error-mapping helpers exist and are the documented reuse point
- [x] Crate docs carry the `# Posture` statements (multi-host,
listener budget, no spawn_blocking)
- [ ] `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 (tests skip)
## References
@@ -106,8 +106,131 @@ name; `synchronous_commit` observable via `SHOW`).
## Notes
> To be filled by implementation agent
Decisions of record made while implementing (the description didn't
pin them):
- **The connection config rides `open`'s first argument, not `PgOpts`**
(the description's bullet list carried it under both; the signature
line `open(config, opts)` and ADR-008 §6's
`alkstore_postgres::open(url, PgOpts)` pin decide it): `PgOpts`
carries schema/max_size/synchronous_commit only. `open` is **async**
(`connect` is; the SQLite twin's sync `open` was file-only). Fail
posture: unparseable config → `pg_error` (a
`tokio_postgres::Error::ConfigParse`); refused connection →
`pg_error` — both typed `Database`, source chain preserved, pinned
by test.
- **`open` connects the listener inline, not inside the task** — an
unreachable server must fail `open` synchronously (the
unreachable-server acceptance row); the forwarder task takes over
the established connection and owns reconnects from there.
- **Forwarder skeleton shape** (`forwarder.rs`, the POC's
~90-line shape restructured): a `LoopCtx` struct carries the loop's
state (clippy's too-many-arguments); per connection generation the
loop spawns a **poll task owning the `Connection`** (`poll_message`
→ bounded broadcast 1024, lag surfaced) before any client query —
pitfall 1's structural exclusion held permanently (the command
select keeps the poll task co-resident with LISTEN/UNLISTEN client
queries every generation); the `Client` is owned by the loop
(pitfall 2). Connection generations carry an `Option<ListenerConnection>`
slot; `ListenerConnection` is the `NoTls`
`Connection<Socket, NoTlsStream>` alias. Commands (LISTEN/UNLISTEN)
ride an mpsc queue the loop issues between polls (reconnect replays
them harmlessly alongside the registry re-issue); the shutdown
watch flips at `close`/`Drop`.
- **Dynamic channel set**: a `BTreeSet`-backed registry
(`Forwarder::{register,unregister}`), registry-write-first ordering
(a registration survives a connection death between write and
LISTEN — the reconnect re-issues from the snapshot).
`register` fails `Database` only mid-reconnect (transient, honest
state surfaced). `subscribe()` on the fanout exposes the raw
notification broadcast — the notify-listen task bridges receivers
onto it. `RawNotification { channel }` is the skeleton's broadcast
content (the `Wake` mapping is that task's).
- **The synthetic reconnect-wake is broadcast by the skeleton** (the
reserved channel `__alkstore_listener_reconnected__` after every
successful re-LISTEN except the first) — the notify-listen task owns
the wake-content/receiver semantics; the skeleton's fanout already
carries it (pinned by the backend-kill test).
- **Listener `application_name` is per-instance**:
`alkstore-pg-listener-{pid}-{seq}` — the deployment.md ops
kill-targetability with a unique suffix so parallel
stores/tests are distinguishable in `pg_stat_activity` (the
outside-the-pool and teardown assertions target the store's own
instance). Exposed at `PgStore::listener_application_name()` for
the notify-listen task's backend-kill tests.
- **Liveness is verified, not assumed**: `connected` starts true (the
handed-in connection is fresh) and a no-channel generation probes
with `SELECT 1` (a dead connection would otherwise report live).
- **Empty `PgOpts::schema` falls back to `DEFAULT_SCHEMA`** (defensive;
the documented default is `.default()`-driven).
- **Error mappings** (`seam.rs`, the mechanism tasks' reuse surface):
`pg_error(tokio_postgres::Error)`, `pool_error(PoolError)`,
`build_error(BuildError)` (the pool-build failure arm —
`PoolBuilder::build` returns `BuildError`, not
`CreatePoolError`), `database_error(msg)` (the SQLite twin's
io-Error string helper).
- **No runtime dep additions beyond the existing set** (tokio-postgres,
deadpool-postgres, tokio, thiserror already carried by the schema
task; `serde_json` added — core's payload type appears in the trait
stubs' signatures). TLS: `NoTls` on both paths (the consumer's
sslmode/TLS story is a deployment concern; noted in the forwarder
module docs, the notify-listen task revisits listener TLS if a task
or ADR ever demands it).
- **Trait stubs** carry the closed-store check *ahead of* the stub
error (`close()`/`Drop` make every arm fail closed immediately —
the post-close fails-closed row is pinned from open onward). Stub
messages name their landing task; their text rides the source
chain (`Database` Display is opaque — asserted via the chain).
- **Machinery accessors are `pub(crate)`** (`pool`, `forwarder`,
`schema`, `listener_application_name`) — the mechanism tasks reach
them; tests exercise them and the lib-build dead-code gates are
`#[cfg_attr(not(test), allow(dead_code))]`-carried until the wiring
tasks land (the wave-3 intermediate posture).
- **Harness reuse**: the schema task's env-carried DSN convention
(`ALKSTORE_PG_HOST/PORT/USER/PASSWORD/DB`), schema-per-test
isolation, skip-clean server-less; the shared default schema is
never dropped by tests; the POC's fanout shape (broadcast 1024,
lag-surfaced) and backoff constants (50 ms → 2 s) are re-owned as
named constants.
## Summary
> To be filled on completion
Stood up `alkstore-postgres`'s public surface: `PgOpts` (schema —
default `alkstore` per ADR-010 §8; `max_size` — default
`DEFAULT_MAX_SIZE = 8`, the deployment budget line with the +1
listener outside the pool; `synchronous_commit` — default on, wired as
connect-options `-c` per-session SET per the POC's verified mechanics;
not `#[non_exhaustive]`, `Debug`+`Clone`+`Default` per ADR-017 §3's
opts exemption), the `async open(config: &str, opts: PgOpts) ->
Result<Box<dyn Store>>` constructor (parse → pool build with
`RecyclingMethod::Fast` → pool-checkout schema bootstrap → dedicated
non-pooled listener connect + forwarder spawn; every failure a typed
`Error::Database` with the source chain preserved), the concrete
`PgStore` re-exported (pool + forwarder handle + schema name held
behind the `Store` trait; `close()`/`Drop` shut the forwarder down
and close the pool — post-close ops fail closed), the LISTEN
forwarder skeleton (`forwarder.rs`: dedicated non-pooled connection,
poll-first structure excluding both POC-pinned deadlock pitfalls,
bounded 1024 broadcast fanout with lag surfaced, dynamic channel set
with registry-write-first recovery ordering, reconnect loop with the
50 ms → 2 s exponential backoff, synthetic reconnect-wake broadcast on
the reserved channel), the seam error mappings (`pg_error`,
`pool_error`, `build_error`, `database_error` — the mechanism tasks'
documented reuse point), and the full `Store` trait as wave-3-style
`Database("… wiring lands with the … task")` stubs with closed-store
fail-closed checks. Crate docs carry the `# Posture` statements
(multi-host, the listener budget line, natively async / no
spawn_blocking). Verified: 12 new open/opts tests
(`src/store/open_tests.rs`) + the 9 schema tests green against the
harness server (boot round-trip, pool bounds + pool saturation leaving
the listener untouched, opts flow: schema name in the created schema +
`synchronous_commit` observable via `SHOW` on both settings, pitfall
pins — immediate post-open LISTEN + long-after-open delivery — and the
outside-the-pool accounting via per-instance `pg_stat_activity`,
backend-kill reconnect with reserved-wake + post-reconnect delivery,
close/drop teardown with server-side session termination + closed
pool + fail-closed ops, unparseable-config and unreachable-server
`Database` failures, the stub surface, the seam mappings' source
chains); workspace `cargo test` green server-less (236 tests, pg tests
skipping per convention); `cargo clippy --all-targets -- -D warnings`
and `cargo fmt --check` clean workspace-wide.