Files
alkstore/docs/research/reference-honker-machinery.md
glm-5.3-flash befbe2e714 OQ-06 resolved: honker-core quality read fires the fork trigger (ADR-011)
Quality read of honker-core's watcher/transactional core cross-checked
against the published crates.io artifact: the core itself is clean
(Writer/Readers, polling-watcher failure handling, WatcherDeathGuard
all verified), but published 0.5.0 predates upstream's unreleased fix
train carrying the issue-#133 savepoint hardening (silent job loss in
the dead-letter paths) and five .ok() error swallows — and ADR-010's
queue depth requires engine-owned queue SQL in any posture. Resolution:
fork honker-core at the reference revision, inherit the clean machinery
and test suites, re-derive queue ops on contract v1, rename tables to
__alkstore_*.

- docs/research/quality-read-honker-core.md — full evidence
- docs/architecture/decisions/011-sqlite-substrate-fork.md — decision
- OQ-06 resolved in open-questions.md; ADR-003/005/008, engine-sqlite,
  queues, README annotated for consistency
2026-10-05 03:39:46 +00:00

22 KiB

status, last_updated
status last_updated
draft 2026-10-05

Reference read: honker queue / scheduler / outbox / lock maintenance machinery

Per-agents conventions: reference checkouts are read freely, never wired in; checkout state noted when findings depend on code specifics.

Checkout: /workspace/honker @ f4e53c6 (node-v0.5.1-10-gf4e53c6), matching the AGENTS.md reference revision. Versions: honker-core 0.5.0 (honker-core/Cargo.toml), honker Rust wrapper 0.5.0 (packages/honker-rs/Cargo.toml). The Rust binding is excluded from the root workspace; all real machinery is in honker-core/src/{lib.rs, honker_ops.rs, cron.rs}.

Architecture in one sentence: SQLite tables (_honker_live, _honker_dead, _honker_locks, _honker_scheduler_tasks, _honker_results, _honker_notifications, rate-limit/stream tables) + SQL scalar functions (honker_*) implemented once in Rust (honker_ops::attach_honker_functions), consumed identically by Python/Node/Rust bindings; cross-process wake is PRAGMA data_version polling (1 ms default) driving an in-process channel — not LISTEN/NOTIFY.

1. Queues

1.1 Job table schema

BOOTSTRAP_HONKER_SQL — honker-core/src/lib.rs:352-432:

_honker_live (lib.rs:353-367) — pending and processing both live here:

column type notes
id INTEGER PK AUTOINCREMENT
queue TEXT NOT NULL namespace within one table
payload TEXT NOT NULL stored as JSON string
state TEXT DEFAULT 'pending' two states: 'pending', 'processing'
priority INTEGER DEFAULT 0 higher = earlier (ORDER BY priority DESC)
run_at INTEGER DEFAULT (unixepoch()) readiness deadline
worker_id TEXT current claim holder
claim_expires_at INTEGER visibility deadline; moves on heartbeat
attempts INTEGER DEFAULT 0 bumped on every claim, including reclaims
max_attempts INTEGER DEFAULT 3 per-row, frozen at enqueue
created_at INTEGER DEFAULT (unixepoch())
expires_at INTEGER job-level expiry; NULL = never
claimed_at INTEGER, nullable start of current attempt; migration added (lib.rs:486-508)

Partial indexes tuned to the hot paths (lib.rs:368-376): claim index (queue, priority DESC, run_at, id) WHERE state IN ('pending','processing'); pending index (queue, run_at); processing index (queue, claim_expires_at). Dead rows never touch the claim index.

_honker_dead (lib.rs:377-388): id, queue, payload, priority, run_at, attempts, max_attempts, last_error, created_at, died_at. Move, not flag: dead-lettering physically DELETEs from _honker_live and INSERTs into _honker_dead. No retention/expiry on dead rows — "Never scanned by the claim path; retention policy is the user's problem" (lib.rs:345-347). No DELETE FROM _honker_dead anywhere in the repo.

1.2 Enqueue and its options

Core signature (honker_ops.rs:984-993):

