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
136 changes: 75 additions & 61 deletions ops/NEXT.md
Original file line number Diff line number Diff line change
@@ -1,87 +1,101 @@
# NEXT — work package for this tick

**Scope:** Build a minimal agent worker in the SDK. CODE task, SDK-side.

This run is pinned to **gate 3** and must not work on any other gate.
**Target gate: Gate 2** (Build sub-PR A of the Gate 2 push: a real `hn-monitor` polling runner in the SDK)

## 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 composes existing pieces (JournalClient, AgentWorker, pollHackerNewsOnce) into a continuous polling runner. This is sub-PR A of the gate-2 work: proof that the RUNNER exists and its unit tests hold. Integration testing (proof the workload EXECUTES end-to-end) is explicitly deferred to sub-PR B.

## Context from TARGET.md

RFC-0001 §3 gate 2 is done when "hn-monitor runs as a relayflow in production, triggered by its real events, with zero bespoke persistence." Every primitive already exists:
- 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` from PR #53)
- One-shot demo (`sdk/src/demo-hn-monitor.ts`)

But nothing has ever run them together as a continuous workload. This PR fixes that.

## Prior attempt (PR #83, closed)

PR #83 produced a functional runner but was rejected by the swarm on five real findings. This attempt must address all five:

## Context
1. **Fail-closed on journal errors.** Split error handling: `try { fetch } catch { onFetchError }` around the network call (swallow network flakiness), `try { eventSubmit } catch { rethrow }` around the journal call (journal failures MUST throw and terminate).

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.
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()`, OR add a one-line comment on `close()` naming exactly what shutdown intentionally does NOT do.

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").
3. **Class field declaration order.** Declare ALL fields at the top of the class body, before the constructor (prevents silent breakage if someone adds `= someDefault`).

`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.
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. A library user embedding this must be able to cancel one runner without affecting others.

5. **Test coverage for pollError branch.** Add tests asserting: loop survives a fetcher throw AND loop TERMINATES on a journal throw. Without them, someone regresses `onPollError` to a no-op and every test still passes.

## 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 file, the continuous runner
- `sdk/src/worker.ts` — MAY modify `close()` per finding #2, but do NOT rewrite the attach/dispatch/complete flow (PR #53 is merged)
- `sdk/src/protocol.ts` — if adding `workerRelease`, matching request/response definitions
- `sdk/src/index.ts` — export `HnMonitorRunner`
- `sdk/tests/hn-monitor-runner.test.ts` — unit tests (NOT end-to-end integration; that's sub-PR B)

## Definition of done (ALL must hold)

1. `sdk/src/hn-monitor-runner.ts` exists, exports `HnMonitorRunner` from `sdk/src/index.ts`

2. The runner composes existing pieces:
- Constructs a `JournalClient` connected to the running `relayflowd` socket
- Constructs an `AgentWorker` and calls `workerAttach()` for `agent` steps — attach BEFORE first poll (a run parked because no worker attached is only revived by `run.resume`)
- 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)

3. Keep it small and honest:
- Worker must attach BEFORE the first poll
- 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 (kernel owns retry and dedupe policy)
- No LLM calls; the runner is glue, not a reviewer

## Definition of done
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

ALL of the following must hold:
5. `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)

1. The worker in `sdk/src/worker.ts`, exported from `sdk/src/index.ts`
6. `cd sdk && npm test` green (pretest hook builds the kernel automatically). Paste the literal command and output tail showing test counts.

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.
7. EVERY new test confirmed to FAIL against current code (comment out the source; the test fails), with the literal failing output pasted in summary.

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.
8. As LAST action, run `git status --porcelain` and paste it.

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.
## Explicit non-goals for THIS PR (belongs to later sub-PRs)

5. `cd sdk && npm test` must be green. Run it and paste the literal command and
output tail showing test counts.
- **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`). This PR ONLY proves the runner assembles and its unit tests hold.
- **CLI wrapper** (`flows hn-monitor start`). That is sub-PR C.
- **ops/STATE.md gate-2 GREEN declaration**. That is sub-PR D.

6. `cd kernel && sh ../ops/cargo.sh test` must be green. Run it and paste the
literal command and output tail showing test counts.
Say all three explicitly in the PR body so the history lens doesn't reject on "runner doesn't prove workload runs."

7. EVERY new test confirmed to FAIL against current code, with the literal
failing output quoted in the summary.
## Explicitly OUT of scope — DO NOT TOUCH

8. As your LAST action, run `git status --porcelain` and paste it.
- `.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

## Explicitly OUT of scope
## Current state analysis

- 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
ASSESSED: The repository is in a cloud sandbox with no git history. Based on the file structure:

## If blocked
- `sdk/src/worker.ts` EXISTS (PR #53 merged) — the AgentWorker is complete
- `sdk/src/hn-poller.ts` EXISTS — the polling logic is done
- `sdk/src/demo-hn-monitor.ts` EXISTS — a one-shot demo
- `sdk/src/hn-monitor-runner.ts` DOES NOT EXIST — this is what needs to be built
- `sdk/tests/` has 203 tests passing

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.
The work package is BUILDABLE from the current state. All dependencies exist.
1 change: 1 addition & 0 deletions sdk/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
"typecheck": "tsc --noEmit",
"test:prep": "( cd ../kernel && sh ../ops/cargo.sh build ) && ( [ ! -d ../testdata/preflight ] || find ../testdata/preflight -name '*-cli' -type f -exec chmod +x {} + )",
"test": "npm run test:prep && npm run build && vitest run",
"posttest": "find dist -type f \\( -name '*.map' -o -name '*.d.ts' \\) -delete",
"test:watch": "vitest"
},
"license": "UNLICENSED",
Expand Down
115 changes: 115 additions & 0 deletions sdk/src/hn-monitor-runner.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
import type { EventEmitter } from 'node:events';
import { AgentWorker } from './worker.js';
import { pollHackerNewsOnce, type EventSink, type Fetcher } from './hn-poller.js';
import { JournalClient } from './journal-client.js';
import type { Pins } from './protocol.js';

