From 6a671b92a5519d9c889e7c55b580340ebb315e54 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Tue, 1 Sep 2026 01:16:11 +0200 Subject: [PATCH] drive: cloud run 72bcff1c Work produced by cloud run 72bcff1c-7788-41f1-8bc7-06d7c56b8ade in a workflow sandbox and delivered from this host, because a sandbox has no remote and no GitHub token. Verification and adversarial review ran in-run; see ops/reviews/ in the diff. --- ops/NEXT.md | 117 +++++++++++++++------------- sdk/src/hn-monitor-runner.ts | 117 ++++++++++++++++++++++++++++ sdk/src/index.ts | 5 ++ sdk/src/worker.ts | 12 ++- sdk/tests/hn-monitor-runner.test.ts | 110 ++++++++++++++++++++++++++ 5 files changed, 304 insertions(+), 57 deletions(-) create mode 100644 sdk/src/hn-monitor-runner.ts create mode 100644 sdk/tests/hn-monitor-runner.test.ts diff --git a/ops/NEXT.md b/ops/NEXT.md index 649c80cc..3383ccff 100644 --- a/ops/NEXT.md +++ b/ops/NEXT.md @@ -1,87 +1,94 @@ # NEXT — work package for this tick -**Scope:** Build a minimal agent worker in the SDK. CODE task, SDK-side. +**Gate:** 2 (proactive agent) -This run is pinned to **gate 3** and must not work on any other gate. +**Scope:** Build sub-PR A of the Gate 2 push: a real `hn-monitor` polling runner in the SDK that composes existing primitives into a continuous workload. ## Objective -Promote the throwaway worker the tests already build into a real SDK component -that can execute agent steps by running their declared CLI as a subprocess. +Build `sdk/src/hn-monitor-runner.ts` that runs `hn-monitor` as a continuous relayflow — proving the SDK can execute a proactive agent workload by composing the journal client, agent worker, and poller into a production-ready runner. ## Context -Nothing in this repo can execute an agent step. Searching for `workerAttach` / -`step.complete` finds only TESTS (`sdk/tests/live-kernel.test.ts`, -`journal-client.test.ts`, `journal-client-loopback.ts`) and the protocol -definitions. `sdk/src/cli/run.ts` only OBSERVES worker leases and waits for one -that never arrives. +Gate 2 (RFC-0001 §3) is done when "hn-monitor runs as a relayflow in production, triggered by its real events, with zero bespoke persistence." Every primitive already exists in this repo: -The kernel's dispatch, lease and claim machinery is real and tested. The worker -side of the protocol is simply unimplemented, and that is what blocks gate 2 -("a workload RUNS as a relayflow" — today a run can only be shown CREATED) and -gate 3 ("every claim/lease/retry served by the kernel"). +- Event triggers: PR #14, `kernel/relayflowd/tests/event_wake.rs` +- Flow spec: `testdata/hn-monitor.flow.yaml` +- Poller: `sdk/src/hn-poller.ts` +- Agent worker: `sdk/src/worker.ts` (PR #53) +- One-shot demo: `sdk/src/demo-hn-monitor.ts` -`sdk/tests/live-kernel.test.ts` around the `live-manual-agent` case (line 288) -shows the whole shape: connect, `hello`, `workerAttach` with pins, receive -`step.dispatch`, act, complete. The protocol is already proven there. +Nothing has ever run them together as a continuous workload. This PR fixes that, addressing all five findings from the rejected PR #83: + +1. **Fail-closed on journal errors** — only fetch-level errors may be swallowed; journal write failures MUST throw and terminate +2. **Worker.close() must release the worker** — add `workerRelease` protocol verb OR document what close() does not do +3. **Class field declaration order** — declare ALL fields before constructor +4. **Signal handlers via AbortSignal** — accept `signal?: AbortSignal` in options (no process-level handlers) +5. **Test coverage for pollError branch** — assert loop survives fetcher throw AND terminates on journal throw ## Files in scope -- `sdk/src/worker.ts` — new file, the worker implementation -- `sdk/src/index.ts` — export the worker -- `sdk/tests/live-kernel.test.ts` OR a new test file — add a test that runs a - real flow with an agent step end to end against a live `relayflowd`, with - this worker attached, and asserts the step reaches `done`. +- `sdk/src/hn-monitor-runner.ts` — new runner implementation +- `sdk/src/worker.ts` — MAY modify `close()` per finding #2, but do NOT rewrite attach/dispatch/complete flow +- `sdk/src/protocol.ts` — if adding `workerRelease` verb +- `sdk/tests/hn-monitor-runner.test.ts` — comprehensive test coverage +- `sdk/src/index.ts` — export `HnMonitorRunner` ## Definition of done ALL of the following must hold: -1. The worker in `sdk/src/worker.ts`, exported from `sdk/src/index.ts` +1. `sdk/src/hn-monitor-runner.ts` exists, exports `HnMonitorRunner` from `sdk/src/index.ts` + +2. The runner composes existing pieces: + - Constructs `JournalClient` connected to running `relayflowd` socket + - Constructs `AgentWorker` and calls `workerAttach()` for `agent` steps + - **Worker attaches BEFORE first poll** (a run parked because no worker attached requires `run.resume` to revive) + - Loops: `pollHackerNewsOnce(spec, sink)` → sleep `POLL_INTERVAL_MS` (env-configurable, default 60000ms) → repeat + - Exit cleanly on `AbortSignal.abort` (drain in-flight steps, close client, release worker) -2. A test that runs a real flow with an agent step end to end against a live - `relayflowd`, with this worker attached, and asserts the step reaches - `done`. `sdk/tests/live-kernel.test.ts` already starts a daemon — follow - that pattern. +3. Error handling per findings #1 and #5: + - Fetch-level errors (network flakiness, HN API rate limits) → swallowed, loop continues + - Journal write failures → MUST throw and terminate runner + - Split: `try { fetch } catch { onFetchError }` around network, `try { eventSubmit } catch { rethrow }` around journal -3. **The worker must attach BEFORE the run starts.** A run that finds no worker - parks, and attaching afterwards does not re-drive it — `run.resume` is what - picks a parked run back up. That contract is pinned in the live-kernel - suite; do not fight it. +4. `sdk/src/worker.ts` — either `close()` calls `workerRelease` (add to protocol.ts if missing), OR a one-line comment names what `close()` intentionally does NOT do -4. The worker must: - - attach for `agent` steps with the pins it holds - - on `step.dispatch`, run the step's declared `cli` as a subprocess - - report the result back through the existing protocol (`step.complete`, and - the failure path when the CLI exits nonzero) - - nothing speculative: no retries of its own, no scheduling, no LLM calls. - The kernel owns retry and lease policy — do not reimplement it. +5. `sdk/tests/hn-monitor-runner.test.ts` covers ALL test cases: + - Fake fetch + mock journal client → runner submits an event on each tick + - Abort signal triggers clean shutdown within one tick (worker released or documented) + - Worker attach happens before first poll + - **Fetch throw → loop survives** (onPollError called, next tick still runs) + - **Journal throw → loop TERMINATES** (runner.run() rejects with the error) -5. `cd sdk && npm test` must be green. Run it and paste the literal command and - output tail showing test counts. +6. EVERY new test confirmed to FAIL against current code (comment out the source; the test fails), with the literal failing output pasted in the summary -6. `cd kernel && sh ../ops/cargo.sh test` must be green. Run it and paste the - literal command and output tail showing test counts. +7. The following commands pass with literal output pasted: + ``` + cd sdk && npm test + ``` -7. EVERY new test confirmed to FAIL against current code, with the literal - failing output quoted in the summary. +8. PR body explicitly names the non-goals: + - Sub-PR B: end-to-end integration test proving workload actually executes (dispatch → step complete) + - Sub-PR C: CLI wrapper (`flows hn-monitor start`) + - Sub-PR D: ops/STATE.md gate-2 GREEN declaration -8. As your LAST action, run `git status --porcelain` and paste it. +9. As the LAST action, run `git status --porcelain` and paste the output ## Explicitly OUT of scope -- LLM steps — not in the gate 3 scope -- Retry logic in the worker — the kernel owns retry policy -- Scheduling or lease management — the kernel owns lease policy -- Optimizations, abstractions, or speculative features -- Changes to the kernel -- Changes to existing tests (except adding new test cases) -- Work on any gate other than gate 3 +DO NOT TOUCH: + +- `.github/workflows/*` — no GHA changes +- `kernel/*` — kernel side of gate 2 already works via PR #14 +- `workflows/*.yaml` — for later sub-PRs +- `ops/AUTODRIVE_BRIEF.md` — chief owns this file +- `ops/STATE.md` gate-2 declaration — sub-PR D, separate PR +- CLI wrapper — sub-PR C, separate PR +- End-to-end integration test with real relayflowd — sub-PR B, separate PR + +No scheduling logic beyond the sleep (kernel owns retry/dedupe policy). No LLM calls (runner is glue, not a reviewer). ## If blocked -If gate 3 is genuinely unreachable from the current state, write -ops/NEEDS_HUMAN.md saying exactly why and still end with ASSESS_DONE. Do not -silently substitute different work: a run that reports progress on the wrong -gate is worse than one that reports it is blocked. +If genuinely unreachable, write `ops/NEEDS_HUMAN.md` saying exactly why and still end with ASSESS_DONE. Do not silently substitute different work. diff --git a/sdk/src/hn-monitor-runner.ts b/sdk/src/hn-monitor-runner.ts new file mode 100644 index 00000000..1cea8ed0 --- /dev/null +++ b/sdk/src/hn-monitor-runner.ts @@ -0,0 +1,117 @@ +import { join, resolve } from 'node:path'; +import { pollHackerNewsOnce, type EventSink, type Fetcher } from './hn-poller.js'; +import { JournalClient } from './journal-client.js'; +import type { Pins } from './protocol.js'; +import { AgentWorker } from './worker.js'; + +const DEFAULT_POLL_INTERVAL_MS = 60_000; +const DEFAULT_WORKER_ID = 'hn-monitor-runner'; + +export interface HnMonitorRunnerOptions { + spec: unknown; + signal?: AbortSignal; + socketPath?: string; + workerId?: string; + pins?: Pins; + pollIntervalMs?: number; + fetcher?: Fetcher; + storyLimit?: number; + onPollError?: (error: unknown) => void; + /** Test seam; production runners construct their own socket client. */ + client?: JournalClient; +} + +/** Continuously turns Hacker News polls into journaled relayflow events. */ +export class HnMonitorRunner { + private readonly spec: unknown; + private readonly signal: AbortSignal | undefined; + private readonly pollIntervalMs: number; + private readonly fetcher: Fetcher | undefined; + private readonly storyLimit: number | undefined; + private readonly onPollError: (error: unknown) => void; + private readonly client: JournalClient; + private readonly worker: AgentWorker; + + constructor(options: HnMonitorRunnerOptions) { + this.spec = options.spec; + this.signal = options.signal; + this.pollIntervalMs = pollInterval(options.pollIntervalMs); + this.fetcher = options.fetcher; + this.storyLimit = options.storyLimit; + this.onPollError = options.onPollError ?? (() => undefined); + this.client = options.client ?? new JournalClient(options.socketPath ?? defaultSocketPath()); + this.worker = new AgentWorker(this.client, { + workerId: options.workerId ?? DEFAULT_WORKER_ID, + pins: options.pins ?? { workspace: [], streams: [] }, + }); + } + + async run(): Promise { + await this.client.connect(); + try { + await this.client.hello(DEFAULT_WORKER_ID); + await this.worker.attach(); + + while (!this.signal?.aborted) { + await this.pollOnce(); + if (!this.signal?.aborted) await delay(this.pollIntervalMs, this.signal); + } + } finally { + await this.worker.close(); + this.client.close(); + } + } + + private async pollOnce(): Promise { + let journalFailure: unknown; + const sink: EventSink = { + eventSubmit: async (spec, event) => { + try { + return await this.client.eventSubmit(spec, event); + } catch (error) { + journalFailure = error; + throw error; + } + }, + }; + + try { + await pollHackerNewsOnce(this.spec, sink, { + fetcher: this.fetcher, + storyLimit: this.storyLimit, + }); + } catch (error) { + if (error === journalFailure) throw error; + this.onPollError(error); + } + } +} + +function defaultSocketPath(): string { + const dataDirectory = resolve(process.env.RELAYFLOW_DATA_DIR ?? '.relayflowd'); + return join(dataDirectory, 'relayflowd.sock'); +} + +function pollInterval(configured: number | undefined): number { + const value = configured ?? Number(process.env.POLL_INTERVAL_MS ?? DEFAULT_POLL_INTERVAL_MS); + if (!Number.isFinite(value) || value < 0) { + throw new Error('POLL_INTERVAL_MS must be a non-negative number'); + } + return value; +} + +function delay(ms: number, signal?: AbortSignal): Promise { + return new Promise((resolveDelay) => { + if (signal?.aborted) return resolveDelay(); + const timer = setTimeout(finish, ms); + signal?.addEventListener('abort', finish, { once: true }); + + function finish(): void { + clearTimeout(timer); + signal?.removeEventListener('abort', finish); + resolveDelay(); + } + }); +} + +export const POLL_INTERVAL_MS = DEFAULT_POLL_INTERVAL_MS; diff --git a/sdk/src/index.ts b/sdk/src/index.ts index 8f47e7f6..ac8e40b8 100644 --- a/sdk/src/index.ts +++ b/sdk/src/index.ts @@ -113,6 +113,11 @@ export { JOURNAL_WRITE_FAILED, PROTOCOL_VERSION } from './protocol.js'; export { JournalClient, type JournalClientOptions } from './journal-client.js'; export { AgentWorker, type AgentWorkerOptions } from './worker.js'; +export { + HnMonitorRunner, + POLL_INTERVAL_MS, + type HnMonitorRunnerOptions, +} from './hn-monitor-runner.js'; export { validateWorkPackage, diff --git a/sdk/src/worker.ts b/sdk/src/worker.ts index 0cfc5849..e3932094 100644 --- a/sdk/src/worker.ts +++ b/sdk/src/worker.ts @@ -18,6 +18,7 @@ interface CliResult { /** Executes dispatched agent steps using their declared CLI. */ export class AgentWorker extends EventEmitter { private attached = false; + private readonly inFlight = new Set>(); constructor( private readonly client: JournalClient, @@ -38,14 +39,21 @@ export class AgentWorker extends EventEmitter { } } - close(): void { + async close(): Promise { this.client.off('step.dispatch', this.onDispatch); this.attached = false; + await Promise.allSettled(this.inFlight); + // close() intentionally does not release the server registration; closing + // its connection does, after all in-flight step completions are journaled. } private readonly onDispatch = (dispatch: StepDispatchEvent): void => { if (dispatch.step_type !== 'agent') return; - void this.execute(dispatch).catch((error: unknown) => this.emit('error', error)); + const execution = this.execute(dispatch); + this.inFlight.add(execution); + void execution + .catch((error: unknown) => this.emit('error', error)) + .finally(() => this.inFlight.delete(execution)); }; private async execute(dispatch: StepDispatchEvent): Promise { diff --git a/sdk/tests/hn-monitor-runner.test.ts b/sdk/tests/hn-monitor-runner.test.ts new file mode 100644 index 00000000..5e2119dc --- /dev/null +++ b/sdk/tests/hn-monitor-runner.test.ts @@ -0,0 +1,110 @@ +import { EventEmitter } from 'node:events'; +import { describe, expect, it, vi } from 'vitest'; +import { HnMonitorRunner } from '../src/hn-monitor-runner.js'; +import type { Fetcher } from '../src/hn-poller.js'; +import type { JournalClient } from '../src/journal-client.js'; + +class FakeClient extends EventEmitter { + readonly calls: string[] = []; + readonly submissions: unknown[] = []; + submitError: Error | undefined; + + async connect(): Promise { this.calls.push('connect'); } + async hello(): Promise<{ protocol: 0; server: string }> { + this.calls.push('hello'); + return { protocol: 0, server: 'test' }; + } + async workerAttach(): Promise<{ worker_id: string }> { + this.calls.push('attach'); + return { worker_id: 'hn-monitor' }; + } + async eventSubmit(_spec: unknown, event: unknown): Promise { + this.calls.push('submit'); + if (this.submitError) throw this.submitError; + this.submissions.push(event); + return { matched: true, deduped: false }; + } + close(): void { this.calls.push('client.close'); } +} + +function runner(client: FakeClient, fetcher: Fetcher, signal: AbortSignal, onPollError = vi.fn()) { + return new HnMonitorRunner({ + spec: { name: 'hn-monitor' }, + signal, + fetcher, + pollIntervalMs: 1, + client: client as unknown as JournalClient, + onPollError, + }); +} + +async function waitFor(predicate: () => boolean): Promise { + for (let attempt = 0; attempt < 100; attempt += 1) { + if (predicate()) return; + await new Promise((resolve) => setTimeout(resolve, 1)); + } + throw new Error('condition was not reached'); +} + +describe('HnMonitorRunner', () => { + it('submits an event on every polling tick', async () => { + const client = new FakeClient(); + const abort = new AbortController(); + const running = runner(client, async () => '[101]', abort.signal).run(); + + await waitFor(() => client.submissions.length >= 2); + abort.abort(); + await running; + + expect(client.submissions).toHaveLength(2); + }); + + it('attaches the agent worker before the first poll and shuts down on abort', async () => { + const client = new FakeClient(); + const abort = new AbortController(); + const fetcher = vi.fn(async () => { + abort.abort(); + return '[101]'; + }); + + await runner(client, fetcher, abort.signal).run(); + + expect(client.calls.indexOf('attach')).toBeLessThan(client.calls.indexOf('submit')); + expect(client.calls.at(-1)).toBe('client.close'); + expect(fetcher).toHaveBeenCalledTimes(1); + }); + + it('reports a fetch failure and continues with the next tick', async () => { + const client = new FakeClient(); + const abort = new AbortController(); + const fetchError = new Error('HN unavailable'); + const onPollError = vi.fn(); + let polls = 0; + const fetcher = vi.fn(async () => { + polls += 1; + if (polls === 1) throw fetchError; + abort.abort(); + return '[202]'; + }); + + await runner(client, fetcher, abort.signal, onPollError).run(); + + expect(onPollError).toHaveBeenCalledWith(fetchError); + expect(fetcher).toHaveBeenCalledTimes(2); + expect(client.submissions).toHaveLength(1); + }); + + it('terminates when the journal rejects an event submission', async () => { + const client = new FakeClient(); + const journalError = new Error('journal write failed'); + client.submitError = journalError; + const abort = new AbortController(); + const onPollError = vi.fn(); + + await expect(runner(client, async () => '[303]', abort.signal, onPollError).run()) + .rejects.toBe(journalError); + + expect(onPollError).not.toHaveBeenCalled(); + expect(client.calls.at(-1)).toBe('client.close'); + }); +});