Postgres engine: queues — Queue (validated constructor over the handle's QueueOpts stamps), the FOR UPDATE SKIP LOCKED claim (one statement, the pre-claim exhausted-reclaimable sweep in the atomic claim frame), JobHandle (one-shot ack/retry/fail in tx frames under the uniform validity predicate, absolute-reset heartbeat, engine-side equal-jitter backoff), dead-letter moves transactional (the #133 class excluded), worker-less ack_batch, unconditional cancel, dead-visible get_job, both-states+retention sweep, best-effort queue-channel wake, wake-driven claim loop pinned (task pg-engine-queues)

This commit is contained in:
glm-5.3-flash committed 2026-10-09 08:44:59 +00:00
1 parent cb067bb4be
commit 3dd83791ff
7 files changed
+2640 -53

No files matched your search

+1
View File
@@ -35,6 +35,7 @@
mod forwarder;
mod notify;
mod opts;
mod queue;
mod resolution;
mod schema;
mod seam;
+999
View File
@@ -0,0 +1,999 @@
//! The queue mechanism's Postgres arm (`pg-engine-queues`) — the
//! auto-commit counterpart to the tx seam's enqueue; the re-derived
//! queue machinery (ADR-004/005's greenfield posture, the pg-boss
//! schema family as design reference) with ADR-010's full depth:
//!
//! - `queue(name, opts)` — the validated constructor (construction is
//! free: validation + a closed-store check — queues are names, not
//! registered objects, ADR-010 §3a); the handle carries its
//! [`QueueOpts`] stamps future enqueues resolve over.
//! - `enqueue(payload, opts)` — the one handle-carrying shape: stamps
//! resolve over **this handle's opts** (`EnqueueOpts::max_attempts`
//! the one per-job override — ADR-020 §3a), through the shared
//! resolution arithmetic (one owner per formula — the engine's own
//! `resolution.rs`; one clock read per enqueue, delay-over-`run_at`,
//! relative-`expires`), `encode_payload` at the seam, the INSERT (the
//! tx seam's shared row-write) on a pool connection — then **wake**:
//! `pg_notify` on the queue's wake channel (the
//! mechanism-name-is-the-channel realization; the stream task's twin
//! posture — best-effort, logged and swallowed: the durable row is
//! the truth).
//! - `claim_one` / `claim_batch` — exactly-once handout under
//! concurrency (ADR-010 §1): the `FOR UPDATE SKIP LOCKED` claim on
//! a pool connection, one statement — an `UPDATE … RETURNING` over a
//! `picked` CTE. Ordering `priority DESC, run_at ASC, id`
//! (enqueue order = id); `attempts += 1` per claim (a reclaim
//! consumes an attempt); the claim sets the row's deadline from the
//! job's own stamped `visibility_timeout_s`; the pre-claim
//! dead-letter sweep for exhausted reclaimables runs first (the
//! substrate's D-12 shape — a worker dying right after its last
//! allowed claim dead-letters at the next claim, not strands);
//! `worker_id` stamps the claimant column (consumer-local,
//! non-empty — ADR-019 §2). The **extent guard** (ADR-023 §2):
//! `n <= 0` → the empty `Vec` at trait-impl entry. Claims return
//! boxed [`JobHandle`]s (ops + row value), decoded through the one
//! 18-column decode owner ([`job_from_row`] — the tx reads' shared
//! shape).
//! - The [`JobHandle`] impl: `job()` (the claimed row value),
//! `ack`/`retry`/`fail` (one-shot, `self: Box<Self>` — the consuming
//! shape; `BoxedFuture<'static>`), `heartbeat(extend)` (repeatable,
//! absolute reset). The uniform validity predicate (ADR-010 §2 —
//! `processing` + unexpired claim deadline) is conjunct in every
//! op's SQL (`worker_id` keeps its conjunct on the single-row
//! forms — the D-31 delta removed the filter from the *batch* ack
//! only); refusals are `Ok(false)`, never errors.
//! - **The backoff curve — engine-side, one owner (ADR-010 §3)**:
//! `retry(err, None)` computes the equal-jitter exponential from the
//! job's stamped `backoff_base_s`: `delay ∈ [base·2^(a−1)/2,
//! base·2^(a−1)]`, capped at 1 hour; the attempt index `a` is the
//! row's `attempts` at the retry (the claim-row value — post-claim
//! count, reclaims included). `Some(d)` passes through. The jitter's
//! randomness source is std-only (the SQLite task's reused posture —
//! a fresh [`RandomState`] hash draw, no `rand` dep; flagged in the
//! task Notes).
//! - Dead-letter moves (retry-at-budget, fail, the pre-claim sweep,
//! the expiry sweep) ride one dead-move INSERT; dead rows are
//! `get_job`-visible (ADR-010 §4). The engine-triggered default
//! strings are the contract's pinned ones (`"max attempts
//! exceeded"`, `"failed"`, `"expired"`); caller strings cross
//! verbatim.
//! - `ack_batch(ids)` — the batch ack applies the uniform validity
//! predicate per id (no worker filter — the wave-3 D-31 delta's
//! shape is the contract's); count returned, partial success
//! ordinary (ADR-010 §1).
//! - `cancel(job_id)` — unconditional delete (pending or processing,
//! not an interrupt); `Ok(false)` when row-less.
//! - `get_job(job_id)` — dead-visible pure read on a pool connection
//! (the resolved values — not raw opts — are what the row shows,
//! ADR-020).
//! - `sweep_expired()` — both-states no-stranded-rows move (any
//! past-expiry row → dead with `"expired"`) + retention enforcement
//! (ADR-010 §5), single-statement atomic; returns the moved count.
//! No leader lock required.
use std::collections::hash_map::RandomState;
use std::hash::BuildHasher;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{SystemTime, UNIX_EPOCH};
use deadpool_postgres::Pool;
use serde_json::Value;
use alkstore::{
BoxedFuture, EnqueueOpts, Error, Job, JobHandle, JobState, Queue, QueueOpts, Result,
validate_local_name,
};
use deadpool_postgres::Object;
use crate::resolution::{Stamps, now_unix};
use crate::schema::{QualifiedTable, tables};
use crate::seam::{database_error, pg_error, pool_error};
use crate::tx::{encode_payload_bytes, enqueue_row};
/// The dead-letter default strings (row content, ADR-010 §4 / ADR-019
/// §3's annotation), identical to the SQLite engine's — the
/// engine-default strings land identically on both engines (the
/// backlog row's wording).
pub(crate) const EXHAUSTED_ERROR: &str = "max attempts exceeded";
pub(crate) const FAILED_ERROR: &str = "failed";
pub(crate) const EXPIRED_ERROR: &str = "expired";
/// The backoff cap (ADR-010 §3 — a pinned contract constant, no knob).
const BACKOFF_CAP_S: i64 = 3600;
/// The queue-scoped handle (see the module docs).
pub struct PgQueueHandle {
name: String,
opts: QueueOpts,
pool: Pool,
schema: String,
closed: Arc<AtomicBool>,
}
impl std::fmt::Debug for PgQueueHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PgQueueHandle")
.field("name", &self.name)
.field("opts", &self.opts)
.finish_non_exhaustive()
}
}
impl PgQueueHandle {
pub(crate) fn new(
name: String,
opts: QueueOpts,
pool: Pool,
schema: String,
closed: Arc<AtomicBool>,
) -> Self {
Self {
name,
opts,
pool,
schema,
closed,
}
}
/// The handle's opts, as the rows stamp them (the `max_attempts`
/// override resolves per enqueue over these).
fn base_stamps(&self) -> Stamps {
Stamps {
max_attempts: self.opts.max_attempts,
visibility_timeout_s: self.opts.visibility_timeout_s,
backoff_base_s: self.opts.backoff_base_s,
dead_letter_retention_s: self.opts.dead_letter_retention_s,
}
}
}
/// The equal-jitter exponential backoff (ADR-010 §3 — engine-side, one
/// owner; the SQLite engine's twin computes identically and wave 5
/// pins equivalence): `delay ∈ [base·2^(a−1)/2, base·2^(a−1)]` from
/// the job's stamped `backoff_base_s`, capped at **1 hour**.
///
/// Integerized second-precision: the draw is uniform over the integers
/// in `[ceil(half), cap]` (inclusive both ends) — never below the
/// pinned lower bound, never above the cap. The attempt index `a` is
/// the row's `attempts` at the retry (post-claim count, reclaims
/// included); the exponent clamps at 30 so the cap arithmetic
/// saturates rather than overflows.
///
/// The jitter draw is std-only (no `rand` dependency without a
/// decision — the ADR-005 dependency posture; the SQLite engine's
/// reused approach, flagged in the task Notes): a fresh [`RandomState`]
/// hasher (per-process-random seeds) over the row id, the attempt
/// count, and the subsecond clock spreads the draw uniformly over the
/// range. Hash quality is only asked to spread the draw; the pinned
/// property is the **range**, what the acceptance test pins
/// statistically.
pub(crate) fn backoff_delay_s(base_s: i64, attempts: i64, id: i64) -> i64 {
let exponent: u32 = (attempts.max(1) as u32).saturating_sub(1).min(30);
let doubling: i64 = base_s.saturating_mul(1i64 << exponent).max(1);
let cap = doubling.min(BACKOFF_CAP_S);
let half = ((cap + 1) / 2).max(1);
let span = (cap - half + 1).max(1);
let now_nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.subsec_nanos())
.unwrap_or(0);
let draw = RandomState::new().hash_one((id, attempts, now_nanos));
half + (draw % span as u64) as i64
}
/// Qualified table reference for this engine-owned schema.
fn table(schema: &str, name: &str) -> String {
QualifiedTable { schema, name }.to_string()
}
/// The live-row select's column list (the uniform 18-column shape the
/// dead-table select mirrors with `last_error`/`died_at` in place of
/// the NULL casts) — one shape, shared with the tx reads.
fn live_columns() -> &'static str {
"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"
}
/// The dead-row select's column list (the 18-column live shape with
/// `last_error`/`died_at` riding the NULL casts' slots).
fn dead_columns() -> &'static str {
"id, queue, 'dead' AS 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,
last_error, died_at"
}
/// The 14 columns a dead-move carries (the live row minus state, the
/// deadline, and the claimant-clearing arms — the move deletes the row
/// whole and re-inserts it with its history).
const DEAD_MOVE_COLUMNS: &str = "id, queue, payload, priority, run_at, attempts,
max_attempts, worker_id, claimed_at,
created_at, expires_at, visibility_timeout_s,
backoff_base_s, dead_letter_retention_s";
/// Decode the dead-move row (the [`DEAD_MOVE_COLUMNS`] shape) into the
/// dead table with `last_error`/`died_at` stamped. One owner: every
/// dead-letter move (retry-at-budget, `fail`, the pre-claim sweep, the
/// expiry sweep) rides it (ADR-010 §4).
async fn insert_dead_row_tx(
client: &tokio_postgres::Client,
schema: &str,
row: &tokio_postgres::Row,
last_error: &str,
died_at: i64,
) -> Result<()> {
let params: [&(dyn tokio_postgres::types::ToSql + Sync); 16] = [
&row_get_i64(row, 0),
&row_get_string(row, 1),
&row_get_bytes(row, 2),
&row_get_i64(row, 3),
&row_get_i64(row, 4),
&row_get_i64(row, 5),
&row_get_i64(row, 6),
&row_get_opt_string(row, 7),
&row_get_opt_i64(row, 8),
&row_get_i64(row, 9),
&row_get_opt_i64(row, 10),
&row_get_i64(row, 11),
&row_get_i64(row, 12),
&row_get_opt_i64(row, 13),
&last_error.to_string(),
&died_at,
];
client
.execute(
&format!(
"INSERT INTO {dead}
(id, queue, payload, priority, run_at, attempts,
max_attempts, worker_id, claimed_at, created_at,
expires_at, visibility_timeout_s, backoff_base_s,
dead_letter_retention_s, last_error, died_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10,
$11, $12, $13, $14, $15, $16)",
dead = table(schema, tables::DEAD)
),
&params[..],
)
.await
.map_err(pg_error)?;
Ok(())
}
// The dead-move row's typed column reads (the RETURNING row's decode —
// the columns arrive typed from the live table; these helpers pull
// them out as owned values so they outlive the row borrow across the
// INSERT's params borrow).
fn row_get_i64(row: &tokio_postgres::Row, idx: usize) -> i64 {
row.get(idx)
}
fn row_get_opt_i64(row: &tokio_postgres::Row, idx: usize) -> Option<i64> {
row.get(idx)
}
fn row_get_string(row: &tokio_postgres::Row, idx: usize) -> String {
row.get(idx)
}
fn row_get_opt_string(row: &tokio_postgres::Row, idx: usize) -> Option<String> {
row.get(idx)
}
fn row_get_bytes(row: &tokio_postgres::Row, idx: usize) -> Vec<u8> {
row.get(idx)
}
/// The dead-move / claim transaction frame on one pooled object: an
/// explicit `BEGIN` … body … `COMMIT`/`ROLLBACK` via `batch_execute`
/// on the held client (the seam's pooled-object posture
/// ([`crate::tx::PgTxHandle`]) scoped to one op — the sqlite
/// substrate's savepoint frame's pg twin; every statement the body
/// runs through the client participates in the open frame). The arms:
///
/// - a failed `BEGIN` discards the object (unknowable state —
/// [`Object::take`]; the pool accounting stays exact);
/// - the body's error arm rolls back explicitly and re-pools (a failed
/// `ROLLBACK` discards the object too);
/// - a failed `COMMIT` discards the object (the wave-3 review's
/// finding (c), pre-empted by the seam).
async fn in_tx<T, F>(object: Object, body: F) -> Result<T>
where
T: Send,
F: FnOnce(&tokio_postgres::Client) -> BoxedFuture<'_, Result<T>>,
{
let client: &tokio_postgres::Client = &object;
if let Err(e) = client.batch_execute("BEGIN").await {
drop(Object::take(object));
return Err(pg_error(e));
}
match body(client).await {
Ok(value) => match client.batch_execute("COMMIT").await {
Ok(()) => {
drop(object);
Ok(value)
}
Err(e) => {
drop(Object::take(object));
Err(pg_error(e))
}
},
Err(e) => {
if client.batch_execute("ROLLBACK").await.is_err() {
drop(Object::take(object));
} else {
drop(object);
}
Err(e)
}
}
}
/// The pre-claim dead-letter sweep for exhausted reclaimables (the
/// substrate's D-12 shape — re-owned as owned pg SQL): rows at-or-past
/// budget that the claim predicate would otherwise reach are moved to
/// dead first, so the claim statement never hands out a row whose
/// next ack path is already exhausted. Same transaction frame as the
/// move paths (the sqlite substrate's savepoint-guarded #133 defect
/// class is excluded here the same way: a failing dead INSERT rolls
/// the whole move back — no stranded-in-neither-table state).
async fn sweep_exhausted_reclaimables_tx(
client: &tokio_postgres::Client,
schema: &str,
queue: &str,
now: i64,
) -> Result<()> {
let sql = format!(
"DELETE FROM {job}
WHERE queue = $1
AND attempts >= max_attempts
AND (expires_at IS NULL OR expires_at > $2)
AND ((state = 'pending' AND run_at <= $2)
OR (state = 'processing' AND claim_expires_at < $2))
RETURNING {DEAD_MOVE_COLUMNS}",
job = table(schema, tables::JOB)
);
let rows = client
.query(&sql, &[&queue, &now])
.await
.map_err(pg_error)?;
for row in &rows {
insert_dead_row_tx(client, schema, row, EXHAUSTED_ERROR, now).await?;
}
Ok(())
}
/// The claim statement: one `UPDATE … RETURNING` over the claimable
/// predicate under `FOR UPDATE SKIP LOCKED` (the exactly-once handout,
/// ADR-010 §1). The claim stamps `processing`, the claimant, the
/// deadline-from-the-job's-own-stamp, `claimed_at`, and
/// `attempts += 1`.
///
/// Parameter positions: $1 worker_id, $2 now, $3 queue, $4 extent n.
fn claim_stmt_sql(schema: &str) -> String {
format!(
"UPDATE {job} SET
state = 'processing',
worker_id = $1,
claim_expires_at = $2::bigint + visibility_timeout_s,
claimed_at = $2,
attempts = attempts + 1
WHERE id IN (
SELECT id FROM {job}
WHERE queue = $3
AND state IN ('pending', 'processing')
AND attempts < max_attempts
AND (expires_at IS NULL OR expires_at > $2)
AND ((state = 'pending' AND run_at <= $2)
OR (state = 'processing' AND claim_expires_at < $2))
ORDER BY priority DESC, run_at ASC, id ASC
FOR UPDATE SKIP LOCKED
LIMIT $4
)
RETURNING {live}",
job = table(schema, tables::JOB),
live = live_columns(),
)
}
/// Decode an 18-column job row into the contract value
/// ([`Job::from_row`] — the `#[doc(hidden)]` constructor's field list
/// of record, ADR-021 §2). Malformed rows map to `Error::Codec` (the
/// decode-side posture guard). One decode owner: the tx reads and this
/// module's claim/get rows share the shape.
pub(crate) 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),
))
}
/// The closed-store error shape (the engine-wide posture — one
/// message).
fn closed_error() -> Error {
database_error("the store is closed")
}
impl Queue for PgQueueHandle {
fn name(&self) -> &str {
&self.name
}
fn enqueue<'a>(&'a self, payload: Value, opts: EnqueueOpts) -> BoxedFuture<'a, Result<i64>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_error()) });
}
let name = self.name.clone();
let base = self.base_stamps();
let pool = self.pool.clone();
let schema = self.schema.clone();
Box::pin(async move {
// The seam: encode_payload's typed `Codec` error `?`-ed
// before any round trip (ADR-023 §1).
let bytes = encode_payload_bytes(&payload)?;
let conn = pool.get().await.map_err(pool_error)?;
// The shared resolution arithmetic — one clock read per
// enqueue, delay-over-`run_at`, relative-`expires` — and
// the INSERT (the tx seam's shared row-write; the
// `max_attempts` override rides `stamps_with_override`
// inside `enqueue_row`).
let id = enqueue_row(&conn, &schema, &name, bytes, &opts, base).await?;
// Then wake: pg_notify on the queue's wake channel — the
// mechanism-name-is-the-channel realization (the stream
// task's twin posture). Best-effort on this auto-commit
// path: the durable row is the truth; a failed wake is
// logged and swallowed and the consumer's re-poll safety
// net covers the gap — no wake may ever fail (or stall) a
// committed enqueue.
if let Err(wake_err) = conn.query_one("SELECT pg_notify($1, '')", &[&name]).await {
eprintln!(
"alkstore-postgres: queue enqueue wake failed for {name:?} \
(the durable row landed; the re-poll safety net covers the gap): {wake_err}"
);
}
Ok(id)
})
}
fn claim_one<'a>(
&'a self,
worker_id: &str,
) -> BoxedFuture<'a, Result<Option<Box<dyn JobHandle>>>> {
let worker_id = worker_id.to_string();
Box::pin(async move { Ok(self.claim_batch(&worker_id, 1).await?.into_iter().next()) })
}
fn claim_batch<'a>(
&'a self,
worker_id: &str,
n: i64,
) -> BoxedFuture<'a, Result<Vec<Box<dyn JobHandle>>>> {
if let Err(e) = validate_local_name(worker_id) {
return Box::pin(async move { Err(e) });
}
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_error()) });
}
// The extent guard (ADR-023 §2): `n <= 0` claims nothing — the
// empty `Vec`, a value, never an error — at trait-impl entry,
// before any round trip.
if n <= 0 {
return Box::pin(async { Ok(Vec::new()) });
}
let name = self.name.clone();
let pool = self.pool.clone();
let schema = self.schema.clone();
let worker = worker_id.to_string();
let closed = self.closed.clone();
Box::pin(async move {
let conn = pool.get().await.map_err(pool_error)?;
// The claim's statement frame: the pre-claim sweep (the
// D-12 shape — at-budget rows the claim predicate would
// reach dead-letter first, no row stranded between the
// two statements' failure arms) + the claim (one
// `UPDATE … SKIP LOCKED` — exactly-once handout under
// concurrency; `now` is the one clock read the frame's
// stamping shares), both committing atomically.
in_tx(conn, |client| {
Box::pin(async move {
let now = now_unix();
sweep_exhausted_reclaimables_tx(client, &schema, &name, now).await?;
let rows = client
.query(&claim_stmt_sql(&schema), &[&worker, &now, &name, &n])
.await
.map_err(pg_error)?;
let mut out = Vec::with_capacity(rows.len());
for row in &rows {
let job = job_from_row(row)?;
out.push(Box::new(PgJobHandle {
job,
pool: pool.clone(),
schema: schema.clone(),
closed: closed.clone(),
}) as Box<dyn JobHandle>);
}
Ok(out)
})
})
.await
})
}
fn ack_batch<'a>(&'a self, ids: &[i64]) -> BoxedFuture<'a, Result<i64>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_error()) });
}
let name = self.name.clone();
let pool = self.pool.clone();
let schema = self.schema.clone();
let ids: Vec<i64> = ids.to_vec();
Box::pin(async move {
// An empty batch is vacuously zero, no round trip.
if ids.is_empty() {
return Ok(0);
}
let now = now_unix();
let conn = pool.get().await.map_err(pool_error)?;
// The batch ack carries no worker filter (the wave-3 D-31
// delta's shape is the contract's, ADR-19 §1): the uniform
// validity predicate — processing + unexpired claim
// deadline — per id; lapsed, non-claimed, acked-elsewhere,
// or cancelled ids silently not counted; count returned,
// partial success ordinary (ADR-010 §1). Queue-scoped (the
// handle's name — a queue handle acks its own rows).
let n = conn
.execute(
&format!(
"DELETE FROM {job}
WHERE id = ANY($1)
AND queue = $2
AND state = 'processing'
AND claim_expires_at >= $3",
job = table(&schema, tables::JOB)
),
&[&ids, &name, &now],
)
.await
.map_err(pg_error)?;
Ok(n as i64)
})
}
fn cancel<'a>(&'a self, job_id: i64) -> BoxedFuture<'a, Result<bool>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_error()) });
}
let name = self.name.clone();
let pool = self.pool.clone();
let schema = self.schema.clone();
Box::pin(async move {
let conn = pool.get().await.map_err(pool_error)?;
// Unconditional (ADR-010 §1): either state, any holder —
// not an interrupt; the holder's next op refuses (the row
// just is gone — the same shape as a deadline lapse).
// Queue-scoped (the handle's name).
let n = conn
.execute(
&format!(
"DELETE FROM {job}
WHERE id = $1 AND queue = $2
AND state IN ('pending', 'processing')",
job = table(&schema, tables::JOB)
),
&[&job_id, &name],
)
.await
.map_err(pg_error)?;
Ok(n > 0)
})
}
fn get_job<'a>(&'a self, job_id: i64) -> BoxedFuture<'a, Result<Option<Job>>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_error()) });
}
let name = self.name.clone();
let pool = self.pool.clone();
let schema = self.schema.clone();
Box::pin(async move {
let conn = pool.get().await.map_err(pool_error)?;
// Queue-scoped + dead-visible (the auto-commit shape,
// ADR-010 §1): dead table first (the diagnosis columns
// ride), then live — the tx seam's `get_job_tx` shape,
// two typed selects, one decode owner ([`job_from_row`]).
let row = conn
.query_opt(
&format!(
"SELECT {dead} FROM {deadt} WHERE id = $1 AND queue = $2",
dead = dead_columns(),
deadt = table(&schema, tables::DEAD)
),
&[&job_id, &name],
)
.await
.map_err(pg_error)?;
if let Some(row) = row {
return Ok(Some(job_from_row(&row)?));
}
let row = conn
.query_opt(
&format!(
"SELECT {live} FROM {jobt} WHERE id = $1 AND queue = $2",
live = live_columns(),
jobt = table(&schema, tables::JOB)
),
&[&job_id, &name],
)
.await
.map_err(pg_error)?;
match row {
Some(row) => Ok(Some(job_from_row(&row)?)),
None => Ok(None),
}
})
}
fn sweep_expired<'a>(&'a self) -> BoxedFuture<'a, Result<i64>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_error()) });
}
let name = self.name.clone();
let pool = self.pool.clone();
let schema = self.schema.clone();
Box::pin(async move {
let now = now_unix();
let conn = pool.get().await.map_err(pool_error)?;
// The no-stranded-rows property (ADR-010 §5): every
// past-`expires_at` row (any state, pending and processing
// alike) moves to dead with `"expired"`, and the dead rows
// past their retention TTL delete (the only sweeper dead
// rows ever have — nothing runs without a caller). One
// transaction frame (the move is atomic — a failing dead
// INSERT rolls the whole move back, no stranded-in-
// neither-table state); no leader lock.
in_tx(conn, |client| {
Box::pin(async move {
let move_sql = format!(
"DELETE FROM {job}
WHERE queue = $1
AND expires_at IS NOT NULL
AND expires_at <= $2
RETURNING {DEAD_MOVE_COLUMNS}",
job = table(&schema, tables::JOB)
);
let rows = client
.query(&move_sql, &[&name, &now])
.await
.map_err(pg_error)?;
for row in &rows {
insert_dead_row_tx(client, &schema, row, EXPIRED_ERROR, now).await?;
}
let retained = client
.execute(
&format!(
"DELETE FROM {dead}
WHERE queue = $1
AND dead_letter_retention_s IS NOT NULL
AND died_at <= $2::bigint - dead_letter_retention_s",
dead = table(&schema, tables::DEAD)
),
&[&name, &now],
)
.await
.map_err(pg_error)?;
Ok(rows.len() as i64 + retained as i64)
})
})
.await
})
}
}
/// The claim handle (see the module docs). The one-shot ops consume
/// `self: Box<Self>` (the whole handle moves into the `'static`
/// future); `heartbeat` is the deliberate repeatable `&self`
/// counter-case (ADR-019 §3 — do not "fix" the asymmetry).
pub struct PgJobHandle {
job: Job,
pool: Pool,
schema: String,
closed: Arc<AtomicBool>,
}
impl std::fmt::Debug for PgJobHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PgJobHandle")
.field("job", &self.job)
.finish_non_exhaustive()
}
}
impl PgJobHandle {
/// The stored job's live row re-read under the uniform validity
/// predicate (ADR-010 §2): `processing` + the caller's unexpired
/// claim deadline + the caller's claimant. `None` = refused (a
/// lapse, a reclaim, a cancel, or an ack elsewhere consumed the
/// window) — the ops' `Ok(false)` arm; the row is untouched.
///
/// The move-carrying 14-column shape (the dead-move's decode).
/// Run inside the op's transaction frame (`FOR UPDATE`-locked by
/// the subsequent DELETE/UPDATE of the same row in that frame —
/// the frame is the op's atomicity).
async fn fetch_move_row(
client: &tokio_postgres::Client,
schema: &str,
job_id: i64,
worker_id: &str,
now: i64,
) -> Result<Option<tokio_postgres::Row>> {
let row = client
.query_opt(
&format!(
"SELECT {move_cols} FROM {job}
WHERE id = $1 AND worker_id = $2
AND state = 'processing'
AND claim_expires_at >= $3
FOR UPDATE",
move_cols = DEAD_MOVE_COLUMNS,
job = table(schema, tables::JOB)
),
&[&job_id, &worker_id, &now],
)
.await
.map_err(pg_error)?;
Ok(row)
}
}
impl JobHandle for PgJobHandle {
fn job(&self) -> &Job {
&self.job
}
fn ack(self: Box<Self>) -> BoxedFuture<'static, Result<bool>> {
let PgJobHandle {
job,
pool,
schema,
closed,
} = *self;
let id = job.id;
let worker_id = job.worker_id.unwrap_or_default();
Box::pin(async move {
let _ = closed;
let now = now_unix();
let conn = pool.get().await.map_err(pool_error)?;
// Delete-on-ack (ADR-010 §1) under the validity predicate
// (ADR-010 §2): 0 rows = refused (lapsed, reclaimed,
// cancelled, or acked elsewhere) — a value, never an
// error.
let n = conn
.execute(
&format!(
"DELETE FROM {job}
WHERE id = $1 AND worker_id = $2
AND state = 'processing'
AND claim_expires_at >= $3",
job = table(&schema, tables::JOB)
),
&[&id, &worker_id, &now],
)
.await
.map_err(pg_error)?;
Ok(n > 0)
})
}
fn retry(
self: Box<Self>,
err: Option<String>,
delay: Option<i64>,
) -> BoxedFuture<'static, Result<bool>> {
let PgJobHandle {
job,
pool,
schema,
closed,
} = *self;
let id = job.id;
let worker_id = job.worker_id.unwrap_or_default();
// `None` = the queue's curve (ADR-010 §3) from the job's
// stamped base; the attempt index is the claim-row's
// post-claim count (reclaims included).
let delay_s = match delay {
Some(d) => d,
None => backoff_delay_s(job.backoff_base_s, job.attempts, id),
};
// The exhausted path's string: `None` err = the pinned
// `"max attempts exceeded"` (the engine-observed path's owned
// string wins over the caller-None case, per ADR-019 §3's
// annotation); `Some(err)` rides the row verbatim.
let exhaust_error = err.unwrap_or_else(|| EXHAUSTED_ERROR.to_string());
Box::pin(async move {
let _ = closed;
let conn = pool.get().await.map_err(pool_error)?;
// One transaction frame for the whole op: the validity
// predicate (the row re-read at 0 rows = refused — the row
// keeps its pre-call state; at-least-once redelivery
// governs) + the reschedule or the dead-letter move, all
// committing atomically.
in_tx(conn, |client| {
Box::pin(async move {
let now = now_unix();
let Some(row) =
PgJobHandle::fetch_move_row(client, &schema, id, &worker_id, now).await?
else {
return Ok(false);
};
let attempts: i64 = row.get(5);
let max_attempts: i64 = row.get(6);
if attempts >= max_attempts {
// The exhausted-budget dead-letter move
// (ADR-010 §4): the dead INSERT + the live
// DELETE under the same predicate, atomic in
// the frame (the sqlite substrate's #133
// defect class excluded — a failing INSERT
// rolls the whole op back, no
// stranded-in-neither-table state).
insert_dead_row_tx(client, &schema, &row, &exhaust_error, now).await?;
client
.execute(
&format!(
"DELETE FROM {job}
WHERE id = $1 AND worker_id = $2
AND state = 'processing'
AND claim_expires_at >= $3",
job = table(&schema, tables::JOB)
),
&[&id, &worker_id, &now],
)
.await
.map_err(pg_error)?;
} else {
// Reschedule pending: the curve draw (or the
// caller's override) is the new ready time from
// now; the claim state clears (deadline,
// claimant, claimed_at).
client
.execute(
&format!(
"UPDATE {job} SET
state = 'pending',
run_at = $2::bigint + $3::bigint,
worker_id = NULL,
claim_expires_at = NULL,
claimed_at = NULL
WHERE id = $1 AND worker_id = $4
AND state = 'processing'
AND claim_expires_at >= $2",
job = table(&schema, tables::JOB)
),
&[&id, &now, &delay_s, &worker_id],
)
.await
.map_err(pg_error)?;
}
Ok(true)
})
})
.await
})
}
fn fail(self: Box<Self>, err: Option<String>) -> BoxedFuture<'static, Result<bool>> {
let PgJobHandle {
job,
pool,
schema,
closed,
} = *self;
let id = job.id;
let worker_id = job.worker_id.unwrap_or_default();
let error = err.unwrap_or_else(|| FAILED_ERROR.to_string());
Box::pin(async move {
let _ = closed;
let conn = pool.get().await.map_err(pool_error)?;
// One transaction frame for the whole op: the validity
// predicate first (0 rows = refused), then the immediate
// dead-letter move (ADR-010 §4) — `fail` is the "stop
// retrying this" operator, budget-unaware; the move is
// atomic in the frame.
in_tx(conn, |client| {
Box::pin(async move {
let now = now_unix();
let Some(row) =
PgJobHandle::fetch_move_row(client, &schema, id, &worker_id, now).await?
else {
return Ok(false);
};
insert_dead_row_tx(client, &schema, &row, &error, now).await?;
client
.execute(
&format!(
"DELETE FROM {job}
WHERE id = $1 AND worker_id = $2
AND state = 'processing'
AND claim_expires_at >= $3",
job = table(&schema, tables::JOB)
),
&[&id, &worker_id, &now],
)
.await
.map_err(pg_error)?;
Ok(true)
})
})
.await
})
}
fn heartbeat<'a>(&'a self, extend: i64) -> BoxedFuture<'a, Result<bool>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_error()) });
}
let pool = self.pool.clone();
let schema = self.schema.clone();
let id = self.job.id;
let worker_id = self.job.worker_id.clone().unwrap_or_default();
Box::pin(async move {
// Renewal — an absolute reset of the claim deadline from
// now (the new full deadline, not additive; ADR-010 §2).
// The validity predicate rides (a late heartbeat refuses —
// it can never steal the job back from a reclaimer);
// claimed_at never moves.
let now = now_unix();
let conn = pool.get().await.map_err(pool_error)?;
let n = conn
.execute(
&format!(
"UPDATE {job}
SET claim_expires_at = $3::bigint + $4::bigint
WHERE id = $1 AND worker_id = $2
AND state = 'processing'
AND claim_expires_at >= $3",
job = table(&schema, tables::JOB)
),
&[&id, &worker_id, &now, &extend],
)
.await
.map_err(pg_error)?;
Ok(n > 0)
})
}
}
+19 -4
View File
@@ -26,7 +26,8 @@
//! the trait stubs below are the wave-3 posture — they return
//! `Err(Database("… wiring lands with the … task"))` until the
//! mechanism tasks replace them. `notify`/`listen` are wired (the
//! notify-listen task's) — the wake contract's pg arm rides the
//! notify-listen task's), `stream` (the streams task's), `queue`
//! (the queues task's) — the wake contract's pg arm rides the
//! forwarder.
use alkstore::Store;
@@ -281,13 +282,24 @@ impl Store for PgStore {
fn queue<'a>(
&'a self,
_name: &str,
_opts: alkstore::QueueOpts,
name: &str,
opts: alkstore::QueueOpts,
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::Queue>>> {
if let Err(e) = alkstore::validate_shared_name(name) {
return Box::pin(async move { Err(e) });
}
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("queue wiring lands with the queues task")) })
let name = name.to_string();
let handle = crate::queue::PgQueueHandle::new(
name,
opts,
self.pool.clone(),
self.schema.clone(),
self.closed.clone(),
);
Box::pin(async move { Ok(Box::new(handle) as Box<dyn alkstore::Queue>) })
}
fn outbox<'a>(
@@ -355,6 +367,9 @@ impl Store for PgStore {
#[cfg(test)]
mod open_tests;
#[cfg(test)]
mod queue_tests;
#[cfg(test)]
mod notify_tests;
+15 -5
View File
@@ -663,8 +663,18 @@ async fn store_trait_methods_are_wiring_stubs() {
// named handle — full behavior the stream tests' scope.
let handle = store.stream("s").await.unwrap();
assert_eq!(handle.name(), "s");
let err = stub_err!(store.queue("q", alkstore::QueueOpts::default()).await);
assert!(matches!(err, alkstore::Error::Database(_)));
// The queue constructor is wired (the queues task's): returns a
// named handle whose enqueue lands — full behavior the queue
// tests' scope.
let queue = store
.queue("stub_q_probe", alkstore::QueueOpts::default())
.await
.unwrap();
assert_eq!(queue.name(), "stub_q_probe");
queue
.enqueue(serde_json::json!({}), alkstore::EnqueueOpts::default())
.await
.unwrap();
let err = stub_err!(store.outbox("o").await);
assert!(matches!(err, alkstore::Error::Database(_)));
let err = stub_err!(store.try_lock("l", "owner", 60).await);
@@ -691,8 +701,8 @@ async fn store_trait_methods_are_wiring_stubs() {
// The remaining 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). `queue` sampled as the representative.
let err = stub_err!(store.queue("ch", alkstore::QueueOpts::default()).await);
// detail-carriage). `outbox` sampled as the representative.
let err = stub_err!(store.outbox("ch").await);
match err {
alkstore::Error::Database(source) => {
let msg = source.to_string();
@@ -701,7 +711,7 @@ async fn store_trait_methods_are_wiring_stubs() {
"stub errors name their landing task, got: {msg}"
);
}
other => panic!("queue stub must be Database, got {other:?}"),
other => panic!("outbox stub must be Database, got {other:?}"),
}
drop(store);
File diff suppressed because it is too large. Load diff
+55 -34
View File
@@ -192,20 +192,10 @@ impl PgTxHandle {
.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.
/// qualified job table — the module-level [`enqueue_row`]'s
/// handle-side bridge.
async fn enqueue_row(
&self,
client: &tokio_postgres::Client,
@@ -214,31 +204,62 @@ impl PgTxHandle {
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))
enqueue_row(client, &self.schema, queue, bytes, opts, base).await
}
}
/// The enqueue INSERT — the one owner of the plain-queue row write.
/// 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` ride it; one clock
/// read per enqueue (the `now` the resolution and the write share).
pub(crate) async fn enqueue_row(
client: &tokio_postgres::Client,
schema: &str,
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(
&enqueue_row_sql(schema),
&[
&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))
}
/// The enqueue INSERT's SQL (schema-qualified job table).
fn enqueue_row_sql(schema: &str) -> 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",
QualifiedTable {
schema,
name: tables::JOB
}
)
}
/// 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