From c6a7eeaa45f82e31b6bb98d670ae2ccfa4e7256e Mon Sep 17 00:00:00 2001 From: "glm-5.3-flash" Date: Fri, 9 Oct 2026 06:00:31 +0000 Subject: [PATCH] =?UTF-8?q?Postgres=20engine:=20open=20constructor,=20PgOp?= =?UTF-8?q?ts,=20pool=20+=20listener=20wiring=20=E2=80=94=20forwarder=20sk?= =?UTF-8?q?eleton=20with=20the=20POC-pinned=20pitfalls=20structurally=20ex?= =?UTF-8?q?cluded,=20seam=20error=20mappings,=20wave-3=20stub=20surface=20?= =?UTF-8?q?(task=20pg-engine-open-opts)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Cargo.lock | 1 + alkstore-postgres/Cargo.toml | 5 +- alkstore-postgres/src/forwarder.rs | 463 ++++++++++++++ alkstore-postgres/src/lib.rs | 21 +- alkstore-postgres/src/opts.rs | 72 +++ alkstore-postgres/src/seam.rs | 43 ++ alkstore-postgres/src/store.rs | 342 ++++++++++ alkstore-postgres/src/store/open_tests.rs | 738 ++++++++++++++++++++++ tasks/pg-engine-open-opts.md | 143 ++++- 9 files changed, 1815 insertions(+), 13 deletions(-) create mode 100644 alkstore-postgres/src/forwarder.rs create mode 100644 alkstore-postgres/src/opts.rs create mode 100644 alkstore-postgres/src/seam.rs create mode 100644 alkstore-postgres/src/store.rs create mode 100644 alkstore-postgres/src/store/open_tests.rs diff --git a/Cargo.lock b/Cargo.lock index 4474959..371bc2d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -28,6 +28,7 @@ dependencies = [ "alkstore", "alkstore-contract-suite", "deadpool-postgres", + "serde_json", "thiserror", "tokio", "tokio-postgres", diff --git a/alkstore-postgres/Cargo.toml b/alkstore-postgres/Cargo.toml index 2518a8a..2496dae 100644 --- a/alkstore-postgres/Cargo.toml +++ b/alkstore-postgres/Cargo.toml @@ -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" } \ No newline at end of file diff --git a/alkstore-postgres/src/forwarder.rs b/alkstore-postgres/src/forwarder.rs new file mode 100644 index 0000000..be430e5 --- /dev/null +++ b/alkstore-postgres/src/forwarder.rs @@ -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, + >::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>, +} + +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 { + 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, + /// 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, + /// 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, + /// 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, + shutdown: watch::Sender, +} + +/// 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 { + 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, + reconnect_config: tokio_postgres::Config, + fanout: broadcast::Sender, + channels: Arc, + commands: mpsc::UnboundedReceiver, + connected: Arc, + shutdown: watch::Receiver, +} + +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::(); + 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) { + 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, + } + } +} diff --git a/alkstore-postgres/src/lib.rs b/alkstore-postgres/src/lib.rs index f4e0ed1..f305190 100644 --- a/alkstore-postgres/src/lib.rs +++ b/alkstore-postgres/src/lib.rs @@ -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}; diff --git a/alkstore-postgres/src/opts.rs b/alkstore-postgres/src/opts.rs new file mode 100644 index 0000000..7185712 --- /dev/null +++ b/alkstore-postgres/src/opts.rs @@ -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, + } + } +} diff --git a/alkstore-postgres/src/seam.rs b/alkstore-postgres/src/seam.rs new file mode 100644 index 0000000..56454e7 --- /dev/null +++ b/alkstore-postgres/src/seam.rs @@ -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) -> Error { + Error::database(std::io::Error::other(message.into())) +} diff --git a/alkstore-postgres/src/store.rs b/alkstore-postgres/src/store.rs new file mode 100644 index 0000000..00a8eac --- /dev/null +++ b/alkstore-postgres/src/store.rs @@ -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, + 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, +} + +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 { + &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> { + Ok(Box::new(open_store(config, opts).await?)) +} + +pub(crate) async fn open_store(config: &str, opts: PgOpts) -> alkstore::Result { + 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>> { + 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>> { + 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>> { + 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>> { + 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>> { + 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>>> { + 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> { + 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> { + 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; diff --git a/alkstore-postgres/src/store/open_tests.rs b/alkstore-postgres/src/store/open_tests.rs new file mode 100644 index 0000000..2c77363 --- /dev/null +++ b/alkstore-postgres/src/store/open_tests.rs @@ -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 { + 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 { + 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, + 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(¬ify_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::().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(_))); +} diff --git a/tasks/pg-engine-open-opts.md b/tasks/pg-engine-open-opts.md index e211241..db18fcf 100644 --- a/tasks/pg-engine-open-opts.md +++ b/tasks/pg-engine-open-opts.md @@ -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` + slot; `ListenerConnection` is the `NoTls` + `Connection` 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 \ No newline at end of file +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>` 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. \ No newline at end of file