Wave-3 decomposition: SQLite engine tasks (open/opts, seam+tx, mechanisms, integration, review gate) + core value-constructor pre-work; plan synced for waves-1-2 review + ADR-023 follow-through

This commit is contained in:
glm-5.3-flash committed 2026-10-08 10:57:29 +00:00
1 parent 44637eea5b
commit 1269246faf
11 files changed
+883 -5

No files matched your search

+50 -5
View File
@@ -1,6 +1,6 @@
---
status: draft
last_updated: 2026-10-07 (initial wave plan; waves 1–2 decomposed)
last_updated: 2026-10-08 (waves 1–2 implemented + reviewed; ADR-023 follow-through folded into wave 3; wave 3 decomposed)
---
# alkstore — Implementation plan
@@ -32,9 +32,9 @@ parallel).
| Wave | Contents | Depends on | Status |
|---|---|---|---|
| [1](#wave-1--foundations) | Workspace scaffold; core crate (errors, value types, full trait surface); contract-suite scaffold | — | decomposed (`tasks/`) |
| [2](#wave-2--sqlite-substrate-fork) | honker-core fork into `alkstore-sqlite/src/substrate/`: port, deltas, provenance, test floor | wave 1 (workspace scaffold only) | decomposed (`tasks/`) |
| 3 | SQLite engine: connection architecture, re-derived queue ops on contract v1, scheduler/outbox, tx seam, SQLite backlog column | waves 1 + 2 | not yet decomposed |
| [1](#wave-1--foundations) | Workspace scaffold; core crate (errors, value types, full trait surface); contract-suite scaffold | — | implemented + reviewed (2026-10-08) |
| [2](#wave-2--sqlite-substrate-fork) | honker-core fork into `alkstore-sqlite/src/substrate/`: port, deltas, provenance, test floor | wave 1 (workspace scaffold only) | implemented + reviewed (2026-10-08) |
| [3](#wave-3--sqlite-engine) | SQLite engine: connection architecture, re-derived queue ops on contract v1, scheduler/outbox, tx seam, SQLite backlog column | waves 1 + 2 | decomposed (`tasks/`) |
| 4 | Postgres engine: schema bootstrap, pool/open, listener/forwarder, all mechanisms, tx seam, pg backlog column | wave 1 | not yet decomposed |
| 5 | Contract suite: the cross-engine equivalence properties (core-contract.md §Verification backlog), version-stamped per ADR-017 | waves 3 + 4 | not yet decomposed |
| 6 | Release readiness: crate docs, deployment matrix final pass, README (written last, honestly), publish prep; mem-engine and fuzzing decisions | wave 5 | not yet decomposed |
@@ -68,6 +68,39 @@ the engine layer that maps it onto the core contract is wave 3, not
here. Reviewability against the lineage (ADR-012 §3) is a property the
wave-2 review gate checks explicitly.
## Wave 3 — SQLite engine
The engine layer that maps the forked substrate onto the core
contract: the `open` constructor (connection architecture — writer
slot, reader pool, watcher spawn — per engine-sqlite.md), the
`spawn_blocking` seam, the full `Store`/`TxHandle`/mechanism-handle
trait impls over the substrate's ops, the scheduler leader loop and
outbox helper, and the engine's backlog column in the contract suite
(the ADR-023 rows among them, factory-parameterized so wave 5 runs
them against both engines).
The waves-1–2 general review
(`docs/reviews/2026-10-08-waves-1-2-general-review.md`) and its
resolution ([ADR-023](../architecture/decisions/023-fourth-review-round.md))
shaped this wave's task set:
- **Already landed pre-decomposition** (commit `44637ee`, not wave-3
tasks): `encode_payload` is fallible (`Result<Vec<u8>, Error::Codec>`,
ADR-023 §1); `open_conn` drops `SQLITE_OPEN_URI` (§3, register
D-29); the numeric-argument domain table is pinned in
core-contract.md (§2); the 1 ms watcher default stands with the
cadence documented in deployment.md (§4).
- **Folded into wave 3's tasks** (the review's §7 wiring items): the
trait-impl extent/duration guards (the contract-side domain rule,
enforced at the engine's trait-impl entry — the `sqlite-engine-*`
tasks carry it per mechanism), the two `#![allow]` lint removals in
`substrate/mod.rs` (the integration task — they can only lift once
every substrate surface is wired), and `SqliteOpts::poll_interval`
wiring (the constructor task).
- **Staying put as ordered**: M-1's retention-failure test and N-5's
panic probe → wave 5's contract-suite rows; N-1's growth-posture
review → wave 6.
## Decided points
- **Contract-suite layout — option (a)**: a small internal
@@ -108,4 +141,16 @@ decomposes. Specific gates:
columns each engine owns.
- **Wave 5 review** — the suite as compatibility instrument: every
backlog row present, version-stamped, green on both engines; this is
the gate that flips the engine specs to `stable`.
the gate that flips the engine specs to `stable`.
## Review rounds so far
- **Wave 1 review gate** (`review-wave-1`) — trait surface vs pinned
ADR text; validation coverage fixes landed.
- **Wave 2 review gate** (`review-wave-2`) — lineage diff clean, one
re-derivation defect found and fixed (D-27).
- **General review, waves 1–2** (2026-10-08,
`docs/reviews/2026-10-08-waves-1-2-general-review.md`) — M-1 fixed
inline (sweep savepoint scope); M-2/N-2/N-4/N-6 resolved as ADR-023
pre-decomposition; lint removal + suite adds folded into waves 3/5;
N-1 recorded for wave 6.
+67
View File
@@ -0,0 +1,67 @@
---
id: core-engine-value-constructors
name: Core — engine-side constructors for the non_exhaustive value types
status: pending
depends_on: []
scope: narrow
risk: low
impact: component
level: implementation
tags: [wave-3, core, pre-work]
---
## Description
The engine crates must produce `Job`, `StreamEvent`, `Schedule`, and
`Wake` — but all four are `#[non_exhaustive]` (ADR-017 §3), which bans
struct-literal construction from outside the defining crate (E0639).
Wave 3's engine is the first code that must build these values from
row data, and it cannot. Give the engine crates a sanctioned
construction path without weakening the consumer posture.
Add `#[doc(hidden)]` constructors on each type — e.g.
`Job::from_row(...)`, `StreamEvent::from_row(...)`, `Schedule::new(...)`,
`Wake::new(...)` — taking the full field set, `pub` (engines are
ordinary downstream crates, not in-crate) but `#[doc(hidden)]` so they
stay out of the consumer-facing docs and are visibly not contract
surface. Doc-comment each with its posture: engine-construction-only,
not covered by ADR-017's semver-minor field-addition promise (an
engine breaking on a field addition is a lockstep-duty event, ADR-017
§5 — engines are in-repo and updated with core).
Alternative considered and rejected: making the fields `pub(crate)` +
engine `#[cfg]` tricks, or a builder trait — both add machinery for a
problem a hidden constructor solves; the non_exhaustive posture for
*consumers* is untouched either way.
This is deliberately a separate pre-work task (not folded into the
first engine task): it touches core's contract-adjacent surface, is
consumed by both engines (wave 4 needs the same constructors), and is
cheap to review in isolation.
## Acceptance Criteria
- [ ] `Job`, `StreamEvent`, `Schedule`, `Wake` each carry a
`#[doc(hidden)]` full-field constructor, documented as
engine-construction-only (not consumer API, not semver-covered)
- [ ] No change to the consumer posture: the types stay
`#[non_exhaustive]`; consumers cannot reach the constructors in
docs (spot-check `cargo doc`)
- [ ] Core tests construct through the new constructors where natural
(no behavior change; existing tests keep passing)
- [ ] `cargo build`, `cargo test -p alkstore`, clippy `-D warnings`,
fmt clean
## References
- docs/architecture/decisions/017-contract-versioning.md §3, §5
- docs/architecture/decisions/019-mechanism-handle-surfaces.md §3
- alkstore/src/job.rs, stream.rs, schedule.rs, wake.rs
## Notes
> To be filled by implementation agent
## Summary
> To be filled on completion
+74
View File
@@ -0,0 +1,74 @@
---
id: review-wave-3
name: Review gate — wave 3 (SQLite engine)
status: pending
depends_on: [sqlite-engine-integration]
scope: moderate
risk: low
impact: phase
level: review
tags: [wave-3, review]
---
## Description
Review the SQLite engine before wave 4 (or 5) decomposes. Primary
lens per the plan's review gates: **engine-vs-contract conformance** —
the engine crate now carries the first real consumer surface, so this
is the gate that checks the code against the pinned contract text
line by line (the wave-1 gate's discipline, now against an
implementation that must *honor* the surface, not just declare it).
Check:
- **Contract conformance**: every trait method's behavior against
core-contract.md's per-mechanism text — validation at entry points
(auto-commit AND tx paths), error arms (`Database`-only receivers,
`Codec` at payload seams, `LeadershipLost` vs clean stop), value
shapes (`Job`/`StreamEvent` field-for-field), the no-work-is-a-value
vocabulary (`None`/`false`/empty-`Vec` where pinned, never errors).
- **ADR-023 follow-through**: the extent guards actually sit at the
trait-impl entry (not substrate-side patches); duration guards on
both lock call sites; plain-path open holds; `poll_interval` wired;
`encode_payload` consumed fallibly at every call site.
- **Seam integrity**: drop = rollback under panic; writer-slot lease
released on every path (commit, rollback, drop, error);
`spawn_blocking` used consistently (no blocking call inside async
context, no connection moved across await).
- **Substrate boundary**: still contract-blind (no `alkstore::` under
`src/substrate/`); the lint removal didn't leave dead ported code
behind; any new divergence registered in PROVENANCE.md.
- **Backoff curve**: the engine-side formula matches ADR-010 §3's
pinned range and cap exactly (wave 4 re-derives it; wave 5 pins
equivalence — this gate checks the one implementation against the
text).
- **Backlog column**: the engine's suite rows present, stamped, green.
- **Family standards**: no panics, no unwrap/expect outside tests, no
inline comments (doc comments fine), module-per-file.
- **Gates**: full workspace build/test/clippy/fmt.
## Acceptance Criteria
- [ ] Contract conformance spot-checked against the spec text (not
just the tests)
- [ ] ADR-023's five backlog rows verified implemented as decided
- [ ] Seam/lease/rollback paths verified (code-read, not test-name-trust)
- [ ] Substrate boundary + provenance register still clean
- [ ] All gates green
- [ ] Findings recorded; wave 4/5 decomposition may proceed
## References
- docs/plans/implementation.md (Review gates)
- docs/architecture/core-contract.md
- docs/architecture/engine-sqlite.md
- docs/architecture/decisions/023-fourth-review-round.md
- docs/reviews/2026-10-08-waves-1-2-general-review.md (the posture this gate inherits)
## Notes
> To be filled by implementation agent
## Summary
> To be filled on completion
+90
View File
@@ -0,0 +1,90 @@
---
id: sqlite-engine-integration
name: SQLite engine — integration (lint removal, full-surface wiring, suite adoption)
status: pending
depends_on: [sqlite-engine-notify-listen, sqlite-engine-streams, sqlite-engine-queues, sqlite-engine-locks, sqlite-engine-scheduler-outbox]
scope: moderate
risk: medium
impact: phase
level: implementation
tags: [wave-3, sqlite-engine, integration]
---
## Description
Close the wave's integration gaps once every mechanism is wired:
- **Remove the two lint suppressions** in
`alkstore-sqlite/src/substrate/mod.rs` (`#![allow(dead_code)]`,
`#![allow(unused_imports)]`) — the waves-1–2 review's explicit
wave-3 obligation. With the engine layer wired, every substrate
surface should be reachable; genuinely dead ported code must be
**cut, not suppressed** (and the cut registered in PROVENANCE.md if
it diverges from the lineage's kept set — ADR-018's delta
discipline). Fix whatever the lints surface; do not re-add allows.
- **Full-surface sweep**: every `Store`/`TxHandle`/mechanism trait
method implemented (no `todo!`/`unimplemented!` anywhere); the
engine crate's lib re-exports its public surface (`open`,
`SqliteOpts`, the concrete store type) per the module-per-file
convention.
- **Contract-suite adoption** (the plan's "waves 3 and 4 adopt the
harness"): implement `StoreFactory` for the SQLite engine (fresh
temp-file store per `open`, idempotent teardown) and wire the
engine's backlog column — the rows this engine owns now, per
ADR-023's verification-backlog additions and the engine-scoped
backlog rows:
- Extent-clamp semantics (`claim_batch(n <= 0)`, stream reads
`limit <= 0`) — stamped `ADR-023 §2`
- Duration-refusal (`try_lock`/`renew` with `ttl <= 0`) — stamped
`ADR-023 §2`
- `encode_payload` typed-failure round-trip (enqueue/publish store
+ decode the exact serialization) — stamped `ADR-023 §1`
- Plain-path open (`?name=value` suffix = literal filename) —
stamped `ADR-023 §3`
- Poll-cadence default + knob (config-surface accounting, not
wall-clock) — stamped `ADR-023 §4`
- The engine-scoped rows core-contract.md assigns this engine:
`PayloadTooLarge` never produced (ADR-016 §5), drop = rollback
no-ghosts (ADR-021 §4), in-tx read-your-own-writes (ADR-021 §1),
enqueue-opts resolution (ADR-020 §1–§3), receiver close arms
(ADR-021 §5) — each stamped per the suite's convention.
Rows run against the SQLite factory now; wave 5 runs the same rows
against pg (the factory-parameterized point).
- **Crate docs**: the engine crate's lib-level doc comments state the
single-host posture and the writer-parking honesty note (ADR-016,
ADR-007's negative consequence) — the identity statements this
crate owns.
- **Gates**: full workspace build/test/clippy/fmt; coverage spot-check
on the engine crate (no large uncovered regions outside error arms).
## Acceptance Criteria
- [ ] Both `#![allow]` lints removed from `substrate/mod.rs`; any
code the removal surfaces is cut (and registered) or wired —
clippy `-D warnings` green without them
- [ ] No `todo!`/`unimplemented!` in the engine crate; full trait
surface implemented
- [ ] `StoreFactory` implemented for SQLite; the backlog-column rows
listed above exist, version-stamped, and run green against the
SQLite factory
- [ ] Engine crate lib docs carry the single-host + writer-parking
posture statements
- [ ] Workspace gates green: `cargo build`, `cargo test`, clippy
`-D warnings`, fmt clean
## References
- docs/reviews/2026-10-08-waves-1-2-general-review.md (§2 smells, §7)
- docs/architecture/decisions/023-fourth-review-round.md (Verification backlog additions)
- docs/architecture/core-contract.md (§Verification backlog)
- docs/architecture/decisions/022-contract-suite-layout.md
- docs/architecture/decisions/016-deployment-honesty.md
- alkstore-contract-suite/src/factory.rs
## Notes
> To be filled by implementation agent
## Summary
> To be filled on completion
+64
View File
@@ -0,0 +1,64 @@
---
id: sqlite-engine-locks
name: SQLite engine — named locks (`try_lock`, `Lock` handle, duration guards)
status: pending
depends_on: [sqlite-engine-seam-tx]
scope: narrow
risk: low
impact: component
level: implementation
tags: [wave-3, sqlite-engine]
---
## Description
Implement `Store::try_lock` and the `Lock` trait over the substrate's
lock ops (`lock_acquire`, `lock_release`, `lock_renew`):
- `try_lock(name, owner, ttl)` — validated entry point (name:
shared-namespace rules; owner: non-empty only); `None` = someone
holds it (no-work is a value). Returns the boxed `Lock` handle.
- **Duration guard (ADR-023 §2)**: `ttl <= 0` is a rejection —
`Err`, opaque `Database` shape — enforced at the trait-impl entry,
before the substrate. The substrate's `lock_renew` already guards
`ttl_s <= 0` (its message preserved in the source chain where it
fires); `lock_acquire` does not — the engine guards both call sites
uniformly so the rule is contract-side, not substrate-side (the pg
engine implements the same rule with no inherited guard at all).
- `Lock::renew(ttl)` — same duration guard; a new full TTL window
from now; `false` = lost it (expired and re-acquired elsewhere).
- `Lock::release(self)` — consuming; `true` = was held; `false` =
already expired/re-acquired (no-op, not an error).
- Exclusion lapses silently at TTL expiry (no revocation event) —
the substrate's expired-row delete on next acquire is the
mechanism; pin the lapse-then-reacquire shape in tests.
## Acceptance Criteria
- [ ] `try_lock` grants exclusively (second acquirer gets `None`);
release makes it acquirable again
- [ ] Duration guards: `try_lock`/`renew` with `ttl <= 0` return
opaque `Err` (`Database`) — never a granted lock, never a
panic; substrate guard message preserved in the source chain on
`renew` — pinned per ADR-023 §2's backlog row
- [ ] `renew` extends from now (full window, not additive); foreign
owner cannot renew; `false` on lapsed/re-acquired
- [ ] `release` removes only the owner's row; `false` = no-op shape
- [ ] TTL lapse: after expiry, another owner acquires (silent lapse)
- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
## References
- docs/architecture/core-contract.md (locks section, numeric-domain rules)
- docs/architecture/decisions/008-contract-v1-pinning.md §7
- docs/architecture/decisions/019-mechanism-handle-surfaces.md §1
- docs/architecture/decisions/023-fourth-review-round.md §2
- alkstore-sqlite/src/substrate/ops.rs (lock_* fns)
## Notes
> To be filled by implementation agent
## Summary
> To be filled on completion
+74
View File
@@ -0,0 +1,74 @@
---
id: sqlite-engine-notify-listen
name: SQLite engine — notify/listen (watcher fanout → `WakeReceiver` bridge)
status: pending
depends_on: [sqlite-engine-seam-tx]
scope: narrow
risk: medium
impact: component
level: implementation
tags: [wave-3, sqlite-engine]
---
## Description
Implement `Store::notify` and `Store::listen` — the wake half of the
notify mechanism (the tx half, `notify_tx`, landed with the seam):
- `notify(channel, payload)`: auto-commit path — writer slot, run the
substrate's `notify()` SQL function (or the direct insert it wraps)
in its own transaction, release. No payload limit on SQLite (the
`PayloadTooLarge` asymmetry — this engine never produces the
variant, ADR-016 §5).
- `listen(channel)`: subscribe to the substrate watcher's fanout
(`SharedUpdateWatcher::subscribe`), bridge the sync
`std::sync::mpsc::Receiver<()>` into a `Box<dyn WakeReceiver>`:
- `recv()` → `Some(Wake { channel })` on fanout; `None` when the
watcher died (terminal close — `WatcherDeathGuard` behavior,
ADR-006; the receiver never reopens).
- `try_recv` / `recv_timeout` per the core trait's pinned arms
(`Closed` distinguishable from idle; timeout expiry = `Ok(None)`).
- The bridge thread does `blocking_send` into a tokio channel (one
`spawn_blocking` thread per subscription — engine-sqlite.md's
mapping row); unsubscribe on receiver drop (the substrate prunes
dropped subscribers — verify the bridge doesn't leak subscriptions).
- Wakes are coalesced hints (overtrigger on purpose): the receiver
delivers `Wake { channel }` only — never payloads, never ids.
- Channel validation at the entry point (`InvalidName`/`ReservedName`).
Note the honest posture: the substrate's watcher fans out on *any*
commit (data_version bump), so a wake may arrive for unrelated writes
— that is the contract's best-effort-hint shape, not a bug; the tests
should pin wake-arrives and close-on-death, not wake-exclusivity.
## Acceptance Criteria
- [ ] `notify` delivers commit-atomic auto-commit notifies; payload
crosses the notify table (no size limit — never
`PayloadTooLarge`)
- [ ] `listen` returns a working `WakeReceiver`: wakes arrive after
commits by other connections; `Wake { channel }` carries the
channel name only
- [ ] Watcher death closes the receiver (`recv() -> None`, terminal)
— tested via the substrate's death-signal path
- [ ] `try_recv`/`recv_timeout` arms match the pinned contract shapes
- [ ] Receiver drop unsubscribes (no leaked subscriptions —
`subscriber_count` returns to baseline)
- [ ] Channel validation on both entry points
- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
## References
- docs/architecture/core-contract.md (notify section)
- docs/architecture/decisions/006-wake-and-delivery-contract.md
- docs/architecture/decisions/021-tx-reads-and-value-shape-fixes.md §5
- docs/architecture/engine-sqlite.md (Mapping the contract: notify/listen)
- alkstore-sqlite/src/substrate/watcher.rs, schema.rs (attach_notify)
## Notes
> To be filled by implementation agent
## Summary
> To be filled on completion
+86
View File
@@ -0,0 +1,86 @@
---
id: sqlite-engine-open-opts
name: SQLite engine — `open` constructor, `SqliteOpts`, connection architecture
status: pending
depends_on: [core-engine-value-constructors]
scope: moderate
risk: medium
impact: component
level: implementation
tags: [wave-3, sqlite-engine]
---
## Description
Stand up the engine crate's public surface: the `open` constructor,
`SqliteOpts`, and the connection architecture of engine-sqlite.md's
"Connection architecture" section — the first wiring of the ported
substrate.
- `SqliteOpts` (engine-crate type, ADR-008 §6's config split — engine
opts live engine-side): at minimum `poll_interval:
Option<Duration>` (`None` = the substrate's 1 ms default,
ADR-023 §4 — the shipping default stands; the knob is wiring onto
`WatcherConfig::with_poll_interval`, which rejects zero) and
reader-pool size. `#[non_exhaustive]` policy does NOT apply (opts
structs consumers construct — ADR-017 §3's exemption; this is an
engine opts struct, same posture). Derive/`Default` it.
- `open(path: &str, opts: SqliteOpts) -> Result<Box<dyn Store>>` (or
a concrete `SqliteStore` re-exported as the trait — implementer's
choice, but the trait is what consumers hold). The path is a **plain
filesystem path** (ADR-023 §3 — `open_conn` already dropped the URI
flag; do not re-add it; `:memory:` works without URI mode).
- Connection architecture per engine-sqlite.md: one `Writer` slot
(substrate's `Writer`), a `Readers` pool (substrate's `Readers`),
and the `SharedUpdateWatcher` spawned against the db path — watcher
spawn is fallible (W-2): map the substrate's `Result<_, String>`
into `Error::Database` (the wave-2 review's watch-item; the mapping
is engine-layer work). Open fails if the watcher cannot start.
- Each connection runs the substrate's bootstrap at open
(`open_conn` → pragmas, `attach_notify`, `attach_alkstore_functions`,
`bootstrap_schema`) — the writer and each pooled reader.
- The store struct holds the shared machinery (writer slot, readers,
watcher handle, db path) behind the `Store` trait; `Drop`/close
stops the watcher and closes connections.
- `spawn_blocking` seam posture established here for all later tasks:
a private helper the trait methods use (tokio dep added to the
engine crate — dev-dep today; make it a real dep).
This task does NOT implement any trait method beyond what's needed to
compile (the trait impls come in the mechanism tasks) — but the store
struct, opts, and constructor are real and tested (open a temp file,
verify bootstrap ran, verify the watcher is up, verify opts flow).
## Acceptance Criteria
- [ ] `SqliteOpts` exists with `poll_interval: Option<Duration>`
(None = 1 ms default) and reader-pool sizing; documented
- [ ] `open(path, opts)` boots the full connection architecture:
writer slot, reader pool, watcher (fallible spawn mapped to
`Database`), substrate bootstrap on every connection
- [ ] Plain-path posture holds: no URI flag re-introduced anywhere in
the open path (pinning test: a `?name=value`-suffixed path
creates a file with that literal name — ADR-023 §3's backlog row)
- [ ] `poll_interval` flows to the watcher config (config-surface
accounting test — assert the watcher's configured cadence, not
wall-clock timing; ADR-023 §4's backlog row)
- [ ] Watcher death handling wired: subscribers' receivers close
(mechanism lands with `listen`, but the death-guard plumbing
exists here)
- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
## References
- docs/architecture/engine-sqlite.md (Connection architecture)
- docs/architecture/decisions/023-fourth-review-round.md §3, §4
- docs/architecture/decisions/008-contract-v1-pinning.md §6
- docs/architecture/decisions/007-transactional-seam.md (SQLite arm)
- alkstore-sqlite/src/substrate/schema.rs, watcher.rs
## Notes
> To be filled by implementation agent
## Summary
> To be filled on completion
+91
View File
@@ -0,0 +1,91 @@
---
id: sqlite-engine-queues
name: SQLite engine — queues (`Queue`, claims, job handles, backoff curve)
status: pending
depends_on: [sqlite-engine-seam-tx]
scope: moderate
risk: medium
impact: component
level: implementation
tags: [wave-3, sqlite-engine]
---
## Description
Implement `Store::queue` and the full `Queue` + `JobHandle` traits
over the substrate's re-derived queue ops (`enqueue`, `claim_batch`,
`ack`, `ack_batch`, `heartbeat`, `retry`, `fail`, `cancel`, `get_job`,
`sweep_expired` — wave 2's owned code):
- `queue(name, opts)` — validated constructor; the handle carries the
`QueueOpts` stamps future enqueues resolve over (ADR-010 §3a).
- `enqueue(payload, opts)` — the one handle-carrying shape: resolve
`EnqueueOpts` over **this handle's opts** (delay-over-`run_at`
precedence, relative-`expires` → absolute `expires_at`, per-call
`max_attempts` override — ADR-020 §1–§3); one clock read per
enqueue; `encode_payload` at the seam.
- `claim_one` / `claim_batch` — through the writer slot (claims are
writes); map the substrate's JSON row output into `Job` +
boxed `JobHandle`s. **Extent guard (ADR-023 §2)**: `n <= 0` yields
the empty `Vec` at the trait-impl entry, before the substrate —
`claim_batch(-1)` must never claim the whole queue (the `LIMIT -1`
dialect artifact dies here; the substrate passes `n` through
contract-blind). `worker_id` non-empty (`InvalidName`).
- `JobHandle` impl: `job()` (the row value), `ack`/`retry`/`fail`
(one-shot, `self: Box<Self>`), `heartbeat(extend)` (repeatable,
absolute reset). The uniform validity predicate is the substrate's;
refusals are `Ok(false)`, never errors.
- **The backoff curve — engine-side, this task**: `retry(err, None)`
computes the equal-jitter exponential from the job's stamped
`backoff_base_s`: `delay ∈ [base·2^(a−1)/2, base·2^(a−1)]`, capped
at 1 hour (ADR-010 §3's pinned formula — the engine owns one
implementation; wave 4 re-derives it identically and wave 5 pins
equivalence). Jitter needs a randomness source — use the substrate's
existing posture if one exists, else `std`-only (no `rand` dep
without a decision; a hash-of-id-and-time fallback is acceptable if
documented — flag the choice in Notes). `retry(err, Some(d))` passes
`d` through.
- `ack_batch(ids)` — per-id validity predicate, count returned.
- `cancel(job_id)`, `get_job(job_id)` (sees dead rows — map
`last_error`/`died_at`), `sweep_expired()` (no leader lock
required).
- Exhaustion dead-letters with the substrate's pinned default strings
(`"max attempts exceeded"`; `fail(None)` → `"failed"`).
## Acceptance Criteria
- [ ] Full `Queue` + `JobHandle` impls; enqueue→claim→ack lifecycle
works; stamps visible via `get_job` match the resolution rules
- [ ] Extent guard: `claim_batch(n <= 0)` claims nothing (empty
`Vec`, no error, never the whole queue) — pinned per ADR-023
§2's backlog row
- [ ] Exactly-once handout under concurrent claims (two workers, one
row — the POC-pinned property, engine-tested here)
- [ ] Backoff curve: `retry(err, None)` delays land in the pinned
range (statistical test over many retries, tolerance-bounded);
cap at 1 hour; `Some(d)` overrides
- [ ] Validity predicate: lapsed-deadline ops refuse with `Ok(false)`
(heartbeat, ack, retry, fail); reclaim consumes an attempt
- [ ] Dead-letter: exhaustion moves to dead with the pinned default
strings; `get_job` sees dead rows with `last_error`/`died_at`
- [ ] `sweep_expired` moves expired rows (both states) + applies
retention; returns the moved count
- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
## References
- docs/architecture/core-contract.md (queues section)
- docs/architecture/decisions/010-queue-semantics-depth.md §1–§5
- docs/architecture/decisions/019-mechanism-handle-surfaces.md §2, §3
- docs/architecture/decisions/020-enqueue-opt-semantics-and-bridges.md §1–§3
- docs/architecture/decisions/023-fourth-review-round.md §2
- docs/architecture/queues.md (Retry/backoff, Dead-letter)
- alkstore-sqlite/src/substrate/queue_ops.rs
## Notes
> To be filled by implementation agent
## Summary
> To be filled on completion
+102
View File
@@ -0,0 +1,102 @@
---
id: sqlite-engine-scheduler-outbox
name: SQLite engine — scheduler (`schedule`/`unschedule`/`run_schedules`) and outbox helper
status: pending
depends_on: [sqlite-engine-queues, sqlite-engine-locks]
scope: broad
risk: high
impact: component
level: implementation
tags: [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
> To be filled by implementation agent
## Summary
> To be filled on completion
+102
View File
@@ -0,0 +1,102 @@
---
id: sqlite-engine-seam-tx
name: SQLite engine — `spawn_blocking` seam, `begin_tx`, `TxHandle`, writer-slot lease
status: pending
depends_on: [sqlite-engine-open-opts]
scope: moderate
risk: high
impact: component
level: implementation
tags: [wave-3, sqlite-engine, seam]
---
## Description
Implement the transactional seam — the highest-risk piece of the
engine and the shape every other mechanism task rides:
- The `spawn_blocking` bridge helper: every engine call (auto-commit
and tx ops alike) round-trips a blocking thread; rusqlite
connections are not `Send`-across-await. Establish the pattern once
(error mapping substrate `rusqlite::Error` → `Error::Database` with
the source chain preserved — ADR-008 §5's opaque fallback) and reuse
it everywhere.
- `begin_tx`: acquire the writer slot (substrate `Writer::acquire`),
open `BEGIN IMMEDIATE`, return the boxed `TxHandle`. The handle is
the **writer-slot lease** (ADR-007): holding it across `await`
points parks the writer — that is the design, not a defect.
- `TxHandle` impl: every `*_tx` method routes through the same
connection via `spawn_blocking` (thread-affinity is handled by the
slot lease — the connection moves with the handle, ops serialize on
it):
- Writes: `enqueue_tx` (stamps the engine's derived plain-queue
defaults 300/3/5/none — no-handle-open shape, ADR-020 §3a note),
`publish_tx`, `publish_with_key_tx` (key validation per ADR-015
§2), `notify_tx` (substrate's `notify()` SQL function inside the
tx — delivered only at commit), `save_offset_tx` (monotone —
substrate's upsert already is), `outbox_enqueue_tx` (derives the
`__alkstore_outbox:{name}` backing queue name engine-side and
stamps the outbox's 60/5/5 set — ADR-014 §1).
- Reads (ADR-021 §1, read-your-own-writes — they run inside the
caller's tx): `get_job_tx`, `get_offset_tx`, `read_since_tx`,
`read_from_consumer_tx`.
- `commit(self: Box<Self>)`: `COMMIT` + release the writer slot.
- **Drop = rollback** (ADR-021 §4): the handle's `Drop` impl issues
`ROLLBACK` and releases the writer slot — panic-safe (the drop runs
across unwinds; no panic may escape drop). This is the property
N-5's wave-5 probe will exercise against a real engine for the
first time — get the Drop impl right here.
- Name validation at every entry point via core's `validation` helpers
(shared-namespace kinds vs consumer-local; ADR-008 §4 — the
commit-atomic path must not bypass validation).
- Payload encoding: call core's now-fallible `encode_payload` at the
seam; `?` the `Codec` error (ADR-023 §1 — the engines' first
consumption of the typed signature).
All eleven `*_tx` methods are implemented **for real** in this task —
the substrate surface is fully ported (wave 2), so nothing stubs. The
opts-resolution helpers this introduces (delay-over-`run_at`,
relative-`expires`, the derived-default stamp sets) are the shared
engine-side arithmetic the auto-commit mechanism tasks reuse — one
owner per formula (ADR-012 §2), written once here. The mechanism tasks
then add what the tx path deliberately lacks: the auto-commit handles,
the extent/duration guards, the backoff curve, the receiver bridges,
and the scheduler runner.
## Acceptance Criteria
- [ ] `begin_tx` acquires the writer slot + `BEGIN IMMEDIATE`;
concurrent `begin_tx` blocks (slot lease) — tested
- [ ] All eleven `*_tx` methods implemented with entry-point
validation; writes commit-atomic, reads see the tx's own writes;
opts resolution (delay-over-run_at, relative-expires, derived
default stamp sets 300/3/5/none and 60/5/5) pinned by tests
- [ ] `commit` commits + releases the slot; drop (without commit)
rolls back + releases the slot — no-ghosts tested for job rows,
stream events, notifications, and offset saves
- [ ] Drop is panic-safe (a panic unwinding through a scope holding
the handle still rolls back and releases the slot)
- [ ] `with_tx` (core's provided method) works end-to-end on this
engine: Ok ⇒ commit, Err ⇒ rollback — tested
- [ ] `encode_payload`'s `Codec` error propagates from enqueue/publish
paths (typed, not silent)
- [ ] Substrate errors map to `Error::Database` with source chains
- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
## References
- docs/architecture/decisions/007-transactional-seam.md
- docs/architecture/decisions/021-tx-reads-and-value-shape-fixes.md §1, §4
- docs/architecture/decisions/014-outbox-tx-enqueue.md §1
- docs/architecture/decisions/020-enqueue-opt-semantics-and-bridges.md §1–§4
- docs/architecture/decisions/023-fourth-review-round.md §1
- docs/architecture/engine-sqlite.md (Mapping the contract: begin_tx, handle ops, commit/rollback)
- alkstore/src/tx.rs (the trait to implement)
## Notes
> To be filled by implementation agent
## Summary
> To be filled on completion
+83
View File
@@ -0,0 +1,83 @@
---
id: sqlite-engine-streams
name: SQLite engine — streams (`StreamHandle`, reads, subscribe, trim)
status: pending
depends_on: [sqlite-engine-seam-tx]
scope: moderate
risk: medium
impact: component
level: implementation
tags: [wave-3, sqlite-engine]
---
## Description
Implement `Store::stream` and the full `StreamHandle` trait over the
substrate's stream ops (`stream_publish`, `stream_read_since`,
`stream_save_offset`, `stream_get_offset` — inherited near-verbatim;
the `topic` column carries the contract's `stream` field, the one name
delta, mapped at this layer):
- `stream(name)` — validated constructor returning the boxed handle.
- `publish` / `publish_with_key` — auto-commit publishes through the
writer slot; `encode_payload` at the seam (`Codec` propagates);
return the assigned offset.
- `read_since(offset, limit)` / `read_from_consumer(consumer, limit)`
— reader-pool reads, offset ASC (global FIFO); map rows into
`StreamEvent` via the core constructor task's `from_row`.
**Extent guard (ADR-023 §2)**: `limit <= 0` yields the empty `Vec`
at the trait-impl entry, before the substrate — the `LIMIT -1`
dialect artifact dies here (the substrate passes `limit` through
contract-blind; the guard is the engine's obligation).
- `save_offset(consumer, offset)` / `get_offset(consumer)` — monotone
saves (substrate's upsert already refuses regressions silently);
`get_offset` of an absent consumer = 0.
- `trim_to(horizon)` — `DELETE FROM __alkstore_stream WHERE topic = ?
AND offset <= ?` inside the writer-slot lease; boundary argument is
**total** (ADR-023 §2: a negative horizon deletes nothing,
idempotently — no guard, pass through); returns the deleted count.
- `subscribe(consumer)` — the durable receiver: attach, read to tail
from the stored offset, then wake-driven re-reads off the watcher
fanout (same bridge shape as `listen`, but yielding
`Result<StreamEvent>` per the `EventReceiver` trait — `Err` carries
`Database` only; `None` = closed on watcher death, terminal).
Receiver-form `save_offset()` saves the last-yielded event's offset
through the same monotone op (ADR-019 §6); `offset()` reports the
receiver's position.
- Consumer identifiers (`consumer`) take the non-empty rule only;
stream names the shared-namespace rules.
## Acceptance Criteria
- [ ] Full `StreamHandle` impl; publish/read/save/get/trim/subscribe
all work against a temp-file store
- [ ] Extent guard: `read_since`/`read_from_consumer` with
`limit <= 0` return the empty `Vec` (never the whole stream) —
pinned per ADR-023 §2's backlog row
- [ ] `trim_to` exact-boundary (`<=`), negative horizon deletes
nothing, surviving offsets never renumbered
- [ ] Monotone saves: regression saves are silent no-ops; direct and
receiver forms compose without regression
- [ ] `subscribe` replays from the stored offset, then delivers
post-attach events wake-driven; watcher death closes it
(`None`, terminal); `recv()`'s `Err` carries `Database` only
- [ ] Key round-trips exactly (`None` vs `Some`); `stream` field
carries the contract name (not `topic`)
- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
## References
- docs/architecture/core-contract.md (streams section)
- docs/architecture/decisions/015-streams-depth.md
- docs/architecture/decisions/019-mechanism-handle-surfaces.md §6
- docs/architecture/decisions/023-fourth-review-round.md §2
- docs/architecture/engine-sqlite.md (Mapping the contract: streams)
- alkstore-sqlite/src/substrate/ops.rs (stream_* fns)
## Notes
> To be filled by implementation agent
## Summary
> To be filled on completion