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

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
pg-engine-queues
pg-engine-locks
broad high component implementation
wave-4
postgres-engine

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 (Codec propagation); 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 (InvalidSpec before 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_lock with 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 … COMMIT on the held object, not BEGIN IMMEDIATE — pg's default isolation gives the row-locked fire tx ADR-009 §4 names; the schedule row's SELECT … FOR UPDATE inside 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 the ScheduleOpts resolved 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.
  • Boundary math is this engine's own implementation (ADR-012 §2's one-owner rule — the substrate's scheduler_tick/scheduler_soonest are SQLite-side): the @every interval 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 by enqueue_tx/queue() by design), stamp the 60/5/5 set with the EnqueueOpts::max_attempts override, 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's Delivery trait object — ADR-019 §5's pinned shape), ack on Ok, retry(err, None) (the queue's curve; the delivery error's Display string as the retry's exhaust string — the wave-3 posture) on Err, false on empty. No engine-issued heartbeat inside delivery (honker parity, inherited deliberately — the boxed JobHandle's own heartbeat remains available to the closure); the boxed JobHandle handed to the caller.
  • outbox_enqueue_tx landed 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 → InvalidSpec pre-storage), upsert read-back; the pg-owned parser pinned by tests (accepts/rejects the grammar's exact surface)
  • run_schedules: leadership via __alkstore_scheduler with 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 returns Err(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 by enqueue_tx/queue(); 60/5/5 stamps + max_attempts override; run_once ok/err/false dispositions with the curve gap asserted; in-delivery heartbeat renewal available; exhaustion to dead with the delivery error string
  • cargo 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_acquire directly, not try_lock — try_lock's entry-point shared-name validation rejects the reserved name (__alkstore_scheduler), so the engine's own holder builds the lease via PgLockHandle::new_internal (the SQLite twin's SqliteLockHandle::new_internal shape) 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 expires rides ScheduleOpts::expires as a relative-seconds EnqueueOpts::expires through the shared enqueue_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. The enqueue_row stamp override re-resolves max_attempts no-op-fully (the schedule row's stamps are already resolved and bind directly).
  • An EngineCtx bundle (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_once claim 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 made pub(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 through PgJobHandle::ack/retry (the new constructor made pub(crate) for the rebuild, mirroring the SQLite twin's SqliteJobHandle::new).
  • The best-effort enqueue wake on the backing queue's channel (the mechanism-name-is-the-channel realization) rides the outbox enqueue too — 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/outbox methods are wired, so store_trait_methods_are_wiring_stubs (which asserted the Err(Database(… wiring lands with …)) shape on them) now asserts the wired constructor/register/unregister surface instead; the stub_err! macro left with the removed stub calls.
  • tick is pub(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 grammar parse_every_interval — cron strings reject pre-storage — the upsert over the schema table with the plain-queue derived stamps
    • ScheduleOpts applied, Schedule read-back), unschedule (true/false), run_schedules (leadership via the engine's lock machinery on __alkstore_scheduler with 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-locked FOR UPDATE tick in one pool-checked-out BEGIN/COMMIT frame — fire enqueues + boundary advance + soonest read commit together — the 64-boundary catch-up cap with skip-forward past-now, clean stop Ok(()) releasing the lock, loss Err(LeadershipLost) before any tick); outbox(name) — the validated constructor, enqueue into the derived __alkstore_outbox:{name} (60/5/5 stamps, the max_attempts override, full ADR-020 resolution, best-effort wake) and run_once(worker_id, delivery) (the ordinary claim frame, ack on Ok, retry(err-as-exhaust-string, curve-delay) on Err, false on empty, no engine-issued heartbeat — the boxed handle's own heartbeat is the consumer's).
  • Store trait wiring (store.rs): the four stubs replaced; the stub fn and its last call sites removed. Queue module deltas: in_tx/sweep_exhausted_reclaimables_tx/claim_stmt_sql/ PgJobHandle::new made pub(crate) (one-owned machinery reuse); lock module deltas: lock_acquire made pub(crate) + PgLockHandle::new_internal added.
  • Tests: store/scheduler_tests.rs (9 — the validation triad + exact parser surface, upsert/unregister read-backs, fire + stamp source via get_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) and store/outbox_tests.rs (5 — constructor validation, enqueue_tx unreachability + 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-postgres green against the harness server (111 + 9, six consecutive full runs — stable); workspace cargo test green server-less (the pg skips cleanly); cargo clippy --all-targets -- -D warnings and cargo fmt --check clean workspace-wide. Pre-existing flake (streams' wake-timing deadline race) untouched — none of the six runs tripped it.