pub fn enqueue(
    conn: &Connection, queue: &str, payload: &str,
    run_at: Option<i64>, delay: Option<i64>,
    priority: i64, max_attempts: i64, expires: Option<i64>,
) -> rusqlite::Result<i64>

SQL function honker_enqueue(...) (honker_ops.rs:593-619). Rust-level EnqueueOpts (packages/honker-rs/src/lib.rs:438-444) — note max_attempts is not here; it comes from queue-level QueueOpts (visibility_timeout_s: 300, max_attempts: 3, lib.rs:421-434):

pub struct EnqueueOpts {
    pub delay: Option<i64>,
    pub run_at: Option<i64>,
    pub priority: i64,
    pub expires: Option<i64>,
}

Semantics (honker_ops.rs:976-1000):

  • run_at vs delay precedence: delay wins — `run_at = unixepoch()
    • delayif delay set; else literalrun_at; else now`.
  • expires: relative seconds → expires_at = unixepoch() + expires; NULL = never expires.
  • No synthetic wake row on enqueue — the INSERT advancing data_version is the wake; an earlier design writing a _honker_notifications row per enqueue grew unboundedly (honker_ops.rs:1002-1006).

arg_i64/arg_opt_i64 (honker_ops.rs:40-54) accept REAL-typed whole numbers because better-sqlite3 binds every JS number as REAL — interop lesson (fractional values error with a diagnostic).

1.3 Claim mechanics

pub fn claim_batch(conn: &Connection, queue: &str, worker_id: &str,
                   n: i64, timeout_s: i64) -> rusqlite::Result<String>

(honker_ops.rs:848-914). timeout_s is the visibility timeout, set per-claim (the Rust binding passes static QueueOpts .visibility_timeout_s, honker-rs lib.rs:584-591).

Single atomic UPDATE ... WHERE id IN (SELECT ...) RETURNING (honker_ops.rs:866-886):

  • Claim predicate: state IN ('pending','processing') AND attempts < max_attempts AND (expires_at IS NULL OR expires_at > unixepoch()) AND ((state='pending' AND run_at <= now) OR (state='processing' AND claim_expires_at < now)).
  • Claim action: state='processing', worker_id=worker, claim_expires_at = now + timeout_s, claimed_at = now, attempts += 1.
  • Ordering: priority DESC, run_at ASC, id ASC, LIMIT n.

Reclaim = another attempt. Claiming bumps attempts on reclaims too — visibility timeouts consume the retry budget. claimed_at is refreshed on reclaim but not on heartbeat (doc honker_ops.rs:831-847; test honker_ops.rs:3103-3149).

claim_one is a thin claim_batch(worker_id, 1) (honker-rs lib.rs:606-609).

Pre-claim dead-letter sweep: every claim_batch first runs dead_letter_exhausted_claimable (honker_ops.rs:781-829) — rows at attempts >= max_attempts that are claimable-reclaimable move to _honker_dead with last_error='max attempts exceeded' before the claim UPDATE (covers "worker died right after the last allowed claim").

Idle-wake helper: honker_queue_next_claim_at(queue) -> unix_ts (honker_ops.rs:936-970) — earliest future run_at (pending) or claim_expires_at + 1 (processing with budget); 0 if nothing. Used by the Rust ClaimWaker::next loop (honker-rs lib.rs:805-845): wake on data_version or sleep until next_at, no busy-poll.

1.4 Visibility timeout / heartbeat / renewal

pub fn heartbeat(conn, job_id, worker_id, extend_s) -> Result<i64>

(honker_ops.rs:1306-1323); 1 = extended, 0 = refused.

  • Heartbeat is a renewal (visibility extension), not a progress signal — no progress field exists. Sets claim_expires_at = unixepoch() + extend_s (absolute reset from now, not additive).
  • Guard: WHERE id = ? AND worker_id = ? AND state = 'processing' AND claim_expires_at >= unixepoch() — a late heartbeat after the visibility timeout is refused so it cannot steal the job back from a reclaimer (dual-execution guard, honker_ops.rs:1312-1314).
  • Missed heartbeat ⇒ claim_expires_at passes ⇒ row claimable by the normal predicate; a reclaiming claimer consumes an attempt. Nothing transitions the row on lapse; stale worker_id etc. remain until a next claim overwrites them (honker_ops.rs:843-847).
  • No per-binding heartbeat thread in core; cadence is caller- implemented.

