Postgres engine: scheduler + outbox — schedule (validation triad, the pg-owned @every-only parser, upsert over the schedule table), unschedule, run_schedules (leadership via the engine's lock machinery on __alkstore_scheduler with a per-instance owner token, in-sleep lease renewals, the row-locked FOR UPDATE tick in one pool tx — fire enqueues + boundary advance + soonest read commit together — the 64-boundary catch-up cap with skip-forward, clean stop / Err(LeadershipLost) arms), outbox (validated constructor, enqueue into the derived __alkstore_outbox:{name} with the 60/5/5 stamps + max_attempts override, run_once ack/retry-curve/false over the ordinary claim machinery, no engine-issued heartbeat) (task pg-engine-scheduler-outbox)

This commit is contained in:
glm-5.3-flash committed 2026-10-09 09:53:48 +00:00
1 parent c3591c2d44
commit 77619c5e93
9 files changed
+2547 -78

No files matched your search

+1
View File
@@ -38,6 +38,7 @@ mod notify;
mod opts;
mod queue;
mod resolution;
mod scheduler;
mod schema;
mod seam;
mod store;
+24 -1
View File
@@ -66,7 +66,7 @@ fn table(schema: &str, name: &str) -> String {
/// expiry-delete + insert + read-back race window is closed by the
/// PK constraint — a concurrent grant of the same name either
/// lands first or reads as foreign).
async fn lock_acquire(
pub(crate) async fn lock_acquire(
client: &tokio_postgres::Client,
schema: &str,
name: &str,
@@ -166,6 +166,29 @@ impl std::fmt::Debug for PgLockHandle {
}
}
impl PgLockHandle {
/// The engine-internal constructor (the SQLite twin's
/// `SqliteLockHandle::new_internal` shape): the scheduler's
/// leadership lease on the reserved name — `try_lock`'s entry
/// validation would reject the reserved name, so the engine's own
/// holder builds the handle directly after its own acquire.
pub(crate) fn new_internal(
name: String,
owner: String,
pool: Pool,
schema: String,
closed: Arc<AtomicBool>,
) -> Self {
Self {
name,
owner,
pool,
schema,
closed,
}
}
}
impl Lock for PgLockHandle {
fn name(&self) -> &str {
&self.name
+15 -3
View File
@@ -304,7 +304,7 @@ fn row_get_bytes(row: &tokio_postgres::Row, idx: usize) -> Vec<u8> {
/// `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>
pub(crate) async fn in_tx<T, F>(object: Object, body: F) -> Result<T>
where
T: Send,
F: FnOnce(&tokio_postgres::Client) -> BoxedFuture<'_, Result<T>>,
@@ -344,7 +344,7 @@ where
/// 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(
pub(crate) async fn sweep_exhausted_reclaimables_tx(
client: &tokio_postgres::Client,
schema: &str,
queue: &str,
@@ -377,7 +377,7 @@ async fn sweep_exhausted_reclaimables_tx(
/// `attempts += 1`.
///
/// Parameter positions: $1 worker_id, $2 now, $3 queue, $4 extent n.
fn claim_stmt_sql(schema: &str) -> String {
pub(crate) fn claim_stmt_sql(schema: &str) -> String {
format!(
"UPDATE {job} SET
state = 'processing',
@@ -745,6 +745,18 @@ impl std::fmt::Debug for PgJobHandle {
}
impl PgJobHandle {
/// The engine-internal constructor (the SQLite twin's
/// `SqliteJobHandle::new` shape): the outbox's `run_once` rebuilds
/// the claim handle for the delivery callable. Fields stay private.
pub(crate) fn new(job: Job, pool: Pool, schema: String, closed: Arc<AtomicBool>) -> Self {
Self {
job,
pool,
schema,
closed,
}
}
/// 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
+755
View File
@@ -0,0 +1,755 @@
//! The scheduler mechanism's Postgres arm (`pg-engine-scheduler-outbox`)
//! — the collapse shape (ADR-009) — and the outbox helper (ADR-014's
//! auto-commit half, ADR-019 §5's `run_once` shape; the tx half —
//! `outbox_enqueue_tx` — landed with the seam task). The structural
//! twin of the SQLite engine's `scheduler.rs`/`outbox.rs` re-owned as
//! pg SQL (ADR-012 §2's one-owner rule — the substrate's
//! `scheduler_register`/`scheduler_tick`/`scheduler_soonest` are
//! SQLite-side, so the `@every` grammar, the boundary math, the
//! 64-cap catch-up, and the skip-forward have this engine's own
//! implementation; the contract suite pins the twins' outcomes
//! identical, wave 5).
//!
//! - `schedule(name, spec, queue, payload, opts)` — upsert by name;
//! entry-point validation on all three name-bearing arguments
//! *before any round trip* (the schedule name and the **queue
//! argument** identically per ADR-021 §3 — a schedule cannot target
//! a reserved name, so the outbox's backing queue stays reachable
//! only through the outbox surface, ADR-014 §1's stated guarantee);
//! the pg-owned `@every` parser (`parse_every_interval` below —
//! `@every <n><unit>` only, s|m|h|d; cron strings → `InvalidSpec`
//! pre-storage); payload encoded at the seam (`Codec`); the
//! first-fire boundary = `now + interval` (strictly after now); the
//! resolved stamps over the plain-queue derived defaults
//! ([`crate::resolution::plain_queue_default_stamps`], ADR-020 §3
//! — no handle is open at fire time by design).
//! - `unschedule(name)` — the delete; `Ok(false)` when absent.
//! - `run_schedules(stop)` — the leader loop: leadership via the
//! engine's own lock machinery ([`crate::lock`]) on the reserved
//! name `__alkstore_scheduler` (ADR-009 §4's layer (a)) with a
//! per-runner-instance-unique owner token (pid + subsec-nanos — a
//! shared constant would make two processes each see "self"
//! holding, the wave-3 lesson), TTL engine-internal (10 s), renewed
//! at the top of every loop iteration **and** on a cadence across
//! the sleep (a boundary further out than the TTL must not lapse
//! its own lease by drift; a failed mid-sleep renewal is not acted
//! on at the sleep site — the loop's top-of-iteration renew owns
//! the loss decision). Acquire-time refusal is `Err(LeadershipLost)`
//! — the same matchable loss a mid-run loser returns. The tick:
//! one pool-checked-out transaction (the [`crate::queue::in_tx`]
//! frame), the schedule row's `SELECT … FOR UPDATE` inside the tick
//! tx as the row lock that serializes rogue second tickers (ADR-009
//! §4's layer (b)); each due fire is an ordinary enqueue stamped
//! with the `ScheduleOpts` resolved over the target queue's derived
//! engine defaults; the 64-boundary catch-up cap with skip-forward
//! past the cap (honker parity — the contract constant, not a knob);
//! the soonest-boundary read for the sleep deadline. Wake-driven
//! tick advance: the runner parks on the sleep deadline but the
//! loop's stop-select shortens every sleep slice to ≤ 1 s and both
//! the tick and the soonest are re-read each iteration boundary —
//! a newly-registered schedule is noticed by the re-read no later
//! than the next slice (the deadline always fires; the wake never
//! gates correctness — the SQLite runner's 60 s idle-tick posture
//! is the acceptable floor if the slice cadence were fussy).
//! - Clean stop (`StopToken::cancel`) → `Ok(())`, the leadership lock
//! released cleanly; loss → `Err(LeadershipLost)` before any tick.
//! - `outbox(name)` — the validated constructor (empty →
//! `InvalidName`, reserved-prefixed → `ReservedName`; carries no
//! opts of its own — the backing queue's derived 60/5/5 set is
//! the engine-documented stamp source, ADR-010 §3a). `enqueue` is
//! the auto-commit path into the derived backing queue
//! `__alkstore_outbox:{name}` (unreachable by `enqueue_tx`/`queue()`
//! by the reserved-prefix design — ADR-014 §1), full ADR-020
//! resolution, riding the queues task's `enqueue_row`. `run_once`
//! claims one job via the ordinary claim machinery (the queues
//! task's `FOR UPDATE SKIP LOCKED` frame), runs the delivery
//! callable (core's [`alkstore::Delivery`]), `ack`s on `Ok`,
//! `retry(err, None)` on `Err` (the queue's curve — the delivery
//! error's `Display` string as the retry's exhaust string), `false`
//! on empty. **No engine-issued heartbeat inside delivery** (honker
//! parity — the boxed [`alkstore::JobHandle`]'s own `heartbeat`
//! remains available to the closure).
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use deadpool_postgres::Pool;
use serde_json::Value;
use alkstore::{
BoxedFuture, Delivery, EnqueueOpts, Error, JobHandle, Outbox, Result, Schedule, ScheduleOpts,
StopToken, validate_shared_name,
};
use crate::lock::{PgLockHandle, lock_acquire};
use crate::queue::{PgJobHandle, backoff_delay_s, in_tx};
use crate::resolution::{
Stamps, now_unix, outbox_backing_queue_name, outbox_default_stamps, plain_queue_default_stamps,
};
use crate::schema::{QualifiedTable, tables};
use crate::seam::{database_error, pg_error, pool_error};
use crate::tx::{encode_payload_bytes, enqueue_row};
/// The contract-constant per-task per-tick catch-up cap (ADR-009 §4 —
/// honker parity, not a knob): a downtime longer than 64 interval
/// widths replays exactly 64 fires and skips forward past the cap.
const SCHEDULER_MAX_CATCHUP_FIRES: i64 = 64;
/// The leadership lock's TTL (seconds): engine-internal detail —
/// renewed every loop iteration *and* across the sleeps (a runner
/// that cannot renew has lost the lock, ADR-009 §4 — the lock
/// machinery is the engine's own).
const LEADER_TTL_S: i64 = 10;
/// Seconds between lease renewals during the sleep.
const SLEEP_RENEW_EVERY_S: u64 = 2;
/// The sleep-slice length (seconds): the loop re-reads the stop token
/// and the soonest-driven deadline at this cadence — the bounded-wait
/// interrupt that keeps a clean stop and a mid-sleep loss check
/// prompt without a dedicated wake channel on the scheduler's side.
const SLEEP_SLICE_S: u64 = 1;
/// The idle sleep when nothing is scheduled (soonest = 0).
const IDLE_SLEEP_S: u64 = 60;
/// The leadership lock's reserved name (ADR-009 §4 — the lock
/// machinery's namespace, engine-derived; consumers can never collide
/// because every entry point rejects the reserved prefix).
pub(crate) const SCHEDULER_LEADER_LOCK: &str = concat!("__alkstore_", "scheduler");
/// The leadership-lock owner: unique per runner instance — the owner
/// string is the claim's identity token, and a shared constant across
/// processes would make the acquire's read-back mistake a foreign
/// leader for self (two leaders holding the same name at once — the
/// wave-3 finding's lesson). pid + subsec-nanos.
fn leader_owner() -> String {
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.subsec_nanos())
.unwrap_or(0);
format!("scheduler-leader-{}-{nanos}", std::process::id())
}
fn invalid_spec(spec: &str) -> Error {
Error::InvalidSpec {
spec: spec.to_string(),
}
}
/// The pg engine's own `@every` parser (ADR-012 §2's one-owner rule —
/// the substrate's parser is SQLite-side): exactly `@every <n><unit>`
/// with `n` a positive integer and unit in `s|m|h|d`. Everything else
/// (cron strings, bare `every 5s`, floats, unknown units) rejects.
/// Returns the interval in seconds.
pub(crate) fn parse_every_interval(spec: &str) -> Result<i64> {
let reject = || invalid_spec(spec);
let body = spec.strip_prefix("@every ").ok_or_else(reject)?;
if body.is_empty() {
return Err(reject());
}
let digits = body.chars().take_while(|c| c.is_ascii_digit()).count();
if digits == 0 || digits == body.len() {
return Err(reject());
}
let n: i64 = body[..digits].parse().map_err(|_| reject())?;
if n <= 0 {
return Err(reject());
}
let mult = match &body[digits..] {
"s" => 1,
"m" => 60,
"h" => 3600,
"d" => 86400,
_ => return Err(reject()),
};
n.checked_mul(mult).ok_or_else(|| Error::InvalidSpec {
spec: spec.to_string(),
})
}
/// The next boundary strictly after `after` (the register path's
/// first fire: `now + interval`, honker's shape).
fn next_every_boundary(spec: &str, after: i64) -> Result<i64> {
let interval = parse_every_interval(spec)?;
after
.checked_add(interval)
.ok_or_else(|| invalid_spec(spec))
}
/// The skip-forward past the catch-up cap (ADR-009 §4): from a
/// boundary `now - gap` widths out of date, jump to
/// `next + (gap / interval + 1) * interval` — strictly beyond `now`
/// (the skipped boundaries never enqueue). 128-bit intermediates keep
/// the arithmetic saturation-safe; the result clamps at `i64::MAX`.
fn skip_forward(next_fire_at: i64, now: i64, interval: i64) -> i64 {
let gap = (now - next_fire_at) as i128;
let interval = interval.max(1) as i128;
let next = (next_fire_at as i128) + (gap / interval + 1) * interval;
if next > i64::MAX as i128 {
i64::MAX
} else {
next as i64
}
}
/// The schedule row's stamp shape (the resolved stamps a fire
/// enqueues with) — read back with the row.
struct SchedRow {
name: String,
queue: String,
spec: String,
payload: Vec<u8>,
priority: i64,
expires_s: Option<i64>,
next_fire_at: i64,
stamps: Stamps,
}
fn sched_table(schema: &str) -> String {
QualifiedTable {
schema,
name: tables::SCHEDULE,
}
.to_string()
}
/// The engine-context bundle: the pool, the engine-owned schema, and
/// the closed flag — the wiring every scheduler/outbox path carries.
/// One struct keeps the entry points' argument shapes tidy (the
/// store trait's wiring site constructs it once).
#[derive(Clone)]
pub(crate) struct EngineCtx {
pub pool: Pool,
pub schema: String,
pub closed: Arc<AtomicBool>,
}
impl EngineCtx {
fn closed_error(&self) -> Error {
database_error("the store is closed")
}
fn check_closed(&self) -> Result<()> {
if self.closed.load(Ordering::Acquire) {
return Err(self.closed_error());
}
Ok(())
}
}
/// The auto-commit register path: pre-storage validation on all three
/// name-bearing arguments, the upsert, the [`Schedule`] read-back.
pub(crate) fn schedule(
name: &str,
spec: &str,
queue: &str,
payload: Value,
opts: ScheduleOpts,
ctx: EngineCtx,
) -> BoxedFuture<'static, Result<Schedule>> {
if let Err(e) = validate_shared_name(name) {
return Box::pin(async move { Err(e) });
}
if let Err(e) = validate_shared_name(queue) {
return Box::pin(async move { Err(e) });
}
if parse_every_interval(spec).is_err() {
let spec = spec.to_string();
return Box::pin(async move { Err(invalid_spec(&spec)) });
}
if let Err(e) = ctx.check_closed() {
return Box::pin(async move { Err(e) });
}
let bytes = match encode_payload_bytes(&payload) {
Ok(bytes) => bytes,
Err(e) => return Box::pin(async move { Err(e) }),
};
let base = plain_queue_default_stamps();
let stamps = Stamps {
max_attempts: opts.max_attempts.unwrap_or(base.max_attempts),
..base
};
let stored_name = name.to_string();
let stored_spec = spec.to_string();
let stored_queue = queue.to_string();
let EngineCtx { pool, schema, .. } = ctx;
Box::pin(async move {
let client = pool.get().await.map_err(pool_error)?;
let now = now_unix();
let first_fire = next_every_boundary(&stored_spec, now)?;
let conn: &tokio_postgres::Client = &client;
conn.execute(
&format!(
"INSERT INTO {sched}
(name, queue, spec, payload, priority, expires_s, next_fire_at,
max_attempts, visibility_timeout_s, backoff_base_s,
dead_letter_retention_s)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
ON CONFLICT (name) DO UPDATE SET
queue = EXCLUDED.queue,
spec = EXCLUDED.spec,
payload = EXCLUDED.payload,
priority = EXCLUDED.priority,
expires_s = EXCLUDED.expires_s,
next_fire_at = EXCLUDED.next_fire_at,
max_attempts = EXCLUDED.max_attempts,
visibility_timeout_s = EXCLUDED.visibility_timeout_s,
backoff_base_s = EXCLUDED.backoff_base_s,
dead_letter_retention_s = EXCLUDED.dead_letter_retention_s",
sched = sched_table(&schema)
),
&[
&stored_name,
&stored_queue,
&stored_spec,
&bytes,
&opts.priority,
&opts.expires,
&first_fire,
&stamps.max_attempts,
&stamps.visibility_timeout_s,
&stamps.backoff_base_s,
&stamps.dead_letter_retention_s,
],
)
.await
.map_err(pg_error)?;
Ok(Schedule::new(stored_name, stored_spec, stored_queue, opts))
})
}
/// The auto-commit unregister path (validation mirrors `schedule`'s —
/// a reserved/empty schedule name never round trips).
pub(crate) fn unschedule(name: &str, ctx: EngineCtx) -> BoxedFuture<'static, Result<bool>> {
if let Err(e) = validate_shared_name(name) {
return Box::pin(async move { Err(e) });
}
if let Err(e) = ctx.check_closed() {
return Box::pin(async move { Err(e) });
}
let EngineCtx { pool, schema, .. } = ctx;
let name = name.to_string();
Box::pin(async move {
let client = pool.get().await.map_err(pool_error)?;
let n = client
.execute(
&format!(
"DELETE FROM {sched} WHERE name = $1",
sched = sched_table(&schema)
),
&[&name],
)
.await
.map_err(pg_error)?;
Ok(n > 0)
})
}
/// One tick under one pool-checked-out transaction (the queues task's
/// [`in_tx`] frame — `BEGIN` … the body … `COMMIT`): the fire enqueues
/// and the boundary advance commit together, or the whole tick rolls
/// back (crash mid-tick ⇒ the boundary refires; ADR-009 §4). The
/// schedule row's `SELECT … FOR UPDATE` inside the tick tx is the
/// row lock that serializes rogue second tickers (layer (b)) — a
/// second ticker re-reading under the lock sees the advanced
/// `next_fire_at` and skips. pg's default isolation is exactly the
/// committed-snapshot+row-lock semantics the guarantee names.
/// Returns the soonest next fire across all schedules for the sleep
/// deadline (0 = none registered).
pub(crate) async fn tick(client: &tokio_postgres::Client, schema: &str, now: i64) -> Result<i64> {
let sched = sched_table(schema);
// The row-locked due read (layer (b)): re-checks `next_fire_at <=
// now` under the lock — a rogue second ticker observing the same
// storage serializes here and then observes the advanced boundary.
let rows = client
.query(
&format!(
"SELECT name, queue, spec, payload, priority, expires_s,
next_fire_at, max_attempts, visibility_timeout_s,
backoff_base_s, dead_letter_retention_s
FROM {sched}
WHERE next_fire_at <= $1
ORDER BY next_fire_at ASC
FOR UPDATE",
sched = sched
),
&[&now],
)
.await
.map_err(pg_error)?;
let mut due = Vec::with_capacity(rows.len());
for row in &rows {
due.push(SchedRow {
name: row.get(0),
queue: row.get(1),
spec: row.get(2),
payload: row.get(3),
priority: row.get(4),
expires_s: row.get(5),
next_fire_at: row.get(6),
stamps: Stamps {
max_attempts: row.get(7),
visibility_timeout_s: row.get(8),
backoff_base_s: row.get(9),
dead_letter_retention_s: row.get(10),
},
});
}
for row in &due {
let interval = parse_every_interval(&row.spec)?;
let mut next_fire_at = row.next_fire_at;
let mut fires = 0i64;
// The one-firer-per-boundary loop: at-least-once per elapsed
// boundary (each fire an ordinary enqueue, stamped from the
// row's resolved stamps — ADR-020 §3's derived-defaults shape),
// up to the 64-cap, then skip-forward past now.
while next_fire_at <= now {
if fires >= SCHEDULER_MAX_CATCHUP_FIRES {
next_fire_at = skip_forward(next_fire_at, now, interval);
break;
}
let _fire_at = next_fire_at;
enqueue_row(
client,
schema,
&row.queue,
row.payload.clone(),
&EnqueueOpts {
priority: row.priority,
expires: row.expires_s,
..Default::default()
},
row.stamps,
)
.await?;
fires += 1;
next_fire_at = next_fire_at.saturating_add(interval);
}
client
.execute(
&format!(
"UPDATE {sched} SET next_fire_at = $2 WHERE name = $1",
sched = sched
),
&[&row.name, &next_fire_at],
)
.await
.map_err(pg_error)?;
}
// The soonest-boundary read for the sleep deadline (inside the
// same tx — the advance is committed before the sleep takes it).
let soonest: i64 = client
.query_one(
&format!(
"SELECT COALESCE(MIN(next_fire_at), 0) FROM {sched}",
sched = sched
),
&[],
)
.await
.map_err(pg_error)?
.get(0);
Ok(soonest)
}
/// The leader loop (see the module docs).
pub(crate) fn run_schedules(stop: StopToken, ctx: EngineCtx) -> BoxedFuture<'static, Result<()>> {
let EngineCtx {
pool,
schema,
closed,
} = ctx;
Box::pin(async move {
if closed.load(Ordering::Acquire) {
return Err(database_error("the store is closed"));
}
// Step 1: acquire the leadership lock via the engine's own
// lock machinery. A refused acquire — someone else leads — is
// `Err(LeadershipLost)` at start: no leader loop exists
// without the lock, and the respawn recipe is identical to a
// mid-run loss.
let owner = leader_owner();
// The acquire op rides the lock module's own acquire function
// (try_lock's entry-point name validation would reject the
// reserved leadership name — the engine-internal acquire is
// the SQLite twin's shape: the machinery, not the consumer
// entry point).
let granted = {
let client = pool.get().await.map_err(pool_error)?;
lock_acquire(
&client,
&schema,
SCHEDULER_LEADER_LOCK,
&owner,
LEADER_TTL_S,
)
.await?
};
if !granted {
return Err(Error::LeadershipLost);
}
let mut lease: Box<dyn alkstore::Lock> = Box::new(PgLockHandle::new_internal(
SCHEDULER_LEADER_LOCK.to_string(),
owner,
pool.clone(),
schema.clone(),
closed.clone(),
));
// Step 2: the loop — renew, tick, sleep until soonest or stop.
loop {
if stop.is_cancelled() {
let _ = lease.release().await;
return Ok(());
}
let renewed = lease.renew(LEADER_TTL_S).await.unwrap_or(false);
if !renewed {
// The lock was lost (expired and re-acquired elsewhere
// — exclusion lapses silently at TTL expiry). Return
// before ticking: never fire on a stolen lease.
return Err(Error::LeadershipLost);
}
let soonest = {
let conn = pool.get().await.map_err(pool_error)?;
in_tx(conn, |client| {
let schema = schema.clone();
Box::pin(async move { tick(client, &schema, now_unix()).await })
})
.await?
};
if stop.is_cancelled() {
let _ = lease.release().await;
return Ok(());
}
// The sleep renews the lease across itself — a boundary
// further out than the TTL must not lose the lock by
// drifting past it (the renew-on-cadence-while-holding
// posture). A failed mid-sleep renewal is not acted on at
// the sleep site: the loop's top-of-iteration renew owns
// the loss decision.
keep_awake_renewing(&mut lease, &stop, soonest).await;
}
})
}
/// The sliced sleep until the next due boundary (or idle horizon) or
/// the stop token — the loop's wake path between ticks. The slice is
/// 1 s, so a clean stop ends the sleep promptly and a mid-sleep lease
/// renewal cadence holds; the deadline always fires (correctness never
/// rides the wake path — it only shortens idle latency).
async fn keep_awake_renewing(lease: &mut Box<dyn alkstore::Lock>, stop: &StopToken, soonest: i64) {
let renew_every = Duration::from_secs(SLEEP_RENEW_EVERY_S);
let mut next_renew = tokio::time::Instant::now() + renew_every;
let started = tokio::time::Instant::now();
let wait_s: u64 = if soonest <= 0 {
IDLE_SLEEP_S
} else {
// A boundary that fell behind during the tick's own runtime is
// due immediately (fire on the next iteration, don't sleep).
(soonest - now_unix()).max(0) as u64
};
let deadline = started + Duration::from_secs(wait_s);
let slice = Duration::from_secs(SLEEP_SLICE_S);
loop {
let now = tokio::time::Instant::now();
if stop.is_cancelled() || now >= deadline {
return;
}
if now >= next_renew {
// A failed mid-sleep renewal is not acted on here — the
// loop's own top-of-iteration renew decides (one owner of
// the loss decision, before any tick).
let _ = lease.renew(LEADER_TTL_S).await;
next_renew = tokio::time::Instant::now() + renew_every;
}
let wake = deadline.min(next_renew);
tokio::time::sleep_until(wake.min(now + slice)).await;
}
}
/// The outbox-scoped handle (see the module docs).
pub struct PgOutboxHandle {
name: String,
backing: String,
pool: Pool,
schema: String,
closed: Arc<AtomicBool>,
}
impl std::fmt::Debug for PgOutboxHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PgOutboxHandle")
.field("name", &self.name)
.finish_non_exhaustive()
}
}
/// The validated constructor (the store's `outbox` wiring): shared-
/// namespace validation at the entry point, the derived backing queue.
pub(crate) fn new_outbox_handle(name: &str, ctx: EngineCtx) -> Result<Box<dyn Outbox>> {
validate_shared_name(name)?;
let EngineCtx {
pool,
schema,
closed,
} = ctx;
Ok(Box::new(PgOutboxHandle {
name: name.to_string(),
backing: outbox_backing_queue_name(name),
pool,
schema,
closed,
}))
}
impl Outbox for PgOutboxHandle {
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(database_error("the store is closed")) });
}
let backing = self.backing.clone();
let pool = self.pool.clone();
let schema = self.schema.clone();
Box::pin(async move {
let bytes = encode_payload_bytes(&payload)?;
let conn = pool.get().await.map_err(pool_error)?;
// Full ADR-020 resolution over the outbox's derived 60/5/5
// set (the "no handle was open" stamp source, ADR-010
// §3a) — through the queues task's shared row-write; then
// the best-effort wake on the backing queue's channel (the
// durable row is the truth — the twin posture of the
// handle enqueue).
let id = enqueue_row(
&conn,
&schema,
&backing,
bytes,
&opts,
outbox_default_stamps(),
)
.await?;
if let Err(wake_err) = conn
.query_one("SELECT pg_notify($1, '')", &[&backing])
.await
{
eprintln!(
"alkstore-postgres: outbox enqueue wake failed for {backing:?} \
(the durable row landed; the re-poll safety net covers the gap): {wake_err}"
);
}
Ok(id)
})
}
fn run_once<'a>(
&'a mut self,
worker_id: &str,
delivery: &'a mut dyn Delivery,
) -> BoxedFuture<'a, Result<bool>> {
if let Err(e) = alkstore::validate_local_name(worker_id) {
return Box::pin(async move { Err(e) });
}
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(database_error("the store is closed")) });
}
let backing = self.backing.clone();
let pool = self.pool.clone();
let schema = self.schema.clone();
let closed = self.closed.clone();
let worker_id = worker_id.to_string();
Box::pin(async move {
// Claim through the ordinary claim machinery — the queues
// task's `claim_batch` shape (one statement under the
// SKIP LOCKED frame; the exactly-once handout), scoped to
// the derived backing queue. Rebuilt here over the same
// module pieces the Queue handle runs (one owner per
// formula) with the backing name standing in for the
// handle's.
let conn = pool.get().await.map_err(pool_error)?;
let job = in_tx(conn, |client| {
let backing = backing.clone();
let schema = schema.clone();
let worker = worker_id.clone();
Box::pin(async move {
crate::queue::sweep_exhausted_reclaimables_tx(
client,
&schema,
&backing,
now_unix(),
)
.await?;
let rows = client
.query(
&crate::queue::claim_stmt_sql(&schema),
&[&worker, &now_unix(), &backing, &1i64],
)
.await
.map_err(pg_error)?;
match rows.first() {
Some(row) => Ok(Some(crate::queue::job_from_row(row)?)),
None => Ok(None),
}
})
})
.await?;
let Some(job) = job else {
return Ok(false);
};
// No engine-issued heartbeat inside delivery: the boxed
// handle's own heartbeat is the consumer's renewal path
// (ADR-019 §5, honker parity).
let handle: Box<dyn JobHandle> = Box::new(PgJobHandle::new(
job.clone(),
pool.clone(),
schema.clone(),
closed.clone(),
));
match delivery.deliver(handle).await {
Ok(()) => {
// Ok ⇒ ack: the owner-scoped delete under the
// validity predicate — the queue handle's ack arm
// re-owned through the same handle type (the
// row-scoped delete is the truth-teller).
let handle: Box<dyn JobHandle> = Box::new(PgJobHandle::new(
job,
pool.clone(),
schema.clone(),
closed.clone(),
));
handle.ack().await?;
Ok(true)
}
Err(err) => {
// `retry(Some(e), None)` exactly — the delivery's
// error string is the exhaust string the way
// `JobHandle::retry(Some(e), None)` would carry it
// (surfaced on a dead row), and the delay is the
// queue's curve (the engine's one-owner
// equal-jitter computation, ADR-010 §3).
let exhaust = err.to_string();
let delay = backoff_delay_s(job.backoff_base_s, job.attempts, job.id);
let handle: Box<dyn JobHandle> = Box::new(PgJobHandle::new(
job,
pool.clone(),
schema.clone(),
closed.clone(),
));
handle.retry(Some(exhaust), Some(delay)).await?;
Ok(true)
}
}
})
}
}
+29 -24
View File
@@ -103,6 +103,16 @@ impl PgStore {
&self.schema
}
/// The scheduler/outbox paths' wiring bundle (the pool, the
/// schema, the closed flag).
fn engine_ctx(&self) -> crate::scheduler::EngineCtx {
crate::scheduler::EngineCtx {
pool: self.pool.clone(),
schema: self.schema.clone(),
closed: self.closed.clone(),
}
}
fn closed_check(&self) -> Result<(), alkstore::Error> {
if self.closed.load(Ordering::Acquire) {
return Err(database_error("the store is closed"));
@@ -223,10 +233,6 @@ impl PgStore {
}
}
fn stub(wiring: &'static str) -> alkstore::Error {
database_error(wiring)
}
impl Store for PgStore {
fn begin_tx(
&self,
@@ -304,12 +310,13 @@ impl Store for PgStore {
fn outbox<'a>(
&'a self,
_name: &str,
name: &str,
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::Outbox>>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("outbox wiring lands with the scheduler-outbox task")) })
let handle = crate::scheduler::new_outbox_handle(name, self.engine_ctx());
Box::pin(async move { handle })
}
fn try_lock<'a>(
@@ -330,41 +337,33 @@ impl Store for PgStore {
fn schedule<'a>(
&'a self,
_name: &str,
_spec: &str,
_queue: &str,
_payload: serde_json::Value,
_opts: alkstore::ScheduleOpts,
name: &str,
spec: &str,
queue: &str,
payload: serde_json::Value,
opts: alkstore::ScheduleOpts,
) -> alkstore::BoxedFuture<'a, alkstore::Result<alkstore::Schedule>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async { Err(stub("schedule wiring lands with the scheduler-outbox task")) })
crate::scheduler::schedule(name, spec, queue, payload, opts, self.engine_ctx())
}
fn unschedule<'a>(&'a self, _name: &str) -> alkstore::BoxedFuture<'a, alkstore::Result<bool>> {
fn unschedule<'a>(&'a self, name: &str) -> alkstore::BoxedFuture<'a, alkstore::Result<bool>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async {
Err(stub(
"unschedule wiring lands with the scheduler-outbox task",
))
})
crate::scheduler::unschedule(name, self.engine_ctx())
}
fn run_schedules<'a>(
&'a self,
_stop: alkstore::StopToken,
stop: alkstore::StopToken,
) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> {
if let Err(e) = self.closed_check() {
return Box::pin(async move { Err(e) });
}
Box::pin(async {
Err(stub(
"scheduler wiring lands with the scheduler-outbox task",
))
})
crate::scheduler::run_schedules(stop, self.engine_ctx())
}
}
@@ -385,3 +384,9 @@ mod stream_tests;
#[cfg(test)]
mod tx_tests;
#[cfg(test)]
mod scheduler_tests;
#[cfg(test)]
mod outbox_tests;
+22 -47
View File
@@ -132,18 +132,6 @@ async fn wait_for_channel(
}
}
/// Extract the `Err` arm from a stub call whose `Ok` type is not
/// `Debug` (the boxed-handle returns: `begin_tx`, `listen`, `stream`,
/// `queue`, `outbox`, `try_lock`).
macro_rules! stub_err {
($call:expr) => {
match $call {
Err(e) => e,
Ok(_) => panic!(concat!(stringify!($call), " must be a stub")),
}
};
}
#[tokio::test(flavor = "multi_thread")]
async fn open_boots_pool_bootstrap_and_forwarder() {
let Some(dsn) = harness_dsn() else {
@@ -620,11 +608,13 @@ async fn open_fails_database_on_unparseable_and_unreachable() {
}
}
/// Acceptance: the trait surface — `begin_tx`, `notify`, `listen`,
/// `stream`, `queue`, and `try_lock` are wired; the remaining methods
/// are the wave-3 stub posture — every unwired `Store` method
/// returns `Err(Database(… wiring lands with the … task))`; `with_tx`
/// surfaces the wired `begin_tx` naturally; no panics.
/// Acceptance: the trait surface — every `Store` method is wired (the
/// wave-3 stub posture is fully retired): `outbox` is a validated
/// constructor (the scheduler-outbox task's), `schedule` registers +
/// unschedule removes, `run_schedules` is refused by a live leadership
/// lock with the matchable loss (full behavior the scheduler/outbox
/// tests' scope); `with_tx` surfaces the wired `begin_tx` naturally;
/// no panics.
#[tokio::test(flavor = "multi_thread")]
async fn store_trait_methods_are_wiring_stubs() {
let Some(dsn) = harness_dsn() else {
@@ -684,42 +674,27 @@ async fn store_trait_methods_are_wiring_stubs() {
.expect("the wired try_lock grants");
assert_eq!(lock.name(), "stub_l_probe");
assert!(lock.release().await.unwrap(), "the wired release deletes");
let err = stub_err!(store.outbox("o").await);
assert!(matches!(err, alkstore::Error::Database(_)));
let err = store
// The outbox constructor is wired (the scheduler-outbox task's): a
// valid name returns a named handle — full behavior the outbox
// tests' scope.
let outbox = store.outbox("stub_o_probe").await.unwrap();
assert_eq!(outbox.name(), "stub_o_probe");
// The scheduler registration surface is wired (the scheduler-outbox
// task's): a valid register reads back, unregister removes — full
// behavior the scheduler tests' scope.
let sched = store
.schedule(
"s",
"stub_s_probe",
"@every 10s",
"q",
"stub_q_probe",
serde_json::json!({}),
alkstore::ScheduleOpts::default(),
)
.await
.unwrap_err();
assert!(matches!(err, alkstore::Error::Database(_)));
let err = store.unschedule("s").await.unwrap_err();
assert!(matches!(err, alkstore::Error::Database(_)));
let err = store
.run_schedules(alkstore::StopToken::new())
.await
.unwrap_err();
assert!(matches!(err, alkstore::Error::Database(_)));
// The 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). `outbox` sampled as the representative.
let err = stub_err!(store.outbox("ch").await);
match err {
alkstore::Error::Database(source) => {
let msg = source.to_string();
assert!(
msg.contains("wiring lands with"),
"stub errors name their landing task, got: {msg}"
);
}
other => panic!("outbox stub must be Database, got {other:?}"),
}
.unwrap();
assert_eq!(sched.name, "stub_s_probe");
assert!(store.unschedule("stub_s_probe").await.unwrap());
assert!(!store.unschedule("stub_s_probe").await.unwrap());
drop(store);
let admin = harness_client().await.unwrap();
+510
View File
@@ -0,0 +1,510 @@
//! The outbox mechanism's acceptance tests (`pg-engine-scheduler-outbox`):
//! the derived backing queue is unreachable by `enqueue_tx` (the
//! reserved-prefix rejection — ADR-014 §1's "no third path" guarantee),
//! the outbox's `enqueue` stamps the 60/5/5 derived set with the one
//! `max_attempts` override, `run_once`'s ok/err/false dispositions with
//! the curve gap asserted, in-delivery heartbeat renewal available (no
//! engine-issued heartbeat), and exhaustion to dead with the delivery
//! error string. The commit-atomic tx half pins under tx_tests (landed
//! with the seam task).
use std::time::Duration;
use alkstore::{Delivery, EnqueueOpts, JobHandle, JobState, Store};
use serde_json::json;
use crate::store::{PgStore, open_store};
const ENV_HOST: &str = "ALKSTORE_PG_HOST";
const ENV_PORT: &str = "ALKSTORE_PG_PORT";
const ENV_USER: &str = "ALKSTORE_PG_USER";
const ENV_PASSWORD: &str = "ALKSTORE_PG_PASSWORD";
const ENV_DB: &str = "ALKSTORE_PG_DB";
fn harness_dsn() -> Option<String> {
let host = std::env::var(ENV_HOST).ok()?;
let port: u16 = std::env::var(ENV_PORT).ok()?.parse().ok()?;
let user = std::env::var(ENV_USER).ok()?;
let password = std::env::var(ENV_PASSWORD).ok()?;
let db = std::env::var(ENV_DB).unwrap_or_else(|_| "postgres".to_string());
Some(format!(
"host={host} port={port} user={user} password={password} dbname={db}"
))
}
fn instance_namer(tag: &str) -> impl Fn() -> String + use<'_> {
use std::sync::atomic::{AtomicU64, Ordering};
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<tokio_postgres::Client> {
let (client, connection) = tokio_postgres::connect(&harness_dsn()?, tokio_postgres::NoTls)
.await
.ok()?;
tokio::spawn(async move {
let _ = connection.await;
});
Some(client)
}
async fn drop_schema(client: &tokio_postgres::Client, schema: &str) {
if schema == crate::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) -> crate::opts::PgOpts {
crate::opts::PgOpts {
schema: schema.to_string(),
..crate::opts::PgOpts::default()
}
}
struct Fixture {
store: PgStore,
admin: tokio_postgres::Client,
schema: String,
}
impl Fixture {
async fn open(namer: &impl Fn() -> String) -> Option<Fixture> {
let dsn = harness_dsn()?;
let schema = namer();
let store = open_store(&dsn, test_opts(&schema)).await.unwrap();
let admin = harness_client().await?;
Some(Fixture {
store,
admin,
schema,
})
}
}
impl Drop for Fixture {
fn drop(&mut self) {
self.store.close();
}
}
fn now_unix() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as i64
}
fn job_table(schema: &str) -> String {
crate::schema::QualifiedTable {
schema,
name: crate::schema::tables::JOB,
}
.to_string()
}
fn dead_table(schema: &str) -> String {
crate::schema::QualifiedTable {
schema,
name: crate::schema::tables::DEAD,
}
.to_string()
}
/// Count rows in the outbox's backing queue directly (the reserved
/// name is unreachable by the handle surface by design).
async fn backing_count(admin: &tokio_postgres::Client, schema: &str, outbox: &str) -> i64 {
admin
.query_one(
&format!(
"SELECT count(*) FROM {} WHERE queue = $1",
job_table(schema)
),
&[&format!("__alkstore_outbox:{outbox}")],
)
.await
.unwrap()
.get(0)
}
/// A minimal delivery impl driving `run_once` from the test side. The
/// error it feeds is replayed as the exhaust string (`retry(Some(e),
/// None)` carries the delivery error's `Display`), so the test holds
/// both ends of that mapping.
struct Scripted {
outcome: Result<(), String>,
seen_payloads: Vec<Vec<u8>>,
}
impl Delivery for Scripted {
fn deliver<'a>(
&'a mut self,
job: Box<dyn JobHandle>,
) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> {
self.seen_payloads.push(job.job().payload.clone());
let outcome = self.outcome.clone().map_err(alkstore::Error::Codec);
Box::pin(async move { outcome })
}
}
/// Acceptance: the constructor validates (empty → `InvalidName`,
/// reserved → `ReservedName`) and the handle name reads back.
#[tokio::test(flavor = "multi_thread")]
async fn outbox_constructor_validates() {
let Some(fx) = Fixture::open(&instance_namer("ob_con")).await else {
eprintln!("skip: no harness server");
return;
};
match fx.store.outbox("").await {
Err(alkstore::Error::InvalidName { .. }) => {}
Ok(_) => panic!("outbox(\"\") must fail entry-point validation"),
Err(other) => panic!("outbox(\"\") got the wrong error: {other:?}"),
}
match fx.store.outbox("__alkstore_outbox:steal").await {
Err(alkstore::Error::ReservedName { .. }) => {}
Ok(_) => panic!("outbox(reserved) must fail entry-point validation"),
Err(other) => panic!("outbox(reserved) got the wrong error: {other:?}"),
}
let outbox = fx.store.outbox("billing").await.unwrap();
assert_eq!(outbox.name(), "billing");
drop_schema(&fx.admin, &fx.schema).await;
}
/// Acceptance: the derived backing queue is a reserved name —
/// `enqueue_tx` into it rejects (`ReservedName`; ADR-014 §1's "no
/// third path" guarantee) and `enqueue` stamps the 60/5/5 derived set
/// with the one `max_attempts` override (raw row probes: the backing
/// queue has no `Queue` handle surface).
#[tokio::test(flavor = "multi_thread")]
async fn outbox_enqueue_stamps_the_derived_set_and_backing_is_unreachable() {
let Some(fx) = Fixture::open(&instance_namer("ob_stamp")).await else {
eprintln!("skip: no harness server");
return;
};
let err = fx
.store
.begin_tx()
.await
.unwrap()
.enqueue_tx("__alkstore_outbox:billing", Default::default(), json!({}))
.await
.unwrap_err();
assert!(
matches!(err, alkstore::Error::ReservedName { .. }),
"the derived backing queue name must be unreachable by enqueue_tx: {err:?}"
);
let outbox = fx.store.outbox("billing").await.unwrap();
assert_eq!(outbox.name(), "billing");
let id = outbox
.enqueue(json!({"email": "x@y.z"}), EnqueueOpts::default())
.await
.unwrap();
let stamps: (i64, i64, i64, Option<i64>, Vec<u8>) = {
let row = fx
.admin
.query_one(
&format!(
"SELECT visibility_timeout_s, max_attempts, backoff_base_s, \
dead_letter_retention_s, payload FROM {} WHERE id = $1",
job_table(&fx.schema)
),
&[&id],
)
.await
.unwrap();
(row.get(0), row.get(1), row.get(2), row.get(3), row.get(4))
};
assert_eq!((stamps.0, stamps.1, stamps.2, stamps.3), (60, 5, 5, None));
assert_eq!(stamps.4, json!({"email": "x@y.z"}).to_string().into_bytes());
// The one per-job override resolves over the outbox set.
let id2 = outbox
.enqueue(
json!({}),
EnqueueOpts {
max_attempts: Some(2),
..Default::default()
},
)
.await
.unwrap();
let max2: i64 = fx
.admin
.query_one(
&format!(
"SELECT max_attempts FROM {} WHERE id = $1",
job_table(&fx.schema)
),
&[&id2],
)
.await
.unwrap()
.get::<_, i64>(0);
assert_eq!(max2, 2);
assert_eq!(backing_count(&fx.admin, &fx.schema, "billing").await, 2);
drop_schema(&fx.admin, &fx.schema).await;
}
/// Acceptance: `run_once` — `false` on an empty backing queue; ack on
/// `Ok` (the row is gone from live); `retry(Some(e), None)` on `Err`
/// (the row returns to pending with the queue's curve stamped — the
/// curve gap ≥ half the base asserted; the delivery's error string
/// would surface on a later dead row); worker_id validation.
#[tokio::test(flavor = "multi_thread")]
async fn run_once_acks_on_ok_and_retries_on_err() {
let Some(fx) = Fixture::open(&instance_namer("ob_run")).await else {
eprintln!("skip: no harness server");
return;
};
let mut outbox = fx.store.outbox("mail").await.unwrap();
let id_ok = outbox
.enqueue(json!({"n": 1}), Default::default())
.await
.unwrap();
let id_err = outbox
.enqueue(json!({"n": 2}), Default::default())
.await
.unwrap();
// Empty-behavior first on a different outbox: no work → `false`.
let mut quiet = fx.store.outbox("quiet").await.unwrap();
let mut noop = Scripted {
outcome: Ok(()),
seen_payloads: Vec::new(),
};
assert!(!quiet.run_once("w0", &mut noop).await.unwrap());
// Empty worker_id rejects at the entry point.
let err = outbox.run_once("", &mut noop).await.unwrap_err();
assert!(
matches!(err, alkstore::Error::InvalidName { .. }),
"{err:?}"
);
// Ok ⇒ ack: the claim is processed and the row leaves live.
let mut ok_delivery = Scripted {
outcome: Ok(()),
seen_payloads: Vec::new(),
};
assert!(outbox.run_once("w1", &mut ok_delivery).await.unwrap());
assert_eq!(ok_delivery.seen_payloads.len(), 1);
assert_eq!(
ok_delivery.seen_payloads[0],
json!({"n": 1}).to_string().into_bytes()
);
// Err ⇒ retry with the curve: the row returns to pending, the
// backoff curve's draw must land ≥ base/2 (base 5 → draw ∈ [3, 5]).
let mut err_delivery = Scripted {
outcome: Err("smtp down".to_string()),
seen_payloads: Vec::new(),
};
assert!(outbox.run_once("w2", &mut err_delivery).await.unwrap());
assert_eq!(err_delivery.seen_payloads.len(), 1);
assert_eq!(
err_delivery.seen_payloads[0],
json!({"n": 2}).to_string().into_bytes()
);
// The backing queue is reserved: get job state through the raw row
// (the engine's job is get_job-visible via the queue surface only
// with a handle — the reserved name rejects at construction, so
// the row probe here is the honest reader).
let (state, run_at_gap, attempts): (String, i64, i64) = {
let row = fx
.admin
.query_one(
&format!(
"SELECT state, run_at - $2, attempts FROM {} WHERE id = $1",
job_table(&fx.schema)
),
&[&id_err, &now_unix()],
)
.await
.unwrap();
(row.get(0), row.get(1), row.get(2))
};
assert_eq!(state, "pending", "the Err outcome retries into pending");
assert!(
(1..=6).contains(&run_at_gap),
"the queue curve draw ∈ [2, 5] (+1 clk): got gap {run_at_gap}"
);
assert_eq!(attempts, 1, "the first claim counts an attempt");
assert_eq!(
backing_count(&fx.admin, &fx.schema, "mail").await,
1,
"the acked row left live"
);
let _ = id_ok;
drop_schema(&fx.admin, &fx.schema).await;
}
/// Acceptance: `run_once`'s no-heartbeat posture — the engine issues
/// no heartbeat during delivery; the boxed handle's own `heartbeat` is
/// the consumer's renewal path (ADR-019 §5). Exercised as the reach of
/// the boxed handle: the consumer closure heartbeats inside delivery
/// and the row's claim deadline resets forward absolutely (inspected
/// through the raw row — the claimed `Job` snapshot keeps the
/// claim-time deadline by design).
#[tokio::test(flavor = "multi_thread")]
async fn run_once_delivery_can_heartbeat_the_claim_itself() {
let Some(fx) = Fixture::open(&instance_namer("ob_hb")).await else {
eprintln!("skip: no harness server");
return;
};
let mut outbox = fx.store.outbox("hb").await.unwrap();
outbox
.enqueue(json!("work"), Default::default())
.await
.unwrap();
struct Renew {
schema: String,
admin: std::sync::Arc<tokio_postgres::Client>,
}
impl Delivery for Renew {
fn deliver<'a>(
&'a mut self,
job: Box<dyn JobHandle>,
) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> {
let schema = self.schema.clone();
let admin = self.admin.clone();
Box::pin(async move {
assert!(job.heartbeat(3000).await.unwrap(), "renewal lands");
assert_eq!(job.job().state, JobState::Processing);
// Inspect the row while the claim is held (a separate
// connection): the deadline moved to ≈ now + 3000 — the
// heartbeat's absolute reset, not the stamped 60 s
// visibility.
let (claim_expires_at, db_now): (i64, i64) = {
let row = admin
.query_one(
&format!(
"SELECT claim_expires_at, $2::bigint FROM {} WHERE id = $1",
crate::schema::QualifiedTable {
schema: &schema,
name: crate::schema::tables::JOB,
}
),
&[&job.job().id, &now_unix()],
)
.await
.unwrap();
(row.get::<_, i64>(0), row.get::<_, i64>(1))
};
assert!(
claim_expires_at > db_now + 2000,
"the consumer's heartbeat renewed the deadline forward: \
row {claim_expires_at} vs now {db_now}"
);
Ok(())
})
}
}
let renew_admin = harness_client().await.unwrap();
let mut d = Renew {
schema: fx.schema.clone(),
admin: std::sync::Arc::new(renew_admin),
};
assert!(outbox.run_once("w-hb", &mut d).await.unwrap());
drop_schema(&fx.admin, &fx.schema).await;
}
/// Acceptance: `run_once`'s exhaustion path — a delivery that keeps
/// failing past the outbox set's max_attempts (here the per-job
/// override 2) dead-letters the row with the delivery's error string
/// carried (`retry(Some(e), None)`'s exhaust string) — driven to the
/// dead row via repeated `run_once` calls with deadline waits.
#[tokio::test(flavor = "multi_thread")]
async fn run_once_exhaustion_dead_letters_with_the_delivery_error() {
let Some(fx) = Fixture::open(&instance_namer("ob_dead")).await else {
eprintln!("skip: no harness server");
return;
};
let mut outbox = fx.store.outbox("doom").await.unwrap();
let id = outbox
.enqueue(
json!("boom"),
EnqueueOpts {
max_attempts: Some(2),
..Default::default()
},
)
.await
.unwrap();
let mut failing = Scripted {
outcome: Err("always fails".to_string()),
seen_payloads: Vec::new(),
};
// max_attempts = 2 (the per-job override over the outbox set) → two
// claim+fail cycles; the curve's draw ∈ [2, 5]+ with base 5.
for round in 1..=2 {
let deadline = tokio::time::Instant::now() + Duration::from_secs(15);
loop {
if outbox.run_once("w-doom", &mut failing).await.unwrap() {
break;
}
assert!(
tokio::time::Instant::now() < deadline,
"round {round}: retry never came claimable"
);
tokio::time::sleep(Duration::from_millis(150)).await;
}
}
// The second retry exhausts the budget: the row is dead with the
// delivery's error string carried (`retry(Some(e), None)`'s exhaust
// string — ADR-019 §5's outcome posture through the engine).
let (last_error, died): (String, i64) = {
let row = fx
.admin
.query_one(
&format!(
"SELECT last_error, died_at FROM {} WHERE id = $1",
dead_table(&fx.schema)
),
&[&id],
)
.await
.unwrap();
(row.get::<_, String>(0), row.get::<_, i64>(1))
};
assert_eq!(last_error, "payload codec failure: always fails");
assert!(died > 0);
// Nothing left to do: run_once returns false on the drained queue.
assert!(!outbox.run_once("w-doom", &mut failing).await.unwrap());
assert_eq!(failing.seen_payloads.len(), 2);
drop_schema(&fx.admin, &fx.schema).await;
}
File diff suppressed because it is too large. Load diff
+110 -3
View File
@@ -1,7 +1,7 @@
---
id: pg-engine-scheduler-outbox
name: Postgres engine — scheduler (`schedule`/`unschedule`/`run_schedules`) and outbox helper
status: pending
status: completed
depends_on: [pg-engine-queues, pg-engine-locks]
scope: broad
risk: high
@@ -147,8 +147,115 @@ queues task's claim machinery and the locks task's lock machinery:
## Notes
> To be filled by implementation agent
> Decisions of record made in implementation that the description
> didn't pin:
>
> - **Leadership acquire rides the lock module's `lock_acquire`
> directly, not `try_lock`** — `try_lock`'s entry-point shared-name
> validation rejects the reserved name
> (`__alkstore_scheduler`), so the engine's own holder builds the
> lease via `PgLockHandle::new_internal` (the SQLite twin's
> `SqliteLockHandle::new_internal` shape) after the raw acquire op.
> Same machinery, bypassing the consumer entry gate by construction.
> - **The wake-driven tick advance chose the bounded-wait interrupt**
> (the task's implementer's choice): no dedicated wake channel on
> the scheduler's side — the sleep is sliced at 1 s
> (`SLEEP_SLICE_S`) with the stop token checked and the lease
> renewal cadence (2 s) evaluated per slice, and the soonest/
> boundary state is re-read at the top of every loop iteration. A
> newly-registered schedule is therefore noticed at the next loop
> re-read after its boundary (and by the idle-tick's 60 s re-read
> at the latest) — the deadline always fires; the wake never gates
> correctness. The SQLite runner's idle-tick posture stands as the
> floor.
> - **The tick's fire-time `expires` rides `ScheduleOpts::expires`
> as a relative-seconds `EnqueueOpts::expires` through the shared
> `enqueue_row`** — the resolution (relative → absolute from the
> fire instant) is the resolution module's existing formula (one
> owner per formula, ADR-012 §2), not a duplicated scheduler-side
> arithmetic. The `enqueue_row` stamp override re-resolves
> `max_attempts` no-op-fully (the schedule row's stamps are already
> resolved and bind directly).
> - **An `EngineCtx` bundle (`pool`/`schema`/`closed`) threads the
> store wiring into the scheduler/outbox module's entry points**
> — `schedule`'s engine-side argument list stayed within the
> clippy arg guard and the wiring reads once at each call site.
> - **The outbox's `run_once` claim rebuilds the queues task's claim
> frame from the module's shared pieces** (`in_tx` +
> `sweep_exhausted_reclaimables_tx` + `claim_stmt_sql` +
> `job_from_row`, all made `pub(crate)`) with the backing queue
> name standing in for the handle's — one owner per formula, no
> duplicated claim SQL; the ok/err dispositions route through
> `PgJobHandle::ack`/`retry` (the `new` constructor made
> `pub(crate)` for the rebuild, mirroring the SQLite twin's
> `SqliteJobHandle::new`).
> - **The best-effort enqueue wake on the backing queue's channel**
> (the mechanism-name-is-the-channel realization) rides the outbox
> `enqueue` too — the queue enqueue's twin posture (durable row is
> the truth, the re-poll safety net covers a failed wake).
> - **The open-tests' stub-probe test retired** — all four
> `schedule`/`unschedule`/`run_schedules`/`outbox` methods are
> wired, so `store_trait_methods_are_wiring_stubs` (which asserted
> the `Err(Database(… wiring lands with …))` shape on them) now
> asserts the wired constructor/register/unregister surface
> instead; the `stub_err!` macro left with the removed stub calls.
> - **`tick` is `pub(crate)` (test-visible)** for the rogue-ticker
> row-lock probe — two concurrent pool connections driving the
> engine's own tick op directly (no leadership lock in either
> hand) fire exactly one boundary.
## Summary
> To be filled on completion
> What landed, verified how:
>
> - `alkstore-postgres/src/scheduler.rs` — the collapse shape's pg
> arm + the outbox helper: `schedule` (the validation triad
> pre-storage, the pg-owned `@every`-only grammar
> `parse_every_interval` — cron strings reject pre-storage — the
> upsert over the schema table with the plain-queue derived stamps
> + `ScheduleOpts` applied, `Schedule` read-back), `unschedule`
> (true/false), `run_schedules` (leadership via the engine's lock
> machinery on `__alkstore_scheduler` with a per-instance pid +
> subsec-nanos owner token, TTL 10 s engine-internal, renew at
> every loop top and on a 2 s cadence across the sliced sleep,
> acquire-time refusal = `Err(LeadershipLost)`, the row-locked
> `FOR UPDATE` tick in one pool-checked-out `BEGIN`/`COMMIT` frame
> — fire enqueues + boundary advance + soonest read commit
> together — the 64-boundary catch-up cap with skip-forward
> past-now, clean stop `Ok(())` releasing the lock, loss
> `Err(LeadershipLost)` before any tick); `outbox(name)` — the
> validated constructor, `enqueue` into the derived
> `__alkstore_outbox:{name}` (60/5/5 stamps, the `max_attempts`
> override, full ADR-020 resolution, best-effort wake) and
> `run_once(worker_id, delivery)` (the ordinary claim frame, ack
> on `Ok`, `retry(err-as-exhaust-string, curve-delay)` on `Err`,
> `false` on empty, no engine-issued heartbeat — the boxed
> handle's own `heartbeat` is the consumer's).
> - Store trait wiring (`store.rs`): the four stubs replaced; the
> `stub` fn and its last call sites removed. Queue module deltas:
> `in_tx`/`sweep_exhausted_reclaimables_tx`/`claim_stmt_sql`/
> `PgJobHandle::new` made `pub(crate)` (one-owned machinery
> reuse); lock module deltas: `lock_acquire` made `pub(crate)` +
> `PgLockHandle::new_internal` added.
> - Tests: `store/scheduler_tests.rs` (9 — the validation triad +
> exact parser surface, upsert/unregister read-backs, fire + stamp
> source via `get_job` + live leadership lock row with the
> pid-stamped owner + clean-stop release, second-runner refusal +
> the one-firer band, acquire-time loss with a zero-fire probe,
> in-sleep lease renewals holding across a 40 s boundary on a 10 s
> TTL, the 64-cap + skip-forward with no in-window refire from a
> 200-boundary far-out horizon, the rogue-ticker row-lock
> serialization, and the server-less parser unit pin) and
> `store/outbox_tests.rs` (5 — constructor validation,
> `enqueue_tx` unreachability + 60/5/5 stamps + override, ok/err/
> false dispositions with the curve gap asserted, in-delivery
> heartbeat renewal, exhaustion to dead with the delivery error
> string). `open_tests.rs`'s stub probe updated to the wired
> posture.
> - Verified: `cargo test -p alkstore-postgres` green against the
> harness server (111 + 9, six consecutive full runs — stable);
> workspace `cargo test` green server-less (the pg skips cleanly);
> `cargo clippy --all-targets -- -D warnings` and
> `cargo fmt --check` clean workspace-wide. Pre-existing flake
> (streams' wake-timing deadline race) untouched — none of the six
> runs tripped it.