diff --git a/ops/NEXT.md b/ops/NEXT.md index 649c80cc6..123411100 100644 --- a/ops/NEXT.md +++ b/ops/NEXT.md @@ -1,87 +1,111 @@ -# NEXT — work package for this tick +# Work package — gate 2, sub-PR A: HN monitor polling runner -**Scope:** Build a minimal agent worker in the SDK. CODE task, SDK-side. +**Scope (quoted from target):** -This run is pinned to **gate 3** and must not work on any other gate. +Build sub-PR A of the Gate 2 push: a real `hn-monitor` polling runner in the SDK. CODE task, `sdk/src/`-side. This is a scaffolding PR — proof that the workload EXECUTES end-to-end is deliberately deferred to sub-PR B (integration test). Do not conflate the two. -## Objective +Prior attempt (PR #83, closed) produced a functional runner but was rejected by the swarm on five real findings. Address them in this attempt: -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. +1. **Fail-closed on journal errors.** Split: `try { fetch } catch { onFetchError }` around the network call, `try { eventSubmit } catch { rethrow }` around the journal call. -## Context +2. **AgentWorker.close() must release the worker (or explicitly document it does not).** Either add a `workerRelease` verb to `sdk/src/protocol.ts` and call it from `close()` (preferred), OR add a one-line comment on `close()` naming exactly what shutdown intentionally does NOT do. -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. +3. **Class field declaration order.** Declare ALL fields at the top of the class body, before the constructor. -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"). +4. **Signal handlers must be opt-in via AbortSignal.** Accept `signal?: AbortSignal` in options; the CLI wrapper (sub-PR C) can create + wire a process-signal-driven AbortController. -`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. +5. **Test coverage for pollError branch.** Add both cases: the loop survives a fetcher throw AND the loop TERMINATES on a journal throw. -## Files in scope +## The task + +Add `sdk/src/hn-monitor-runner.ts`. It composes the existing pieces into a continuous runner: + + - constructs a `JournalClient` connected to the running `relayflowd` socket + - constructs an `AgentWorker` (from `sdk/src/worker.ts`) and calls `workerAttach()` for `agent` steps — attach BEFORE first poll + - loops: `pollHackerNewsOnce(spec, sink)` → sleep `POLL_INTERVAL_MS` (env-configurable, default 60000 = 60s) → repeat + - exit cleanly on `AbortSignal.abort` (drain in-flight steps, close client, release worker per finding #2) + - exported from `sdk/src/index.ts` + +Keep it small and honest: + - the worker must attach BEFORE the first poll + - the poller layer handles single-fetch failures with a typed error; the loop just moves to the next tick — but journal errors MUST fail the runner (finding #1) + - no scheduling logic beyond the sleep (the kernel owns retry and dedupe policy) + - no LLM calls; the runner is glue, not a reviewer + +## Explicit non-goals for THIS PR (belongs to later sub-PRs) -- `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`. + - Proving the workload actually executes end-to-end (dispatch → step complete). That is sub-PR B (integration test with real relayflowd + fake HN fetch + assert step reaches `done`). + - CLI wrapper (`flows hn-monitor start`). That is sub-PR C. + - ops/STATE.md gate-2 GREEN declaration. That is sub-PR D. -## Definition of done +## Definition of done (all of it) -ALL of the following must hold: + - `sdk/src/hn-monitor-runner.ts` exists, exports `HnMonitorRunner` from `sdk/src/index.ts` + - `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 + - `sdk/src/protocol.ts` — if you added `workerRelease`, matching request/response definitions + - `sdk/tests/hn-monitor-runner.test.ts` covers ALL of these: + - 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) + - `cd sdk && npm test` green (pretest hook builds the kernel automatically), with literal output pasted + - EVERY new test confirmed to FAIL against current code (comment out the source; the test fails), with the literal failing output pasted + - PR body explicitly names the non-goals (test-actually-runs is sub-PR B; CLI is sub-PR C; gate-2 declaration is sub-PR D) + - as your LAST action, run `git status --porcelain` and paste it -1. The worker in `sdk/src/worker.ts`, exported from `sdk/src/index.ts` +## Out of scope for THIS tick — DO NOT TOUCH -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. + - `.github/workflows/*` — no GHA changes + - `kernel/*` — the kernel side of gate 2 already works via PR #14 + - `workflows/*.yaml` — those are for later sub-PRs + - `ops/AUTODRIVE_BRIEF.md` — chief owns this file, not the drive loop + - CLI wrapper — sub-PR C, separate PR + - end-to-end integration test with real relayflowd — sub-PR B, separate PR + - ops/STATE.md gate-2 declaration — sub-PR D, separate PR + +## Files in scope + + - `sdk/src/hn-monitor-runner.ts` (NEW) + - `sdk/src/worker.ts` (modify `close()` only) + - `sdk/src/protocol.ts` (add `workerRelease` if needed) + - `sdk/src/index.ts` (export the runner) + - `sdk/tests/hn-monitor-runner.test.ts` (NEW) + +## Objective -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. +Compose existing primitives (JournalClient, AgentWorker, pollHackerNewsOnce) into a continuous polling runner that addresses the five review findings from PR #83, with comprehensive test coverage proving fail-closed journal error handling and fetch error survival. -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. +## Current state -5. `cd sdk && npm test` must be green. Run it and paste the literal command and - output tail showing test counts. +Both test suites pass: -6. `cd kernel && sh ../ops/cargo.sh test` must be green. Run it and paste the - literal command and output tail showing test counts. +``` +$ cd /project/workflows/runs/bc7818d9-80e8-405f-8913-9866c78c8ff1/kernel && sh ../ops/cargo.sh test --workspace 2>&1 | tail -20 + Running unittests src/lib.rs (/home/daytona/.relayflows-toolchain/target/326323060/debug/deps/relayflowd_journal-e575ab403beb54a9) -7. EVERY new test confirmed to FAIL against current code, with the literal - failing output quoted in the summary. +running 6 tests +test registry::tests::registry_is_a_rebuildable_run_locator ... ok +test tests::an_unconfirmed_election_is_reclaimed_by_the_next_attempt_not_treated_as_done ... ok +test tests::append_is_durable_and_monotonic_after_reopen ... ok +test tests::effects_are_deduplicated_at_the_journal_boundary ... ok +test tests::failed_commit_is_returned_not_swallowed ... ok +test tests::rollover_is_atomic_scaffolding_for_epoch_resume ... ok -8. As your LAST action, run `git status --porcelain` and paste it. +test result: ok. 6 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out; finished in 0.02s +``` -## Explicitly OUT of scope +``` +$ cd /project/workflows/runs/bc7818d9-80e8-405f-8913-9866c78c8ff1/sdk && npm test 2>&1 | tail -20 + ✓ tests/hn-poller.test.ts (3 tests) 4ms -- 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 + Test Files 15 passed (15) + Tests 203 passed (203) + Start at 05:19:12 + Duration 43.50s (transform 320ms, setup 0ms, collect 613ms, tests 40.86s, environment 2ms, prepare 532ms) +``` -## If blocked +Kernel: 44 tests passed (4+1+26+5+6+doc tests across all crates) +SDK: 203 tests passed across 15 test files -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. +`sdk/src/hn-monitor-runner.ts` does not exist. The work package is to create it along with its tests, addressing the five findings from PR #83. diff --git a/sdk/src/hn-monitor-runner.ts b/sdk/src/hn-monitor-runner.ts new file mode 100644 index 000000000..9a4602704 --- /dev/null +++ b/sdk/src/hn-monitor-runner.ts @@ -0,0 +1,119 @@ +import { JournalClient } from './journal-client.js'; +import { + pollHackerNewsOnce, + type EventSink, + type Fetcher, + type PollOptions, +} from './hn-poller.js'; +import type { Pins } from './protocol.js'; +import { AgentWorker } from './worker.js'; + +const DEFAULT_POLL_INTERVAL_MS = 60_000; + +interface RunnerClient extends EventSink { + connect(): Promise; + close(): void; +} + +interface RunnerWorker { + attach(): Promise; + close(): void; +} + +export interface HnMonitorRunnerOptions extends PollOptions { + socketPath: string; + workerId: string; + pins: Pins; + signal?: AbortSignal; + pollIntervalMs?: number; + onPollError?: (error: HnPollFetchError) => void; + client?: RunnerClient; + worker?: RunnerWorker; +} + +/** A transient network failure; journal/protocol failures are never wrapped. */ +export class HnPollFetchError extends Error { + constructor(cause: unknown) { + super(`HN poll fetch failed: ${cause instanceof Error ? cause.message : String(cause)}`, { cause }); + this.name = 'HnPollFetchError'; + } +} + +/** Continuously submits HN events to relayflowd until its signal is aborted. */ +export class HnMonitorRunner { + private readonly client: RunnerClient; + private readonly fetcher: Fetcher; + private readonly intervalMs: number; + private readonly options: HnMonitorRunnerOptions; + private readonly spec: unknown; + private readonly worker: RunnerWorker; + + constructor(spec: unknown, options: HnMonitorRunnerOptions) { + this.options = options; + this.intervalMs = pollInterval(options.pollIntervalMs); + this.client = options.client ?? new JournalClient(options.socketPath); + this.worker = options.worker + ?? new AgentWorker(this.client as JournalClient, { workerId: options.workerId, pins: options.pins }); + this.fetcher = typedFetcher(options.fetcher ?? fetchText); + this.spec = spec; + } + + async run(): Promise { + await this.client.connect(); + try { + await this.worker.attach(); + while (!this.options.signal?.aborted) { + try { + await pollHackerNewsOnce(this.spec, this.client, { + storyLimit: this.options.storyLimit, + createdBy: this.options.createdBy, + fetcher: this.fetcher, + }); + } catch (error) { + if (!(error instanceof HnPollFetchError)) throw error; + this.options.onPollError?.(error); + } + await waitForNextTick(this.intervalMs, this.options.signal); + } + } finally { + this.worker.close(); + this.client.close(); + } + } +} + +function pollInterval(explicit: number | undefined): number { + const raw = explicit ?? Number(process.env.POLL_INTERVAL_MS ?? DEFAULT_POLL_INTERVAL_MS); + if (!Number.isFinite(raw) || raw < 0) throw new Error('POLL_INTERVAL_MS must be a non-negative number'); + return raw; +} + +function typedFetcher(fetcher: Fetcher): Fetcher { + return async (url) => { + try { + return await fetcher(url); + } catch (error) { + throw new HnPollFetchError(error); + } + }; +} + +async function fetchText(url: string): Promise { + const response = await fetch(url); + if (!response.ok) throw new Error(`HTTP ${response.status}`); + return response.text(); +} + +function waitForNextTick(ms: number, signal?: AbortSignal): Promise { + if (signal?.aborted) return Promise.resolve(); + return new Promise((resolve) => { + const timer = setTimeout(done, ms); + signal?.addEventListener('abort', done, { once: true }); + + function done(): void { + clearTimeout(timer); + signal?.removeEventListener('abort', done); + resolve(); + } + }); +} diff --git a/sdk/src/index.ts b/sdk/src/index.ts index 8f47e7f6f..91f06ff78 100644 --- a/sdk/src/index.ts +++ b/sdk/src/index.ts @@ -144,6 +144,11 @@ export { type Fetcher, type PollOptions, } from './hn-poller.js'; +export { + HnMonitorRunner, + HnPollFetchError, + type HnMonitorRunnerOptions, +} from './hn-monitor-runner.js'; // Directory watcher — second proactive workload for gate 2 primitives. // Non-provider: no HTTP, no API tokens, no gate-6 dependency. Proves the diff --git a/sdk/src/worker.ts b/sdk/src/worker.ts index 0cfc5849b..bbebca801 100644 --- a/sdk/src/worker.ts +++ b/sdk/src/worker.ts @@ -39,6 +39,7 @@ export class AgentWorker extends EventEmitter { } close(): void { + // Does NOT explicitly release the worker; closing its JournalClient releases the connection registration. this.client.off('step.dispatch', this.onDispatch); this.attached = false; } diff --git a/sdk/tests/hn-monitor-runner.test.ts b/sdk/tests/hn-monitor-runner.test.ts new file mode 100644 index 000000000..b401f346d --- /dev/null +++ b/sdk/tests/hn-monitor-runner.test.ts @@ -0,0 +1,87 @@ +import { describe, expect, it, vi } from 'vitest'; +import { HnMonitorRunner, HnPollFetchError } from '../src/hn-monitor-runner.js'; + +function harness(overrides: Record = {}) { + const order: string[] = []; + const client = { + connect: vi.fn(async () => { order.push('connect'); }), + close: vi.fn(() => { order.push('client.close'); }), + eventSubmit: vi.fn(async () => { order.push('submit'); }), + }; + const worker = { + attach: vi.fn(async () => { order.push('attach'); }), + close: vi.fn(() => { order.push('worker.close'); }), + }; + const controller = new AbortController(); + const runner = new HnMonitorRunner({ name: 'hn-monitor' }, { + socketPath: '/unused.sock', + workerId: 'hn-test', + pins: {}, + signal: controller.signal, + pollIntervalMs: 0, + storyLimit: 1, + fetcher: async () => '[1]', + client, + worker, + ...overrides, + }); + return { client, controller, order, runner, worker }; +} + +describe('HnMonitorRunner', () => { + it('attaches before its first poll and submits on every tick', async () => { + const h = harness(); + h.client.eventSubmit.mockImplementation(async () => { + h.order.push('submit'); + if (h.client.eventSubmit.mock.calls.length === 2) h.controller.abort(); + }); + + await h.runner.run(); + + expect(h.client.eventSubmit).toHaveBeenCalledTimes(2); + expect(h.order.indexOf('attach')).toBeLessThan(h.order.indexOf('submit')); + }); + + it('abort triggers clean shutdown without waiting for the next tick', async () => { + const h = harness({ pollIntervalMs: 60_000 }); + h.client.eventSubmit.mockImplementation(async () => { h.controller.abort(); }); + + await h.runner.run(); + + expect(h.worker.close).toHaveBeenCalledOnce(); + expect(h.client.close).toHaveBeenCalledOnce(); + expect(h.order.indexOf('worker.close')).toBeLessThan(h.order.indexOf('client.close')); + }); + + it('survives a fetch throw, reports its typed error, and polls next tick', async () => { + const onPollError = vi.fn(); + let fetches = 0; + const h = harness({ + onPollError, + fetcher: async () => { + fetches += 1; + if (fetches === 1) throw new Error('network down'); + h.controller.abort(); + return '[2]'; + }, + }); + + await h.runner.run(); + + expect(fetches).toBe(2); + expect(onPollError).toHaveBeenCalledWith(expect.any(HnPollFetchError)); + expect(h.client.eventSubmit).toHaveBeenCalledOnce(); + }); + + it('terminates and closes when the journal rejects an event', async () => { + const journalError = new Error('journal write failed'); + const h = harness(); + h.client.eventSubmit.mockRejectedValue(journalError); + + await expect(h.runner.run()).rejects.toBe(journalError); + + expect(h.client.eventSubmit).toHaveBeenCalledOnce(); + expect(h.worker.close).toHaveBeenCalledOnce(); + expect(h.client.close).toHaveBeenCalledOnce(); + }); +});