1.5 Retry behavior — attempts and backoff

pub fn retry(conn, job_id, worker_id, delay_s, error) -> Result<i64>

(honker_ops.rs:1045-1123). Requires a still-valid processing claim; attempts >= max_attempts ⇒ dead-letter (savepoint-wrapped DELETE→INSERT, honker_ops.rs:1080-1107); else pending with run_at = unixepoch() + delay_s, claim fields cleared.

Backoff is NOT in the queue core. retry takes a caller-chosen delay_s. Curves live in wrappers:

  • Rust Outbox::retry_delay: base_backoff_s * 2^(attempts-1), saturating (honker-rs lib.rs:543-549).
  • Python _compute_delay: retry_delay * backoff**(attempts-1) (packages/honker/python/honker/_worker.py:108-114); @task defaults retry_delay_s=60, backoff=1.0 (constant).
  • No jitter and no cap in either (Rust caps only via saturating_mul).

1.6 Dead-letter handling

  • Dedicated _honker_dead; move semantics (§1.1). Triggers: (a) retry() at budget (honker_ops.rs:1080-1107), (b) explicit fail() (honker_ops.rs:1135-1192), (c) dead_letter_exhausted_claimable pre-claim (honker_ops.rs:781-829), (d) sweep_expired expiration (honker_ops.rs:1336-1381).
  • last_error strings: caller error, 'max attempts exceeded', 'expired'.
  • Retention: none. No sweep/TTL/pruning of dead rows exists.
  • All four move-sites are wrapped in a shared savepoint helper in_savepoint (honker_ops.rs:118-254) — DELETE ... RETURNING + decode + INSERT had a half-failing window; a failing decode/INSERT previously produced live=0, dead=0 silent job loss (CHANGELOG "Core SQLite error propagation"; UnwindUndo undoes the frame if the body panics). The savepoint hardening is very recent, post-0.5.1 work in the CHANGELOG "Unreleased" section. Most instructive correctness pattern in the codebase.

1.7 sweep_expired

pub fn sweep_expired(conn, queue) -> Result<i64>

(honker_ops.rs:1336-1381). Only sweeps state='pending' rows with expires_at IS NOT NULL AND expires_at <= now → _honker_dead with last_error='expired'. Does not touch processing rows — expired claims are handled lazily by the claim predicate. Python doc: "The claim path already ignores expired rows, so sweep is cleanup-only — not correctness-critical" (_honker.py:387-394).

1.8 cancel / get_job / ack_batch

pub fn cancel(conn, job_id) -> Result<i64>          // honker_ops.rs:1203-1209
pub fn get_job(conn, job_id) -> Result<String>      // :1220-1295
pub fn ack(conn, job_id, worker_id) -> Result<i64>  // :1028-1035
pub fn ack_batch(conn, ids_json, worker_id) -> Result<i64> // :920-934
  • cancel: unconditional DELETE ... WHERE id=? AND state IN ('pending','processing') regardless of which worker holds it. Not an interrupt — the holder's next ack/heartbeat returns 0 (honker-rs docs lib.rs:625-631).
  • get_job: reads _honker_live only → dead jobs are invisible; JSON on hit, empty string on miss (honker-rs maps to None, lib.rs:642-650). claimed_at: null for never-claimed jobs.
  • ack/ack_batch: DELETE ... WHERE id IN (...) AND worker_id=? AND claim_expires_at >= now RETURNING id, returning count. Deliberately savepoint-free (honker_ops.rs:916-919). Ack fails (0) if the claim expired. Note: ack doesn't check state='processing' — keys off worker_id + claim_expires_at only (harmless today, implicit).

2. Scheduler

2.1 Storage schema

