Postgres engine: transactional seam — begin_tx, PgTxHandle, all eleven *_tx methods with drop=rollback detached teardown, probe-pinned unknowable-state discard arms, engine-side resolution arithmetic (task pg-engine-seam-tx)

This commit is contained in:
glm-5.3-flash committed 2026-10-09 06:50:53 +00:00
1 parent c6a7eeaa45
commit cf5ceea70e
7 files changed
+2318 -29

No files matched your search

+3
View File
@@ -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;
+227
View File
@@ -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<i64>,
}
/// 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<i64>,
}
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));
}
}
+4 -1
View File
@@ -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;
+16 -17
View File
@@ -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);
File diff suppressed because it is too large. Load diff
+661
View File
@@ -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<Object>,
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<dyn TxHandle + Send>>> {
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<dyn TxHandle + Send>)
})
}
/// 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<PgTxHandle>> {
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<u8>,
opts: &EnqueueOpts,
base: crate::resolution::Stamps,
) -> Result<i64> {
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<Vec<u8>> {
alkstore::encode_payload(payload)
}
impl TxHandle for PgTxHandle {
fn enqueue_tx<'a>(
&'a mut self,
queue: &str,
opts: EnqueueOpts,
payload: Value,
) -> BoxedFuture<'a, Result<i64>> {
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<i64>> {
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<String>,
payload: Value,
) -> BoxedFuture<'a, Result<i64>> {
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<Option<Job>>> {
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<i64>> {
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<Vec<StreamEvent>>> {
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<Vec<StreamEvent>>> {
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<i64>> {
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<Self>) -> 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<Job> {
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<Vec<u8>> = 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<Vec<StreamEvent>> {
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let key: Option<String> = row
.try_get(1)
.map_err(|e| Error::Codec(format!("stream row key must be text: {e}")))?;
let payload: Option<Vec<u8>> = 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)
}