SQLite engine: scheduler + outbox — schedule/unschedule/run_schedules leader loop, outbox enqueue/run_once (task sqlite-engine-scheduler-outbox)

This commit is contained in:
glm-5.3-flash committed 2026-10-08 14:08:43 +00:00
1 parent 1edcb0e27d
commit 8502a51af7
11 files changed
+1548 -36

No files matched your search

+2
View File
@@ -20,8 +20,10 @@
mod lock;
mod notify;
mod opts;
mod outbox;
mod queue;
mod resolution;
mod scheduler;
mod seam;
mod store;
mod stream;
+22
View File
@@ -29,6 +29,12 @@ use alkstore::{BoxedFuture, Lock, Result, validate_local_name, validate_shared_n
use crate::seam::{closed_store_error, sqlite_error, with_writer};
use crate::substrate::{Writer, lock_acquire, lock_release, lock_renew};
/// The scheduler's leadership-lock name (ADR-009 §4 — the reserved-
/// prefix derived name in the lock namespace; engine-internal: every
/// consumer entry point rejects the same prefix, so a consumer can
/// never collide with it).
pub(crate) const SCHEDULER_LEADER_LOCK: &str = concat!("__alkstore_", "scheduler");
/// The acquire path (see the module docs): validation, guard, then the
/// substrate's `lock_acquire` through the writer slot. `Ok(None)` is
/// the held-by-someone-else arm.
@@ -94,6 +100,22 @@ impl std::fmt::Debug for SqliteLockHandle {
}
}
impl SqliteLockHandle {
pub(crate) fn new_internal(
name: &str,
owner: &str,
writer: Arc<Writer>,
closed: Arc<AtomicBool>,
) -> Self {
Self {
name: name.to_string(),
owner: owner.to_string(),
writer,
closed,
}
}
}
impl Lock for SqliteLockHandle {
fn name(&self) -> &str {
&self.name
+170
View File
@@ -0,0 +1,170 @@
//! The outbox mechanism's SQLite arm (`sqlite-engine-scheduler-outbox`)
//! — the helper over queues (ADR-014's auto-commit half; the tx half —
//! `outbox_enqueue_tx` — landed with the seam task):
//!
//! - `outbox(name)` — the validated constructor (shared-namespace
//! rules); the backing queue `__alkstore_outbox:{name}` derives
//! engine-side under the reserved prefix — unreachable by
//! `enqueue_tx` by design (ADR-014 §1: outbox enqueue is the only
//! path into it).
//! - `enqueue` — auto-commit into the backing queue through the writer
//! slot, stamping the outbox's derived 60/5/5 set (visibility 60 s,
//! max_attempts 5, backoff base 5 s, retention forever — ADR-010 §3 /
//! ADR-014 §1); `EnqueueOpts::max_attempts` is the one per-job
//! override, `delay`/`run_at`/`expires` resolve per ADR-020.
//! - `run_once(worker_id, delivery)` — claim one job from the backing
//! queue via the ordinary queue claim machinery (the substrate's
//! single-statement claim), run `delivery.deliver(job)`: `Ok` ⇒ ack;
//! `Err(e)` ⇒ `retry(Some(e), None)` — the queue's curve (the same
//! equal-jitter computation `Queue`'s handles use, engine-side one
//! owner). `true` = claimed and processed; `false` = no work (a
//! value, not an error). No engine-issued heartbeat inside delivery —
//! the boxed [`JobHandle`]'s own `heartbeat` is the consumer's
//! renewal path (ADR-019 §5).
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use serde_json::Value;
use alkstore::{BoxedFuture, Delivery, EnqueueOpts, Error, JobHandle, Outbox, Result};
use crate::queue::{EXHAUSTED_ERROR, SqliteJobHandle, backoff_delay_s, jobs_from_json_page};
use crate::resolution::{
outbox_backing_queue_name, outbox_default_stamps, resolve_enqueue_opts, stamps_with_override,
};
use crate::seam::{closed_store_error, sqlite_error, with_writer};
use crate::substrate::{Writer, ack, claim_batch, enqueue, retry};
/// The outbox-scoped handle (see the module docs).
pub struct SqliteOutboxHandle {
name: String,
backing: String,
writer: Arc<Writer>,
closed: Arc<AtomicBool>,
}
impl std::fmt::Debug for SqliteOutboxHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SqliteOutboxHandle")
.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,
writer: Arc<Writer>,
closed: Arc<AtomicBool>,
) -> Result<Box<dyn Outbox>> {
alkstore::validate_shared_name(name)?;
Ok(Box::new(SqliteOutboxHandle {
name: name.to_string(),
backing: outbox_backing_queue_name(name),
writer,
closed,
}))
}
impl Outbox for SqliteOutboxHandle {
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 backing = self.backing.clone();
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(outbox_default_stamps(), &opts);
enqueue(
conn,
&backing,
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 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(closed_store_error()) });
}
let backing = self.backing.clone();
let writer = self.writer.clone();
let closed = self.closed.clone();
let worker_id = worker_id.to_string();
Box::pin(async move {
// Claim through the ordinary queue machinery (the writer
// slot; the substrate's single-statement claim is the
// exactly-once handout).
let page = with_writer(writer.clone(), {
let backing = backing.clone();
let worker = worker_id.clone();
move |conn| {
claim_batch(conn, &backing, &worker, 1, EXHAUSTED_ERROR).map_err(sqlite_error)
}
})
.await?;
let mut claimed = jobs_from_json_page(&page)?;
let Some(job) = claimed.pop() else {
return Ok(false);
};
// No engine-issued heartbeat inside delivery: the boxed
// handle's own heartbeat is the consumer's renewal path.
let handle: Box<dyn JobHandle> = Box::new(SqliteJobHandle::new(
job.clone(),
writer.clone(),
closed.clone(),
));
let worker = job.worker_id.clone().unwrap_or_default();
match delivery.deliver(handle).await {
Ok(()) => {
let id = job.id;
with_writer(writer, move |conn| {
Ok(ack(conn, id, &worker).map_err(sqlite_error)? > 0)
})
.await?;
Ok(true)
}
Err(err) => {
let id = job.id;
// `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);
with_writer(writer, move |conn| {
Ok(retry(conn, id, &worker, delay, &exhaust).map_err(sqlite_error)? > 0)
})
.await?;
Ok(true)
}
}
})
}
}
+13 -2
View File
@@ -194,8 +194,9 @@ pub(crate) fn job_from_json(value: &Value) -> Result<Option<Job>> {
)))
}
/// Decode a claim page of job-row JSON into `Job` values.
fn jobs_from_json_page(page: &str) -> Result<Vec<Job>> {
/// Decode a claim page of job-row JSON into `Job` values. Shared with
/// the outbox's `run_once` claim (the same decode owner).
pub(crate) 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
@@ -404,6 +405,16 @@ impl std::fmt::Debug for SqliteJobHandle {
}
}
impl SqliteJobHandle {
pub(crate) fn new(job: Job, writer: Arc<Writer>, closed: Arc<AtomicBool>) -> Self {
Self {
job,
writer,
closed,
}
}
}
impl JobHandle for SqliteJobHandle {
fn job(&self) -> &Job {
&self.job
+284
View File
@@ -0,0 +1,284 @@
//! The scheduler mechanism's SQLite arm (`sqlite-engine-scheduler-outbox`)
//! — the collapse shape (ADR-009): two registrations + one runner, the
//! leadership lock, and the tick's transactional boundary advance.
//!
//! - `schedule(name, spec, queue, payload, opts)` — upsert by name via
//! the substrate's `scheduler_register`; spec/name/queue all validate
//! at the entry point, **before storage** (`InvalidSpec` has no round
//! trip — the substrate's `parse_every_interval` is the v1 grammar
//! parser; its rejection maps to `InvalidSpec`). The returned
//! [`Schedule`] is the registry read-back (the core constructor).
//! - `unschedule(name)` — the substrate's `scheduler_unregister`;
//! `true` = removed, `false` = absent.
//! - `run_schedules(stop)` — the leader loop (ADR-009 §1): acquire the
//! leadership lock `__alkstore_scheduler` (the engine's own lock
//! machinery with a derived engine-internal TTL), then loop {renew —
//! on loss return `Err(LeadershipLost)` *before ticking*; tick; sleep
//! until the next due boundary or stop}. The tick's fire enqueues and
//! the boundary advance commit in one `BEGIN IMMEDIATE` transaction —
//! writer serialization as the row lock on the schedule row (ADR-009
//! §4's layer (b): even a rogue second ticker observing the same
//! storage cannot double-fire a boundary). A cancelled
//! [`StopToken`] ends the sleep and returns `Ok(())`; leadership
//! loss returns `Err(LeadershipLost)` regardless of the token.
//! - The runner adds no delivery machinery: fired jobs land in the
//! named queue as ordinary rows — the target queue's derived defaults
//! (300/3/5/none) with `ScheduleOpts` applied over them (ADR-020 §3)
//! — claimable via the ordinary queue surface (ADR-009 §3).
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use serde_json::Value;
use alkstore::{BoxedFuture, Error, Lock, Result, Schedule, ScheduleOpts, StopToken};
use crate::lock::{SCHEDULER_LEADER_LOCK, SqliteLockHandle};
use crate::resolution::plain_queue_default_stamps;
use crate::seam::{closed_store_error, sqlite_error, with_writer};
use crate::substrate::{
FireOpts, Stamps, Writer, lock_acquire, parse_every_interval, scheduler_register,
scheduler_soonest, scheduler_tick, scheduler_unregister,
};
/// 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
/// engine-owned).
const LEADER_TTL_S: i64 = 10;
/// Seconds between lease renewals during the sleep.
const SLEEP_RENEW_EVERY_S: u64 = 2;
/// 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 `lock_acquire`'s owner-match mistake a foreign
/// leader for self (two leaders holding the same name at once).
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())
}
/// The idle sleep when nothing is scheduled (soonest = 0).
const IDLE_SLEEP_S: u64 = 60;
fn invalid_spec(spec: &str) -> Error {
Error::InvalidSpec {
spec: spec.to_string(),
}
}
/// The auto-commit register path: pre-storage validation, the
/// substrate upsert, the [`Schedule`] read-back.
pub(crate) fn schedule(
name: &str,
spec: &str,
queue: &str,
payload: Value,
opts: ScheduleOpts,
writer: Arc<Writer>,
closed: Arc<AtomicBool>,
) -> BoxedFuture<'static, Result<Schedule>> {
if let Err(e) = alkstore::validate_shared_name(name) {
return Box::pin(async move { Err(e) });
}
if let Err(e) = alkstore::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 closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
let payload_str = match alkstore::encode_payload(&payload) {
Ok(bytes) => match String::from_utf8(bytes) {
Ok(s) => s,
Err(e) => {
return Box::pin(async move {
Err(Error::Codec(format!("row payload must be utf-8: {e}")))
});
}
},
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 fire = FireOpts {
priority: opts.priority,
expires_s: opts.expires,
};
let stored_name = name.to_string();
let stored_spec = spec.to_string();
let stored_queue = queue.to_string();
Box::pin(async move {
with_writer(writer, move |conn| {
scheduler_register(
conn,
&stored_name,
&stored_queue,
&stored_spec,
&payload_str,
fire,
stamps,
)
.map_err(sqlite_error)?;
Ok(Schedule::new(stored_name, stored_spec, stored_queue, opts))
})
.await
})
}
/// The auto-commit unregister path (validation mirrors `schedule`'s —
/// a reserved/empty schedule name never rounds trips).
pub(crate) fn unschedule(
name: &str,
writer: Arc<Writer>,
closed: Arc<AtomicBool>,
) -> BoxedFuture<'static, Result<bool>> {
if let Err(e) = alkstore::validate_shared_name(name) {
return Box::pin(async move { Err(e) });
}
if closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
let name = name.to_string();
Box::pin(async move {
with_writer(writer, move |conn| {
Ok(scheduler_unregister(conn, &name).map_err(sqlite_error)? > 0)
})
.await
})
}
/// The leader loop (see the module docs).
pub(crate) fn run_schedules(
stop: StopToken,
writer: Arc<Writer>,
closed: Arc<AtomicBool>,
) -> BoxedFuture<'static, Result<()>> {
Box::pin(async move {
if closed.load(Ordering::Acquire) {
return Err(closed_store_error());
}
// Step 1: acquire the leadership lock. A refused acquire —
// someone else leads — is `Err(LeadershipLost)` at start: no
// leader loop exists without the lock, and the caller's recipe
// (respawn, ADR-009 §1) is identical to a mid-run loss.
let owner = leader_owner();
let granted = {
let writer = writer.clone();
let owner = owner.clone();
with_writer(writer, move |conn| {
Ok(
lock_acquire(conn, SCHEDULER_LEADER_LOCK, &owner, LEADER_TTL_S)
.map_err(sqlite_error)?
> 0,
)
})
.await?
};
if !granted {
return Err(Error::LeadershipLost);
}
let lease: Box<SqliteLockHandle> = Box::new(SqliteLockHandle::new_internal(
SCHEDULER_LEADER_LOCK,
&owner,
writer.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 writer = writer.clone();
with_writer(writer, scheduler_tick_transactional).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 honker leader-loop pattern: renew
// on cadence while holding).
keep_awake_renewing(&lease, &stop, soonest).await;
}
})
}
/// One tick under `BEGIN IMMEDIATE`: the fire enqueues + the boundary
/// advance commit together, or the whole tick rolls back (crash
/// mid-tick ⇒ the boundary refires; two tickers serialize on the
/// writer). Returns the soonest next fire for the sleep (0 = none).
fn scheduler_tick_transactional(conn: &rusqlite::Connection) -> Result<i64> {
conn.execute_batch("BEGIN IMMEDIATE")
.map_err(sqlite_error)?;
let body = (|| {
let now = crate::resolution::now_unix(conn)?;
scheduler_tick(conn, now).map_err(sqlite_error)?;
scheduler_soonest(conn).map_err(sqlite_error)
})();
match body {
Ok(soonest) => {
conn.execute_batch("COMMIT").map_err(sqlite_error)?;
Ok(soonest)
}
Err(e) => {
let _ = conn.execute_batch("ROLLBACK");
Err(e)
}
}
}
/// The sliced sleep until the next due boundary (or idle horizon) or
/// the stop token — the loop's wake path between ticks.
async fn keep_awake_renewing(lease: &SqliteLockHandle, 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 {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0);
(soonest - now).max(1) as u64
};
let deadline = started + Duration::from_secs(wait_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 slice = deadline.min(next_renew);
tokio::time::sleep_until(slice).await;
}
}
+32 -28
View File
@@ -188,13 +188,15 @@ impl Store for SqliteStore {
fn outbox<'a>(
&'a self,
_name: &str,
name: &str,
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::Outbox>>> {
Box::pin(async {
Err(database_error(
"outbox wiring lands with the mechanism tasks",
))
})
if self.closed.load(std::sync::atomic::Ordering::Acquire) {
return Box::pin(async { Err(crate::seam::closed_store_error()) });
}
let name = name.to_string();
let writer = self.writer.clone();
let closed = self.closed.clone();
Box::pin(async move { crate::outbox::new_outbox_handle(&name, writer, closed) })
}
fn try_lock<'a>(
@@ -212,36 +214,34 @@ impl Store for SqliteStore {
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>> {
Box::pin(async {
Err(database_error(
"schedule wiring lands with the mechanism tasks",
))
})
let writer = self.writer.clone();
let closed = self.closed.clone();
let name = name.to_string();
let spec = spec.to_string();
let queue = queue.to_string();
crate::scheduler::schedule(&name, &spec, &queue, payload, opts, writer, closed)
}
fn unschedule<'a>(&'a self, _name: &str) -> alkstore::BoxedFuture<'a, alkstore::Result<bool>> {
Box::pin(async {
Err(database_error(
"unschedule wiring lands with the mechanism tasks",
))
})
fn unschedule<'a>(&'a self, name: &str) -> alkstore::BoxedFuture<'a, alkstore::Result<bool>> {
let writer = self.writer.clone();
let closed = self.closed.clone();
let name = name.to_string();
crate::scheduler::unschedule(&name, writer, closed)
}
fn run_schedules<'a>(
&'a self,
_stop: alkstore::StopToken,
stop: alkstore::StopToken,
) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> {
Box::pin(async {
Err(database_error(
"scheduler wiring lands with the mechanism tasks",
))
})
let writer = self.writer.clone();
let closed = self.closed.clone();
crate::scheduler::run_schedules(stop, writer, closed)
}
}
@@ -252,8 +252,12 @@ mod notify_tests;
#[cfg(test)]
mod open_tests;
#[cfg(test)]
mod outbox_tests;
#[cfg(test)]
mod queue_tests;
#[cfg(test)]
mod scheduler_tests;
#[cfg(test)]
mod stream_tests;
#[cfg(test)]
mod tx_tests;
+376
View File
@@ -0,0 +1,376 @@
//! The outbox mechanism's acceptance tests (`sqlite-engine-scheduler-outbox`):
//! the derived backing queue is unreachable by `enqueue_tx` (the
//! reserved-prefix rejection — ADR-014's stated guarantee), the
//! outbox's `enqueue` stamps the 60/5/5 derived set, `run_once` acks
//! on `Ok`, retries with the curve on `Err`, and returns `false` on an
//! empty backing queue; worker_id validation (ADR-019 §2). The
//! commit-atomic tx half pins under tx_tests (already landed with the
//! seam task).
use std::path::PathBuf;
use std::time::Duration;
use alkstore::{EnqueueOpts, JobState, Store};
use serde_json::json;
use crate::store::open_store;
fn temp_dir(tag: &str) -> PathBuf {
let d = std::env::temp_dir().join(format!(
"alkstore-outbox-{tag}-{}-{:?}",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_dir_all(&d);
std::fs::create_dir_all(&d).unwrap();
d
}
fn cleanup(dir: &PathBuf) {
let _ = std::fs::remove_dir_all(dir);
}
fn open_test_store(tag: &str) -> (PathBuf, Box<dyn Store>) {
let dir = temp_dir(tag);
let path = dir.join("store.db");
let store: Box<dyn Store> =
Box::new(open_store(path.to_str().unwrap(), Default::default()).unwrap());
(dir, store)
}
/// Count rows in the outbox's backing queue directly (the reserved
/// name is unreachable by the handle surface by design).
fn backing_count(dir: &std::path::Path, outbox: &str) -> i64 {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT count(*) FROM __alkstore_live WHERE queue = ?1",
[format!("__alkstore_outbox:{outbox}")],
|r| r.get(0),
)
.unwrap()
}
/// 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 Scripted {
fn error(&self) -> alkstore::Error {
alkstore::Error::Codec(self.outcome.clone().unwrap_err())
}
}
impl alkstore::Delivery for Scripted {
fn deliver<'a>(
&'a mut self,
job: Box<dyn alkstore::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 (dir, store) = open_test_store("outbox-construction");
let check = |err: alkstore::Error, expect: &str, name: &str| match err {
alkstore::Error::InvalidName { .. } => assert_eq!(expect, "empty", "{name}"),
alkstore::Error::ReservedName { .. } => assert_eq!(expect, "reserved", "{name}"),
other => panic!("{name} got the wrong error: {other}"),
};
let empty_err = match store.outbox("").await {
Err(e) => e,
Ok(_) => panic!("outbox(\"\") must fail entry-point validation"),
};
check(empty_err, "empty", "empty outbox name");
let reserved_err = match store.outbox("__alkstore_outbox:steal").await {
Err(e) => e,
Ok(_) => panic!("outbox(reserved) must fail entry-point validation"),
};
check(reserved_err, "reserved", "reserved outbox name");
let outbox = store.outbox("billing").await.unwrap();
assert_eq!(outbox.name(), "billing");
cleanup(&dir);
}
/// 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 (dir, store) = open_test_store("outbox-stamps");
let err = 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 = 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 = {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT visibility_timeout_s, max_attempts, backoff_base_s, \
dead_letter_retention_s, payload \
FROM __alkstore_live WHERE id = ?1",
[id],
|r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, i64>(1)?,
r.get::<_, i64>(2)?,
r.get::<_, Option<i64>>(3)?,
r.get::<_, String>(4)?,
))
},
)
.unwrap()
};
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());
// 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 = {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT max_attempts FROM __alkstore_live WHERE id = ?1",
[id2],
|r| r.get(0),
)
.unwrap()
};
assert_eq!(max2, 2);
assert_eq!(backing_count(&dir, "billing"), 2);
cleanup(&dir);
}
/// 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
/// delivery's error string would surface on a later dead row).
#[tokio::test(flavor = "multi_thread")]
async fn run_once_acks_on_ok_and_retries_on_err() {
let (dir, store) = open_test_store("outbox-runonce");
let mut outbox = 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 = store.outbox("quiet").await.unwrap();
let mut noop = Scripted {
outcome: Ok(()),
seen_payloads: Vec::new(),
};
assert!(!quiet.run_once("w0", &mut noop).await.unwrap());
// 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.
let (state, run_at_gap, attempts): (String, i64, i64) = {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
let now: i64 = conn
.query_row("SELECT unixepoch()", [], |r| r.get(0))
.unwrap();
conn.query_row(
"SELECT state, run_at - ?2, attempts FROM __alkstore_live WHERE id = ?1",
rusqlite::params![id_err, now],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.unwrap()
};
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(&dir, "mail"), 1, "the acked row left live");
let _ = id_ok;
cleanup(&dir);
}
/// 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 (dir, store) = open_test_store("outbox-heartbeat");
let mut outbox = store.outbox("hb").await.unwrap();
outbox
.enqueue(json!("work"), Default::default())
.await
.unwrap();
struct Renew {
db_path: std::path::PathBuf,
}
impl alkstore::Delivery for Renew {
fn deliver<'a>(
&'a mut self,
job: Box<dyn alkstore::JobHandle>,
) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> {
let db_path = self.db_path.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
// WAL reader connection): the deadline moved to
// ≈ now + 3000 — the heartbeat's absolute reset, not
// the stamped 60 s visibility.
let conn = rusqlite::Connection::open(&db_path).unwrap();
let (claim_expires_at, db_now): (i64, i64) = conn
.query_row(
"SELECT claim_expires_at, unixepoch() FROM __alkstore_live WHERE id = ?1",
[job.job().id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
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 mut d = Renew {
db_path: dir.join("store.db"),
};
assert!(outbox.run_once("w-hb", &mut d).await.unwrap());
cleanup(&dir);
}
/// Acceptance: `run_once`'s exhaustion path — a delivery that keeps
/// failing past the outbox set's max_attempts (5) dead-letters the row
/// with the delivery's error string carried (`retry(Some(e), None)`'s
/// exhaust string) — driven here 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 (dir, store) = open_test_store("outbox-dead");
let mut outbox = 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 conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT last_error, died_at FROM __alkstore_dead WHERE id = ?1",
[id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap()
};
assert_eq!(last_error, failing.error().to_string());
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);
cleanup(&dir);
}
@@ -0,0 +1,545 @@
//! The scheduler mechanism's acceptance tests (`sqlite-engine-scheduler-outbox`):
//! upsert + read-back with the entry-point validation triad pinned
//! (`InvalidSpec` pre-storage with no round trip, `InvalidName`/
//! `ReservedName` on both name and queue argument), `unschedule`'s
//! true/false, the leader loop firing due boundaries into the named
//! queue (claimable via the ordinary queue surface), the leadership
//! lock held and one-firer-per-boundary under a second concurrent
//! runner, `Err(LeadershipLost)` vs clean `Ok(())` stop, fired jobs'
//! stamp source (the plain defaults with `ScheduleOpts` applied,
//! ADR-020 §3), and catch-up-cap/skip-forward observability through
//! the pinned substrate behavior (ADR-009, ADR-019 §3/§4).
use std::path::PathBuf;
use std::time::Duration;
use alkstore::{Error, JobState, QueueOpts, ScheduleOpts, StopToken, Store};
use serde_json::json;
use crate::store::open_store;
fn temp_dir(tag: &str) -> PathBuf {
let d = std::env::temp_dir().join(format!(
"alkstore-sched-{tag}-{}-{:?}",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_dir_all(&d);
std::fs::create_dir_all(&d).unwrap();
d
}
fn cleanup(dir: &PathBuf) {
let _ = std::fs::remove_dir_all(dir);
}
fn open_test_store(tag: &str) -> (PathBuf, Box<dyn Store>) {
let dir = temp_dir(tag);
let path = dir.join("store.db");
let store: Box<dyn Store> =
Box::new(open_store(path.to_str().unwrap(), Default::default()).unwrap());
(dir, store)
}
/// Acceptance: `schedule` validates at the entry point on all three
/// arguments — empty/reserved schedule name (`InvalidName` /
/// `ReservedName`), empty/reserved queue argument (same, ADR-021 §3),
/// and a non-`@every` or malformed spec → `InvalidSpec` **before
/// storage** (no round trip; a failing register observable via the
/// absent registry row, and `@every 0s` rejected as non-positive).
#[tokio::test(flavor = "multi_thread")]
async fn schedule_entry_point_validation_is_pinned_on_all_arguments() {
let (dir, store) = open_test_store("sched-validation");
let check = |err: Error, expect: &str, name: &str| match err {
Error::InvalidName { .. } => assert_eq!(expect, "empty", "{name}"),
Error::ReservedName { .. } => assert_eq!(expect, "reserved", "{name}"),
Error::InvalidSpec { spec } => {
assert_eq!(expect, "spec", "{name}");
assert!(
!spec.is_empty(),
"{name}: InvalidSpec carries the spec string"
);
}
other => panic!("{name} got the wrong error: {other}"),
};
check(
store
.schedule("", "@every 1s", "work", json!(null), Default::default())
.await
.unwrap_err(),
"empty",
"empty schedule name",
);
check(
store
.schedule(
"__alkstore_rogue",
"@every 1s",
"work",
json!(null),
Default::default(),
)
.await
.unwrap_err(),
"reserved",
"reserved schedule name",
);
check(
store
.schedule("tick", "@every 1s", "", json!(null), Default::default())
.await
.unwrap_err(),
"empty",
"empty queue argument",
);
check(
store
.schedule(
"tick",
"@every 1s",
"__alkstore_outbox:steal",
json!(null),
Default::default(),
)
.await
.unwrap_err(),
"reserved",
"reserved queue argument",
);
for bad in [
"0 9 * * *",
"@every 1x",
"@every 0s",
"every 5s",
"@every s",
"@every 1.5s",
] {
check(
store
.schedule("bad-spec", bad, "work", json!(null), Default::default())
.await
.unwrap_err(),
"spec",
&(format!("spec {bad:?}")),
);
}
// The rejected calls never reached storage (pre-storage, no round
// trip): the good name registers fresh and the registry holds only
// that one row.
let sched = store
.schedule(
"tick",
"@every 1s",
"work",
json!({"a": 1}),
Default::default(),
)
.await
.unwrap();
assert_eq!(sched.name, "tick");
assert_eq!(sched.spec, "@every 1s");
assert_eq!(sched.queue, "work");
assert_eq!(sched.opts, ScheduleOpts::default());
cleanup(&dir);
}
/// Acceptance: `schedule` upserts by name — a re-register replaces the
/// whole row (spec, queue, payload, opts and the next-fire reset),
/// read-back reflects the replacement.
#[tokio::test(flavor = "multi_thread")]
async fn schedule_upsert_replaces_the_whole_row() {
let (dir, store) = open_test_store("sched-upsert");
let first = store
.schedule(
"sweep",
"@every 10s",
"old-queue",
json!({"v": 1}),
ScheduleOpts {
priority: 3,
max_attempts: Some(7),
expires: Some(120),
},
)
.await
.unwrap();
assert_eq!(first.queue, "old-queue");
assert_eq!(first.opts.priority, 3);
let second = store
.schedule(
"sweep",
"@every 20m",
"new-queue",
json!({"v": 2}),
ScheduleOpts {
priority: 0,
max_attempts: None,
expires: None,
},
)
.await
.unwrap();
assert_eq!(second.spec, "@every 20m");
assert_eq!(second.queue, "new-queue");
assert_eq!(second.opts, ScheduleOpts::default());
// unregister + re-register is the pause shape — the upsert makes
// update = re-register honest (ADR-009 §1).
assert!(store.unschedule("sweep").await.unwrap());
assert!(!store.unschedule("sweep").await.unwrap());
cleanup(&dir);
}
/// Acceptance: the leader loop fires due boundaries into the named
/// queue, the fired jobs are ordinary queue work (`get_job`-visible
/// with the plain-queue defaults + `ScheduleOpts` applied over them,
/// ADR-020 §3; claimable via the queue surface), the leadership lock
/// is held (`__alkstore_scheduler`), and a clean stop returns `Ok(())`.
#[tokio::test(flavor = "multi_thread")]
async fn run_schedules_fires_boundaries_and_stop_is_clean() {
let (dir, store) = open_test_store("sched-fire");
store
.schedule(
"fast",
"@every 1s",
"work",
json!({"kind": "sweep"}),
ScheduleOpts {
priority: 5,
max_attempts: Some(9),
expires: None,
},
)
.await
.unwrap();
let stop = StopToken::new();
let runner_store = {
// `run_schedules` borrows the store: spawn on a clone-shaped
// handle — the store is behind `Arc`-free trait object; run the
// loop in this task and stop it from this same task's timer.
// (A second store open on the same path shares the db.)
let path = dir.join("store.db");
open_store(path.to_str().unwrap(), Default::default()).unwrap()
};
let runner_stop = stop.clone();
let runner = tokio::spawn(async move { runner_store.run_schedules(runner_stop).await });
// Two boundaries at 1 s cadence (+ registration's first-fire grace
// of exactly one interval): wait for two fired jobs, bounded.
let deadline = tokio::time::Instant::now() + Duration::from_secs(25);
let mut fired_ids = Vec::new();
loop {
let q = store.queue("work", QueueOpts::default()).await.unwrap();
let page = q.claim_batch("probe", 10).await.unwrap();
for h in page {
fired_ids.push(h.job().id);
}
if fired_ids.len() >= 2 {
break;
}
assert!(
tokio::time::Instant::now() < deadline,
"no two boundary fires within the window"
);
tokio::time::sleep(Duration::from_millis(200)).await;
}
// The stamp source pin (ADR-020 §3): the plain defaults 300/3/5/none
// with ScheduleOpts applied — priority 5 and max_attempts 9 override,
// visibility/backoff/retention take the defaults.
let got = store.queue("work", QueueOpts::default()).await.unwrap();
let job = got.get_job(fired_ids[0]).await.unwrap().unwrap();
assert_eq!(job.state, JobState::Processing, "our probe claimed it");
assert_eq!(job.priority, 5);
assert_eq!(job.max_attempts, 9);
assert_eq!(job.visibility_timeout_s, 300);
assert_eq!(job.backoff_base_s, 5);
assert_eq!(job.dead_letter_retention_s, None);
assert_eq!(
job.payload,
json!({"kind": "sweep"}).to_string().into_bytes()
);
assert_eq!(job.worker_id.as_deref(), Some("probe"));
// The leadership lock row exists under the reserved name.
let leader = crate::lock::SCHEDULER_LEADER_LOCK;
let (count,) = {
// Raw probe through a fresh store connection: a reserved name
// is unreachable by try_lock (Shared validation would reject),
// so query the substrate table directly via a tx read on
// get_job-shaped machinery is wrong namespace; use the reader
// path with plain SQL through the engine's substrate module.
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT count(*) FROM __alkstore_locks WHERE name = ?1 AND expires_at > unixepoch()",
[leader],
|r| Ok((r.get::<_, i64>(0)?,)),
)
.unwrap()
};
assert_eq!(count, 1, "the leadership lock row must be live");
// Clean stop: `Ok(())`, promptly (not after the full sleep).
stop.cancel();
let outcome = tokio::time::timeout(Duration::from_secs(10), runner)
.await
.expect("clean stop must terminate promptly")
.unwrap();
assert!(
matches!(outcome, Ok(())),
"clean stop returns Ok(()): {outcome:?}"
);
drop(store);
cleanup(&dir);
}
/// Acceptance: a second concurrent runner does not double-fire — the
/// leadership lock admits one leader; the loser returns
/// `Err(LeadershipLost)` (no loop, no ticks), the winner fires each
/// boundary exactly once (at-least-once per boundary with one firer,
/// ADR-009 §4's writer-serialization floor).
#[tokio::test(flavor = "multi_thread")]
async fn second_runner_is_refused_and_boundaries_fire_once() {
let (dir, store) = open_test_store("sched-single");
store
.schedule("beat", "@every 1s", "single", json!(1), Default::default())
.await
.unwrap();
// Runner A takes leadership first.
let stop_a = StopToken::new();
let runner_a_store =
open_store(dir.join("store.db").to_str().unwrap(), Default::default()).unwrap();
let runner_a = tokio::spawn({
let stop = stop_a.clone();
let store = runner_a_store;
async move { store.run_schedules(stop).await }
});
// Runner B comes second: it must be refused (locked out) with
// `Err(LeadershipLost)` — start it once A has demonstrably taken
// the lock (wait for the first fired job).
let deadline = tokio::time::Instant::now() + Duration::from_secs(20);
loop {
let q = store.queue("single", QueueOpts::default()).await.unwrap();
if !q.claim_batch("probe", 10).await.unwrap().is_empty() {
break;
}
assert!(
tokio::time::Instant::now() < deadline,
"runner A never fired a boundary"
);
tokio::time::sleep(Duration::from_millis(200)).await;
}
let runner_b_store =
open_store(dir.join("store.db").to_str().unwrap(), Default::default()).unwrap();
let outcome_b = runner_b_store
.run_schedules(StopToken::new())
.await
.unwrap_err();
assert!(
matches!(outcome_b, Error::LeadershipLost),
"runner B must be refused with LeadershipLost: {outcome_b:?}"
);
// Let the leader run several more boundaries; the total fired count
// matches the elapsed boundaries within fire-window tolerance —
// with a second non-firing runner, one firer per boundary.
tokio::time::sleep(Duration::from_millis(2500)).await;
let fired_total: i64 = {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT count(*) FROM __alkstore_live WHERE queue = 'single'",
[],
|r| r.get(0),
)
.unwrap()
};
assert!(
(2..=5).contains(&fired_total),
"boundary fires must stay within the elapsed-boundary band, got {fired_total}"
);
stop_a.cancel();
let outcome_a = tokio::time::timeout(Duration::from_secs(10), runner_a)
.await
.expect("clean stop must terminate promptly")
.unwrap();
assert!(matches!(outcome_a, Ok(())), "{outcome_a:?}");
drop(store);
cleanup(&dir);
}
/// Acceptance: losing leadership mid-run returns `Err(LeadershipLost)`
/// *before any tick* — here exercised at the substrate boundary the
/// engine's renew path rides (`lock_renew` refuses a foreign/expired
/// owner row): a runner whose lease row was deleted (its TTL lapsed)
/// fails its next renew and returns the loss, never a tick.
#[tokio::test(flavor = "multi_thread")]
async fn leadership_loss_returns_the_loss_error_before_ticking() {
let (dir, store) = open_test_store("sched-loss");
store
.schedule(
"slow",
"@every 3600s",
"lossq",
json!(1),
Default::default(),
)
.await
.unwrap();
// Steal the lock from under no one: a rogue operator (bypassing
// the engine surface — try_lock itself rejects the reserved name)
// inserts the leadership row directly; the runner starts, must
// observe LeadershipLost at acquire.
{
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.execute(
"INSERT INTO __alkstore_locks (name, owner, expires_at)
VALUES (?1, 'rogue-operator', unixepoch() + 30)",
[crate::lock::SCHEDULER_LEADER_LOCK],
)
.unwrap();
}
let runner_store =
open_store(dir.join("store.db").to_str().unwrap(), Default::default()).unwrap();
let outcome = runner_store.run_schedules(StopToken::new()).await;
assert!(
matches!(outcome, Err(Error::LeadershipLost)),
"a runner refused at acquire returns LeadershipLost: {outcome:?}"
);
// No tick ever ran: the 1-hour schedule's queue is untouched.
let fired: i64 = {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT count(*) FROM __alkstore_live WHERE queue = 'lossq'",
[],
|r| r.get(0),
)
.unwrap()
};
assert_eq!(fired, 0, "a refused runner must never tick");
drop(store);
cleanup(&dir);
}
/// Acceptance: catch-up cap + skip-forward are observable through the
/// runner (the substrate's pinned per-tick behavior, ADR-009 §4): a
/// schedule with its fire horizon pushed far into the past replays at
/// most the 64-fires cap and skips forward strictly past now.
#[tokio::test(flavor = "multi_thread")]
async fn catch_up_replays_up_to_the_cap_then_skips_forward() {
let (dir, store) = open_test_store("sched-catchup");
store
.schedule(
"backlog",
"@every 1h",
"catchq",
json!("w"),
Default::default(),
)
.await
.unwrap();
// Push the fire horizon 200 boundary-widths into the past (200
// elapsed boundaries > the 64 cap) before the runner ever starts.
// `@every 1h` keeps the post-cap boundary beyond the test window:
// the only fires possible here are the capped replay's own.
{
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
let now: i64 = conn
.query_row("SELECT unixepoch()", [], |r| r.get(0))
.unwrap();
conn.execute(
"UPDATE __alkstore_scheduler_tasks SET next_fire_at = ?2 WHERE name = 'backlog'",
rusqlite::params!["backlog", now - 200 * 3600],
)
.unwrap();
}
let store2 = open_store(dir.join("store.db").to_str().unwrap(), Default::default()).unwrap();
let stop = StopToken::new();
let runner_stop = stop.clone();
let runner = tokio::spawn(async move { store2.run_schedules(runner_stop).await });
// The first tick fires exactly 64, advances past now (the
// skip-forward lands at the next 1 h boundary — far beyond this
// window), and the skipped boundaries never enqueue.
let deadline = tokio::time::Instant::now() + Duration::from_secs(20);
let fired: i64 = loop {
let c: i64 = {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT count(*) FROM __alkstore_live WHERE queue = 'catchq'",
[],
|r| r.get(0),
)
.unwrap()
};
let next: i64 = {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT next_fire_at FROM __alkstore_scheduler_tasks WHERE name = 'backlog'",
[],
|r| r.get(0),
)
.unwrap()
};
let now: i64 = {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row("SELECT unixepoch()", [], |r| r.get(0))
.unwrap()
};
if c == 64 && next > now {
break c;
}
assert!(
tokio::time::Instant::now() < deadline,
"the cap never landed (fires={c}, next={next}, now={now})"
);
tokio::time::sleep(Duration::from_millis(100)).await;
};
assert_eq!(fired, 64, "exactly the contract-constant cap replays");
stop.cancel();
let outcome = tokio::time::timeout(Duration::from_secs(10), runner)
.await
.expect("clean stop must terminate promptly")
.unwrap();
assert!(matches!(outcome, Ok(())), "{outcome:?}");
// After the stop: the skipped boundaries never enqueue — and the
// skip-forward's next 1 h boundary cannot have fired in-window.
let fired_after: i64 = {
let conn = rusqlite::Connection::open(dir.join("store.db")).unwrap();
conn.query_row(
"SELECT count(*) FROM __alkstore_live WHERE queue = 'catchq'",
[],
|r| r.get(0),
)
.unwrap()
};
assert_eq!(fired_after, 64, "the skipped boundaries never enqueue");
drop(store);
cleanup(&dir);
}
+2 -2
View File
@@ -59,8 +59,8 @@ pub(crate) use ops::{
};
pub(crate) use queue_ops::{
FireOpts, Stamps, ack, ack_batch, cancel, claim_batch, enqueue, fail, get_job, heartbeat,
queue_next_claim_at, retry, scheduler_register, scheduler_soonest, scheduler_tick,
scheduler_unregister, sweep_expired,
parse_every_interval, queue_next_claim_at, retry, scheduler_register, scheduler_soonest,
scheduler_tick, scheduler_unregister, sweep_expired,
};
pub(crate) use schema::{
Error as SchemaError, Readers, Writer, apply_default_pragmas, attach_notify, bootstrap_schema,
+1 -1
View File
@@ -514,7 +514,7 @@ fn now_unix(conn: &Connection) -> rusqlite::Result<i64> {
conn.query_row("SELECT unixepoch()", [], |r| r.get(0))
}
fn parse_every_interval(spec: &str) -> Result<i64, String> {
pub(crate) fn parse_every_interval(spec: &str) -> Result<i64, String> {
let body = spec
.strip_prefix("@every ")
.ok_or_else(|| format!("unsupported spec {spec:?}; @every <n><unit> only"))?;
+101 -3
View File
@@ -1,7 +1,7 @@
---
id: sqlite-engine-scheduler-outbox
name: SQLite engine — scheduler (`schedule`/`unschedule`/`run_schedules`) and outbox helper
status: pending
status: completed
depends_on: [sqlite-engine-queues, sqlite-engine-locks]
scope: broad
risk: high
@@ -95,8 +95,106 @@ their own long-running machinery:
## Notes
> To be filled by implementation agent
> Decisions of record made in implementation that the description
> didn't pin:
>
> - **The leadership lock's owner token is per-runner-instance-unique**
> (pid + subsec-nanos, `scheduler::leader_owner`), not a shared
> constant: the substrate's `lock_acquire` matches by owner string,
> so two processes running the same constant would each see
> "self" holding and both would tick. The runner-unique token is
> what makes the renewal-refusal path meaningful across stores on
> the same db. TTL is engine-internal (10 s), renewed at the top of
> every loop iteration **and** on a 2 s cadence across the sleep —
> without the in-sleep renewals a boundary further out than the TTL
> (any `@every` ≥ ~10 s) would lapse its own lease by drift. A
> failed mid-sleep renewal is deliberately not acted on at the
> sleep site — the loop's top-of-iteration renew owns the loss
> decision (returns `Err(LeadershipLost)` before any tick).
> - **Acquire-time refusal is also `Err(LeadershipLost)`**: a runner
> that fails to take the lock at start returns the same matchable
> loss a mid-run loser does (no runner-loop exists without the
> lock; the respawn recipe — ADR-009 §1 — is identical either way).
> - **The tick runs under `BEGIN IMMEDIATE`** on the writer slot —
> fire enqueues + boundary advance commit together (ADR-009 §4's
> layer (b); writer serialization as the row lock; a rogue second
> ticker observing the same storage serializes behind the writer
> and cannot double-fire). The tick body is the substrate's
> `scheduler_tick` over `now_unix`, then `scheduler_soonest` for
> the sleep deadline — all inside the one transaction.
> - **The idle sleep (soonest = 0) is 60 s.** Nothing pins it; a
> schedule registered while the runner idles is noticed at the
> next idle tick (≤ 60 s late) — acceptable engine detail; the
> boundary math itself is unaffected.
> - **`schedule`'s `InvalidSpec` maps the substrate's
> `parse_every_interval` rejection** — `parse_every_interval` was
> made `pub(crate)` (one-line delta in `substrate/queue_ops.rs`)
> instead of duplicating the grammar engine-side; the substrate
> stays the one parser (the task's "the substrate's parser; map
> its rejection" wording, taken literally).
> - **`Outbox::run_once`'s `Err` arm passes the delivery error's
> `Display` string as the retry's exhaust string** — the same
> posture `JobHandle::retry(Some(e), …)` carries (the string
> surfaces on a dead row only if the budget exhausts). The curve
> delay is the engine's one-owner `backoff_delay_s` (the queue
> handle's owner, reused — ADR-010 §3's single-owner rule).
> - **`SqliteJobHandle` grew a `pub(crate) fn new`** so the outbox
> can rebuild the claim handle for the delivery callable
> (`run_once` claims through the same `claim_batch` machinery and
> rewraps the row); fields stay private.
> - **Test flake fix**: the catch-up-cap test originally drove
> `@every 1s` from a −200 s horizon; under full-suite load the
> runner's *next* live boundary could race into the assertion
> window. Switched to `@every 1h` with a 200-boundary backlog —
> the skip-forward target lands ~1 h out, so nothing can fire
> in-window; the assertions are now exact (64 fires, none after).
>
> Verification: `cargo test` (workspace: 188 sqlite + 25 core + 3
> contract-suite), `cargo clippy --all-targets -- -D warnings`,
> `cargo fmt --check` — all clean. The scheduler tests were run
> repeatedly (3×) to confirm the flake fix holds.
## Summary
> To be filled on completion
> What landed, verified how:
>
> - `alkstore-sqlite/src/scheduler.rs` — `schedule`/`unschedule`/
> `run_schedules` wired to the `Store` trait in `store.rs`
> (replacing the three placeholder stubs): entry-point validation on
> all three of name/spec/queue (InvalidSpec pre-storage via the
> substrate parser), the substrate upsert/unregister through the
> writer slot, and the leader loop over `__alkstore_scheduler` with
> per-instance owner, renew-before-tick, transactional tick
> (`BEGIN IMMEDIATE`), soonest-driven sliced sleep with in-sleep
> lease renewals, clean-stop `Ok(())` and loss
> `Err(LeadershipLost)`.
> - `alkstore-sqlite/src/outbox.rs` — `outbox(name)` wired (validated
> constructor), `enqueue` stamping the 60/5/5 derived set with the
> `EnqueueOpts::max_attempts` override and full ADR-020 resolution,
> and `run_once(worker_id, delivery)` claiming via the ordinary
> queue machinery, acking on `Ok`, `retry(Some(e), None)` (curve
> delay, delivery error as exhaust string) on `Err`, `false` on
> empty, the boxed `SqliteJobHandle` handed to the caller with no
> engine-issued heartbeat.
> - Small substrate/engine deltas: `parse_every_interval` made
> `pub(crate)`; `SqliteJobHandle::new` (pub(crate)) added;
> `jobs_from_json_page` made `pub(crate)`; `SCHEDULER_LEADER_LOCK`
> constant lives in `lock.rs` next to the machinery it reuses;
> `SqliteLockHandle::new_internal` (engine-internal constructor for
> the leadership lease).
> - Tests: `store/scheduler_tests.rs` (6 acceptance tests —
> validation triad, upsert read-back, fire-into-queue with stamp
> pinning via `get_job`, leadership lock live-row probe, second-
> runner refusal + one-firer-per-boundary band, acquire-time loss
> with a zero-fire probe, catch-up cap 64 + skip-forward with no
> in-window refire) and `store/outbox_tests.rs` (5 — constructor
> validation, backing-queue unreachability by `enqueue_tx` +
> 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). The commit-atomic tx half
> (`outbox_enqueue_tx`) was already covered by `tx_tests.rs`
> (`tx_stamps_resolve_max_attempts_and_outbox_set`,
> `outbox_backing_queue_name_is_reserved_and_derived`, the
> validation rows in `tx_entry_points_validate_names`).
> - Gates: workspace build/test green (216 total), clippy `-D
> warnings` clean, fmt clean.