Fork re-derivation: queue ops on contract v1 (stamps, per-row claim visibility, savepoint-guarded dead-letter, both-states sweep, dead-visible get_job, @every scheduler) — ADR-010 §1–§5/§3a/§8, ADR-009 §2–§4, ADR-011/012, task fork-rederive-queue-ops
This commit is contained in:
1 parent
43a135c453
commit
915641bfe6
6 files changed
+2188
-23
No files matched your search
@@ -74,8 +74,8 @@ port-surfaced adaptations the plan didn't anticipate are appended
|
||||
|
||||
The entries below are the fork-port deltas as landed by
|
||||
`fork-port-connection-watcher` (paths finalized against the ported
|
||||
tree). The re-derivation entry (D-12) is still pre-declared —
|
||||
`fork-rederive-queue-ops` lands that half. `fork-provenance-and-floor`
|
||||
tree). D-12 landed as the re-derivation's `queue_ops.rs` module +
|
||||
schema extensions with `fork-rederive-queue-ops`. `fork-provenance-and-floor`
|
||||
verifies the full register at gate time. Cherry-picks: empty at
|
||||
scaffold, append-only forever.
|
||||
|
||||
@@ -92,7 +92,7 @@ scaffold, append-only forever.
|
||||
| D-09 | port delta | W-3 — `data_version` u32 wrap: recorded on `poll_data_version`, no action (a wrap fires one spurious wake; wakes are re-read hints) | watcher machinery (notes on `poll_data_version`; `watcher.rs`) | ADR-012 §4 (not actionable — recorded) | ours |
|
||||
| D-10 | port delta | table family renamed `_honker_*` → `__alkstore_*` across the ported storage surface; no online rename migration (fresh bootstrap only; leftover `_honker_*` orphans are inert and untouched — pinned by test) | bootstrap / schema + notify/stream/lock SQL (`schema.rs`, `ops.rs`) | ADR-011 (naming); ADR-010 §8 (pre-authorization); ADR-008 §4 (storage-internal); ADR-012 §5 | ours |
|
||||
| D-11 | port delta | bootstrap race swallow re-keyed — on `ALTER TABLE` duplicate-column failure, verify via `pragma_table_info` (present ⇒ benign race swallowed; absent ⇒ propagate); upstream's error-string matching not inherited | bootstrap / schema machinery (`column_present` / `add_column_if_absent`; `schema.rs`) | ADR-012 §5; quality-read §4 (schema brittleness) | ours |
|
||||
| D-12 | re-derivation | queue ops re-derived on contract v1 (new code, contract-derived names): enqueue + per-job option stamping; single-statement claim with per-row visibility from stamps; savepoint-guarded retry/fail/dead-letter; both-states no-stranded-rows `sweep_expired` with retention deletion; dead-visible `get_job`; scheduler tick with `@every` boundary math; stamp columns + `claimed_at` added to the schema | queue-op modules (new) + schema (stamp columns) | ADR-011 scope (re-derive); ADR-010 §3a/§5/§1; quality-read D-1–D-3, D-5–D-7 | ours |
|
||||
| D-12 | re-derivation | queue ops re-derived on contract v1 (new code, contract-derived names): enqueue + per-job option stamping (`Stamps` — max_attempts, visibility, backoff base, retention — resolved engine-side, applied per row); single-statement `claim_batch` with per-row visibility from the job's own stamps in the claim UPDATE (`claim_expires_at = unixepoch() + visibility_timeout_s`), ordering `priority DESC, run_at ASC, id ASC`, `attempts += 1` per claim, pre-claim dead-letter sweep for exhausted reclaimables; `ack`/`ack_batch`/`heartbeat`/`retry`/`fail` carry the uniform validity predicate (processing + unexpired claim deadline; refusal = 0/false, not error); all dead-letter moves (retry-at-budget, fail, pre-claim sweep, `sweep_expired`) savepoint-guarded via `in_savepoint` (the #133 defect class); both-states no-stranded-rows `sweep_expired` with `dead_letter_retention_s` deletion; dead-visible `get_job` returning the full stamp field list + `last_error`/`died_at`; `cancel` unconditional; default error strings (`"max attempts exceeded"`, `"expired"`) passed in, never owned here; scheduler register/tick/soonest/unregister over `__alkstore_scheduler_tasks` with `@every`-only boundary math (`parse_every_interval`, s/m/h/d), 64-cap catch-up per tick with skip-forward (`SCHEDULER_MAX_CATCHUP_FIRES`), fire enqueue + row advance in one caller transaction; stamp columns + `claimed_at` added to live/dead/scheduler-table schemas with append-column migrations | queue-op modules (`queue_ops.rs`) + schema (stamp columns, dead-table extensions, retention index) | ADR-011 scope (re-derive); ADR-010 §3a/§5/§1; quality-read D-1–D-3, D-5–D-7 | ours |
|
||||
| D-13 | hygiene | the substrate stays sync — family porting is the discipline deltas, not an asyncification port; the bridged seam at the engine layer is the async story | whole subtree | ADR-012 §4; ADR-003 (seam) | ours |
|
||||
| D-14 | hygiene | no panics in library code (see D-07 for the one substantive instance), no `unwrap()`/`expect()` outside tests | whole subtree | AGENTS.md code conventions; ADR-011 (family standard) | ours |
|
||||
| D-15 | hygiene | no comments in code (doc comments fine) — ported upstream comments elided | whole subtree | AGENTS.md code conventions; ADR-011 (family standard) | ours |
|
||||
|
||||
@@ -37,7 +37,8 @@
|
||||
//! of the fork, ported by `fork-port-connection-watcher`) lives under
|
||||
//! this module below — `schema.rs`, `watcher.rs`, `ops.rs` in
|
||||
//! upstream's `lib.rs` / `honker_ops.rs` order. The re-derived queue
|
||||
//! ops (`fork-rederive-queue-ops`) land beside them.
|
||||
//! ops (`fork-rederive-queue-ops`, register D-12) live in
|
||||
//! `queue_ops.rs` — owned contract-v1 code, not lineage body.
|
||||
//!
|
||||
//! Port state (wave 2): the machinery below is ported but not yet
|
||||
//! wired by the engine layer (wave 3 wires `open`, the trait impl, the
|
||||
@@ -48,6 +49,7 @@
|
||||
#![allow(unused_imports)]
|
||||
|
||||
mod ops;
|
||||
mod queue_ops;
|
||||
mod schema;
|
||||
mod watcher;
|
||||
|
||||
|
||||
@@ -62,7 +62,7 @@ fn to_sql_err<E: std::fmt::Display>(e: E) -> rusqlite::Error {
|
||||
rusqlite::Error::UserFunctionError(Box::new(std::io::Error::other(e.to_string())))
|
||||
}
|
||||
|
||||
fn in_savepoint<T>(
|
||||
pub(crate) fn in_savepoint<T>(
|
||||
conn: &Connection,
|
||||
name: &str,
|
||||
body: impl FnOnce() -> rusqlite::Result<T>,
|
||||
@@ -887,10 +887,49 @@ mod schema_migration_tests {
|
||||
.unwrap()
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.unwrap();
|
||||
assert_eq!(live_cols.len(), 13);
|
||||
assert_eq!(live_cols.len(), 16);
|
||||
assert!(live_cols.contains(&"expires_at".to_string()));
|
||||
assert!(live_cols.contains(&"claimed_at".to_string()));
|
||||
|
||||
let dead_cols: Vec<String> = conn
|
||||
.prepare("SELECT name FROM pragma_table_info('__alkstore_dead')")
|
||||
.unwrap()
|
||||
.query_map([], |r| r.get::<_, String>(0))
|
||||
.unwrap()
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.unwrap();
|
||||
for stamp in [
|
||||
"visibility_timeout_s",
|
||||
"backoff_base_s",
|
||||
"dead_letter_retention_s",
|
||||
"claimed_at",
|
||||
"expires_at",
|
||||
] {
|
||||
assert!(
|
||||
live_cols.contains(&stamp.to_string()) && dead_cols.contains(&stamp.to_string()),
|
||||
"{stamp} must be stamped on both live and dead rows"
|
||||
);
|
||||
}
|
||||
|
||||
let sched_cols: Vec<String> = conn
|
||||
.prepare("SELECT name FROM pragma_table_info('__alkstore_scheduler_tasks')")
|
||||
.unwrap()
|
||||
.query_map([], |r| r.get::<_, String>(0))
|
||||
.unwrap()
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.unwrap();
|
||||
for stamp in [
|
||||
"visibility_timeout_s",
|
||||
"backoff_base_s",
|
||||
"dead_letter_retention_s",
|
||||
"max_attempts",
|
||||
] {
|
||||
assert!(
|
||||
sched_cols.contains(&stamp.to_string()),
|
||||
"{stamp} must exist on the schedule row"
|
||||
);
|
||||
}
|
||||
|
||||
let sc_cols: Vec<String> = conn
|
||||
.prepare("SELECT name FROM pragma_table_info('__alkstore_stream_consumers')")
|
||||
.unwrap()
|
||||
@@ -1027,9 +1066,9 @@ mod schema_migration_tests {
|
||||
);
|
||||
assert_eq!(
|
||||
migrated_cols.last().map(String::as_str),
|
||||
Some("claimed_at"),
|
||||
"claimed_at is appended by ALTER TABLE, so the fresh CREATE TABLE \
|
||||
must put it last too"
|
||||
Some("dead_letter_retention_s"),
|
||||
"the stamp columns are appended by ALTER TABLE, so the fresh \
|
||||
CREATE TABLE must put them last too"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
File diff suppressed because it is too large.
Load diff
@@ -127,7 +127,10 @@ pub(crate) const BOOTSTRAP_ALKSTORE_SQL: &str = "
|
||||
max_attempts INTEGER NOT NULL DEFAULT 3,
|
||||
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||
expires_at INTEGER,
|
||||
claimed_at INTEGER
|
||||
claimed_at INTEGER,
|
||||
visibility_timeout_s INTEGER NOT NULL DEFAULT 300,
|
||||
backoff_base_s INTEGER NOT NULL DEFAULT 5,
|
||||
dead_letter_retention_s INTEGER
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS __alkstore_live_claim
|
||||
ON __alkstore_live(queue, priority DESC, run_at, id)
|
||||
@@ -148,8 +151,15 @@ pub(crate) const BOOTSTRAP_ALKSTORE_SQL: &str = "
|
||||
max_attempts INTEGER NOT NULL DEFAULT 0,
|
||||
last_error TEXT,
|
||||
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||
died_at INTEGER NOT NULL DEFAULT (unixepoch())
|
||||
died_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||
visibility_timeout_s INTEGER NOT NULL DEFAULT 300,
|
||||
backoff_base_s INTEGER NOT NULL DEFAULT 5,
|
||||
dead_letter_retention_s INTEGER,
|
||||
claimed_at INTEGER,
|
||||
expires_at INTEGER
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS __alkstore_dead_retention
|
||||
ON __alkstore_dead(queue, died_at);
|
||||
CREATE TABLE IF NOT EXISTS __alkstore_locks (
|
||||
name TEXT PRIMARY KEY,
|
||||
owner TEXT NOT NULL,
|
||||
@@ -164,7 +174,10 @@ pub(crate) const BOOTSTRAP_ALKSTORE_SQL: &str = "
|
||||
expires_s INTEGER,
|
||||
next_fire_at INTEGER NOT NULL,
|
||||
enabled INTEGER NOT NULL DEFAULT 1,
|
||||
max_attempts INTEGER NOT NULL DEFAULT 3
|
||||
max_attempts INTEGER NOT NULL DEFAULT 3,
|
||||
visibility_timeout_s INTEGER NOT NULL DEFAULT 300,
|
||||
backoff_base_s INTEGER NOT NULL DEFAULT 5,
|
||||
dead_letter_retention_s INTEGER
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS __alkstore_stream (
|
||||
offset INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
@@ -235,7 +248,63 @@ pub(crate) fn bootstrap_schema(conn: &Connection) -> Result<(), Error> {
|
||||
"max_attempts",
|
||||
"max_attempts INTEGER NOT NULL DEFAULT 3",
|
||||
)?;
|
||||
add_column_if_absent(
|
||||
conn,
|
||||
"__alkstore_scheduler_tasks",
|
||||
"visibility_timeout_s",
|
||||
"visibility_timeout_s INTEGER NOT NULL DEFAULT 300",
|
||||
)?;
|
||||
add_column_if_absent(
|
||||
conn,
|
||||
"__alkstore_scheduler_tasks",
|
||||
"backoff_base_s",
|
||||
"backoff_base_s INTEGER NOT NULL DEFAULT 5",
|
||||
)?;
|
||||
add_column_if_absent(
|
||||
conn,
|
||||
"__alkstore_scheduler_tasks",
|
||||
"dead_letter_retention_s",
|
||||
"dead_letter_retention_s INTEGER",
|
||||
)?;
|
||||
add_column_if_absent(conn, "__alkstore_live", "claimed_at", "claimed_at INTEGER")?;
|
||||
add_column_if_absent(
|
||||
conn,
|
||||
"__alkstore_live",
|
||||
"visibility_timeout_s",
|
||||
"visibility_timeout_s INTEGER NOT NULL DEFAULT 300",
|
||||
)?;
|
||||
add_column_if_absent(
|
||||
conn,
|
||||
"__alkstore_live",
|
||||
"backoff_base_s",
|
||||
"backoff_base_s INTEGER NOT NULL DEFAULT 5",
|
||||
)?;
|
||||
add_column_if_absent(
|
||||
conn,
|
||||
"__alkstore_live",
|
||||
"dead_letter_retention_s",
|
||||
"dead_letter_retention_s INTEGER",
|
||||
)?;
|
||||
add_column_if_absent(
|
||||
conn,
|
||||
"__alkstore_dead",
|
||||
"visibility_timeout_s",
|
||||
"visibility_timeout_s INTEGER NOT NULL DEFAULT 300",
|
||||
)?;
|
||||
add_column_if_absent(
|
||||
conn,
|
||||
"__alkstore_dead",
|
||||
"backoff_base_s",
|
||||
"backoff_base_s INTEGER NOT NULL DEFAULT 5",
|
||||
)?;
|
||||
add_column_if_absent(
|
||||
conn,
|
||||
"__alkstore_dead",
|
||||
"dead_letter_retention_s",
|
||||
"dead_letter_retention_s INTEGER",
|
||||
)?;
|
||||
add_column_if_absent(conn, "__alkstore_dead", "claimed_at", "claimed_at INTEGER")?;
|
||||
add_column_if_absent(conn, "__alkstore_dead", "expires_at", "expires_at INTEGER")?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
---
|
||||
id: fork-rederive-queue-ops
|
||||
name: Fork re-derivation — queue ops on contract v1 (stamps, claim, sweep, get_job)
|
||||
status: pending
|
||||
status: completed
|
||||
depends_on: [fork-port-connection-watcher]
|
||||
scope: broad
|
||||
risk: high
|
||||
@@ -76,22 +76,22 @@ Plus the inherited queue-machinery tests adapted where applicable.
|
||||
|
||||
## Acceptance Criteria
|
||||
|
||||
- [ ] All re-derived ops above present with contract-blind primitive
|
||||
- [x] All re-derived ops above present with contract-blind primitive
|
||||
APIs (no contract types, no error-taxonomy types in the
|
||||
substrate)
|
||||
- [ ] Savepoint hardening proven by test: a forced mid-flight error in
|
||||
- [x] Savepoint hardening proven by test: a forced mid-flight error in
|
||||
the dead-letter move strands no row in either table
|
||||
- [ ] No `.ok()`-style error swallows (the D-class defects); all
|
||||
- [x] No `.ok()`-style error swallows (the D-class defects); all
|
||||
substrate errors propagate typed
|
||||
- [ ] Claim exclusivity under concurrent claims (multi-connection
|
||||
- [x] Claim exclusivity under concurrent claims (multi-connection
|
||||
test); reclaim consumes an attempt; validity predicate refuses
|
||||
lapsed-deadline ops as values (false), not errors
|
||||
- [ ] `sweep_expired` moves both states; retention deletion enforced;
|
||||
- [x] `sweep_expired` moves both states; retention deletion enforced;
|
||||
no-stranded-rows property test green
|
||||
- [ ] Scheduler: boundary advance + fire atomic under row lock;
|
||||
- [x] Scheduler: boundary advance + fire atomic under row lock;
|
||||
64-cap catch-up; `@every` math unit-tested (s|m|h|d)
|
||||
- [ ] `get_job` sees dead rows with full stamp/error fields
|
||||
- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
|
||||
- [x] `get_job` sees dead rows with full stamp/error fields
|
||||
- [x] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
|
||||
|
||||
## References
|
||||
|
||||
@@ -104,8 +104,94 @@ Plus the inherited queue-machinery tests adapted where applicable.
|
||||
|
||||
## Notes
|
||||
|
||||
> To be filled by implementation agent
|
||||
- **Module shape**: the re-derivation is one new substrate module,
|
||||
`queue_ops.rs` (beside the ported `schema.rs` / `watcher.rs` /
|
||||
`ops.rs`), with contract-derived names throughout per ADR-012 §3 —
|
||||
the re-derived half is new code, so no lineage-name fidelity
|
||||
obligation applies. It reuses the ported `in_savepoint` machinery
|
||||
(made `pub(crate)` for the new module — the only touch to the kept
|
||||
half besides schema).
|
||||
- **API shape**: stamp values bundle into a `Stamps { max_attempts,
|
||||
visibility_timeout_s, backoff_base_s, dead_letter_retention_s }`
|
||||
struct (clippy arg-count floor of 7 forced the bundling; the bundle
|
||||
is also the honest shape — ADR-010 §3a's four stamps land together).
|
||||
Scheduler registration takes `FireOpts { priority, expires_s }` +
|
||||
`Stamps`. All resolution inputs (ready time as an absolute `run_at`,
|
||||
retry delay as a literal, the default error strings
|
||||
`"max attempts exceeded"` / `"expired"`) are engine-layer — the
|
||||
substrate takes already-resolved primitives. Cron machinery absent:
|
||||
`parse_every_interval` is the whole spec grammar (ADR-009 §2), and a
|
||||
non-`@every` spec found in storage is rejected (pinned by test).
|
||||
- **Claim mechanics**: single `UPDATE … WHERE id IN (WITH picked AS
|
||||
select-ordered-limit)` statement — per-row visibility via
|
||||
`claim_expires_at = unixepoch() + visibility_timeout_s` (the §3a
|
||||
bridge: the deadline computes from the row's own stamp, not a
|
||||
uniform call value). Ordering `priority DESC, run_at ASC, id ASC`
|
||||
in the picked subquery. The pre-claim dead-letter sweep for
|
||||
exhausted reclaimables is retained from the lineage (savepoint-
|
||||
guarded) with the engine-supplied exhaust string.
|
||||
- **Schema deltas**: stamp columns on `__alkstore_live`
|
||||
(`visibility_timeout_s`, `backoff_base_s`, `dead_letter_retention_s`)
|
||||
alongside the ported `claimed_at`; the dead table gains the same
|
||||
stamps plus `claimed_at`/`expires_at` (ADR-019 §3's dead-row field
|
||||
list); the scheduler table carries the stamps too (schedule fires
|
||||
enqueue through the same stamping path); a `__alkstore_dead(queue,
|
||||
died_at)` index serves the retention sweep. All append-column
|
||||
migrations, so an existing `__alkstore_*` file migrates in place.
|
||||
- **`retry` shape**: the retry branch is deliberately *not*
|
||||
savepoint-wrapped (it is a single UPDATE; there is no move), while
|
||||
the exhaustion branch is (DELETE→INSERT) — mirroring the lineage's
|
||||
post-#133 discipline, with the caller's resolved delay and the
|
||||
engine's exhaust string passed in.
|
||||
- **Boundary-advance transactionality**: `scheduler_tick` keeps the
|
||||
lineage's shape — enqueue + advance in the caller's (writer)
|
||||
transaction, so crash mid-tick rolls both back; per-tick row-lock
|
||||
re-check rides SQLite's writer serialization exactly as ADR-009 §4
|
||||
pins for this engine. The 64-cap (`SCHEDULER_MAX_CATCHUP_FIRES`) and
|
||||
skip-forward (next boundary strictly after `now`) are pinned by
|
||||
tests, including the never-doubles-fire re-tick check.
|
||||
- **Test floor**: the inherited queue-op suites adapted (savepoint
|
||||
rollback paths through the Rust API; the scalar-function-path
|
||||
variants are not re-created — the substrate's SQL-function surface
|
||||
for queues gets attached at wave-3 wiring, where those adapters
|
||||
reproduce mechanically if needed), plus the contract-property suite
|
||||
(41 new tests): claim exclusivity under 6 concurrent connections ×
|
||||
120 jobs (zero double-handout, integrity intact), stamps
|
||||
immutability across claim/heartbeat, lapsed-deadline refusal as
|
||||
values for all four handle ops, reclaim-eats-attempt, ack-vs-reclaim
|
||||
race resolution, both-states sweep + retention TTL, the
|
||||
no-stranded-rows property including the zombie shape, savepoint
|
||||
hardening at every dead-letter site via a forced mid-flight error
|
||||
and a blocked dead INSERT, boundary/catch-up/skip-forward
|
||||
scheduler math, `@every` s/m/h/d unit math, `get_job` dead
|
||||
visibility with the full stamp list.
|
||||
- **Register**: D-12 updated to the landed state (the only pre-declared
|
||||
entry this task owned); no new port-surfaced register entries arose —
|
||||
the deltas above are all within D-12's declared scope.
|
||||
|
||||
## Summary
|
||||
|
||||
> To be filled on completion
|
||||
The queue half of the fork's re-derivation landed as
|
||||
`alkstore-sqlite/src/substrate/queue_ops.rs` (~560 lines of owned
|
||||
contract-v1 code + a 41-test property suite; 84 sqlite-crate tests
|
||||
green total): contract-blind `enqueue` with per-job stamping
|
||||
(`Stamps`), single-statement `claim_batch` with per-row visibility
|
||||
from the job's own stamps + the retained pre-claim dead-letter sweep,
|
||||
`ack`/`ack_batch`/`heartbeat`/`retry`/`fail`/`cancel` under the
|
||||
uniform validity predicate (refusals as values), savepoint-guarded
|
||||
dead-letter moves at every DELETE→INSERT site (forced-error and
|
||||
blocked-INSERT tests pin no-strand-in-either-table), the both-states
|
||||
no-stranded-rows `sweep_expired` with `dead_letter_retention_s`
|
||||
enforcement, dead-visible `get_job` carrying the full ADR-019 §3 field
|
||||
list (stamps, `claimed_at`, `last_error`, `died_at`), and the
|
||||
scheduler surface (`register` upsert, `tick` with `@every` boundary
|
||||
math, 64-cap catch-up + skip-forward, `soonest`, `unregister`) with
|
||||
fire+advance atomic under the writer transaction. Schema: stamp
|
||||
columns + `claimed_at`/`expires_at` on live and dead, stamps on the
|
||||
schedule table, retention index on dead — all append-column
|
||||
migrations. No contract or taxonomy types in the substrate; the
|
||||
default error strings are call-site parameters. PROVENANCE.md D-12
|
||||
recorded as landed. Verified: `cargo build`, workspace `cargo test`
|
||||
(84 sqlite + 26 elsewhere), `cargo clippy --all-targets -- -D
|
||||
warnings`, `cargo fmt --check` all clean; zero tokio, no
|
||||
panics/unwrap/expect outside tests.
|
||||
Reference in new issue
Block a user