_honker_scheduler_tasks (lib.rs:400-410): name PK, queue, cron_expr, payload, priority (default 0), expires_s (nullable), next_fire_at NOT NULL, enabled INTEGER DEFAULT 1, max_attempts INTEGER DEFAULT 3. Runtime ALTER migrations tolerate the concurrent-bootstrap "duplicate column" race (lib.rs:443-508). No last_fired_at, no tz column — cron arithmetic is system local time.

2.2 Surface

Core (honker_ops.rs): scheduler_register (:1490-1536), scheduler_unregister (:1538-1551), scheduler_tick (:1583-1654), scheduler_soonest (:1656-1664), scheduler_pause/resume (:1669-1689, toggling enabled), scheduler_list (:1694-1752), scheduler_update (:1759-1844). Plus honker_cron_next_after(expr, from_unix) DETERMINISTIC SQL function (:683-692).

Rust binding (packages/honker-rs/src/lib.rs:1365-1584): add, add_with_max_attempts, remove, tick, soonest, pause, resume, list, update (with ScheduleUpdate { cron_expr, payload, priority, expires_s: Option<Option<i64>> } — inner Option distinguishes "clear" from "leave alone"), update_max_attempts, and the blocking run(stop, owner) leader loop. Python mirrors (_scheduler.py:113-437).

Register semantics: upsert by name, replaces entirely; next_fire_at = next cron boundary strictly after now; clamps max_attempts < 1 to 1.

2.3 Tick and catch-up

schema_tick(conn, now_unix) -> JSON [{name, queue, fire_at, job_id}] (honker_ops.rs:1583-1654):

  • Selects WHERE next_fire_at <= now AND enabled = 1.
  • Per task, loops while next_fire_at <= now: enqueue (run_at NULL = claim immediately, task-level expires_s/max_attempts applied), then next_fire_at = cron_next_after(expr, next_fire_at) — one fire per missed boundary, minute-by-minute catch-up.
  • SCHEDULER_MAX_CATCHUP_FIRES = 64 cap per task per tick (honker_ops.rs:1564-1575); beyond the cap remaining boundaries skip — next_fire_at jumps to next_after(now). Documented as semantic: "Run the scheduler continuously, use coarser schedules, or raise this constant."
  • Final next_fire_at persisted in a separate UPDATE (honker_ops.rs :1647-1651) within the caller's transaction — enqueue+advance commit together.
  • Dedup: no DB-level claim token; the advance-then-return contract under BEGIN IMMEDIATE means at most one ticker observes an unfired boundary — pinned by a 10-thread race test (tests/test_scheduler.py:397-470). The leader lock is the production gate; the SQL is "still safe if someone forgets it".

Guarantee: at-least-once per boundary, up to 64-boundary backlog; beyond that, gap-skip by design. Crash mid-tick rolls back both enqueues and the advance (boundary refires).

2.4 Leader loop / election

Election is a named TTL lock in _honker_locks, not SQLite advisory APIs:

  • Rust Scheduler::run (lib.rs:1509-1584): LOCK_NAME="honker-scheduler", TTL 60s, heartbeat every 20s (Python _scheduler.py:135-137, TTL 60, heartbeat 30). Non-leaders poll every 5s.
  • Leader loop (lib.rs:1539-1584): renew lock → on failure exit before ticking (no dual-fire alongside the thief) → tick() → soonest → sleep min(heartbeat_interval, until_soonest).
  • Python _main_loop additionally re-proves ownership before every fire (_scheduler.py:382-437), runs tick+soonest in one writer tx. Losing the lock raises LeadershipLost.
  • Wake-on-register: register/unregister/pause/resume/update rely on data_version advancing (a former synthetic notification row was removed as unbounded-growth — scheduler_wake is now a no-op, honker_ops.rs:1553-1562).

2.5 Cron parsing (cron.rs, 473 lines)

  • 5-field cron (min hour dom month dow), 6-field with seconds, and @every <n><unit> (s/m/h/d) (cron.rs:1-14, 97-129). Full syntax: *, ranges, lists, steps; dow 0=Sunday; field validation with good errors (cron.rs:131-178).
  • Next-boundary search: iterative scan, bounded to 100 years (cron.rs:207-315). DST handled: spring-forward gap skips to real time, ambiguous fall-back picks the first instantiation (cron.rs:330-358; tests cron.rs:456-472).
  • Brittleness — system local timezone dependence: no TZ parameter in the public surface (next_after_unix uses chrono::Local, cron.rs:319-321; tz-parameterized variant is pub(crate) test-only). A server TZ change silently re-times every cron task; stored schedules are not TZ-qualified.

