Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
117 changes: 62 additions & 55 deletions ops/NEXT.md
Original file line number Diff line number Diff line change
@@ -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.
117 changes: 117 additions & 0 deletions sdk/src/hn-monitor-runner.ts
Original file line number Diff line number Diff line change
@@ -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<void> {
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<void> {
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<void> {
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;
5 changes: 5 additions & 0 deletions sdk/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
12 changes: 10 additions & 2 deletions sdk/src/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Promise<void>>();

constructor(
private readonly client: JournalClient,
Expand All @@ -38,14 +39,21 @@ export class AgentWorker extends EventEmitter {
}
}

close(): void {
async close(): Promise<void> {
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<void> {
Expand Down
Loading