diff --git a/alkstore-postgres/src/lib.rs b/alkstore-postgres/src/lib.rs index f305190..ebd11c8 100644 --- a/alkstore-postgres/src/lib.rs +++ b/alkstore-postgres/src/lib.rs @@ -34,9 +34,11 @@ mod forwarder; mod opts; +mod resolution; mod schema; mod seam; mod store; +mod tx; pub use opts::{DEFAULT_MAX_SIZE, PgOpts}; pub use schema::{ @@ -44,3 +46,4 @@ pub use schema::{ tables, }; pub use store::{PgStore, open}; +pub use tx::PgTxHandle; diff --git a/alkstore-postgres/src/resolution.rs b/alkstore-postgres/src/resolution.rs new file mode 100644 index 0000000..184c7be --- /dev/null +++ b/alkstore-postgres/src/resolution.rs @@ -0,0 +1,227 @@ +//! The engine-side opts-resolution arithmetic (ADR-020 §1–§3) and the +//! derived stamp sets — the shared formulas the tx seam stamps rows +//! with and the auto-commit mechanism tasks reuse, one owner per +//! formula (ADR-012 §2 — the pg engine's own `resolution.rs`; the +//! SQLite one is reference not dependency). +//! +//! # The clock posture +//! +//! `now` is the **one `std::time` read per enqueue** (the single clock +//! read, second-precision) — not a SQL-side `EXTRACT(EPOCH …)`: +//! the resolution arithmetic (delay-over-`run_at`, relative-`expires`) +//! is owned here in Rust, so the same instant stamps `run_at`/ +//! `created_at`/`expires_at` without a second read, and the timestamps +//! leave the engine as plain `i64` parameters the SQL binds. One +//! read, one instant, every column resolved from it. + +use std::time::{SystemTime, UNIX_EPOCH}; + +use alkstore::{EnqueueOpts, QueueOpts}; + +/// The engine's second-precision clock — the single read the +/// resolution and row write share (ADR-020 §1: "one enqueue, one +/// clock read"). No SQL round trip: a plain `std::time` read. +pub(crate) fn now_unix() -> i64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_secs() as i64) + .unwrap_or(0) +} + +/// The plain queue's derived `QueueOpts` defaults (300/3/5/none, +/// honker parities — ADR-010 §3) — the stamp source for the +/// no-handle-open shapes (`enqueue_tx`, scheduler boundary fires, +/// ADR-020 §3 / §3a note). +/// +/// Named (not an inline `QueueOpts::default()` call) so the "derived +/// defaults" formula has one owner: a future default change lands +/// here, once, engine-wide. +pub(crate) fn plain_queue_default_stamps() -> Stamps { + let QueueOpts { + visibility_timeout_s, + max_attempts, + backoff_base_s, + dead_letter_retention_s, + } = QueueOpts::default(); + Stamps { + max_attempts, + visibility_timeout_s, + backoff_base_s, + dead_letter_retention_s, + } +} + +/// The outbox backing queue's derived `QueueOpts` set (60/5/5 — ADR-010 +/// §3, visibility 60 s, max_attempts 5, backoff base 5 s, retention +/// forever). The stamp source for `outbox_enqueue_tx` / +/// `Outbox::enqueue`. +pub(crate) fn outbox_default_stamps() -> Stamps { + Stamps { + visibility_timeout_s: 60, + max_attempts: 5, + backoff_base_s: 5, + dead_letter_retention_s: None, + } +} + +/// The derived backing-queue name for an outbox (the reserved +/// namespace — ADR-008 §4). Engine-internal derivation: the reserved +/// prefix makes the name unreachable by `enqueue_tx` by design +/// (ADR-014 §1) — outbox enqueue is the only path into it. +pub(crate) fn outbox_backing_queue_name(outbox: &str) -> String { + format!("{}outbox:{outbox}", alkstore::RESERVED_PREFIX) +} + +/// The one per-enqueue stamp override (ADR-010 §3a): +/// `EnqueueOpts::max_attempts` is the only per-job override — the +/// other stamps (visibility/backoff/retention) always take the base's. +/// +/// One helper both enqueue shapes share: the tx paths (`enqueue_tx` +/// over the plain-queue defaults, `outbox_enqueue_tx` over the +/// outbox's 60/5/5 set) and the auto-commit `Queue::enqueue` over the +/// handle's `QueueOpts` — resolved once, engine-wide (ADR-012 §2). +pub(crate) fn stamps_with_override(base: Stamps, opts: &EnqueueOpts) -> Stamps { + Stamps { + max_attempts: opts.max_attempts.unwrap_or(base.max_attempts), + ..base + } +} + +/// The enqueue-time stamp set a row carries (ADR-010 §3a) — +/// column-bound primitives; the engine layer resolves +/// contract semantics into these (the substrate boundary posture, +/// engine-owned here since there is no forked substrate). +#[derive(Clone, Copy, Debug)] +pub(crate) struct Stamps { + pub max_attempts: i64, + pub visibility_timeout_s: i64, + pub backoff_base_s: i64, + pub dead_letter_retention_s: Option, +} + +/// Resolve [`EnqueueOpts`] against the enqueue instant: the row's ready +/// time (`delay` → `now + delay`, wins over `run_at` (absolute); +/// neither → `now`) and the absolute row expiry (`expires` relative +/// seconds from enqueue; `None` = never) — ADR-020 §1/§2. One clock +/// read per enqueue: the caller's `now` stamps everything. +pub(crate) struct ResolvedEnqueue { + pub run_at: i64, + pub expires_at: Option, +} + +pub(crate) fn resolve_enqueue_opts(now: i64, opts: &EnqueueOpts) -> ResolvedEnqueue { + let run_at = match opts.delay { + Some(delay) => now.saturating_add(delay), + None => opts.run_at.unwrap_or(now), + }; + let expires_at = opts.expires.map(|s| now.saturating_add(s)); + ResolvedEnqueue { run_at, expires_at } +} + +#[cfg(test)] +mod resolution_tests { + use super::*; + + #[test] + fn plain_queue_default_stamps_are_300_3_5_none() { + let stamps = plain_queue_default_stamps(); + assert_eq!( + ( + stamps.visibility_timeout_s, + stamps.max_attempts, + stamps.backoff_base_s, + stamps.dead_letter_retention_s + ), + (300, 3, 5, None) + ); + } + + #[test] + fn outbox_default_stamps_are_60_5_5_forever() { + let stamps = outbox_default_stamps(); + assert_eq!( + ( + stamps.visibility_timeout_s, + stamps.max_attempts, + stamps.backoff_base_s, + stamps.dead_letter_retention_s + ), + (60, 5, 5, None) + ); + } + + #[test] + fn outbox_backing_queue_name_is_reserved_and_derived() { + assert_eq!( + outbox_backing_queue_name("billing"), + "__alkstore_outbox:billing" + ); + } + + #[test] + fn stamp_override_changes_only_max_attempts() { + let base = plain_queue_default_stamps(); + let opts = EnqueueOpts { + max_attempts: Some(9), + ..Default::default() + }; + let stamped = stamps_with_override(base, &opts); + assert_eq!(stamped.max_attempts, 9); + assert_eq!( + ( + stamped.visibility_timeout_s, + stamped.backoff_base_s, + stamped.dead_letter_retention_s + ), + (300, 5, None) + ); + } + + #[test] + fn delay_wins_over_run_at() { + let opts = EnqueueOpts { + delay: Some(120), + run_at: Some(100_000_000), + ..Default::default() + }; + let resolved = resolve_enqueue_opts(1_000, &opts); + assert_eq!(resolved.run_at, 1_120, "delay wins over run_at"); + assert_eq!(resolved.expires_at, None); + } + + #[test] + fn run_at_alone_is_literal_and_neither_is_now() { + let absolute = EnqueueOpts { + run_at: Some(1_500_000_000), + ..Default::default() + }; + assert_eq!(resolve_enqueue_opts(1_000, &absolute).run_at, 1_500_000_000); + assert_eq!( + resolve_enqueue_opts(1_000, &Default::default()).run_at, + 1_000 + ); + } + + #[test] + fn expires_is_relative_resolved_to_absolute() { + let opts = EnqueueOpts { + expires: Some(90), + ..Default::default() + }; + let resolved = resolve_enqueue_opts(1_000, &opts); + assert_eq!(resolved.expires_at, Some(1_090)); + } + + #[test] + fn saturating_resolution_never_panics() { + let opts = EnqueueOpts { + delay: Some(i64::MAX), + run_at: Some(i64::MAX), + expires: Some(i64::MAX), + ..Default::default() + }; + let resolved = resolve_enqueue_opts(i64::MAX, &opts); + assert_eq!(resolved.run_at, i64::MAX); + assert_eq!(resolved.expires_at, Some(i64::MAX)); + } +} diff --git a/alkstore-postgres/src/store.rs b/alkstore-postgres/src/store.rs index 00a8eac..f2b5cd1 100644 --- a/alkstore-postgres/src/store.rs +++ b/alkstore-postgres/src/store.rs @@ -231,7 +231,7 @@ impl Store for PgStore { 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")) }) + crate::tx::PgTxHandle::begin(self.pool.clone(), self.schema.clone()) } fn notify<'a>( @@ -340,3 +340,6 @@ impl Store for PgStore { #[cfg(test)] mod open_tests; + +#[cfg(test)] +mod tx_tests; diff --git a/alkstore-postgres/src/store/open_tests.rs b/alkstore-postgres/src/store/open_tests.rs index 2c77363..0a75b84 100644 --- a/alkstore-postgres/src/store/open_tests.rs +++ b/alkstore-postgres/src/store/open_tests.rs @@ -541,7 +541,10 @@ async fn close_and_drop_teardown_cleanly() { matches!(err, alkstore::Error::Database(_)), "post-close notify fails closed, got: {err:?}" ); - let err = stub_err!(store.begin_tx().await); + let err = match store.begin_tx().await { + Err(e) => e, + Ok(_) => panic!("post-close begin_tx must fail closed"), + }; assert!( matches!(err, alkstore::Error::Database(_)), "post-close begin_tx fails closed, got: {err:?}" @@ -627,22 +630,15 @@ async fn store_trait_methods_are_wiring_stubs() { 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" - ); + // `begin_tx` is wired (the seam task's): the seam boots and + // commits — its full behavior the tx tests' scope. + let tx = store.begin_tx().await.unwrap(); + tx.commit().await.unwrap(); + 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" - ); + .await; + assert!(err.is_ok(), "with_tx surfaces the wired begin_tx naturally"); let err = store .notify("ch", serde_json::json!("x")) .await @@ -686,7 +682,10 @@ async fn store_trait_methods_are_wiring_stubs() { // 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); + let err = store + .notify("ch", serde_json::json!("x")) + .await + .unwrap_err(); match err { alkstore::Error::Database(source) => { let msg = source.to_string(); @@ -695,7 +694,7 @@ async fn store_trait_methods_are_wiring_stubs() { "stub errors name their landing task, got: {msg}" ); } - other => panic!("begin_tx stub must be Database, got {other:?}"), + other => panic!("notify stub must be Database, got {other:?}"), } drop(store); diff --git a/alkstore-postgres/src/store/tx_tests.rs b/alkstore-postgres/src/store/tx_tests.rs new file mode 100644 index 0000000..c945bd4 --- /dev/null +++ b/alkstore-postgres/src/store/tx_tests.rs @@ -0,0 +1,1263 @@ +//! The transactional seam's acceptance tests (`pg-engine-seam-tx`): +//! `begin_tx` + pooled-object checkout, all eleven `*_tx` methods, the +//! opts-resolution pinning (delay-over-`run_at`, relative-`expires`, +//! the derived 300/3/5/none and 60/5/5 stamp sets), commit/drop +//! dispositions (no-ghosts across job rows, events, notifications, and +//! offset saves), drop panic-safety + the unknowable-state discard +//! arms, ops-after-consume → `Closed`, `with_tx` end-to-end, the +//! 8000-byte notify boundary, and pool accounting under concurrency. +//! +//! Harness convention (the schema/open tasks'): connection settings +//! ride the environment (`ALKSTORE_PG_HOST/PORT/USER/PASSWORD/DB`), +//! never hardcoded; tests without a reachable server skip cleanly. +//! Isolation is a fresh unique schema per test. + +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +use alkstore::EnqueueOpts; +use alkstore::{Error, JobState, Store, StreamEvent, TxHandle}; + +use crate::opts::PgOpts; +use crate::schema::DEFAULT_SCHEMA; +use crate::store::open_store; +use crate::tx::PgTxHandle; + +const ENV_HOST: &str = "ALKSTORE_PG_HOST"; +const ENV_PORT: &str = "ALKSTORE_PG_PORT"; +const ENV_USER: &str = "ALKSTORE_PG_USER"; +const ENV_PASSWORD: &str = "ALKSTORE_PG_PASSWORD"; +const ENV_DB: &str = "ALKSTORE_PG_DB"; + +fn harness_dsn() -> Option { + let host = std::env::var(ENV_HOST).ok()?; + let port: u16 = std::env::var(ENV_PORT).ok()?.parse().ok()?; + let user = std::env::var(ENV_USER).ok()?; + let password = std::env::var(ENV_PASSWORD).ok()?; + let db = std::env::var(ENV_DB).unwrap_or_else(|_| "postgres".to_string()); + Some(format!( + "host={host} port={port} user={user} password={password} dbname={db}" + )) +} + +fn instance_namer(tag: &str) -> impl Fn() -> String + use<'_> { + let counter = AtomicU64::new(0); + move || { + format!( + "{tag}_{}_{}_{}", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos(), + counter.fetch_add(1, Ordering::SeqCst), + ) + } +} + +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 { + return; + } + let sql = format!( + "DROP SCHEMA IF EXISTS {} CASCADE", + crate::schema::quote_identifier(schema) + ); + let _ = client.batch_execute(&sql).await; +} + +fn test_opts(schema: &str) -> PgOpts { + PgOpts { + schema: schema.to_string(), + ..PgOpts::default() + } +} + +fn now_unix() -> i64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_secs() as i64 +} + +/// Count the job rows in a queue via a raw admin client. +async fn raw_job_count(client: &tokio_postgres::Client, schema: &str, queue: &str) -> i64 { + client + .query_one( + &format!( + "SELECT count(*) FROM {} WHERE queue = $1", + crate::schema::QualifiedTable { + schema, + name: crate::schema::tables::JOB, + } + ), + &[&queue], + ) + .await + .unwrap() + .get(0) +} + +/// Count the stream events for a stream via a raw admin client. +async fn raw_event_count(client: &tokio_postgres::Client, schema: &str, stream: &str) -> i64 { + client + .query_one( + &format!( + "SELECT count(*) FROM {} WHERE stream = $1", + crate::schema::QualifiedTable { + schema, + name: crate::schema::tables::EVENTS, + } + ), + &[&stream], + ) + .await + .unwrap() + .get(0) +} + +/// Count the offset checkpoints via a raw admin client. +async fn raw_offset_count(client: &tokio_postgres::Client, schema: &str) -> i64 { + client + .query_one( + &format!( + "SELECT count(*) FROM {}", + crate::schema::QualifiedTable { + schema, + name: crate::schema::tables::OFFSETS, + } + ), + &[], + ) + .await + .unwrap() + .get(0) +} + +/// The forwarder's fanout subscription surface — the listener hears +/// what the server delivers; the listener-side probe for the +/// no-ghost-notification assertion (a listener hears nothing from a +/// rolled-back `notify_tx`). +async fn listener_hears( + rx: &mut tokio::sync::broadcast::Receiver, + channel: &str, + within: Duration, +) -> bool { + let deadline = tokio::time::Instant::now() + within; + 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 == channel => 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, + } + } +} + +/// Acceptance: `begin_tx` checks out a pooled client + `BEGIN`; +/// concurrent `begin_tx` draws distinct pooled objects (pool-bounded); +/// pool exhaustion blocks (deadpool's acquire semantics). +#[tokio::test(flavor = "multi_thread")] +async fn begin_tx_checks_out_and_concurrent_begins_bound_by_the_pool() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("begin")(); + let store = open_store( + &dsn, + PgOpts { + max_size: 1, + ..test_opts(&schema) + }, + ) + .await + .unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let id = tx + .enqueue_tx("q", EnqueueOpts::default(), serde_json::json!({})) + .await + .unwrap(); + assert!(id > 0, "the checked-out client op-ens"); + + // max_size = 1: the second begin_tx blocks on pool exhaustion. + let schema2 = schema.clone(); + let blocked = tokio::spawn(async move { + let store = open_store( + &dsn, + PgOpts { + schema: schema2, + max_size: 1, + ..PgOpts::default() + }, + ) + .await + .unwrap(); + store.begin_tx().await + }); + tokio::time::sleep(Duration::from_millis(200)).await; + assert!( + !blocked.is_finished(), + "pool exhaustion must block begin_tx (deadpool acquire semantics)" + ); + + // Distinct pooled objects: releasing the exhaustion-blocker (the + // blocked begin draws its own distinct object) and concurrently + // issuing ops on both handles never interleaves rows. + drop(tx); + let second = tokio::time::timeout(Duration::from_secs(5), blocked) + .await + .expect("the blocked begin must resume when the pool frees") + .unwrap(); + assert!( + second.is_ok(), + "the exhausted-pool begin must succeed after release" + ); + + store.close(); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: `enqueue_tx` stamps the derived plain-queue defaults +/// 300/3/5/none and resolves opts (delay-over-`run_at`, +/// relative-`expires`) with one clock read (ADR-020 §1–§3a). +#[tokio::test(flavor = "multi_thread")] +async fn enqueue_tx_stamps_defaults_and_resolves_opts() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("stamps")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let job_id = tx + .enqueue_tx( + "emails", + EnqueueOpts { + delay: Some(120), + run_at: Some(100_000_000), + priority: 7, + max_attempts: None, + expires: Some(60), + }, + serde_json::json!({"a": 1}), + ) + .await + .unwrap(); + tx.commit().await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let job = tx.get_job_tx("emails", job_id).await.unwrap().unwrap(); + tx.commit().await.unwrap(); + + let now = now_unix(); + assert_eq!(job.id, job_id); + assert_eq!(job.queue, "emails"); + assert_eq!(job.state, JobState::Pending); + assert_eq!( + ( + job.visibility_timeout_s, + job.max_attempts, + job.backoff_base_s, + job.dead_letter_retention_s + ), + (300, 3, 5, None), + "the plain-queue derived stamps 300/3/5/none" + ); + assert_ne!(job.run_at, 100_000_000, "delay wins over run_at"); + assert!( + (job.run_at - (now + 118)).abs() <= 3, + "the row stamps the resolved ready time" + ); + assert!( + job.expires_at.is_some_and(|e| (e - (now + 60)).abs() <= 3), + "the resolved absolute expiry is stamped" + ); + assert_eq!(job.priority, 7); + assert!( + job.created_at.abs() - job.run_at.abs() <= 130, + "one clock read stamps the row family" + ); + assert_eq!( + job.payload, + serde_json::json!({"a": 1}).to_string().into_bytes() + ); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: `run_at` alone is literal; neither field → ready now +/// (ADR-020 §1), and the outbox shape stamps 60/5/5 + the derived +/// backing-queue name (unreachable by `enqueue_tx`, ADR-014 §1). +#[tokio::test(flavor = "multi_thread")] +async fn run_at_literal_and_outbox_stamps_60_5_5() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("outbox-stamps")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let absolute = tx + .enqueue_tx( + "q", + EnqueueOpts { + run_at: Some(1_500_000_000), + ..Default::default() + }, + serde_json::json!({}), + ) + .await + .unwrap(); + let outbox = tx + .outbox_enqueue_tx( + "billing", + EnqueueOpts { + max_attempts: Some(2), + ..Default::default() + }, + serde_json::json!({"o": 1}), + ) + .await + .unwrap(); + tx.commit().await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let a = tx.get_job_tx("q", absolute).await.unwrap().unwrap(); + tx.commit().await.unwrap(); + assert_eq!(a.run_at, 1_500_000_000, "run_at alone is literal"); + + // The outbox backing queue's row: the reserved-prefixed derived + // name (raw-sql read — no contract surface reads it by name; the + // derived name is unreachable by `enqueue_tx` by design). + let admin = harness_client().await.unwrap(); + let row = admin + .query_one( + &format!( + "SELECT queue, visibility_timeout_s, max_attempts, + backoff_base_s, dead_letter_retention_s, payload + FROM {} WHERE id = $1", + crate::schema::QualifiedTable { + schema: &schema, + name: crate::schema::tables::JOB, + } + ), + &[&outbox], + ) + .await + .unwrap(); + let queue: String = row.get(0); + assert_eq!( + queue, "__alkstore_outbox:billing", + "the derived backing name" + ); + let vish: i64 = row.get(1); + let max_attempts: i64 = row.get(2); + let backoff: i64 = row.get(3); + let retention: Option = row.get(4); + assert_eq!( + (vish, max_attempts, backoff, retention), + (60, 2, 5, None), + "the outbox's derived set 60/5/5 with the override" + ); + let payload: Vec = row.get(5); + assert_eq!( + payload, + serde_json::json!({"o": 1}).to_string().into_bytes() + ); + drop_schema(&admin, &schema).await; + drop(store); +} + +/// Acceptance: read-your-own-writes inside the tx, both directions +/// (ADR-021 §1) — `get_job_tx`/`read_since_tx`/`get_offset_tx` see the +/// tx's own writes pre-commit; committed rows visible post-commit from +/// other connections. +#[tokio::test(flavor = "multi_thread")] +async fn tx_reads_see_own_writes_and_commits_land() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("ryow")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let admin = harness_client().await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let id = tx + .enqueue_tx("q", EnqueueOpts::default(), serde_json::json!({"n": 9})) + .await + .unwrap(); + let off = tx + .publish_tx("events", serde_json::json!({"n": 1})) + .await + .unwrap(); + // In-tx reads see the tx's own writes. + let seen = tx.get_job_tx("q", id).await.unwrap().unwrap(); + assert_eq!(seen.id, id); + let page = tx.read_since_tx("events", 0, 10).await.unwrap(); + assert_eq!(page.len(), 1, "the tx's own publish reads back"); + assert_eq!(page[0].offset, off); + assert_eq!(page[0].stream, "events"); + // Post-commit reads from *other pool connections* see the commit. + tx.commit().await.unwrap(); + assert_eq!(raw_job_count(&admin, &schema, "q").await, 1); + assert_eq!(raw_event_count(&admin, &schema, "events").await, 1); + + // The read-committed post-commit visibility (POC-pinned, other + // direction): a fresh pool tx reads the committed rows. + let mut tx = store.begin_tx().await.unwrap(); + let job = tx.get_job_tx("q", id).await.unwrap().unwrap(); + let page = tx.read_since_tx("events", 0, 10).await.unwrap(); + tx.commit().await.unwrap(); + assert_eq!( + job.payload, + serde_json::json!({"n": 9}).to_string().into_bytes() + ); + assert_eq!(page.len(), 1); + drop_schema(&admin, &schema).await; + drop(store); +} + +/// Acceptance: `publish_with_key_tx` round-trips the key; offsets are +/// monotone; `save_offset_tx` is monotone (a save below the stored +/// checkpoint is a silent no-op — ADR-019 §6). +#[tokio::test(flavor = "multi_thread")] +async fn keyed_publish_round_trips_and_offsets_are_monotone() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("monotone")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let o1 = tx + .publish_with_key_tx("s", Some("k".to_string()), serde_json::json!({})) + .await + .unwrap(); + let plain = tx.publish_tx("s", serde_json::json!({})).await.unwrap(); + assert!(o1 < plain, "offsets are monotone"); + + tx.save_offset_tx("s", "alice", plain).await.unwrap(); + tx.save_offset_tx("s", "alice", 1).await.unwrap(); + let checkpoint = tx.get_offset_tx("s", "alice").await.unwrap(); + assert_eq!(checkpoint, plain, "a save below the checkpoint is a no-op"); + + let from = tx.read_from_consumer_tx("s", "alice", 10).await.unwrap(); + assert_eq!(from.len(), 0, "everything is at or below the checkpoint"); + tx.commit().await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let page = tx.read_since_tx("s", 0, 10).await.unwrap(); + tx.commit().await.unwrap(); + let keyed: &StreamEvent = page.iter().find(|e| e.key.is_some()).unwrap(); + assert_eq!(keyed.key.as_deref(), Some("k"), "the key round-trips"); + + // The absent consumer reads 0 (ADR-019 §1). + let mut tx = store.begin_tx().await.unwrap(); + let absent = tx.get_offset_tx("s", "nobody").await.unwrap(); + tx.commit().await.unwrap(); + assert_eq!(absent, 0); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: no ghosts after drop-rollback — job rows, stream +/// events, notifications, and offset saves drop together (ADR-021 §4's +/// uniform no-ghosts list); a listener hears nothing from a +/// rolled-back `notify_tx`, and the committed positive control +/// delivers (native commit-atomicity, ADR-007). +#[tokio::test(flavor = "multi_thread")] +async fn rollback_ghosts_nothing_across_all_writes() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("no-ghosts")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let admin = harness_client().await.unwrap(); + + // The rollback arm's channel — registered and subscribed FIRST + // (the forwarder's LISTEN issued before any notify can ride), and + // notified ONLY by rolled-back writes in this test: the only + // delivery that could arrive is a ghost. + let rollback_channel = instance_namer("rollback_ch")(); + store.forwarder().register(&rollback_channel).unwrap(); + let mut rollback_rx = store.forwarder().subscribe(); + + let mut tx = store.begin_tx().await.unwrap(); + let ghost_id = tx + .enqueue_tx("q", EnqueueOpts::default(), serde_json::json!({})) + .await + .unwrap(); + tx.publish_tx("s", serde_json::json!({})).await.unwrap(); + tx.notify_tx(&rollback_channel, serde_json::json!({})) + .await + .unwrap(); + tx.save_offset_tx("s", "bob", 1).await.unwrap(); + drop(tx); + assert_eq!(raw_job_count(&admin, &schema, "q").await, 0); + + // …and again, explicit-drop (the same ghost set, exercising the + // same RAII path): every count stays zero, and no delivery rides + // the rollback channel. + let mut tx = store.begin_tx().await.unwrap(); + tx.enqueue_tx("q", EnqueueOpts::default(), serde_json::json!({})) + .await + .unwrap(); + tx.publish_tx("s", serde_json::json!({})).await.unwrap(); + tx.notify_tx(&rollback_channel, serde_json::json!({})) + .await + .unwrap(); + tx.save_offset_tx("s", "bob", 2).await.unwrap(); + drop(tx); + + // Give the detached rollback + any would-be ghost time to surface, + // then assert still-one and silence. + tokio::time::sleep(Duration::from_millis(400)).await; + assert_eq!( + raw_job_count(&admin, &schema, "q").await, + 0, + "no ghost job rows" + ); + assert_eq!( + raw_event_count(&admin, &schema, "s").await, + 0, + "no ghost events" + ); + assert_eq!( + raw_offset_count(&admin, &schema).await, + 0, + "no ghost checkpoints" + ); + let _ = ghost_id; + let heard = listener_hears( + &mut rollback_rx, + &rollback_channel, + Duration::from_millis(300), + ) + .await; + assert!( + !heard, + "a listener must hear nothing after a rolled-back notify_tx" + ); + + // The committed-side positive control on its OWN channel: a + // committed notify_tx delivers at commit. + let commit_channel = instance_namer("commit_ch")(); + store.forwarder().register(&commit_channel).unwrap(); + let mut commit_rx = store.forwarder().subscribe(); + let mut tx = store.begin_tx().await.unwrap(); + tx.notify_tx(&commit_channel, serde_json::json!({})) + .await + .unwrap(); + tx.commit().await.unwrap(); + let heard = listener_hears(&mut commit_rx, &commit_channel, Duration::from_secs(5)).await; + assert!(heard, "a committed notify_tx delivers at commit"); + + drop_schema(&admin, &schema).await; + drop(store); +} + +/// Acceptance: pool accounting stays exact through the unknowable-state +/// arms — a failed COMMIT discards the object (induced: constraint +/// violation poisons the tx server-side), and the next draw gets a +/// fresh object. +#[tokio::test(flavor = "multi_thread")] +async fn failed_commit_discards_and_the_pool_stays_exactly_bounded() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("commit-err")(); + let store = open_store( + &dsn, + PgOpts { + max_size: 2, + ..test_opts(&schema) + }, + ) + .await + .unwrap(); + let admin = harness_client().await.unwrap(); + + // Poison the tx server-side with raw in-tx SQL through the + // concrete-typed handle (the boxed trait object cannot downcast — + // ADR-007; the tests exercise the concrete surface): a *deferred* + // constraint makes COMMIT itself fail (the unknowable-state arm + // the wave-3 review's finding (c) names — an aborted tx commits + // fine server-side by design, so the FAILING-commit arm is + // induced with a deferred FK). + let quoted = |t: &str| { + crate::schema::QualifiedTable { + schema: &schema, + name: t, + } + .to_string() + }; + let (parent, child) = (quoted("tx_parent"), quoted("tx_child")); + admin + .batch_execute(&format!( + "DROP TABLE IF EXISTS {child}; + CREATE TABLE {parent} (id int primary key); + CREATE TABLE {child} (id int primary key, + p int REFERENCES {parent} DEFERRABLE INITIALLY DEFERRED); + INSERT INTO {parent} VALUES (1); DELETE FROM {parent}" + )) + .await + .unwrap(); + let pool = store.pool().clone(); + let tx = PgTxHandle::begin_concrete(pool, schema.clone()) + .await + .unwrap(); + let client_ref: &tokio_postgres::Client = tx.pooled_client(); + let _ = client_ref + .execute(&format!("INSERT INTO {child} (id, p) VALUES (1, 99)"), &[]) + .await; + let err = Box::new(tx) + .commit() + .await + .expect_err("a failing commit must error"); + match &err { + Error::Database(source) => { + assert!(!source.to_string().is_empty(), "source chain non-empty"); + } + other => panic!("failed commit must be Database, got {other:?}"), + } + + // Accounting still exact: a fresh draw after the discard proceeds + // (a pool slot freed), and usable. + let c1 = store.pool().get().await.unwrap(); + let one: i32 = c1.query_one("SELECT 1", &[]).await.unwrap().get(0); + assert_eq!(one, 1); + drop(c1); + let _ = admin + .batch_execute(&format!( + "DROP TABLE IF EXISTS {child}; DROP TABLE IF EXISTS {parent}" + )) + .await; + drop_schema(&admin, &schema).await; + drop(store); +} + +/// Acceptance: ops after commit (the consumed shell) fail closed with +/// `Closed`; the committed shell's drop is harmless (idempotent +/// teardown); the next begin_tx is fresh. +#[tokio::test(flavor = "multi_thread")] +async fn committed_shell_fails_closed_and_teardown_is_clean() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("one-shot")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let mut tx: Box = store.begin_tx().await.unwrap(); + tx.enqueue_tx("q", Default::default(), serde_json::json!({})) + .await + .unwrap(); + tx.commit().await.unwrap(); + + // The next begin is fresh (the re-pooled client serves it). + let mut next = store.begin_tx().await.unwrap(); + let id = next + .enqueue_tx("q", Default::default(), serde_json::json!({"n": 2})) + .await + .unwrap(); + next.commit().await.unwrap(); + assert!(id > 0); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: drop (without commit) rolls back + re-pools — no-ghosts +/// verified for job rows, events, and offsets, and the pool is +/// immediately re-drawable (the re-pool accounting). +#[tokio::test(flavor = "multi_thread")] +async fn drop_rolls_back_and_re_pools() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("drop-rollback")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let admin = harness_client().await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let id = tx + .enqueue_tx("q", EnqueueOpts::default(), serde_json::json!({})) + .await + .unwrap(); + tx.publish_tx("s", serde_json::json!({})).await.unwrap(); + tx.save_offset_tx("s", "bob", 1).await.unwrap(); + drop(tx); + + // The detached rollback: give it a beat, then assert no-ghosts and + // that a fresh begin draws (the object re-pooled by the rollback). + let mut waited = 0; + let mut fresh = loop { + tokio::time::sleep(Duration::from_millis(50)).await; + waited += 50; + if let Ok(tx2) = tokio::time::timeout(Duration::from_millis(50), store.begin_tx()).await { + break tx2.unwrap(); + } + assert!( + waited < 5000, + "the dropped handle's object must return to the pool" + ); + }; + let ghost = fresh.get_job_tx("q", id).await.unwrap(); + let events = fresh.read_since_tx("s", 0, 100).await.unwrap(); + let checkpoint = fresh.get_offset_tx("s", "bob").await.unwrap(); + fresh.commit().await.unwrap(); + assert_eq!(ghost, None, "no ghost job rows after drop-rollback"); + assert!(events.is_empty(), "no ghost events after drop-rollback"); + assert_eq!(checkpoint, 0, "no ghost checkpoints after drop-rollback"); + + drop_schema(&admin, &schema).await; + drop(store); +} + +/// Acceptance: drop is panic-safe — a panic unwinding through a scope +/// holding the handle rolls back (no-ghosts) without a panic escaping +/// drop; the drop-future arm (abort mid-tx) rolls back too. +#[tokio::test(flavor = "multi_thread")] +async fn panic_and_cancellation_through_the_handle_rolls_back() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("panic")(); + let admin = harness_client().await.unwrap(); + let panicked = tokio::spawn({ + let dsn = dsn.clone(); + let schema = schema.clone(); + async move { + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let mut tx = store.begin_tx().await.unwrap(); + tx.enqueue_tx("q", EnqueueOpts::default(), serde_json::json!({})) + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(50)).await; + panic!("unwinding through the tx scope"); + } + }) + .await; + assert!(panicked.is_err(), "the spawned task must have panicked"); + + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let mut tx = store.begin_tx().await.unwrap(); + let count = tx + .read_since_tx("unused_stream_for_probe", 0, 1) + .await + .unwrap() + .len(); + tx.commit().await.unwrap(); + assert_eq!(count, 0); + assert_eq!( + raw_job_count(&admin, &schema, "q").await, + 0, + "no ghosts after the panic's rollback" + ); + drop_schema(&admin, &schema).await; + + // Cancellation arm: abort mid-tx → the future drops → the handle's + // Drop rolls back. + let store2 = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let task = tokio::spawn(async move { + let mut tx = store2.begin_tx().await.unwrap(); + tx.enqueue_tx("q2", EnqueueOpts::default(), serde_json::json!({})) + .await + .unwrap(); + tokio::time::sleep(Duration::from_secs(30)).await; + tx.commit().await.unwrap(); + }); + tokio::time::sleep(Duration::from_millis(200)).await; + task.abort(); + + let admin = harness_client().await.unwrap(); + let mut tx = store.begin_tx().await.unwrap(); + let _ = tx + .read_since_tx("unused_stream_for_probe", 0, 1) + .await + .unwrap(); + tx.commit().await.unwrap(); + assert_eq!( + raw_job_count(&admin, &schema, "q2").await, + 0, + "no ghosts from the cancelled tx" + ); + drop_schema(&admin, &schema).await; + drop(store); +} + +/// Acceptance: ops on a consumed (committed/dropped) shell fail closed +/// with `Closed` — the fail-closed guard arm (the wave-3 shape), on +/// every op the shell can receive. Pinned via the concrete handle's +/// test-only consume (the boxed trait object cannot op after its move +/// into `commit`). +#[tokio::test(flavor = "multi_thread")] +async fn consumed_shell_ops_fail_closed() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("closed-shell")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let pool = store.pool().clone(); + + let mut tx = PgTxHandle::begin_concrete(pool, schema.clone()) + .await + .unwrap(); + tx.consume_shell(); + + // The op surfaces through the trait methods on the concrete type: + let err = tx + .enqueue_tx("q", Default::default(), serde_json::json!({})) + .await + .unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + let err = tx.publish_tx("s", serde_json::json!({})).await.unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + let err = tx.notify_tx("ch", serde_json::json!({})).await.unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + let err = tx.save_offset_tx("s", "c", 1).await.unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + let err = tx.get_job_tx("q", 1).await.unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + let err = tx.get_offset_tx("s", "c").await.unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + let err = tx.read_since_tx("s", 0, 1).await.unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + let err = tx.read_from_consumer_tx("s", "c", 1).await.unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + let err = tx + .outbox_enqueue_tx("o", Default::default(), serde_json::json!({})) + .await + .unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + // Commit on the drained shell errors Closed (one-shot shape). + let err = Box::new(tx).commit().await.unwrap_err(); + assert!(matches!(err, Error::Closed), "{err:?}"); + + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: `with_tx` works end-to-end — Ok ⇒ commit, Err ⇒ +/// rollback (the provided method's disposition, ADR-007/ADR-021 §4). +#[tokio::test(flavor = "multi_thread")] +async fn with_tx_commits_on_ok_and_rolls_back_on_err() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("with-tx")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let admin = harness_client().await.unwrap(); + + store + .with_tx(Box::new(|tx| { + Box::pin(async move { + tx.enqueue_tx("ok-q", EnqueueOpts::default(), serde_json::json!({})) + .await + .unwrap(); + Ok(()) + }) + })) + .await + .unwrap(); + + let err = store + .with_tx(Box::new(|tx| { + Box::pin(async move { + tx.enqueue_tx("err-q", EnqueueOpts::default(), serde_json::json!({})) + .await + .unwrap(); + Err::<(), _>(Error::database(std::io::Error::other("caller error"))) + }) + })) + .await; + assert!(err.is_err()); + + assert_eq!( + raw_job_count(&admin, &schema, "ok-q").await, + 1, + "Ok commits" + ); + assert_eq!( + raw_job_count(&admin, &schema, "err-q").await, + 0, + "Err rolls back" + ); + drop_schema(&admin, &schema).await; + drop(store); +} + +/// Acceptance: the 8000-byte boundary — typed `PayloadTooLarge +/// { limit: 8000 }` client-side before any round trip; ≤ 8000 crosses +/// (ADR-016 §5). +#[tokio::test(flavor = "multi_thread")] +async fn notify_payload_boundary_is_client_side_typed() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("boundary")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let admin = harness_client().await.unwrap(); + + let channel = instance_namer("boundary_ch")(); + let mut tx = store.begin_tx().await.unwrap(); + + // The check's quantity is the serde_json serialization (contract- + // pinned): a bare string Value serializes to `""` (n + 2 + // bytes). The typed arm mirrors the server exactly — the pg wire + // budget is 8000 bytes counting the NUL terminator: ≤ 7999-byte + // texts deliver (the largest crossing serialization is a + // 7997-char string, 7999 bytes), anything ≥ 8000 is rejected + // client-side with the typed variant — probe-pinned (7999 ok, + // 8000 both-protocol-rejected) so no oversized payload ever + // surfaces as an opaque server error. + let under = serde_json::json!("b".repeat(7997)); // 7999-byte serialization + tx.notify_tx(&channel, under).await.unwrap(); + + let at = "a".repeat(7998); // 8000-byte serialization + let err = tx + .notify_tx(&channel, serde_json::json!(at)) + .await + .unwrap_err(); + match err { + Error::PayloadTooLarge { limit } => assert_eq!(limit, 8000), + other => panic!("expect PayloadTooLarge, got {other:?}"), + } + let over = "a".repeat(8000); // 8002-byte serialization + let err = tx + .notify_tx(&channel, serde_json::json!(over)) + .await + .unwrap_err(); + assert!( + matches!(err, Error::PayloadTooLarge { limit: 8000 }), + "the typed arm at any over-size, got {err:?}" + ); + // The check fires before any round trip: a rejected notify_tx + // does not poison the tx (no statement rode the oversized + // payload) — the tx commits cleanly afterward. + tx.commit().await.unwrap(); + + drop_schema(&admin, &schema).await; + drop(store); +} + +/// Acceptance: entry-point validation at every `*_tx` surface, +/// identically on names (shared vs local kinds — ADR-008 §4); the +/// failed validation never enters the engine. +#[tokio::test(flavor = "multi_thread")] +async fn tx_entry_points_validate_names() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("validation")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let check = |err: Error, expect: &str| match err { + Error::InvalidName { .. } => assert_eq!(expect, "empty"), + Error::ReservedName { .. } => assert_eq!(expect, "reserved"), + other => panic!("{expect} case got the wrong error: {other}"), + }; + + check( + tx.enqueue_tx("", Default::default(), serde_json::json!({})) + .await + .unwrap_err(), + "empty", + ); + check( + tx.enqueue_tx( + "__alkstore_outbox:steal", + Default::default(), + serde_json::json!({}), + ) + .await + .unwrap_err(), + "reserved", + ); + check( + tx.publish_tx("", serde_json::json!({})).await.unwrap_err(), + "empty", + ); + check( + tx.publish_with_key_tx("", None, serde_json::json!({})) + .await + .unwrap_err(), + "empty", + ); + check( + tx.notify_tx("", serde_json::json!({})).await.unwrap_err(), + "empty", + ); + check(tx.save_offset_tx("", "c", 1).await.unwrap_err(), "empty"); + check(tx.save_offset_tx("s", "", 1).await.unwrap_err(), "empty"); + check(tx.get_job_tx("", 0).await.unwrap_err(), "empty"); + check(tx.get_offset_tx("", "c").await.unwrap_err(), "empty"); + check(tx.read_since_tx("", 0, 1).await.unwrap_err(), "empty"); + check( + tx.read_from_consumer_tx("", "", 1).await.unwrap_err(), + "empty", + ); + check( + tx.outbox_enqueue_tx("", Default::default(), serde_json::json!({})) + .await + .unwrap_err(), + "empty", + ); + // Whitespace-only is rejected identically. + check( + tx.publish_tx(" ", serde_json::json!({})) + .await + .unwrap_err(), + "empty", + ); + // The tx keeps working: the failed entry points never entered the + // engine and left no residue. + let id = tx + .enqueue_tx("q", Default::default(), serde_json::json!({})) + .await + .unwrap(); + tx.commit().await.unwrap(); + assert!(id > 0); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: empty-`Some`-key → `InvalidName` (ADR-015 §2 / +/// annotated ADR-008 §5) on the tx path. +#[tokio::test(flavor = "multi_thread")] +async fn publish_with_key_tx_rejects_empty_keys() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("key-val")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + for key in [Some(String::new()), Some(" ".to_string())] { + let err = tx + .publish_with_key_tx("s", key, serde_json::json!({})) + .await + .unwrap_err(); + assert!( + matches!(err, Error::InvalidName { .. }), + "empty present key is InvalidName, got {err:?}" + ); + } + tx.commit().await.unwrap(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: `get_job_tx` is queue-scoped — a found row's queue +/// column must match the argument (the wave-3 seam's scoping rule); +/// dead rows are get-job-visible too (the `get_job` shape with +/// `last_error`/`died_at`). +#[tokio::test(flavor = "multi_thread")] +async fn get_job_tx_scopes_to_the_named_queue_and_sees_dead_rows() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("scope")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let admin = harness_client().await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let id = tx + .enqueue_tx("mine", Default::default(), serde_json::json!({})) + .await + .unwrap(); + tx.commit().await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let other = tx.get_job_tx("other-q", id).await.unwrap(); + assert_eq!( + other, None, + "a row from another queue must not surface through a foreign queue name" + ); + // Dead rows carry the diagnosis fields (raw-plant a dead row via + // typed binds — the simple protocol carries no parameters). + admin + .execute( + &format!( + "INSERT INTO {} (id, queue, payload, created_at, died_at) + VALUES ($1, 'mine', $2, 0, 1)", + crate::schema::QualifiedTable { + schema: &schema, + name: crate::schema::tables::DEAD, + } + ), + &[&id, &serde_json::json!({}).to_string().into_bytes()], + ) + .await + .unwrap(); + let dead = tx + .get_job_tx("mine", id) + .await + .unwrap() + .expect("the dead row must surface"); + assert_eq!(dead.state, JobState::Dead); + assert_eq!(dead.died_at, Some(1)); + assert!(dead.last_error.is_none()); + tx.commit().await.unwrap(); + drop_schema(&admin, &schema).await; + drop(store); +} + +/// Acceptance: the extent guard on the tx stream reads (ADR-023 §2) — +/// `limit <= 0` reads nothing, a value not an error, at trait-impl +/// entry. +#[tokio::test(flavor = "multi_thread")] +async fn tx_read_extents_refuse_non_positive_limits_as_empty() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("extents")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + for i in 0..3 { + tx.publish_tx("s", serde_json::json!({"i": i})) + .await + .unwrap(); + } + tx.commit().await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + assert!(tx.read_since_tx("s", 0, 0).await.unwrap().is_empty()); + assert!(tx.read_since_tx("s", 0, -1).await.unwrap().is_empty()); + assert!( + tx.read_from_consumer_tx("s", "c", 0) + .await + .unwrap() + .is_empty() + ); + let one = tx.read_since_tx("s", 0, 2).await.unwrap(); + assert_eq!(one.len(), 2, "a positive extent reads that many"); + tx.commit().await.unwrap(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Interleaved ops in one tx all land in that tx in order (the pooled +/// client's per-connection pipeline — the tx is the serialization). +#[tokio::test(flavor = "multi_thread")] +async fn interleaved_tx_ops_serialize_in_order() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("serialize")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let mut offsets = Vec::new(); + for i in 0..20 { + offsets.push( + tx.publish_tx("s", serde_json::json!({"i": i})) + .await + .unwrap(), + ); + } + tx.commit().await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let page = tx.read_since_tx("s", 0, 100).await.unwrap(); + tx.commit().await.unwrap(); + let seen: Vec = page.iter().map(|e| e.offset).collect(); + assert_eq!(seen, offsets, "the tx's own writes serialize in order"); + + // A failed op inside a tx keeps the handle usable (a poisoned tx + // is pg's abort; a *validation*-failed op is engine-side only). + let mut tx = store.begin_tx().await.unwrap(); + let err = tx + .enqueue_tx("", Default::default(), serde_json::json!({})) + .await + .unwrap_err(); + assert!(matches!(err, Error::InvalidName { name } if name.is_empty())); + let id = tx + .enqueue_tx("q", Default::default(), serde_json::json!({})) + .await + .unwrap(); + assert!(id > 0, "the failed op did not poison the handle"); + drop(tx); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Many txs from many tasks each commit exactly once — the pool +/// serializes the draws, the rows all land (the pool-bounded seam's +/// fan-out posture). +#[tokio::test(flavor = "multi_thread")] +async fn many_txs_commit_exactly_once_through_the_pool() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("fanout")(); + let store = Arc::new(open_store(&dsn, test_opts(&schema)).await.unwrap()); + let admin = harness_client().await.unwrap(); + + let mut joins = Vec::new(); + for i in 0..8 { + let store = Arc::clone(&store); + joins.push(tokio::spawn(async move { + let mut tx = store.begin_tx().await.unwrap(); + tx.enqueue_tx("q", Default::default(), serde_json::json!({"t": i})) + .await + .unwrap(); + tx.commit().await.unwrap(); + })); + } + for j in joins { + j.await.unwrap(); + } + assert_eq!( + raw_job_count(&admin, &schema, "q").await, + 8, + "every tx landed its row exactly once" + ); + drop_schema(&admin, &schema).await; + Arc::into_inner(store).unwrap().close(); +} diff --git a/alkstore-postgres/src/tx.rs b/alkstore-postgres/src/tx.rs new file mode 100644 index 0000000..c7c70c8 --- /dev/null +++ b/alkstore-postgres/src/tx.rs @@ -0,0 +1,661 @@ +//! The transactional seam's pg arm (ADR-007) — `begin_tx`, the +//! [`PgTxHandle`] (the pooled-object handle), and all eleven `*_tx` +//! methods (`pg-engine-seam-tx`). +//! +//! # No bridge, no `spawn_blocking` +//! +//! The handle holds the pooled object directly: tokio-postgres's +//! client is `Send + Sync` (POC #2 compile-probe — ADR-007's pg arm), +//! so every op is a straight `.await` through the held +//! [`deadpool_postgres::Object`], no hop cost, no thread rigging. +//! `begin_tx` checks the pooled client out, issues `BEGIN`, and hands +//! the caller the boxed handle; `commit` and the drop path conclude +//! the transaction and return (or, when the state is unknowable, +//! discard) the object — the pool accounting stays exact on every +//! arm. +//! +//! # The unknowable-state disposition +//! +//! A failed `BEGIN`/`COMMIT`/`ROLLBACK` leaves the client's session +//! state unknowable: the object is **discarded, not re-pooled** — +//! [`Object::take`](deadpool_postgres::Object::take) takes it from +//! the pool permanently (the size accounting stays exact; the next +//! checkout draws a fresh object). Same posture for an already-closed +//! client (`is_closed`). +//! +//! # Drop = rollback (ADR-021 §4) +//! +//! A handle dropped without `commit` rolls back: `Drop` schedules the +//! rollback on the construction runtime (async-native drop posture — +//! `Drop` cannot `.await`; a detached task carries the `ROLLBACK` +//! then re-pools the object, and a failed rollback discards it). The +//! detached task never races ghosts: `tokio::spawn`'s failure arm +//! (runtime shut down before the task runs — the one case the +//! rollback does not run in) drops the held object, whose connection +//! closes server-side and the server itself aborts the open +//! transaction (POC ground: no ghost survives a closed session). Drop +//! is panic-safe (`Handle::spawn` never panics — the handle is +//! captured, not current-context resolved) and never blocks. +//! +//! Ops after `commit`/drop-consumption fail closed with +//! [`Error::Closed`] (the wave-3 fail-closed guard shape) — one-shot +//! handles cannot op past their terminal transition. +//! +//! # Entry-point validation + payload encoding +//! +//! Every name-bearing method validates via core's `validation` +//! helpers (shared-namespace kinds vs consumer-local — ADR-008 §4) +//! **before any round trip**; payloads encode through core's fallible +//! [`alkstore::encode_payload`] with the typed `Codec` error `?`-ed +//! at the seam (ADR-023 §1). `notify_tx` additionally checks the +//! 8000-byte payload boundary client-side, before any round trip +//! (`PayloadTooLarge { limit: 8000 }` — ADR-016 §5); the measured +//! quantity is the serde_json serialization of the payload `Value` — +//! the same byte string ADR-020 §4 pins as the stored row bytes — +//! which is valid UTF-8 by construction. +//! +//! # Engine-side arithmetic +//! +//! The stamp/derivation arithmetic lives in [`crate::resolution`] +//! (one owner per formula, ADR-012 §2) — the plain-queue 300/3/5/none +//! set on `enqueue_tx`, the outbox's 60/5/5 set and the +//! `__alkstore_outbox:{name}` derivation on `outbox_enqueue_tx`, and +//! `resolve_enqueue_opts`' delay-over-`run_at` + relative-`expires` +//! resolution on both enqueue shapes, one clock read per enqueue +//! (the resolution's `now` is the std-time read, document there). + +use deadpool_postgres::Object; +use serde_json::Value; +use tokio::runtime::Handle; + +use alkstore::{ + BoxedFuture, EnqueueOpts, Error, Job, JobState, Result, StreamEvent, TxHandle, + validate_local_name, validate_shared_name, +}; + +use crate::resolution::{ + ResolvedEnqueue, outbox_backing_queue_name, outbox_default_stamps, plain_queue_default_stamps, + resolve_enqueue_opts, stamps_with_override, +}; +use crate::schema::{QualifiedTable, tables}; +use crate::seam::{pg_error, pool_error}; + +/// The notification payload boundary (the contract's pinned limit — +/// client-side checked, typed before any round trip; the predicate +/// rejects at `len >= limit`, the server's NUL-inclusive wire budget). +pub(crate) const NOTIFY_PAYLOAD_LIMIT: usize = 8000; + +/// The caller-held pg transaction handle: the pooled-object handle. +/// Constructed only through [`crate::PgStore`]'s `begin_tx` — never +/// held across an engine `Store` trait object (the trait is what +/// consumers see). +pub struct PgTxHandle { + object: Option, + schema: String, + runtime: Handle, +} + +impl std::fmt::Debug for PgTxHandle { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("PgTxHandle") + .field("open", &self.object.is_some()) + .finish_non_exhaustive() + } +} + +impl PgTxHandle { + /// `begin_tx`: pool checkout + `BEGIN`, returning the boxed + /// handle. Pool exhaustion blocks (deadpool's acquire semantics); + /// a closed pool (or store) fails its checkout. A failed `BEGIN` + /// discards the object (session state unknowable) — checkout + + /// begin fail as one typed `Database` error, no stranded slot. + pub(crate) fn begin( + pool: deadpool_postgres::Pool, + schema: String, + ) -> BoxedFuture<'static, Result>> { + Box::pin(async move { + let object = pool.get().await.map_err(pool_error)?; + let client: &tokio_postgres::Client = &object; + if let Err(e) = client.batch_execute("BEGIN").await { + drop(Object::take(object)); + return Err(pg_error(e)); + } + Ok(Box::new(PgTxHandle { + object: Some(object), + schema, + runtime: Handle::current(), + }) as Box) + }) + } + + /// The concrete-typed constructor — the tests' fault-injection + /// surface (poisoning the tx server-side needs raw in-tx SQL, and + /// the boxed trait object cannot downcast, ADR-007). Never + /// consumer API. + #[cfg(test)] + pub(crate) fn begin_concrete( + pool: deadpool_postgres::Pool, + schema: String, + ) -> BoxedFuture<'static, Result> { + Box::pin(async move { + let object = pool.get().await.map_err(pool_error)?; + let client: &tokio_postgres::Client = &object; + if let Err(e) = client.batch_execute("BEGIN").await { + drop(Object::take(object)); + return Err(pg_error(e)); + } + Ok(PgTxHandle { + object: Some(object), + schema, + runtime: Handle::current(), + }) + }) + } + + /// The private op bridge: borrow the held object for one `.await` + /// op; ops after a terminal disposition fail closed with + /// [`Error::Closed`]. + fn client(&self) -> Result<&tokio_postgres::Client> { + self.object + .as_ref() + .map(|o| { + let client: &tokio_postgres::Client = o; + client + }) + .ok_or(Error::Closed) + } + + /// The held pooled client — the fault-injection surface the + /// tests exercise (server-side tx poisoning); never consumer API. + #[cfg(test)] + pub(crate) fn pooled_client(&self) -> &tokio_postgres::Client { + self.client().expect("tests drive open handles") + } + + /// Drain the held object, simulating a consumed (committed/dropped) + /// shell — the tests' entry into the fail-closed guard arm. + #[cfg(test)] + pub(crate) fn consume_shell(&mut self) { + self.object = None; + } + + /// Schema-qualify a table reference for this handle's engine-owned + /// schema. + fn table(&self, name: &str) -> String { + QualifiedTable { + schema: &self.schema, + name, + } + .to_string() + } + + fn enqueue_row_sql(&self) -> String { + format!( + "INSERT INTO {} (queue, payload, run_at, priority, max_attempts, + expires_at, created_at, visibility_timeout_s, + backoff_base_s, dead_letter_retention_s) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) + RETURNING id", + self.table(tables::JOB) + ) + } + + /// The plain-queue enqueue's shared write (the `enqueue_tx` / + /// `outbox_enqueue_tx` INSERT), called with this handle's + /// qualified job table. + async fn enqueue_row( + &self, + client: &tokio_postgres::Client, + queue: &str, + bytes: Vec, + opts: &EnqueueOpts, + base: crate::resolution::Stamps, + ) -> Result { + let now = crate::resolution::now_unix(); + let ResolvedEnqueue { run_at, expires_at } = resolve_enqueue_opts(now, opts); + let stamps = stamps_with_override(base, opts); + let row = client + .query_one( + &self.enqueue_row_sql(), + &[ + &queue, + &bytes, + &run_at, + &opts.priority, + &stamps.max_attempts, + &expires_at, + &now, + &stamps.visibility_timeout_s, + &stamps.backoff_base_s, + &stamps.dead_letter_retention_s, + ], + ) + .await + .map_err(pg_error)?; + Ok(row.get(0)) + } +} + +/// Encode a trait-crossing payload `Value` into the exact stored bytes +/// (the serde_json serialization — ADR-020 §4). The notify carriage +/// reuses it: the JSON serialization is valid UTF-8 text by +/// construction (ADR-016 §5). +fn encode_payload_bytes(payload: &Value) -> Result> { + alkstore::encode_payload(payload) +} + +impl TxHandle for PgTxHandle { + fn enqueue_tx<'a>( + &'a mut self, + queue: &str, + opts: EnqueueOpts, + payload: Value, + ) -> BoxedFuture<'a, Result> { + if let Err(e) = validate_shared_name(queue) { + return Box::pin(async move { Err(e) }); + } + let queue = queue.to_string(); + Box::pin(async move { + let client = self.client()?; + let bytes = encode_payload_bytes(&payload)?; + self.enqueue_row(client, &queue, bytes, &opts, plain_queue_default_stamps()) + .await + }) + } + + fn publish_tx<'a>(&'a mut self, stream: &str, payload: Value) -> BoxedFuture<'a, Result> { + let stream = stream.to_string(); + Box::pin(async move { self.publish_with_key_tx(&stream, None, payload).await }) + } + + fn publish_with_key_tx<'a>( + &'a mut self, + stream: &str, + key: Option, + payload: Value, + ) -> BoxedFuture<'a, Result> { + if let Err(e) = validate_shared_name(stream) { + return Box::pin(async move { Err(e) }); + } + let key = match key { + Some(k) if k.is_empty() || k.trim().is_empty() => { + return Box::pin(async move { Err(Error::InvalidName { name: k }) }); + } + other => other, + }; + let stream = stream.to_string(); + Box::pin(async move { + let client = self.client()?; + let bytes = encode_payload_bytes(&payload)?; + let created_at = crate::resolution::now_unix(); + let row = client + .query_one( + &format!( + "INSERT INTO {} (stream, \"key\", payload, created_at) + VALUES ($1, $2, $3, $4) + RETURNING id", + self.table(tables::EVENTS) + ), + &[&stream, &key, &bytes, &created_at], + ) + .await + .map_err(pg_error)?; + Ok(row.get(0)) + }) + } + + fn notify_tx<'a>(&'a mut self, channel: &str, payload: Value) -> BoxedFuture<'a, Result<()>> { + if let Err(e) = validate_shared_name(channel) { + return Box::pin(async move { Err(e) }); + } + let channel = channel.to_string(); + Box::pin(async move { + let bytes = encode_payload_bytes(&payload)?; + // The 8000-byte boundary: client-side, typed, before any + // round trip (ADR-016 §5). The pinned quantity is the + // serde_json serialization — the same byte string the + // stored-row bytes pin; the predicate mirrors the server + // exactly (the pg payload budget is the wire's 8000 bytes + // counting the NUL terminator — probe-pinned: ≤ 7999-byte + // texts deliver, 8000 is server-rejected) so every + // too-big payload stays typed, never an opaque + // server-side reject. + if bytes.len() >= NOTIFY_PAYLOAD_LIMIT { + return Err(Error::PayloadTooLarge { + limit: NOTIFY_PAYLOAD_LIMIT, + }); + } + let client = self.client()?; + let text = String::from_utf8(bytes) + .map_err(|e| Error::Codec(format!("notify payload must be utf-8 text: {e}")))?; + client + .query_one("SELECT pg_notify($1, $2)", &[&channel, &text]) + .await + .map_err(pg_error)?; + Ok(()) + }) + } + + fn save_offset_tx<'a>( + &'a mut self, + stream: &str, + consumer: &str, + offset: i64, + ) -> BoxedFuture<'a, Result<()>> { + if let Err(e) = validate_shared_name(stream) { + return Box::pin(async move { Err(e) }); + } + if let Err(e) = validate_local_name(consumer) { + return Box::pin(async move { Err(e) }); + } + let stream = stream.to_string(); + let consumer = consumer.to_string(); + Box::pin(async move { + let client = self.client()?; + // Monotone upsert (one owner): a first save seeds the + // checkpoint — clamped at 0, the absent-consumer=0 rule + // (ADR-019 §1) — and a save at-or-below the stored one + // updates nothing (ADR-019 §6's silent no-op). + client + .execute( + &format!( + "INSERT INTO {} AS o (stream, consumer, \"offset\") + VALUES ($1, $2, GREATEST($3::bigint, 0)) + ON CONFLICT (stream, consumer) DO UPDATE + SET \"offset\" = EXCLUDED.\"offset\" + WHERE EXCLUDED.\"offset\" > o.\"offset\"", + self.table(tables::OFFSETS) + ), + &[&stream, &consumer, &offset], + ) + .await + .map_err(pg_error)?; + Ok(()) + }) + } + + fn get_job_tx<'a>( + &'a mut self, + queue: &str, + job_id: i64, + ) -> BoxedFuture<'a, Result>> { + if let Err(e) = validate_shared_name(queue) { + return Box::pin(async move { Err(e) }); + } + let queue = queue.to_string(); + Box::pin(async move { + let client = self.client()?; + // Queue-scoped (the seam's scoping rule): the id + // namespace is global across queues — the queue column + // must match the argument or the row does not surface. + // Dead rows are get-job-visible too (the auto-commit + // `get_job` shape; the diagnosis columns ride). + let row = client + .query_opt( + &format!( + "SELECT id, queue, 'dead', payload, priority, run_at, attempts, + max_attempts, worker_id, claimed_at, claim_expires_at, + created_at, expires_at, visibility_timeout_s, + backoff_base_s, dead_letter_retention_s, + last_error, died_at + FROM {} WHERE id = $1 AND queue = $2", + self.table(tables::DEAD) + ), + &[&job_id, &queue], + ) + .await + .map_err(pg_error)?; + if let Some(row) = row { + return Ok(Some(job_from_row(&row)?)); + } + let row = client + .query_opt( + &format!( + "SELECT id, queue, state, payload, priority, run_at, attempts, + max_attempts, worker_id, claimed_at, claim_expires_at, + created_at, expires_at, visibility_timeout_s, + backoff_base_s, dead_letter_retention_s, + NULL::text, NULL::bigint + FROM {} WHERE id = $1 AND queue = $2", + self.table(tables::JOB) + ), + &[&job_id, &queue], + ) + .await + .map_err(pg_error)?; + match row { + Some(row) => Ok(Some(job_from_row(&row)?)), + None => Ok(None), + } + }) + } + + fn get_offset_tx<'a>( + &'a mut self, + stream: &str, + consumer: &str, + ) -> BoxedFuture<'a, Result> { + if let Err(e) = validate_shared_name(stream) { + return Box::pin(async move { Err(e) }); + } + if let Err(e) = validate_local_name(consumer) { + return Box::pin(async move { Err(e) }); + } + let stream = stream.to_string(); + let consumer = consumer.to_string(); + Box::pin(async move { + let client = self.client()?; + let row = client + .query_opt( + &format!( + "SELECT \"offset\" FROM {} WHERE stream = $1 AND consumer = $2", + self.table(tables::OFFSETS) + ), + &[&stream, &consumer], + ) + .await + .map_err(pg_error)?; + Ok(row.map(|r| r.get(0)).unwrap_or(0)) + }) + } + + fn read_since_tx<'a>( + &'a mut self, + stream: &str, + offset: i64, + limit: i64, + ) -> BoxedFuture<'a, Result>> { + if let Err(e) = validate_shared_name(stream) { + return Box::pin(async move { Err(e) }); + } + let stream = stream.to_string(); + Box::pin(async move { + // The extent guard (ADR-023 §2): `limit <= 0` reads + // nothing — a value, never an error — at trait-impl + // entry, before any round trip. + if limit <= 0 { + return Ok(Vec::new()); + } + let client = self.client()?; + let rows = client + .query( + &format!( + "SELECT id, \"key\", payload, created_at + FROM {} + WHERE stream = $1 AND id > $2 + ORDER BY id ASC + LIMIT $3", + self.table(tables::EVENTS) + ), + &[&stream, &offset, &limit], + ) + .await + .map_err(pg_error)?; + stream_events_from_rows(&rows, stream) + }) + } + + fn read_from_consumer_tx<'a>( + &'a mut self, + stream: &str, + consumer: &str, + limit: i64, + ) -> BoxedFuture<'a, Result>> { + if let Err(e) = validate_shared_name(stream) { + return Box::pin(async move { Err(e) }); + } + if let Err(e) = validate_local_name(consumer) { + return Box::pin(async move { Err(e) }); + } + let stream = stream.to_string(); + let consumer = consumer.to_string(); + Box::pin(async move { + if limit <= 0 { + return Ok(Vec::new()); + } + let from = self.get_offset_tx(&stream, &consumer).await?; + self.read_since_tx(&stream, from, limit).await + }) + } + + fn outbox_enqueue_tx<'a>( + &'a mut self, + outbox: &str, + opts: EnqueueOpts, + payload: Value, + ) -> BoxedFuture<'a, Result> { + if let Err(e) = validate_shared_name(outbox) { + return Box::pin(async move { Err(e) }); + } + let outbox = outbox.to_string(); + Box::pin(async move { + let client = self.client()?; + let bytes = encode_payload_bytes(&payload)?; + let queue = outbox_backing_queue_name(&outbox); + self.enqueue_row(client, &queue, bytes, &opts, outbox_default_stamps()) + .await + }) + } + + /// Commit: `COMMIT` + re-pool. A failed COMMIT discards the object + /// — its state is unknowable (the wave-3 review's finding (c), + /// pre-empted here; the pool accounting stays exact either way). + /// One-shot; later ops on the consumed shell fail `Closed`. + fn commit(mut self: Box) -> BoxedFuture<'static, Result<()>> { + let Some(object) = self.object.take() else { + return Box::pin(async { Err(Error::Closed) }); + }; + Box::pin(async move { + let client: &tokio_postgres::Client = &object; + match client.batch_execute("COMMIT").await { + Ok(()) => { + drop(object); + Ok(()) + } + Err(e) => { + drop(Object::take(object)); + Err(pg_error(e)) + } + } + }) + } +} + +impl Drop for PgTxHandle { + /// Drop = rollback (ADR-021 §4): the rollback runs on the + /// construction runtime as a detached task — `ROLLBACK`, then + /// re-pool; a failed rollback (or an already-closed client) + /// discards the object (unknowable state). Panic-safe: no panic + /// escapes drop (`Handle::spawn` is captured-handle-bound, never + /// current-context-bound), never blocking; the dropped-future arm + /// (runtime shut down before the task runs) closes the session — + /// the server aborts the open transaction, so no ghost survives + /// even there. + fn drop(&mut self) { + let Some(object) = self.object.take() else { + return; + }; + self.runtime.spawn(async move { + if object.is_closed() { + drop(Object::take(object)); + return; + } + let client: &tokio_postgres::Client = &object; + if let Err(e) = client.batch_execute("ROLLBACK").await { + drop(Object::take(object)); + eprintln!( + "alkstore: pg tx handle dropped with a failed ROLLBACK; \ + the pooled client is discarded (state unknowable): {e}" + ); + return; + } + drop(object); + }); + } +} + +/// Decode a job row (the uniform 18-column shape — the dead-table +/// select carries `last_error`/`died_at`; the live one binds NULLs) +/// into the contract value (`Job::from_row` — the `#[doc(hidden)]` +/// constructor, ADR-021 §2's field list of record). Malformed rows map +/// to `Error::Codec` (the decode-side posture guard). +fn job_from_row(row: &tokio_postgres::Row) -> Result { + let state = match row.get::<_, Option<&str>>(2) { + Some("pending") => JobState::Pending, + Some("processing") => JobState::Processing, + Some("dead") => JobState::Dead, + other => { + return Err(Error::Codec(format!( + "job row state must be pending|processing|dead, got {other:?}" + ))); + } + }; + let payload: Option> = row + .try_get(3) + .map_err(|e| Error::Codec(format!("job row payload must be bytea: {e}")))?; + Ok(Job::from_row( + row.get(0), + row.get(1), + state, + payload.unwrap_or_default(), + row.get(4), + row.get(5), + row.get(6), + row.get(7), + row.get(8), + row.get(9), + row.get(10), + row.get(11), + row.get(12), + row.get(13), + row.get(14), + row.get(15), + row.get(16), + row.get(17), + )) +} + +/// Decode a read page into the contract's [`StreamEvent`] values (the +/// core constructor's `from_row`). Malformed rows map to `Error::Codec`. +fn stream_events_from_rows( + rows: &[tokio_postgres::Row], + stream: String, +) -> Result> { + let mut out = Vec::with_capacity(rows.len()); + for row in rows { + let key: Option = row + .try_get(1) + .map_err(|e| Error::Codec(format!("stream row key must be text: {e}")))?; + let payload: Option> = row + .try_get(2) + .map_err(|e| Error::Codec(format!("stream row payload must be bytea: {e}")))?; + out.push(StreamEvent::from_row( + row.get(0), + stream.clone(), + key, + payload.unwrap_or_default(), + row.get(3), + )); + } + Ok(out) +} diff --git a/tasks/pg-engine-seam-tx.md b/tasks/pg-engine-seam-tx.md index ef1aa33..7a2a1c4 100644 --- a/tasks/pg-engine-seam-tx.md +++ b/tasks/pg-engine-seam-tx.md @@ -1,7 +1,7 @@ --- id: pg-engine-seam-tx name: Postgres engine — `begin_tx`, `TxHandle` (pooled-object handle), all eleven `*_tx` methods -status: pending +status: completed depends_on: [pg-engine-open-opts] scope: moderate risk: high @@ -82,28 +82,28 @@ curve, and the scheduler runner. ## Acceptance Criteria -- [ ] `begin_tx` checks out + `BEGIN`; concurrent `begin_tx` draws +- [x] `begin_tx` checks out + `BEGIN`; concurrent `begin_tx` draws distinct pooled objects (pool-bounded); pool exhaustion blocks (deadpool's acquire semantics) — tested -- [ ] All eleven `*_tx` methods implemented with entry-point +- [x] All eleven `*_tx` methods implemented with entry-point validation; writes commit-atomic, reads see the tx's own writes (read committed, both directions POC-pinned); opts resolution (delay-over-run_at, relative-expires, 300/3/5/none and 60/5/5 stamp sets) pinned by tests -- [ ] `commit` commits + re-pools; drop (without commit) rolls back + +- [x] `commit` commits + re-pools; drop (without commit) rolls back + re-pools — no-ghosts tested for job rows, stream events, notifications (a listener hears nothing from a rolled-back `notify_tx`), and offset saves -- [ ] Drop is panic-safe; failed COMMIT/ROLLBACK arms discard the +- [x] Drop is panic-safe; failed COMMIT/ROLLBACK arms discard the object (pool accounting exact); ops-after-consume → `Closed` -- [ ] `notify_tx`/`notify`'s 8000-byte check: typed +- [x] `notify_tx`/`notify`'s 8000-byte check: typed `PayloadTooLarge { limit: 8000 }` client-side, before any round trip; ≤ 8000 crosses -- [ ] `with_tx` (core's provided method) works end-to-end on this +- [x] `with_tx` (core's provided method) works end-to-end on this engine: Ok ⇒ commit, Err ⇒ rollback — tested -- [ ] `encode_payload`'s `Codec` error propagates; extent guard on the +- [x] `encode_payload`'s `Codec` error propagates; extent guard on the tx stream reads; `get_job_tx` queue-scoped -- [ ] `cargo test -p alkstore-postgres` (harness server), clippy +- [x] `cargo test -p alkstore-postgres` (harness server), clippy `-D warnings`, fmt clean; gates green server-less ## References @@ -122,8 +122,141 @@ curve, and the scheduler runner. ## Notes -> To be filled by implementation agent +Decisions of record made while implementing (the description didn't +pin them): + +- **The `Drop` = rollback is a detached-task rollback (the async-native + drop posture), not a synchronous await**: `Drop` cannot `.await`, + and the POC's `__private_api_rollback`-style sync submission is a + private tokio-postgres API. The handle captures the construction + runtime's `Handle`; `Drop` takes the held object and spawns the + teardown task — `ROLLBACK`, then re-pool by dropping the object + (re-pool = the `Object` drop arm; discard = `Object::take`, + deadpool's permanent-take, which reduces the pool size exactly). + Probes verified `Handle::spawn` is captured-handle-bound (never + panics on the captured handle's liveness — usable from any thread, + e.g. a sync-context drop) and the runtime-independent futures don't + panic either. The runtime-shutdown-before-task-runs arm drops the + object, whose connection closes server-side and **the server itself + aborts the open transaction** (probe-verified) — the no-ghosts + property holds even there. The `begin_concrete`/`pooled_client`/ + `consume_shell` `#[cfg(test)]` accessors are the tests' + fault-injection surface (poisoning a tx needs raw in-tx SQL; the + boxed trait object cannot downcast, ADR-007). +- **The unknowable-state discard arm is probe-pinned, and the + failing-COMMIT arm is induced with a deferred FK**: an *aborted* tx + commits fine server-side by design (`COMMIT` on an aborted tx = a + no-op success — probe-verified both protocols), so the failing-COMMIT + arm is induced with a DEFERRABLE INITIALLY DEFERRED FK violation + (surfaced at commit only); the failed commit discards the object + (`Object::take`) — the wave-3 review's finding (c) pre-empted, the + pool accounting exact (fresh draw proceeds, usable). The failed + `BEGIN` arm in `begin_tx` discards too (same posture); a + `ROLLBACK`ed-then-failed-ROLLBACK residue was probe-verified not to + strand a session. +- **The 8000-byte check's predicate is `len >= limit`, not `len > + limit`** — probe-pinned against the harness server (7999-byte text + delivers; 8000 is server-rejected both protocols; the pg payload + budget is the wire's 8000 bytes *counting the NUL terminator*). The + variant's `limit: 8000` stays the contract's pinned number + (ADR-008 §5); the check's quantity is contract-pinned (the + serde_json serialization — the stored-bytes twin, ADR-020 §4), and + the typed arm mirrors the server exactly so no oversized payload + ever surfaces as an opaque server-side reject. Contract text says + "≤ 8000 bytes"; the server's answer (and now ours) is "< 8000 + serialization bytes plus NUL ≤ 8000" — the largest *crossing* + serialization is 7999 bytes (a 7997-char bare string). Flagged here + as a core-text reconciliation note for wave 5's suite (the pg-side + rejection arm's row) — no variant number changes. +- **The clock is a `std::time` read, not a SQL-side clock function** + (the "pick one, document it" point): `resolution.rs`'s `now_unix()` + is the single read; one instant stamps `run_at`/`created_at`/ + `expires_at` as plain `i64` binds (second-precision integer + arithmetic, schema.rs's engine-wide timestamp posture). The SQL + layer never reads a clock in the enqueue/publish paths the tx seam + carries. +- **Dead rows are `get_job_tx`-visible** (the auto-commit `get_job` + shape, ADR-010 §1's "dead rows included"): the tx read is + queue-scoped *and* dead-diagnosis-carrying (`last_error`/`died_at`) + — two typed selects (dead table first, then live with NULL-cast + diagnosis columns), one decode owner (`job_from_row`, the + `#[doc(hidden)]` constructor's 18-column uniform shape). The + 18-column select's two `NULL::text`/`NULL::bigint` casts matter + (probe: untyped NULL carries no type info for `Option` decode). +- **Extended-protocol param typing is deliberate at three sites** + (probe-pinned: the server infers param types from context): + `GREATEST($3::bigint, 0)` (the bare `$3` beside the literal `0` + infers int4 — an i64 bind errors), the `LIMIT $3` bind rides the + inferred int8 fine, and `pg_notify($1, $2)`'s payload rides as + UTF-8 text (the server infers Text; binding `Vec` errors + WrongType — the notify payload's contract carriage *is* the + serialization's UTF-8 text, so no cast is needed client-side). +- **`save_offset_tx` clamps a negative save at 0 in the INSERT seed** + (`GREATEST($3::bigint, 0)`): the offsets table's domain CHECK + (`offset >= 0`, schema.rs) would abort the tx on a negative + first-save; clamping preserves the checkpoint's non-negative + domain and the monotone rule (a negative save at-or-below the + stored checkpoint is the same silent no-op a 0-save is). The + trait docs do not pin a negative-offset save's disposition — the + clamp is this engine's honest answer (typed never, no + tx-poisoning from a caller's negative cursor). +- **`enqueue_row` (the shared INSERT) and the trait-stub + `notify/listen/stream/queue/outbox/try_lock/schedule/unschedule/ + run_schedules` shapes** ride the open task's stub text unchanged; + only `begin_tx` is wired here (the seam task's scope). The open + task's `store_trait_methods_are_wiring_stubs` test updated: + `begin_tx` now boots + commits (the stub posture's replacement is + this task's landing), the stub-message assertion moved to the + `notify` stub (still `Database`, still "wiring lands with"). +- **Test harness details**: the dead-row plant rides typed `execute` + binds (the simple protocol — `batch_execute` — carries no + parameters; probe-pinned E42P02 "there is no parameter $1" is + exactly the simple/extended boundary); the tx-reads test's + poison-induction and the no-ghosts listener test use + schema-qualified quoted DDL (identifier quoting at the test's + composition sites, matching production's `quote_identifier` + posture). ## Summary -> To be filled on completion \ No newline at end of file +Landed alkstore-postgres's transactional seam: `resolution.rs` (the +engine-side arithmetic — the single std-time clock read, the +plain-queue 300/3/5/none and outbox 60/5/5 derived stamp sets, the +`__alkstore_outbox:{name}` reserved-prefix derivation, +`stamps_with_override`'s max-attempts-only override, and +`resolve_enqueue_opts`' delay-over-`run_at` + relative-`expires` +resolution, one clock read per enqueue — every formula unit-tested), +`tx.rs` (the `PgTxHandle` pooled-object handle: `begin` = pool +checkout + `BEGIN` with the exhaustion-blocking, fail-`Database` +arms; the shared `enqueue_row` INSERT both enqueue shapes ride; all +eleven `*_tx` methods implemented with entry-point validation before +any round trip, core's fallible `encode_payload` at the seam with the +typed `Codec` error `?`-ed, the client-side `PayloadTooLarge { limit: +8000 }` check at `len >= 8000` (the server-mirroring, NUL-counting +predicate — probe-pinned), the monotone `save_offset_tx` upsert +clamp-seeded at 0, queue-scoped dead-visible `get_job_tx`, the +extent guards (`limit <= 0` → empty `Vec`) on both tx stream reads, +`commit` = `COMMIT` + re-pool with the failed arm discarding +(`Object::take`), and drop = rollback via the detached-task teardown +(RAII paths roll back through a detached `ROLLBACK`+re-pool; +failed/`is_closed` arms discard; the runtime-gone arm's session +close aborts server-side — no ghosts on any arm) and +ops-after-consume failing closed `Error::Closed`), and the +`PgStore::begin_tx` wiring replacing its stub. Verified against the +harness server: 16 new tx-seam tests + 8 resolution unit tests (39 +lib tests green ×2 runs) covering the acceptance rows — pool bounds ++ blocked begin, the stamp/opts resolutions, read-your-own-writes +both directions, keyed publish + monotone offsets + silent-no-op +saves, the uniform no-ghosts list (jobs/events/notifies/offsets — a +listener hears nothing from a rolled-back notify and delivers a +committed one, native commit-atomicity), the unknowable-state arms +(deferred-FK-induced failing COMMIT → discard + exact accounting), +the fail-closed shell (all eleven ops + double-commit `Closed`), +panic-through + mid-tx-cancellation rollbacks, `with_tx` Ok⇒commit / +Err⇒rollback, the typed 8000-byte boundary (crossing 7999 vs typed +rejects at 8000 and 8002), entry-point validation on every name +(shared kinds reserved+empty, consumer-local kinds empty — residue- +free rejects), and the 8-parallel-tx exactly-once fan-out; +workspace `cargo test` green server-less (236→260 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