pg: scheduler runner resilience — quarantine bad rows, retry transient ticks

Review 002 Finding 5 (+ Finding 6's scheduler doc bullet):

- tick: a due row whose stored spec fails @every re-parse is
  quarantined, not fatal — logged with the schedule name, boundary
  advanced strictly past now (skip-forward at min interval), tick tx
  still commits, remaining due rows proceed
- run_schedules: a pool/database tick failure retries 3x with a short
  doubling backoff (250ms -> 1s cap) before the loop exits Err; the
  top-of-iteration renew keeps owning the lease-loss decision
- exiting errors carry schedule-name/tick-phase context in the source
  chain (ErrorContext wrapper; Database's Display is opaque)
- module docs: the Err-exit TTL-lapse posture stated; the false
  per-slice soonest re-read claim corrected to the honest slice/idle
  posture (60s idle floor)

Tests: behavioral quarantine pin (tampered via direct SQL UPDATE;
runner survives to clean Ok(()), other schedule fires, bad row never
fires, boundary advanced) and a server-less retry-policy pin.
Verified: cargo test -p alkstore-postgres green against the harness
server (121+10+9), build/clippy -D warnings/fmt green server-less.
This commit is contained in:
glm-5.3-flash committed 2026-10-10 04:58:43 +00:00
1 parent ace94f0e13
commit 12d0497b1c
3 files changed
+419 -30

No files matched your search

+185 -19
View File
@@ -43,14 +43,39 @@
//! with the `ScheduleOpts` resolved over the target queue's derived
//! engine defaults; the 64-boundary catch-up cap with skip-forward
//! past the cap (honker parity — the contract constant, not a knob);
//! the soonest-boundary read for the sleep deadline. Wake-driven
//! tick advance: the runner parks on the sleep deadline but the
//! loop's stop-select shortens every sleep slice to ≤ 1 s and both
//! the tick and the soonest are re-read each iteration boundary —
//! a newly-registered schedule is noticed by the re-read no later
//! than the next slice (the deadline always fires; the wake never
//! gates correctness — the SQLite runner's 60 s idle-tick posture
//! is the acceptable floor if the slice cadence were fussy).
//! the soonest-boundary read for the sleep deadline. A stored spec
//! that fails re-parse on a due row (`@every` grammar re-parsed from
//! the stored text each tick; specs are pre-validated at
//! `schedule()`, so an unparseable row is tampered or foreign) is
//! **quarantined, not fatal**: logged (the engine's diagnostic
//! posture) with the schedule name, its boundary advanced past now
//! by the skip-forward helper (the row cannot spin), and the
//! remaining due rows proceed — the tick tx still commits (the
//! quarantine advance is an ordinary write). The quarantined row re-
//! enters the due set on every subsequent tick (advanced strictly
//! past now each time), so it never fires and never fails the
//! runner while its stored spec is broken. Wake-driven tick
//! advance: the runner parks on the sleep deadline but the loop's
//! stop-select shortens every sleep slice to ≤ 1 s; the slices renew
//! the lease and check the stop token only — the soonest is *not*
//! re-read mid-sleep. A newly-registered schedule is noticed when
//! the runner's current soonest deadline arrives; an idle runner
//! (soonest = 0) notices up to `IDLE_SLEEP_S` (60 s) late — the
//! deadline always fires and the wake never gates correctness, so
//! the 60 s idle-tick floor is the accepted posture.
//! - Transient-error retry: a tick failed by a pool/database error
//! retries `TICK_RETRIES` (3) times with a short doubling backoff
//! (250 ms → 1 s cap) before the loop exits `Err` — a single hiccup
//! no longer kills the runner (the consumer's respawn recipe remains
//! the last resort). The retry never masks a lost lease: the loop's
//! top-of-iteration renew owns the loss decision, and a lease that
//! lapses mid-retries is caught by that renew before the next tick.
//! The exiting error's source chain carries the schedule-name /
//! tick-phase context (Database's Display is opaque; the chain is
//! the carrier). On the `Err` exit paths (renewal loss, tick retries
//! exhausted) the leadership lock row is **not** released — its TTL
//! lapse (10 s) returns the name to acquire candidates; that
//! documented TTL-lapse posture is the contract.
//! - Clean stop (`StopToken::cancel`) → `Ok(())`, the leadership lock
//! released cleanly; loss → `Err(LeadershipLost)` before any tick.
//! - `outbox(name)` — the validated constructor (empty →
@@ -114,6 +139,50 @@ const SLEEP_SLICE_S: u64 = 1;
/// The idle sleep when nothing is scheduled (soonest = 0).
const IDLE_SLEEP_S: u64 = 60;
/// The tick's transient-error retry budget: a tick failed by a
/// pool/database error (pool checkout, `BEGIN`, a statement, `COMMIT`)
/// retries this many times before the loop exits `Err` (the module
/// docs' retry posture — a single hiccup no longer kills the runner;
/// leadership loss stays the top-of-iteration renew's decision).
pub(crate) const TICK_RETRIES: u64 = 3;
/// The short backoff between the tick's retries: doubling from 250 ms,
/// capped at 1 s (retries 1..=3 wait 250 ms, 500 ms, 1 s; beyond the
/// budget the cap holds). Total added stop latency across the full
/// retry window is ≤ 1.75 s.
pub(crate) fn tick_retry_backoff(retry: u64) -> Duration {
let exp = (retry.max(1) - 1).min(2);
Duration::from_millis(250 << exp)
}
/// The tick-site context wrapper (the task's error-context shape):
/// the exiting error's source chain carries the schedule-name /
/// tick-phase context — Database's own Display is opaque ("database
/// error"), so the message rides the source chain (`Database` →
/// `io::Error` carrying this message → the original error → the
/// driver chain).
#[derive(Debug)]
struct ErrorContext {
message: String,
source: Error,
}
impl std::fmt::Display for ErrorContext {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.message)
}
}
impl std::error::Error for ErrorContext {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
Some(&self.source)
}
}
fn context_error(message: String, source: Error) -> Error {
Error::database(std::io::Error::other(ErrorContext { message, source }))
}
/// The leadership lock's reserved name (ADR-009 §4 — the lock
/// machinery's namespace, engine-derived; consumers can never collide
/// because every entry point rejects the reserved prefix).
@@ -400,7 +469,45 @@ pub(crate) async fn tick(client: &tokio_postgres::Client, schema: &str, now: i64
}
for row in &due {
let interval = parse_every_interval(&row.spec)?;
// Quarantine, not fatality (the module docs' posture): a
// stored spec that fails re-parse is a tampered or
// future-foreign row — log it with the schedule name, advance
// its boundary strictly past now (the skip-forward helper at
// the minimum interval), and continue with the remaining due
// rows. The advance is an ordinary write in the tick tx.
let interval = match parse_every_interval(&row.spec) {
Ok(interval) => interval,
Err(bad) => {
eprintln!(
"alkstore-postgres: schedule {name:?} spec {spec:?} failed to re-parse ({bad}) \
— quarantining: boundary advanced past now, schedule skipped until \
the stored spec is repaired or removed",
name = row.name,
spec = row.spec,
);
let name = row.name.clone();
let advanced = skip_forward(row.next_fire_at, now, 1);
client
.execute(
&format!(
"UPDATE {sched} SET next_fire_at = $2 WHERE name = $1",
sched = sched
),
&[&name, &advanced],
)
.await
.map_err(|e| {
context_error(
format!(
"scheduler tick boundary advance (spec quarantine) for \
schedule {name:?}"
),
pg_error(e),
)
})?;
continue;
}
};
let mut next_fire_at = row.next_fire_at;
let mut fires = 0i64;
@@ -426,7 +533,16 @@ pub(crate) async fn tick(client: &tokio_postgres::Client, schema: &str, now: i64
},
row.stamps,
)
.await?;
.await
.map_err(|e| {
context_error(
format!(
"scheduler tick fire (enqueue) for schedule {name:?}",
name = row.name
),
e,
)
})?;
fires += 1;
next_fire_at = next_fire_at.saturating_add(interval);
}
@@ -440,7 +556,15 @@ pub(crate) async fn tick(client: &tokio_postgres::Client, schema: &str, now: i64
&[&row.name, &next_fire_at],
)
.await
.map_err(pg_error)?;
.map_err(|e| {
context_error(
format!(
"scheduler tick boundary advance for schedule {name:?}",
name = row.name
),
pg_error(e),
)
})?;
}
// The soonest-boundary read for the sleep deadline (inside the
@@ -516,14 +640,7 @@ pub(crate) fn run_schedules(stop: StopToken, ctx: EngineCtx) -> BoxedFuture<'sta
// before ticking: never fire on a stolen lease.
return Err(Error::LeadershipLost);
}
let soonest = {
let conn = pool.get().await.map_err(pool_error)?;
in_tx(conn, |client| {
let schema = schema.clone();
Box::pin(async move { tick(client, &schema, now_unix()).await })
})
.await?
};
let soonest = tick_with_retries(&pool, &schema).await?;
if stop.is_cancelled() {
let _ = lease.release().await;
return Ok(());
@@ -539,6 +656,55 @@ pub(crate) fn run_schedules(stop: StopToken, ctx: EngineCtx) -> BoxedFuture<'sta
})
}
/// One tick attempt over a fresh pool checkout (the loop body's
/// original shape, hoisted for the retry wrapper).
async fn tick_once(pool: &Pool, schema: &str) -> Result<i64> {
let conn = pool.get().await.map_err(pool_error)?;
in_tx(conn, |client| {
let schema = schema.to_string();
Box::pin(async move { tick(client, &schema, now_unix()).await })
})
.await
}
/// The loop's tick with the transient-error retry budget: a
/// pool/database failure retries [`TICK_RETRIES`] times with the short
/// [`tick_retry_backoff`] between attempts before the loop exits `Err`
/// (the module docs' retry posture — every tick error at this site is
/// pool/database-shaped, so every error is retryable; a lost lease is
/// never an error here and the loop's next top-of-iteration renew owns
/// the loss decision regardless). Each failed attempt logs (the
/// engine's diagnostic posture); the exhausting exit wraps the last
/// error with the tick-phase context in the source chain
/// ([`context_error`]).
async fn tick_with_retries(pool: &Pool, schema: &str) -> Result<i64> {
let mut retries = 0u64;
loop {
match tick_once(pool, schema).await {
Ok(soonest) => return Ok(soonest),
Err(e) if retries < TICK_RETRIES => {
retries += 1;
let delay = tick_retry_backoff(retries);
eprintln!(
"alkstore-postgres: scheduler tick failed (retry {retries}/{TICK_RETRIES} \
after {delay:?}): {e:?}"
);
tokio::time::sleep(delay).await;
}
Err(e) => {
return Err(context_error(
format!(
"scheduler tick failed after {retries} retries against transient \
pool/database errors (tick phase: the leader loop; the leadership \
lock row is left to lapse at its TTL)"
),
e,
));
}
}
}
}
/// The sliced sleep until the next due boundary (or idle horizon) or
/// the stop token — the loop's wake path between ticks. The slice is
/// 1 s, so a clean stop ends the sleep promptly and a mid-sleep lease
+158 -1
View File
@@ -10,7 +10,10 @@
//! one-firer-per-boundary, acquire-time loss before any tick, in-sleep
//! lease renewals holding across boundaries further out than the TTL,
//! the 64-cap catch-up + skip-forward with no in-window refire, the
//! tick's row-locked transactional advance, and clean-stop Ok(()).
//! tick's row-locked transactional advance, the tampered-row quarantine
//! (the runner survives a bad spec row, other schedules keep firing,
//! the bad boundary advances) and the tick retry policy pin, and
//! clean-stop Ok(()).
//!
//! Harness convention (the schema/open tasks'): connection settings
//! ride the environment (`ALKSTORE_PG_HOST/PORT/USER/PASSWORD/DB`),
@@ -1036,6 +1039,134 @@ async fn rogue_concurrent_ticks_cannot_double_fire_one_boundary() {
drop_schema(&fx.admin, &fx.schema).await;
}
/// Acceptance: a tampered spec row is *quarantined*, not fatal (the
/// scheduler-resilience posture): a row whose stored `@every` spec no
/// longer parses (tampered via direct SQL — specs are pre-validated at
/// `schedule()`, so only a tampered/foreign row is unparseable on read)
/// makes the tick log it, advance its boundary strictly past now, and
/// proceed with the remaining due rows. The runner survives (a stop
/// token still ends it `Ok(())`), the other schedule keeps firing, and
/// the bad row never fires.
#[tokio::test(flavor = "multi_thread")]
async fn tampered_spec_row_is_quarantined_and_the_runner_survives() {
let Some(fx) = Fixture::open(&instance_namer("so_quar")).await else {
eprintln!("skip: no harness server");
return;
};
fx.store
.schedule("good", "@every 1s", "work", json!(1), Default::default())
.await
.unwrap();
fx.store
.schedule(
"corrupt",
"@every 3600s",
"corruptq",
json!(1),
Default::default(),
)
.await
.unwrap();
// Tamper via direct SQL: an unparseable spec + a boundary far into
// the past (the row is due on the runner's first tick).
let before = now_unix();
fx.admin
.execute(
&format!(
"UPDATE {} SET spec = '0 9 * * *', next_fire_at = $1 WHERE name = 'corrupt'",
crate::schema::QualifiedTable {
schema: &fx.schema,
name: crate::schema::tables::SCHEDULE,
}
),
&[&(before - 3600)],
)
.await
.unwrap();
let stop = StopToken::new();
let runner_store = open_store(&harness_dsn().unwrap(), test_opts(&fx.schema))
.await
.unwrap();
let runner_stop = stop.clone();
let runner = tokio::spawn(async move { runner_store.run_schedules(runner_stop).await });
// The runner survives the tampered row: the good schedule keeps
// firing across several ticks (each of which re-quarantines the
// bad row as it re-enters the due set).
let deadline = tokio::time::Instant::now() + Duration::from_secs(25);
let fired_good: i64 = loop {
let c: i64 = fx
.admin
.query_one(
&format!(
"SELECT count(*) FROM {} WHERE queue = 'work'",
crate::schema::QualifiedTable {
schema: &fx.schema,
name: crate::schema::tables::JOB,
}
),
&[],
)
.await
.unwrap()
.get(0);
if c >= 2 {
break c;
}
assert!(
tokio::time::Instant::now() < deadline,
"the good schedule stopped firing — quarantine was fatal"
);
tokio::time::sleep(Duration::from_millis(200)).await;
};
assert!(fired_good >= 2);
// The quarantined row never fires.
let fired_corrupt: i64 = fx
.admin
.query_one(
&format!(
"SELECT count(*) FROM {} WHERE queue = 'corruptq'",
crate::schema::QualifiedTable {
schema: &fx.schema,
name: crate::schema::tables::JOB,
}
),
&[],
)
.await
.unwrap()
.get(0);
assert_eq!(fired_corrupt, 0, "a quarantined row must never fire");
// The quarantine advanced the bad row's boundary past the pre-run
// now (its backdated `before - 3600` never survives a tick), and
// the stored spec is untouched.
let (next, spec, _) = sched_row(&fx.admin, &fx.schema, "corrupt").await;
assert_eq!(spec, "0 9 * * *", "quarantine must not rewrite the spec");
assert!(
next > before,
"the quarantined boundary must advance past now (next={next}, before={before})"
);
// The runner survives to the end: a clean stop is `Ok(())` — an
// early `Err` exit (fatality) would surface here instead.
stop.cancel();
let outcome = tokio::time::timeout(Duration::from_secs(10), runner)
.await
.expect("clean stop must terminate promptly")
.unwrap();
assert!(
matches!(outcome, Ok(())),
"the runner must survive the quarantined row: {outcome:?}"
);
drop_schema(&fx.admin, &fx.schema).await;
}
/// The pg engine's own `@every` parser's unit pins (server-less —
/// ADR-012 §2's one-owner rule: this engine owns its grammar; the
/// contract suite pins the twins' outcomes identical, wave 5).
@@ -1079,3 +1210,29 @@ fn parse_every_interval_pins_the_grammar_surface() {
_ => unreachable!("cron must map InvalidSpec"),
}
}
/// The tick retry policy's pin (server-less — the transient-error
/// injection's unit arm): the budget is exactly `TICK_RETRIES` retries
/// beyond the first attempt, and the backoff shape is the short
/// doubling (250 ms → 500 ms → 1 s), capped at the 1 s ceiling.
#[test]
fn the_tick_retry_policy_is_pinned() {
use crate::scheduler::{TICK_RETRIES, tick_retry_backoff};
assert_eq!(TICK_RETRIES, 3);
assert_eq!(tick_retry_backoff(1), Duration::from_millis(250));
assert_eq!(tick_retry_backoff(2), Duration::from_millis(500));
assert_eq!(tick_retry_backoff(3), Duration::from_millis(1_000));
assert_eq!(
tick_retry_backoff(3),
tick_retry_backoff(300),
"the backoff is capped at the 1 s ceiling"
);
for r in 1..=TICK_RETRIES {
assert!(
tick_retry_backoff(r) <= Duration::from_secs(1),
"retry {r} backoff must stay short (≤ 1 s): {:?}",
tick_retry_backoff(r)
);
}
}
+76 -10
View File
@@ -1,7 +1,7 @@
---
id: pg-fix-scheduler-resilience
name: Scheduler runner resilience — quarantine bad rows, retry transient errors (review 002 Finding 5)
status: pending
status: completed
depends_on: []
scope: narrow
risk: medium
@@ -76,20 +76,20 @@ code.
## Acceptance Criteria
- [ ] A tampered/foreign spec row is quarantined (logged with name,
- [x] A tampered/foreign spec row is quarantined (logged with name,
boundary advanced past now, other schedules still fire, runner
survives) — pinned by behavioral test
- [ ] A transient tick error is retried N times (pinned N) with
- [x] A transient tick error is retried N times (pinned N) with
backoff before the loop returns `Err` — policy pinned by test
- [ ] The exiting error's chain carries schedule-name/tick-phase
- [x] The exiting error's chain carries schedule-name/tick-phase
context
- [ ] The error-exit leadership posture (TTL lapse covers the
- [x] The error-exit leadership posture (TTL lapse covers the
unreleased lock) stated in the module docs
- [ ] The `scheduler.rs:46-53` doc claim corrected to the honest
- [x] The `scheduler.rs:46-53` doc claim corrected to the honest
slice/idle posture
- [ ] Existing scheduler tests stay green (clean stop,
- [x] Existing scheduler tests stay green (clean stop,
`LeadershipLost`, catch-up cap, boundary math)
- [ ] `cargo test -p alkstore-postgres` (harness server), clippy
- [x] `cargo test -p alkstore-postgres` (harness server), clippy
`-D warnings`, fmt clean; gates green server-less
## References
@@ -105,8 +105,74 @@ code.
## Notes
> To be filled by implementation agent
- **Retry budget pinned at N = 3** (`TICK_RETRIES`), retry *beyond*
the first attempt (4 total); backoff doubling 250 ms → 500 ms → 1 s
(`tick_retry_backoff`, capped at 1 s; the full retry window adds ≤
1.75 s of stop latency). Every tick error at the loop's tick site is
pool/database-shaped (`Error::Database`), so the retry is
unconditional — `LeadershipLost` can never originate there, and a
lost lease stays the top-of-iteration renew's decision (documented
explicitly: a mid-retry lapse is caught by the renew on the next
iteration, before any tick).
- **Quarantine advance mechanism**: the task said "the tick's
skip-forward helper" — since an unparseable spec yields no interval,
the helper runs at the minimum interval
(`skip_forward(next_fire_at, now, 1)`), so the advance is strictly
past now and monotone (no within-tick spin). Consequence noted in
the module docs: the quarantined row re-enters the due set on every
subsequent tick and is re-quarantined (and re-logged) each time —
the runner keeps a ≥1 s cadence while the stored spec is broken
rather than failing; repairing/unscheduling the row stops the noise.
- **Error-context carrier**: `Database`'s Display is opaque, so a
plain `database_error(msg)` would drop the driver detail along with
losing the source chain. Implemented as `context_error`: the exiting
error is `Database → io::Error (context message) → original error →
driver chain` — the message carries schedule name / tick phase AND
the full driver chain is preserved. Wrapped sites: the fire enqueue
and the boundary advance inside `tick` (schedule name + phase), and
the loop's retry-exhausted exit (tick phase).
- **Exit-path leadership**: documentation only (the TTL-lapse posture
is the bar); no opportunistic release — keeping the exiting path
simple per the task's fine print.
- **Per-slice soonest re-check not implemented** (the task's optional
arm): the module doc now states the honest posture (1 s slices only
renew the lease and check the stop token; new schedules noticed at
the current soonest deadline, or up to 60 s late when idle).
- **No behavioral transient-error injection test**: no deterministic
arrangement exists that wouldn't depend on timing (server kills race
the retry window); the unit-level pin of count/backoff shape is the
bar and landed (`the_tick_retry_policy_is_pinned`, server-less).
- **SQLite twin parity**: `alkstore-sqlite/src/scheduler.rs` retains
the original propagate-straight-out posture (its substrate tick has
the same re-parse-from-storage exposure). The task scoped the fix to
the pg engine; if the twins' error postures must agree at wave 5's
cross-engine suite, follow up there.
## Summary
> To be filled on completion
Landed the scheduler runner resilience fix in
`alkstore-postgres/src/scheduler.rs` (review 002 Finding 5 + the
scheduler bullet of Finding 6): (1) bad-row quarantine — a due row
whose stored spec fails `@every` re-parse is logged with its name,
has its boundary advanced strictly past now via the skip-forward
helper at minimum interval, and the remaining due rows proceed (the
tick tx still commits); (2) transient-error retry — the loop's tick
(`tick_once` hoisted from the inline body) retries 3× with a short
doubling backoff (250 ms → 1 s cap) before the loop exits `Err`;
(3) error context — schedule-name/tick-phase context rides the exiting
error's source chain via the `ErrorContext` wrapper (source chain
preserved, unlike a plain message remap); (4) exit-path leadership
honesty — the module docs state the TTL-lapse posture for the
unreleased lock; (5) the false `scheduler.rs:46-53` doc claim
corrected (slices renew the lease/check the stop token only; new
schedules noticed at the current soonest deadline or up to 60 s late
when idle). Tests: behavioral quarantine pin (`tampered_spec_row_
is_quarantined_and_the_runner_survives` — tampered via direct SQL
UPDATE, runner survives to a clean `Ok(())`, other schedule fires, bad
row never fires, boundary advanced) and the server-less retry-policy
pin (`the_tick_retry_policy_is_pinned`). Verified: full
`cargo test -p alkstore-postgres` green against the harness server
(121 lib + 10 contract_suite + 9 schema, including all pre-existing
scheduler tests), workspace `cargo build`/`cargo test`,
`cargo clippy --all-targets -- -D warnings`, `cargo fmt --check` all
green server-less.