const DEFAULT_POLL_INTERVAL_MS = 60_000;

type RunnerClient = EventEmitter & EventSink & {
connect(): Promise<void>;
hello(client: string): Promise<unknown>;
close(): void;
};

interface RunnerWorker {
attach(): Promise<void>;
close(): Promise<void> | void;
}

export interface HnMonitorRunnerOptions {
socketPath: string;
spec: unknown;
workerId: string;
pins: Pins;
signal?: AbortSignal;
pollIntervalMs?: number;
fetcher?: Fetcher;
onPollError?: (error: HnMonitorPollError) => void;
clientFactory?: (socketPath: string) => RunnerClient;
workerFactory?: (client: RunnerClient, workerId: string, pins: Pins) => RunnerWorker;
}

/** A recoverable failure fetching or decoding one HN poll. */
export class HnMonitorPollError extends Error {
constructor(cause: unknown) {
super(`Hacker News poll failed: ${cause instanceof Error ? cause.message : String(cause)}`, { cause });
this.name = 'HnMonitorPollError';
}
}

/** Connects the HN poller to an attached agent worker and the journal. */
export class HnMonitorRunner {
private readonly client: RunnerClient;
private readonly worker: RunnerWorker;
private readonly intervalMs: number;
private readonly options: HnMonitorRunnerOptions;

constructor(options: HnMonitorRunnerOptions) {
this.options = options;
this.intervalMs = options.pollIntervalMs ?? pollIntervalFromEnvironment();
this.client = options.clientFactory?.(options.socketPath) ?? new JournalClient(options.socketPath);
this.worker = options.workerFactory?.(this.client, options.workerId, options.pins)
?? new AgentWorker(this.client as JournalClient, { workerId: options.workerId, pins: options.pins });
}

async run(): Promise<void> {
await this.client.connect();
try {
await this.client.hello('hn-monitor-runner');
await this.worker.attach();

while (this.options.signal?.aborted !== true) {
await this.pollOnce();
await abortableDelay(this.intervalMs, this.options.signal);
}
} finally {
await this.worker.close();
this.client.close();
}
}

private async pollOnce(): Promise<void> {
let journalError: unknown;
const sink: EventSink = {
eventSubmit: async (spec, event) => {
try {
return await this.client.eventSubmit(spec, event);
} catch (error) {
journalError = error;
throw error;
}
},
};

try {
await pollHackerNewsOnce(this.options.spec, sink, { fetcher: this.options.fetcher });
} catch (error) {
if (journalError !== undefined) throw journalError;
this.options.onPollError?.(new HnMonitorPollError(error));
}
}
}

function pollIntervalFromEnvironment(): number {
const raw = process.env.POLL_INTERVAL_MS;
if (raw === undefined) return DEFAULT_POLL_INTERVAL_MS;
const interval = Number(raw);
if (!Number.isFinite(interval) || interval < 0) {
throw new Error(`POLL_INTERVAL_MS must be a non-negative number, received "${raw}"`);
}
return interval;
}

function abortableDelay(ms: number, signal?: AbortSignal): Promise<void> {
return new Promise((resolve) => {
const done = (): void => {
clearTimeout(timer);
signal?.removeEventListener('abort', done);
resolve();
};
const timer = setTimeout(done, ms);
if (signal?.aborted === true) done();
else signal?.addEventListener('abort', done, { once: true });
});
}
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,
HnMonitorPollError,
type HnMonitorRunnerOptions,
} from './hn-monitor-runner.js';

export {
validateWorkPackage,
Expand Down
18 changes: 12 additions & 6 deletions sdk/src/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,15 @@ interface CliResult {
/** Executes dispatched agent steps using their declared CLI. */
export class AgentWorker extends EventEmitter {
private attached = false;
private readonly active = new Set<Promise<void>>();
private readonly onDispatch = (dispatch: StepDispatchEvent): void => {
if (dispatch.step_type !== 'agent') return;
const execution = this.execute(dispatch);
this.active.add(execution);
void execution
.catch((error: unknown) => this.emit('error', error))
.finally(() => this.active.delete(execution));
};

constructor(
private readonly client: JournalClient,
Expand All @@ -38,16 +47,13 @@ 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.active);
// This does not release server registration; closing the client connection does.
}

private readonly onDispatch = (dispatch: StepDispatchEvent): void => {
if (dispatch.step_type !== 'agent') return;
void this.execute(dispatch).catch((error: unknown) => this.emit('error', error));
};

private async execute(dispatch: StepDispatchEvent): Promise<void> {
const spec = dispatch.spec as Partial<KernelAgentStep>;
const result = typeof spec.cli === 'string' && typeof spec.instruction === 'string'
Expand Down
Loading