3. Outbox

Pattern: an outbox is a regular queue named _outbox:{name} — no separate table, no separate delivery semantics (honker-rs lib.rs:478-492; _honker.py:865-934).

Rust shape (OutboxOpts: visibility_timeout_s: 60, max_attempts: 5, base_backoff_s: 5, lib.rs:454-469):

pub fn run_once<F, E>(&self, worker_id: &str, delivery: F) -> Result<bool>
where F: FnMut(serde_json::Value) -> Result<(), E>

(lib.rs:516-541): claim_one → parse payload → delivery(payload) → Ok ⇒ job.ack() (ack failure surfaced as error, :527-529) ⇒ Err ⇒ job.retry(base_backoff_s * 2^(attempts-1), err) (:531-538). A pull worker — app calls run_once in its own loop.

Transactional enqueue: outbox.enqueue_tx(&tx, ...) routes through the caller's Transaction (lib.rs:506-513); single mutex-guarded-connection means same-thread in-tx ops must use _tx variants (lib.rs:241-250). Failure path: at-least-once with exponential backoff, dead-lettering via normal budget mechanics. No heartbeat inside run_once — delivery slower than visibility_timeout_s (default 60s) can have its claim expire mid-delivery and be redelivered (honker's own outbox is exposed to the dual-execution window).

4. Locks (maintenance-relevant)

_honker_locks (name TEXT PK, owner TEXT NOT NULL, expires_at INTEGER NOT NULL) (lib.rs:389-393).

pub fn lock_acquire(conn, name, owner, ttl_s) -> Result<i64>   // 1/0  honker_ops.rs:1387-1415
pub fn lock_release(conn, name, owner) -> Result<i64>          // :1417-1423
pub fn lock_renew(conn, name, owner, ttl_s) -> Result<i64>     // 1/0  :1431-1442
  • lock_acquire: opportunistic per-name expiry sweep (DELETE WHERE name=? AND expires_at <= now) + INSERT OR IGNORE + read-back. The PK on name prevents dual acquisition.
  • INSERT OR IGNORE does not refresh TTL on same-owner re-acquire — hence lock_renew as a separate function (honker_ops.rs:1425-1431; lib.rs SQL comment :329-332); acquire returns 1 while silently keeping the old expiry. lock_renew requires (name, owner) match and positive ttl.
  • Lock::heartbeat(ttl_s) RAII wrapper (honker-rs lib.rs:1682-1691), with the honest caveat "holding the Lock value alone does not guarantee you still own the lock".
  • No general lock-expiry sweep — expired locks for names nobody acquires again persist ("Expired rows ... are opportunistically pruned on every acquire attempt", _honker.py:962-965). Locks also serve as operator mutexes (skip-overlapping-runs recipe, _honker.py:1160-1170).

5. Notifications pruning

_honker_notifications (id, channel, payload, created_at) (lib.rs:308-315). Never auto-pruned — "no magic timer" (lib.rs :303-305; tests pin that notify() never prunes, _honker.py :1182-1208).

  • Rust prune_notifications(older_than_s) (lib.rs:323-332) and prune_notifications_keep_latest(max_keep) (lib.rs:336-351 — rank-based; OFFSET trick correct across id gaps; max_keep < 0 errors).
  • Python combines both conditions with OR semantics — aggressive older_than_s is not blocked by max_keep; max_keep "is not a floor" (_honker.py:1182-1245).

6. Result storage (brief)

_honker_results (job_id INTEGER PK, value TEXT, created_at, expires_at) (lib.rs:411-416).

result_save (ttl ≤ 0 ⇒ NULL), result_get (expired reads as None), result_sweep (honker_ops.rs:1850-1900). Python adds get_result -> (found, value) and wait_result (_honker.py:431-498); the Python worker saves results before ack and swallows result-persistence failures so side effects aren't duplicated (_worker.py:75-83).

