10 KiB
10 KiB
id, name, status, depends_on, scope, risk, impact, level, tags
| id | name | status | depends_on | scope | risk | impact | level | tags | ||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|
| sqlite-engine-scheduler-outbox | SQLite engine — scheduler (`schedule`/`unschedule`/`run_schedules`) and outbox helper | completed |
|
broad | high | component | implementation |
|
Description
Implement the scheduler collapse (ADR-009) and the outbox helper (ADR-019 §1, ADR-014) — the last mechanisms, and the only ones with their own long-running machinery:
Scheduler:
schedule(name, spec, queue, payload, opts)— upsert by name via the substrate'sscheduler_register; validate name AND queue argument (shared-namespace rules, ADR-021 §3); validate the spec against the v1 grammar (@every <n><unit>, s|m|h|d) before storage —InvalidSpec, no round trip (the substrate'sparse_every_intervalis the parser; map its rejection toInvalidSpec). Returns theScheduleread-back (via the core constructor).unschedule(name)— substrate'sscheduler_unregister;true= removed,false= absent.run_schedules(stop)— the leader loop (ADR-009 §1), engine-side async over the sync substrate tick:- Acquire the leadership lock (
__alkstore_scheduler— the reserved-prefix derived name; use the engine's own lock machinery with a derived TTL, engine-internal detail). - Loop: renew the lock — on loss, return
Err(LeadershipLost)before ticking (never fire on a stolen lease); callscheduler_tick(fires due boundaries, catch-up-capped at 64, skip-forward — all substrate-side); sleep untilscheduler_soonestorstop. - A cancelled
StopTokenends the sleep and returns the cleanOk(()). The loop runs inspawn_blocking(it's a blocking sleep loop) with the stop token checked between iterations — or a tokio-select shape over a watch channel wrapping the token (implementer's choice; the token's flip-everywhere semantics are core's).
- Acquire the leadership lock (
- The runner must be honest about the wake path: after a tick fires jobs, queue claims proceed through the ordinary queue machinery — the scheduler adds no delivery machinery (ADR-009 §3).
Outbox:
outbox(name)— validated constructor; derives the backing queue name__alkstore_outbox:{name}(reserved prefix — unreachable byenqueue_txby design).Outbox::enqueue— auto-commit enqueue into the backing queue, stamping the outbox's derived 60/5/5 set (ADR-014 §1) — the one non-plain defaults set.Outbox::run_once(worker_id, delivery)— claim one job from the backing queue, rundelivery.deliver(job):Ok⇒ ack;Err(e)⇒retry(Some(e), None)(the queue's curve).true= claimed and processed;false= no work. No engine-issued heartbeat inside delivery (the boxedJobHandle's own heartbeat is the consumer's renewal path).
Acceptance Criteria
scheduleupserts; spec/name/queue validation all pinned (InvalidSpecpre-storage,InvalidName/ReservedNameon both name and queue argument);Scheduleread-back correctunschedulereturns true/false correctlyrun_schedules: fires due boundaries into the named queue (jobs claimable via the ordinary queue surface); leadership lock held; a second concurrent runner does not double-fire (one-firer-per-boundary — the writer serialization floor)- Leadership loss returns
Err(LeadershipLost)before any tick; clean stop returnsOk(()); catch-up cap + skip-forward observable through the substrate's pinned behavior - Scheduler-fired jobs stamp the target queue's derived defaults
(300/3/5/none) with
ScheduleOptsapplied over them (ADR-020 §3) —get_job-visible - Outbox: derived backing queue unreachable by
enqueue_tx(reserved-prefix rejection);outbox_enqueue_tx(from the seam task) commit-atomic;enqueuestamps 60/5/5;run_onceacks onOk, retries onErr, returns false on empty cargo test -p alkstore-sqlite, clippy-D warnings, fmt clean
References
- docs/architecture/decisions/009-scheduler-collapse.md
- docs/architecture/decisions/014-outbox-tx-enqueue.md
- docs/architecture/decisions/019-mechanism-handle-surfaces.md §1, §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-sqlite.md (Mapping the contract: scheduler/outbox)
- alkstore-sqlite/src/substrate/queue_ops.rs (scheduler_* fns)
Notes
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'slock_acquirematches 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 (returnsErr(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 IMMEDIATEon 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'sscheduler_tickovernow_unix, thenscheduler_soonestfor 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'sInvalidSpecmaps the substrate'sparse_every_intervalrejection —parse_every_intervalwas madepub(crate)(one-line delta insubstrate/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'sErrarm passes the delivery error'sDisplaystring as the retry's exhaust string — the same postureJobHandle::retry(Some(e), …)carries (the string surfaces on a dead row only if the budget exhausts). The curve delay is the engine's one-ownerbackoff_delay_s(the queue handle's owner, reused — ADR-010 §3's single-owner rule).SqliteJobHandlegrew apub(crate) fn newso the outbox can rebuild the claim handle for the delivery callable (run_onceclaims through the sameclaim_batchmachinery and rewraps the row); fields stay private.- Test flake fix: the catch-up-cap test originally drove
@every 1sfrom a −200 s horizon; under full-suite load the runner's next live boundary could race into the assertion window. Switched to@every 1hwith 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
What landed, verified how:
alkstore-sqlite/src/scheduler.rs—schedule/unschedule/run_scheduleswired to theStoretrait instore.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_schedulerwith per-instance owner, renew-before-tick, transactional tick (BEGIN IMMEDIATE), soonest-driven sliced sleep with in-sleep lease renewals, clean-stopOk(())and lossErr(LeadershipLost).alkstore-sqlite/src/outbox.rs—outbox(name)wired (validated constructor),enqueuestamping the 60/5/5 derived set with theEnqueueOpts::max_attemptsoverride and full ADR-020 resolution, andrun_once(worker_id, delivery)claiming via the ordinary queue machinery, acking onOk,retry(Some(e), None)(curve delay, delivery error as exhaust string) onErr,falseon empty, the boxedSqliteJobHandlehanded to the caller with no engine-issued heartbeat.- Small substrate/engine deltas:
parse_every_intervalmadepub(crate);SqliteJobHandle::new(pub(crate)) added;jobs_from_json_pagemadepub(crate);SCHEDULER_LEADER_LOCKconstant lives inlock.rsnext 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 viaget_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) andstore/outbox_tests.rs(5 — constructor validation, backing-queue unreachability byenqueue_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 bytx_tests.rs(tx_stamps_resolve_max_attempts_and_outbox_set,outbox_backing_queue_name_is_reserved_and_derived, the validation rows intx_entry_points_validate_names).- Gates: workspace build/test green (216 total), clippy
-D warningsclean, fmt clean.