diff --git a/ops/reviews/20260903-scheduled-trigger-design.md b/ops/reviews/20260903-scheduled-trigger-design.md new file mode 100644 index 000000000..3315915ad --- /dev/null +++ b/ops/reviews/20260903-scheduled-trigger-design.md @@ -0,0 +1,563 @@ +# 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** | 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 +> 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) ++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: **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. + +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 +``` + +--- + +## 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 + 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.** *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(…))` + 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..a684f80bc 100644 --- a/sdk/src/index.ts +++ b/sdk/src/index.ts @@ -162,3 +162,23 @@ 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, + TickEmitError, + 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..2f9bb5e07 --- /dev/null +++ b/sdk/src/tick-source.ts @@ -0,0 +1,318 @@ +/** + * 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)}`; +} + +/** + * 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. + * + * 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'); + requireNonNegativeInteger(schedule.epochMs ?? 0, 'epochMs'); + requireNonNegativeInteger(nowMs, 'nowMs'); + 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; + + // 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 = { + schedule_id: schedule.scheduleId, + slot, + scheduled_for_ms: scheduledFor, + interval_ms: schedule.intervalMs, + emitted_at_ms: nowMs, + lag_ms: nowMs - scheduledFor, + }; + 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; + result.emittedSlots.push(slot); + } + + return result; +} diff --git a/sdk/src/validate.ts b/sdk/src/validate.ts index 9358c9dcd..64386c9a9 100644 --- a/sdk/src/validate.ts +++ b/sdk/src/validate.ts @@ -71,8 +71,27 @@ 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. + * + * 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 = Number.MAX_SAFE_INTEGER; + class Validator { private errors: string[] = []; private ids = new Set(); @@ -217,6 +236,32 @@ 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 ${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/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..549c5888f --- /dev/null +++ b/sdk/tests/tick-source.test.ts @@ -0,0 +1,409 @@ +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, + TickEmitError, + 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' }); + }); +}); + +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 35e1d91b6..1c9f3c608 100644 --- a/sdk/tests/validate.test.ts +++ b/sdk/tests/validate.test.ts @@ -324,6 +324,40 @@ 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'], + // 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 + // 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