Files
alkstore/tasks/sqlite-engine-scheduler-outbox.md

10 KiB
Raw Permalink Blame History

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
sqlite-engine-queues
sqlite-engine-locks
broad high component implementation
wave-3
sqlite-engine

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's scheduler_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's parse_every_interval is the parser; map its rejection to InvalidSpec). Returns the Schedule read-back (via the core constructor).
  • unschedule(name) — substrate's scheduler_unregister; true = removed, false = absent.
  • run_schedules(stop) — the leader loop (ADR-009 §1), engine-side async over the sync substrate tick:
    1. 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).
    2. Loop: renew the lock — on loss, return Err(LeadershipLost) before ticking (never fire on a stolen lease); call scheduler_tick (fires due boundaries, catch-up-capped at 64, skip-forward — all substrate-side); sleep until scheduler_soonest or stop.
    3. A cancelled StopToken ends the sleep and returns the clean Ok(()). The loop runs in spawn_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).
  • 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 by enqueue_tx by 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, run delivery.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 boxed JobHandle's own heartbeat is the consumer's renewal path).

Acceptance Criteria

  • schedule upserts; spec/name/queue validation all pinned (InvalidSpec pre-storage, InvalidName/ReservedName on both name and queue argument); Schedule read-back correct
  • unschedule returns true/false correctly
  • run_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 returns Ok(()); 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 ScheduleOpts applied 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; enqueue stamps 60/5/5; run_once acks on Ok, retries on Err, 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'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

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.