SQLite engine: queues — Queue/JobHandle over the writer slot, engine-side backoff curve, extent guard (task sqlite-engine-queues)

queue.rs: SqliteQueueHandle (QueueOpts-carrying, stamp-resolving
enqueue over the seam's shared resolution arithmetic), claim_one/
claim_batch through the writer slot with ADR-023 §2's extent guard
(n <= 0 -> empty Vec at the trait-impl entry), the full JobHandle impl
(one-shot ack/retry/fail as 'static boxed futures, repeatable absolute
reset heartbeat, substrate's uniform validity predicate), and the
engine-owned equal-jitter exponential backoff (std-only RandomState
hash jitter, integerized inclusive [ceil(half), cap], 1-hour cap;
no rand dep - documented).

ack_batch loses its substrate worker filter (register D-31 - ADR-019
§1's worker-less batch form); job_from_json decode + stamps_with_
override moved to their one owners (queue.rs / resolution.rs);
reader-pool helpers lifted from stream.rs into seam.rs. Store::queue
wired; 10 acceptance tests (lifecycle/stamps, extent guard,
exactly-once under concurrency, backoff range + cap + override,
validity predicate + reclaim, dead-letter defaults + get_job
visibility, cancel, ack_batch, sweep both-states + retention,
validation + closed-store). Workspace 25+3+171 green, clippy
-D warnings, fmt clean; sqlite suite 4x green.
This commit is contained in:
glm-5.3-flash committed 2026-10-08 12:59:57 +00:00
1 parent 4913302b47
commit 513df0b311
11 files changed
+1708 -195

No files matched your search

+1
View File
@@ -19,6 +19,7 @@
mod notify;
mod opts;
mod queue;
mod resolution;
mod seam;
mod store;
+472
View File
@@ -0,0 +1,472 @@
//! The queue mechanism's SQLite arm (`sqlite-engine-queues`) — the
//! auto-commit counterpart to the tx substrate's queue ops:
//!
//! - `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 the
//! [`QueueOpts`] stamps future enqueues resolve over.
//! - `enqueue(payload, opts)` — the one handle-carrying shape: stamps
//! resolve over **this handle's opts** (ADR-020 §3a's handle
//! annotation — the handle's visibility/backoff/retention stamp the
//! row directly; `EnqueueOpts::max_attempts` is the one per-job
//! override), through the same resolution arithmetic the tx path
//! uses (one clock read per enqueue, delay-over-`run_at`,
//! relative-`expires`, `encode_payload` at the seam).
//! - `claim_one` / `claim_batch` — through the writer slot (claims
//! are writes; the substrate's single-statement claim is the
//! exactly-once handout, ADR-010 §1). The **extent guard (ADR-023
//! §2)** short-circuits `n <= 0` to the empty `Vec` at the
//! trait-impl entry, before the substrate — the `LIMIT -1` dialect
//! artifact dies here (the substrate passes `n` through
//! contract-blind). `worker_id` is consumer-local, non-empty only
//! (ADR-019 §2 — `InvalidName`). Claim rows decode into `Job` +
//! boxed [`JobHandle`]s via the shared `job_from_json` decode owner.
//! - 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 is
//! the substrate's (processing + unexpired claim deadline, ADR-010
//! §2); 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 (ADR-023 §2's
//! boundary class — any `i64` is total).
//! - `ack_batch(ids)` — the contract's batch ack carries no
//! `worker_id` (ADR-019 §1's surface): the substrate's per-id
//! validity predicate (processing + unexpired deadline) counts
//! landed ids; partial success is ordinary (ADR-010 §1).
//! - `cancel(job_id)` (unconditional delete), `get_job(job_id)`
//! (sees dead rows — `last_error`/`died_at` carried), and
//! `sweep_expired()` (both-states move + retention enforcement; no
//! leader lock required, ADR-010 §5) round out the maintenance
//! surface.
//! - Dead-letter strings (row content, not error variants): the
//! pinned defaults for the engine-triggered paths
//! (`"max attempts exceeded"` at exhaustion, `"failed"` for
//! `fail(None)` per ADR-019 §3's annotation, `"expired"` on sweeps,
//! ADR-010 §4); a caller-supplied `Some(err)` carries the caller's
//! string on the path it names.
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 serde_json::Value;
use alkstore::{
BoxedFuture, EnqueueOpts, Error, Job, JobHandle, JobState, Queue, QueueOpts, Result,
validate_local_name,
};
use crate::resolution::{resolve_enqueue_opts, stamps_with_override};
use crate::seam::{closed_store_error, sqlite_error, with_reader, with_writer};
use crate::substrate::{
Readers, Stamps, Writer, ack, ack_batch, cancel, claim_batch, enqueue, fail, get_job,
heartbeat, retry, sweep_expired,
};
/// The pinned dead-letter default strings (row content, ADR-010 §4 /
/// ADR-019 §3's annotation — the engine passes them; the substrate
/// stores them verbatim).
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 SqliteQueueHandle {
name: String,
opts: QueueOpts,
writer: Arc<Writer>,
readers: Arc<Readers>,
closed: Arc<AtomicBool>,
}
impl std::fmt::Debug for SqliteQueueHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SqliteQueueHandle")
.field("name", &self.name)
.field("opts", &self.opts)
.finish_non_exhaustive()
}
}
impl SqliteQueueHandle {
pub(crate) fn new(
name: String,
opts: QueueOpts,
writer: Arc<Writer>,
readers: Arc<Readers>,
closed: Arc<AtomicBool>,
) -> Self {
Self {
name,
opts,
writer,
readers,
closed,
}
}
/// The handle's opts, as the substrate stamps 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,
}
}
}
/// Decode the substrate's job-row JSON into the contract value
/// (`Job::from_row` — the engine-side constructor, ADR-021 §2's
/// `claimed_at` surfaced). One decode owner: the tx reads (`tx.rs`)
/// and this module's claim/get rows share it; malformed rows map to
/// `Error::Codec` (the decode-side posture guard — the fork's writer
/// always wrote well-formed rows).
pub(crate) fn job_from_json(value: &Value) -> Result<Option<Job>> {
let state = match value["state"].as_str() {
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: Vec<u8> = value["payload"]
.as_str()
.ok_or_else(|| Error::Codec("job row payload must be a string".to_string()))?
.as_bytes()
.to_vec();
let field = |name: &str| -> Result<i64> {
serde_json::from_value(value[name].clone())
.map_err(|e| Error::Codec(format!("job row {name} must be i64: {e}")))
};
let opt_field = |name: &str| -> Result<Option<i64>> {
if value[name].is_null() {
Ok(None)
} else {
serde_json::from_value(value[name].clone())
.map_err(|e| Error::Codec(format!("job row {name} must be i64: {e}")))
}
};
let opt_str = |name: &str| -> Result<Option<String>> {
if value[name].is_null() {
Ok(None)
} else {
serde_json::from_value(value[name].clone())
.map_err(|e| Error::Codec(format!("job row {name} must be a string: {e}")))
}
};
Ok(Some(Job::from_row(
field("id")?,
value["queue"]
.as_str()
.ok_or_else(|| Error::Codec("job row queue must be a string".to_string()))?
.to_string(),
state,
payload,
field("priority")?,
field("run_at")?,
field("attempts")?,
field("max_attempts")?,
opt_str("worker_id")?,
opt_field("claimed_at")?,
opt_field("claim_expires_at")?,
field("created_at")?,
opt_field("expires_at")?,
field("visibility_timeout_s")?,
field("backoff_base_s")?,
opt_field("dead_letter_retention_s")?,
opt_str("last_error")?,
opt_field("died_at")?,
)))
}
/// Decode a claim page of job-row JSON into `Job` values.
fn jobs_from_json_page(page: &str) -> Result<Vec<Job>> {
let array: Value = serde_json::from_str(page)
.map_err(|e| Error::Codec(format!("claim page must be json: {e}")))?;
let rows = array
.as_array()
.ok_or_else(|| Error::Codec("claim page must be a json array".to_string()))?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
out.push(
job_from_json(row)?
.ok_or_else(|| Error::Codec("claim page row must decode to a job".to_string()))?,
);
}
Ok(out)
}
/// The equal-jitter exponential backoff (ADR-010 §3 — engine-side,
/// one owner; wave 4's pg engine re-derives it 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 jitter draw is std-only (no `rand` dependency without a
/// decision — the ADR-005 dependency posture; documented 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
}
impl Queue for SqliteQueueHandle {
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_store_error()) });
}
let name = self.name.clone();
let base = self.base_stamps();
Box::pin(async move {
with_writer(self.writer.clone(), move |conn| {
let bytes = alkstore::encode_payload(&payload)?;
let resolved = resolve_enqueue_opts(conn, &opts)?;
let stamps = stamps_with_override(base, &opts);
enqueue(
conn,
&name,
std::str::from_utf8(&bytes)
.map_err(|e| Error::Codec(format!("row payload must be utf-8: {e}")))?,
resolved.run_at,
opts.priority,
resolved.expires_at,
stamps,
)
.map_err(sqlite_error)
})
.await
})
}
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_store_error()) });
}
// The extent guard (ADR-023 §2) — before the substrate, whose
// `LIMIT ?` passes `n` through contract-blind (SQLite's
// `LIMIT -1` = unbounded is exactly the artifact that dies
// here).
if n <= 0 {
return Box::pin(async { Ok(Vec::new()) });
}
let name = self.name.clone();
let writer = self.writer.clone();
let closed = self.closed.clone();
let worker = worker_id.to_string();
Box::pin(async move {
with_writer(writer.clone(), move |conn| {
let page =
claim_batch(conn, &name, &worker, n, EXHAUSTED_ERROR).map_err(sqlite_error)?;
Ok(jobs_from_json_page(&page)?
.into_iter()
.map(|job| {
Box::new(SqliteJobHandle {
job,
writer: writer.clone(),
closed: closed.clone(),
}) as Box<dyn JobHandle>
})
.collect())
})
.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_store_error()) });
}
let ids_json =
serde_json::to_string(ids).map_err(|e| Error::Codec(format!("ids must encode: {e}")));
Box::pin(async move {
let ids_json = ids_json?;
with_writer(self.writer.clone(), move |conn| {
ack_batch(conn, &ids_json).map_err(sqlite_error)
})
.await
})
}
fn cancel<'a>(&'a self, job_id: i64) -> BoxedFuture<'a, Result<bool>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
Box::pin(async move {
with_writer(self.writer.clone(), move |conn| {
Ok(cancel(conn, job_id).map_err(sqlite_error)? > 0)
})
.await
})
}
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_store_error()) });
}
let readers = self.readers.clone();
Box::pin(async move {
with_reader(readers, move |conn| {
let row = get_job(conn, job_id).map_err(sqlite_error)?;
match row {
Some(json) => {
let value: Value = serde_json::from_str(&json)
.map_err(|e| Error::Codec(format!("job row must be json: {e}")))?;
job_from_json(&value)
}
None => Ok(None),
}
})
.await
})
}
fn sweep_expired<'a>(&'a self) -> BoxedFuture<'a, Result<i64>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
let name = self.name.clone();
Box::pin(async move {
with_writer(self.writer.clone(), move |conn| {
sweep_expired(conn, &name, EXPIRED_ERROR).map_err(sqlite_error)
})
.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 SqliteJobHandle {
job: Job,
writer: Arc<Writer>,
closed: Arc<AtomicBool>,
}
impl std::fmt::Debug for SqliteJobHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SqliteJobHandle")
.field("job", &self.job)
.finish_non_exhaustive()
}
}
impl JobHandle for SqliteJobHandle {
fn job(&self) -> &Job {
&self.job
}
fn ack(self: Box<Self>) -> BoxedFuture<'static, Result<bool>> {
let SqliteJobHandle { job, writer, .. } = *self;
let id = job.id;
let worker_id = job.worker_id.unwrap_or_default();
Box::pin(async move {
with_writer(writer, move |conn| {
Ok(ack(conn, id, &worker_id).map_err(sqlite_error)? > 0)
})
.await
})
}
fn retry(
self: Box<Self>,
err: Option<String>,
delay: Option<i64>,
) -> BoxedFuture<'static, Result<bool>> {
let SqliteJobHandle { job, writer, .. } = *self;
let id = job.id;
let worker_id = job.worker_id.unwrap_or_default();
let exhaust_error = err.unwrap_or_else(|| EXHAUSTED_ERROR.to_string());
let delay_s = match delay {
Some(d) => d,
None => backoff_delay_s(job.backoff_base_s, job.attempts, id),
};
Box::pin(async move {
with_writer(writer, move |conn| {
Ok(retry(conn, id, &worker_id, delay_s, &exhaust_error).map_err(sqlite_error)? > 0)
})
.await
})
}
fn fail(self: Box<Self>, err: Option<String>) -> BoxedFuture<'static, Result<bool>> {
let SqliteJobHandle { job, writer, .. } = *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 {
with_writer(writer, move |conn| {
Ok(fail(conn, id, &worker_id, &error).map_err(sqlite_error)? > 0)
})
.await
})
}
fn heartbeat<'a>(&'a self, extend: i64) -> BoxedFuture<'a, Result<bool>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
let writer = self.writer.clone();
let id = self.job.id;
let worker_id = self.job.worker_id.clone().unwrap_or_default();
Box::pin(async move {
with_writer(writer, move |conn| {
Ok(heartbeat(conn, id, &worker_id, extend).map_err(sqlite_error)? > 0)
})
.await
})
}
}
+15
View File
@@ -64,6 +64,21 @@ pub(crate) fn outbox_backing_queue_name(outbox: &str) -> String {
format!("{}outbox:{outbox}", alkstore::RESERVED_PREFIX)
}
/// The one per-enqueue stamp override (ADR-010 §3a):
/// `EnqueueOpts::max_attempts` is the only per-job override — the
/// other stamps (visibility/backoff/retention) always take the base's.
///
/// One helper both enqueue shapes share: the tx paths (`enqueue_tx`
/// over the plain-queue defaults, `outbox_enqueue_tx` over the
/// outbox's 60/5/5 set) and the auto-commit `Queue::enqueue` over the
/// handle's `QueueOpts` — resolved once, engine-wide (ADR-012 §2).
pub(crate) fn stamps_with_override(base: Stamps, opts: &EnqueueOpts) -> Stamps {
Stamps {
max_attempts: opts.max_attempts.unwrap_or(base.max_attempts),
..base
}
}
/// Resolve [`EnqueueOpts`] against the enqueue instant: the row's ready
/// time (`delay` → `now + delay`, wins over `run_at` (absolute);
/// neither → `now`) and the absolute row expiry (`expires` relative
+85 -4
View File
@@ -11,18 +11,19 @@
//! inside a bridged closure is the `JoinError` arm's `Database`
//! mapping, family-standard no-panics discipline maintained.
use alkstore::Error;
use alkstore::{Error, Result};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use crate::substrate::Writer;
use crate::substrate::{Readers, SchemaError, Writer};
/// The blocking bridge every trait method's substrate round trip goes
/// through. Private: the seam is engine-internal — consumers meet it
/// only through the async trait surface.
pub(crate) async fn blocking<T, F>(f: F) -> Result<T, Error>
pub(crate) async fn blocking<T, F>(f: F) -> Result<T>
where
T: Send + 'static,
F: FnOnce() -> Result<T, Error> + Send + 'static,
F: FnOnce() -> Result<T> + Send + 'static,
{
match tokio::task::spawn_blocking(f).await {
Ok(result) => result,
@@ -75,3 +76,83 @@ where
})
.await
}
/// A reader-pool read through the `spawn_blocking` seam: acquire a
/// pooled reader, run the body, release on both arms (a failed SELECT
/// leaves the pooled connection healthy). A closed pool (store closed)
/// fails closed with `Database`.
pub(crate) async fn with_reader<T, F>(readers: Arc<Readers>, f: F) -> Result<T>
where
T: Send + 'static,
F: FnOnce(&rusqlite::Connection) -> Result<T> + Send + 'static,
{
blocking(move || {
let conn = acquire_reader(&readers)?;
let out = f(&conn);
readers.release(conn);
out
})
.await
}
pub(crate) fn acquire_reader(readers: &Arc<Readers>) -> Result<rusqlite::Connection> {
match readers.acquire() {
Ok(conn) => Ok(conn),
Err(SchemaError::Sqlite(e)) => {
if is_closed_err(&e) {
Err(closed_store_error())
} else {
Err(sqlite_error(e))
}
}
}
}
/// A synchronous pooled read used by bridge threads (the stream and
/// queue subscribe/bridge shapes): the closed arms are distinguished —
/// a closed store exits the loop, an open store's failure is transient
/// (`Failed` carries the mapped error).
pub(crate) fn blocking_read<T>(
readers: &Arc<Readers>,
store_closed: &Arc<AtomicBool>,
f: impl FnOnce(&rusqlite::Connection) -> Result<T>,
) -> std::result::Result<T, ReaderLoopExit> {
if store_closed.load(Ordering::Acquire) {
return Err(ReaderLoopExit::Closed);
}
let conn = match readers.acquire() {
Ok(conn) => conn,
Err(e) => match e {
SchemaError::Sqlite(inner) if is_closed_err(&inner) => {
return Err(ReaderLoopExit::Closed);
}
SchemaError::Sqlite(inner) => {
return Err(ReaderLoopExit::Failed(sqlite_error(inner)));
}
},
};
let out = f(&conn);
readers.release(conn);
match out {
Ok(value) => Ok(value),
Err(_e) if store_closed.load(Ordering::Acquire) => Err(ReaderLoopExit::Closed),
Err(e) => Err(ReaderLoopExit::Failed(e)),
}
}
/// The blocked-read exit arms (see [`blocking_read`]).
pub(crate) enum ReaderLoopExit {
Closed,
Failed(Error),
}
pub(crate) fn is_closed_err(e: &rusqlite::Error) -> bool {
matches!(
e,
rusqlite::Error::SqliteFailure(_, Some(msg)) if msg.contains("Database is closed")
)
}
pub(crate) fn closed_store_error() -> Error {
Error::database(std::io::Error::other("the store is closed"))
}
+18 -6
View File
@@ -166,13 +166,23 @@ impl Store for SqliteStore {
fn queue<'a>(
&'a self,
_name: &str,
_opts: alkstore::QueueOpts,
name: &str,
opts: alkstore::QueueOpts,
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::Queue>>> {
Box::pin(async {
Err(database_error(
"queue wiring lands with the mechanism tasks",
))
if let Err(e) = alkstore::validate_shared_name(name) {
return Box::pin(async move { Err(e) });
}
if self.closed.load(std::sync::atomic::Ordering::Acquire) {
return Box::pin(async { Err(crate::seam::database_error("the store is closed")) });
}
let name = name.to_string();
let writer = self.writer.clone();
let readers = self.readers.clone();
let closed = self.closed.clone();
Box::pin(async move {
Ok(Box::new(crate::queue::SqliteQueueHandle::new(
name, opts, writer, readers, closed,
)) as Box<dyn alkstore::Queue>)
})
}
@@ -236,6 +246,8 @@ mod notify_tests;
#[cfg(test)]
mod open_tests;
#[cfg(test)]
mod queue_tests;
#[cfg(test)]
mod stream_tests;
#[cfg(test)]
mod tx_tests;
+999
View File
@@ -0,0 +1,999 @@
//! The queue mechanism's acceptance tests (`sqlite-engine-queues`):
//! the enqueue→claim→ack lifecycle works end-to-end over a temp-file
//! store with stamps visible via `get_job` matching the resolution
//! rules, the extent guard (ADR-023 §2) claims nothing on `n <= 0`,
//! the exactly-once handout under concurrent claims (the POC-pinned
//! property, engine-tested here), the backoff curve's range pin
//! (statistical, tolerance-bounded) with the 1-hour cap and the
//! `Some(d)` override, the uniform validity predicate refusing lapsed
//! ops with `Ok(false)` and a reclaim consuming an attempt, the
//! dead-letter defaults and dead-row `get_job` visibility, cancel,
//! `ack_batch`'s per-id predicate without a worker filter, and
//! `sweep_expired`'s both-states move + retention count (ADR-010;
//! ADR-019 §1–§3; ADR-020 §1–§3; ADR-023 §2).
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use alkstore::{EnqueueOpts, JobState, QueueOpts, Store};
use serde_json::{Value, json};
use crate::queue::{EXHAUSTED_ERROR, EXPIRED_ERROR, FAILED_ERROR, backoff_delay_s};
use crate::store::open_store;
fn temp_dir(tag: &str) -> PathBuf {
let d = std::env::temp_dir().join(format!(
"alkstore-queue-{tag}-{}-{:?}",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_dir_all(&d);
std::fs::create_dir_all(&d).unwrap();
d
}
fn temp_path(tag: &str) -> PathBuf {
temp_dir(tag).join("store.db")
}
fn opts() -> QueueOpts {
QueueOpts {
visibility_timeout_s: 300,
max_attempts: 3,
backoff_base_s: 5,
dead_letter_retention_s: None,
}
}
/// Acceptance: the full lifecycle — enqueue stamps the resolved values
/// (resolved `run_at` after delay-over-`run_at`, absolute `expires_at`
/// after relative-`expires`, the handle's `QueueOpts` stamps with the
/// per-job `max_attempts` override), claim returns a boxed handle
/// carrying the row, ack deletes, and `get_job` sees the states along
/// the way (ADR-010 §1/§3a; ADR-020 §1–§2; ADR-019 §1).
#[tokio::test(flavor = "multi_thread")]
async fn enqueue_claim_ack_lifecycle_with_resolved_stamps() {
let dir = temp_dir("lifecycle");
let path = temp_path("lifecycle");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let qopts = QueueOpts {
visibility_timeout_s: 120,
max_attempts: 4,
backoff_base_s: 9,
dead_letter_retention_s: Some(7200),
};
let queue = store.queue("emails", qopts.clone()).await.unwrap();
assert_eq!(queue.name(), "emails");
// Delay wins over run_at (ADR-020 §1); relative expires resolves
// to the absolute row expiry (§2).
let before = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as i64;
let id = queue
.enqueue(
json!({"to": "a@example.com"}),
EnqueueOpts {
delay: Some(60),
run_at: Some(before + 10_000),
priority: 3,
max_attempts: Some(7),
expires: Some(1200),
},
)
.await
.unwrap();
let job = queue.get_job(id).await.unwrap().expect("row exists");
assert_eq!(job.id, id);
assert_eq!(job.queue, "emails");
assert_eq!(job.state, JobState::Pending);
assert_eq!(job.priority, 3);
assert!(
(before + 60..=before + 62).contains(&job.run_at),
"ready time is now + delay, delay wins over run_at: {}",
job.run_at
);
assert_eq!(
job.expires_at,
Some(before + 1200),
"relative expires resolves to the absolute row expiry"
);
assert_eq!(job.max_attempts, 7, "the per-job override stamps the row");
assert_eq!(job.visibility_timeout_s, 120, "the handle's stamp");
assert_eq!(job.backoff_base_s, 9, "the handle's stamp");
assert_eq!(
job.dead_letter_retention_s,
Some(7200),
"the handle's stamp"
);
assert_eq!(job.attempts, 0);
assert_eq!(job.worker_id, None);
assert_eq!(job.claimed_at, None);
assert_eq!(job.claim_expires_at, None);
assert_eq!(job.last_error, None);
assert_eq!(job.died_at, None);
let decoded: Value = job.payload_as().unwrap();
assert_eq!(decoded, json!({"to": "a@example.com"}));
// Not yet claimable — ready time is a minute out (the resolved
// delay stamps the row).
assert!(queue.claim_one("w1").await.unwrap().is_none());
// A second enqueue with run_at in the past is claimable now, the
// one-clocked shape: relative expires resolves over the same
// instant as the ready time.
let id2 = queue
.enqueue(
json!({"n": 2}),
EnqueueOpts {
run_at: Some(before - 10),
..Default::default()
},
)
.await
.unwrap();
let job2 = queue.get_job(id2).await.unwrap().unwrap();
assert!(
(before - 10..=before - 9).contains(&job2.run_at),
"run_at is used literally when delay is unset"
);
// Claim: the handle carries the row, the claim counts an attempt.
let handle = queue.claim_one("w1").await.unwrap().expect("claimable");
let job = handle.job();
assert_eq!(job.id, id2);
assert_eq!(job.state, JobState::Processing);
assert_eq!(job.attempts, 1);
assert_eq!(job.worker_id.as_deref(), Some("w1"));
assert!(job.claimed_at.is_some());
assert_eq!(
job.claim_expires_at,
Some(job.claimed_at.unwrap() + 120),
"the deadline comes from the job's own visibility stamp"
);
// get_job reflects the claim.
let seen = queue.get_job(id2).await.unwrap().unwrap();
assert_eq!(seen.state, JobState::Processing);
assert_eq!(seen.worker_id.as_deref(), Some("w1"));
// A second worker cannot claim it (exactly-once handout); the
// delayed row is still not claimable either.
assert!(queue.claim_one("w2").await.unwrap().is_none());
// Ack deletes the row (ADR-010 §1).
assert!(handle.ack().await.unwrap());
assert!(queue.get_job(id2).await.unwrap().is_none());
assert!(queue.claim_one("w2").await.unwrap().is_none());
assert!(
queue.get_job(id).await.unwrap().is_some(),
"the delayed row survives untouched"
);
store.close();
let _ = std::fs::remove_dir_all(dir);
}
/// Acceptance: the extent guard (ADR-023 §2) — `claim_batch(n <= 0)`
/// claims nothing (empty `Vec`, no error, never the whole queue — the
/// `LIMIT -1` dialect artifact is dead) and `claim_one` on an empty
/// queue is `None` (no-work is a value).
#[tokio::test(flavor = "multi_thread")]
async fn extent_guard_claims_nothing_on_non_positive_n() {
let dir = temp_dir("extent");
let path = temp_path("extent");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let queue = store.queue("work", opts()).await.unwrap();
for i in 0..5 {
queue
.enqueue(json!({"i": i}), EnqueueOpts::default())
.await
.unwrap();
}
for n in [-1, 0] {
let got = queue.claim_batch("w1", n).await.unwrap();
assert!(
got.is_empty(),
"claim_batch({n}) must claim nothing (empty Vec, no error)"
);
}
// Nothing was destroyed by the guarded call.
let seen = queue.get_job(1).await.unwrap().unwrap();
assert_eq!(seen.state, JobState::Pending);
// A positive claim on the same state still works and hands out
// exactly the requested extent.
let got = queue.claim_batch("w1", 2).await.unwrap();
assert_eq!(got.len(), 2);
store.close();
let _ = std::fs::remove_dir_all(dir);
}
/// Acceptance: exactly-once handout under concurrent claims — two
/// workers racing claim on one row; exactly one wins (the POC-pinned
/// property, engine-tested here through the writer slot), and a
/// many-worker many-row drain hands every row out exactly once.
#[tokio::test(flavor = "multi_thread")]
async fn concurrent_claims_hand_each_row_out_exactly_once() {
let dir = temp_dir("exclusive");
let path = temp_path("exclusive");
let store: Arc<crate::store::SqliteStore> =
Arc::new(open_store(path.to_str().unwrap(), Default::default()).unwrap());
// Two workers, one row: exactly one claim wins.
{
let queue = store.queue("solo", opts()).await.unwrap();
let id = queue
.enqueue(json!({"only": true}), EnqueueOpts::default())
.await
.unwrap();
let q1 = store.queue("solo", opts()).await.unwrap();
let q2 = store.queue("solo", opts()).await.unwrap();
let (a, b) = tokio::join!(q1.claim_one("w1"), q2.claim_one("w2"));
let winners = [a.unwrap(), b.unwrap()]
.into_iter()
.flatten()
.collect::<Vec<_>>();
assert_eq!(winners.len(), 1, "exactly one worker got the row");
assert_eq!(winners[0].job().id, id);
}
// Many workers drain a whole queue with no double handout.
const JOBS: i64 = 40;
const WORKERS: usize = 6;
{
let queue = store.queue("shared", opts()).await.unwrap();
for i in 0..JOBS {
queue
.enqueue(json!({"i": i}), EnqueueOpts::default())
.await
.unwrap();
}
let claimed = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
for w in 0..WORKERS {
let store = store.clone();
let claimed = claimed.clone();
handles.push(tokio::spawn(async move {
let queue = store.queue("shared", opts()).await.unwrap();
let mut got = 0usize;
loop {
let batch = queue
.claim_batch(&format!("w{w}"), 4)
.await
.expect("claim must not error under contention");
if batch.is_empty() {
break;
}
got += batch.len();
// Hold the rows (no ack) — the drain semantics are
// what is under test.
}
claimed.fetch_add(got, Ordering::SeqCst);
got
}));
}
let mut shares = Vec::new();
for h in handles {
shares.push(h.await.unwrap());
}
assert_eq!(
claimed.load(Ordering::SeqCst),
JOBS as usize,
"every job handed out exactly once"
);
assert!(
shares.iter().any(|s| *s > 0) && shares.iter().any(|s| *s < JOBS as usize),
"work was shared across workers, not drained by one"
);
// No id in processing twice; every job is processing now.
let queue = store.queue("shared", opts()).await.unwrap();
let mut distinct = std::collections::HashSet::new();
for id in 2..(JOBS + 2) {
if let Some(job) = queue.get_job(id).await.unwrap() {
assert_eq!(
job.state,
JobState::Processing,
"job {id} must be processing"
);
distinct.insert(job.id);
}
}
assert_eq!(distinct.len(), JOBS as usize, "no id claimed twice");
}
store.close();
let _ = std::fs::remove_dir_all(dir);
}
/// Acceptance: the backoff curve — `retry(err, None)` computes the
/// equal-jitter exponential `delay ∈ [base·2^(a−1)/2, base·2^(a−1)]`
/// from the job's stamped `backoff_base_s`, capped at 1 hour
/// (ADR-010 §3): statistical pin over many draws, tolerance-bounded;
/// the end-to-end row's `run_at` lands in the range; `Some(d)`
/// overrides.
#[tokio::test(flavor = "multi_thread")]
async fn backoff_curve_lands_in_the_pinned_range() {
// Statistical pin of the pure formula (a = attempts index):
// base 5, a=1 → ceil(5/2)=3..5; draws must stay in the inclusive
// integer range and spread across most of it.
for (base, attempts, lo, hi) in [(5, 1, 3, 5), (5, 2, 5, 10), (5, 3, 10, 20), (1, 1, 1, 1)] {
let mut lo_seen = i64::MAX;
let mut hi_seen = i64::MIN;
for id in 0..400 {
let d = backoff_delay_s(base, attempts, id);
assert!(
(lo..=hi).contains(&d),
"delay {d} outside [{lo}, {hi}] for base={base} a={attempts}"
);
lo_seen = lo_seen.min(d);
hi_seen = hi_seen.max(d);
}
if lo == hi {
assert_eq!(lo_seen, hi_seen, "degenerate range pins to the value");
} else {
assert!(
lo_seen <= lo + (hi - lo) / 2 && hi_seen >= hi - (hi - lo) / 2,
"draws must spread over the range ([{lo_seen}, {hi_seen}] vs [{lo}, {hi}])"
);
}
}
// Cap at 1 hour: a huge doubling exponent still lands in
// [1800, 3600].
for id in 0..200 {
let d = backoff_delay_s(5, 40, id);
assert!((1800..=3600).contains(&d), "capped delay {d} out of range");
}
// End-to-end: the row's run_at after retry(err, None) is now plus
// a curve draw from the job's stamped base.
let dir = temp_dir("backoff");
let path = temp_path("backoff");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let queue = store
.queue(
"retrying",
QueueOpts {
visibility_timeout_s: 300,
max_attempts: 5,
backoff_base_s: 8,
dead_letter_retention_s: None,
},
)
.await
.unwrap();
let id = queue
.enqueue(json!({"task": 1}), EnqueueOpts::default())
.await
.unwrap();
let handle = queue.claim_one("w1").await.unwrap().unwrap();
assert_eq!(handle.job().attempts, 1, "a = the post-claim count");
let before = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as i64;
assert!(handle.retry(None, None).await.unwrap());
let job = queue.get_job(id).await.unwrap().unwrap();
assert_eq!(job.state, JobState::Pending);
let span_lo = before - 1 + 8 / 2;
let span_hi = before + 1 + 8;
assert!(
(span_lo..=span_hi).contains(&job.run_at),
"run_at {} must be now + a draw from [4, 8] (base 8, a=1)",
job.run_at
);
// Some(d) overrides the curve: wait out the curve draw, claim
// again, retry with the explicit delay.
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30);
let override_landed = loop {
match queue.claim_one("w1").await {
Ok(Some(h)) => {
let ok = h
.retry(Some("still broken".to_string()), Some(123))
.await
.unwrap();
assert!(ok, "the override retry lands");
break true;
}
Ok(None) => {}
Err(e) => panic!("claim errored: {e}"),
}
assert!(
tokio::time::Instant::now() < deadline,
"curve draw never elapsed"
);
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
};
assert!(override_landed);
let job = queue.get_job(id).await.unwrap().unwrap();
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as i64;
assert!(
(now + 123 - 2..=now + 123 + 2).contains(&job.run_at) && job.run_at >= before + 120,
"the explicit delay passes through: {} (now {now})",
job.run_at
);
store.close();
let _ = std::fs::remove_dir_all(dir);
}
/// Wait until the claim deadline of the claimed row `id` on `queue`
/// has lapsed — the bounded-wait posture (no tight-race assertions).
async fn until_deadline_lapsed(queue: &dyn alkstore::Queue, id: i64) {
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(6);
loop {
let lapsed = match queue.get_job(id).await {
Ok(Some(job)) => {
job.claim_expires_at.is_some()
&& job.claim_expires_at.unwrap()
< std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as i64
}
_ => false,
};
if lapsed {
return;
}
assert!(
tokio::time::Instant::now() < deadline,
"deadline: never lapsed"
);
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
}
/// Acceptance: the uniform validity predicate — once the claim
/// deadline has lapsed, **all** handle ops refuse with `Ok(false)`
/// (heartbeat, ack, retry, fail — never errors), the row is untouched,
/// and a reclaim consumes an attempt (ADR-010 §2). Healthy renewal is
/// pinned first: `heartbeat` is the repeatable `&self` case (ADR-019
/// §3), an absolute reset; the one-shot ops consume their handles
/// whatever the row's answer (ADR-019 §3's annotation).
#[tokio::test(flavor = "multi_thread")]
async fn lapsed_deadline_refuses_all_handle_ops_and_reclaim_consumes_an_attempt() {
let dir = temp_dir("validity");
let path = temp_path("validity");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
// Healthy phase (visibility 300): renewal lands and is repeatable.
let queue = store.queue("healthy", opts()).await.unwrap();
let id = queue
.enqueue(json!({"work": true}), EnqueueOpts::default())
.await
.unwrap();
let handle = queue.claim_one("w1").await.unwrap().unwrap();
let first_deadline = handle.job().claim_expires_at.unwrap();
assert!(handle.heartbeat(5000).await.unwrap(), "renewal lands");
assert!(handle.heartbeat(5000).await.unwrap(), "repeatable");
let renewed = queue.get_job(id).await.unwrap().unwrap();
assert!(
renewed.claim_expires_at.unwrap() >= first_deadline + 4500,
"renewal is an absolute reset from now: {} vs {first_deadline}",
renewed.claim_expires_at.unwrap()
);
assert!(
renewed.claimed_at == handle.job().claimed_at,
"heartbeat never touches claimed_at"
);
assert!(handle.ack().await.unwrap(), "healthy ack lands");
// Lapsed phase (visibility 1): every op refuses with Ok(false).
let queue = store
.queue(
"lapsed",
QueueOpts {
visibility_timeout_s: 1,
max_attempts: 5,
backoff_base_s: 5,
dead_letter_retention_s: None,
},
)
.await
.unwrap();
let mut ids = Vec::new();
for _ in 0..4 {
ids.push(
queue
.enqueue(json!({"work": true}), EnqueueOpts::default())
.await
.unwrap(),
);
}
let mut claims: Vec<Box<dyn alkstore::JobHandle>> = Vec::new();
for id in &ids {
let h = queue
.claim_batch("w1", 1)
.await
.unwrap()
.into_iter()
.next()
.expect("claimable");
assert_eq!(h.job().id, *id);
claims.push(h);
}
for id in &ids {
until_deadline_lapsed(queue.as_ref(), *id).await;
}
// The lapsed handles refuse with Ok(false) — op by op.
assert!(
!claims[0].heartbeat(5000).await.unwrap(),
"late heartbeat refuses"
);
assert!(!claims[0].heartbeat(5000).await.unwrap(), "still refuses");
assert!(!claims.remove(1).ack().await.unwrap(), "lapsed ack refuses");
assert!(
!claims
.remove(1)
.retry(Some("boom".to_string()), Some(10))
.await
.unwrap(),
"lapsed retry refuses"
);
assert!(
!claims
.remove(1)
.fail(Some("boom".to_string()))
.await
.unwrap(),
"lapsed fail refuses"
);
// The rows are untouched by the refusals.
for id in &ids {
let job = queue.get_job(*id).await.unwrap().unwrap();
assert_eq!(job.state, JobState::Processing);
assert_eq!(job.attempts, 1);
}
// A reclaim wins and consumes an attempt.
let reclaimed = queue.claim_one("w2").await.unwrap().expect("reclaimable");
assert_eq!(reclaimed.job().attempts, 2, "the reclaim counts an attempt");
assert_eq!(reclaimed.job().worker_id.as_deref(), Some("w2"));
// The original worker's side effects may complete, but its op
// cannot land now (reclaim won atomically): post-reclaim, the
// original claim's ops refuse.
assert!(
!claims[0].heartbeat(5000).await.unwrap(),
"the original claim's ops refuse post-reclaim"
);
store.close();
let _ = std::fs::remove_dir_all(dir);
}
/// Acceptance: dead-letter — `retry` at exhausted budget moves the row
/// to dead with the pinned default string (`"max attempts exceeded"`),
/// `fail(None)` records `"failed"`, `get_job` sees dead rows with
/// `last_error`/`died_at` carrying the full stamp list (ADR-010 §1/§4;
/// ADR-019 §3's annotation).
#[tokio::test(flavor = "multi_thread")]
async fn exhaustion_and_fail_dead_letter_with_the_pinned_strings() {
let dir = temp_dir("dead-letter");
let path = temp_path("dead-letter");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let queue = store
.queue(
"dying",
QueueOpts {
visibility_timeout_s: 300,
max_attempts: 2,
backoff_base_s: 5,
dead_letter_retention_s: None,
},
)
.await
.unwrap();
// Exhaustion via retry: claim (attempts=1), retry under budget →
// pending, reclaim (attempts=2 == max) → the retry dead-letters.
let id1 = queue
.enqueue(json!({"n": 1}), EnqueueOpts::default())
.await
.unwrap();
let h = queue.claim_one("w1").await.unwrap().unwrap();
assert!(
h.retry(None, Some(0)).await.unwrap(),
"under budget: reschedules"
);
let job = queue.get_job(id1).await.unwrap().unwrap();
assert_eq!(job.state, JobState::Pending);
assert_eq!(job.attempts, 1);
let h = queue.claim_one("w1").await.unwrap().unwrap();
assert_eq!(h.job().attempts, 2, "a reclaim counts too");
assert!(
h.retry(Some("caller string".to_string()), None)
.await
.unwrap(),
"the exhausted-budget retry lands as the dead-letter move"
);
let dead = queue
.get_job(id1)
.await
.unwrap()
.expect("dead rows visible");
assert_eq!(dead.state, JobState::Dead);
assert_eq!(
dead.last_error.as_deref(),
Some("caller string"),
"retry(Some(err)) carries the caller's string on the dead row"
);
assert!(dead.died_at.is_some());
assert_eq!(dead.max_attempts, 2, "the dead row carries its stamps");
assert_eq!(dead.visibility_timeout_s, 300);
assert_eq!(dead.backoff_base_s, 5);
assert_eq!(dead.attempts, 2);
assert!(
dead.claimed_at.is_some(),
"the dead row keeps claim history"
);
assert_eq!(dead.worker_id, None);
// retry(None) exhaustion pins the default string.
let id2 = queue
.enqueue(json!({"n": 2}), EnqueueOpts::default())
.await
.unwrap();
let h = queue.claim_one("w1").await.unwrap().unwrap();
assert_eq!(h.job().attempts, 1, "a fresh claim is one attempt");
assert!(
h.retry(None, Some(0)).await.unwrap(),
"attempts 1 < 2: reschedules"
);
let h = queue.claim_one("w1").await.unwrap().unwrap();
assert_eq!(h.job().attempts, 2);
assert!(
h.retry(None, None).await.unwrap(),
"the at-budget retry dead-letters"
);
let dead = queue.get_job(id2).await.unwrap().unwrap();
assert_eq!(dead.state, JobState::Dead);
assert_eq!(
dead.last_error.as_deref(),
Some(EXHAUSTED_ERROR),
"retry(None) exhaustion records the pinned default string"
);
assert!(dead.died_at.is_some());
store.close();
let _ = std::fs::remove_dir_all(dir);
// fail(None) pins "failed"; fail(Some(e)) carries the caller's.
let dir2 = temp_dir("dead-letter-2");
let path2 = temp_path("dead-letter-2");
let store = open_store(path2.to_str().unwrap(), Default::default()).unwrap();
let queue = store.queue("dying", opts()).await.unwrap();
let id = queue
.enqueue(json!({"n": 3}), EnqueueOpts::default())
.await
.unwrap();
let h = queue.claim_one("w1").await.unwrap().unwrap();
assert!(h.fail(None).await.unwrap(), "fail lands");
let dead = queue.get_job(id).await.unwrap().unwrap();
assert_eq!(dead.state, JobState::Dead);
assert_eq!(
dead.last_error.as_deref(),
Some(FAILED_ERROR),
"fail(None) records the pinned default string"
);
let id = queue
.enqueue(json!({"n": 4}), EnqueueOpts::default())
.await
.unwrap();
let h = queue.claim_one("w1").await.unwrap().unwrap();
assert!(h.fail(Some("explicit reason".to_string())).await.unwrap());
let dead = queue.get_job(id).await.unwrap().unwrap();
assert_eq!(dead.last_error.as_deref(), Some("explicit reason"));
// The pinning constants are the contract's defaults.
assert_eq!(EXHAUSTED_ERROR, "max attempts exceeded");
assert_eq!(FAILED_ERROR, "failed");
assert_eq!(EXPIRED_ERROR, "expired");
store.close();
let _ = std::fs::remove_dir_all(dir2);
}
/// Acceptance: `cancel` is unconditional (either state, any holder),
/// `ack_batch` applies the per-id validity predicate without a worker
/// filter (partial success ordinary), and missing ids are silent
/// (ADR-010 §1; ADR-019 §1).
#[tokio::test(flavor = "multi_thread")]
async fn cancel_and_ack_batch_behaviors() {
let dir = temp_dir("cancel-ackbatch");
let path = temp_path("cancel-ackbatch");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let queue = store.queue("ops", opts()).await.unwrap();
// cancel: pending and processing, holder-independent. The first
// claim takes the head of the FIFO order (equal run_at, id ASC →
// the first-enqueued row) — observe, don't assume.
let pending = queue
.enqueue(json!({"p": 1}), EnqueueOpts::default())
.await
.unwrap();
let processing = queue
.enqueue(json!({"p": 2}), EnqueueOpts::default())
.await
.unwrap();
let h = queue.claim_one("w1").await.unwrap().unwrap();
let claimed_first = h.job().id;
assert!(claimed_first == pending || claimed_first == processing);
let still_live = if claimed_first == pending {
processing
} else {
pending
};
assert!(queue.cancel(claimed_first).await.unwrap());
assert!(
queue.cancel(still_live).await.unwrap(),
"cancel is unconditional"
);
assert!(!queue.cancel(still_live).await.unwrap(), "gone = false");
assert!(queue.get_job(still_live).await.unwrap().is_none());
assert!(!h.ack().await.unwrap(), "the holder's next op refuses");
// ack_batch: per-id predicate, no worker filter, count returned.
let a = queue
.enqueue(json!({"a": 1}), EnqueueOpts::default())
.await
.unwrap();
let b = queue
.enqueue(json!({"b": 1}), EnqueueOpts::default())
.await
.unwrap();
let c = queue
.enqueue(json!({"c": 1}), EnqueueOpts::default())
.await
.unwrap();
let ha = queue.claim_one("w1").await.unwrap().unwrap();
let hb = queue.claim_one("w2").await.unwrap().unwrap();
assert_eq!(ha.job().id, a, "FIFO: the first-enqueued row claims first");
assert_eq!(hb.job().id, b);
assert!(queue.cancel(c).await.unwrap());
// a is w1's claim, b is w2's — both count (no worker filter on the
// batch form); the cancelled c and the missing id do not.
let count = queue.ack_batch(&[a, b, c, 999_999]).await.unwrap();
assert_eq!(count, 2, "only the valid ids count; partial is ordinary");
assert!(queue.get_job(a).await.unwrap().is_none());
assert!(queue.get_job(b).await.unwrap().is_none());
// An empty batch is vacuously zero.
assert_eq!(queue.ack_batch(&[]).await.unwrap(), 0);
store.close();
let _ = std::fs::remove_dir_all(dir);
}
/// Acceptance: `sweep_expired` moves every past-`expires_at` row (both
/// states — the no-stranded-rows property, ADR-010 §5) to dead with
/// `"expired"`, and returns the moved count; unexpired rows are
/// untouched; no leader lock required (the bare call just runs).
#[tokio::test(flavor = "multi_thread")]
async fn sweep_expired_moves_both_states_and_applies_retention() {
let dir = temp_dir("sweep");
let path = temp_path("sweep");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let queue = store.queue("sweeping", opts()).await.unwrap();
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as i64;
// Expired pending + expired processing + unexpired: only the two
// expired rows move; the count is 2. `expires: Some(0)` resolves
// to expires_at = enqueue-instant — second-precision, so the sweep
// predicate (`<= unixepoch()`) sees it as lapsed from the next
// second on (the wait below makes that definite). The wait also
// holds off claims so the expired-pending row cannot be claimed
// first: only the not-yet-expired rows are claimable at that
// point.
let expired_pending = queue
.enqueue(
json!({"k": 1}),
EnqueueOpts {
expires: Some(3600),
..Default::default()
},
)
.await
.unwrap();
let expired_processing = queue
.enqueue(
json!({"k": 2}),
EnqueueOpts {
expires: Some(1),
..Default::default()
},
)
.await
.unwrap();
let unexpired = queue
.enqueue(
json!({"k": 3}),
EnqueueOpts {
expires: Some(3600),
..Default::default()
},
)
.await
.unwrap();
// Claim the head (the expired-pending row, FIFO) then back-date
// both rows' expiry directly — the honest cross-connection probe —
// so the sweep's both-states move is exercised deterministically.
let h = queue.claim_one("w1").await.unwrap().unwrap();
assert_eq!(h.job().id, expired_pending);
{
let conn = rusqlite::Connection::open(&path).unwrap();
for id in [expired_pending, expired_processing] {
conn.execute(
"UPDATE __alkstore_live SET expires_at = unixepoch() - 1 WHERE id = ?1",
[id],
)
.unwrap();
}
}
let moved = queue.sweep_expired().await.unwrap();
assert_eq!(moved, 2, "both expired rows moved (pending + processing)");
for id in [expired_pending, expired_processing] {
let dead = queue.get_job(id).await.unwrap().expect("moved to dead");
assert_eq!(dead.state, JobState::Dead);
assert_eq!(dead.last_error.as_deref(), Some(EXPIRED_ERROR));
assert!(dead.died_at.is_some() && dead.died_at.unwrap() >= now);
}
let live = queue.get_job(unexpired).await.unwrap().unwrap();
assert_eq!(live.state, JobState::Pending, "unexpired rows untouched");
// A bare sweep on an empty-expiry queue is 0 and not an error.
assert_eq!(queue.sweep_expired().await.unwrap(), 0);
store.close();
let _ = std::fs::remove_dir_all(dir);
}
/// Acceptance: retention — dead rows past their
/// `dead_letter_retention_s` stamp are deleted by the sweep (the only
/// sweeper dead rows ever have), rows with `None` retention live
/// forever; the deleted count reports through (ADR-010 §4/§5).
#[tokio::test(flavor = "multi_thread")]
async fn sweep_enforces_dead_letter_retention() {
let dir = temp_dir("retention");
let path = temp_path("retention");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let keeping = store
.queue(
"keeping",
QueueOpts {
dead_letter_retention_s: None,
..opts()
},
)
.await
.unwrap();
let expiring = store
.queue(
"expiring",
QueueOpts {
dead_letter_retention_s: Some(1),
..opts()
},
)
.await
.unwrap();
let keep_id = keeping
.enqueue(json!({"k": 1}), EnqueueOpts::default())
.await
.unwrap();
let expire_id = expiring
.enqueue(json!({"k": 2}), EnqueueOpts::default())
.await
.unwrap();
let h = keeping.claim_one("w1").await.unwrap().unwrap();
assert!(h.fail(None).await.unwrap());
let h = expiring.claim_one("w1").await.unwrap().unwrap();
assert!(h.fail(None).await.unwrap());
// Move the expired-retention row's death back past its TTL (the
// honest probe: retention reads died_at vs now).
{
let conn = rusqlite::Connection::open(&path).unwrap();
conn.execute(
"UPDATE __alkstore_dead SET died_at = unixepoch() - 10 WHERE id = ?1",
[expire_id],
)
.unwrap();
}
// The sweep on the expiring queue deletes the TTL'd row; the
// keeping queue's row (retention None) survives.
assert_eq!(expiring.sweep_expired().await.unwrap(), 1);
assert!(expiring.get_job(expire_id).await.unwrap().is_none());
assert_eq!(keeping.sweep_expired().await.unwrap(), 0);
assert!(
keeping.get_job(keep_id).await.unwrap().is_some(),
"retention None = dead rows live forever"
);
store.close();
let _ = std::fs::remove_dir_all(dir);
}
/// Acceptance: entry-point validation — the constructor rejects empty
/// and reserved names (`InvalidName`/`ReservedName`), claims reject
/// empty `worker_id`, and ops against a closed store fail closed with
/// the engine-wide `Database` shape (ADR-008 §4; ADR-019 §2).
#[tokio::test(flavor = "multi_thread")]
async fn validation_and_closed_store_posture() {
let dir = temp_dir("validation");
let path = temp_path("validation");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
match store.queue("", opts()).await {
Err(alkstore::Error::InvalidName { .. }) => {}
_other => panic!("empty queue name must be InvalidName"),
}
match store.queue(alkstore::RESERVED_PREFIX, opts()).await {
Err(alkstore::Error::ReservedName { .. }) => {}
_other => panic!("reserved queue name must be ReservedName"),
}
let queue = store.queue("valid", opts()).await.unwrap();
match queue.claim_batch("", 1).await {
Err(alkstore::Error::InvalidName { .. }) => {}
_other => panic!("empty worker_id must be InvalidName"),
}
store.close();
// Closed store: all queue ops fail closed with Database.
let err = store
.queue("valid", opts())
.await
.err()
.expect("closed store refuses construction");
assert!(matches!(err, alkstore::Error::Database(_)), "got {err:?}");
let err = queue
.enqueue(json!({}), EnqueueOpts::default())
.await
.expect_err("closed store refuses enqueue");
assert!(matches!(err, alkstore::Error::Database(_)), "got {err:?}");
let err = queue
.claim_batch("w1", 1)
.await
.err()
.expect("closed store refuses claim");
assert!(matches!(err, alkstore::Error::Database(_)), "got {err:?}");
let err = queue.get_job(1).await.expect_err("closed refuses get");
assert!(matches!(err, alkstore::Error::Database(_)), "got {err:?}");
let err = queue
.sweep_expired()
.await
.expect_err("closed refuses sweep");
assert!(matches!(err, alkstore::Error::Database(_)), "got {err:?}");
let _ = std::fs::remove_dir_all(dir);
}
+14 -87
View File
@@ -69,10 +69,13 @@ use alkstore::{
validate_shared_name,
};
use crate::seam::{blocking, sqlite_error, with_writer};
use crate::seam::{
ReaderLoopExit as LoopExit, blocking_read, closed_store_error, sqlite_error, with_reader,
with_writer,
};
use crate::substrate::{
Readers, SchemaError, SharedUpdateWatcher, Writer, stream_get_offset, stream_publish,
stream_read_since, stream_save_offset,
Readers, SharedUpdateWatcher, Writer, stream_get_offset, stream_publish, stream_read_since,
stream_save_offset,
};
/// How many events a wake-driven re-read fetches per page while
@@ -120,59 +123,6 @@ fn publish_with_key(
})
}
/// A reader-pool read through the `spawn_blocking` seam: acquire a
/// pooled reader, run the body, release on both arms (a failed SELECT
/// leaves the pooled connection healthy). A closed pool (store closed)
/// fails closed with `Database`.
async fn with_reader<T, F>(readers: Arc<Readers>, f: F) -> Result<T>
where
T: Send + 'static,
F: FnOnce(&rusqlite::Connection) -> Result<T> + Send + 'static,
{
blocking(move || {
let conn = acquire_reader(&readers)?;
let out = f(&conn);
readers.release(conn);
out
})
.await
}
fn acquire_reader(readers: &Arc<Readers>) -> Result<rusqlite::Connection> {
match readers.acquire() {
Ok(conn) => Ok(conn),
Err(SchemaError::Sqlite(e)) => {
if is_closed_err(&e) {
Err(closed_store_error())
} else {
Err(sqlite_error(e))
}
}
}
}
fn is_closed_err(e: &rusqlite::Error) -> bool {
matches!(
e,
rusqlite::Error::SqliteFailure(_, Some(msg)) if msg.contains("Database is closed")
)
}
fn closed_store_error() -> Error {
Error::database(std::io::Error::other("the store is closed"))
}
/// `recv()`'s `Err` arm carries `Database` only (ADR-021 §5): a
/// read error that mapped elsewhere (the decode-side `Codec` posture
/// guard) is remapped into the opaque fallback with the original error
/// preserved as the source.
fn database_only(e: Error) -> Error {
match e {
already @ Error::Database(_) => already,
other => Error::database(other),
}
}
/// One decoded page from the reader pool: `offset > from`, ASC, up to
/// `limit` (the extent guard is the caller's — this helper is
/// contract-blind to it).
@@ -509,37 +459,14 @@ fn run_subscribe_loop(
}
}
enum LoopExit {
Closed,
Failed(Error),
}
/// A synchronous pooled read used by the bridge thread. The closed
/// arms are distinguished: a closed store exits the loop, an open
/// store's failure is transient (`Failed` carries the mapped error).
fn blocking_read<T>(
readers: &Arc<Readers>,
store_closed: &Arc<AtomicBool>,
f: impl FnOnce(&rusqlite::Connection) -> Result<T>,
) -> std::result::Result<T, LoopExit> {
if store_closed.load(Ordering::Acquire) {
return Err(LoopExit::Closed);
}
let conn = match readers.acquire() {
Ok(conn) => conn,
Err(e) => match e {
SchemaError::Sqlite(inner) if is_closed_err(&inner) => return Err(LoopExit::Closed),
SchemaError::Sqlite(inner) => {
return Err(LoopExit::Failed(sqlite_error(inner)));
}
},
};
let out = f(&conn);
readers.release(conn);
match out {
Ok(value) => Ok(value),
Err(_e) if store_closed.load(Ordering::Acquire) => Err(LoopExit::Closed),
Err(e) => Err(LoopExit::Failed(e)),
/// `recv()`'s `Err` arm carries `Database` only (ADR-021 §5): a
/// read error that mapped elsewhere (the decode-side `Codec` posture
/// guard) is remapped into the opaque fallback with the original error
/// preserved as the source.
fn database_only(e: Error) -> Error {
match e {
already @ Error::Database(_) => already,
other => Error::database(other),
}
}
@@ -119,6 +119,7 @@ fixes. Cherry-picks: empty at scaffold, append-only forever.
| D-28 | hygiene | wave-2 review fix — the cross-mechanism pressure test's lock-churn branch exercised a per-producer lock name (literal `"p{producer}"` had no interpolation, so all producers contended on one name) with its result swallowed via `.unwrap_or(0)`; now interpolates, releases periodically, propagates errors | `ops.rs` (pressure test, lock branch) | test-quality only (no lineage counterpart, no library code touched); `fork-provenance-and-floor`'s ported-surface pressure intent | ours |
| D-29 | port delta | `SQLITE_OPEN_URI` dropped from `open_conn`'s flags (the watcher's own opens never carried it — the inherited posture was internally inconsistent); the path argument is a plain filesystem path, `?name=value` suffixes are literal filenames, no connection-semantics mutation outside the engine constructor's opts; pinned by test (`open_conn_treats_uri_shaped_path_as_plain_filename`) | `schema.rs` (`open_conn`) | ADR-023 §3 (plain-path open — the URI parameter surface is SQLite's, not this engine's; no consumer-inventory row names URI features) | ours |
| D-30 | port delta | `open_conn_bootstrapped` — new function stacking the engine's full per-connection open posture (pragmas, notify, `attach_alkstore_functions`, `bootstrap_schema`), and `Readers::acquire`'s open path switched from bare `open_conn` to it (fresh connections only — the pool never re-bootstraps a reused connection). The lineage's pools opened readers pragmas+notify only, the function set + schema living only on the writer; this engine's reader-slot operations (claim/ack, stream reads, lock ops, scheduler checks, leadership probes) run this substrate's SQL functions on pooled connections, so every pooled connection carries the full surface | `schema.rs` (`open_conn_bootstrapped`, `Readers::acquire` open path) | engine-sqlite.md's connection architecture ("each connection runs the substrate's bootstrap at open — the writer and each pooled reader") — engine-layer obligation of `sqlite-engine-open-opts`; landed with the engine's `open` wiring | ours |
| D-31 | port delta | `ack_batch(ids_json)` drops the `worker_id` filter — the D-12 predicate keeps `state = 'processing' AND claim_expires_at >= unixepoch()` per id, but the batch form loses its lineage-shaped `worker_id = ?2` conjunct: the contract's `Queue::ack_batch(ids)` carries no worker identity (ADR-019 §1's surface — the batch ack is the batch form of *ack*, its outcomes are per-id independent, and a queue-scoped handle does not re-bind a claimant at call time). The single-row `ack(job_id, worker_id)` keeps its worker conjunct unchanged | `queue_ops.rs` (`ack_batch` + its substrate test's call shape) | ADR-019 §1 (the `ack_batch(ids) -> count` surface, worker-less); ADR-010 §1 (batch form of ack, per-id predicate); ADR-012 §4 (delta discipline — D-12's re-derivation adjusted) | ours |
### Cherry-picks
+4 -8
View File
@@ -123,19 +123,15 @@ pub(crate) fn ack(conn: &Connection, job_id: i64, worker_id: &str) -> rusqlite::
Ok(deleted as i64)
}
pub(crate) fn ack_batch(
conn: &Connection,
ids_json: &str,
worker_id: &str,
) -> rusqlite::Result<i64> {
pub(crate) fn ack_batch(conn: &Connection, ids_json: &str) -> rusqlite::Result<i64> {
let mut stmt = conn.prepare_cached(
"DELETE FROM __alkstore_live
WHERE id IN (SELECT value FROM json_each(?1))
AND state = 'processing' AND worker_id = ?2
AND state = 'processing'
AND claim_expires_at >= unixepoch()
RETURNING id",
)?;
let mut rows = stmt.query(rusqlite::params![ids_json, worker_id])?;
let mut rows = stmt.query(rusqlite::params![ids_json])?;
let mut count = 0;
while rows.next()?.is_some() {
count += 1;
@@ -1224,7 +1220,7 @@ mod queue_tests {
let canceled = claimed(&conn, "emails", "w1", 5);
cancel(&conn, canceled).unwrap();
let ids = format!("[{a},{b},{canceled},999999]");
assert_eq!(ack_batch(&conn, &ids, "w1").unwrap(), 2);
assert_eq!(ack_batch(&conn, &ids).unwrap(), 2);
assert_eq!(live_count(&conn, a), 0);
assert_eq!(live_count(&conn, b), 0);
}
+4 -79
View File
@@ -24,17 +24,18 @@
use serde_json::Value;
use std::sync::Arc;
use alkstore::{BoxedFuture, EnqueueOpts, Error, Job, JobState, Result, StreamEvent, TxHandle};
use alkstore::{BoxedFuture, EnqueueOpts, Error, Job, Result, StreamEvent, TxHandle};
use alkstore::{validate_local_name, validate_shared_name};
use crate::queue::job_from_json;
use crate::resolution::{
outbox_backing_queue_name, outbox_default_stamps, plain_queue_default_stamps,
resolve_enqueue_opts,
resolve_enqueue_opts, stamps_with_override,
};
use crate::seam::{blocking, sqlite_error};
use crate::stream::stream_events_from_json_page;
use crate::substrate::{
Stamps, Writer, enqueue, get_job, stream_get_offset, stream_publish, stream_read_since,
Writer, enqueue, get_job, stream_get_offset, stream_publish, stream_read_since,
stream_save_offset,
};
@@ -164,18 +165,6 @@ impl Drop for ConnLease {
}
}
/// Stamp resolution for the tx path over the handle-less default
/// shapes: `EnqueueOpts::max_attempts` is the one per-job override,
/// resolved over the base stamp set (ADR-010 §3a; the other stamps
/// take the base untouched — the derived defaults carry
/// visibility/backoff/retention).
fn stamps_with_override(base: Stamps, opts: &EnqueueOpts) -> Stamps {
Stamps {
max_attempts: opts.max_attempts.unwrap_or(base.max_attempts),
..base
}
}
impl TxHandle for SqliteTxHandle {
fn enqueue_tx<'a>(
&'a mut self,
@@ -428,67 +417,3 @@ impl Drop for SqliteTxHandle {
}
}
}
/// Decode the substrate's job-row JSON into the contract value
/// (`Job::from_row` — the engine-side constructor, ADR-021 §2's
/// `claimed_at` surfaced).
fn job_from_json(value: &serde_json::Value) -> Result<Option<Job>> {
let state = match value["state"].as_str() {
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: Vec<u8> = value["payload"]
.as_str()
.ok_or_else(|| Error::Codec("job row payload must be a string".to_string()))?
.as_bytes()
.to_vec();
let field = |name: &str| -> Result<i64> {
serde_json::from_value(value[name].clone())
.map_err(|e| Error::Codec(format!("job row {name} must be i64: {e}")))
};
let opt_field = |name: &str| -> Result<Option<i64>> {
if value[name].is_null() {
Ok(None)
} else {
serde_json::from_value(value[name].clone())
.map_err(|e| Error::Codec(format!("job row {name} must be i64: {e}")))
}
};
let opt_str = |name: &str| -> Result<Option<String>> {
if value[name].is_null() {
Ok(None)
} else {
serde_json::from_value(value[name].clone())
.map_err(|e| Error::Codec(format!("job row {name} must be a string: {e}")))
}
};
Ok(Some(Job::from_row(
field("id")?,
value["queue"]
.as_str()
.ok_or_else(|| Error::Codec("job row queue must be a string".to_string()))?
.to_string(),
state,
payload,
field("priority")?,
field("run_at")?,
field("attempts")?,
field("max_attempts")?,
opt_str("worker_id")?,
opt_field("claimed_at")?,
opt_field("claim_expires_at")?,
field("created_at")?,
opt_field("expires_at")?,
field("visibility_timeout_s")?,
field("backoff_base_s")?,
opt_field("dead_letter_retention_s")?,
opt_str("last_error")?,
opt_field("died_at")?,
)))
}