Fork port: connection architecture + watcher machinery into the substrate (ADR-011/012 §3–§5, task fork-port-connection-watcher)

Kept half of the honker-core fork lands in
alkstore-sqlite/src/substrate/ as schema.rs / watcher.rs / ops.rs
(register D-17): PRAGMA/WAL open posture + set_journal_mode_wal retry,
Writer, Readers, the polling watcher family (SharedUpdateWatcher,
WatcherDeathGuard, stat_identity dead-man's switch), in_savepoint/
UnwindUndo mutation discipline, REAL-coercion arg helpers, notify
scalar + notifications table with the ADR-010 §6 at-attach pruning
cap, stream functions, lock functions.

Port deltas (ADR-012 §4): W-1 bounded reconnect backoff
(MAX_RECONNECT_TICKS=100), W-2 fallible watcher spawn (Result; engine
maps to Database at open in wave 3), dead-man's-switch panic replaced
by log-and-exit through the ordinary death path — death still closes
every subscriber (pinned by test, join now Ok). Table family
_honker_* -> __alkstore_* (D-10); duplicate-column race swallow
re-keyed to pragma_table_info (D-11); scheduler cron_expr -> spec;
fresh-only bootstrap, append-column migrations kept, column-order
equality pinned. Drops confirmed absent: cron, kernel/shm backends,
rate-limit/result tables, superseded queue functions (D-01..D-04).
file-id retained for the kept dead-man's switch (D-18).

43 engine-crate tests green (adapted inherited suites + delta tests +
cross-mechanism pressure); cargo build/clippy -D warnings/fmt clean.
PROVENANCE.md register updated to the landed state (D-01..D-20).
This commit is contained in:
glm-5.3-flash committed 2026-10-08 03:25:48 +00:00
1 parent 2d855b9546
commit 43a135c453
8 files changed
+2751 -30

No files matched your search

Generated
+90 -3
View File
@@ -38,7 +38,11 @@ version = "0.1.0"
dependencies = [
"alkstore",
"alkstore-contract-suite",
"file-id",
"parking_lot",
"rusqlite",
"serde_json",
"thiserror",
]
[[package]]
@@ -220,6 +224,15 @@ version = "0.1.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7360491ce676a36bf9bb3c56c1aa791658183a54d2744120f27285738d90465a"
[[package]]
name = "file-id"
version = "0.2.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e1fc6a637b6dc58414714eddd9170ff187ecb0933d4c7024d1abbd23a3cc26e9"
dependencies = [
"windows-sys 0.60.2",
]
[[package]]
name = "find-msvc-tools"
version = "0.1.14"
@@ -414,7 +427,7 @@ checksum = "1788edb87fdc09c7e26304471e2f5be8cdefb1b6930d6e3985fc02ff53bf86ee"
dependencies = [
"libc",
"wasi 0.11.1+wasi-snapshot-preview1",
"windows-sys",
"windows-sys 0.61.2",
]
[[package]]
@@ -701,7 +714,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c3d1e2c7f27f8d4cb10542a02c49005dbd6e93095799d6f3be745fae9f8fedd4"
dependencies = [
"libc",
"windows-sys",
"windows-sys 0.61.2",
]
[[package]]
@@ -787,7 +800,7 @@ dependencies = [
"pin-project-lite",
"socket2",
"tokio-macros",
"windows-sys",
"windows-sys 0.61.2",
]
[[package]]
@@ -1018,6 +1031,15 @@ version = "0.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5"
[[package]]
name = "windows-sys"
version = "0.60.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f2f500e4d28234f72040990ec9d39e3a6b950f9f22d3dba18416c35882612bcb"
dependencies = [
"windows-targets",
]
[[package]]
name = "windows-sys"
version = "0.61.2"
@@ -1027,6 +1049,71 @@ dependencies = [
"windows-link",
]
[[package]]
name = "windows-targets"
version = "0.53.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4945f9f551b88e0d65f3db0bc25c33b8acea4d9e41163edf90dcd0b19f9069f3"
dependencies = [
"windows-link",
"windows_aarch64_gnullvm",
"windows_aarch64_msvc",
"windows_i686_gnu",
"windows_i686_gnullvm",
"windows_i686_msvc",
"windows_x86_64_gnu",
"windows_x86_64_gnullvm",
"windows_x86_64_msvc",
]
[[package]]
name = "windows_aarch64_gnullvm"
version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a9d8416fa8b42f5c947f8482c43e7d89e73a173cead56d044f6a56104a6d1b53"
[[package]]
name = "windows_aarch64_msvc"
version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b9d782e804c2f632e395708e99a94275910eb9100b2114651e04744e9b125006"
[[package]]
name = "windows_i686_gnu"
version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "960e6da069d81e09becb0ca57a65220ddff016ff2d6af6a223cf372a506593a3"
[[package]]
name = "windows_i686_gnullvm"
version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fa7359d10048f68ab8b09fa71c3daccfb0e9b559aed648a8f95469c27057180c"
[[package]]
name = "windows_i686_msvc"
version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1e7ac75179f18232fe9c285163565a57ef8d3c89254a30685b57d83a38d326c2"
[[package]]
name = "windows_x86_64_gnu"
version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9c3842cdd74a865a8066ab39c8a7a473c0778a3f29370b5fd6b4b9aa7df4a499"
[[package]]
name = "windows_x86_64_gnullvm"
version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0ffa179e2d07eee8ad8f57493436566c7cc30ac536a3379fdf008f47f6bb7ae1"
[[package]]
name = "windows_x86_64_msvc"
version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d6bbff5f0aada427a1e5a6da5f1f98158182f26556f345ac9e04d36d0ebed650"
[[package]]
name = "wit-bindgen"
version = "0.57.1"
+10 -1
View File
@@ -7,7 +7,16 @@ repository.workspace = true
[dependencies]
alkstore = { version = "0.1", path = "../alkstore" }
rusqlite = { version = "0.40", features = ["bundled"] }
parking_lot = "0.12"
rusqlite = { version = "0.40", features = ["bundled", "functions"] }
serde_json = "1"
thiserror = "2"
# file-id only supports unix and windows. Other targets (WASI, Redox,
# illumos, etc.) get the `(0, 0)` fallback in substrate::watcher's
# stat_identity.
[target.'cfg(any(unix, windows))'.dependencies]
file-id = "0.2"
[dev-dependencies]
alkstore-contract-suite = { path = "../alkstore-contract-suite" }
+23 -12
View File
@@ -67,10 +67,17 @@ contract-v1 code), `drop` (not ported), `hygiene`
The entries below are the initial entries fixed by the retained
records (ADR-011's scope register, ADR-012 §4–§5) as
**pre-declared checklist rows** — they name deltas the port tasks will
make. `fork-provenance-and-floor` completes/verifies them against the
actual ported tree at gate time (paths finalized, any port-surfaced
adaptation appended). Cherry-picks: empty at scaffold, append-only
forever.
make. At the port tasks' landing, entries are updated in place to the
landed state (paths finalized, what/why made precise); entries for
port-surfaced adaptations the plan didn't anticipate are appended
(D-17..D-20, landed with `fork-port-connection-watcher`).
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`
verifies the full register at gate time. Cherry-picks: empty at
scaffold, append-only forever.
| ID | Category | What | Where | Why | Lineage |
|----|----------|------|-------|-----|---------|
@@ -78,18 +85,22 @@ forever.
| D-02 | drop | experimental watcher backends (`kernel-watcher`, `shm-fast-path`) and their optional deps (`notify`, `memmap2`, `libc`, `file-id`) not ported | — | ADR-011 scope (drop); quality-read §6; cuts optional deps | ours |
| D-03 | drop | rate-limit and result tables (cut-flags) not ported | — | ADR-002 cut-flags; ADR-011 scope (drop) | ours |
| D-04 | drop | superseded queue functions not ported | — | ADR-011 scope (drop); the queue half is re-derived (see R-01) | ours |
| D-05 | port delta | W-1 — bounded reconnect backoff in the watcher loop (a vanished db file no longer produces ~1000 open attempts/sec) | watcher machinery (`run_poll_loop` region) | ADR-012 §4; quality-read W-1 | ours |
| D-06 | port delta | W-2 — watcher spawn becomes fallible; a store whose watcher could not start fails at open time | watcher machinery (`spawn_with_config` region) | ADR-012 §4; quality-read W-2 | ours |
| D-07 | port delta | dead-man's switch — the db-file-identity-change panic replaced by a deliberate watcher-fatal death (log the precise diagnostic, exit through the ordinary death path); `WatcherDeathGuard`'s death-closes-subscribers behavior unchanged | `stat_identity` / `WatcherDeathGuard` region | ADR-012 §4 (family no-panics standard; observable semantics unchanged) | ours |
| D-08 | port delta | W-4 — watcher connection opens RW: deliberate inheritance, keep RW | `run_poll_loop` region | ADR-012 §4 (keep-RW recorded; not actionable as a change — recorded where the read put it) | ours |
| D-09 | port delta | W-3 — `data_version` u32 wrap: recorded, no action | watcher machinery (notes) | ADR-012 §4 (not actionable — recorded) | ours |
| D-10 | port delta | table family renamed `_honker_*` → `__alkstore_*` across the storage surface; no online rename migration (fresh bootstrap only; leftover `_honker_*` orphans are inert and untouched) | bootstrap / schema + all queue-op SQL + notify table | 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 | ADR-012 §5; quality-read §4 (schema brittleness) | ours |
| D-05 | port delta | W-1 — bounded reconnect backoff in the watcher reconnect path (`MAX_RECONNECT_TICKS` = 100 poll ticks between attempts; the lineage reconnected on every tick) | watcher machinery (`run_poll_loop` region; `watcher.rs`) | ADR-012 §4; quality-read W-1 | ours |
| D-06 | port delta | W-2 — watcher spawn becomes fallible (`UpdateWatcher::spawn*` / `SharedUpdateWatcher::new*` return `Result`); the engine surfaces the failure at open time | watcher machinery (`spawn_with_config` / `spawn_thread` region; `watcher.rs`) | ADR-012 §4; quality-read W-2 | ours |
| D-07 | port delta | dead-man's switch — the db-file-identity-change panic replaced by a deliberate watcher-fatal death (log the precise diagnostic, exit through the ordinary death path); `WatcherDeathGuard`'s death-closes-subscribers behavior unchanged | `stat_identity` / `WatcherDeathGuard` region; `watcher.rs` (`run_poll_loop`, death guard) | ADR-012 §4 (family no-panics standard; observable semantics unchanged) | ours |
| D-08 | port delta | W-4 — watcher connection opens RW: deliberate inheritance, keep RW | `run_poll_loop` region; `watcher.rs` | ADR-012 §4 (keep-RW recorded; not actionable as a change — recorded where the read put it) | ours |
| 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-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 |
| D-16 | port delta | notify machinery — `notify()` scalar + notifications table kept with engine-internal hygiene, plus at-attach pruning cap implementing ADR-010 §6 | notify machinery (`attach_notify` region) | ADR-011 scope; ADR-010 §6 | ours |
| D-16 | port delta | notify machinery — `notify()` scalar + `__alkstore_notifications` table kept, upstream's no-pruning posture replaced by an at-attach pruning cap (`NOTIFY_ATTACH_MAX_ROWS` = 10k rows, oldest-first trim, implementing ADR-010 §6's engine-internal hygiene) | notify machinery (`attach_notify`; `schema.rs`) | ADR-011 scope; ADR-010 §6 | ours |
| D-17 | port delta | the lineage module structure (everything in `lib.rs` + `honker_ops.rs`) mapped to the substrate subtree's modules: `schema.rs` (pragmas/open, notify, bootstrap, `Writer`, `Readers`), `watcher.rs` (the polling watcher family), `ops.rs` (savepoint discipline, arg coercion, stream/lock ops, function attachments); upstream names preserved inside each | `schema.rs` / `watcher.rs` / `ops.rs` (subtree layout) | ADR-012 §3's fidelity posture applies to names and kept-body shape; the ADR-013 fold makes the module file the unit, and a ~4k-line flat file does not fit a bounded module subtree (AGENTS.md module-per-file); diff reviewability is preserved register-entry by register-entry against upstream's file+line region | ours |
| D-18 | port delta | `file-id` retained (target-gated `cfg(any(unix, windows))` dep, as upstream) — the register's D-02 drop-list names it among the *experimental backends'* optional deps, but it drives the kept dead-man's switch (`stat_identity`), so it stays | `watcher.rs` (`stat_identity`); engine manifest | the kept-half scope (quality-read §6 keeps the dead-man's switch verbatim modulo deltas) overrides D-02's enumeration for this dep; the experimental backends' true deps (`notify`, `memmap2`, `libc`) are cut | ours |
| D-19 | hygiene | scalar-function attachment sites renamed `attach_honker_functions` → `attach_alkstore_functions`, `honker_bootstrap` → `alkstore_bootstrap` (the surface carries the fork's identity in its function names; the coercion helpers keep upstream names `arg_i64`/`arg_opt_i64` per the fidelity posture) | `ops.rs` (`attach_alkstore_functions`) | ADR-011 (fork re-owning names); D-10's naming delta at the SQL-function layer | ours |
| D-20 | port delta | `WatcherConfig` slims to the polling backend's fields (`poll_interval` only; `WatcherBackend` and its parse/probe machinery dropped with D-02); the default poll interval (1 ms) and the 100 ms identity-check interval are inherited verbatim | `watcher.rs` (`WatcherConfig`) | ADR-012 §4 (polling is the only backend); ADR-012 §3 (kept values stay) | ours |
### Cherry-picks
+20 -5
View File
@@ -33,8 +33,23 @@
//! this directory) is the preservation obligation's carrier;
//! provenance register: [`PROVENANCE.md`](PROVENANCE.md).
//!
//! Scaffold state: this module is intentionally empty of machinery.
//! The port tasks (`fork-port-connection-watcher`,
//! `fork-rederive-queue-ops`) land the forked lineage and the
//! re-derived queue ops here; the register and the license notice are
//! already in place so the ports land into a finished frame.
//! Scaffold state: the connection/watcher machinery (the kept half
//! 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.
//!
//! 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
//! seam). A subset of the ported surface is therefore unreachable
//! until that wiring lands.
#![allow(dead_code)]
#![allow(unused_imports)]
mod ops;
mod schema;
mod watcher;
pub(crate) use schema::{attach_notify, bootstrap_schema, open_conn};
pub(crate) use watcher::SharedUpdateWatcher;
File diff suppressed because it is too large. Load diff
+534
View File
@@ -0,0 +1,534 @@
//! Connection open posture, notify machinery, bootstrap, and the
//! writer/reader connection architecture — ported from honker-core
//! `lib.rs` @ `f4e53c6` (upstream names preserved; the table family
//! renamed `_honker_*` → `__alkstore_*` per ADR-011; the bootstrap's
//! duplicate-column race swallow re-keyed to `pragma_table_info`
//! per ADR-012 §5; the at-attach notifications pruning cap per
//! ADR-010 §6 / register D-16).
use parking_lot::{Condvar, Mutex};
use rusqlite::functions::FunctionFlags;
use rusqlite::{Connection, OpenFlags};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
#[derive(thiserror::Error, Debug)]
pub(crate) enum Error {
#[error("Database error: {0}")]
Sqlite(#[from] rusqlite::Error),
}
pub(crate) const DEFAULT_PRAGMAS: &str = "PRAGMA busy_timeout = 5000;
PRAGMA synchronous = NORMAL;
PRAGMA foreign_keys = ON;
PRAGMA cache_size = -32000;
PRAGMA temp_store = MEMORY;
PRAGMA wal_autocheckpoint = 10000;";
pub(crate) fn apply_default_pragmas(conn: &Connection) -> rusqlite::Result<()> {
conn.execute_batch("PRAGMA busy_timeout = 5000;")?;
set_journal_mode_wal(conn)?;
conn.execute_batch(DEFAULT_PRAGMAS)
}
fn set_journal_mode_wal(conn: &Connection) -> rusqlite::Result<()> {
let already_wal = |conn: &Connection| -> rusqlite::Result<bool> {
let mode: String = conn.query_row("PRAGMA journal_mode", [], |row| row.get(0))?;
Ok(mode.eq_ignore_ascii_case("wal"))
};
if already_wal(conn)? {
return Ok(());
}
let deadline = std::time::Instant::now() + std::time::Duration::from_millis(5000);
let mut backoff = std::time::Duration::from_millis(1);
loop {
match conn.execute_batch("PRAGMA journal_mode = WAL;") {
Ok(()) => return Ok(()),
Err(e) => {
if already_wal(conn).unwrap_or(false) {
return Ok(());
}
if std::time::Instant::now() >= deadline {
return Err(e);
}
std::thread::sleep(backoff);
backoff = (backoff * 2).min(std::time::Duration::from_millis(50));
}
}
}
}
/// Upper bound on `__alkstore_notifications` rows before the next
/// attach prunes oldest-first past it (register D-16 — ADR-010 §6's
/// at-attach cap; engine-internal hygiene, not consumer-visible).
///
/// Not upstream. Upstream's `attach_notify` never prunes ("no magic
/// timer" — callers invoke `prune_notifications` explicitly); this
/// substrate's postures pins hygiene at attach so the engine grows
/// nothing unbounded by deferring to a caller who never prunes.
/// Rows are only trimmed while their senders' transactions are
/// committed — pruning runs on a fresh connection before its
/// `notify()` function is registered, so no in-flight wake row can
/// race the DELETE.
pub(crate) const NOTIFY_ATTACH_MAX_ROWS: i64 = 10_000;
/// Install the `__alkstore_notifications` table and the
/// `notify(channel, payload)` SQL scalar function on `conn`.
/// Idempotent.
pub(crate) fn attach_notify(conn: &Connection) -> Result<(), Error> {
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS __alkstore_notifications (
id INTEGER PRIMARY KEY AUTOINCREMENT,
channel TEXT NOT NULL,
payload TEXT NOT NULL,
created_at INTEGER NOT NULL DEFAULT (unixepoch())
);
CREATE INDEX IF NOT EXISTS __alkstore_notifications_recent
ON __alkstore_notifications(channel, id);",
)?;
conn.execute(
"DELETE FROM __alkstore_notifications
WHERE id IN (
SELECT id FROM __alkstore_notifications
ORDER BY id DESC
LIMIT -1 OFFSET ?1
)",
[NOTIFY_ATTACH_MAX_ROWS],
)?;
conn.create_scalar_function("notify", 2, FunctionFlags::SQLITE_UTF8, |ctx| {
let channel: String = ctx.get(0)?;
let payload: String = ctx.get(1)?;
let db = unsafe { ctx.get_connection() }?;
let mut ins = db.prepare_cached(
"INSERT INTO __alkstore_notifications (channel, payload) VALUES (?1, ?2)",
)?;
let id = ins.insert(rusqlite::params![channel, payload])?;
Ok(id)
})?;
Ok(())
}
pub(crate) const BOOTSTRAP_ALKSTORE_SQL: &str = "
CREATE TABLE IF NOT EXISTS __alkstore_live (
id INTEGER PRIMARY KEY AUTOINCREMENT,
queue TEXT NOT NULL,
payload TEXT NOT NULL,
state TEXT NOT NULL DEFAULT 'pending',
priority INTEGER NOT NULL DEFAULT 0,
run_at INTEGER NOT NULL DEFAULT (unixepoch()),
worker_id TEXT,
claim_expires_at INTEGER,
attempts INTEGER NOT NULL DEFAULT 0,
max_attempts INTEGER NOT NULL DEFAULT 3,
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
expires_at INTEGER,
claimed_at INTEGER
);
CREATE INDEX IF NOT EXISTS __alkstore_live_claim
ON __alkstore_live(queue, priority DESC, run_at, id)
WHERE state IN ('pending', 'processing');
CREATE INDEX IF NOT EXISTS __alkstore_live_pending_deadline
ON __alkstore_live(queue, run_at)
WHERE state = 'pending';
CREATE INDEX IF NOT EXISTS __alkstore_live_processing_deadline
ON __alkstore_live(queue, claim_expires_at)
WHERE state = 'processing';
CREATE TABLE IF NOT EXISTS __alkstore_dead (
id INTEGER PRIMARY KEY,
queue TEXT NOT NULL,
payload TEXT NOT NULL,
priority INTEGER NOT NULL DEFAULT 0,
run_at INTEGER NOT NULL DEFAULT 0,
attempts INTEGER NOT NULL DEFAULT 0,
max_attempts INTEGER NOT NULL DEFAULT 0,
last_error TEXT,
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
died_at INTEGER NOT NULL DEFAULT (unixepoch())
);
CREATE TABLE IF NOT EXISTS __alkstore_locks (
name TEXT PRIMARY KEY,
owner TEXT NOT NULL,
expires_at INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS __alkstore_scheduler_tasks (
name TEXT PRIMARY KEY,
queue TEXT NOT NULL,
spec TEXT NOT NULL,
payload TEXT NOT NULL,
priority INTEGER NOT NULL DEFAULT 0,
expires_s INTEGER,
next_fire_at INTEGER NOT NULL,
enabled INTEGER NOT NULL DEFAULT 1,
max_attempts INTEGER NOT NULL DEFAULT 3
);
CREATE TABLE IF NOT EXISTS __alkstore_stream (
offset INTEGER PRIMARY KEY AUTOINCREMENT,
topic TEXT NOT NULL,
key TEXT,
payload TEXT NOT NULL,
created_at INTEGER NOT NULL DEFAULT (unixepoch())
);
CREATE INDEX IF NOT EXISTS __alkstore_stream_topic
ON __alkstore_stream(topic, offset);
CREATE TABLE IF NOT EXISTS __alkstore_stream_consumers (
name TEXT NOT NULL,
topic TEXT NOT NULL,
offset INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (name, topic)
);
";
fn column_present(conn: &Connection, table: &str, column: &str) -> bool {
let mut stmt = match conn.prepare(&format!(
"SELECT 1 FROM pragma_table_info('{table}') WHERE name = ?1"
)) {
Ok(stmt) => stmt,
Err(_) => return false,
};
stmt.query_row([column], |_| Ok(true)).unwrap_or(false)
}
fn add_column_if_absent(
conn: &Connection,
table: &str,
column: &str,
ddl: &str,
) -> Result<(), Error> {
if column_present(conn, table, column) {
return Ok(());
}
match conn.execute(&format!("ALTER TABLE {table} ADD COLUMN {ddl}"), []) {
Ok(_) => Ok(()),
Err(e) => {
if column_present(conn, table, column) {
Ok(())
} else {
Err(e.into())
}
}
}
}
/// Install the alkstore queue schema on `conn`. Idempotent.
///
/// Fresh bootstrap only (ADR-012 §5): a file carrying honker's legacy
/// `_honker_*` tables gets inert orphans — never touched, never read,
/// never migrated. The migrations below exist for the same-schema
/// case: a database this bootstrap already created in an earlier
/// revision, missing a column appended since.
pub(crate) fn bootstrap_schema(conn: &Connection) -> Result<(), Error> {
conn.execute_batch(BOOTSTRAP_ALKSTORE_SQL)?;
add_column_if_absent(
conn,
"__alkstore_scheduler_tasks",
"enabled",
"enabled INTEGER NOT NULL DEFAULT 1",
)?;
add_column_if_absent(
conn,
"__alkstore_scheduler_tasks",
"max_attempts",
"max_attempts INTEGER NOT NULL DEFAULT 3",
)?;
add_column_if_absent(conn, "__alkstore_live", "claimed_at", "claimed_at INTEGER")?;
Ok(())
}
pub(crate) fn open_conn(path: &str, install_notify: bool) -> Result<Connection, Error> {
let conn = Connection::open_with_flags(
path,
OpenFlags::SQLITE_OPEN_READ_WRITE
| OpenFlags::SQLITE_OPEN_CREATE
| OpenFlags::SQLITE_OPEN_URI,
)?;
apply_default_pragmas(&conn)?;
if install_notify {
attach_notify(&conn)?;
}
Ok(conn)
}
pub(crate) struct Writer {
slot: Mutex<Option<Connection>>,
available: Condvar,
closed: AtomicBool,
}
impl Writer {
pub(crate) fn new(conn: Connection) -> Self {
Self {
slot: Mutex::new(Some(conn)),
available: Condvar::new(),
closed: AtomicBool::new(false),
}
}
pub(crate) fn acquire(&self) -> Option<Connection> {
let mut guard = self.slot.lock();
loop {
if self.closed.load(Ordering::Acquire) {
return None;
}
if let Some(c) = guard.take() {
return Some(c);
}
self.available.wait(&mut guard);
}
}
pub(crate) fn try_acquire(&self) -> Option<Connection> {
if self.closed.load(Ordering::Acquire) {
return None;
}
self.slot.lock().take()
}
pub(crate) fn release(&self, conn: Connection) {
if self.closed.load(Ordering::Acquire) {
return;
}
let mut guard = self.slot.lock();
*guard = Some(conn);
self.available.notify_one();
}
pub(crate) fn close(&self) {
self.closed.store(true, Ordering::Release);
let mut guard = self.slot.lock();
guard.take();
self.available.notify_all();
}
}
pub(crate) struct Readers {
pool: Mutex<Vec<Connection>>,
outstanding: Mutex<usize>,
available: Condvar,
path: String,
max: usize,
closed: AtomicBool,
}
impl Readers {
pub(crate) fn new(path: String, max: usize) -> Self {
Self {
pool: Mutex::new(Vec::new()),
outstanding: Mutex::new(0),
available: Condvar::new(),
path,
max: max.max(1),
closed: AtomicBool::new(false),
}
}
pub(crate) fn acquire(&self) -> Result<Connection, Error> {
loop {
if self.closed.load(Ordering::Acquire) {
return Err(closed_err());
}
let mut pool = self.pool.lock();
if let Some(c) = pool.pop() {
return Ok(c);
}
let mut out = self.outstanding.lock();
if *out < self.max {
*out += 1;
drop(out);
drop(pool);
match open_conn(&self.path, false) {
Ok(conn) => {
if self.closed.load(Ordering::Acquire) {
*self.outstanding.lock() -= 1;
drop(conn);
return Err(closed_err());
}
return Ok(conn);
}
Err(e) => {
*self.outstanding.lock() -= 1;
return Err(e);
}
}
}
drop(out);
self.available.wait(&mut pool);
}
}
pub(crate) fn release(&self, conn: Connection) {
if self.closed.load(Ordering::Acquire) {
return;
}
let mut pool = self.pool.lock();
pool.push(conn);
self.available.notify_one();
}
pub(crate) fn close(&self) {
self.closed.store(true, Ordering::Release);
self.pool.lock().clear();
self.available.notify_all();
}
}
fn closed_err() -> Error {
Error::Sqlite(rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_MISUSE),
Some("Database is closed".to_string()),
))
}
#[cfg(test)]
mod writer_reader_tests {
use super::*;
fn mem() -> Connection {
Connection::open_in_memory().unwrap()
}
fn temp_db(name: &str) -> std::path::PathBuf {
let p = std::env::temp_dir().join(format!(
"alkstore-{name}-{}-{:?}.db",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_file(&p);
let _ = std::fs::remove_file(format!("{}-wal", p.display()));
let _ = std::fs::remove_file(format!("{}-shm", p.display()));
p
}
#[test]
fn writer_try_acquire_returns_none_when_held() {
let w = Writer::new(mem());
let conn = w.acquire().unwrap();
assert!(w.try_acquire().is_none());
w.release(conn);
assert!(w.try_acquire().is_some());
}
#[test]
fn writer_close_drops_idle_connection() {
let w = Writer::new(mem());
w.close();
assert!(w.acquire().is_none());
assert!(w.try_acquire().is_none());
}
#[test]
fn writer_close_drops_returned_connection() {
let w = Writer::new(mem());
let conn = w.acquire().unwrap();
w.close();
w.release(conn);
assert!(w.try_acquire().is_none());
}
#[test]
fn readers_close_returns_closed_err() {
let tmp = temp_db("readers-close");
let _ = std::fs::remove_file(&tmp);
let c = Connection::open(&tmp).unwrap();
c.execute_batch("PRAGMA journal_mode=WAL;").unwrap();
drop(c);
let r = Readers::new(tmp.to_string_lossy().into_owned(), 4);
let c = r.acquire().unwrap();
r.release(c);
r.close();
match r.acquire() {
Err(Error::Sqlite(rusqlite::Error::SqliteFailure(_, Some(msg)))) => {
assert!(msg.contains("Database is closed"));
}
other => panic!("expected closed err, got {other:?}"),
}
let _ = std::fs::remove_file(&tmp);
}
#[test]
fn readers_open_failure_does_not_leak_capacity() {
let r = Arc::new(Readers::new(
"/this/parent/does/not/exist/alkstore-readers-leak.db".into(),
2,
));
for i in 0..5 {
let r = r.clone();
let handle = std::thread::spawn(move || r.acquire());
let result = handle.join().unwrap();
assert!(
result.is_err(),
"attempt {i}: expected open failure, got Ok"
);
}
}
#[test]
fn open_conn_applies_wal_and_default_pragmas() {
let tmp = temp_db("open-pragmas");
let conn = open_conn(tmp.to_str().unwrap(), false).unwrap();
let mode: String = conn
.pragma_query_value(None, "journal_mode", |r| r.get(0))
.unwrap();
assert!(mode.eq_ignore_ascii_case("wal"));
let busy: i64 = conn
.pragma_query_value(None, "busy_timeout", |r| r.get(0))
.unwrap();
assert_eq!(busy, 5000);
drop(conn);
let _ = std::fs::remove_file(&tmp);
let _ = std::fs::remove_file(format!("{}-wal", tmp.display()));
let _ = std::fs::remove_file(format!("{}-shm", tmp.display()));
}
#[test]
fn concurrent_opens_never_return_database_is_locked() {
const ROUNDS: usize = 8;
const OPENERS: usize = 16;
for round in 0..ROUNDS {
let dir = std::env::temp_dir().join(format!(
"alkstore-open-race-{}-{round}-{:?}",
std::process::id(),
std::thread::current().id()
));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("pressure.db");
let barrier = Arc::new(std::sync::Barrier::new(OPENERS));
let handles: Vec<_> = (0..OPENERS)
.map(|_| {
let path = path.clone();
let barrier = Arc::clone(&barrier);
std::thread::spawn(move || {
let conn = Connection::open(&path)?;
barrier.wait();
apply_default_pragmas(&conn)?;
let mode: String =
conn.query_row("PRAGMA journal_mode", [], |row| row.get(0))?;
assert!(
mode.eq_ignore_ascii_case("wal"),
"expected WAL after open, got {mode}"
);
Ok::<_, rusqlite::Error>(())
})
})
.collect();
let failures: Vec<String> = handles
.into_iter()
.filter_map(|h| h.join().unwrap().err().map(|e| e.to_string()))
.collect();
std::fs::remove_dir_all(&dir).ok();
assert!(
failures.is_empty(),
"round {round}: concurrent opens failed: {failures:?}"
);
}
}
}
+888
View File
@@ -0,0 +1,888 @@
//! The polling update watcher — ported from honker-core `lib.rs`
//! @ `f4e53c6` (upstream names preserved). Port deltas (ADR-012 §4):
//! W-1 bounded reconnect backoff in the reconnect path; W-2 fallible
//! spawn (`Result`); the dead-man's switch's db-file-identity-change
//! panic replaced by a deliberate watcher-fatal death that logs the
//! precise diagnostic and exits through the ordinary death path —
//! `WatcherDeathGuard` still closes every subscriber on death.
//! The experimental backend machinery (`kernel_watcher`, `shm_watcher`,
//! their deps) is dropped — the polling backend is the only backend.
use parking_lot::Mutex;
use rusqlite::{Connection, OpenFlags, ffi};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::mpsc::{SyncSender, TrySendError};
use std::time::{Duration, Instant};
pub(crate) const DEFAULT_WATCHER_POLL_INTERVAL: Duration = Duration::from_millis(1);
#[derive(Debug, Clone)]
pub(crate) struct WatcherConfig {
pub(crate) poll_interval: Duration,
}
impl Default for WatcherConfig {
fn default() -> Self {
Self {
poll_interval: DEFAULT_WATCHER_POLL_INTERVAL,
}
}
}
impl WatcherConfig {
pub(crate) fn with_poll_interval(mut self, poll_interval: Duration) -> Result<Self, String> {
if poll_interval.is_zero() {
return Err("watcher poll interval must be positive".to_string());
}
self.poll_interval = poll_interval;
Ok(self)
}
}
/// Platform-specific file identity: `(dev, ino)` on Unix,
/// `(volume_serial, file_index)` on Windows. Used to detect when the
/// database file has been replaced underneath us (atomic rename,
/// litestream restore, volume remount).
///
/// Uses the `file-id` crate on unix and windows for stable Rust
/// support without nightly features. Falls back to `(0, 0)` on other
/// targets — on those targets the dead-man's switch is a no-op
/// (every `stat_identity` returns `(0, 0)` so the equality check
/// never trips); replacement detection is disabled but the watcher
/// still functions.
#[cfg(any(unix, windows))]
pub(crate) fn stat_identity(path: &Path) -> std::io::Result<(u64, u64)> {
let id = file_id::get_file_id(path)?;
match id {
file_id::FileId::Inode {
device_id,
inode_number,
} => Ok((device_id, inode_number)),
file_id::FileId::LowRes {
volume_serial_number,
file_index,
} => Ok((u64::from(volume_serial_number), file_index)),
file_id::FileId::HighRes {
volume_serial_number,
file_id,
} => Ok(fold_high_res(volume_serial_number, file_id)),
}
}
#[cfg(any(unix, windows))]
fn fold_high_res(volume_serial_number: u64, file_id: u128) -> (u64, u64) {
let file_index = ((file_id >> 64) as u64) ^ (file_id as u64);
(volume_serial_number, file_index)
}
#[cfg(not(any(unix, windows)))]
pub(crate) fn stat_identity(_path: &Path) -> std::io::Result<(u64, u64)> {
Ok((0, 0))
}
/// Read the pager's `data_version` counter via `PRAGMA data_version`.
/// Returns a u32 incremented on every commit by any connection (and on
/// checkpoint).
///
/// W-3 (register D-09 — recorded, no action): the counter is a u32 and
/// wraps after 2³² commits. A wrap reads as a version change, firing
/// one spurious wake; wakes are re-read hints, so a spurious wake is
/// harmless. No wrap tracking is added.
pub(crate) fn poll_data_version(conn: &Connection) -> Result<u32, String> {
conn.pragma_query_value(None, "data_version", |row| row.get(0))
.map_err(|e| e.to_string())
}
fn is_transient_lock_error(e: &rusqlite::Error) -> bool {
matches!(
e,
rusqlite::Error::SqliteFailure(
ffi::Error {
code: ffi::ErrorCode::DatabaseBusy | ffi::ErrorCode::DatabaseLocked,
..
},
_,
)
)
}
/// Upper bound (in poll ticks) between reconnect attempts when the
/// watcher's connection is down (W-1 — ADR-012 §4).
const MAX_RECONNECT_TICKS: u64 = 100;
/// Open flags shared by the watcher's initial open and its reconnect
/// opens; test instrumentation hooks the loop's opens here.
type ReconnectOpenFn = fn(&Path) -> rusqlite::Result<Connection>;
fn open_watcher_conn(path: &Path) -> rusqlite::Result<Connection> {
Connection::open_with_flags(
path,
OpenFlags::SQLITE_OPEN_READ_WRITE | OpenFlags::SQLITE_OPEN_NO_MUTEX,
)
}
/// Polling loop body shared by [`UpdateWatcher`].
///
/// Three layers:
///
/// 1. **Fast path:** `PRAGMA data_version`. Compare the
/// integer to last seen value. Notify on change.
/// 2. **Error recovery:** If the query fails, reconnect the SQLite
/// connection and force one wake — with bounded backoff so a
/// vanished db file produces bounded open attempts (W-1), not the
/// lineage's one attempt per millisecond.
/// 3. **Identity check (about every 100 ms):** `stat(db_path)` to
/// compare `(dev, ino)`. If the file was replaced, log the precise
/// diagnostic and exit through the ordinary death path (the fork's
/// dead-man's-switch delta — the lineage panicked here); continuing
/// would silently watch stale data. `WatcherDeathGuard` closes
/// every subscriber exactly as it does for any watcher death.
pub(crate) fn run_poll_loop<F>(
db_path: PathBuf,
on_change: F,
stop: Arc<AtomicBool>,
ready: std::sync::mpsc::SyncSender<()>,
poll_interval: Duration,
open_conn_fn: ReconnectOpenFn,
) where
F: Fn(),
{
let mut conn = match open_conn_fn(&db_path) {
Ok(c) => Some(c),
Err(e) => {
eprintln!("alkstore: failed to open watcher connection: {e}");
None
}
};
let mut last_version = conn
.as_ref()
.and_then(|c| poll_data_version(c).ok())
.unwrap_or(0);
let initial_identity = match stat_identity(&db_path) {
Ok(id) => id,
Err(e) => {
eprintln!("alkstore: failed to stat database for identity check: {e}");
(0, 0)
}
};
let mut next_identity_check = Instant::now() + UPDATE_WATCHER_IDENTITY_INTERVAL;
let mut reconnect_ticks: u64 = MAX_RECONNECT_TICKS;
let _ = ready.send(());
drop(ready);
while !stop.load(Ordering::Acquire) {
std::thread::sleep(poll_interval);
if let Some(ref c) = conn {
match c.pragma_query_value(None, "data_version", |row| row.get::<_, u32>(0)) {
Ok(version) => {
if version != last_version {
last_version = version;
on_change();
}
}
Err(e) if is_transient_lock_error(&e) => {}
Err(e) => {
eprintln!("alkstore: data_version poll failed: {e}");
conn = None;
on_change();
}
}
} else {
reconnect_ticks = reconnect_ticks.saturating_add(1);
if reconnect_ticks >= MAX_RECONNECT_TICKS {
reconnect_ticks = 0;
match open_conn_fn(&db_path) {
Ok(c) => {
last_version = poll_data_version(&c).unwrap_or(0);
conn = Some(c);
}
Err(e) => {
eprintln!("alkstore: reconnect failed: {e}");
}
}
}
}
let now = Instant::now();
if now >= next_identity_check {
next_identity_check = now + UPDATE_WATCHER_IDENTITY_INTERVAL;
match stat_identity(&db_path) {
Ok(current) => {
if current != initial_identity {
eprintln!(
"alkstore: database file replaced: \
expected (dev={}, ino={}), \
found (dev={}, ino={}) at {:?}. \
The watcher cannot recover; it has stopped and \
closed all subscribers; close the store and reopen.",
initial_identity.0, initial_identity.1, current.0, current.1, db_path
);
return;
}
}
Err(e) => {
eprintln!("alkstore: stat identity check failed: {e}");
conn = None;
on_change();
}
}
}
}
}
pub(crate) struct UpdateWatcher {
stop: Arc<AtomicBool>,
handle: Option<std::thread::JoinHandle<()>>,
}
const UPDATE_WATCHER_IDENTITY_INTERVAL: Duration = Duration::from_millis(100);
impl UpdateWatcher {
/// Spawn a watcher thread on `db_path`. `on_change` is called once
/// per observed commit. The thread runs until [`UpdateWatcher`] is
/// dropped or [`stop`](Self::stop) is called.
///
/// Fallible (W-2 — ADR-012 §4): errors if the thread cannot be
/// spawned; the engine surfaces the failure at store-open time.
/// Baseline-capture failures inside the thread (open failures, a
/// vanished db file) leave the sender dropped and the loop itself
/// retrying per W-1's bounded backoff — not a spawn failure.
pub(crate) fn spawn<F>(db_path: PathBuf, on_change: F) -> Result<Self, String>
where
F: Fn() + Send + 'static,
{
Self::spawn_with_config(db_path, on_change, WatcherConfig::default())
}
/// Like [`spawn`](Self::spawn) but with an explicit poll cadence.
pub(crate) fn spawn_with_config<F>(
db_path: PathBuf,
on_change: F,
config: WatcherConfig,
) -> Result<Self, String>
where
F: Fn() + Send + 'static,
{
let stop = Arc::new(AtomicBool::new(false));
let stop_t = stop.clone();
let (ready_tx, ready_rx) = std::sync::mpsc::sync_channel::<()>(1);
let handle = Self::spawn_thread(
db_path,
on_change,
stop_t,
ready_tx,
config.poll_interval,
None,
)?;
let _ = ready_rx.recv();
Ok(Self {
stop,
handle: Some(handle),
})
}
/// The thread-spawn step, isolated so the fallible boundary (W-2)
/// has one owned site. `stack_size` overrides the std default;
/// production passes None — tests force a real spawn failure
/// through it.
fn spawn_thread<F>(
db_path: PathBuf,
on_change: F,
stop_t: Arc<AtomicBool>,
ready_tx: std::sync::mpsc::SyncSender<()>,
poll_interval: Duration,
stack_size: Option<usize>,
) -> Result<std::thread::JoinHandle<()>, String>
where
F: Fn() + Send + 'static,
{
let builder = std::thread::Builder::new().name("alkstore-update-poll".into());
let builder = match stack_size {
Some(sz) => builder.stack_size(sz),
None => builder,
};
builder
.spawn(move || {
run_poll_loop(
db_path,
on_change,
stop_t,
ready_tx,
poll_interval,
open_watcher_conn,
)
})
.map_err(|e| format!("failed to spawn alkstore update-watcher thread: {e}"))
}
pub(crate) fn stop(&self) {
self.stop.store(true, Ordering::Release);
}
/// Stop the watcher and wait for the thread to exit.
pub(crate) fn join(mut self) -> std::thread::Result<()> {
self.stop();
match self.handle.take() {
Some(h) => h.join(),
None => Ok(()),
}
}
}
impl Drop for UpdateWatcher {
fn drop(&mut self) {
self.stop();
}
}
struct WatcherDeathGuard {
senders: Arc<Mutex<HashMap<u64, SyncSender<()>>>>,
}
impl Drop for WatcherDeathGuard {
fn drop(&mut self) {
self.senders.lock().clear();
}
}
pub(crate) struct SharedUpdateWatcher {
watcher: Mutex<Option<UpdateWatcher>>,
senders: Arc<Mutex<HashMap<u64, SyncSender<()>>>>,
next_id: AtomicU64,
}
impl SharedUpdateWatcher {
/// Spawn the shared poll thread for `db_path`.
///
/// Fallible (W-2 — ADR-012 §4): the engine surfaces the failure at
/// store-open time.
pub(crate) fn new(db_path: PathBuf) -> Result<Self, String> {
Self::new_with_config(db_path, WatcherConfig::default())
}
/// Like [`new`](Self::new) but with an explicit poll cadence.
pub(crate) fn new_with_config(db_path: PathBuf, config: WatcherConfig) -> Result<Self, String> {
let senders: Arc<Mutex<HashMap<u64, SyncSender<()>>>> =
Arc::new(Mutex::new(HashMap::new()));
let senders_t = senders.clone();
let death_guard = WatcherDeathGuard {
senders: senders.clone(),
};
let watcher = UpdateWatcher::spawn_with_config(
db_path,
move || {
let _ = &death_guard;
let mut list = senders_t.lock();
list.retain(|_id, s| match s.try_send(()) {
Ok(()) | Err(TrySendError::Full(_)) => true,
Err(TrySendError::Disconnected(_)) => false,
});
},
config,
)?;
Ok(Self {
watcher: Mutex::new(Some(watcher)),
senders,
next_id: AtomicU64::new(0),
})
}
/// Subscribe. Returns a subscriber id and a [`Receiver<()>`] that
/// sees one tick per observed database update. Callers MUST
/// [`unsubscribe`](Self::unsubscribe) the returned id when done —
/// otherwise the sender stays in the map and a bridge thread
/// blocking on `recv()` will never see a disconnect.
///
/// Channel capacity is 1: bursts coalesce into one wake per drain
/// cycle. Wakes are "go re-read state" signals.
pub(crate) fn subscribe(&self) -> (u64, std::sync::mpsc::Receiver<()>) {
let id = self.next_id.fetch_add(1, Ordering::Relaxed);
let (tx, rx) = std::sync::mpsc::sync_channel(1);
self.senders.lock().insert(id, tx);
(id, rx)
}
/// Remove a subscriber. The corresponding receiver sees
/// `Err(RecvError)` on its next blocking `recv()`, letting a
/// bridge thread exit cleanly.
pub(crate) fn unsubscribe(&self, id: u64) {
self.senders.lock().remove(&id);
}
pub(crate) fn subscriber_count(&self) -> usize {
self.senders.lock().len()
}
/// Disconnect all subscribers and synchronously join the poll
/// thread. The thread owns the watcher's connection; joining drops
/// that connection and releases the file handle. Idempotent — safe
/// to call more than once.
pub(crate) fn close(&self) -> std::thread::Result<()> {
self.senders.lock().clear();
match self.watcher.lock().take() {
Some(watcher) => watcher.join(),
None => Ok(()),
}
}
}
impl Drop for SharedUpdateWatcher {
fn drop(&mut self) {
self.senders.lock().clear();
if let Some(watcher) = self.watcher.get_mut().take() {
drop(watcher);
}
}
}
#[cfg(test)]
mod watcher_tests {
use super::*;
fn temp_db(name: &str) -> PathBuf {
let p = std::env::temp_dir().join(format!(
"alkstore-{name}-{}-{:?}.db",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_file(&p);
let _ = std::fs::remove_file(format!("{}-wal", p.display()));
let _ = std::fs::remove_file(format!("{}-shm", p.display()));
p
}
fn wal_db(path: &Path) {
let conn = Connection::open(path).unwrap();
conn.execute_batch("PRAGMA journal_mode = WAL;").unwrap();
}
/// W-2: a thread-spawn failure surfaces as `Err` from the spawn
/// path — rather than the lineage's `expect` panic — with the
/// mapped error naming the watcher. Driven through the real spawn
/// site with an oversized stack reservation, which reliably fails
/// on every target.
#[test]
fn spawn_fails_when_the_thread_cannot_start() {
let (tx, _rx) = std::sync::mpsc::sync_channel::<()>(1);
let err = UpdateWatcher::spawn_thread(
PathBuf::from("/nonexistent/alkstore-w2.db"),
|| {},
Arc::new(AtomicBool::new(false)),
tx,
Duration::from_millis(1),
Some(usize::MAX),
)
.expect_err("an oversized stack reservation must fail to spawn");
assert!(
err.contains("failed to spawn alkstore update-watcher thread"),
"the W-2 error must name the watcher spawn: {err}"
);
}
#[test]
fn watcher_config_rejects_zero_poll_interval() {
let err = WatcherConfig::default()
.with_poll_interval(Duration::from_millis(0))
.unwrap_err();
assert_eq!(err, "watcher poll interval must be positive");
assert!(
WatcherConfig::default()
.with_poll_interval(Duration::from_millis(5))
.is_ok()
);
}
#[test]
fn shared_update_watcher_fans_out_to_many_subscribers() {
let tmp = temp_db("shared-fanout");
wal_db(&tmp);
let shared = SharedUpdateWatcher::new(tmp.clone()).unwrap();
let subs: Vec<(u64, std::sync::mpsc::Receiver<()>)> =
(0..50).map(|_| shared.subscribe()).collect();
let writer = Connection::open(&tmp).unwrap();
writer
.execute(
"CREATE TABLE IF NOT EXISTS _test_trigger(id INTEGER PRIMARY KEY)",
[],
)
.unwrap();
for i in 0..5 {
std::thread::sleep(Duration::from_millis(5));
writer
.execute("INSERT INTO _test_trigger(id) VALUES (?)", [i])
.unwrap();
}
std::thread::sleep(Duration::from_millis(50));
for (i, (_id, rx)) in subs.iter().enumerate() {
let mut got_any = false;
while rx.try_recv().is_ok() {
got_any = true;
}
assert!(got_any, "subscriber {i} saw no ticks");
}
shared.close().unwrap();
let _ = std::fs::remove_file(&tmp);
let _ = std::fs::remove_file(format!("{}-wal", tmp.display()));
let _ = std::fs::remove_file(format!("{}-shm", tmp.display()));
}
#[test]
fn shared_update_watcher_explicit_unsubscribe_disconnects_receiver() {
let tmp = temp_db("unsub");
wal_db(&tmp);
let shared = SharedUpdateWatcher::new(tmp.clone()).unwrap();
let (id, rx) = shared.subscribe();
assert_eq!(shared.subscriber_count(), 1);
shared.unsubscribe(id);
assert_eq!(shared.subscriber_count(), 0);
assert!(rx.recv().is_err());
shared.close().unwrap();
let _ = std::fs::remove_file(&tmp);
let _ = std::fs::remove_file(format!("{}-wal", tmp.display()));
let _ = std::fs::remove_file(format!("{}-shm", tmp.display()));
}
/// Subscribers must learn that the watcher thread died — not just
/// stop receiving wakes silently. We force the watcher's death
/// via the dead-man's switch (file replacement) and assert that an
/// already-subscribed receiver returns `Err(RecvError)` on its
/// next blocking `recv()`. Without WatcherDeathGuard this test
/// hangs (subscriber blocks forever) and times out.
///
/// Also asserts the port's dead-man's-switch delta (ADR-012 §4):
/// the thread exits *without panicking* — `join()` is `Ok(())`
/// where the lineage's `join()` returned the panic payload.
#[test]
#[cfg(unix)]
fn shared_update_watcher_signals_subscribers_on_watcher_death() {
let tmp = temp_db("death-signal");
wal_db(&tmp);
let shared = SharedUpdateWatcher::new(tmp.clone()).unwrap();
let (_id, rx) = shared.subscribe();
// Let the watcher snapshot the initial identity.
std::thread::sleep(Duration::from_millis(200));
let other = temp_db("death-other");
let _ = std::fs::File::create(&other).unwrap();
std::fs::rename(&other, &tmp).unwrap();
let deadline = std::time::Instant::now() + Duration::from_secs(2);
loop {
match rx.recv_timeout(Duration::from_millis(100)) {
Ok(()) => continue,
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => break,
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {}
}
assert!(
std::time::Instant::now() < deadline,
"watcher died but subscriber's channel never disconnected — \
WatcherDeathGuard didn't fire?"
);
}
shared.close().unwrap();
let _ = std::fs::remove_file(&tmp);
let _ = std::fs::remove_file(format!("{}-wal", tmp.display()));
let _ = std::fs::remove_file(format!("{}-shm", tmp.display()));
}
/// The port's dead-man's-switch delta observed directly at the
/// watcher: a db-file replacement produces a clean (non-panic)
/// thread exit. The lineage panicked here; ADR-012 §4 replaces
/// the unwind with an ordinary death.
#[test]
#[cfg(unix)]
fn update_watcher_dies_cleanly_on_file_replacement() {
let tmp = temp_db("watcher-replace");
wal_db(&tmp);
let watcher = UpdateWatcher::spawn(tmp.clone(), || {}).unwrap();
std::thread::sleep(Duration::from_millis(200));
let tmp2 = temp_db("watcher-replace-new");
wal_db(&tmp2);
std::fs::rename(&tmp2, &tmp).unwrap();
std::thread::sleep(Duration::from_millis(500));
let result = watcher.join();
assert!(
result.is_ok(),
"watcher should have exited cleanly on file replacement \
(dead-man's switch must not panic, ADR-012 §4), got {result:?}"
);
let _ = std::fs::remove_file(&tmp);
let _ = std::fs::remove_file(format!("{}-wal", tmp.display()));
let _ = std::fs::remove_file(format!("{}-shm", tmp.display()));
}
#[test]
fn shared_update_watcher_prunes_subscribers_when_receiver_dropped() {
let tmp = temp_db("prune");
wal_db(&tmp);
let shared = SharedUpdateWatcher::new(tmp.clone()).unwrap();
{
let _subs: Vec<_> = (0..10).map(|_| shared.subscribe()).collect();
assert_eq!(shared.subscriber_count(), 10);
}
let writer = Connection::open(&tmp).unwrap();
writer
.execute(
"CREATE TABLE IF NOT EXISTS _test_prune(id INTEGER PRIMARY KEY)",
[],
)
.unwrap();
let deadline = std::time::Instant::now() + Duration::from_secs(2);
while shared.subscriber_count() != 0 && std::time::Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(5));
writer
.execute("INSERT INTO _test_prune(id) VALUES (random())", [])
.unwrap();
}
assert_eq!(shared.subscriber_count(), 0);
shared.close().unwrap();
let _ = std::fs::remove_file(&tmp);
let _ = std::fs::remove_file(format!("{}-wal", tmp.display()));
let _ = std::fs::remove_file(format!("{}-shm", tmp.display()));
}
#[test]
fn data_version_detects_commits_and_ignores_rollbacks() {
let tmp = temp_db("dv");
let watcher = Connection::open(&tmp).unwrap();
watcher.execute_batch("PRAGMA journal_mode = WAL;").unwrap();
let writer = Connection::open(&tmp).unwrap();
let v0 = poll_data_version(&watcher).unwrap();
writer.execute("CREATE TABLE t(x INTEGER)", []).unwrap();
let v1 = poll_data_version(&watcher).unwrap();
assert!(v1 > v0, "commit should increment data_version");
writer.execute_batch("BEGIN IMMEDIATE;").unwrap();
writer.execute("INSERT INTO t VALUES (1)", []).unwrap();
writer.execute_batch("ROLLBACK;").unwrap();
let v2 = poll_data_version(&watcher).unwrap();
assert_eq!(v2, v1, "rollback should not increment data_version");
let _ = std::fs::remove_file(&tmp);
}
#[test]
fn data_version_survives_wal_checkpoint() {
let tmp = temp_db("dv-ckpt");
let watcher = Connection::open(&tmp).unwrap();
watcher.execute_batch("PRAGMA journal_mode = WAL;").unwrap();
let w0 = poll_data_version(&watcher).unwrap();
let writer = Connection::open(&tmp).unwrap();
writer.execute("CREATE TABLE t(x INTEGER)", []).unwrap();
let w1 = poll_data_version(&watcher).unwrap();
assert!(
w1 > w0,
"commit from other conn should increment data_version"
);
writer
.execute_batch("PRAGMA wal_checkpoint(TRUNCATE);")
.unwrap();
let w2 = poll_data_version(&watcher).unwrap();
assert!(
w2 > w1,
"checkpoint from other conn should increment data_version"
);
writer.execute("INSERT INTO t VALUES (1)", []).unwrap();
let w3 = poll_data_version(&watcher).unwrap();
assert!(
w3 > w2,
"post-checkpoint commit should increment data_version"
);
let _ = std::fs::remove_file(&tmp);
}
fn poll_data_version_works_in_journal_mode(mode: &str) {
let tmp = temp_db(&format!("jm-{}", mode.to_ascii_lowercase()));
let watcher = Connection::open(&tmp).unwrap();
watcher
.execute_batch(&format!("PRAGMA journal_mode = {mode};"))
.unwrap();
let actual: String = watcher
.pragma_query_value(None, "journal_mode", |r| r.get(0))
.unwrap();
assert_eq!(
actual.to_ascii_uppercase(),
mode.to_ascii_uppercase(),
"PRAGMA journal_mode = {mode} silently fell back to {actual}"
);
let writer = Connection::open(&tmp).unwrap();
let v0 = poll_data_version(&watcher).unwrap();
writer.execute("CREATE TABLE t(x INTEGER)", []).unwrap();
let v1 = poll_data_version(&watcher).unwrap();
assert!(
v1 > v0,
"journal_mode={mode}: cross-conn commit should bump \
data_version; saw {v0} -> {v1}"
);
writer.execute_batch("BEGIN IMMEDIATE;").unwrap();
writer.execute("INSERT INTO t VALUES (1)", []).unwrap();
writer.execute_batch("ROLLBACK;").unwrap();
let v2 = poll_data_version(&watcher).unwrap();
assert_eq!(
v2, v1,
"journal_mode={mode}: rollback should not bump data_version"
);
let _ = std::fs::remove_file(&tmp);
let _ = std::fs::remove_file(format!("{}-wal", tmp.display()));
let _ = std::fs::remove_file(format!("{}-shm", tmp.display()));
let _ = std::fs::remove_file(format!("{}-journal", tmp.display()));
}
#[test]
fn poll_data_version_works_in_wal() {
poll_data_version_works_in_journal_mode("WAL");
}
#[test]
fn poll_data_version_works_in_delete() {
poll_data_version_works_in_journal_mode("DELETE");
}
#[test]
fn poll_data_version_works_in_truncate() {
poll_data_version_works_in_journal_mode("TRUNCATE");
}
#[test]
fn poll_data_version_works_in_persist() {
poll_data_version_works_in_journal_mode("PERSIST");
}
#[cfg(any(unix, windows))]
#[test]
fn stat_identity_detects_file_replacement() {
let tmp = temp_db("id-test");
let tmp2 = temp_db("id-test2");
std::fs::write(&tmp, b"original").unwrap();
std::fs::write(&tmp2, b"replacement").unwrap();
let id1 = stat_identity(&tmp).unwrap();
let id2 = stat_identity(&tmp2).unwrap();
assert_ne!(id1, id2, "different files should have different identities");
std::fs::rename(&tmp2, &tmp).unwrap();
let id3 = stat_identity(&tmp).unwrap();
assert_eq!(
id3, id2,
"renamed file should carry the replacement's identity"
);
let _ = std::fs::remove_file(&tmp);
}
#[cfg(any(unix, windows))]
#[test]
fn fold_high_res_uses_both_halves() {
let (vsn, idx) = fold_high_res(0xAABB, 0x0000_0000_0000_0000_DEAD_BEEF_CAFE_F00D);
assert_eq!(vsn, 0xAABB);
assert_eq!(idx, 0xDEAD_BEEF_CAFE_F00D);
let upper = 0x1111_2222_3333_4444u64;
let lower = 0x5555_6666_7777_8888u64;
let file_id = ((upper as u128) << 64) | (lower as u128);
let (vsn, idx) = fold_high_res(0xCCDD, file_id);
assert_eq!(vsn, 0xCCDD);
assert_eq!(idx, upper ^ lower);
let same = 0xDEAD_BEEF_CAFE_F00Du64;
let file_id = ((same as u128) << 64) | (same as u128);
let (_, idx) = fold_high_res(0, file_id);
assert_eq!(idx, 0);
}
/// W-1: with the watcher connection down and the db file vanished,
/// the reconnect path is throttled — not one open attempt per
/// poll tick. Driven directly: a missing db path fails every open,
/// the loop runs its reconnect schedule against it, and the
/// instrumented open counter counts attempts.
#[test]
fn reconnect_attempts_are_bounded_by_backoff() {
static ATTEMPTS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
fn counting_open(path: &Path) -> rusqlite::Result<Connection> {
ATTEMPTS.fetch_add(1, Ordering::Relaxed);
open_watcher_conn(path)
}
let tick_ms = 10u64;
const WINDOW_MS: u64 = 2_600;
let ticks_in_window = WINDOW_MS / tick_ms;
let db_path = PathBuf::from(format!(
"/this/parent/does/not/exist/alkstore-w1-{}.db",
std::process::id()
));
let stop = Arc::new(AtomicBool::new(false));
let stop_t = stop.clone();
let (ready_tx, ready_rx) = std::sync::mpsc::sync_channel::<()>(1);
let handle = std::thread::spawn(move || {
run_poll_loop(
db_path,
|| {},
stop_t,
ready_tx,
Duration::from_millis(tick_ms),
counting_open,
)
});
ready_rx
.recv_timeout(Duration::from_secs(2))
.expect("watcher loop baseline capture");
std::thread::sleep(Duration::from_millis(WINDOW_MS));
stop.store(true, Ordering::Release);
handle.join().unwrap();
let attempts = ATTEMPTS.load(Ordering::Relaxed);
assert!(
attempts > 0,
"the loop must keep trying to reconnect while the db is missing"
);
// The backoff bound: at tick=tick_ms, the reconnect schedule
// admits ticks_in_window / MAX_RECONNECT_TICKS (+1 in-flight)
// attempts in the window, not ticks_in_window (the lineage's
// one-per-tick schedule). Scheduler jitter can compress the
// effective window slightly; one extra attempt of headroom is
// the margin for that, not for a broken backoff — the lineage's
// one-per-tick schedule would produce ~260 attempts here.
let bound = ticks_in_window / MAX_RECONNECT_TICKS + 2;
assert!(
attempts <= bound,
"W-1 violated: {attempts} reconnect attempts in a {WINDOW_MS}ms \
window at {tick_ms}ms/tick; the backoff admits at most {bound}"
);
}
}
+108 -9
View File
@@ -1,7 +1,7 @@
---
id: fork-port-connection-watcher
name: Fork port — connection architecture + watcher machinery
status: pending
status: completed
depends_on: [fork-substrate-scaffold]
scope: broad
risk: medium
@@ -73,17 +73,17 @@ is not separately testable through a public seam).
## Acceptance Criteria
- [ ] Kept machinery ported with upstream names/structure preserved;
- [x] Kept machinery ported with upstream names/structure preserved;
diff vs. lineage reviewable (the wave-2 review checks this)
- [ ] Three watcher deltas applied; dead-man's switch exits without
- [x] Three watcher deltas applied; dead-man's switch exits without
panicking; `WatcherDeathGuard` still closes all subscribers
- [ ] Dropped machinery absent (no cron, no experimental watchers, no
- [x] Dropped machinery absent (no cron, no experimental watchers, no
rate-limit/result tables, no superseded queue functions)
- [ ] Table family is `__alkstore_*`; bootstrap fresh-only; race
- [x] Table family is `__alkstore_*`; bootstrap fresh-only; race
swallow re-keyed to `pragma_table_info`
- [ ] Inherited test suites green (adapted); watcher-death-closes-
- [x] Inherited test suites green (adapted); watcher-death-closes-
subscribers test present
- [ ] Substrate is sync (no tokio imports); no panics/`unwrap()` in
- [x] Substrate is sync (no tokio imports); no panics/`unwrap()` in
library code; clippy `-D warnings`, fmt clean
## References
@@ -95,8 +95,107 @@ is not separately testable through a public seam).
## Notes
> To be filled by implementation agent
- **Module layout**: the lineage's two files (`lib.rs` ~4k lines,
`honker_ops.rs` ~3.3k) map to three substrate modules in upstream's
order — `schema.rs` (pragmas/open, notify, bootstrap, `Writer`,
`Readers`), `watcher.rs` (the polling watcher family),
`ops.rs` (savepoint discipline, arg coercion, stream/lock ops, the
scalar-function attachments). Upstream names preserved on the kept
body; the mapping is register entry D-17 (the flat 4k `lib.rs`
doesn't fit the module-per-file standard; the ADR-013 fold makes the
module file the unit).
- **`file-id` retained** (D-18): the task's drop list names it among the
experimental backends' optional deps, but it drives the *kept*
dead-man's switch (`stat_identity`) — the kept-half scope wins for
this dep; target-gated `cfg(any(unix, windows))` exactly as upstream
carries it.
- **W-1 shape**: bounded reconnect backoff as a tick-counter
(`MAX_RECONNECT_TICKS` = 100 poll ticks between attempts → ≤10
open attempts/sec at the default 1 ms cadence, vs the lineage's
~1000/sec). Implemented inside `run_poll_loop`; the open step is a
parameter (default = the real open) so the test drives the loop
directly with a counting open.
- **W-2 shape**: `SharedUpdateWatcher::new*` and `UpdateWatcher::spawn*`
return `Result<_, String>` (upstream's spawn used `.expect(...)`).
The thread-build step is isolated in `UpdateWatcher::spawn_thread`
with a test-reachable stack-size knob so the `Err` path is provable
(an oversized stack reservation fails reliably; a nonexistent db path
does NOT fail thread spawn — baseline capture inside the loop is a
retry path, not a spawn failure). The engine maps the String into the
taxonomy's `Database` in wave 3.
- **Dead-man's switch delta**: identity mismatch logs the precise
diagnostic (upstream's panic message text, repointed) and `return`s
through the ordinary death path; the death-closes-subscribers test
additionally asserts `join()` is now `Ok` (the lineage's returned
the panic payload).
- **Bootstrap**: the scheduler's `cron_expr` column is carried as
`spec` in `__alkstore_scheduler_tasks` (the storage the scheduler
collapse lands in; cron strings are rejected at the engine layer —
ADR-009; naming is storage-internal so the column rename is
registration D-10's surface). The three append-column migrations
stay (enabled/max_attempts/claimed_at — the schema D-12's stamps
extend); the duplicate-column race swallow is re-keyed to
`pragma_table_info` exactly per ADR-012 §5.
- **Bootstrap schema scope**: the bootstrap carries the full
`__alkstore_*` family including the live/dead table shapes the
re-derivation extends — the tables must exist for the engine's
`alkstore_bootstrap()` to be total even before wave 3 wires the
queue ops; the re-derivation task alters/extends as its SQL needs.
- **Notify at-attach cap (D-16)**: `NOTIFY_ATTACH_MAX_ROWS` = 10_000,
oldest-first trim below the cap on every attach — replacing
upstream's "no pruning" posture per ADR-010 §6 (the fix for
unbounded notifications growth must not become a caller's chore).
- **Test scope boundary**: honker's queue-op suites (savepoint
dead-letter paths, claim pressure, `claimed_at` transitions,
scheduler) test the *re-derivation's* machinery — they are
`fork-rederive-queue-ops`'s floor, not this task's. This task
inherits: PRAGMA/WAL open posture (incl. the concurrent-open race
suite), watcher lifecycle + failure handling (fan-out, unsubscribe,
subscriber prune, death signal, all four journal modes, checkpoint,
stat identity, XOR-fold unit tests), savepoint machinery tests
(rollback-failure reporting, panic unwinding, scalar-path frames),
arg-coercion tests (driven through a probe scalar function — the
queue-function shape they exist for is re-derived next), notify
tests, lock tests, a new stream-ops suite, and a cross-mechanism
pressure test (streams+notify+locks, 5 producers) — the pressure
shape minus queue ops. Added: the W-1 backoff-bound test, the W-2
spawn-failure test, the panic-free death test, the
column-order-equality migration test, the legacy-orphan-inertness
test. 43 tests green.
- **Port state**: the substrate compiles as engine-crate modules and
is test-covered, but nothing re-exports through the engine's public
API yet — wave 3 wires `open` + the trait impl (`mod.rs` carries
`#![allow(dead_code)]` marked as lift-on-wiring).
## Summary
> To be filled on completion
The kept half of the honker-core fork landed in
`alkstore-sqlite/src/substrate/` as `schema.rs` + `watcher.rs` +
`ops.rs` (~1k lines ported machinery + adapted inherited suites, 43
tests): PRAGMA/WAL open posture with `set_journal_mode_wal` retry
(+ the concurrent-open race suite), `Writer`, `Readers` (incl. the
open-failure capacity-leak test), the polling watcher +
`SharedUpdateWatcher` + `WatcherDeathGuard` + `stat_identity`
dead-man's switch (all four journal modes + checkpoint + fan-out +
prune suites), `in_savepoint`/`UnwindUndo` (+ panic-discipline
tests), REAL-coercion arg helpers (+ coercion tests), the notify
scalar + table (now with the ADR-010 §6 at-attach pruning cap), the
stream functions (+ read/publish/offset suite), and the lock functions
(+ grant/refuse/renew/release suites). Deltas: W-1 bounded reconnect
backoff (tick-counting, test-pinned at ≤ 2 attempts in a 2.6 s window
vs the lineage's ~260), W-2 fallible watcher spawn (Result, mapped
error names the watcher), the dead-man's-switch panic replaced by
log-and-exit (join is Ok; death still closes every subscriber — the
adapted death-signal test pins both), `_honker_*` → `__alkstore_*`
(legacy tables pinned inert by test), duplicate-column race swallow
re-keyed to `pragma_table_info`, scheduler `cron_expr` → `spec`,
fresh-only bootstrap with append-column migrations (column-order
equality pinned). Dropped fully: cron, kernel/shm watcher backends,
rate-limit/result tables, superseded queue functions (the pressure
suite's queue half moved to the re-derivation task). `file-id` kept
(dead-man's-switch dep; register D-18 corrects D-02's enumeration).
PROVENANCE.md register updated to the landed state (D-01..D-20).
Verified: `cargo build`, `cargo test` (workspace; 43 sqlite-crate
tests green), `cargo clippy --all-targets -- -D warnings`, `cargo fmt
--check` all clean. Subtree stays sync (zero tokio); no
panics/unwrap/expect outside tests.