14 KiB
14 KiB
id, name, status, depends_on, scope, risk, impact, level, tags
| id | name | status | depends_on | scope | risk | impact | level | tags | ||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|
| pg-engine-scheduler-outbox | Postgres engine — scheduler (`schedule`/`unschedule`/`run_schedules`) and outbox helper | completed |
|
broad | high | component | implementation |
|
Description
Implement the scheduler collapse's pg arm (ADR-009) and the outbox helper (ADR-014/019 §5) — the wave's last mechanism pair, riding the queues task's claim machinery and the locks task's lock machinery:
Scheduler (over the schema task's schedule table):
schedule(name, spec, queue, payload, opts)— upsert by name; entry-point validation on all three name-bearing arguments (schedule name like a queue name; the queue argument identically per ADR-021 §3 — a schedule cannot target a reserved name, so the outbox backing queue stays reachable only through the outbox surface); payload encoded at the seam (Codecpropagation); the spec grammar is@every <n><unit>only (s|m|h|d) — the pg engine owns its parser (no shared substrate parser on this side; the grammar is small and pinned; cron strings →InvalidSpec). Rejection is pre-storage (InvalidSpecbefore any round trip).unschedule(name)— the unregister;Ok(false)when absent.run_schedules(stop: StopToken)— the leader loop, the structural twin of the SQLite engine's runner with the pg deltas:- Leadership via the engine's own lock machinery on the
reserved name
__alkstore_scheduler(ADR-009 §1) —try_lockwith a per-runner-instance-unique owner token (the wave-3 finding's lesson: a shared constant would make two processes each see "self" holding; the token is pid + subsec-nanos or equivalent). TTL engine-internal (10 s), renewed at the top of every loop iteration and on a cadence across the sleep (the wave-3 posture: without in-sleep renewals, a boundary further out than the TTL lapses its own lease by drift). A failed mid-sleep renewal is not acted on at the sleep site — the loop's top-of-iteration renew owns the loss decision. - Acquire-time refusal is
Err(LeadershipLost)— the same matchable loss a mid-run loser returns (the wave-3 posture: no runner-loop exists without the lock; the respawn recipe is identical either way). - The tick: fire enqueues + boundary advance commit together —
on pg this is one pool-checked-out transaction (
BEGIN…COMMITon the held object, notBEGIN IMMEDIATE— pg's default isolation gives the row-locked fire tx ADR-009 §4 names; the schedule row'sSELECT … FOR UPDATEinside the tick tx is the row lock that serializes rogue second tickers). The tick body: boundary math over the registered specs (the one-firer-per-boundary guarantee — at-least-once per elapsed boundary, the 64-boundary catch-up cap with skip-forward past the cap, honker parity), each fire an ordinary enqueue stamped with theScheduleOptsresolved over the target queue's derived engine defaults (no handle is open — ADR-020 §3), then the soonest-boundary read for the sleep deadline. - Wake-driven tick advance (the pg delta over the SQLite runner's pure-sleep shape): the runner parks on the sleep deadline, but an enqueue-wake on the scheduler's own wake channel (or a plain bounded-wait interrupt — implementer's choice, documented) may shorten the sleep so a newly-registered schedule is noticed promptly; correctness never rides the wake (the deadline always fires; the wake only shortens idle latency — the SQLite runner's 60 s idle-tick posture is the acceptable floor if the wake path proves fussy; flag the choice in Notes).
- Clean stop (
StopToken::cancel) →Ok(()), the leadership lock released cleanly; loss →Err(LeadershipLost)before any tick.
- Leadership via the engine's own lock machinery on the
reserved name
- Boundary math is this engine's own implementation (ADR-012 §2's
one-owner rule — the substrate's
scheduler_tick/scheduler_soonestare SQLite-side): the@everyinterval arithmetic, the boundary advance, the 64-cap catch-up, and the skip-forward. The contract suite (wave 5) pins the two engines' boundary behavior to identical outcomes.
Outbox (a helper over queues — no separate mechanism):
Store::outbox(name)— validated constructor (empty →InvalidName, reserved-prefixed →ReservedName); carries no opts of its own (the backing queue's derived 60/5/5 set is engine-documented, ADR-010 §3a).enqueue(payload, opts)— the auto-commit path into the derived backing queue: derive__alkstore_outbox:{name}(the reserved prefix makes it unreachable byenqueue_tx/queue()by design), stamp the 60/5/5 set with theEnqueueOpts::max_attemptsoverride, full ADR-020 resolution, ride the queues task's enqueue machinery.run_once(worker_id, delivery) -> bool— the pull op: claim one job from the backing queue via the ordinary claim machinery (the queues task's), run the delivery callable (core'sDeliverytrait object — ADR-019 §5's pinned shape),ackonOk,retry(err, None)(the queue's curve; the delivery error'sDisplaystring as the retry's exhaust string — the wave-3 posture) onErr,falseon empty. No engine-issued heartbeat inside delivery (honker parity, inherited deliberately — the boxedJobHandle's ownheartbeatremains available to the closure); the boxedJobHandlehanded to the caller.outbox_enqueue_txlanded with the seam task; this task's tests may extend its coverage (the commit-atomicity row is suite-owned in wave 5).
Acceptance Criteria
schedule: validation triad (name/spec/queue),@every-only grammar (cron →InvalidSpecpre-storage), upsert read-back; the pg-owned parser pinned by tests (accepts/rejects the grammar's exact surface)run_schedules: leadership via__alkstore_schedulerwith a per-instance owner token; acquire-time refusal =Err(LeadershipLost); tick fires + boundary advance in one transaction (row-locked — a rogue second ticker cannot double-fire); loss returnsErr(LeadershipLost)before any tick; clean stop releases the lock; in-sleep lease renewals hold across boundaries further out than the TTL- One-firer-per-boundary: two runners on the same db, exactly one fire per boundary (the band test); 64-cap catch-up + skip-forward with no in-window refire
- Fired jobs stamp the derived defaults resolved over
ScheduleOpts(get_job-visible) outbox: constructor validation; backing queue unreachable byenqueue_tx/queue(); 60/5/5 stamps +max_attemptsoverride;run_onceok/err/false dispositions with the curve gap asserted; in-delivery heartbeat renewal available; exhaustion to dead with the delivery error stringcargo test -p alkstore-postgres(harness server), clippy-D warnings, fmt clean; gates green server-less
References
- docs/architecture/decisions/009-scheduler-collapse.md (§1–§6)
- docs/architecture/decisions/014-outbox-tx-enqueue.md
- docs/architecture/decisions/019-mechanism-handle-surfaces.md §4, §5
- docs/architecture/decisions/020-enqueue-opt-semantics-and-bridges.md §3
- docs/architecture/decisions/021-tx-reads-and-value-shape-fixes.md §3
- docs/architecture/engine-postgres.md (Mapping the contract: scheduler/outbox)
- docs/architecture/queues.md (scheduler section)
- alkstore-sqlite/src/scheduler.rs, outbox.rs (the structural twin — same semantics, own boundary math; its Notes carry the flake-fix lessons: drive catch-up tests from far-out horizons so no live boundary races the assertion window)
- alkstore/src/stop_token.rs, outbox.rs (the traits)
Notes
Decisions of record made in implementation that the description didn't pin:
- Leadership acquire rides the lock module's
lock_acquiredirectly, nottry_lock—try_lock's entry-point shared-name validation rejects the reserved name (__alkstore_scheduler), so the engine's own holder builds the lease viaPgLockHandle::new_internal(the SQLite twin'sSqliteLockHandle::new_internalshape) after the raw acquire op. Same machinery, bypassing the consumer entry gate by construction.- The wake-driven tick advance chose the bounded-wait interrupt (the task's implementer's choice): no dedicated wake channel on the scheduler's side — the sleep is sliced at 1 s (
SLEEP_SLICE_S) with the stop token checked and the lease renewal cadence (2 s) evaluated per slice, and the soonest/ boundary state is re-read at the top of every loop iteration. A newly-registered schedule is therefore noticed at the next loop re-read after its boundary (and by the idle-tick's 60 s re-read at the latest) — the deadline always fires; the wake never gates correctness. The SQLite runner's idle-tick posture stands as the floor.- The tick's fire-time
expiresridesScheduleOpts::expiresas a relative-secondsEnqueueOpts::expiresthrough the sharedenqueue_row— the resolution (relative → absolute from the fire instant) is the resolution module's existing formula (one owner per formula, ADR-012 §2), not a duplicated scheduler-side arithmetic. Theenqueue_rowstamp override re-resolvesmax_attemptsno-op-fully (the schedule row's stamps are already resolved and bind directly).- An
EngineCtxbundle (pool/schema/closed) threads the store wiring into the scheduler/outbox module's entry points —schedule's engine-side argument list stayed within the clippy arg guard and the wiring reads once at each call site.- The outbox's
run_onceclaim rebuilds the queues task's claim frame from the module's shared pieces (in_tx+sweep_exhausted_reclaimables_tx+claim_stmt_sql+job_from_row, all madepub(crate)) with the backing queue name standing in for the handle's — one owner per formula, no duplicated claim SQL; the ok/err dispositions route throughPgJobHandle::ack/retry(thenewconstructor madepub(crate)for the rebuild, mirroring the SQLite twin'sSqliteJobHandle::new).- The best-effort enqueue wake on the backing queue's channel (the mechanism-name-is-the-channel realization) rides the outbox
enqueuetoo — the queue enqueue's twin posture (durable row is the truth, the re-poll safety net covers a failed wake).- The open-tests' stub-probe test retired — all four
schedule/unschedule/run_schedules/outboxmethods are wired, sostore_trait_methods_are_wiring_stubs(which asserted theErr(Database(… wiring lands with …))shape on them) now asserts the wired constructor/register/unregister surface instead; thestub_err!macro left with the removed stub calls.tickispub(crate)(test-visible) for the rogue-ticker row-lock probe — two concurrent pool connections driving the engine's own tick op directly (no leadership lock in either hand) fire exactly one boundary.
Summary
What landed, verified how:
alkstore-postgres/src/scheduler.rs— the collapse shape's pg arm + the outbox helper:schedule(the validation triad pre-storage, the pg-owned@every-only grammarparse_every_interval— cron strings reject pre-storage — the upsert over the schema table with the plain-queue derived stamps
ScheduleOptsapplied,Scheduleread-back),unschedule(true/false),run_schedules(leadership via the engine's lock machinery on__alkstore_schedulerwith a per-instance pid + subsec-nanos owner token, TTL 10 s engine-internal, renew at every loop top and on a 2 s cadence across the sliced sleep, acquire-time refusal =Err(LeadershipLost), the row-lockedFOR UPDATEtick in one pool-checked-outBEGIN/COMMITframe — fire enqueues + boundary advance + soonest read commit together — the 64-boundary catch-up cap with skip-forward past-now, clean stopOk(())releasing the lock, lossErr(LeadershipLost)before any tick);outbox(name)— the validated constructor,enqueueinto the derived__alkstore_outbox:{name}(60/5/5 stamps, themax_attemptsoverride, full ADR-020 resolution, best-effort wake) andrun_once(worker_id, delivery)(the ordinary claim frame, ack onOk,retry(err-as-exhaust-string, curve-delay)onErr,falseon empty, no engine-issued heartbeat — the boxed handle's ownheartbeatis the consumer's).- Store trait wiring (
store.rs): the four stubs replaced; thestubfn and its last call sites removed. Queue module deltas:in_tx/sweep_exhausted_reclaimables_tx/claim_stmt_sql/PgJobHandle::newmadepub(crate)(one-owned machinery reuse); lock module deltas:lock_acquiremadepub(crate)+PgLockHandle::new_internaladded.- Tests:
store/scheduler_tests.rs(9 — the validation triad + exact parser surface, upsert/unregister read-backs, fire + stamp source viaget_job+ live leadership lock row with the pid-stamped owner + clean-stop release, second-runner refusal + the one-firer band, acquire-time loss with a zero-fire probe, in-sleep lease renewals holding across a 40 s boundary on a 10 s TTL, the 64-cap + skip-forward with no in-window refire from a 200-boundary far-out horizon, the rogue-ticker row-lock serialization, and the server-less parser unit pin) andstore/outbox_tests.rs(5 — constructor validation,enqueue_txunreachability + 60/5/5 stamps + override, ok/err/ false dispositions with the curve gap asserted, in-delivery heartbeat renewal, exhaustion to dead with the delivery error string).open_tests.rs's stub probe updated to the wired posture.- Verified:
cargo test -p alkstore-postgresgreen against the harness server (111 + 9, six consecutive full runs — stable); workspacecargo testgreen server-less (the pg skips cleanly);cargo clippy --all-targets -- -D warningsandcargo fmt --checkclean workspace-wide. Pre-existing flake (streams' wake-timing deadline race) untouched — none of the six runs tripped it.