From 969ed44e21263000cd1948a77e99087f1d50656a Mon Sep 17 00:00:00 2001 From: kjgbot Date: Thu, 3 Sep 2026 14:19:36 +0200 Subject: [PATCH 1/2] =?UTF-8?q?feat(sdk):=20a=20relayflow=20can=20be=20sch?= =?UTF-8?q?eduled=20=E2=80=94=20tick=20event=20source?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A relayflow could not be scheduled. `grep -rniE "cron|schedule|interval|timer"` over sdk/src/spec.ts and sdk/src/compile.ts returned nothing; every shipped trigger was fed by a poller reacting to something external. RFC-0001 already decided the shape, so a schedule is an EVENT SOURCE, not a kernel feature: gate 2 "proves triggers are entry conditions, not schedulers", and the 2026-08-27 dogfood run is on record for "a cron trigger reported `succeeded` into a void with no worker enrolled". There is no `cron:` field on the spec and no Rust changed. sdk/src/tick-source.ts sits beside hn-poller and dir-watcher-poller and submits `flows.tick` through the same event.submit path, so it inherits rather than reimplements the kernel's dedupe claim and its liveness sweep. Duplicates and skips are separate mechanisms, and neither covers for the other: - The dedupe key is derived from (schedule_id, scheduled_for_ms) — the slot's grid instant, never wall-clock-at-emit. The kernel's (flow_key, subscription_id, dedupe_key) claim then makes a double-fire, a re-delivery, two racing pollers, and a poller restart with a lost cursor all idempotent. Restart-idempotency belongs to the key, not to the cursor. - The cursor emits every slot between the last emitted and now, so a sleeping poller backfills instead of silently dropping slots. It advances only after a successful submit — dir-watcher's `seen` discipline. Slots beyond maxCatchUp are reported in `skippedSlots`, not dropped quietly. Liveness: the kernel already implements the RFC's RelayCron claim + stale_after reconciliation; what was missing was the authoring half. `staleAfterMs` was unauthorable through the SDK, so every flow silently inherited the 5-minute engine default. It is now declarable and bounds-checked against the same i64 limit the kernel enforces. Not closed: a schedule provisioned but never fired is still undetected (the kernel's own documented gap), and nothing restarts a dead schedule. Also fixes a pre-existing P0 this work could not proceed without: toKernelSpec spread `flow.triggers` through untouched, so every event subscription reached the kernel in camelCase and relayflowd — deny_unknown_fields over snake_case — refused the spec outright: malformed run spec: unknown field `dedupeKeyTemplate`, expected one of `id`, `executor`, `event_type`, `pattern`, `dedupe_key_template`, `stale_after_ms` Every event-triggered flow in testdata/ was unauthorable through the supported SDK path; the committed snake_case fixtures hid it because no test compiled a triggered flow and compared it to one. The fixed compiler now reproduces those fixtures byte-for-byte, and spec-parity pins all of them. Worked example: testdata/tick-heartbeat.flow.yaml with fixtures generated through the SDK compiler. Four missed slots produce four distinct successful runs, each reporting its own grid instant and lag; a re-delivery produces none. Full evidence, both gate baselines and mutation verification in ops/reviews/20260903-scheduled-trigger-design.md. Gates (RELAYFLOWD_BIN pinned to this worktree's build): tsc --noEmit 0 -> 0 tsc -p tsconfig.tests.json 0 -> 0 vitest run 370 passed/3 skipped -> 408 passed/3 skipped cargo test --workspace 100 passed -> 100 passed (identical test set) 38 test names added, 0 removed. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01FtQSAcGDta5VH9xiZFT4sR Session-Id: c228933d-4f94-4d83-9a9a-daf3c83b94f1 --- .../20260903-scheduled-trigger-design.md | 459 ++++++++++++++++++ sdk/src/compile.ts | 63 ++- sdk/src/index.ts | 19 + sdk/src/spec.ts | 37 +- sdk/src/tick-source.ts | 253 ++++++++++ sdk/src/validate.ts | 35 ++ sdk/tests/live-kernel.test.ts | 200 ++++++++ sdk/tests/spec-parity.test.ts | 35 ++ sdk/tests/tick-source.test.ts | 309 ++++++++++++ sdk/tests/validate.test.ts | 29 ++ testdata/preflight/tick-slot-report-cli | 29 ++ testdata/tick-heartbeat.flow.yaml | 56 +++ testdata/tick-heartbeat.spec.canonical.json | 1 + testdata/tick-heartbeat.spec.sha256 | 1 + 14 files changed, 1523 insertions(+), 3 deletions(-) create mode 100644 ops/reviews/20260903-scheduled-trigger-design.md create mode 100644 sdk/src/tick-source.ts create mode 100644 sdk/tests/tick-source.test.ts create mode 100755 testdata/preflight/tick-slot-report-cli create mode 100644 testdata/tick-heartbeat.flow.yaml create mode 100644 testdata/tick-heartbeat.spec.canonical.json create mode 100644 testdata/tick-heartbeat.spec.sha256 diff --git a/ops/reviews/20260903-scheduled-trigger-design.md b/ops/reviews/20260903-scheduled-trigger-design.md new file mode 100644 index 000000000..093b80abe --- /dev/null +++ b/ops/reviews/20260903-scheduled-trigger-design.md @@ -0,0 +1,459 @@ +# Scheduled triggers: a relayflow can be scheduled + +2026-09-03 · branch `feat/scheduled-trigger-0903` · base `origin/main` `990093b` + +A relayflow could not be scheduled. `grep -rniE "cron|schedule|interval|timer"` +over `sdk/src/spec.ts` and `sdk/src/compile.ts` returned nothing; every shipped +trigger is fed by a poller reacting to something external (`hn-poller.ts`, +`dir-watcher-poller.ts`); `agent-relay cloud schedules` reports "No workflow +schedules found." "Run this flow every ten minutes" had no expression anywhere +in the authoring dialect. + +This PR builds that primitive. **No Rust changed.** + +--- + +## 1. The shape, and why + +RFC-0001 had already decided it, so the work was to follow the decision rather +than make one: + +- gate 2 "**proves: triggers are entry conditions, not schedulers**"; +- the first dogfood run, 2026-08-27: "a cron trigger reported `succeeded` into + a void with no worker enrolled" — a scheduler *inside* the kernel is exactly + what produced that; +- "the trigger plane is **liveness-checked** … **because a flow that is never + triggered is silently zero — Native's silent-death problem**." + +So a schedule is an **event source**, not a kernel feature. There is no `cron:` +field on the spec. `sdk/src/tick-source.ts` sits beside the directory watcher +and the HN poller, submits `flows.tick` through the same `event.submit` path +every other source uses, and thereby inherits — rather than reimplements — the +kernel's dedupe claim and its liveness sweep. + +The kernel learns about time the way it learns about everything else: as an +event. + +### The slot grid + +Time is divided into fixed slots anchored at a declared `epochMs`: + +``` +slot(t) = floor((t - epochMs) / intervalMs) +scheduledForMs(n) = epochMs + n * intervalMs +``` + +A slot is an interval **of the grid**, not the moment a poller happened to wake +up. That distinction is the whole design: `scheduledForMs` is a pure function of +the grid, so one scheduled instant has one identity no matter when — or how many +times — a poller notices it. + +Payload of one tick: + +```json +{"schedule_id":"heartbeat-1m","slot":29400001,"scheduled_for_ms":1764000060000, + "interval_ms":60000,"emitted_at_ms":1764000257000,"lag_ms":197000} +``` + +`lag_ms` exists so a backfilled run can tell that it is running for a slot from +the past rather than for now. + +### Why `intervalMs` and not cron syntax + +Not a shortcut — a boundary. A cron expression is a *surface* concern +(RFC settled decision 13: "the kernel vocabulary is closed; the surface is +open"). `TickSchedule` is the primitive a cron parser would compile *to*: any +expression that can name a sequence of instants can drive `emitDueTicks`. +Shipping a parser now would have been the speculative abstraction AGENTS.md §6 +forbids. Not built, deliberately. + +--- + +## 2. Dedupe by construction + +**Duplicates and skips are two different failure modes with two different +mechanisms, and neither covers for the other.** Getting that split right is what +makes the primitive honest. + +### Duplicates — the dedupe key + +The key is derived from `(schedule_id, scheduled_for_ms)` through the flow's +declared template: + +``` +dedupeKeyTemplate: '{{event.type}}:{{payload.schedule_id}}:{{payload.scheduled_for_ms}}' +``` + +`scheduled_for_ms` — **not** `emitted_at_ms`, which is carried in the payload for +observability and deliberately excluded from the key. The kernel then applies its +existing `(flow_key, subscription_id, dedupe_key)` claim +(`kernel/relayflowd-journal/src/registry.rs`), so the second delivery of a slot +returns `deduped: true, run: null` and creates no journal. + +This one mechanism covers every duplicate shape: + +| Scenario | Why it is idempotent | +|---|---| +| Double-fire inside one slot | Both emissions derive the same key | +| Two pollers racing | Same grid, same key, kernel claim picks one | +| Re-delivery of an old tick | Key is a function of the slot, not of now | +| **Poller restart with a lost cursor** | Re-emitting slot N derives N's key again | + +The last row is the important one: **restart-idempotency belongs to the dedupe +key, not to the cursor.** A caller that never persists the cursor still cannot +double-run a slot. It only loses backfill. + +### Skips — the cursor + +`emitDueTicks` emits **every** slot between the last one it emitted and now, not +just the current one. A poller asleep across three slots backfills three ticks +into three distinct runs rather than silently dropping two. + +The cursor advances **only after a successful submit**, exactly as +`pollDirectoryOnce` adds to its `seen` set only after a successful submit. A +journal failure mid-backfill leaves the remaining slots due, so the next poll +retries them. The cursor is a plain JSON object owned by the caller, precisely so +a caller that wants restart-safe backfill can persist it. + +### The catch-up bound + +A poller down for a week on a one-minute schedule has ~10,000 outstanding slots. +Replaying all of them is a stampede, not a recovery. `maxCatchUp` (default 60) +caps one poll's backfill to the **newest** slots — the current slot's work is the +relevant work, the oldest is the most stale. + +Slots beyond the bound are **returned in `skippedSlots`**, not dropped quietly. +They are returned rather than thrown so a week-long outage recovers to the +current slot instead of wedging, and returned rather than ignored so the skip +cannot be silent (AGENTS.md §4, fail closed / no silent fallbacks). The cursor +advances past them so each skip is reported exactly once instead of forever. + +--- + +## 3. Liveness — what is guaranteed, and what is not + +### What I provide + +The kernel **already has** a full trigger-liveness plane and I did not need to +build one — `kernel/relayflowd/src/server/liveness.rs` implements exactly the +RelayCron deterministic-id claim + `stale_after` reconciliation the RFC names +(`sweep_claims` single-winner bucket election, `detect_stale` → journal → latch +ordering, CAS-guarded latch, at-least-once-per-silence). + +What was missing was the **authoring half**. The kernel's `TriggerSpec` has +carried `stale_after_ms` all along, but the SDK's `TriggerSpec`, `TRIGGER_KEYS` +and `KernelTriggerSpec` did not, so **a flow author could not declare a silence +budget through the supported path**. Every flow silently inherited the 5-minute +engine default — a decision no author made. + +This PR adds `staleAfterMs` to the authoring dialect, validates it against the +same i64 bound the kernel enforces (`TriggerSpec::effective_stale_after_ms`), and +`testdata/tick-heartbeat.flow.yaml` declares three slots' worth of it. + +Observed end to end against a real daemon (§5.3): a tick schedule that stops +firing produces a journaled `subscription.stale` entry **and** a greppable +stderr line naming the flow, the subscription, the last event time and the +budget. "This schedule last fired at T" is recorded in the `subscriptions` table +by every successful match, deduped or not. + +### What I did NOT do — stated plainly + +1. **Provisioned-but-never-fired is still undetected.** A schedule whose source + never submits a single tick has no `subscriptions` row and no `event_dedupe` + row, so the sweep sees nothing. This is the kernel's own documented "Known + gap" (`server/liveness.rs`), it predates this PR, and closing it means + pre-registering spec triggers at spec-observation time — a kernel change, + which this PR is scoped out of. **A tick source that dies after firing at + least once is detected; one that never starts is not.** +2. **Nothing restarts a dead schedule.** Detection is journaled and logged; the + sweep does not re-provision, page, or escalate. The observable hook is left, + the action is not built. +3. **No `flows tick start` CLI.** `emitDueTicks` is a pure function of + `(schedule, cursor, now)` — the timer that calls it is the caller's. Adding a + `flows tick` subcommand alongside `flows hn-monitor start` is the obvious next + rung and is outside the four items I was given. Consequence: nothing in-repo + currently *runs* a schedule continuously, so `agent-relay cloud schedules` + will still report nothing after this merges. **This PR builds the primitive + and proves it; it does not put a schedule into production.** + +--- + +## 4. The unavoidable defect I had to fix first + +`toKernelSpec` did **not** lower trigger keys to the kernel dialect. It spread +`flow.triggers` through untouched, so every event subscription reached the kernel +in camelCase — and `relayflowd`'s `TriggerSpec` is `#[serde(deny_unknown_fields)]` +over snake_case. Captured, against the committed fixture, on **unmodified +`origin/main`**: + +``` +$ node -e "... compileYamlToCanonicalJson(testdata/event-triggered-flow.yaml) ..." +SDK canonical: ...,"triggers":[{"dedupeKeyTemplate":"{{event.type}}:{{payload.message}}","eventType":"test.ping",...}],... +fixture : ...,"triggers":[{"dedupe_key_template":"{{event.type}}:{{payload.message}}","event_type":"test.ping",...}],... + +$ relayflowd --data-dir $D run $D/sdk-compiled.json +Error: parse run spec /tmp/tickproof.3ndq/sdk-compiled.json + +Caused by: + malformed run spec: unknown field `dedupeKeyTemplate`, expected one of `id`, `executor`, `event_type`, `pattern`, `dedupe_key_template`, `stale_after_ms` + +$ relayflowd --data-dir $D2 run testdata/event-triggered-flow.spec.canonical.json +{"run_id":"01M1KJGFDKCYS08H53QQ39Z1E3","status":"parked","completion_reason":null,"completed_steps":0} +``` + +**Every event-triggered flow in `testdata/` was unauthorable through the +supported SDK path.** It went unnoticed because the committed +`*.spec.canonical.json` fixtures are snake_case (kernel-produced or hand-written) +and `spec-parity.test.ts` only exercised the four `hello-*` fixtures, none of +which has a trigger. + +This was unavoidable for item 4, which requires the tick fixtures be generated +**through the SDK compiler**. It is SDK-side, not kernel-side. The fix is +`toKernelTrigger` / `kernelTriggerToAuthoring` in `sdk/src/compile.ts`. + +The strongest evidence it is correct: the fixed compiler now reproduces the +pre-existing snake_case fixtures **byte-for-byte** — those fixtures were the +kernel's truth all along, and the compiler simply was not producing them: + +``` +event-triggered-flow.yaml MATCHES committed fixture +hn-monitor.flow.yaml MATCHES committed fixture +``` + +`spec-parity.test.ts` now pins all three (plus the new tick fixture) so this +cannot silently regress again. + +`sdk/src/spec.ts`'s `TriggerSpec` also said "Inert gate-1 trigger declaration" +with only `id` and `executor` while three shipped specs used `eventType`, +`pattern` and `dedupeKeyTemplate`. It now describes what actually ships. + +--- + +## 5. The worked example, run for real + +`testdata/tick-heartbeat.flow.yaml` + `.spec.canonical.json` + `.spec.sha256`, +following `hn-monitor` and `dir-watcher` exactly. Both fixtures generated through +`compileYamlToCanonicalJson` / `compileAndHash` from `sdk/dist/compile.js` — never +by hand. `spec_hash` = `0c3d089f0075c53442c5c2241ada734af5e450a2b558c16030aaf1d1e29127ff`. + +`testdata/preflight/tick-slot-report-cli` is the step's agent CLI, modelled on +`wake-context-probe-cli`: it reads the tick out of `$RELAYFLOW_WAKE_CONTEXT` and +emits the JSON the flow's `json_schema` gate requires. Deterministic on purpose — +the subject under test is the trigger plane, not a model. + +### 5.1 A real run: four backfilled slots → four runs, then a deduped re-delivery + +Poller asleep across slots 29400001–29400003, wakes 17s into 29400004; then a +second poller re-delivers 29400004 with a lost cursor. + +``` +$ relayflowd --data-dir /tmp/tickdemo.kTUR serve & +$ node /tmp/tick-worked-example.mjs /tmp/tickdemo.kTUR + +emittedSlots: [29400001,29400002,29400003,29400004] +skippedSlots: [] + slot 29400001 -> matched=true deduped=false run=01M1KK6XV9M5DMJQH3V7HPK8AJ + slot 29400002 -> matched=true deduped=false run=01M1KK6XVJT7CMM7T5H5SZSXYA + slot 29400003 -> matched=true deduped=false run=01M1KK6XVR8TTFQ7JH5KKEG7V3 + slot 29400004 -> matched=true deduped=false run=01M1KK6XVYTT6N791BWEX8D1FH +replay of slot 29400004: matched=true deduped=true run=null +run 01M1KK6XV9M5DMJQH3V7HPK8AJ slot 29400001: reason=success output={"lag_ms":197000,"schedule_id":"heartbeat-1m","scheduled_for_ms":1764000060000,"slot":29400001} +run 01M1KK6XVJT7CMM7T5H5SZSXYA slot 29400002: reason=success output={"lag_ms":137000,"schedule_id":"heartbeat-1m","scheduled_for_ms":1764000120000,"slot":29400002} +run 01M1KK6XVR8TTFQ7JH5KKEG7V3 slot 29400003: reason=success output={"lag_ms":77000,"schedule_id":"heartbeat-1m","scheduled_for_ms":1764000180000,"slot":29400003} +run 01M1KK6XVYTT6N791BWEX8D1FH slot 29400004: reason=success output={"lag_ms":17000,"schedule_id":"heartbeat-1m","scheduled_for_ms":1764000240000,"slot":29400004} +``` + +Four missed slots produced four distinct successful runs, each reporting its own +grid instant and its own lag; the re-delivery produced **zero**. The script is +reproduced in §7 so this is re-runnable. + +### 5.2 The dedupe identity in the kernel's own journal + +`event_key` is the scheduled instant, not the emit time: + +``` +subscription.matched {"event_key":"flows.tick:heartbeat-1m:1764000000000","subscription_id":"every-minute", ... + "triggering_event":{"payload":{"emitted_at_ms":1764000000000,"interval_ms":60000,"lag_ms":0, + "schedule_id":"heartbeat-1m","scheduled_for_ms":1764000000000,"slot":29400000},"type":"flows.tick"}} +``` + +### 5.3 The liveness sweep firing for a tick schedule + +Same flow with `staleAfterMs: 1000`, one tick submitted, then silence. Daemon +stderr: + +``` +relayflowd: subscription.stale flow="b356ed314641457efb46f578b5f3ae23a0ff819a95a0573a8bf881115ee2a79b" sub="every-minute" event_type="flows.tick" last_event_at_ms=1788437770130 stale_after_ms=1000 detected_at_ms=1788437791447 +``` + +And the journal for that run: + +``` +subscription.registered {"effective_stale_after_ms":1000,"event_type":"flows.tick","executor":"agent-worker","subscription_id":"every-minute"} +subscription.matched {"event_key":"flows.tick:heartbeat-1m:1764000000000","subscription_id":"every-minute", ...} +subscription.stale {"detected_at_ms":1788437791447,"event_type":"flows.tick","flow_key":"b356ed31...","last_event_at_ms":1788437770130,"stale_after_ms":1000,"subscription_id":"every-minute"} +``` + +`effective_stale_after_ms: 1000` is the flow's **declared** budget, not the +5-minute engine default — which is what proves `staleAfterMs` survives the +authoring → kernel lowering. + +--- + +## 6. Gates + +Baseline measured on the worktree **before any edit**, at `origin/main` +`990093b`. `RELAYFLOWD_BIN` pinned for every SDK run to +`/Users/khaliqgant/.relayflows-toolchain/target/173824371/debug/relayflowd` +— this worktree's own build (`173824371` is `cksum` of this worktree's path). +Pinned because `locateRelayflowd` otherwise picks the newest daemon by mtime +across every worktree's target tree. + +| Gate | Baseline (`990093b`) | Head | +|---|---|---| +| `tsc --noEmit` | exit 0 | exit 0 | +| `tsc -p tsconfig.tests.json` | exit 0 | exit 0 | +| `vitest run` | 370 passed, 3 skipped, **0 failed** | 408 passed, 3 skipped, **0 failed** | +| `cargo test --workspace` | 100 passed, 0 failed | 100 passed, 0 failed | + +> The first baseline attempt reported 12 failures. Those were an artifact of not +> having run `test:prep` — no `sdk/dist`, no built kernel, no chmod on +> `testdata/preflight/*-cli`. After the documented prep the baseline is clean, and +> that clean run is the one compared above. + +### Per-file test accounting (set-diff, not counts) + +``` ++6 sdk/tests/live-kernel.test.ts (21 -> 27) ++6 sdk/tests/spec-parity.test.ts (15 -> 21) ++6 sdk/tests/validate.test.ts (36 -> 42) ++20 sdk/tests/tick-source.test.ts (0 -> 20) + total 373 -> 411 (370 passed + 3 skipped -> 408 passed + 3 skipped) +``` + +Full-name set-diff: **38 test names added, 0 removed, 0 renamed.** +`git diff --stat sdk/tests/live-kernel.test.ts` is `200 ++++` with **zero +deletions** — the new block is a pure insertion, no existing case touched. + +Kernel: `diff` of the sorted `test … ok` line set between baseline and head is +empty — **identical kernel test set**, consistent with no Rust changing. + +### Not trusting green CI + +flows CI runs only `linux-x64-artifact` and `packed-consumer`; neither the +kernel nor the SDK suite. Everything above is local output. + +--- + +## 7. Mutation verification + +Both mutations: reverted the specific behaviour, ran the specific tests, +captured the failure, restored **byte-for-byte** (`sha256(tick-source.ts)` = +`3edce9029eb6e679ad539b413e2d0251bf3f0cd4a6848b5eea7c00fdc8e87b5e` before and +after both), re-ran, captured the pass. + +### M1 — dedupe key from wall clock instead of the scheduled instant + +`scheduled_for_ms: scheduledFor` → `scheduled_for_ms: nowMs`. Ticks still fire; +only the *bound* is removed. + +``` +× tick source: a double-fire produces ONE run > two emissions of one scheduled instant claim the same key, so one run + → expected 'flows.tick:heartbeat-1m:6000001' to be 'flows.tick:heartbeat-1m:6000000' +× tick source: a double-fire produces ONE run > a poller RESTART that loses its cursor re-emits the slot but does not re-run it + → expected 'flows.tick:heartbeat-1m:6005000' to be 'flows.tick:heartbeat-1m:6000000' +× a relayflow can be scheduled: ... > TWO ticks for ONE scheduled instant produce exactly ONE run + → the kernel spawned a second run for one scheduled instant: expected { deduped: false, matched: true, …(2) } to match object { matched: true, deduped: true, …(1) } +Tests 3 failed | 1 passed | 43 skipped (47) +``` + +The third line is the one that matters: a **real relayflowd** spawned a second +run for one scheduled instant. Restored → `Tests 4 passed | 43 skipped`. + +### M2 — backfill removed (emit only the current slot) + +``` +× backfills every slot a sleeping poller passed over → expected [ 104 ] to deeply equal [ 101, 102, 103, 104 ] +× a backfilled tick carries its own lag ... → expected [ +0 ] to deeply equal [ 120000, 60000, +0 ] +× does NOT advance past a slot whose submit failed ... → promise resolved "{ emittedSlots: [ 104 ], …(2) }" instead of rejecting +× reports slots dropped by the catch-up bound ... → expected [ 110 ] to deeply equal [ 108, 109, 110 ] +× reports each skipped slot exactly once ... → expected [] to deeply equal [ 101, 102, 103, 104 ] +× caps at DEFAULT_MAX_CATCH_UP ... → expected [ 500 ] to have a length of 60 but got 1 +× emits nothing when the current slot has already been emitted → expected [ 100 ] to deeply equal [] +× emits nothing when the clock moves backwards ... → expected [ 90 ] to deeply equal [] +× a MISSED interval is backfilled into its own run ... → expected [ 29400004 ] to deeply equal [ 29400001, 29400002, 29400003, …(1) ] +Tests 9 failed | 38 passed (47) +``` + +Restored → `tests/tick-source.test.ts (20 tests) ✓`, +`tests/live-kernel.test.ts (27 tests) ✓`, `Tests 47 passed (47)`. + +### The tests are bounds, not mechanisms + +Deliberately **not** written: "a tick fired". That passes under M1 and would +have shipped a schedule that double-runs every slot. What is pinned instead: + +- two ticks for one scheduled instant → **one** run (live kernel); +- a poller restart re-emitting a slot → **no** second run (live kernel); +- successive slots are **not** deduped away — the mirror test, without which a + constant key would pass everything above and the flow would run once, ever; +- four missed slots → **four distinct** runs, not one collapsed run; +- a failed submit does **not** advance the cursor past its slot; +- slots beyond the catch-up bound are **reported**, exactly once; +- a tick for a different `schedule_id` **does not wake** the flow; +- the journal carries the **declared** budget, not the engine default. + +--- + +## 8. Reproducing §5.1 + +```js +// /tmp/tick-worked-example.mjs — node /tmp/tick-worked-example.mjs +import { readFileSync } from 'node:fs'; +const SDK = '/sdk/dist', ROOT = ''; +const { JournalClient } = await import(`${SDK}/journal-client.js`); +const { AgentWorker } = await import(`${SDK}/worker.js`); +const { emitDueTicks } = await import(`${SDK}/tick-source.js`); + +const dataDir = process.argv[2]; +const spec = JSON.parse(readFileSync(`${ROOT}/testdata/tick-heartbeat.spec.canonical.json`, 'utf8')); +for (const s of spec.steps) if (s.id === 'report-slot') s.cli = `${ROOT}/testdata/preflight/tick-slot-report-cli`; + +const client = new JournalClient(`${dataDir}/relayflowd.sock`, { requestTimeoutMs: 5000 }); +await client.connect(); await client.hello('tick-worked-example'); +const worker = new AgentWorker(client, { workerId: 'tick-worked-example-worker', + pins: { workspace: [{ surface: 'repo', revision_id: 'rev-a' }], streams: [] } }); +await worker.attach(); + +const schedule = { scheduleId: 'heartbeat-1m', intervalMs: 60_000, epochMs: 0 }; +const cursor = { lastEmittedSlot: 29_400_000 }; +const nowMs = 29_400_004 * 60_000 + 17_000; // asleep 3 slots, wakes 17s into 29400004 +const result = await emitDueTicks(spec, client, { schedule, cursor, nowMs }); +const replay = await emitDueTicks(spec, client, { schedule, cursor: {}, nowMs: nowMs + 9_000 }); +// ... print result.emittedSlots / outcomes, then journalRead each run's step.completed +``` + +--- + +## 9. What I could not verify + +- **Provisioned-but-never-fired liveness.** Not attempted; it needs a kernel + change (§3). The failure mode remains open exactly as + `server/liveness.rs` documents it. +- **A schedule running in production.** No CLI runner ships here, so + `agent-relay cloud schedules` will still report "No workflow schedules found" + after this merges. The primitive is proven; it is not deployed. +- **Behaviour across a daemon restart mid-backfill.** The unit test covers a + failed submit mid-backfill with a fake sink; I did not kill a live daemon + between two slots of one backfill. The dedupe claim is durable in SQLite so I + expect it to hold, but I did not observe it and am not claiming it. +- **Clock skew between two hosts running the same schedule.** Two pollers on the + same grid dedupe correctly (proved), but I did not test hosts whose clocks + disagree by more than one interval — that would put them in different slots + and produce two runs. Real, unaddressed, out of scope here. +- **`maxCatchUp = 60` as the right default.** Chosen as an hour of one-minute + slots. Not empirically justified. +- **`hn-monitor.flow.yaml` round-trip** through `kernelToAuthoring(toKernelSpec(…))` + differs under `JSON.stringify` but is equal under `toEqual` — key ordering only, + in the `output` → `verification` sugar, unrelated to triggers and unchanged by + this PR. diff --git a/sdk/src/compile.ts b/sdk/src/compile.ts index 35df7b388..18c31a146 100644 --- a/sdk/src/compile.ts +++ b/sdk/src/compile.ts @@ -23,11 +23,13 @@ import type { KernelRunSpec, KernelStepCommon, KernelStepSpec, + KernelTriggerSpec, KernelVerificationSpec, LlmStepSpec, NamedAgentSpec, StepSpec, StepType, + TriggerSpec, } from './spec.js'; import { SPEC_SCHEMA_VERSION } from './spec.js'; import { canonicalize, specHash } from './canonical.js'; @@ -195,7 +197,7 @@ export function toKernelSpec(flow: FlowSpec): KernelRunSpec { ...(flow.name !== undefined ? { name: flow.name } : {}), ...(flow.description !== undefined ? { description: flow.description } : {}), ...(flow.cli !== undefined ? { cli: flow.cli } : {}), - ...(flow.triggers?.length ? { triggers: flow.triggers } : {}), + ...(flow.triggers?.length ? { triggers: flow.triggers.map(toKernelTrigger) } : {}), steps: flow.steps.map((step) => toKernelStep(resolveNamedAgent(step, flow.agents))), ...(flow.budget !== undefined ? { @@ -222,8 +224,15 @@ export function kernelToAuthoring(value: unknown): unknown { ); const steps = requireKernelArray(root['steps'], 'spec.steps') .map((step, index) => kernelStepToAuthoring(step, `spec.steps[${index}]`)); + const triggers = root['triggers']; return { - ...copyDefined(root, ['version', 'name', 'description', 'cli', 'triggers']), + ...copyDefined(root, ['version', 'name', 'description', 'cli']), + ...(triggers !== undefined + ? { + triggers: requireKernelArray(triggers, 'spec.triggers') + .map((trigger, index) => kernelTriggerToAuthoring(trigger, `spec.triggers[${index}]`)), + } + : {}), steps, ...(root['budget'] !== undefined ? { budget: kernelBudgetToAuthoring(root['budget'], 'spec.budget') } @@ -231,6 +240,56 @@ export function kernelToAuthoring(value: unknown): unknown { }; } +/** + * Lower one authoring trigger into the kernel dialect. + * + * This mapping was missing entirely: `toKernelSpec` used to spread + * `flow.triggers` through untouched, so every event subscription reached the + * kernel in camelCase and `relayflowd` — whose `TriggerSpec` is + * `#[serde(deny_unknown_fields)]` over snake_case — refused the spec outright: + * + * malformed run spec: unknown field `dedupeKeyTemplate`, expected one of + * `id`, `executor`, `event_type`, `pattern`, `dedupe_key_template`, + * `stale_after_ms` + * + * The committed `testdata/*.spec.canonical.json` fixtures are snake_case and + * the kernel accepts them, which is why nothing noticed: no test compiled a + * triggered flow through this function and compared it to a fixture. Every + * triggered flow in `testdata/` was therefore unauthorable through the + * supported SDK path. `tests/spec-parity.test.ts` now pins the mapping. + */ +function toKernelTrigger(trigger: TriggerSpec): KernelTriggerSpec { + return { + id: trigger.id, + executor: trigger.executor, + ...(trigger.eventType !== undefined ? { event_type: trigger.eventType } : {}), + ...(trigger.pattern !== undefined ? { pattern: trigger.pattern } : {}), + ...(trigger.dedupeKeyTemplate !== undefined + ? { dedupe_key_template: trigger.dedupeKeyTemplate } + : {}), + ...(trigger.staleAfterMs !== undefined ? { stale_after_ms: trigger.staleAfterMs } : {}), + }; +} + +/** Inverse of `toKernelTrigger`. Kernel-only keys are refused, never dropped. */ +function kernelTriggerToAuthoring(value: unknown, at: string): unknown { + const trigger = requireKernelObject( + value, + ['id', 'executor', 'event_type', 'pattern', 'dedupe_key_template', 'stale_after_ms'], + at, + ); + return { + id: trigger['id'], + executor: trigger['executor'], + ...(trigger['event_type'] !== undefined ? { eventType: trigger['event_type'] } : {}), + ...(trigger['pattern'] !== undefined ? { pattern: trigger['pattern'] } : {}), + ...(trigger['dedupe_key_template'] !== undefined + ? { dedupeKeyTemplate: trigger['dedupe_key_template'] } + : {}), + ...(trigger['stale_after_ms'] !== undefined ? { staleAfterMs: trigger['stale_after_ms'] } : {}), + }; +} + function kernelStepToAuthoring(value: unknown, at: string): unknown { const unionKeys = [ 'id', 'type', 'depends_on', 'max_iterations', 'retry', 'verification', diff --git a/sdk/src/index.ts b/sdk/src/index.ts index 44f587d51..7f147660a 100644 --- a/sdk/src/index.ts +++ b/sdk/src/index.ts @@ -162,3 +162,22 @@ export { type DirLister, type PollOptions as DirWatcherPollOptions, } from './dir-watcher-poller.js'; + +// Tick event source — a relayflow can be scheduled. A schedule is an event +// source subject to the same liveness sweep as any other subscription, not a +// scheduler inside the kernel (RFC-0001 gate 2: "triggers are entry +// conditions, not schedulers"). +export { + emitDueTicks, + scheduledForMs, + slotFor, + tickDedupeKey, + DEFAULT_MAX_CATCH_UP, + TICK_DEDUPE_KEY_TEMPLATE, + TICK_EVENT_TYPE, + type EventSink as TickEventSink, + type TickCursor, + type TickEmitResult, + type TickPayload, + type TickSchedule, +} from './tick-source.js'; diff --git a/sdk/src/spec.ts b/sdk/src/spec.ts index 76d0a9e83..d56a00fdb 100644 --- a/sdk/src/spec.ts +++ b/sdk/src/spec.ts @@ -178,11 +178,42 @@ export interface NamedAgentSpec { model: string; } -/** Inert gate-1 trigger declaration. Matching and dispatch belong to gate 2. */ +/** + * A trigger is an entry condition, not a scheduler (RFC-0001 gate 2). It names + * the event type that wakes the flow, the payload subset that must match, the + * template that derives the dedupe key, and the silence budget after which the + * kernel's liveness sweep declares the subscription dead. + * + * The event-subscription fields were shipped in `testdata/` long before this + * interface described them: `hn-monitor.flow.yaml`, `dir-watcher.flow.yaml` + * and `event-triggered-flow.yaml` all carry `eventType`, `pattern` and + * `dedupeKeyTemplate`, and `validate.ts` has always accepted them. The type + * still said "inert gate-1 declaration" with only `id` and `executor`, so the + * authoring dialect disagreed with both the shipped specs and the kernel. + */ export interface TriggerSpec { id: string; /** Executor registration required before this trigger may start a run. */ executor: string; + /** Event type this trigger subscribes to. Lowers to `event_type`. */ + eventType?: string; + /** Recursive-subset match against the event payload. Lowers to `pattern`. */ + pattern?: Record; + /** Derives the dedupe key. Lowers to `dedupe_key_template`. */ + dedupeKeyTemplate?: string; + /** + * Silence budget in milliseconds. When no matching event arrives inside it, + * the kernel's liveness sweep journals `subscription.stale` and emits a + * `relayflowd: subscription.stale ...` line + * (`kernel/relayflowd/src/server/liveness.rs`). Omitted means the engine + * default (`DEFAULT_STALE_AFTER_MS`, 5 minutes) applies — which is a + * decision the author did not make, not the absence of a budget. + * + * A flow that is never triggered is silently zero (RFC-0001 §"Trigger + * liveness"), so declaring this is how a schedule stops being able to die + * quietly. Lowers to `stale_after_ms`. + */ + staleAfterMs?: number; } /** @@ -290,6 +321,10 @@ export interface KernelBudgetSpec { export interface KernelTriggerSpec { id: string; executor: string; + event_type?: string; + pattern?: Record; + dedupe_key_template?: string; + stale_after_ms?: number; } /** The compiled spec as the kernel parses, journals, and hashes it. */ diff --git a/sdk/src/tick-source.ts b/sdk/src/tick-source.ts new file mode 100644 index 000000000..dc7a94e63 --- /dev/null +++ b/sdk/src/tick-source.ts @@ -0,0 +1,253 @@ +/** + * Scheduled ticks -> relayflow events. + * + * A relayflow could not be scheduled. `grep -rniE "cron|schedule|interval"` over + * `spec.ts` and `compile.ts` returned nothing, and every shipped trigger is fed + * by a poller reacting to something external (`hn-poller.ts`, + * `dir-watcher-poller.ts`). "Run this flow every ten minutes" had no expression. + * + * RFC-0001 already decided the shape, so this is deliberately NOT a `cron:` + * field on the kernel spec: + * + * - gate 2 "proves: triggers are entry conditions, not schedulers"; + * - the first dogfood run (2026-08-27) is on record for "a cron trigger + * reported `succeeded` into a void with no worker enrolled" — a scheduler + * inside the kernel is exactly what produced that; + * - the trigger plane is liveness-checked "because a flow that is never + * triggered is silently zero — Native's silent-death problem". + * + * So a schedule is an EVENT SOURCE, sitting beside the directory watcher and + * the HN poller, speaking the same `event.submit` path, and subject to the same + * liveness sweep as any other subscription. The kernel learns about time the + * way it learns about everything else: as an event. + * + * ## The slot grid + * + * Time is divided into fixed slots anchored at `epochMs`: + * + * slot(t) = floor((t - epochMs) / intervalMs) + * scheduledForMs(n) = epochMs + n * intervalMs + * + * A slot is an interval of the grid, not a moment the poller happened to wake + * up. That distinction is the whole design. `scheduledForMs` is a pure function + * of the grid, so the same slot has the same identity no matter when — or how + * many times — a poller notices it. Wall-clock-at-emit would give two different + * identities to one scheduled instant and produce two runs. + * + * ## Two failure modes, two different mechanisms + * + * They are separate on purpose, and neither one covers for the other: + * + * - **Duplicates** are prevented by the dedupe key, which is derived from + * `(schedule_id, scheduled_for_ms)` through the flow's `dedupeKeyTemplate`. + * The kernel's `(flow_key, subscription_id, dedupe_key)` claim then makes + * the second delivery of a slot a no-op. This holds for a double-fire, a + * re-delivery, two pollers racing, and a poller that restarts with a lost + * cursor and re-emits a slot it already emitted. + * + * - **Skips** are prevented by the cursor. `emitDueTicks` emits every slot + * between the last one it emitted and now, not just the current one, so a + * poller that was asleep across three slots backfills three ticks rather + * than silently dropping two. The cursor advances only after a successful + * submit, so a journal failure mid-backfill leaves the rest for the next + * poll — the same discipline `dir-watcher-poller` applies to its `seen` set. + * + * The cursor is caller-owned (a plain JSON-serializable object) precisely so a + * caller that wants restart-safe backfill can persist it. A caller that does + * not persist it loses backfill across a restart but CANNOT double-run a slot, + * because that bound belongs to the dedupe key rather than to the cursor. + */ + +/** Anything that can submit an event through the journal protocol. */ +export interface EventSink { + eventSubmit(spec: unknown, event: { type: string; payload?: unknown; key?: string }): Promise; +} + +/** The event type every tick carries. */ +export const TICK_EVENT_TYPE = 'flows.tick'; + +/** + * The dedupe key template a tick-triggered flow must declare. Exported so a + * spec and this source cannot drift: `testdata/tick-heartbeat.flow.yaml` uses + * this exact string and `tests/tick-source.test.ts` asserts they match. + * + * `scheduled_for_ms` — not `emitted_at_ms` — is what makes the key idempotent. + */ +export const TICK_DEDUPE_KEY_TEMPLATE = + '{{event.type}}:{{payload.schedule_id}}:{{payload.scheduled_for_ms}}'; + +/** Default bound on how many missed slots one poll will backfill. */ +export const DEFAULT_MAX_CATCH_UP = 60; + +/** A declared schedule. Pure data — no timers, no I/O, no ambient clock. */ +export interface TickSchedule { + /** + * Stable identity of this schedule. It is half the dedupe key, so changing + * it re-runs every slot; two schedules on one flow must differ here. + */ + scheduleId: string; + /** Slot width in milliseconds. Must be a positive integer. */ + intervalMs: number; + /** + * Grid anchor. Slots are measured from here, so this is what decides whether + * an hourly schedule fires on the hour or at seven minutes past. Defaults to + * 0 (the Unix epoch), which puts an hourly schedule on the hour in UTC. + */ + epochMs?: number; + /** + * Upper bound on slots backfilled in a single poll. A poller down for a week + * on a one-minute schedule has ten thousand outstanding slots, and replaying + * all of them would be a stampede, not a recovery. Slots beyond the bound are + * REPORTED in the result rather than dropped quietly (see `TickEmitResult`). + */ + maxCatchUp?: number; +} + +/** + * Caller-owned cursor. Plain JSON so a caller can persist it across restarts. + * `emitDueTicks` mutates it in place, exactly as `pollDirectoryOnce` mutates + * its `seen` set. + */ +export interface TickCursor { + /** Highest slot successfully submitted, or undefined before the first poll. */ + lastEmittedSlot?: number; +} + +/** The payload of one `flows.tick` event. */ +export interface TickPayload { + schedule_id: string; + /** Slot index on the grid. Monotonic, and stable across restarts. */ + slot: number; + /** The scheduled instant this tick stands for. The dedupe identity. */ + scheduled_for_ms: number; + interval_ms: number; + /** + * Wall clock when the tick was submitted. Observability only — it is + * deliberately NOT part of the dedupe key, because it differs between a + * first delivery and a re-delivery of the same slot. + */ + emitted_at_ms: number; + /** + * How far behind the grid this emission was, in milliseconds + * (`emitted_at_ms - scheduled_for_ms`). A catch-up tick carries a large + * value; a punctual one carries roughly zero. Lets a flow tell "I am running + * for a slot from an hour ago" from "I am running for now". + */ + lag_ms: number; +} + +/** What one poll did, including what it deliberately did not do. */ +export interface TickEmitResult { + /** Slots submitted this poll, oldest first. */ + emittedSlots: number[]; + /** The submit outcomes, index-aligned with `emittedSlots`. */ + outcomes: unknown[]; + /** + * Slots that were due but fell outside `maxCatchUp`, oldest first. Non-empty + * means real scheduled work was passed over: the caller MUST surface it. It + * is returned rather than thrown so a poller that was down for a week still + * recovers to the current slot instead of wedging, and it is returned rather + * than ignored so the skip cannot be silent. + */ + skippedSlots: number[]; +} + +/** Slot index containing `nowMs` on this schedule's grid. */ +export function slotFor(schedule: TickSchedule, nowMs: number): number { + return Math.floor((nowMs - (schedule.epochMs ?? 0)) / schedule.intervalMs); +} + +/** The scheduled instant of a slot. Pure function of the grid, never of `now`. */ +export function scheduledForMs(schedule: TickSchedule, slot: number): number { + return (schedule.epochMs ?? 0) + slot * schedule.intervalMs; +} + +/** + * The dedupe key the kernel will derive for a slot, computed here so a test can + * assert the identity directly without a live kernel. Kept in lockstep with + * `TICK_DEDUPE_KEY_TEMPLATE`; `tests/tick-source.test.ts` pins the agreement. + */ +export function tickDedupeKey(schedule: TickSchedule, slot: number): string { + return `${TICK_EVENT_TYPE}:${schedule.scheduleId}:${scheduledForMs(schedule, slot)}`; +} + +function requirePositiveInteger(value: number, field: string): void { + if (!Number.isInteger(value) || value <= 0) { + throw new Error(`tick schedule: ${field} must be a positive integer, got ${String(value)}`); + } +} + +/** + * Submit a `flows.tick` event for every slot that has come due since the cursor + * last advanced, and move the cursor. + * + * The FIRST poll on a fresh cursor emits only the current slot. Backfilling + * from the grid anchor instead would replay every slot since the Unix epoch on + * the first tick of a new schedule. + * + * Journal errors from `eventSubmit` propagate, with the cursor left pointing at + * the last slot that actually reached the kernel. + */ +export async function emitDueTicks( + spec: unknown, + sink: EventSink, + options: { schedule: TickSchedule; cursor: TickCursor; nowMs: number }, +): Promise { + const { schedule, cursor, nowMs } = options; + requirePositiveInteger(schedule.intervalMs, 'intervalMs'); + const maxCatchUp = schedule.maxCatchUp ?? DEFAULT_MAX_CATCH_UP; + requirePositiveInteger(maxCatchUp, 'maxCatchUp'); + if (schedule.scheduleId === '') { + throw new Error('tick schedule: scheduleId must be a non-empty string'); + } + + const currentSlot = slotFor(schedule, nowMs); + const firstDue = cursor.lastEmittedSlot === undefined + ? currentSlot + : cursor.lastEmittedSlot + 1; + + // Clock went backwards, or the cursor is ahead of the grid. Emitting nothing + // is correct: those slots are already claimed, and re-emitting them would be + // deduped anyway. + if (firstDue > currentSlot) { + return { emittedSlots: [], outcomes: [], skippedSlots: [] }; + } + + const due: number[] = []; + for (let slot = firstDue; slot <= currentSlot; slot++) due.push(slot); + + // Over the bound: keep the NEWEST slots. The current slot is the one whose + // work is still relevant; the oldest are the most stale. Report the rest. + const skippedSlots = due.length > maxCatchUp ? due.slice(0, due.length - maxCatchUp) : []; + const toEmit = due.length > maxCatchUp ? due.slice(due.length - maxCatchUp) : due; + + // A skipped slot never reaches the kernel, so nothing downstream would ever + // record it. Advancing the cursor past it here is what makes the skip appear + // exactly once in `skippedSlots` instead of being re-reported forever. + if (skippedSlots.length > 0) { + cursor.lastEmittedSlot = skippedSlots[skippedSlots.length - 1]; + } + + const emittedSlots: number[] = []; + const outcomes: unknown[] = []; + for (const slot of toEmit) { + const scheduledFor = scheduledForMs(schedule, slot); + const payload: TickPayload = { + schedule_id: schedule.scheduleId, + slot, + scheduled_for_ms: scheduledFor, + interval_ms: schedule.intervalMs, + emitted_at_ms: nowMs, + lag_ms: nowMs - scheduledFor, + }; + outcomes.push(await sink.eventSubmit(spec, { type: TICK_EVENT_TYPE, payload })); + // Advance ONLY after the submit succeeded. A journal failure must leave + // this slot due so the next poll retries it — the same rule + // `pollDirectoryOnce` applies to its `seen` set, and the reason a crash + // mid-backfill cannot swallow a slot. + cursor.lastEmittedSlot = slot; + emittedSlots.push(slot); + } + + return { emittedSlots, outcomes, skippedSlots }; +} diff --git a/sdk/src/validate.ts b/sdk/src/validate.ts index 9358c9dcd..19e45d248 100644 --- a/sdk/src/validate.ts +++ b/sdk/src/validate.ts @@ -71,8 +71,18 @@ const TRIGGER_KEYS = [ 'eventType', 'pattern', 'dedupeKeyTemplate', + 'staleAfterMs', ] as const; +/** + * The kernel stores a silence budget as SQLite `INTEGER` and compares it in + * `i64` (`TriggerSpec::effective_stale_after_ms`), refusing anything wider. + * Mirror the bound here so an unrepresentable budget is a compile error rather + * than an engine error at submit time — a budget the sweep cannot represent + * fails OPEN, which is the silent death this field exists to prevent. + */ +const MAX_STALE_AFTER_MS = 9_223_372_036_854_775_807; + class Validator { private errors: string[] = []; private ids = new Set(); @@ -217,6 +227,31 @@ class Validator { if (!isNonEmptyString(candidate.executor)) { this.fail(`${at}.executor: expected a non-empty string`); } + this.validateStaleAfterMs(candidate.staleAfterMs, `${at}.staleAfterMs`); + } + } + + /** + * A silence budget must be a positive, i64-representable whole number of + * milliseconds. Zero is refused rather than treated as "no budget": a + * zero-length budget marks the subscription stale on the very next sweep, + * which reads as a permanently-broken schedule and trains an operator to + * ignore the alert. + */ + private validateStaleAfterMs(value: unknown, at: string): void { + if (value === undefined) return; + if (typeof value !== 'number' || !Number.isInteger(value)) { + this.fail(`${at}: expected an integer number of milliseconds`); + return; + } + if (value <= 0) { + this.fail(`${at}: expected a positive number of milliseconds, got ${value}`); + return; + } + if (value > MAX_STALE_AFTER_MS) { + this.fail( + `${at}: ${value} exceeds the i64 range the kernel's liveness sweep can represent`, + ); } } diff --git a/sdk/tests/live-kernel.test.ts b/sdk/tests/live-kernel.test.ts index fff05665c..577752a4d 100644 --- a/sdk/tests/live-kernel.test.ts +++ b/sdk/tests/live-kernel.test.ts @@ -22,6 +22,7 @@ import { JournalClient } from '../src/journal-client.js'; import type { StepDispatchEvent } from '../src/protocol.js'; import { AgentWorker } from '../src/worker.js'; import { resolveSpecCliPaths } from '../src/cli/hn-monitor.js'; +import { emitDueTicks, type TickCursor } from '../src/tick-source.js'; const ROOT = join(dirname(fileURLToPath(import.meta.url)), '..', '..'); const SDK = join(ROOT, 'sdk'); @@ -1360,6 +1361,194 @@ steps: }); }); +describe('a relayflow can be scheduled: tick source against live relayflowd', () => { + // The worked example. The primitive under test is sdk/src/tick-source.ts; + // what makes this acceptance evidence rather than a plumbing demo is that a + // real relayflowd holds the dedupe claim and spawns (or refuses to spawn) + // the runs. + const TICK_SPEC = () => JSON.parse( + readFileSync(join(TESTDATA, 'tick-heartbeat.spec.canonical.json'), 'utf8'), + ) as Parameters[0] & { steps: { id: string; cli?: string }[] }; + + function specWithReportCli() { + const spec = TICK_SPEC(); + for (const step of spec.steps) { + if (step.id === 'report-slot') step.cli = join(TESTDATA, 'preflight', 'tick-slot-report-cli'); + } + return spec; + } + + const MINUTE = 60_000; + const schedule = { scheduleId: 'heartbeat-1m', intervalMs: MINUTE, epochMs: 0 }; + + it('a tick spawns a real run whose step reports the SCHEDULED instant', async () => { + const dataDir = temporaryDirectory('flows-live-tick-'); + await startDaemon(dataDir); + const client = await connectClient(dataDir); + await client.hello('live-tick'); + const worker = new AgentWorker(client, { + workerId: 'live-tick-worker', + pins: { workspace: [{ surface: 'repo', revision_id: 'rev-a' }], streams: [] }, + }); + await worker.attach(); + + const spec = specWithReportCli(); + const cursor: TickCursor = { lastEmittedSlot: 29_399_999 }; + // Emit 43s into slot 29_400_000. The step must report the slot boundary, + // not 29_400_000 * MINUTE + 43_000. + const nowMs = 29_400_000 * MINUTE + 43_000; + const result = await emitDueTicks(spec, client, { schedule, cursor, nowMs }); + + expect(result.emittedSlots).toEqual([29_400_000]); + expect(result.outcomes[0]).toMatchObject({ matched: true, deduped: false }); + const runId = (result.outcomes[0] as { run: { run_id: string } }).run.run_id; + + expect(await waitForStep(client, runId, 'report-slot', 'done', 20_000)).toMatchObject({ + type: 'agent', + state: 'done', + }); + + const entries = (await client.journalRead(runId)).entries; + const completed = entries.find( + (entry) => isObject(entry) && entry['entry_type'] === 'step.completed' + && entry['step_id'] === 'report-slot', + ) as { payload: { output: Record } } | undefined; + expect(completed, 'the tick-woken step never completed').toBeDefined(); + // The bound: the run reports the grid instant and its own lag, so a + // backfilled run can tell it is running for a slot from the past. + expect(completed!.payload.output).toEqual({ + schedule_id: 'heartbeat-1m', + slot: 29_400_000, + scheduled_for_ms: 29_400_000 * MINUTE, + lag_ms: 43_000, + }); + + await worker.close(); + }, 45_000); + + it('TWO ticks for ONE scheduled instant produce exactly ONE run', async () => { + // The gate. A test asserting "a tick fired" would pass with a dedupe key + // derived from wall clock; this one would not. + const dataDir = temporaryDirectory('flows-live-tick-dedupe-'); + await startDaemon(dataDir); + const client = await connectClient(dataDir); + await client.hello('live-tick-dedupe'); + const spec = specWithReportCli(); + + // Two pollers, independent cursors, same slot, different emit instants. + const first = await emitDueTicks(spec, client, { + schedule, + cursor: { lastEmittedSlot: 29_399_999 }, + nowMs: 29_400_000 * MINUTE + 1_000, + }); + const second = await emitDueTicks(spec, client, { + schedule, + cursor: { lastEmittedSlot: 29_399_999 }, + nowMs: 29_400_000 * MINUTE + 52_000, + }); + + expect(first.outcomes[0]).toMatchObject({ matched: true, deduped: false }); + expect(second.outcomes[0], 'the kernel spawned a second run for one scheduled instant') + .toMatchObject({ matched: true, deduped: true, run: null }); + // One run journal on disk, not two. + expect(runJournals(dataDir)).toHaveLength(1); + }, 45_000); + + it('a poller RESTART re-emitting a slot does not re-run it', async () => { + const dataDir = temporaryDirectory('flows-live-tick-restart-'); + await startDaemon(dataDir); + const client = await connectClient(dataDir); + await client.hello('live-tick-restart'); + const spec = specWithReportCli(); + + const before: TickCursor = { lastEmittedSlot: 29_399_999 }; + await emitDueTicks(spec, client, { schedule, cursor: before, nowMs: 29_400_000 * MINUTE }); + expect(runJournals(dataDir)).toHaveLength(1); + + // Cursor lost. The restarted poller treats the current slot as due. + const afterRestart: TickCursor = {}; + const replay = await emitDueTicks(spec, client, { + schedule, cursor: afterRestart, nowMs: 29_400_000 * MINUTE + 30_000, + }); + + expect(replay.emittedSlots).toEqual([29_400_000]); + expect(replay.outcomes[0]).toMatchObject({ matched: true, deduped: true }); + expect(runJournals(dataDir), 'a restart inside one slot produced a second run').toHaveLength(1); + }, 45_000); + + it('a MISSED interval is backfilled into its own run, not collapsed into the current one', async () => { + // The other half of the gate. Dedupe that keyed on the schedule rather + // than the instant would collapse all four slots into one run. + const dataDir = temporaryDirectory('flows-live-tick-backfill-'); + await startDaemon(dataDir); + const client = await connectClient(dataDir); + await client.hello('live-tick-backfill'); + const spec = specWithReportCli(); + + const cursor: TickCursor = { lastEmittedSlot: 29_400_000 }; + const result = await emitDueTicks(spec, client, { + schedule, cursor, nowMs: 29_400_004 * MINUTE, + }); + + expect(result.emittedSlots).toEqual([29_400_001, 29_400_002, 29_400_003, 29_400_004]); + const runIds = result.outcomes.map((o) => (o as { run: { run_id: string } }).run.run_id); + expect(new Set(runIds).size, 'backfilled slots collapsed into fewer runs').toBe(4); + expect(runJournals(dataDir)).toHaveLength(4); + }, 45_000); + + it('a tick for a DIFFERENT schedule id does not wake this flow', async () => { + // Without the trigger's `pattern`, every schedule in the process would + // wake every tick-triggered flow, since they share one event type. + const dataDir = temporaryDirectory('flows-live-tick-pattern-'); + await startDaemon(dataDir); + const client = await connectClient(dataDir); + await client.hello('live-tick-pattern'); + const spec = specWithReportCli(); + + const other = await emitDueTicks(spec, client, { + schedule: { ...schedule, scheduleId: 'some-other-schedule' }, + cursor: { lastEmittedSlot: 29_399_999 }, + nowMs: 29_400_000 * MINUTE, + }); + + expect(other.outcomes[0]).toMatchObject({ matched: false, deduped: false, run: null }); + expect(runJournals(dataDir)).toHaveLength(0); + }, 45_000); + + it('journals the declared silence budget, so a dead schedule is not silently zero', async () => { + // Liveness. `subscription.registered` carries the budget the sweep will + // actually apply — which is how an operator can tell a declared budget + // from the engine default. Without a declared budget this flow would + // inherit 5 minutes without its author ever choosing it. + const dataDir = temporaryDirectory('flows-live-tick-liveness-'); + await startDaemon(dataDir); + const client = await connectClient(dataDir); + await client.hello('live-tick-liveness'); + const spec = specWithReportCli(); + + const result = await emitDueTicks(spec, client, { + schedule, cursor: { lastEmittedSlot: 29_399_999 }, nowMs: 29_400_000 * MINUTE, + }); + const runId = (result.outcomes[0] as { run: { run_id: string } }).run.run_id; + + const entries = (await client.journalRead(runId)).entries; + const registered = entries.find( + (entry) => isObject(entry) && entry['entry_type'] === 'subscription.registered', + ) as { payload: Record } | undefined; + expect(registered, 'no subscription.registered entry — the sweep has nothing to key on') + .toBeDefined(); + expect(registered!.payload).toMatchObject({ + subscription_id: 'every-minute', + event_type: 'flows.tick', + // 180_000 is the flow's declared budget; 300_000 is the engine default. + // Asserting the declared value is what proves staleAfterMs survives the + // authoring -> kernel lowering rather than being dropped. + effective_stale_after_ms: 180_000, + }); + }, 45_000); +}); + + function requireExecutable(path: string, source: string, buildCommand: string): void { try { accessSync(path, constants.X_OK); @@ -1489,6 +1678,17 @@ function runArtifacts(dataDir: string): string[] { return readdirSync(dataDir).filter((name) => name !== 'relayflowd.sock').sort(); } +/** + * Per-run journal files. `runArtifacts` lists the data dir itself, which is + * constant regardless of how many runs exist — counting runs needs this. + */ +function runJournals(dataDir: string): string[] { + const runs = join(dataDir, 'runs'); + return existsSync(runs) + ? readdirSync(runs).filter((name) => name.endsWith('.sqlite3')).sort() + : []; +} + function eventOnce(client: JournalClient, event: string): Promise { return new Promise((resolveEvent) => client.once(event, (value) => resolveEvent(value as T))); } diff --git a/sdk/tests/spec-parity.test.ts b/sdk/tests/spec-parity.test.ts index c7403f1da..662f47b85 100644 --- a/sdk/tests/spec-parity.test.ts +++ b/sdk/tests/spec-parity.test.ts @@ -60,6 +60,41 @@ describe('spec parity: one dialect at the SDK<->kernel boundary', () => { expect(() => kernelToAuthoring(kernel)).toThrow('retry.initial_backoff_ms'); }); + // The trigger dialect. `toKernelSpec` used to spread `flow.triggers` through + // untouched, so an event subscription reached the kernel in camelCase and + // `relayflowd` -- `#[serde(deny_unknown_fields)]` over snake_case -- refused + // the whole spec: + // malformed run spec: unknown field `dedupeKeyTemplate`, expected one of + // `id`, `executor`, `event_type`, `pattern`, `dedupe_key_template`, + // `stale_after_ms` + // Nothing caught it because the committed fixtures are snake_case and no test + // compiled a triggered flow through the SDK and compared the two. These do. + for (const [flowFile, canonicalFile] of [ + ['event-triggered-flow.yaml', 'event-triggered-flow.spec.canonical.json'], + ['hn-monitor.flow.yaml', 'hn-monitor.spec.canonical.json'], + ['tick-heartbeat.flow.yaml', 'tick-heartbeat.spec.canonical.json'], + ] as const) { + it(`lowers ${flowFile}'s triggers to the kernel's snake_case dialect`, () => { + expect(compileYamlToCanonicalJson(fixture(flowFile))).toBe(fixture(canonicalFile).trim()); + }); + } + + it('hashes tick-heartbeat to the pinned spec_hash', () => { + expect(compileAndHash(fixture('tick-heartbeat.flow.yaml')).hash) + .toBe(fixture('tick-heartbeat.spec.sha256').trim()); + }); + + it('round-trips every trigger field across the kernel dialect', () => { + const flow = compileYaml(fixture('tick-heartbeat.flow.yaml')); + expect(kernelToAuthoring(toKernelSpec(flow))).toEqual(flow); + }); + + it('refuses a kernel trigger key the authoring dialect cannot represent', () => { + const kernel = toKernelSpec(compileYaml(fixture('tick-heartbeat.flow.yaml'))); + (kernel.triggers![0] as Record)['schedule'] = '*/5 * * * *'; + expect(() => kernelToAuthoring(kernel)).toThrow('spec.triggers[0]'); + }); + it('round-trips flow, trigger, and step CLI declarations', () => { const flow = compileYaml(` version: '0.1.0' diff --git a/sdk/tests/tick-source.test.ts b/sdk/tests/tick-source.test.ts new file mode 100644 index 000000000..c6e58cdb9 --- /dev/null +++ b/sdk/tests/tick-source.test.ts @@ -0,0 +1,309 @@ +import { readFileSync } from 'node:fs'; +import { dirname, join } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { describe, expect, it } from 'vitest'; +import { + DEFAULT_MAX_CATCH_UP, + TICK_DEDUPE_KEY_TEMPLATE, + TICK_EVENT_TYPE, + emitDueTicks, + scheduledForMs, + slotFor, + tickDedupeKey, + type TickCursor, + type TickPayload, + type TickSchedule, +} from '../src/tick-source.js'; +import { compileYaml, toKernelSpec } from '../src/compile.js'; + +const TESTDATA = join(dirname(fileURLToPath(import.meta.url)), '..', '..', 'testdata'); + +/** + * Stands in for the kernel's `(flow_key, subscription_id, dedupe_key)` claim. + * `tests/live-kernel.test.ts` proves the real kernel behaves this way; this + * fake lets the bound tests below run without a daemon. + */ +function dedupingSink(schedule: TickSchedule) { + const claimed = new Set(); + const runs: TickPayload[] = []; + const submitted: TickPayload[] = []; + return { + runs, + submitted, + claimed, + async eventSubmit(_spec: unknown, event: { type: string; payload?: unknown }) { + const payload = event.payload as TickPayload; + submitted.push(payload); + const key = `${event.type}:${payload.schedule_id}:${payload.scheduled_for_ms}`; + // Cross-check: the key the kernel would derive from the flow's declared + // template must equal the one this source claims it derives. + expect(key).toBe(tickDedupeKey(schedule, payload.slot)); + if (claimed.has(key)) return { matched: true, deduped: true }; + claimed.add(key); + runs.push(payload); + return { matched: true, deduped: false }; + }, + }; +} + +const MINUTE = 60_000; +const schedule = (over: Partial = {}): TickSchedule => ({ + scheduleId: 'heartbeat-1m', + intervalMs: MINUTE, + epochMs: 0, + ...over, +}); + +describe('tick source: the slot grid', () => { + it('derives the scheduled instant from the grid, NOT from the emit time', async () => { + // The bound: a tick emitted 37s into its slot must still carry the slot's + // boundary. If `scheduled_for_ms` tracked wall clock, this fails. + const s = schedule(); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 99 }; + + await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 100 * MINUTE + 37_000 }); + + expect(sink.submitted).toHaveLength(1); + expect(sink.submitted[0]!.slot).toBe(100); + expect(sink.submitted[0]!.scheduled_for_ms).toBe(100 * MINUTE); + expect(sink.submitted[0]!.emitted_at_ms).toBe(100 * MINUTE + 37_000); + expect(sink.submitted[0]!.lag_ms).toBe(37_000); + }); + + it('anchors slots at epochMs so the grid is a declared choice', () => { + const anchored = schedule({ epochMs: 30_000 }); + expect(slotFor(anchored, 30_000)).toBe(0); + expect(slotFor(anchored, 89_999)).toBe(0); + expect(slotFor(anchored, 90_000)).toBe(1); + expect(scheduledForMs(anchored, 5)).toBe(30_000 + 5 * MINUTE); + }); +}); + +describe('tick source: a double-fire produces ONE run', () => { + it('two emissions of one scheduled instant claim the same key, so one run', async () => { + // The gate. "A tick fired" proves nothing; this is the real bound. + const s = schedule(); + const sink = dedupingSink(s); + + // Two independent pollers, each with its own cursor, both awake in slot 100 + // at DIFFERENT wall-clock instants inside the slot. + const cursorA: TickCursor = { lastEmittedSlot: 99 }; + const cursorB: TickCursor = { lastEmittedSlot: 99 }; + await emitDueTicks({}, sink, { schedule: s, cursor: cursorA, nowMs: 100 * MINUTE + 1 }); + await emitDueTicks({}, sink, { schedule: s, cursor: cursorB, nowMs: 100 * MINUTE + 45_000 }); + + expect(sink.submitted).toHaveLength(2); + expect(sink.runs).toHaveLength(1); + expect(sink.runs[0]!.slot).toBe(100); + }); + + it('a poller RESTART that loses its cursor re-emits the slot but does not re-run it', async () => { + const s = schedule(); + const sink = dedupingSink(s); + + const before: TickCursor = { lastEmittedSlot: 99 }; + await emitDueTicks({}, sink, { schedule: s, cursor: before, nowMs: 100 * MINUTE + 5_000 }); + expect(sink.runs).toHaveLength(1); + + // Restart: cursor is gone, the process wakes up mid-slot-100 and treats the + // current slot as due. + const afterRestart: TickCursor = {}; + await emitDueTicks({}, sink, { schedule: s, cursor: afterRestart, nowMs: 100 * MINUTE + 50_000 }); + + expect(sink.submitted).toHaveLength(2); + expect(sink.submitted[1]!.slot).toBe(100); + expect(sink.runs, 'a restart inside one slot produced a second run').toHaveLength(1); + }); + + it('successive slots are NOT deduped away — the schedule keeps running', async () => { + // The mirror of the dedupe test: if the key were constant per schedule, + // every test above would still pass and the flow would run exactly once, + // ever. This is what proves the key is per-instant, not per-schedule. + const s = schedule(); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 99 }; + + await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 100 * MINUTE }); + await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 101 * MINUTE }); + await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 102 * MINUTE }); + + expect(sink.runs.map((r) => r.slot)).toEqual([100, 101, 102]); + }); +}); + +describe('tick source: a missed interval does not silently vanish', () => { + it('backfills every slot a sleeping poller passed over', async () => { + // The bound: down for four slots, back up. Emitting only the current slot + // would lose three scheduled runs and report success. + const s = schedule(); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + const result = await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 104 * MINUTE + 10 }); + + expect(result.emittedSlots).toEqual([101, 102, 103, 104]); + expect(sink.runs.map((r) => r.slot)).toEqual([101, 102, 103, 104]); + expect(result.skippedSlots).toEqual([]); + expect(cursor.lastEmittedSlot).toBe(104); + }); + + it('a backfilled tick carries its own lag, so a stale run can tell it is stale', async () => { + const s = schedule(); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 103 * MINUTE }); + + expect(sink.runs.map((r) => r.lag_ms)).toEqual([2 * MINUTE, 1 * MINUTE, 0]); + }); + + it('does NOT advance past a slot whose submit failed — the next poll retries it', async () => { + // Crash-window contract, matching `pollDirectoryOnce`'s `seen` discipline. + const s = schedule(); + let calls = 0; + const failing = { + async eventSubmit(_spec: unknown, event: { type: string; payload?: unknown }) { + calls++; + if ((event.payload as TickPayload).slot === 102) { + throw new Error('journal client: connection closed'); + } + return { matched: true, deduped: false }; + }, + }; + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + await expect( + emitDueTicks({}, failing, { schedule: s, cursor, nowMs: 104 * MINUTE }), + ).rejects.toThrow(/journal client/); + + expect(calls).toBe(2); + // 101 landed, 102 did not. The cursor must sit at 101 so 102, 103 and 104 + // are still due. + expect(cursor.lastEmittedSlot).toBe(101); + + const recovering = dedupingSink(s); + const result = await emitDueTicks({}, recovering, { schedule: s, cursor, nowMs: 104 * MINUTE }); + expect(result.emittedSlots).toEqual([102, 103, 104]); + }); + + it('reports slots dropped by the catch-up bound instead of dropping them quietly', async () => { + // The bound: a week of downtime must not replay ten thousand runs, and it + // must not pretend nothing was missed either. + const s = schedule({ maxCatchUp: 3 }); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + const result = await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 110 * MINUTE }); + + expect(result.emittedSlots).toEqual([108, 109, 110]); + expect(result.skippedSlots).toEqual([101, 102, 103, 104, 105, 106, 107]); + expect(cursor.lastEmittedSlot).toBe(110); + }); + + it('reports each skipped slot exactly once, not on every later poll', async () => { + const s = schedule({ maxCatchUp: 2 }); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + const first = await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 106 * MINUTE }); + expect(first.skippedSlots).toEqual([101, 102, 103, 104]); + + const second = await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 107 * MINUTE }); + expect(second.skippedSlots).toEqual([]); + expect(second.emittedSlots).toEqual([107]); + }); + + it('caps at DEFAULT_MAX_CATCH_UP when the schedule declares no bound', async () => { + const s = schedule(); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 0 }; + + const result = await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 500 * MINUTE }); + + expect(result.emittedSlots).toHaveLength(DEFAULT_MAX_CATCH_UP); + expect(result.skippedSlots).toHaveLength(500 - DEFAULT_MAX_CATCH_UP); + }); +}); + +describe('tick source: edges', () => { + it('the first poll on a fresh cursor emits ONE tick, not every slot since the epoch', async () => { + const s = schedule(); + const sink = dedupingSink(s); + const cursor: TickCursor = {}; + + const result = await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 100 * MINUTE }); + + expect(result.emittedSlots).toEqual([100]); + expect(result.skippedSlots).toEqual([]); + }); + + it('emits nothing when the current slot has already been emitted', async () => { + const s = schedule(); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + const result = await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 100 * MINUTE + 59_999 }); + + expect(result.emittedSlots).toEqual([]); + expect(sink.submitted).toHaveLength(0); + }); + + it('emits nothing when the clock moves backwards, and resumes at the next real slot', async () => { + const s = schedule(); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + const backwards = await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 90 * MINUTE }); + expect(backwards.emittedSlots).toEqual([]); + expect(cursor.lastEmittedSlot, 'a backwards clock rewound the cursor').toBe(100); + + const forward = await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 101 * MINUTE }); + expect(forward.emittedSlots).toEqual([101]); + }); + + it('refuses a non-positive interval rather than dividing by zero', async () => { + await expect( + emitDueTicks({}, dedupingSink(schedule()), { + schedule: schedule({ intervalMs: 0 }), + cursor: {}, + nowMs: 1, + }), + ).rejects.toThrow(/intervalMs must be a positive integer/); + }); + + it('refuses an empty schedule id, which would collide with another schedule', async () => { + await expect( + emitDueTicks({}, dedupingSink(schedule()), { + schedule: schedule({ scheduleId: '' }), + cursor: {}, + nowMs: 1, + }), + ).rejects.toThrow(/scheduleId must be a non-empty string/); + }); +}); + +describe('tick source and the worked example agree', () => { + const flow = compileYaml(readFileSync(join(TESTDATA, 'tick-heartbeat.flow.yaml'), 'utf8')); + const trigger = flow.triggers![0]!; + + it('the spec subscribes to the event type the source emits', () => { + expect(trigger.eventType).toBe(TICK_EVENT_TYPE); + }); + + it('the spec dedupes on the scheduled instant, not on the emit time', () => { + expect(trigger.dedupeKeyTemplate).toBe(TICK_DEDUPE_KEY_TEMPLATE); + expect(trigger.dedupeKeyTemplate).not.toContain('emitted_at_ms'); + }); + + it('the spec declares a silence budget, so the schedule cannot die quietly', () => { + // Without this the engine default applies and the flow's own author has + // made no statement about how long silence is acceptable. + expect(trigger.staleAfterMs).toBeGreaterThan(0); + expect(toKernelSpec(flow).triggers![0]!.stale_after_ms).toBe(trigger.staleAfterMs); + }); + + it('the spec narrows to one schedule id, so sibling schedules cannot wake it', () => { + expect(trigger.pattern).toEqual({ schedule_id: 'heartbeat-1m' }); + }); +}); diff --git a/sdk/tests/validate.test.ts b/sdk/tests/validate.test.ts index 35e1d91b6..c0c64640e 100644 --- a/sdk/tests/validate.test.ts +++ b/sdk/tests/validate.test.ts @@ -324,6 +324,35 @@ describe('validate: preflight declarations', () => { expect(result.errors).toContain('spec.triggers[1].id: duplicate trigger id "hourly"'); }); + it('accepts a declared silence budget', () => { + const result = validateSpec({ + version: '0.1.0', + triggers: [{ id: 'tick', executor: 'agent-worker', staleAfterMs: 180_000 }], + steps: [{ id: 'ready', type: 'deterministic', command: 'true' }], + }); + expect(result).toEqual({ ok: true, errors: [] }); + }); + + it.each([ + [0, 'expected a positive number'], + [-1, 'expected a positive number'], + [1.5, 'expected an integer'], + ['180000', 'expected an integer'], + [Number.MAX_VALUE, 'exceeds the i64 range'], + ])('rejects staleAfterMs %p', (value, fragment) => { + // A budget the kernel cannot represent fails OPEN — the sweep can never + // mark the subscription stale — so it must not compile. Zero fails the + // other way: it alerts on every sweep and trains an operator to ignore it. + const result = validateSpec({ + version: '0.1.0', + triggers: [{ id: 'tick', executor: 'agent-worker', staleAfterMs: value }], + steps: [{ id: 'ready', type: 'deterministic', command: 'true' }], + }); + expect(result.ok).toBe(false); + expect(result.errors.join(' ')).toContain('spec.triggers[0].staleAfterMs'); + expect(result.errors.join(' ')).toContain(fragment); + }); + it('rejects malformed and unknown trigger/CLI fields fail-closed', () => { const result = validateSpec({ version: '0.1.0', diff --git a/testdata/preflight/tick-slot-report-cli b/testdata/preflight/tick-slot-report-cli new file mode 100755 index 000000000..f9ae33f3a --- /dev/null +++ b/testdata/preflight/tick-slot-report-cli @@ -0,0 +1,29 @@ +#!/usr/bin/env node +// Agent CLI for the `tick-heartbeat` worked example. Reads the tick that woke +// this run out of $RELAYFLOW_WAKE_CONTEXT and emits the JSON the flow's +// json_schema gate requires. +// +// Deterministic on purpose. The point of the worked example is that the +// SCHEDULE fired the run and that one scheduled instant produces one run; a +// model in the loop would add nondeterminism to a test whose subject is the +// trigger plane, not the model. +import { receiveWrapperRequest } from './wrapper-session.mjs'; + +await receiveWrapperRequest(); + +const raw = process.env.RELAYFLOW_WAKE_CONTEXT; +if (raw === undefined) { + process.stderr.write('tick-slot-report-cli: no RELAYFLOW_WAKE_CONTEXT — this step was not woken by a tick\n'); + process.exit(1); +} +const payload = JSON.parse(raw)?.triggering_event?.payload; +if (payload === undefined || payload === null) { + process.stderr.write(`tick-slot-report-cli: wake context carried no triggering event payload: ${raw}\n`); + process.exit(1); +} +process.stdout.write(JSON.stringify({ + schedule_id: payload.schedule_id, + slot: payload.slot, + scheduled_for_ms: payload.scheduled_for_ms, + lag_ms: payload.lag_ms, +})); diff --git a/testdata/tick-heartbeat.flow.yaml b/testdata/tick-heartbeat.flow.yaml new file mode 100644 index 000000000..588cce124 --- /dev/null +++ b/testdata/tick-heartbeat.flow.yaml @@ -0,0 +1,56 @@ +version: '0.1.0' +name: tick-heartbeat +description: >- + Run one step on a schedule. Third proactive workload on gate 2 primitives + (hn-monitor is the first, dir-watcher the second) — and the one that proves + a relayflow can be SCHEDULED. There is no `cron:` field: the schedule is an + event source (sdk/src/tick-source.ts) that submits `flows.tick` through the + same `event.submit` path every other source uses, so it inherits the kernel's + dedupe claim and its liveness sweep instead of inventing either. +triggers: + - id: every-minute + executor: agent-worker + eventType: flows.tick + # Narrow to ONE schedule. Without this, every schedule in the process would + # wake this flow, since they all share the `flows.tick` event type. + pattern: + schedule_id: heartbeat-1m + # The scheduled instant, never the wall clock at emit. Two deliveries of one + # slot derive the same key and the kernel's + # (flow_key, subscription_id, dedupe_key) claim spawns one run — so a + # double-fire, a re-delivery, and a poller restart that re-emits a slot are + # all idempotent by construction. Keep in lockstep with + # TICK_DEDUPE_KEY_TEMPLATE in sdk/src/tick-source.ts. + dedupeKeyTemplate: '{{event.type}}:{{payload.schedule_id}}:{{payload.scheduled_for_ms}}' + # Silence budget: three slots. A schedule whose source dies stops being + # silently zero — the kernel's sweep journals `subscription.stale` and + # logs it (kernel/relayflowd/src/server/liveness.rs). Declared rather than + # inherited, because the 5-minute engine default is a decision this flow's + # author did not make. + staleAfterMs: 180000 +steps: + - id: report-slot + type: agent + instruction: >- + A scheduled tick fired. Read the wake context, which carries the + triggering event's payload, and output a JSON object naming: the + schedule id, the slot index, the scheduled instant in epoch + milliseconds, and how far behind schedule this run started + (the payload's lag_ms). + recoveryMode: reset + output: + type: object + required: + - schedule_id + - slot + - scheduled_for_ms + - lag_ms + properties: + schedule_id: + type: string + slot: + type: integer + scheduled_for_ms: + type: integer + lag_ms: + type: integer diff --git a/testdata/tick-heartbeat.spec.canonical.json b/testdata/tick-heartbeat.spec.canonical.json new file mode 100644 index 000000000..aceb55ac7 --- /dev/null +++ b/testdata/tick-heartbeat.spec.canonical.json @@ -0,0 +1 @@ +{"description":"Run one step on a schedule. Third proactive workload on gate 2 primitives (hn-monitor is the first, dir-watcher the second) — and the one that proves a relayflow can be SCHEDULED. There is no `cron:` field: the schedule is an event source (sdk/src/tick-source.ts) that submits `flows.tick` through the same `event.submit` path every other source uses, so it inherits the kernel's dedupe claim and its liveness sweep instead of inventing either.","name":"tick-heartbeat","steps":[{"depends_on":[],"id":"report-slot","instruction":"A scheduled tick fired. Read the wake context, which carries the triggering event's payload, and output a JSON object naming: the schedule id, the slot index, the scheduled instant in epoch milliseconds, and how far behind schedule this run started (the payload's lag_ms).","max_iterations":1,"recovery_mode":"reset","retry":{"initial_backoff_ms":100,"jitter_percent":20,"max_backoff_ms":60000,"multiplier":2},"type":"agent","verification":{"json_schema":{"properties":{"lag_ms":{"type":"integer"},"schedule_id":{"type":"string"},"scheduled_for_ms":{"type":"integer"},"slot":{"type":"integer"}},"required":["schedule_id","slot","scheduled_for_ms","lag_ms"],"type":"object"}}}],"triggers":[{"dedupe_key_template":"{{event.type}}:{{payload.schedule_id}}:{{payload.scheduled_for_ms}}","event_type":"flows.tick","executor":"agent-worker","id":"every-minute","pattern":{"schedule_id":"heartbeat-1m"},"stale_after_ms":180000}],"version":"0.1.0"} diff --git a/testdata/tick-heartbeat.spec.sha256 b/testdata/tick-heartbeat.spec.sha256 new file mode 100644 index 000000000..6350ee30b --- /dev/null +++ b/testdata/tick-heartbeat.spec.sha256 @@ -0,0 +1 @@ +0c3d089f0075c53442c5c2241ada734af5e450a2b558c16030aaf1d1e29127ff From ad4b4312f34e433d40a8cee44a84fa28877bf91e Mon Sep 17 00:00:00 2001 From: kjgbot Date: Thu, 3 Sep 2026 14:52:50 +0200 Subject: [PATCH 2/2] fix(sdk): close three silently-zero paths in the tick source MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signoff returned REVIEW_PASSED at 969ed44 with no P0 and no P1. These three are fixed anyway because all three reproduce the exact failure class this PR exists to address — RFC-0001's "a flow that is never triggered is silently zero". A scheduled-trigger primitive with paths that produce a silently-zero schedule undercuts the thing it is for. P2-1 — a skip could vanish when the same poll failed. The cursor was advanced past `skippedSlots` BEFORE the emit loop, so a submit failure discarded the returned result (the only place a skipped slot is ever recorded, since it never reaches the kernel) while the cursor had already moved past the evidence. Seven skipped slots, no report anywhere, not re-derivable next poll. The module's own guarantee is that a skip cannot be silent; in that path it was. The pre-advance is gone, and a failure now throws TickEmitError carrying emittedSlots, outcomes and skippedSlots. Three paths account for a skip and none drops it: a submit succeeds and the result carries it; a submit throws and the error carries it; nothing was submitted, so the cursor never moved and the next poll re-derives the identical due range. The reviewer noted this becomes P1 the moment the CLI runner lands and must be fixed as its prerequisite — so it is now met. P2-2 — a NaN grid made a schedule permanently and silently zero. `epochMs` and `nowMs` were unvalidated while `intervalMs` was, and the asymmetry was the bug: a NaN made slotFor() return NaN, every comparison against it false, and the poll returned an empty result — no submit, no skip, no throw. An infinity died with a raw `Invalid array length`. Both now go through requireNonNegativeInteger, which covers NaN, both infinities, negatives and non-integers in one check. A grid that cannot be computed refuses at the call, by name. P3 — `9_223_372_036_854_775_807` does not express i64::MAX in a double; it rounds UP to 2^63, so `value > MAX_STALE_AFTER_MS` admitted exactly the one value the kernel refuses. The SDK/kernel-agreement failure this file exists to prevent, in miniature. The bound is now Number.MAX_SAFE_INTEGER: above 2^53 a JS number cannot name a specific integer, so a larger budget could not cross the boundary faithfully even if the kernel would take it. 2^53 ms is ~285,000 years. Red captured first for all three. Mutation-verified, both files restored byte-for-byte (tick-source.ts a0f2d1e4..., validate.ts 9986dc54...): M3 re-add the cursor pre-advance -> 1 failed (expected 107 to be 100) M3b throw cause, not TickEmitError -> 1 failed (not an instance of TickEmitError) M4 drop the epochMs/nowMs guards -> 8 failed M5 restore the i64 literal bound -> 2 failed (expected true to be false) M3 and M3b failing one test each, and different tests, shows the two halves of the P2-1 fix are independently load-bearing rather than one mechanism counted twice. Gates (RELAYFLOWD_BIN pinned to this worktree's build): tsc --noEmit exit 0 tsc -p tsconfig.tests.json exit 0 vitest run 423 passed, 3 skipped, 0 failed (was 408/3/0) cargo test --workspace 100 passed, identical test set to baseline +15 test names, 0 removed. Fixtures untouched: `git diff --stat testdata/` is empty and spec_hash is unchanged at 0c3d089f..., so the compiled dialect is identical to the signed-off head. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01FtQSAcGDta5VH9xiZFT4sR Session-Id: c228933d-4f94-4d83-9a9a-daf3c83b94f1 --- .../20260903-scheduled-trigger-design.md | 130 ++++++++++++++++-- sdk/src/index.ts | 1 + sdk/src/tick-source.ts | 89 ++++++++++-- sdk/src/validate.ts | 14 +- sdk/tests/tick-source.test.ts | 100 ++++++++++++++ sdk/tests/validate.test.ts | 7 +- 6 files changed, 313 insertions(+), 28 deletions(-) diff --git a/ops/reviews/20260903-scheduled-trigger-design.md b/ops/reviews/20260903-scheduled-trigger-design.md index 093b80abe..3315915ad 100644 --- a/ops/reviews/20260903-scheduled-trigger-design.md +++ b/ops/reviews/20260903-scheduled-trigger-design.md @@ -313,7 +313,7 @@ across every worktree's target tree. |---|---|---| | `tsc --noEmit` | exit 0 | exit 0 | | `tsc -p tsconfig.tests.json` | exit 0 | exit 0 | -| `vitest run` | 370 passed, 3 skipped, **0 failed** | 408 passed, 3 skipped, **0 failed** | +| `vitest run` | 370 passed, 3 skipped, **0 failed** | 423 passed, 3 skipped, **0 failed** | | `cargo test --workspace` | 100 passed, 0 failed | 100 passed, 0 failed | > The first baseline attempt reported 12 failures. Those were an artifact of not @@ -326,12 +326,12 @@ across every worktree's target tree. ``` +6 sdk/tests/live-kernel.test.ts (21 -> 27) +6 sdk/tests/spec-parity.test.ts (15 -> 21) -+6 sdk/tests/validate.test.ts (36 -> 42) -+20 sdk/tests/tick-source.test.ts (0 -> 20) - total 373 -> 411 (370 passed + 3 skipped -> 408 passed + 3 skipped) ++8 sdk/tests/validate.test.ts (36 -> 44) ++33 sdk/tests/tick-source.test.ts (0 -> 33) + total 373 -> 426 (370 passed + 3 skipped -> 423 passed + 3 skipped) ``` -Full-name set-diff: **38 test names added, 0 removed, 0 renamed.** +Full-name set-diff: **53 test names added, 0 removed, 0 renamed.** `git diff --stat sdk/tests/live-kernel.test.ts` is `200 ++++` with **zero deletions** — the new block is a pure insertion, no existing case touched. @@ -435,6 +435,107 @@ const replay = await emitDueTicks(spec, client, { schedule, cursor: {}, nowMs: n --- +## 8a. Post-signoff fixes (P2-1, P2-2, P3) + +Signoff returned REVIEW_PASSED at `969ed44` with no P0 and no P1. These three +were fixed anyway, at the lead's gate, because all three reproduce **the exact +failure class this PR exists to address** — RFC-0001's "a flow that is never +triggered is silently zero". A scheduled-trigger primitive with paths that +produce a silently-zero schedule undercuts the thing it is for. + +Red captured first for all three, then fixed, then mutation-verified. + +### P2-1 — a skip could vanish when the same poll failed + +The cursor was advanced past `skippedSlots` **before** the emit loop. If a +submit then threw, the returned result — the only place `skippedSlots` ever +lived, since a skipped slot never reaches the kernel — was discarded, while the +cursor had already moved past the evidence. Seven skipped slots, no report +anywhere, and not re-derivable on the next poll. The module's own stated +guarantee is that the skip cannot be silent; in that path it was. + +Red, before the fix: + +``` +× leaves the skipped slots re-derivable when the FIRST submit fails + → cursor advanced past slots that never reached the kernel: expected 107 to be 100 +× carries the skipped slots on the error when a later submit fails + → the failing submit did not throw: The instanceof assertion needs a constructor but undefined was given. +``` + +Fix: the pre-advance is gone, and a submit failure now throws `TickEmitError` +carrying `emittedSlots`, `outcomes` and `skippedSlots`. Three independent paths +now account for a skip and no path drops it: + +- a submit succeeds → the cursor advances past the skipped slots as a side + effect and the **result** carries them (happy path, still exactly once); +- a submit throws → **`TickEmitError`** carries them to the caller; +- nothing was submitted → the cursor never moved and the next poll + **re-derives** the identical due range. + +The reviewer noted this "becomes P1 the moment the CLI runner lands, and must be +fixed as its prerequisite" — so it is now a prerequisite that is already met. + +### P2-2 — a NaN grid made a schedule permanently, silently zero + +`epochMs` and `nowMs` were unvalidated while `intervalMs` was validated, and the +asymmetry was the bug. A `NaN` made `slotFor` return `NaN`, every comparison +against it false, and the poll returned an empty result: **no submit, no skip, +no throw.** An infinity was no better as a diagnosis — the backfill loop died +with a raw `Invalid array length`. + +Red, before the fix (8 cases): + +``` +× refuses epochMs = NaN → promise resolved "{ emittedSlots: [], outcomes: [], …(1) }" instead of rejecting +× refuses epochMs = Infinity → ... but got 'Invalid array length' +× refuses nowMs = -1 → promise resolved "{ emittedSlots: [ -1 ], …(2) }" instead of rejecting +× refuses nowMs = 1.5 → promise resolved "{ emittedSlots: [ +0 ], …(2) }" instead of rejecting +``` + +Fix: `requireNonNegativeInteger` on both, mirroring `requirePositiveInteger`. +`Number.isInteger` is false for `NaN` and both infinities, so one check covers +all three. A grid that cannot be computed refuses at the call, by name. + +### P3 — the i64 literal admitted exactly the value the kernel refuses + +`const MAX_STALE_AFTER_MS = 9_223_372_036_854_775_807` does not express `i64::MAX` +in a double — it rounds **up** to 2^63, i.e. `i64::MAX + 1`. So `value > MAX` +admitted precisely the one value `relayflowd` rejects. The SDK/kernel-agreement +failure this file exists to prevent, in miniature. + +Red: `rejects staleAfterMs 9223372036854776000 → expected true to be false`. + +Fix: the bound is `Number.MAX_SAFE_INTEGER`. Above 2^53 a JS number cannot name +a specific integer at all, so a larger budget could not cross the boundary +faithfully even if the kernel would take it. 2^53 ms is ~285,000 years. + +### Mutation verification of the fixes + +Pre-mutation `sha256`: `tick-source.ts` `a0f2d1e4874bf164a0cd657028c89636ebf2d4adf4977f55f7441e795b0dd3ad`, +`validate.ts` `9986dc54c79476d1355cbfadcf83c6356ffe6d476ca4ef8286c23c6ba6cf4876`. +Both restored to those exact hashes afterwards. + +| Mutation | Reverted | Result | +|---|---|---| +| **M3** | re-add the cursor pre-advance | 1 failed — `expected 107 to be 100` | +| **M3b** | `throw cause` instead of `TickEmitError` | 1 failed — `expected Error: journal client: closed to be an instance of TickEmitError` | +| **M4** | drop both `requireNonNegativeInteger` calls | **8 failed** | +| **M5** | restore the i64 literal bound | 2 failed — `expected true to be false` | + +M3 and M3b failing **one test each, different tests**, is the useful result: the +two halves of the P2-1 fix are independently load-bearing, not one mechanism +double-counted. + +### Fixtures untouched + +`git diff --stat testdata/` is empty and `spec_hash` is unchanged at +`0c3d089f0075c53442c5c2241ada734af5e450a2b558c16030aaf1d1e29127ff` — these +fixes are behavioural and validation-side only, so the compiled dialect and its +hash are identical to the signed-off head. + +--- + ## 9. What I could not verify - **Provisioned-but-never-fired liveness.** Not attempted; it needs a kernel @@ -443,14 +544,17 @@ const replay = await emitDueTicks(spec, client, { schedule, cursor: {}, nowMs: n - **A schedule running in production.** No CLI runner ships here, so `agent-relay cloud schedules` will still report "No workflow schedules found" after this merges. The primitive is proven; it is not deployed. -- **Behaviour across a daemon restart mid-backfill.** The unit test covers a - failed submit mid-backfill with a fake sink; I did not kill a live daemon - between two slots of one backfill. The dedupe claim is durable in SQLite so I - expect it to hold, but I did not observe it and am not claiming it. -- **Clock skew between two hosts running the same schedule.** Two pollers on the - same grid dedupe correctly (proved), but I did not test hosts whose clocks - disagree by more than one interval — that would put them in different slots - and produce two runs. Real, unaddressed, out of scope here. +- **Behaviour across a daemon restart mid-backfill.** *I* did not verify this — + the unit test covers a failed submit mid-backfill with a fake sink, and I did + not kill a live daemon between two slots of one backfill. The independent + signoff did, on its own live harness: dedupe survives, 3 slots to 3 journals. + Recorded here as the reviewer's evidence, not mine. +- **Clock skew between two hosts running the same schedule.** I flagged this as + a suspected duplicate-run risk. The signoff checked it and the fear was + **wrong**: skew shifts *when* a slot fires, not how many times it fires, + because the slot's identity is a property of the grid and not of either + host's clock. Left here corrected rather than deleted, since the original + claim is in the PR description. - **`maxCatchUp = 60` as the right default.** Chosen as an hour of one-minute slots. Not empirically justified. - **`hn-monitor.flow.yaml` round-trip** through `kernelToAuthoring(toKernelSpec(…))` diff --git a/sdk/src/index.ts b/sdk/src/index.ts index 7f147660a..a684f80bc 100644 --- a/sdk/src/index.ts +++ b/sdk/src/index.ts @@ -170,6 +170,7 @@ export { export { emitDueTicks, scheduledForMs, + TickEmitError, slotFor, tickDedupeKey, DEFAULT_MAX_CATCH_UP, diff --git a/sdk/src/tick-source.ts b/sdk/src/tick-source.ts index dc7a94e63..2f9bb5e07 100644 --- a/sdk/src/tick-source.ts +++ b/sdk/src/tick-source.ts @@ -171,12 +171,60 @@ export function tickDedupeKey(schedule: TickSchedule, slot: number): string { return `${TICK_EVENT_TYPE}:${schedule.scheduleId}:${scheduledForMs(schedule, slot)}`; } +/** + * Error thrown when a submit inside `emitDueTicks` fails, carrying the poll's + * accounting so far. + * + * A plain rethrow discarded `skippedSlots` — the ONLY record that real + * scheduled work had been passed over, since a skipped slot never reaches the + * kernel and nothing downstream would ever see it. That made the skip silent + * in exactly the failure path where an operator most needs it, contradicting + * this module's own guarantee and reproducing the silent-death class the whole + * primitive exists to prevent. + */ +export class TickEmitError extends Error { + readonly emittedSlots: number[]; + readonly outcomes: unknown[]; + readonly skippedSlots: number[]; + constructor(result: TickEmitResult, cause: unknown) { + const detail = cause instanceof Error ? cause.message : String(cause); + super( + `tick emit failed after ${result.emittedSlots.length} slot(s)` + + `${result.skippedSlots.length > 0 ? `, with ${result.skippedSlots.length} slot(s) skipped by the catch-up bound` : ''}` + + `: ${detail}`, + { cause }, + ); + this.name = 'TickEmitError'; + this.emittedSlots = result.emittedSlots; + this.outcomes = result.outcomes; + this.skippedSlots = result.skippedSlots; + } +} + function requirePositiveInteger(value: number, field: string): void { + // `Number.isInteger` is false for NaN and for both infinities, so this one + // check covers all three. if (!Number.isInteger(value) || value <= 0) { throw new Error(`tick schedule: ${field} must be a positive integer, got ${String(value)}`); } } +/** + * `epochMs` and `nowMs` reach arithmetic that decides whether ANY slot is due. + * They were unvalidated while `intervalMs` was not, and the asymmetry was the + * bug: a NaN made `slotFor` return NaN, every comparison against it false, and + * the poll returned an empty result — no submit, no skip, no throw. A schedule + * permanently and silently zero. An infinity was worse than quiet but no + * better as a diagnosis: the backfill loop died with `Invalid array length`. + * + * A grid that cannot be computed must refuse at the call, loudly and by name. + */ +function requireNonNegativeInteger(value: number, field: string): void { + if (!Number.isInteger(value) || value < 0) { + throw new Error(`tick schedule: ${field} must be a non-negative integer, got ${String(value)}`); + } +} + /** * Submit a `flows.tick` event for every slot that has come due since the cursor * last advanced, and move the cursor. @@ -197,6 +245,8 @@ export async function emitDueTicks( requirePositiveInteger(schedule.intervalMs, 'intervalMs'); const maxCatchUp = schedule.maxCatchUp ?? DEFAULT_MAX_CATCH_UP; requirePositiveInteger(maxCatchUp, 'maxCatchUp'); + requireNonNegativeInteger(schedule.epochMs ?? 0, 'epochMs'); + requireNonNegativeInteger(nowMs, 'nowMs'); if (schedule.scheduleId === '') { throw new Error('tick schedule: scheduleId must be a non-empty string'); } @@ -221,15 +271,22 @@ export async function emitDueTicks( const skippedSlots = due.length > maxCatchUp ? due.slice(0, due.length - maxCatchUp) : []; const toEmit = due.length > maxCatchUp ? due.slice(due.length - maxCatchUp) : due; - // A skipped slot never reaches the kernel, so nothing downstream would ever - // record it. Advancing the cursor past it here is what makes the skip appear - // exactly once in `skippedSlots` instead of being re-reported forever. - if (skippedSlots.length > 0) { - cursor.lastEmittedSlot = skippedSlots[skippedSlots.length - 1]; - } - - const emittedSlots: number[] = []; - const outcomes: unknown[] = []; + // The cursor is NOT advanced past `skippedSlots` here. It used to be, and + // that lost the skip outright: if a submit then threw, the returned result + // — the only place `skippedSlots` lived — was discarded, while the cursor + // had already moved past the evidence, so the next poll could not re-derive + // it either. Seven skipped slots could vanish with no report anywhere. + // + // Instead the skip is accounted for by whichever of these happens: + // - a submit succeeds, advancing the cursor past the skipped slots as a + // side effect, and the result carries `skippedSlots` (the happy path, + // still reported exactly once); + // - a submit throws, and `TickEmitError` carries `skippedSlots` to the + // caller; + // - nothing was submitted at all, so the cursor never moved and the next + // poll re-derives the identical due range. + // No path drops it. + const result: TickEmitResult = { emittedSlots: [], outcomes: [], skippedSlots }; for (const slot of toEmit) { const scheduledFor = scheduledForMs(schedule, slot); const payload: TickPayload = { @@ -240,14 +297,22 @@ export async function emitDueTicks( emitted_at_ms: nowMs, lag_ms: nowMs - scheduledFor, }; - outcomes.push(await sink.eventSubmit(spec, { type: TICK_EVENT_TYPE, payload })); + let outcome: unknown; + try { + outcome = await sink.eventSubmit(spec, { type: TICK_EVENT_TYPE, payload }); + } catch (cause) { + // Fail closed, but never quietly: the partial accounting travels with + // the failure instead of dying with the discarded return value. + throw new TickEmitError(result, cause); + } + result.outcomes.push(outcome); // Advance ONLY after the submit succeeded. A journal failure must leave // this slot due so the next poll retries it — the same rule // `pollDirectoryOnce` applies to its `seen` set, and the reason a crash // mid-backfill cannot swallow a slot. cursor.lastEmittedSlot = slot; - emittedSlots.push(slot); + result.emittedSlots.push(slot); } - return { emittedSlots, outcomes, skippedSlots }; + return result; } diff --git a/sdk/src/validate.ts b/sdk/src/validate.ts index 19e45d248..64386c9a9 100644 --- a/sdk/src/validate.ts +++ b/sdk/src/validate.ts @@ -80,8 +80,17 @@ const TRIGGER_KEYS = [ * Mirror the bound here so an unrepresentable budget is a compile error rather * than an engine error at submit time — a budget the sweep cannot represent * fails OPEN, which is the silent death this field exists to prevent. + * + * The bound is `Number.MAX_SAFE_INTEGER`, NOT `i64::MAX`. Writing the i64 + * bound as a JS literal does not express it: `9_223_372_036_854_775_807` + * rounds UP to 2^63 in a double, so `value > MAX` then ADMITTED exactly the + * one value the kernel refuses — the SDK/kernel-agreement failure this whole + * file exists to prevent, in miniature. Above 2^53 a JS number cannot name a + * specific integer at all, so any larger budget could not be transmitted + * faithfully even if the kernel would take it. 2^53 ms is ~285,000 years; + * nothing real is lost by refusing beyond it. */ -const MAX_STALE_AFTER_MS = 9_223_372_036_854_775_807; +const MAX_STALE_AFTER_MS = Number.MAX_SAFE_INTEGER; class Validator { private errors: string[] = []; @@ -250,7 +259,8 @@ class Validator { } if (value > MAX_STALE_AFTER_MS) { this.fail( - `${at}: ${value} exceeds the i64 range the kernel's liveness sweep can represent`, + `${at}: ${value} exceeds ${MAX_STALE_AFTER_MS}, the largest budget that survives the ` + + `SDK -> kernel boundary exactly (the sweep stores it as i64)`, ); } } diff --git a/sdk/tests/tick-source.test.ts b/sdk/tests/tick-source.test.ts index c6e58cdb9..549c5888f 100644 --- a/sdk/tests/tick-source.test.ts +++ b/sdk/tests/tick-source.test.ts @@ -10,6 +10,7 @@ import { scheduledForMs, slotFor, tickDedupeKey, + TickEmitError, type TickCursor, type TickPayload, type TickSchedule, @@ -307,3 +308,102 @@ describe('tick source and the worked example agree', () => { expect(trigger.pattern).toEqual({ schedule_id: 'heartbeat-1m' }); }); }); + +describe('P2-1: a skip is never lost, even when the same poll fails', () => { + // The module promises "the skip cannot be silent". It broke that promise in + // exactly one path: the cursor was advanced past the skipped slots BEFORE + // the emit loop, so a submit failure threw away the only report of them and + // left the cursor past the evidence. Silently zero — the failure class this + // whole PR exists to address. + + it('carries the skipped slots on the error when a later submit fails', async () => { + const s = schedule({ maxCatchUp: 3 }); + const failing = { + async eventSubmit(_spec: unknown, event: { payload?: unknown }) { + if ((event.payload as TickPayload).slot === 109) throw new Error('journal client: closed'); + return { matched: true, deduped: false }; + }, + }; + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + // due = 101..110; skipped = 101..107; toEmit = 108,109,110. 108 lands, 109 dies. + const error = await emitDueTicks({}, failing, { schedule: s, cursor, nowMs: 110 * MINUTE }) + .then(() => undefined, (e: unknown) => e); + + expect(error, 'the failing submit did not throw').toBeInstanceOf(TickEmitError); + const thrown = error as TickEmitError; + expect(thrown.skippedSlots, 'seven skipped slots vanished with no report anywhere') + .toEqual([101, 102, 103, 104, 105, 106, 107]); + expect(thrown.emittedSlots).toEqual([108]); + expect(thrown.cause).toBeInstanceOf(Error); + }); + + it('leaves the skipped slots re-derivable when the FIRST submit fails', async () => { + // Nothing reached the kernel, so the cursor must not have moved at all — + // the next poll has to be able to re-derive the whole due range. + const s = schedule({ maxCatchUp: 3 }); + const failing = { async eventSubmit() { throw new Error('journal client: closed'); } }; + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + await expect(emitDueTicks({}, failing, { schedule: s, cursor, nowMs: 110 * MINUTE })) + .rejects.toThrow(/journal client/); + expect(cursor.lastEmittedSlot, 'cursor advanced past slots that never reached the kernel') + .toBe(100); + + // Recovery poll re-derives the same accounting. + const recovering = dedupingSink(s); + const retry = await emitDueTicks({}, recovering, { schedule: s, cursor, nowMs: 110 * MINUTE }); + expect(retry.skippedSlots).toEqual([101, 102, 103, 104, 105, 106, 107]); + expect(retry.emittedSlots).toEqual([108, 109, 110]); + }); + + it('still reports a skip exactly once on the success path', async () => { + // Guard against over-correcting: the fix must not re-report skips forever. + const s = schedule({ maxCatchUp: 2 }); + const sink = dedupingSink(s); + const cursor: TickCursor = { lastEmittedSlot: 100 }; + + expect((await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 106 * MINUTE })).skippedSlots) + .toEqual([101, 102, 103, 104]); + expect((await emitDueTicks({}, sink, { schedule: s, cursor, nowMs: 107 * MINUTE })).skippedSlots) + .toEqual([]); + }); +}); + +describe('P2-2: an uncomputable grid refuses instead of going quiet', () => { + // A NaN epochMs or nowMs made slotFor() return NaN, so firstDue > currentSlot + // was false-y in the wrong direction and the poll returned an empty result: + // no submit, no skip, no throw. A schedule permanently and silently zero, + // which is precisely Native's silent-death problem. + it.each([ + ['epochMs', Number.NaN], + ['epochMs', Number.POSITIVE_INFINITY], + ['epochMs', -1], + ['epochMs', 1.5], + ['nowMs', Number.NaN], + ['nowMs', Number.POSITIVE_INFINITY], + ['nowMs', -1], + ['nowMs', 1.5], + ])('refuses %s = %p', async (field, value) => { + const s = schedule(field === 'epochMs' ? { epochMs: value } : {}); + const nowMs = field === 'nowMs' ? value : 100 * MINUTE; + await expect( + emitDueTicks({}, dedupingSink(schedule()), { schedule: s, cursor: {}, nowMs }), + ).rejects.toThrow(new RegExp(`${field} must be a non-negative integer`)); + }); + + it('still accepts the zero epoch and a zero now', async () => { + const s = schedule({ epochMs: 0 }); + const sink = dedupingSink(s); + const result = await emitDueTicks({}, sink, { schedule: s, cursor: {}, nowMs: 0 }); + expect(result.emittedSlots).toEqual([0]); + }); + + it('refuses a NaN maxCatchUp too', async () => { + await expect( + emitDueTicks({}, dedupingSink(schedule()), { + schedule: schedule({ maxCatchUp: Number.NaN }), cursor: {}, nowMs: 100 * MINUTE, + }), + ).rejects.toThrow(/maxCatchUp must be a positive integer/); + }); +}); diff --git a/sdk/tests/validate.test.ts b/sdk/tests/validate.test.ts index c0c64640e..1c9f3c608 100644 --- a/sdk/tests/validate.test.ts +++ b/sdk/tests/validate.test.ts @@ -338,7 +338,12 @@ describe('validate: preflight declarations', () => { [-1, 'expected a positive number'], [1.5, 'expected an integer'], ['180000', 'expected an integer'], - [Number.MAX_VALUE, 'exceeds the i64 range'], + [Number.MAX_VALUE, 'exceeds'], + // P3: 9_223_372_036_854_775_807 is not representable in JS -- the literal + // rounds UP to 2^63, i.e. i64::MAX + 1. A `> MAX_STALE_AFTER_MS` bound + // built from it therefore ADMITS exactly the one value the kernel refuses. + [2 ** 63, 'exceeds'], + [Number.MAX_SAFE_INTEGER + 1, 'exceeds'], ])('rejects staleAfterMs %p', (value, fragment) => { // A budget the kernel cannot represent fails OPEN — the sweep can never // mark the subscription stale — so it must not compile. Zero fails the