7. Maintenance defaults: what runs automatically?

Nothing. Everything is consumer-driven. No background sweeper thread in honker-core (only the UpdateWatcher data-version poller, lib.rs:942-1026). sweep_expired, result_sweep, rate_limit_sweep, prune_notifications*, scheduler_tick, lock_renew all invoked by user code; the scheduler tick only enqueues. Implicit maintenance: (a) pre-claim dead-lettering inside claim_batch (honker_ops.rs :855-859), (b) opportunistic expired-lock deletion inside lock_acquire (honker_ops.rs:1393-1397). Python workers get idle-poll fallback (default 5s) alongside data_version wake (_honker.py:374-385).

8. Observed defects / brittleness (feeds OQ-06)

  1. Doc drift in the Rust binding's sweep_expired: "Sweep expired processing rows back to pending" (lib.rs:652) — actual semantics is pending → _honker_dead (last_error='expired'), never processing→pending (honker_ops.rs:1336-1381).
  2. Zombie processing rows: a processing row whose job-level expires_at passed is unreachable by every path if its worker died — claim predicate requires future expiry (honker_ops.rs:878), pre-claim dead-lettering likewise (honker_ops.rs:792), sweep_expired is pending-only (:1346). Sits in _honker_live forever unless canceled or a still-live worker retries it. Narrow but real.
  3. Silent-boundary skip after a scheduler outage — cap 64 then jump (honker_ops.rs:1612-1621). Deliberate and documented, but part of the guarantee.
  4. Local-timezone cron — chrono::Local hard-wired in the public surface (cron.rs:319-321).
  5. No dead-row retention anywhere — _honker_dead, _honker_notifications, _honker_stream grow unbounded unless the user schedules pruning (lib.rs:345-347, _honker.py:88-90).
  6. lock_acquire doesn't refresh TTL for same-owner re-acquire (honker_ops.rs:1398-1401); _honker_locks has no standalone sweep.
  7. Backoff has no jitter/cap (Python plain 2**(attempts-1), _honker.py:930; Python no overflow guard, Rust saturates).
  8. Reclaim consumes the attempt budget (honker_ops.rs:870-872) — coherent with dead-lettering but a footgun ("timeout ≠ attempt" assumptions are wrong).
  9. Heartbeat refusal after expiry leaves an at-least-once dual-execution window by design (honker_ops.rs:1312-1321) — downstream idempotency is the user's job; docs admit it (_scheduler.py:122-133).
  10. get_job cannot see dead jobs (honker_ops.rs:1237-1240) — post-mortem diagnosis is SQL-only.
  11. Savepoint-class defects (issue #133, post-0.5.1 hardening): any DELETE ... RETURNING + decode + INSERT can strand rows in neither table on a mid-flight error; fixed recently with in_savepoint (honker_ops.rs:90-117, tests :2340-2700).
  12. ack/ack_batch don't check state='processing' (honker_ops.rs:922-926) — implicit rather than enforced.

9. Patterns worth carrying over

  • One shared SQL-scalar implementation for all bindings, tested once.
  • Partial indexes matching the claim hot path; dead rows out of it.
  • unixepoch() second-precision timestamps (single clock source); simple, multiprocess-safe.
  • Wake = PRAGMA data_version delta, coalesced; no wake rows in the DB (the "no synthetic notification" fix appears three times: honker_ops.rs:1002, 1119-1120, 1553-1562).
  • queue_next_claim_at as the sleep-computing claim counterpart (zero-poll idle workers).
  • Savepoint-safe multi-statement mutations.

10. Aftermath — where these findings landed

OQ-06 consumed this document and went further: a follow-on quality read (quality-read-honker-core.md) cross-checked the published crates.io artifact against this checkout and found the savepoint/error-propagation fix train (§8 items 11, the CHANGELOG "Unreleased" work) unpublished — firing ADR-005's fork trigger (ADR-011). The defect list in §8 is now the fork's fix/inheritance register; this document remains the accurate read of the reference revision.