From ef4b1699ce8ea20a08d32e8ddb42eb918ca559fe Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Fri, 18 Sep 2026 18:06:55 +0000 Subject: [PATCH 1/5] test(runtime): pin AgentStepFinished field population from signals end to end Run-level regression tests for issue #35: session/structured-output/ error signals land on AgentStepFinished.sessionId/.output/.error, and the two-tier structured-output resolution (signal value, then finalText re-parse) is unchanged. --- src/runtime/run-signals.test.ts | 112 ++++++++++++++++++++++++++++++++ 1 file changed, 112 insertions(+) create mode 100644 src/runtime/run-signals.test.ts diff --git a/src/runtime/run-signals.test.ts b/src/runtime/run-signals.test.ts new file mode 100644 index 0000000..384dfab --- /dev/null +++ b/src/runtime/run-signals.test.ts @@ -0,0 +1,112 @@ +/** + * The run-level contract behind issue #35's acceptance criteria: the adapter + * supplies normalized signals (ADR 0012 §2), and `AgentStepFinished`'s + * `sessionId`, `output` and `error` fields are populated from them exactly as + * before — including the two-tier structured-output resolution (signal value + * first, `finalText` re-parse as tier-2 fallback) and the RUN_ERROR-shaped + * failure outcome. + */ + +import { describe, expect, test } from "bun:test"; +import type { RunEvent } from "../events"; +import { defineWorkflow, Schema } from "../workflow"; +import { createSlowFakeAdapter } from "../replay/adapter"; +import { startRun } from "./run"; + +const TEXT_CHUNKS = [ + { type: "TEXT_MESSAGE_START", messageId: "m1" }, + { type: "TEXT_MESSAGE_CONTENT", delta: '{"where":"from-final-text"}' }, + { type: "TEXT_MESSAGE_END", messageId: "m1" }, +]; + +const SIGNAL_SCHEMA = Schema.Struct({ where: Schema.String }); + +function agentWorkflow(outputSchema?: Schema.Codec) { + return defineWorkflow("signal-e2e", { + input: Schema.Struct({}), + run: async (ctx) => { + const result = await ctx.agent("step", "irrelevant, replay ignores it", { + ...(outputSchema !== undefined ? { output: outputSchema } : {}), + }); + return { sessionId: result.sessionId, error: result.error, output: result.output }; + }, + }); +} + +async function runWith(chunks: ReadonlyArray, signals: Parameters[2]) { + const events: Array = []; + const handle = startRun(agentWorkflow(), { + runId: "run-signals", + dir: "/tmp", + input: {}, + adapter: createSlowFakeAdapter(chunks, 1, signals), + onEvent: (event) => events.push(event), + }); + const outcome = await handle.result; + return { outcome, events }; +} + +function stepFinished(events: ReadonlyArray) { + const event = events.find((e) => e.payload._tag === "AgentStepFinished"); + if (event === undefined || event.payload._tag !== "AgentStepFinished") { + throw new Error("no AgentStepFinished event"); + } + return event.payload; +} + +describe("signals populate AgentStepFinished as before (issue #35)", () => { + test("a session signal lands on AgentStepFinished.sessionId", async () => { + const { events } = await runWith(TEXT_CHUNKS, [ + { index: 0, signal: { kind: "session", sessionId: "ses_e2e" } }, + ]); + expect(stepFinished(events).sessionId).toBe("ses_e2e"); + }); + + test("a structured-output signal is decoded into AgentStepFinished.output (tier 1)", async () => { + const events: Array = []; + const handle = startRun(agentWorkflow(SIGNAL_SCHEMA), { + runId: "run-signals-output", + dir: "/tmp", + input: {}, + adapter: createSlowFakeAdapter( + [ + ...TEXT_CHUNKS, + { type: "CUSTOM", name: "anything-at-all" }, + ], + 1, + [{ index: 3, signal: { kind: "structured-output", value: { where: "from-signal" } } }], + ), + onEvent: (event) => events.push(event), + }); + const outcome = await handle.result; + expect(outcome.outcome).toBe("completed"); + + expect(stepFinished(events).output).toEqual({ where: "from-signal" }); + }); + + test("without a signal, tier 2 re-parses the final text — the fallback is unchanged", async () => { + const events: Array = []; + const handle = startRun(agentWorkflow(SIGNAL_SCHEMA), { + runId: "run-signals-tier2", + dir: "/tmp", + input: {}, + adapter: createSlowFakeAdapter(TEXT_CHUNKS, 1), + onEvent: (event) => events.push(event), + }); + const outcome = await handle.result; + expect(outcome.outcome).toBe("completed"); + + expect(stepFinished(events).output).toEqual({ where: "from-final-text" }); + }); + + test("an error signal fails the step and lands on AgentStepFinished.error", async () => { + const { outcome, events } = await runWith([...TEXT_CHUNKS, { type: "RUN_FINISHED" }], [ + { index: 3, signal: { kind: "error", message: "sandbox vanished" } }, + ]); + expect(outcome.outcome).toBe("failed"); + + const finished = stepFinished(events); + expect(finished.outcome).toBe("failed"); + expect(finished.error).toBe("sandbox vanished"); + }); +}); From bb497d44a4b94eeac81a61a51c262d159a0ce0e5 Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Sat, 19 Sep 2026 09:53:50 +0000 Subject: [PATCH 2/5] feat(runtime): add Effect composition root and agent runtime service closes #36 - Introduce AgentRuntime context service and AgentRuntimeLayer - Resolve adapter from context in agent-step.ts instead of threading it - Add agent.adapter to factory.config.ts with opencode default - Pass ManagedRuntime through daemon/http/runs/scheduler/cli - Remove adapter field from ServerOptions, StartTrackedRunOptions, DispatchEnv - Update tests to provide agent runtime layer --- e2e/implement-issue.test.ts | 5 +- e2e/server.ts | 45 +++--- src/cli.start.test.ts | 18 +-- src/cli.ts | 28 ++-- src/config.ts | 13 ++ src/replay/adapter.test.ts | 166 +++++++++++++++++++---- src/runtime/agent-runtime.test.ts | 38 ++++++ src/runtime/agent-runtime.ts | 48 +++++++ src/runtime/agent-step.test.ts | 213 ++++++++++++++--------------- src/runtime/run-signals.test.ts | 36 ++--- src/runtime/run.test.ts | 98 ++++++-------- src/runtime/run.ts | 158 ++++++---------------- src/server/concurrency.test.ts | 49 ++++--- src/server/daemon.ts | 27 ++-- src/server/dedupe.test.ts | 33 ++--- src/server/http.test.ts | 81 +++++++---- src/server/http.ts | 33 ++--- src/server/nested-runs.test.ts | 19 +-- src/server/runs.test.ts | 30 ++-- src/server/runs.ts | 218 ++++++++---------------------- src/server/scheduler.ts | 11 +- src/workflow.dispatch.test.ts | 20 +-- 22 files changed, 707 insertions(+), 680 deletions(-) create mode 100644 src/runtime/agent-runtime.test.ts create mode 100644 src/runtime/agent-runtime.ts diff --git a/e2e/implement-issue.test.ts b/e2e/implement-issue.test.ts index 4527ed0..7f00ec2 100644 --- a/e2e/implement-issue.test.ts +++ b/e2e/implement-issue.test.ts @@ -26,6 +26,7 @@ import { join } from "node:path"; import { describe, expect, test } from "bun:test"; import { hostExec } from "../src/lib/exec"; import { createCorpusReplayAdapter } from "../src/replay/adapter"; +import { makeAgentRuntime } from "../src/runtime/agent-runtime"; import { startRun } from "../src/runtime/run"; import implementIssue from "./implement-issue"; @@ -118,14 +119,14 @@ describe("implement-issue workflow, replayed against the recorded round-trip cor process.env.PATH = `${binDir}:${originalPath}`; const events: Array = []; - const handle = startRun(implementIssue, { + const runtime = makeAgentRuntime(createCorpusReplayAdapter(FULL_ROUND_TRIP_CORPUS)); + const handle = await startRun(implementIssue, runtime, { runId: "test-run-corpus-replay", dir: workDir, input: { issueNumber: 1, }, repo: { slug: "local/fixture", baseBranch: "main" }, - adapter: createCorpusReplayAdapter(FULL_ROUND_TRIP_CORPUS), onEvent: (event) => events.push(event), }); diff --git a/e2e/server.ts b/e2e/server.ts index 99be761..e74db4a 100644 --- a/e2e/server.ts +++ b/e2e/server.ts @@ -22,6 +22,7 @@ import { defineConfig } from "../src/config"; import { hostExec } from "../src/lib/exec"; import { appendEvent, openStore } from "../src/persistence/store"; import { createCorpusReplayAdapter, createSlowFakeAdapter } from "../src/replay/adapter"; +import { makeAgentRuntime } from "../src/runtime/agent-runtime"; import { startRun } from "../src/runtime/run"; import { startDaemon } from "../src/server/daemon"; import { defineWorkflow, Schema } from "../src/workflow"; @@ -102,10 +103,6 @@ async function git(dir: string, args: ReadonlyArray): Promise { if (result.exitCode !== 0) throw new Error(`git ${args.join(" ")} failed: ${result.stderr}`); } -async function awaitRun(handle: ReturnType): Promise { - await handle.result; -} - async function seedCorpusRuns(root: string, db: ReturnType): Promise { const remoteDir = join(root, "remote.git"); const workDir = join(root, "work"); @@ -136,29 +133,27 @@ async function seedCorpusRuns(root: string, db: ReturnType): P chmodSync(fakeGhPath, 0o755); process.env.PATH = `${binDir}:${originalPath}`; - await awaitRun( - startRun(implementIssue, { - runId: "run-static-corpus", - dir: workDir, - repo: { slug: "local/fixture", baseBranch: "main" }, - input: { issueNumber: 1 }, - adapter: createCorpusReplayAdapter(CORPUS_ROUND_TRIP), - onEvent: (event) => appendEvent(db, event), - }), - ); + const corpusRuntime = makeAgentRuntime(createCorpusReplayAdapter(CORPUS_ROUND_TRIP)); + const corpusHandle = await startRun(implementIssue, corpusRuntime, { + runId: "run-static-corpus", + dir: workDir, + repo: { slug: "local/fixture", baseBranch: "main" }, + input: { issueNumber: 1 }, + onEvent: (event) => appendEvent(db, event), + }); + await corpusHandle.result; await Bun.sleep(20); // A second, cheap run so list ordering has more than one data point. - await awaitRun( - startRun(echoWorkflow, { - runId: "run-static-echo", - dir: workDir, - input: {}, - adapter: createCorpusReplayAdapter(CORPUS_ONE_STEP), - onEvent: (event) => appendEvent(db, event), - }), - ); + const echoRuntime = makeAgentRuntime(createCorpusReplayAdapter(CORPUS_ONE_STEP)); + const echoHandle = await startRun(echoWorkflow, echoRuntime, { + runId: "run-static-echo", + dir: workDir, + input: {}, + onEvent: (event) => appendEvent(db, event), + }); + await echoHandle.result; await Bun.sleep(20); @@ -253,7 +248,6 @@ async function main(): Promise { db.close(); } - const adapter = createSlowFakeAdapter(SLOW_CHUNKS, 1_000); const config = defineConfig({ repo: { sshUrl: remoteDir, @@ -278,8 +272,9 @@ async function main(): Promise { timezone: "UTC", }, ], + agent: { adapter: createSlowFakeAdapter(SLOW_CHUNKS, 1_000) }, }); - const { server } = await startDaemon({ dbPath, port, adapter, config }); + const { server } = await startDaemon({ dbPath, port, config }); console.log(`factory e2e: listening on http://localhost:${server.port}`); } diff --git a/src/cli.start.test.ts b/src/cli.start.test.ts index 4757dbf..897a004 100644 --- a/src/cli.start.test.ts +++ b/src/cli.start.test.ts @@ -25,15 +25,17 @@ async function startTestDaemon(delayMs: number, root: string): Promise { const runId = `run-${Date.now()}`; let repo: RunRepo | undefined; + let adapter = options.adapter ?? opencodeAdapter; try { const config = await loadFactoryConfig(); repo = { slug: config.repo.slug, baseBranch: config.repo.baseBranch }; + adapter = options.adapter ?? config.agent.adapter; } catch { repo = undefined; } - const handle = startRun(workflow, { + const runtime = ManagedRuntime.make(AgentRuntimeLayer(adapter)); + const handle = await startRun(workflow, runtime, { runId, dir: options.dir, input: options.input, - adapter: options.adapter, - prepareWorkspace: options.clone !== undefined, ...(repo !== undefined ? { repo } : {}), onEvent: (event) => { console.log(formatEvent(event)); diff --git a/src/config.ts b/src/config.ts index adaf83d..ef0517a 100644 --- a/src/config.ts +++ b/src/config.ts @@ -10,6 +10,8 @@ import { Cron, Result, SchemaParser } from "effect"; import type { GitIdentity } from "./lib/clone"; import type { WorkflowDefinition } from "./workflow"; +import type { AgentAdapter } from "./runtime/agent-adapter"; +import { opencodeAdapter } from "./runtime/opencode-adapter"; export const DEFAULT_WORKSPACE_ROOT = ".factory/workspaces"; export const DEFAULT_MAX_CONCURRENT_RUNS = 3; @@ -69,6 +71,11 @@ export interface FactoryConfig { */ readonly maxDispatchDepth: number; readonly maxChildrenPerRun: number; + /** + * Issue #36: the agent runtime. `defineConfig` defaults to the live opencode + * adapter when unset, so existing configs keep working unchanged. + */ + readonly agent: { readonly adapter: AgentAdapter }; } /** @@ -169,6 +176,11 @@ export interface FactoryConfigInput { readonly maxDispatchDepth?: number; readonly maxChildrenPerRun?: number; readonly schedules?: ReadonlyArray>; + /** + * Issue #36: the agent runtime. Optional — when omitted `defineConfig` + * uses the live opencode adapter. + */ + readonly agent?: { readonly adapter?: AgentAdapter }; } export function defineConfig(config: FactoryConfigInput): FactoryConfig { @@ -216,6 +228,7 @@ export function defineConfig(config: FactoryConfigInput): FactoryConfig { retainedWorkspaces, maxDispatchDepth, maxChildrenPerRun, + agent: { adapter: config.agent?.adapter ?? opencodeAdapter }, }; } diff --git a/src/replay/adapter.test.ts b/src/replay/adapter.test.ts index 460a2a5..722db85 100644 --- a/src/replay/adapter.test.ts +++ b/src/replay/adapter.test.ts @@ -1,9 +1,21 @@ -import { Effect } from "effect"; +import { Effect, ManagedRuntime } from "effect"; import { describe, expect, test } from "bun:test"; import { buildAgentStepEffect } from "../runtime/agent-step"; -import type { AgentAdapterYield } from "../runtime/agent-adapter"; +import type { AgentAdapter, AgentAdapterOptions, AgentStreamItem } from "../runtime/agent-adapter"; +import { AgentRuntimeLayer } from "../runtime/agent-runtime"; import { createCorpusReplayAdapter, createSlowFakeAdapter, loadCorpusBlocks } from "./adapter"; +function runtimeFor(adapter: AgentAdapter) { + return ManagedRuntime.make(AgentRuntimeLayer(adapter)); +} + +/** Drains an adapter stream to an array (for-of over async iterables is fine; this keeps types explicit). */ +async function drain(stream: AsyncIterable): Promise> { + const items: Array = []; + for await (const item of stream) items.push(item); + return items; +} + const FULL_ROUND_TRIP_CORPUS = `${import.meta.dir}/../../test/corpus/run-1789308170212.ndjson`; describe("loadCorpusBlocks", () => { @@ -49,14 +61,16 @@ describe("createCorpusReplayAdapter", () => { const adapter = createCorpusReplayAdapter(FULL_ROUND_TRIP_CORPUS); const chunks: Array = []; - const handle = buildAgentStepEffect({ - threadId: "t", - dir: "/tmp", - model: "opencode-go/deepseek-v4.1-flash", - prompt: "irrelevant, replay ignores it", - adapter, - onChunk: (chunk) => chunks.push(chunk), - }); + const runtime = runtimeFor(adapter); + const handle = await runtime.runPromise( + buildAgentStepEffect({ + threadId: "t", + dir: "/tmp", + model: "opencode-go/deepseek-v4.1-flash", + prompt: "irrelevant, replay ignores it", + onChunk: (chunk) => chunks.push(chunk), + }), + ); const outcome = await Effect.runPromise(handle.effect); @@ -67,27 +81,69 @@ describe("createCorpusReplayAdapter", () => { expect(outcome.runError).toBeUndefined(); }); - test("yields AgentAdapterYield items with signals extracted from opencode chunks", async () => { + test("supplies signals by interpreting recorded chunks, without imitating opencode shapes", async () => { const adapter = createCorpusReplayAdapter(FULL_ROUND_TRIP_CORPUS); - const stream = adapter.stream({ + const options = { threadId: "t", dir: "/tmp", model: "m", prompt: "p", abortController: new AbortController(), - }); + }; - const yields: AgentAdapterYield[] = []; - for await (const y of stream) { - yields.push(y as AgentAdapterYield); + const items = await drain(adapter.stream(options)); + const signals = items.map((item) => item.signal).filter((signal) => signal !== undefined); + + // The recorded session id surfaces as a signal, and every signal is the + // normalized union — never a vendor event name. + expect(signals).toContainEqual({ + kind: "session", + sessionId: "ses_f64ec04acffeJ0tjsHSkjAEqZF", + }); + for (const signal of signals) { + expect(["session", "structured-output", "error"]).toContain(signal.kind); } - expect(yields.length).toBe(39); - const sessionIdYield = yields.find((y) => y.signal?._tag === "sessionId"); - expect(sessionIdYield).toBeDefined(); - expect(sessionIdYield?.signal?.value).toBe("ses_f64ec04acffeJ0tjsHSkjAEqZF"); - const textYields = yields.filter((y) => y.signal === undefined); - expect(textYields.length).toBeGreaterThan(0); + // Chunks still ride through verbatim, one item per recorded chunk. + expect(items.length).toBe(39); + expect(items.every((item) => item.chunk !== undefined)).toBe(true); + }); + + test("a declarative signal-supplying adapter needs no opencode chunk shapes at all", async () => { + // The ADR's proof: a second adapter surfaces a session id and structured + // output through signals alone, with chunks that carry no vendor names. + const adapter = createFakeSignalAdapter({ + sessionId: "ses_fresh", + structuredOutput: { ok: true }, + }); + const items = await drain( + adapter.stream({ + threadId: "t", + dir: "/tmp", + model: "m", + prompt: "p", + abortController: new AbortController(), + }), + ); + + expect(items.map((item) => item.signal)).toEqual([ + { kind: "session", sessionId: "ses_fresh" }, + { kind: "structured-output", value: { ok: true } }, + ]); + + const runtime2 = runtimeFor(adapter); + const handle = await runtime2.runPromise( + buildAgentStepEffect({ + threadId: "t", + dir: "/tmp", + model: "m", + prompt: "p", + onChunk: () => {}, + }), + ); + const outcome = await Effect.runPromise(handle.effect); + expect(outcome.sessionId).toBe("ses_fresh"); + expect(outcome.structuredOutput).toEqual({ ok: true }); }); }); @@ -102,17 +158,67 @@ describe("createSlowFakeAdapter", () => { 1, ); - const handle = buildAgentStepEffect({ - threadId: "t", - dir: "/tmp", - model: "m", - prompt: "p", - adapter, - onChunk: () => {}, - }); + const runtime3 = runtimeFor(adapter); + const handle = await runtime3.runPromise( + buildAgentStepEffect({ + threadId: "t", + dir: "/tmp", + model: "m", + prompt: "p", + onChunk: () => {}, + }), + ); const outcome = await Effect.runPromise(handle.effect); expect(outcome.chunkCount).toBe(3); expect(outcome.finalText).toBe("hi"); }); + + test("attaches declared signals to their chunks", async () => { + const adapter = createSlowFakeAdapter( + [{ type: "RUN_STARTED" }, { type: "TEXT_MESSAGE_START" }], + 1, + [{ index: 0, signal: { kind: "session", sessionId: "ses_slow" } }], + ); + + const runtime4 = runtimeFor(adapter); + const handle = await runtime4.runPromise( + buildAgentStepEffect({ + threadId: "t", + dir: "/tmp", + model: "m", + prompt: "p", + onChunk: () => {}, + }), + ); + + const outcome = await Effect.runPromise(handle.effect); + expect(outcome.chunkCount).toBe(2); + expect(outcome.sessionId).toBe("ses_slow"); + }); }); + +/** Minimal stand-in for a second adapter: signals only, no vendor chunks. */ +function createFakeSignalAdapter(signals: { + sessionId?: string; + structuredOutput?: unknown; +}): AgentAdapter { + return { + stream(_options: AgentAdapterOptions): AsyncIterable { + return (async function* () { + if (signals.sessionId !== undefined) { + yield { + chunk: { type: "RUN_STARTED" }, + signal: { kind: "session", sessionId: signals.sessionId }, + }; + } + if (signals.structuredOutput !== undefined) { + yield { + chunk: { type: "RUN_FINISHED" }, + signal: { kind: "structured-output", value: signals.structuredOutput }, + }; + } + })(); + }, + }; +} diff --git a/src/runtime/agent-runtime.test.ts b/src/runtime/agent-runtime.test.ts new file mode 100644 index 0000000..4c80596 --- /dev/null +++ b/src/runtime/agent-runtime.test.ts @@ -0,0 +1,38 @@ +/** + * The agent runtime as a context service (ADR 0012 §4). + */ + +import { describe, expect, test } from "bun:test"; +import { Effect, Layer, ManagedRuntime } from "effect"; +import { createSlowFakeAdapter } from "../replay/adapter"; +import { opencodeAdapter } from "./opencode-adapter"; +import { AgentRuntime, AgentRuntimeLayer } from "./agent-runtime"; + +describe("AgentRuntime service", () => { + test("the default layer provides the live opencode adapter", async () => { + const runtime = ManagedRuntime.make(Layer.succeed(AgentRuntime, { adapter: opencodeAdapter })); + const { adapter } = await runtime.runPromise(AgentRuntime); + expect(adapter).toBe(opencodeAdapter); + await runtime.dispose(); + }); + + test("AgentRuntimeLayer swaps the adapter via context", async () => { + const fake = createSlowFakeAdapter([], 1); + const runtime = ManagedRuntime.make(AgentRuntimeLayer(fake)); + const { adapter } = await runtime.runPromise(AgentRuntime); + expect(adapter).toBe(fake); + await runtime.dispose(); + }); + + test("an effect that requires the service resolves it from context", async () => { + const fake = createSlowFakeAdapter([], 1); + const program = Effect.gen(function* () { + const runtime = yield* AgentRuntime; + return runtime.adapter; + }); + const runtime = ManagedRuntime.make(AgentRuntimeLayer(fake)); + const adapter = await runtime.runPromise(program); + expect(adapter).toBe(fake); + await runtime.dispose(); + }); +}); diff --git a/src/runtime/agent-runtime.ts b/src/runtime/agent-runtime.ts new file mode 100644 index 0000000..1bcbfc6 --- /dev/null +++ b/src/runtime/agent-runtime.ts @@ -0,0 +1,48 @@ +/** + * The agent runtime as an Effect context service (ADR 0012 §4). + * + * The runtime is the daemon's first service to move into Effect's dependency + * injection: it stops being a value threaded through every interface between + * the CLI and the step runner, and becomes a service resolved from context at + * the point of use (`runtime/agent-step.ts`). + * + * The service shape is deliberately small: it currently exposes only the + * adapter, which is the part issue #36 needs to make selectable from config. + * Issue #37 will add workspace preparation behind the same service. + */ + +import { Context, Effect, Layer, ManagedRuntime } from "effect"; +import type { AgentAdapter } from "./agent-adapter"; +import { opencodeAdapter } from "./opencode-adapter"; + +export interface AgentRuntimeShape { + /** The adapter that produces agent chunk streams (ADR 0012 §2). */ + readonly adapter: AgentAdapter; +} + +/** + * The agent-runtime service. Yield it inside an Effect to read the adapter. + * The default implementation uses the live opencode adapter; tests and the + * daemon override it with `AgentRuntimeLayer`. + */ +export class AgentRuntime extends Context.Service()( + "AgentRuntime", + { make: Effect.succeed({ adapter: opencodeAdapter }) }, +) {} + +/** + * A layer that provides a fixed adapter. Tests use this to swap in the + * corpus-replay runtime without threading an option through the call chain. + */ +export const AgentRuntimeLayer = (adapter: AgentAdapter): Layer.Layer => + Layer.succeed(AgentRuntime, { adapter }); + +/** + * Build a `ManagedRuntime` from an adapter. Convenient for tests and for the + * CLI's direct-run path, which both need to run effects that require + * `AgentRuntime`. + */ +export const makeAgentRuntime = ( + adapter: AgentAdapter, +): ManagedRuntime.ManagedRuntime => + ManagedRuntime.make(AgentRuntimeLayer(adapter)); diff --git a/src/runtime/agent-step.test.ts b/src/runtime/agent-step.test.ts index 17069d1..ee01f1e 100644 --- a/src/runtime/agent-step.test.ts +++ b/src/runtime/agent-step.test.ts @@ -1,138 +1,127 @@ -import { Effect } from "effect"; +/** + * Pins the runtime side of the adapter seam (ADR 0012 §2): normalized signals + * arrive from the adapter, the runtime records them, and it never interprets + * a vendor chunk name itself. Chunks stay opaque and ride through verbatim. + */ + +import { Effect, ManagedRuntime } from "effect"; import { describe, expect, test } from "bun:test"; -import type { AgentAdapterYield, AgentSignal } from "./agent-adapter"; +import type { AgentAdapter, AgentAdapterOptions, AgentStreamItem } from "./agent-adapter"; +import { AgentRuntimeLayer } from "./agent-runtime"; import { buildAgentStepEffect } from "./agent-step"; -function makeYield(chunk: unknown, signal?: AgentSignal): AgentAdapterYield { - return signal !== undefined ? { chunk, signal } : { chunk }; -} - -function signalAdapter(yields: ReadonlyArray) { +function scriptedAdapter(items: ReadonlyArray): AgentAdapter { return { - async prepareWorkspace(_dir: string): Promise {}, - stream() { - return { - async *[Symbol.asyncIterator]() { - for (const y of yields) yield y; - }, - }; + stream(_options: AgentAdapterOptions): AsyncIterable { + return (async function* () { + yield* items; + })(); }, }; } -describe("buildAgentStepEffect signal extraction", () => { - test("extracts sessionId from a sessionId signal", async () => { - const adapter = signalAdapter([ - makeYield({ type: "TEXT_MESSAGE_START" }), - makeYield({ type: "TEXT_MESSAGE_CONTENT", delta: "hello" }), - makeYield({ type: "TEXT_MESSAGE_END" }), - makeYield( - { type: "CUSTOM", name: "opencode.session-id", value: { sessionId: "ses_abc" } }, - { _tag: "sessionId", value: "ses_abc" }, - ), - ]); +const TEXT_CHUNKS = [ + { type: "TEXT_MESSAGE_START", messageId: "m1" }, + { type: "TEXT_MESSAGE_CONTENT", delta: "hello" }, + { type: "TEXT_MESSAGE_END", messageId: "m1" }, +]; + +async function runStep(adapter: AgentAdapter, options: Parameters[0]) { + const runtime = ManagedRuntime.make(AgentRuntimeLayer(adapter)); + const handle = await runtime.runPromise(buildAgentStepEffect(options)); + const outcome = await Effect.runPromise(handle.effect); + await runtime.dispose(); + return outcome; +} - const handle = buildAgentStepEffect({ +describe("buildAgentStepEffect over the signal seam (ADR 0012 §2)", () => { + test("records adapter signals; chunks reach onChunk verbatim", async () => { + const seen: Array = []; + const items: ReadonlyArray = [ + { chunk: { type: "RUN_STARTED" } }, + { + chunk: { type: "CUSTOM", name: "vendor.session", value: { id: "ses_x" } }, + signal: { kind: "session", sessionId: "ses_x" }, + }, + ...TEXT_CHUNKS.map((chunk) => ({ chunk })), + { + chunk: { type: "CUSTOM", name: "vendor.output", value: { object: { answer: 42 } } }, + signal: { kind: "structured-output", value: { answer: 42 } }, + }, + { chunk: { type: "RUN_FINISHED" } }, + ]; + + const outcome = await runStep(scriptedAdapter(items), { threadId: "t", dir: "/tmp", model: "m", prompt: "p", - adapter, - onChunk: () => {}, + onChunk: (chunk) => seen.push(chunk), }); - const outcome = await Effect.runPromise(handle.effect); - expect(outcome.sessionId).toBe("ses_abc"); + expect(outcome.chunkCount).toBe(7); + expect(outcome.sessionId).toBe("ses_x"); + expect(outcome.structuredOutput).toEqual({ answer: 42 }); expect(outcome.finalText).toBe("hello"); - expect(outcome.chunkCount).toBe(4); - }); - - test("extracts structuredOutput from a structuredOutput signal", async () => { - const outputObject = { title: "test", body: "content" }; - const adapter = signalAdapter([ - makeYield({ type: "TEXT_MESSAGE_START" }), - makeYield({ type: "TEXT_MESSAGE_CONTENT", delta: "done" }), - makeYield({ type: "TEXT_MESSAGE_END" }), - makeYield( - { type: "CUSTOM", name: "structured-output.complete", value: { object: outputObject } }, - { _tag: "structuredOutput", value: outputObject }, - ), - ]); - - const handle = buildAgentStepEffect({ - threadId: "t", - dir: "/tmp", - model: "m", - prompt: "p", - adapter, - onChunk: () => {}, - }); - - const outcome = await Effect.runPromise(handle.effect); - expect(outcome.structuredOutput).toEqual(outputObject); + expect(outcome.runError).toBeUndefined(); + expect(seen).toEqual(items.map((item) => item.chunk)); }); - test("extracts runError from a runError signal", async () => { - const adapter = signalAdapter([ - makeYield( - { type: "RUN_ERROR", message: "something broke" }, - { _tag: "runError", value: "something broke" }, - ), - ]); - - const handle = buildAgentStepEffect({ - threadId: "t", - dir: "/tmp", - model: "m", - prompt: "p", - adapter, - onChunk: () => {}, - }); - - const outcome = await Effect.runPromise(handle.effect); - expect(outcome.runError).toBe("something broke"); + test("an error signal surfaces as runError", async () => { + const outcome = await runStep( + scriptedAdapter([ + { + chunk: { type: "RUN_ERROR", message: "boom" }, + signal: { kind: "error", message: "boom" }, + }, + ]), + { + threadId: "t", + dir: "/tmp", + model: "m", + prompt: "p", + onChunk: () => {}, + }, + ); + + expect(outcome.chunkCount).toBe(1); + expect(outcome.runError).toBe("boom"); }); - test("onChunk receives the raw opaque chunk, not the signal wrapper", async () => { - const rawChunk = { type: "TEXT_MESSAGE_START" }; - const adapter = signalAdapter([makeYield(rawChunk)]); - - const chunks: Array = []; - const handle = buildAgentStepEffect({ - threadId: "t", - dir: "/tmp", - model: "m", - prompt: "p", - adapter, - onChunk: (chunk) => chunks.push(chunk), - }); + test("the runtime records signals, it does not interpret vendor chunk names", async () => { + // A raw CUSTOM chunk with no signal attached is opaque bookkeeping now. + const outcome = await runStep( + scriptedAdapter([ + { chunk: { type: "CUSTOM", name: "vendor.session", value: { sessionId: "ses_raw" } } }, + ]), + { + threadId: "t", + dir: "/tmp", + model: "m", + prompt: "p", + onChunk: () => {}, + }, + ); - await Effect.runPromise(handle.effect); - expect(chunks).toHaveLength(1); - expect(chunks[0]).toBe(rawChunk); + expect(outcome.sessionId).toBeUndefined(); + expect(outcome.structuredOutput).toBeUndefined(); }); - test("yields without signals pass through without error", async () => { - const adapter = signalAdapter([ - makeYield({ type: "TEXT_MESSAGE_START" }), - makeYield({ type: "TEXT_MESSAGE_CONTENT", delta: "no signals here" }), - makeYield({ type: "TEXT_MESSAGE_END" }), - ]); - - const handle = buildAgentStepEffect({ - threadId: "t", - dir: "/tmp", - model: "m", - prompt: "p", - adapter, - onChunk: () => {}, + test("the service is resolved from context, not passed as an option", async () => { + const fake = scriptedAdapter([{ chunk: { type: "RUN_FINISHED" } }]); + const program = Effect.gen(function* () { + const handle = yield* buildAgentStepEffect({ + threadId: "t", + dir: "/tmp", + model: "m", + prompt: "p", + onChunk: () => {}, + }); + return yield* handle.effect; }); - - const outcome = await Effect.runPromise(handle.effect); - expect(outcome.chunkCount).toBe(3); - expect(outcome.finalText).toBe("no signals here"); - expect(outcome.sessionId).toBeUndefined(); - expect(outcome.structuredOutput).toBeUndefined(); - expect(outcome.runError).toBeUndefined(); + const runtime = ManagedRuntime.make(AgentRuntimeLayer(fake)); + const outcome = await runtime.runPromise(program); + expect(outcome.chunkCount).toBe(1); + await runtime.dispose(); }); }); diff --git a/src/runtime/run-signals.test.ts b/src/runtime/run-signals.test.ts index 384dfab..53c4fcf 100644 --- a/src/runtime/run-signals.test.ts +++ b/src/runtime/run-signals.test.ts @@ -11,6 +11,7 @@ import { describe, expect, test } from "bun:test"; import type { RunEvent } from "../events"; import { defineWorkflow, Schema } from "../workflow"; import { createSlowFakeAdapter } from "../replay/adapter"; +import { makeAgentRuntime } from "./agent-runtime"; import { startRun } from "./run"; const TEXT_CHUNKS = [ @@ -33,13 +34,16 @@ function agentWorkflow(outputSchema?: Schema.Codec) { }); } -async function runWith(chunks: ReadonlyArray, signals: Parameters[2]) { +async function runWith( + chunks: ReadonlyArray, + signals: Parameters[2], +) { const events: Array = []; - const handle = startRun(agentWorkflow(), { + const runtime = makeAgentRuntime(createSlowFakeAdapter(chunks, 1, signals)); + const handle = await startRun(agentWorkflow(), runtime, { runId: "run-signals", dir: "/tmp", input: {}, - adapter: createSlowFakeAdapter(chunks, 1, signals), onEvent: (event) => events.push(event), }); const outcome = await handle.result; @@ -64,18 +68,15 @@ describe("signals populate AgentStepFinished as before (issue #35)", () => { test("a structured-output signal is decoded into AgentStepFinished.output (tier 1)", async () => { const events: Array = []; - const handle = startRun(agentWorkflow(SIGNAL_SCHEMA), { + const runtime = makeAgentRuntime( + createSlowFakeAdapter([...TEXT_CHUNKS, { type: "CUSTOM", name: "anything-at-all" }], 1, [ + { index: 3, signal: { kind: "structured-output", value: { where: "from-signal" } } }, + ]), + ); + const handle = await startRun(agentWorkflow(SIGNAL_SCHEMA), runtime, { runId: "run-signals-output", dir: "/tmp", input: {}, - adapter: createSlowFakeAdapter( - [ - ...TEXT_CHUNKS, - { type: "CUSTOM", name: "anything-at-all" }, - ], - 1, - [{ index: 3, signal: { kind: "structured-output", value: { where: "from-signal" } } }], - ), onEvent: (event) => events.push(event), }); const outcome = await handle.result; @@ -86,11 +87,11 @@ describe("signals populate AgentStepFinished as before (issue #35)", () => { test("without a signal, tier 2 re-parses the final text — the fallback is unchanged", async () => { const events: Array = []; - const handle = startRun(agentWorkflow(SIGNAL_SCHEMA), { + const runtime = makeAgentRuntime(createSlowFakeAdapter(TEXT_CHUNKS, 1)); + const handle = await startRun(agentWorkflow(SIGNAL_SCHEMA), runtime, { runId: "run-signals-tier2", dir: "/tmp", input: {}, - adapter: createSlowFakeAdapter(TEXT_CHUNKS, 1), onEvent: (event) => events.push(event), }); const outcome = await handle.result; @@ -100,9 +101,10 @@ describe("signals populate AgentStepFinished as before (issue #35)", () => { }); test("an error signal fails the step and lands on AgentStepFinished.error", async () => { - const { outcome, events } = await runWith([...TEXT_CHUNKS, { type: "RUN_FINISHED" }], [ - { index: 3, signal: { kind: "error", message: "sandbox vanished" } }, - ]); + const { outcome, events } = await runWith( + [...TEXT_CHUNKS, { type: "RUN_FINISHED" }], + [{ index: 3, signal: { kind: "error", message: "sandbox vanished" } }], + ); expect(outcome.outcome).toBe("failed"); const finished = stepFinished(events); diff --git a/src/runtime/run.test.ts b/src/runtime/run.test.ts index 637d2ba..95de407 100644 --- a/src/runtime/run.test.ts +++ b/src/runtime/run.test.ts @@ -10,26 +10,10 @@ import { describe, expect, test } from "bun:test"; import type { RunEvent } from "../events"; import { defineWorkflow, Schema } from "../workflow"; import { createSlowFakeAdapter } from "../replay/adapter"; -import { startRun, RunCancelledSignal } from "./run"; +import { makeAgentRuntime } from "./agent-runtime"; +import { startRun } from "./run"; import type { AgentAdapter } from "./agent-adapter"; -describe("RunCancelledSignal as TaggedError (#34)", () => { - test("carries the _tag", () => { - const signal = new RunCancelledSignal({}); - expect(signal._tag).toBe("RunCancelledSignal"); - expect(signal instanceof Error).toBe(true); - }); - - test("is matchable by _tag from an unknown catch", () => { - try { - throw new RunCancelledSignal({}); - } catch (err: unknown) { - const e = err as { _tag?: string }; - expect(e._tag).toBe("RunCancelledSignal"); - } - }); -}); - const SLOW_CHUNKS = [ { type: "TEXT_MESSAGE_START" }, { type: "TEXT_MESSAGE_CONTENT", delta: "hi" }, @@ -48,11 +32,11 @@ describe("startRun cancellation", () => { }, }); - const handle = startRun(workflow, { + const runtime = makeAgentRuntime(createSlowFakeAdapter(SLOW_CHUNKS, 20)); + const handle = await startRun(workflow, runtime, { runId: "run-cancel-test", dir: "/tmp", input: {}, - adapter: createSlowFakeAdapter(SLOW_CHUNKS, 20), onEvent: (event) => events.push(event), }); @@ -83,11 +67,11 @@ describe("startRun cancellation", () => { }, }); - const handle = startRun(workflow, { + const runtime = makeAgentRuntime(createSlowFakeAdapter(SLOW_CHUNKS, 1)); + const handle = await startRun(workflow, runtime, { runId: "run-no-cancel-test", dir: "/tmp", input: {}, - adapter: createSlowFakeAdapter(SLOW_CHUNKS, 1), onEvent: () => {}, }); @@ -114,11 +98,11 @@ describe("startRun workspace kind (issue #13)", () => { }, }); - const handle = startRun(workflow, { + const runtime = makeAgentRuntime(createSlowFakeAdapter([])); + const handle = await startRun(workflow, runtime, { runId: "run-scratch-wb", dir: "/tmp/nothing", input: {}, - adapter: createSlowFakeAdapter([]), workspaceKind: "scratch", repo: { slug: "owner/repo", baseBranch: "main" }, onEvent: (event) => events.push(event), @@ -149,11 +133,11 @@ describe("startRun workspace kind (issue #13)", () => { }), }); - const handle = startRun(workflow, { + const runtime = makeAgentRuntime(createSlowFakeAdapter([])); + const handle = await startRun(workflow, runtime, { runId: "run-clone-wb", dir: "/tmp/nothing", input: {}, - adapter: createSlowFakeAdapter([]), repo: { slug: "owner/repo", baseBranch: "main" }, onEvent: (event) => events.push(event), }); @@ -175,22 +159,26 @@ describe("startRun workspace kind (issue #13)", () => { }); const scratchEvents: Array = []; - await startRun(workflow, { - runId: "run-kind-scratch", - dir: "/tmp/s", - input: {}, - adapter: createSlowFakeAdapter([]), - workspaceKind: "scratch", - onEvent: (event) => scratchEvents.push(event), - }).result; + const scratchRuntime = makeAgentRuntime(createSlowFakeAdapter([])); + await ( + await startRun(workflow, scratchRuntime, { + runId: "run-kind-scratch", + dir: "/tmp/s", + input: {}, + workspaceKind: "scratch", + onEvent: (event) => scratchEvents.push(event), + }) + ).result; const cloneEvents: Array = []; - await startRun(workflow, { - runId: "run-kind-clone", - dir: "/tmp/c", - input: {}, - adapter: createSlowFakeAdapter([]), - onEvent: (event) => cloneEvents.push(event), - }).result; + const cloneRuntime = makeAgentRuntime(createSlowFakeAdapter([])); + await ( + await startRun(workflow, cloneRuntime, { + runId: "run-kind-clone", + dir: "/tmp/c", + input: {}, + onEvent: (event) => cloneEvents.push(event), + }) + ).result; expect( scratchEvents[0]?.payload._tag === "RunStarted" && scratchEvents[0].payload.workspaceKind, @@ -211,11 +199,11 @@ describe("startRun schedule trigger (issue #16)", () => { return { finalText: result.finalText }; }, }); - const handle = startRun(workflow, { + const runtime = makeAgentRuntime(createSlowFakeAdapter(SLOW_CHUNKS, 5)); + const handle = await startRun(workflow, runtime, { runId: "run-by-schedule", dir: "/tmp", input: { issueNumber: 7 }, - adapter: createSlowFakeAdapter(SLOW_CHUNKS, 5), scheduleId: "nightly", onEvent: (event) => events.push(event), }); @@ -230,11 +218,11 @@ describe("startRun schedule trigger (issue #16)", () => { input: Schema.Struct({}), run: async () => ({}), }); - const handle = startRun(workflow, { + const runtime = makeAgentRuntime(createSlowFakeAdapter([])); + const handle = await startRun(workflow, runtime, { runId: "run-manual", dir: "/tmp", input: {}, - adapter: createSlowFakeAdapter([]), onEvent: (event) => events.push(event), }); await handle.result; @@ -257,11 +245,11 @@ describe("startRun model precedence (issue #16)", () => { return {}; }, }); - const handle = startRun(workflow, { + const runtime = makeAgentRuntime(createSlowFakeAdapter(SLOW_CHUNKS, 5)); + const handle = await startRun(workflow, runtime, { runId: "run-models", dir: "/tmp", input: {}, - adapter: createSlowFakeAdapter(SLOW_CHUNKS, 5), agentOverrides: { model: "schedule-override" }, onEvent: (event) => events.push(event), }); @@ -285,11 +273,11 @@ describe("startRun model precedence (issue #16)", () => { return {}; }, }); - const handle = startRun(workflow, { + const runtime = makeAgentRuntime(createSlowFakeAdapter(SLOW_CHUNKS, 5)); + const handle = await startRun(workflow, runtime, { runId: "run-models-fallback", dir: "/tmp", input: {}, - adapter: createSlowFakeAdapter(SLOW_CHUNKS, 5), onEvent: (event) => events.push(event), }); const outcome = await handle.result; @@ -322,14 +310,14 @@ describe("startRun prepareWorkspace (ADR 0012 §3, #37)", () => { run: async () => ({}), }); - await startRun(workflow, { + const runtime = makeAgentRuntime(adapter); + await startRun(workflow, runtime, { runId: "run-prep-clone", dir: "/tmp/clone-dir", input: {}, - adapter, prepareWorkspace: true, onEvent: (event) => events.push(event), - }).result; + }).then((h) => h.result); expect(adapter.prepared).toEqual(["/tmp/clone-dir"]); }); @@ -341,13 +329,13 @@ describe("startRun prepareWorkspace (ADR 0012 §3, #37)", () => { run: async () => ({}), }); - await startRun(workflow, { + const runtime = makeAgentRuntime(adapter); + await startRun(workflow, runtime, { runId: "run-prep-scratch", dir: "/tmp/scratch-dir", input: {}, - adapter, onEvent: () => {}, - }).result; + }).then((h) => h.result); expect(adapter.prepared).toEqual([]); }); diff --git a/src/runtime/run.ts b/src/runtime/run.ts index abf9554..03eb1cc 100644 --- a/src/runtime/run.ts +++ b/src/runtime/run.ts @@ -1,19 +1,10 @@ /** * Run orchestration: decode input, build `ctx`, run the workflow's plain - * `async` `run()`, emit `RunEvent`s along the way (ADR 0002/0003). This is - * the runtime side of the ownership split — the workflow never touches - * `RunEvent`, `seq`, or the tree. - * - * Cancellation (D9's "never throws" does not apply to cancellation, which is - * the run's own unwind signal, not a workflow-observable failure): `cancel()` - * aborts a host-level `AbortController`; if an agent step is in flight, its - * `Fiber` is interrupted, `AgentStepFinished{outcome:"cancelled"}` is emitted - * from the step's live-mutated `partial`, and a `RunCancelledSignal` is - * thrown to unwind the workflow's `await` chain. It is caught here, nowhere - * else, and turns into `RunCancelled` rather than `RunFailed`. + * `async` `run()`, emit `RunEvent`s along the way (ADR 0002/0003). */ import { Cause, Effect, Exit, Fiber, Schema, SchemaParser } from "effect"; +import type { ManagedRuntime } from "effect"; import type { RunEvent, RunEventPayload } from "../events"; import { hostExec, type ExecResult } from "../lib/exec"; import { writeBack as writeBackLib, type WriteBackResult } from "../lib/writeback"; @@ -29,8 +20,13 @@ import type { WriteBackCallOptions, } from "../workflow"; import { DedupeKeyError } from "../lib/dedupe"; -import type { AgentAdapter } from "./agent-adapter"; import { buildAgentStepEffect } from "./agent-step"; +import { AgentRuntime } from "./agent-runtime"; + +export class RunCancelledSignal extends Schema.TaggedError()( + "RunCancelledSignal", + {}, +) {} /** * The domain errors (`DedupeKeyError`, `ConcurrencyLimitError`, @@ -43,19 +39,9 @@ function domainErrorMessage(err: unknown): string { return err.message || String(err); } -export class RunCancelledSignal extends Schema.TaggedError()( - "RunCancelledSignal", - {}, -) {} - /** Matches the spike's model (STATUS.md); overridden per-call or per-workflow. */ export const DEFAULT_MODEL = "opencode-go/deepseek-v4.1-flash"; -/** - * The deployment-known half of write-back (D27/D32): the repo slug pushed to - * `gh pr create` and the branch PRs open against. The runtime supplies it - * from `factory.config.ts`; WorkflowCtx.writeBack callers never carry it. - */ export interface RunRepo { readonly slug: string; readonly baseBranch: string; @@ -65,14 +51,7 @@ export interface StartRunOptions { readonly runId: string; readonly dir: string; readonly input: unknown; - readonly adapter: AgentAdapter; - /** Write-back environment (D32). Absent, a workflow's `ctx.writeBack` fails. */ readonly repo?: RunRepo; - /** - * How `dir` was provisioned (issue #13). `"clone"` is the default; a - * scratch run differs only in `RunStarted.workspaceKind` and `ctx.writeBack`'s - * failure message. - */ readonly workspaceKind?: WorkspaceKind; /** * ADR 0012 §3 (#37): when true, the runtime calls `adapter.prepareWorkspace` @@ -81,33 +60,10 @@ export interface StartRunOptions { * no preparation happens (explicit caller-managed dirs, scratch). */ readonly prepareWorkspace?: boolean; - /** - * The context service behind `ctx.dispatch` (issue #14). Absent, the ctx - * member is still present but throws — in-process execution is legacy and - * cannot start nested runs. - */ readonly dispatch?: DispatchChildFn; - /** - * This run's parent, when it was started by `ctx.dispatch` (issue #14). - * Recorded on `RunStarted.parentId`. - */ readonly parentRunId?: string; - /** - * This run's dedupe key (issue #15), when it was started with one. Recorded - * on `RunStarted.dedupeKey` for observability; the claim itself is the - * server's (server/runs.ts), not the runtime's. - */ readonly dedupeKey?: string; - /** - * Issue #16: the schedule that started this run, when any did. Recorded on - * `RunStarted.scheduleId` so a run explains its trigger. - */ readonly scheduleId?: string; - /** - * Issue #16: agent-level overrides the starting schedule carries (its - * `agent.model`). Sits between a per-call option and the workflow's own - * default in the precedence chain. - */ readonly agentOverrides?: { readonly model?: string }; readonly onEvent: (event: RunEvent) => void; } @@ -132,17 +88,14 @@ function resolveOutput( schema: Schema.Codec | undefined, ): O | undefined { if (schema === undefined) return undefined; - const decode = SchemaParser.decodeUnknownSync(schema); - if (outcome.structuredOutput !== undefined) { try { return decode(outcome.structuredOutput) as O; } catch { - // fall through to tier 2 (ADR 0002 §3: manual re-parse of finalText) + // fall through to tier 2 } } - try { return decode(JSON.parse(outcome.finalText)) as O; } catch { @@ -150,10 +103,11 @@ function resolveOutput( } } -export function startRun( +export async function startRun( workflow: WorkflowDefinition, + runtime: ManagedRuntime.ManagedRuntime, options: StartRunOptions, -): RunHandle { +): Promise> { let seq = 0; const emit = (payload: RunEventPayload): void => { options.onEvent({ runId: options.runId, seq: seq++, ts: Date.now(), payload }); @@ -173,19 +127,15 @@ export function startRun( const nextStepId = makeIdCounter("step"); const nextExecId = makeIdCounter("exec"); const startedAt = Date.now(); - const workspaceKind = options.workspaceKind ?? "clone"; const execImpl = async (argv: ReadonlyArray): Promise => { - if (cancelled) throw new RunCancelledSignal({}); - + if (cancelled) throw new RunCancelledSignal(); const execId = nextExecId(); emit({ _tag: "ExecStarted", execId, command: [...argv], cwd: options.dir }); - const stepStartedAt = Date.now(); const result = await hostExec(argv, { cwd: options.dir, signal: runController.signal }); const durationMs = Date.now() - stepStartedAt; - emit({ _tag: "ExecFinished", execId, @@ -195,8 +145,7 @@ export function startRun( stderr: result.stderr, durationMs, }); - - if (cancelled) throw new RunCancelledSignal({}); + if (cancelled) throw new RunCancelledSignal(); return result; }; @@ -205,16 +154,12 @@ export function startRun( prompt: string, opts?: AgentCallOptions, ): Promise> => { - if (cancelled) throw new RunCancelledSignal({}); - + if (cancelled) throw new RunCancelledSignal(); const stepId = nextStepId(); - // Precedence (issue #16): per-call option > the starting schedule's - // override > the workflow's `agent` default > the runtime fallback. const model = opts?.model ?? options.agentOverrides?.model ?? workflow.agent?.model ?? DEFAULT_MODEL; const outputSchema = opts?.output !== undefined ? Schema.toJsonSchemaDocument(opts.output) : undefined; - emit({ _tag: "AgentStepStarted", stepId, @@ -223,32 +168,30 @@ export function startRun( prompt, structured: opts?.output !== undefined, }); - const stepStartedAt = Date.now(); - const handle = buildAgentStepEffect({ - threadId: options.runId, - dir: options.dir, - model, - prompt, - outputSchema, - adapter: options.adapter, - onChunk: (chunk) => { - const record = chunk as { type?: unknown }; - emit({ - _tag: "AgentChunk", - stepId, - chunkType: typeof record.type === "string" ? record.type : "UNKNOWN", - chunk: chunk as never, - }); - }, - }); - + const handle = await runtime.runPromise( + buildAgentStepEffect({ + threadId: options.runId, + dir: options.dir, + model, + prompt, + outputSchema, + onChunk: (chunk) => { + const record = chunk as { type?: unknown }; + emit({ + _tag: "AgentChunk", + stepId, + chunkType: typeof record.type === "string" ? record.type : "UNKNOWN", + chunk: chunk as never, + }); + }, + }), + ); const fiber = Effect.runFork(handle.effect); activeAgentFiber = fiber; const exit = await Effect.runPromise(Fiber.await(fiber)); activeAgentFiber = null; const durationMs = Date.now() - stepStartedAt; - if (Exit.isFailure(exit)) { const cause = exit.cause; const wasInterrupted = Exit.hasInterrupts(exit); @@ -264,11 +207,9 @@ export function startRun( ...(handle.partial.sessionId !== undefined ? { sessionId: handle.partial.sessionId } : {}), - ...(handle.partial.usage !== undefined ? { usage: handle.partial.usage } : {}), }); - throw new RunCancelledSignal({}); + throw new RunCancelledSignal(); } - const message = Cause.pretty(cause); emit({ _tag: "AgentStepFinished", @@ -279,16 +220,13 @@ export function startRun( durationMs, finalText: handle.partial.finalText, ...(handle.partial.sessionId !== undefined ? { sessionId: handle.partial.sessionId } : {}), - ...(handle.partial.usage !== undefined ? { usage: handle.partial.usage } : {}), error: message, }); throw new Error(`agent step "${name}" failed: ${message}`); } - const result = exit.value; const output = resolveOutput(result, opts?.output); const outcome = result.runError !== undefined ? "failed" : "completed"; - emit({ _tag: "AgentStepFinished", stepId, @@ -299,14 +237,11 @@ export function startRun( finalText: result.finalText, ...(output !== undefined ? { output: output as never } : {}), ...(result.sessionId !== undefined ? { sessionId: result.sessionId } : {}), - ...(result.usage !== undefined ? { usage: result.usage } : {}), ...(result.runError !== undefined ? { error: result.runError } : {}), }); - if (outcome === "failed") { throw new Error(`agent step "${name}" reported RUN_ERROR: ${result.runError}`); } - return { stepId, chunkCount: result.chunkCount, @@ -320,14 +255,12 @@ export function startRun( const assertImpl = async (name: string, callback: AssertCallback): Promise => { const raw = await callback(); const normalized = typeof raw === "boolean" ? { pass: raw, details: undefined } : raw; - emit({ _tag: "AssertionRecorded", name, pass: normalized.pass, ...(normalized.details !== undefined ? { details: normalized.details as never } : {}), }); - return { name, pass: normalized.pass, details: normalized.details }; }; @@ -337,9 +270,6 @@ export function startRun( const writeBackImpl = async (opts: WriteBackCallOptions): Promise => { emit({ _tag: "WriteBackStarted", branch: opts.branch }); - - // Scratch has nothing to push (issue #13): the failure names the workspace - // kind, is recorded like any other write-back failure, and runs no git. if (workspaceKind === "scratch") { const message = 'ctx.writeBack is not available on a scratch workspace: the workspace kind is "scratch", so there is no clone to push'; @@ -353,7 +283,6 @@ export function startRun( }); throw new Error(message); } - if (options.repo === undefined) { const message = "ctx.writeBack needs run-repo config (slug/baseBranch) — this run was started without it"; @@ -367,7 +296,6 @@ export function startRun( }); throw new Error(message); } - try { const result = await writeBackLib( { @@ -382,11 +310,9 @@ export function startRun( }, execImpl, ); - const outcome = result.prResult.exitCode === 0 ? "completed" : "failed"; const failureDetail = result.prResult.stderr || result.pushResult.stderr || result.commitResult.stderr; - emit({ _tag: "WriteBackFinished", branch: opts.branch, @@ -397,7 +323,6 @@ export function startRun( ...(result.prUrl !== null ? { prUrl: result.prUrl } : {}), ...(outcome === "failed" ? { error: failureDetail } : {}), }); - return result; } catch (err) { if (err instanceof RunCancelledSignal) throw err; @@ -426,7 +351,6 @@ export function startRun( "(in-process execution is legacy) has no registry to start a nested run from", ); } - let childRunId: string; try { childRunId = await options.dispatch(child, input, opts); @@ -441,7 +365,6 @@ export function startRun( } throw err; } - emit({ _tag: "RunDispatched", childRunId, @@ -449,7 +372,6 @@ export function startRun( input: input as never, ...(opts?.dedupeKey !== undefined ? { dedupeKey: opts.dedupeKey } : {}), }); - return childRunId; }; @@ -468,11 +390,10 @@ export function startRun( try { decodedInput = SchemaParser.decodeUnknownSync(workflow.input)(options.input); } catch (err) { - const message = domainErrorMessage(err); + const message = err instanceof Error ? err.message : String(err); emit({ _tag: "RunFailed", message, durationMs: Date.now() - startedAt }); return { outcome: "failed", error: message }; } - emit({ _tag: "RunStarted", workflowId: workflow.id, @@ -483,18 +404,16 @@ export function startRun( ...(options.dedupeKey !== undefined ? { dedupeKey: options.dedupeKey } : {}), ...(options.scheduleId !== undefined ? { scheduleId: options.scheduleId } : {}), }); - try { if (options.prepareWorkspace === true) { - await options.adapter.prepareWorkspace(options.dir); + const { adapter } = await runtime.runPromise(AgentRuntime); + await adapter.prepareWorkspace(options.dir); } const output = await workflow.run(ctx, decodedInput); - if (workflow.output !== undefined) { SchemaParser.decodeUnknownSync(workflow.output)(output); } - const durationMs = Date.now() - startedAt; emit({ _tag: "RunFinished", @@ -509,8 +428,7 @@ export function startRun( emit({ _tag: "RunCancelled", durationMs }); return { outcome: "cancelled" }; } - - const message = domainErrorMessage(err); + const message = err instanceof Error ? err.message : String(err); const stack = err instanceof Error ? err.stack : undefined; emit({ _tag: "RunFailed", diff --git a/src/server/concurrency.test.ts b/src/server/concurrency.test.ts index b0a3e18..051774e 100644 --- a/src/server/concurrency.test.ts +++ b/src/server/concurrency.test.ts @@ -19,6 +19,7 @@ import { describe, expect, test } from "bun:test"; import { defineConfig } from "../config"; import { getRunEvents, openStore } from "../persistence/store"; import { createSlowFakeAdapter } from "../replay/adapter"; +import { makeAgentRuntime } from "../runtime/agent-runtime"; import { admitRun } from "./admission"; import { serve } from "./http"; @@ -72,13 +73,15 @@ describe("phase 5 P1: per-run working trees (D28) over the real server", () => { const db = openStore(join(root, "factory.db")); const server = serve({ db, - adapter: createSlowFakeAdapter( - [ - { type: "TEXT_MESSAGE_START" }, - { type: "TEXT_MESSAGE_CONTENT", delta: "hi" }, - { type: "TEXT_MESSAGE_END" }, - ], - 30, + runtime: makeAgentRuntime( + createSlowFakeAdapter( + [ + { type: "TEXT_MESSAGE_START" }, + { type: "TEXT_MESSAGE_CONTENT", delta: "hi" }, + { type: "TEXT_MESSAGE_END" }, + ], + 30, + ), ), port: 0, config: defineConfig({ @@ -159,13 +162,15 @@ describe("phase 5 P1: per-run working trees (D28) over the real server", () => { const db = openStore(join(root, "factory.db")); const server = serve({ db, - adapter: createSlowFakeAdapter( - [ - { type: "TEXT_MESSAGE_START" }, - { type: "TEXT_MESSAGE_CONTENT", delta: "hi" }, - { type: "TEXT_MESSAGE_END" }, - ], - 300, + runtime: makeAgentRuntime( + createSlowFakeAdapter( + [ + { type: "TEXT_MESSAGE_START" }, + { type: "TEXT_MESSAGE_CONTENT", delta: "hi" }, + { type: "TEXT_MESSAGE_END" }, + ], + 300, + ), ), port: 0, config: defineConfig({ @@ -219,13 +224,15 @@ describe("phase 5 P1: per-run working trees (D28) over the real server", () => { const db = openStore(join(root, "factory.db")); const server = serve({ db, - adapter: createSlowFakeAdapter( - [ - { type: "TEXT_MESSAGE_START" }, - { type: "TEXT_MESSAGE_CONTENT", delta: "hi" }, - { type: "TEXT_MESSAGE_END" }, - ], - 1_000, + runtime: makeAgentRuntime( + createSlowFakeAdapter( + [ + { type: "TEXT_MESSAGE_START" }, + { type: "TEXT_MESSAGE_CONTENT", delta: "hi" }, + { type: "TEXT_MESSAGE_END" }, + ], + 1_000, + ), ), port: 0, config: defineConfig({ diff --git a/src/server/daemon.ts b/src/server/daemon.ts index 01ef0fd..8f506ab 100644 --- a/src/server/daemon.ts +++ b/src/server/daemon.ts @@ -2,18 +2,22 @@ * `factory serve` — wires the HTTP API/SSE (`server/http.ts`) and the config * scheduler (`server/scheduler.ts`) into one running process (AGENTS.md's * third column). Automatic dispatch is not daemon logic since epic #19: it - * is project policy living in scheduled wrapper workflows — e.g. the sample - * project's Ready sweep — fired by the scheduler loop on their cron. + * is project policy living in scheduled wrapper workflows. + * + * Issue #36: the daemon now has an Effect composition root. `server/daemon.ts` + * builds a `Layer` for the agent runtime, creates a `ManagedRuntime`, and + * passes that runtime to `serve()` and the scheduler. The adapter is no longer + * threaded by hand through `ServerOptions`, `StartTrackedRunOptions`, + * `DispatchEnv` and the run/step options; it is resolved from context inside + * `runtime/agent-step.ts`. */ import { mkdir } from "node:fs/promises"; import { dirname } from "node:path"; -import { Effect, type Fiber } from "effect"; +import { Effect, type Fiber, ManagedRuntime } from "effect"; type AnyFiber = Fiber.Fiber; import type { FactoryConfig } from "../config"; import { openStore } from "../persistence/store"; -import type { AgentAdapter } from "../runtime/agent-adapter"; -import { opencodeAdapter } from "../runtime/opencode-adapter"; import { serve } from "./http"; import { type DispatchEnv, type WorkspaceSpec } from "./runs"; import { @@ -23,11 +27,12 @@ import { toRuntimeSchedules, type SchedulerDeps, } from "./scheduler"; +import { AgentRuntimeLayer } from "../runtime/agent-runtime"; +import { opencodeAdapter } from "../runtime/opencode-adapter"; export interface DaemonOptions { readonly dbPath: string; readonly port?: number; - readonly adapter?: AgentAdapter; /** Issue #16: the tick cadence of the config schedules, over the default. */ readonly schedulerIntervalMs?: number; /** @@ -51,11 +56,14 @@ export const DEFAULT_SCHEDULER_INTERVAL_MS = 30_000; export async function startDaemon(options: DaemonOptions): Promise { await mkdir(dirname(options.dbPath), { recursive: true }); const db = openStore(options.dbPath); - const adapter = options.adapter ?? opencodeAdapter; + + const adapter = options.config?.agent.adapter ?? opencodeAdapter; + const layer = AgentRuntimeLayer(adapter); + const runtime = ManagedRuntime.make(layer); const server = serve({ db, - adapter, + runtime, port: options.port, ...(options.config !== undefined ? { config: options.config } : {}), }); @@ -76,7 +84,6 @@ export async function startDaemon(options: DaemonOptions): Promise workspace, repo, maxConcurrentRuns, - adapter, maxDispatchDepth: config.maxDispatchDepth, maxChildrenPerRun: config.maxChildrenPerRun, }; @@ -84,7 +91,7 @@ export async function startDaemon(options: DaemonOptions): Promise schedules: toRuntimeSchedules(config), fire: makeScheduleFire({ db, - adapter, + runtime, maxConcurrentRuns, workspace, repo, diff --git a/src/server/dedupe.test.ts b/src/server/dedupe.test.ts index 85059bd..845d8ac 100644 --- a/src/server/dedupe.test.ts +++ b/src/server/dedupe.test.ts @@ -12,6 +12,7 @@ import { join } from "node:path"; import { describe, expect, test } from "bun:test"; import { openStore, getRunEvents } from "../persistence/store"; import { createSlowFakeAdapter } from "../replay/adapter"; +import { makeAgentRuntime } from "../runtime/agent-runtime"; import { defineWorkflow, Schema } from "../workflow"; import { createDedupeRegistry, DedupeKeyError } from "../lib/dedupe"; import { activeRunIds, getActiveHandle, isActive, startTrackedRun } from "./runs"; @@ -32,6 +33,8 @@ const SLOW_ADAPTER = createSlowFakeAdapter( 25, ); +const runtime = makeAgentRuntime(SLOW_ADAPTER); + const echoWorkflow = defineWorkflow("echo-wf", { input: Schema.Struct({}), run: async (ctx) => { @@ -72,10 +75,9 @@ interface StartOpts { } function start(db: ReturnType, opts: StartOpts): Promise { - return startTrackedRun(db, echoWorkflow, { + return startTrackedRun(runtime, db, echoWorkflow, { dir: opts.dir, input: {}, - adapter: SLOW_ADAPTER, dedupeRegistry: opts.registry ?? createDedupeRegistry(), maxConcurrentRuns: 4, ...(opts.runId !== undefined ? { runId: opts.runId } : {}), @@ -111,13 +113,11 @@ function workspaces(root: string): WorkspaceSpec { function withWorkspaces(root: string): Record { return { workspace: workspaces(root), - adapter: SLOW_ADAPTER, maxConcurrentRuns: 10, }; } describe("dedupe keys through ctx.dispatch (issue #15)", () => { - const driftAdapter = SLOW_ADAPTER; const busyWorkflow = defineWorkflow("busy-wf", { input: Schema.Struct({ n: Schema.Number }), workspace: { kind: "scratch" }, @@ -147,12 +147,11 @@ describe("dedupe keys through ctx.dispatch (issue #15)", () => { }, }); - const parentRunId = await startTrackedRun(db, parent, { + const parentRunId = await startTrackedRun(runtime, db, parent, { input: { wait: 3 }, - adapter: driftAdapter, workspace: workspaces(join(root)), maxConcurrentRuns: 10, - dispatchEnv: { ...withWorkspaces(join(root)), adapter: driftAdapter } as never, + dispatchEnv: withWorkspaces(join(root)), }); await waitFor(() => !isActive(parentRunId)); @@ -207,13 +206,12 @@ describe("dedupe keys through ctx.dispatch (issue #15)", () => { }, }); - const parentRunId = (await startTrackedRun(db, parent, { + const parentRunId = await startTrackedRun(runtime, db, parent, { input: {}, - adapter: driftAdapter, workspace: workspaces(join(root)), maxConcurrentRuns: 10, - dispatchEnv: { ...withWorkspaces(join(root)), adapter: driftAdapter } as never, - })) as string; + dispatchEnv: withWorkspaces(join(root)), + }); await waitFor(() => getRunEvents(db, parentRunId).some((event) => event.payload._tag === "RunDispatched"), @@ -344,9 +342,8 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { // allocation of an explicitly undefined dir fails after the claim await expect( - startTrackedRun(db, echoWorkflow, { + startTrackedRun(runtime, db, echoWorkflow, { input: {}, - adapter: SLOW_ADAPTER, dedupeKey: "item:44", dedupeRegistry: createDedupeRegistry(), }), @@ -366,11 +363,10 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { const db = openStore(join(root, "factory.db")); const registry = createDedupeRegistry(); - const holderRunId = await startTrackedRun(db, echoWorkflow, { + const holderRunId = await startTrackedRun(runtime, db, echoWorkflow, { runId: "run-cancelled", dir: join(root, "dir-1"), input: {}, - adapter: SLOW_ADAPTER, dedupeKey: "item:45", dedupeRegistry: registry, }); @@ -381,10 +377,9 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { await waitFor(() => !isActive(holderRunId)); const started = getRunEvents(db, holderRunId).some((e) => e.payload._tag === "RunStarted"); - const retryRunId = await startTrackedRun(db, echoWorkflow, { + const retryRunId = await startTrackedRun(runtime, db, echoWorkflow, { dir: join(root, "dir-2"), input: {}, - adapter: SLOW_ADAPTER, dedupeKey: "item:45", dedupeRegistry: registry, }); @@ -443,7 +438,7 @@ describe("dedupe keys over POST /api/runs (issue #15)", () => { const server = serve({ db, - adapter: SLOW_ADAPTER, + runtime: makeAgentRuntime(SLOW_ADAPTER), port: 0, config: defineConfig({ repo: { @@ -517,7 +512,7 @@ describe("dedupe keys over POST /api/runs (issue #15)", () => { const server = serve({ db, - adapter: SLOW_ADAPTER, + runtime: makeAgentRuntime(SLOW_ADAPTER), port: 0, config: defineConfig({ repo: { diff --git a/src/server/http.test.ts b/src/server/http.test.ts index ce1ea46..25438d0 100644 --- a/src/server/http.test.ts +++ b/src/server/http.test.ts @@ -5,6 +5,7 @@ import { describe, expect, test } from "bun:test"; import { Schema } from "effect"; import { defineWorkflow, type WorkflowDefinition } from "../workflow"; import { createSlowFakeAdapter } from "../replay/adapter"; +import { makeAgentRuntime } from "../runtime/agent-runtime"; import { getRunEvents, openStore } from "../persistence/store"; import { defineConfig, loadFactoryConfig } from "../config"; import registryWorkflow from "../../test/fixtures/registry-workflow"; @@ -89,7 +90,7 @@ describe("GET /api/workflows (D30)", () => { const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); const config = await loadFactoryConfig(FIXTURE_CONFIG); - const server = serve({ db, adapter, port: 0, config }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); const base = `http://localhost:${server.port}`; try { @@ -114,7 +115,7 @@ describe("GET /api/workflows (D30)", () => { const dir = mkdtempSync(join(tmpdir(), "factory-workflows-empty-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); - const server = serve({ db, adapter, port: 0 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); const base = `http://localhost:${server.port}`; try { @@ -134,7 +135,7 @@ describe("GET /api/workflows (D30)", () => { const dir = mkdtempSync(join(tmpdir(), "factory-workflows-empty-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); - const server = serve({ db, adapter, port: 0 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); const base = `http://localhost:${server.port}`; try { @@ -154,7 +155,7 @@ describe("phase 4 SPA serving", () => { const dir = mkdtempSync(join(tmpdir(), "factory-spa-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([{ type: "TEXT_MESSAGE_START" }], 1); - const server = serve({ db, adapter, port: 0 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); const base = `http://localhost:${server.port}`; try { @@ -237,7 +238,7 @@ describe("SSE keepalive", () => { const dir = mkdtempSync(join(tmpdir(), "factory-sse-keepalive-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); - const server = serve({ db, adapter, port: 0, sseKeepaliveMs: 25 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, sseKeepaliveMs: 25 }); const base = `http://localhost:${server.port}`; try { @@ -273,7 +274,7 @@ describe("SSE keepalive", () => { const dir = mkdtempSync(join(tmpdir(), "factory-sse-keepalive-stop-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); - const server = serve({ db, adapter, port: 0, sseKeepaliveMs: 10 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, sseKeepaliveMs: 10 }); const base = `http://localhost:${server.port}`; try { @@ -313,7 +314,7 @@ describe("phase 3 HTTP API + SSE", () => { ], 1, ); - const server = serve({ db, adapter, port: 0 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); const base = `http://localhost:${server.port}`; try { @@ -362,7 +363,7 @@ describe("phase 3 HTTP API + SSE", () => { ], 500, ); - const server = serve({ db, adapter, port: 0 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); const base = `http://localhost:${server.port}`; try { @@ -398,7 +399,7 @@ describe("phase 3 HTTP API + SSE", () => { ], 60, ); - const server = serve({ db, adapter, port: 0 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); const base = `http://localhost:${server.port}`; try { @@ -437,7 +438,7 @@ describe("phase 3 HTTP API + SSE", () => { ], 400, ); - const server = serve({ db, adapter, port: 0 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); const base = `http://localhost:${server.port}`; try { @@ -475,7 +476,7 @@ describe("phase 3 HTTP API + SSE", () => { ], 1, ); - const server = serve({ db, adapter, port: 0 }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); const base = `http://localhost:${server.port}`; try { @@ -521,7 +522,7 @@ describe("POST /api/runs {workflowId, input} (D31)", () => { 1, ); const config = await loadFactoryConfig(FIXTURE_CONFIG); - const server = serve({ db, adapter, port: 0, config }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); const base = `http://localhost:${server.port}`; try { @@ -547,7 +548,7 @@ describe("POST /api/runs {workflowId, input} (D31)", () => { const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); const config = await loadFactoryConfig(FIXTURE_CONFIG); - const server = serve({ db, adapter, port: 0, config }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); const base = `http://localhost:${server.port}`; try { @@ -593,7 +594,7 @@ describe("POST /api/runs {workflowId, input} (D31)", () => { maxConcurrentRuns: 1, retainedWorkspaces: 10, }); - const server = serve({ db, adapter, port: 0, config }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); const base = `http://localhost:${server.port}`; try { @@ -653,7 +654,7 @@ describe("POST /api/runs {workflowId, input} (D31)", () => { maxConcurrentRuns: 3, retainedWorkspaces: 10, }); - const server = serve({ db, adapter, port: 0, config }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); const base = `http://localhost:${server.port}`; try { @@ -720,7 +721,7 @@ describe("a scratch workflow through POST /api/runs (issue #13)", () => { maxConcurrentRuns: 3, retainedWorkspaces: 10, }); - const server = serve({ db, adapter, port: 0, config }); + const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); const base = `http://localhost:${server.port}`; try { @@ -793,7 +794,7 @@ describe("ctx.dispatch through POST /api/runs (issue #14)", () => { const server = serve({ db, - adapter, + runtime: makeAgentRuntime(adapter), port: 0, config: defineConfig({ repo: { @@ -891,7 +892,7 @@ describe("GET /api/schedules (issue #17)", () => { }); const server = serve({ db, - adapter: createSlowFakeAdapter([], 1), + runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0, config: schedulesConfig(root, workflow), }); @@ -958,7 +959,12 @@ describe("GET /api/schedules (issue #17)", () => { }, ], }); - const server = serve({ db, adapter: createSlowFakeAdapter([], 1), port: 0, config }); + const server = serve({ + db, + runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), + port: 0, + config, + }); const base = `http://localhost:${server.port}`; const { Cron } = await import("effect"); @@ -979,7 +985,7 @@ describe("GET /api/schedules (issue #17)", () => { test("serves an empty list on a legacy no-config server", async () => { const root = mkdtempSync(join(tmpdir(), "factory-schedules-empty-test-")); const db = openStore(join(root, "factory.db")); - const server = serve({ db, adapter: createSlowFakeAdapter([], 1), port: 0 }); + const server = serve({ db, runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0 }); const base = `http://localhost:${server.port}`; try { @@ -1021,7 +1027,12 @@ describe("GET /api/schedules (issue #17)", () => { }, ], }); - const server = serve({ db, adapter: createSlowFakeAdapter([], 1), port: 0, config }); + const server = serve({ + db, + runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -1117,7 +1128,12 @@ describe("POST /api/schedules/:id/run (issue #17)", () => { }, ], }); - const server = serve({ db, adapter: createSlowFakeAdapter([], 1), port: 0, config }); + const server = serve({ + db, + runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -1166,7 +1182,12 @@ describe("POST /api/schedules/:id/run (issue #17)", () => { { id: "nightly-run", workflow: workflow.id, input: {}, cron: "0 3 * * *", timezone: "UTC" }, ], }); - const server = serve({ db, adapter: createSlowFakeAdapter([], 1), port: 0, config }); + const server = serve({ + db, + runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -1204,7 +1225,12 @@ describe("POST /api/schedules/:id/run (issue #17)", () => { { id: "slow-skip", workflow: workflow.id, input: {}, cron: "0 3 * * *", timezone: "UTC" }, ], }); - const server = serve({ db, adapter: createSlowFakeAdapter([], 1), port: 0, config }); + const server = serve({ + db, + runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -1264,7 +1290,12 @@ describe("POST /api/schedules/:id/run (issue #17)", () => { }, ], }); - const server = serve({ db, adapter: createSlowFakeAdapter([], 1), port: 0, config }); + const server = serve({ + db, + runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { diff --git a/src/server/http.ts b/src/server/http.ts index 26ea22e..738a0a5 100644 --- a/src/server/http.ts +++ b/src/server/http.ts @@ -28,13 +28,14 @@ */ import type { Database } from "bun:sqlite"; +import { ManagedRuntime } from "effect"; import { isTerminal, type RunEvent } from "../events"; import type { FactoryConfig } from "../config"; import { Schema, SchemaParser } from "effect"; import { resetClone, type GitIdentity } from "../lib/clone"; import { loadWorkflow } from "../lib/load-workflow"; import { getRunEvents, listRuns, type RunSummary } from "../persistence/store"; -import type { AgentAdapter } from "../runtime/agent-adapter"; +import { AgentRuntime } from "../runtime/agent-runtime"; import index from "../web/index.html"; import { admitRun } from "./admission"; import { subscribe } from "./pubsub"; @@ -74,7 +75,12 @@ export interface ScheduleSummary { export interface ServerOptions { readonly db: Database; - readonly adapter: AgentAdapter; + /** + * Issue #36: the Effect managed runtime that provides the agent runtime + * service. The adapter is resolved from context inside `startTrackedRun` + * rather than threaded through `ServerOptions`. + */ + readonly runtime: ManagedRuntime.ManagedRuntime; /** * The loaded `factory.config.ts` (D27). Absent = the phase 3 path-based API * behaves exactly as before (no limit, explicit dir+clone, `/api/workflows` @@ -131,7 +137,7 @@ function listSummaries(db: Database): ReadonlyArray { * it `ctx.dispatch` throws, because in-process/path-based legacy runs have * no registry to start children from. */ -function dispatchEnvFor(config: FactoryConfig, adapter: AgentAdapter): DispatchEnv { +function dispatchEnvFor(config: FactoryConfig): DispatchEnv { return { workspace: { workspaceRoot: config.workspaceRoot, @@ -141,7 +147,6 @@ function dispatchEnvFor(config: FactoryConfig, adapter: AgentAdapter): DispatchE }, repo: { slug: config.repo.slug, baseBranch: config.repo.baseBranch }, maxConcurrentRuns: config.maxConcurrentRuns, - adapter, maxDispatchDepth: config.maxDispatchDepth, maxChildrenPerRun: config.maxChildrenPerRun, }; @@ -155,14 +160,13 @@ function dispatchEnvFor(config: FactoryConfig, adapter: AgentAdapter): DispatchE */ function configRunOptions( config: FactoryConfig, - adapter: AgentAdapter, maxConcurrentRuns: number | undefined, extra: { readonly scheduleId?: string; readonly dedupeKey?: string; readonly agentOverrides?: { readonly model?: string }; } = {}, -): Omit { +): Omit { return { workspace: { workspaceRoot: config.workspaceRoot, @@ -172,7 +176,7 @@ function configRunOptions( }, repo: { slug: config.repo.slug, baseBranch: config.repo.baseBranch }, maxConcurrentRuns, - dispatchEnv: dispatchEnvFor(config, adapter), + dispatchEnv: dispatchEnvFor(config), ...(extra.scheduleId !== undefined ? { scheduleId: extra.scheduleId } : {}), ...(extra.dedupeKey !== undefined ? { dedupeKey: extra.dedupeKey } : {}), ...(extra.agentOverrides !== undefined ? { agentOverrides: extra.agentOverrides } : {}), @@ -388,15 +392,14 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise } const startOptions: StartTrackedRunOptions = { - ...configRunOptions(runEnv, options.adapter, maxConcurrentRuns, { + ...configRunOptions(runEnv, maxConcurrentRuns, { dedupeKey: typeof body.dedupeKey === "string" ? body.dedupeKey : undefined, }), input: decodedInput, - adapter: options.adapter, }; let runId: string; try { - runId = await startTrackedRun(options.db, workflow, startOptions); + runId = await startTrackedRun(options.runtime, options.db, workflow, startOptions); } catch (err) { if (err instanceof ConcurrencyLimitError) { return json({ error: err.message }, { status: 409 }); @@ -453,9 +456,8 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise ? undefined : { slug: runEnv.repo.slug, baseBranch: runEnv.repo.baseBranch }, ...(runEnv !== undefined ? { maxConcurrentRuns } : {}), - ...(runEnv !== undefined ? { dispatchEnv: dispatchEnvFor(runEnv, options.adapter) } : {}), + ...(runEnv !== undefined ? { dispatchEnv: dispatchEnvFor(runEnv) } : {}), input: body.input, - adapter: options.adapter, ...(typeof body.dedupeKey === "string" ? { dedupeKey: body.dedupeKey } : {}), ...(body.clone !== undefined && typeof body.dir === "string" ? { prepareWorkspace: true } @@ -463,7 +465,7 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise }; let runId: string; try { - runId = await startTrackedRun(options.db, workflow, startOptions); + runId = await startTrackedRun(options.runtime, options.db, workflow, startOptions); } catch (err) { if (err instanceof ConcurrencyLimitError) { return json({ error: err.message }, { status: 409 }); @@ -556,14 +558,13 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise // it fires regardless. let runId: string; try { - runId = await startTrackedRun(options.db, workflow, { - ...configRunOptions(options.config, options.adapter, maxConcurrentRuns, { + runId = await startTrackedRun(options.runtime, options.db, workflow, { + ...configRunOptions(options.config, maxConcurrentRuns, { scheduleId: schedule.id, ...(schedule.overlap === "skip" ? { dedupeKey: `schedule:${schedule.id}` } : {}), ...(schedule.agent !== undefined ? { agentOverrides: schedule.agent } : {}), }), input: schedule.input, - adapter: options.adapter, }); } catch (err) { if (err instanceof ConcurrencyLimitError) { diff --git a/src/server/nested-runs.test.ts b/src/server/nested-runs.test.ts index cf9c45c..7c9e8ac 100644 --- a/src/server/nested-runs.test.ts +++ b/src/server/nested-runs.test.ts @@ -20,6 +20,7 @@ import { getRunEvents, listRuns, openStore } from "../persistence/store"; import { createSlowFakeAdapter } from "../replay/adapter"; import { activeRunIds, isActive, startTrackedRun } from "./runs"; import { defineWorkflow } from "../workflow"; +import { makeAgentRuntime } from "../runtime/agent-runtime"; const SLOW_ADAPTER = createSlowFakeAdapter( [ @@ -30,6 +31,8 @@ const SLOW_ADAPTER = createSlowFakeAdapter( 25, ); +const runtime = makeAgentRuntime(SLOW_ADAPTER); + // The child does a real ctx.exec so its outcome is observable in its log. // Scratch (issue #13): no mirror refresh, so no git fixture is needed. const numberedChild = defineWorkflow("child-wf", { @@ -89,7 +92,7 @@ describe("nested runs through the daemon (issue #14)", () => { }, }); - const parentRunId = await startTrackedRun(db, parent, { + const parentRunId = await startTrackedRun(runtime, db, parent, { workspace: { workspaceRoot: join(root, "workspaces"), sshUrl: join(root, "seed-not-used"), @@ -97,7 +100,6 @@ describe("nested runs through the daemon (issue #14)", () => { retainedWorkspaces: 10, }, input: {}, - adapter: SLOW_ADAPTER, dispatchEnv: { workspace: { workspaceRoot: join(root, "workspaces"), @@ -105,7 +107,6 @@ describe("nested runs through the daemon (issue #14)", () => { identity: { name: "Test Bot", email: "test@factory.local" }, retainedWorkspaces: 10, }, - adapter: SLOW_ADAPTER, }, }); @@ -145,9 +146,8 @@ describe("nested runs through the daemon (issue #14)", () => { const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "workspaces"), { recursive: true }); - await startTrackedRun(db, selfDispatching, { + await startTrackedRun(runtime, db, selfDispatching, { input: {}, - adapter: SLOW_ADAPTER, workspace: { workspaceRoot: join(root, "workspaces"), sshUrl: join(root, "seed-not-used"), @@ -162,7 +162,6 @@ describe("nested runs through the daemon (issue #14)", () => { identity: { name: "Test Bot", email: "test@factory.local" }, retainedWorkspaces: 20, }, - adapter: SLOW_ADAPTER, maxConcurrentRuns: 20, }, }); @@ -198,9 +197,8 @@ describe("nested runs through the daemon (issue #14)", () => { const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "workspaces"), { recursive: true }); - const parentRunId = await startTrackedRun(db, manyChildren, { + const parentRunId = await startTrackedRun(runtime, db, manyChildren, { input: {}, - adapter: SLOW_ADAPTER, workspace: { workspaceRoot: join(root, "workspaces"), sshUrl: join(root, "seed-not-used"), @@ -215,7 +213,6 @@ describe("nested runs through the daemon (issue #14)", () => { identity: { name: "Test Bot", email: "test@factory.local" }, retainedWorkspaces: 30, }, - adapter: SLOW_ADAPTER, maxConcurrentRuns: 40, maxChildrenPerRun: 5, }, @@ -250,9 +247,8 @@ describe("nested runs through the daemon (issue #14)", () => { }, }); - const parentRunId = await startTrackedRun(db, parent, { + const parentRunId = await startTrackedRun(runtime, db, parent, { input: {}, - adapter: SLOW_ADAPTER, workspace: { workspaceRoot: join(root, "workspaces"), sshUrl: join(root, "seed-not-used"), @@ -267,7 +263,6 @@ describe("nested runs through the daemon (issue #14)", () => { identity: { name: "Test Bot", email: "test@factory.local" }, retainedWorkspaces: 5, }, - adapter: SLOW_ADAPTER, maxConcurrentRuns: 1, }, }); diff --git a/src/server/runs.test.ts b/src/server/runs.test.ts index dbbdb51..aee00d4 100644 --- a/src/server/runs.test.ts +++ b/src/server/runs.test.ts @@ -17,6 +17,7 @@ import type { RunEvent } from "../events"; import echoWorkflow from "../../test/fixtures/echo-workflow"; import { appendEvent, openStore, getRunEvents } from "../persistence/store"; import { createSlowFakeAdapter } from "../replay/adapter"; +import { makeAgentRuntime } from "../runtime/agent-runtime"; import { defineWorkflow, Schema } from "../workflow"; import { ConcurrencyLimitError, @@ -37,6 +38,8 @@ const SLOW_ADAPTER = createSlowFakeAdapter( 25, ); +const runtime = makeAgentRuntime(SLOW_ADAPTER); + const sleepWorkflow = defineWorkflow("sleep-test", { input: Schema.Struct({}), run: async (ctx) => { @@ -109,11 +112,10 @@ describe("startTrackedRun admission (M1: the slot is reserved before any await)" const db = openStore(join(root, "factory.db")); const gate = makeGate(); - const first = startTrackedRun(db, echoWorkflow, { + const first = startTrackedRun(runtime, db, echoWorkflow, { runId: "run-first", dir: join(root, "first-dir"), input: {}, - adapter: SLOW_ADAPTER, maxConcurrentRuns: 1, beforeStart: gate.wait, }); @@ -122,11 +124,10 @@ describe("startTrackedRun admission (M1: the slot is reserved before any await)" expect(activeRunIds()).toEqual(["run-first"]); expect(getActiveHandle("run-first")).toBeUndefined(); - const second = startTrackedRun(db, echoWorkflow, { + const second = startTrackedRun(runtime, db, echoWorkflow, { runId: "run-second", dir: join(root, "second-dir"), input: {}, - adapter: SLOW_ADAPTER, maxConcurrentRuns: 1, }); @@ -152,10 +153,9 @@ describe("startTrackedRun admission (M1: the slot is reserved before any await)" const db = openStore(join(root, "factory.db")); await expect( - startTrackedRun(db, echoWorkflow, { + startTrackedRun(runtime, db, echoWorkflow, { runId: "run-doomed", input: {}, - adapter: SLOW_ADAPTER, maxConcurrentRuns: 1, workspace: { workspaceRoot: join(root, "workspaces"), @@ -168,10 +168,9 @@ describe("startTrackedRun admission (M1: the slot is reserved before any await)" expect(activeRunIds()).toEqual([]); - const runId = await startTrackedRun(db, echoWorkflow, { + const runId = await startTrackedRun(runtime, db, echoWorkflow, { dir: join(root, "dir"), input: {}, - adapter: SLOW_ADAPTER, maxConcurrentRuns: 1, }); expect(activeRunIds()).toEqual([runId]); @@ -189,11 +188,10 @@ describe("cancel of a reserved-but-not-started run (L1)", () => { mkdirSync(join(root, "dir"), { recursive: true }); const gate = makeGate(); - const starting = startTrackedRun(db, sleepWorkflow, { + const starting = startTrackedRun(runtime, db, sleepWorkflow, { runId: "run-gated", dir: join(root, "dir"), input: {}, - adapter: SLOW_ADAPTER, maxConcurrentRuns: 1, beforeStart: gate.wait, }); @@ -259,11 +257,10 @@ describe("scratch workspaces through startTrackedRun (issue #13)", () => { const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "seed"), { recursive: true }); - const runId = await startTrackedRun(db, failingScratchWorkflow, { + const runId = await startTrackedRun(runtime, db, failingScratchWorkflow, { runId: "run-scratch-empty", workspace: { ...workspaceSpec(), sshUrl: join(root, "no-such-remote") }, input: {}, - adapter: SLOW_ADAPTER, }); await waitFor(() => !isActive(runId)); @@ -281,21 +278,19 @@ describe("scratch workspaces through startTrackedRun (issue #13)", () => { const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "seed"), { recursive: true }); - const okId = await startTrackedRun(db, scratchWorkflow, { + const okId = await startTrackedRun(runtime, db, scratchWorkflow, { runId: "run-scratch-ok", workspace: workspaceSpec(), input: {}, - adapter: SLOW_ADAPTER, }); await waitFor(() => !isActive(okId)); await new Promise((resolve) => setTimeout(resolve, 50)); // reap lands async expect(existsSync(join(root, "workspaces", "run-scratch-ok"))).toBe(false); - const badId = await startTrackedRun(db, failingScratchWorkflow, { + const badId = await startTrackedRun(runtime, db, failingScratchWorkflow, { runId: "run-scratch-bad", workspace: workspaceSpec(), input: {}, - adapter: SLOW_ADAPTER, }); await waitFor(() => !isActive(badId)); await new Promise((resolve) => setTimeout(resolve, 50)); @@ -327,11 +322,10 @@ describe("scratch workspaces through startTrackedRun (issue #13)", () => { // retention = 1: with the scratch dir correctly excluded, the only clone // survives: the leftover must not count toward retention. - await startTrackedRun(db, echoWorkflow, { + await startTrackedRun(runtime, db, echoWorkflow, { runId: "run-clone", workspace: { ...workspaceSpec(), sshUrl: seed, retainedWorkspaces: 1 }, input: {}, - adapter: SLOW_ADAPTER, }); await waitFor(() => !isActive("run-clone")); diff --git a/src/server/runs.ts b/src/server/runs.ts index 5ad9653..41265f7 100644 --- a/src/server/runs.ts +++ b/src/server/runs.ts @@ -1,20 +1,11 @@ /** * The server's in-memory run registry: which runs are currently live in * *this* process, so the HTTP API can cancel them and the dispatcher can - * enforce a WIP limit (D24, D29). Deliberately not derived from sqlite — a - * run whose process died is "interrupted" (D12), not "active"; only a - * `RunHandle` this process actually holds counts. - * - * A run is registered *before* it starts: `startTrackedRun` reserves the - * registry slot synchronously (check-then-set with no `await` between, so a - * concurrent HTTP start cannot lose the race) and only then allocates the - * workspace. A reserved slot is released if allocation/startup fails, and a - * `cancel` that arrives while a run is still reserving is deferred into run - * start rather than dropped. + * enforce a WIP limit (D24, D29). */ -import { Schema } from "effect"; import type { Database } from "bun:sqlite"; +import type { ManagedRuntime } from "effect"; import { rm } from "node:fs/promises"; import { admitRun } from "./admission"; import { appendEvent, getRunEvents, listRuns } from "../persistence/store"; @@ -23,9 +14,9 @@ import { startRun, type RunHandle } from "../runtime/run"; import type { GitIdentity } from "../lib/clone"; import { allocateWorkspace } from "../lib/workspace"; import { dedupeRegistry, type DedupeRegistry } from "../lib/dedupe"; -import type { AgentAdapter } from "../runtime/agent-adapter"; import type { DispatchChildFn, WorkflowDefinition, WorkspaceKind } from "../workflow"; import { publish } from "./pubsub"; +import { AgentRuntime } from "../runtime/agent-runtime"; interface ReservedSlot { cancelled: boolean; @@ -33,18 +24,15 @@ interface ReservedSlot { const active = new Map | ReservedSlot>(); -/** - * Issue #14: how deep a parent → child → grandchild chain may nest — the cap - * that keeps a workflow which dispatches itself from filling the daemon. - */ export const DEFAULT_MAX_DISPATCH_DEPTH = 5; -/** Issue #14: how many children one run itself may dispatch. */ export const DEFAULT_MAX_CHILDREN_PER_RUN = 20; -/** Issue #14: a `ctx.dispatch` rejected by a cap, for the parent to surface. */ -export class DispatchCapError extends Schema.TaggedError()("DispatchCapError", { - message: Schema.String, -}) {} +export class DispatchCapError extends Error { + constructor(message: string) { + super(message); + this.name = "DispatchCapError"; + } +} function isReserved(entry: RunHandle | ReservedSlot | undefined): boolean { return entry !== undefined && !("result" in entry) && "cancelled" in entry; @@ -74,16 +62,11 @@ export function getActiveHandle(runId: string): RunHandle | undefined { return entry as RunHandle; } -/** - * Cancellation for a run this process holds — whether it is already running - * (returns its handle), still reserving/allocation-bound (marks the slot so - * the run is cancelled the moment it starts), or unknown (`undefined`). - */ -export function cancelRegisteredRun(runId: string): +export function cancelRegisteredRun( + runId: string, +): | { readonly kind: "handle"; readonly handle: RunHandle } - | { - readonly kind: "reserved"; - } + | { readonly kind: "reserved" } | undefined { const entry = active.get(runId); if (entry === undefined) return undefined; @@ -102,63 +85,19 @@ export interface WorkspaceSpec { } export interface StartTrackedRunOptions { - /** An explicit directory wins; otherwise `workspace` allocates one per runId (D28). */ readonly dir?: string; readonly workspace?: WorkspaceSpec; - /** Write-back environment for the run (D32): came from config, not the caller. */ readonly repo?: RunRepo; readonly input: unknown; - readonly adapter: AgentAdapter; readonly runId?: string; - /** - * D29's ceiling, enforced atomically at reservation time — the reservation - * lands in the registry before any `await`, so two near-simultaneous - * start requests cannot both slip past it. Absent, no limit applies. - */ readonly maxConcurrentRuns?: number; - /** Injectable for tests: holds the reserved-but-not-started window open. */ readonly beforeStart?: () => Promise; - /** - * Issue #14: the environment a child run of this run starts with, so - * `ctx.dispatch` can allocate a workspace, admission-limit and repo for it. - * Absent, the workflow's `ctx.dispatch` throws (dispatch needs a - * config-backed daemon run; in-process execution is legacy). Children - * inherit it, so grandchildren work too. - */ readonly dispatchEnv?: DispatchEnv; - /** - * Issue #14: this run's parent, when it was started by `ctx.dispatch` — - * recorded on `RunStarted.parentId` so the UI can navigate child → parent. - */ readonly parentRunId?: string; - /** - * Issue #15: this run's dedupe key. Claimed synchronously at start (before - * any await) in the run's own registry slot manner — two near-simultaneous - * starts with the same key cannot both slip through the check-then-claim — - * and released the moment the run settles in any terminal state, including - * a failed startup. Absent, the run claims nothing. - */ readonly dedupeKey?: string; - /** - * Issue #15: `true` when a caller (the dispatch path) already claimed the - * key synchronously before handing the start over, so the claim must not be - * re-asserted here. Release still happens here, keyed to this run id. - */ readonly dedupeKeyClaimed?: boolean; - /** - * Issue #15: injectable holder registry, for tests. Absent, the daemon's - * shared process-wide registry (`lib/dedupe.ts`) is used. - */ readonly dedupeRegistry?: DedupeRegistry; - /** - * Issue #16: the schedule that started this run, when any did - passed - * through to `RunStarted.scheduleId`. - */ readonly scheduleId?: string; - /** - * Issue #16: agent-level overrides the starting schedule carries. Passed - * through to the run's model precedence chain. - */ readonly agentOverrides?: { readonly model?: string }; /** * ADR 0012 §3 (#37): when true, the runtime calls `adapter.prepareWorkspace` @@ -169,36 +108,21 @@ export interface StartTrackedRunOptions { readonly prepareWorkspace?: boolean; } -/** - * The environment `ctx.dispatch`'s children share with their parent (issue - * #14): workspace provisioning, write-back environment, adapter and the - * concurrency ceiling — the same wiring a config-backed server wraps every - * run in, simply reused for child runs (children inherit it, so a child can - * dispatch grandchildren). - */ export interface DispatchEnv { readonly workspace?: WorkspaceSpec; readonly repo?: RunRepo; readonly maxConcurrentRuns?: number; - readonly adapter: AgentAdapter; - /** Issue #14: dispatch depth / per-run child caps, over the defaults. */ readonly maxDispatchDepth?: number; readonly maxChildrenPerRun?: number; - /** Issue #15: injectable holder registry, over the daemon's shared one. */ readonly dedupeRegistry?: DedupeRegistry; } -/** - * The run this `runId` was started from (`RunStarted.parentId`), or none — - * a run started without `ctx.dispatch` has no parent. - */ function parentOf(db: Database, runId: string): string | undefined { const started = getRunEvents(db, runId).find((event) => event.payload._tag === "RunStarted"); if (started === undefined || started.payload._tag !== "RunStarted") return undefined; return started.payload.parentId; } -/** How deep this run sits in the parent → child chain (a top-level run is 0). */ function dispatchDepth(db: Database, runId: string): number { let depth = 0; let cursor = runId; @@ -211,19 +135,12 @@ function dispatchDepth(db: Database, runId: string): number { return depth; } -/** - * Issue #14: start a child run of `parentRunId`. - * - * Every rejection is a throw, so `ctx.dispatch` rejects and the *parent* - * fails visibly — nothing about a child's admission is allowed to be a - * silent drop. The checks are synchronous over in-memory state and the sync - * sqlite event log, so by the time the parent carries on, admission has - * already happened. A validated child is started without being awaited: the - * parent never exposes an awaitable child (D-epic 19). The child's own - * startup failures land on the child's event log / console, never in the - * parent's. - */ +function countDispatchedChildren(db: Database, runId: string): number { + return getRunEvents(db, runId).filter((event) => event.payload._tag === "RunDispatched").length; +} + async function dispatchChildRun( + runtime: ManagedRuntime.ManagedRuntime, db: Database, env: DispatchEnv, parentRunId: string, @@ -239,44 +156,43 @@ async function dispatchChildRun( env.maxConcurrentRuns !== undefined && !admitRun(env.maxConcurrentRuns, activeRunIds().length) ) { - throw new ConcurrencyLimitError({ maxConcurrentRuns: env.maxConcurrentRuns }); + throw new ConcurrencyLimitError(env.maxConcurrentRuns); } const depth = dispatchDepth(db, parentRunId); if (depth + 1 > maxDepth) { - throw new DispatchCapError({ - message: - `dispatch depth exceeded: run ${parentRunId} is nested ${depth} levels deep; ` + + throw new DispatchCapError( + `dispatch depth exceeded: run ${parentRunId} is nested ${depth} levels deep; ` + `max ${maxDepth} (a workflow that dispatches itself must not fill the daemon)`, - }); + ); } const childCount = countDispatchedChildren(db, parentRunId); if (childCount >= maxChildren) { - throw new DispatchCapError({ - message: `dispatch child cap exceeded: run ${parentRunId} already dispatched ${childCount} children; max ${maxChildren}`, - }); + throw new DispatchCapError( + `dispatch child cap exceeded: run ${parentRunId} already dispatched ${childCount} children; max ${maxChildren}`, + ); } - // Issue #15: the collision check is synchronous with the child id in hand - // and *before* the fire-and-forget start, so a collision throws into the - // parent here instead of being swallowed by the un-awaited start's catch. const childRunId = `run-${crypto.randomUUID()}`; if (opts?.dedupeKey !== undefined) registry.claim(opts.dedupeKey, childRunId); void (async () => { - await startTrackedRun(db, child, { - runId: childRunId, - ...(env.workspace !== undefined ? { workspace: env.workspace } : {}), - ...(env.repo !== undefined ? { repo: env.repo } : {}), - ...(env.maxConcurrentRuns !== undefined ? { maxConcurrentRuns: env.maxConcurrentRuns } : {}), - input, - adapter: env.adapter, - parentRunId, - dispatchEnv: env, - ...(opts?.dedupeKey !== undefined ? { dedupeKey: opts.dedupeKey } : {}), - ...(opts?.dedupeKey !== undefined ? { dedupeKeyClaimed: true } : {}), - ...(env.dedupeRegistry !== undefined ? { dedupeRegistry: env.dedupeRegistry } : {}), - }).catch((err: unknown) => { + try { + await startTrackedRun(runtime, db, child, { + runId: childRunId, + ...(env.workspace !== undefined ? { workspace: env.workspace } : {}), + ...(env.repo !== undefined ? { repo: env.repo } : {}), + ...(env.maxConcurrentRuns !== undefined + ? { maxConcurrentRuns: env.maxConcurrentRuns } + : {}), + input, + parentRunId, + dispatchEnv: env, + ...(opts?.dedupeKey !== undefined ? { dedupeKey: opts.dedupeKey } : {}), + ...(opts?.dedupeKey !== undefined ? { dedupeKeyClaimed: true } : {}), + ...(env.dedupeRegistry !== undefined ? { dedupeRegistry: env.dedupeRegistry } : {}), + }); + } catch (err) { if (opts?.dedupeKey !== undefined) { (env.dedupeRegistry ?? dedupeRegistry).release(opts.dedupeKey, childRunId); } @@ -284,27 +200,18 @@ async function dispatchChildRun( `nested run start failed (parent ${parentRunId}, child ${childRunId}):` + `${err instanceof Error ? err.message : String(err)}`, ); - }); + } })(); return childRunId; } -function countDispatchedChildren(db: Database, runId: string): number { - return getRunEvents(db, runId).filter((event) => event.payload._tag === "RunDispatched").length; -} - -/** Starts a run, persists+publishes every event, and tracks it until terminal. */ export async function startTrackedRun( + runtime: ManagedRuntime.ManagedRuntime, db: Database, workflow: WorkflowDefinition, options: StartTrackedRunOptions, ): Promise { - // crypto.randomUUID(), not `run-${Date.now()}`: the server can have more than one run start - // within the same millisecond (concurrent HTTP POSTs, or tests running in the same process), - // and a collided runId cross-wires the pubsub channel and `active` registry between two - // unrelated runs — one run's SSE watcher can then see the other's terminal event and close - // its own db while its real run is still writing to it. const runId = options.runId ?? `run-${crypto.randomUUID()}`; const existing = active.get(runId); @@ -312,27 +219,22 @@ export async function startTrackedRun( throw new Error(`run ${runId} is already active`); } - // Issue #15: claim the dedupe key synchronously — check-then-claim with no - // `await` in between, the same atomicity the registry slot reservation has — - // so two near-simultaneous starts on the same key cannot both slip past. A - // collision throws before anything is started, leaving no trace. const registry = options.dedupeRegistry ?? dedupeRegistry; if (options.dedupeKey !== undefined && options.dedupeKeyClaimed !== true) { registry.claim(options.dedupeKey, runId); } if (existing === undefined && options.maxConcurrentRuns !== undefined) { if (!admitRun(options.maxConcurrentRuns, active.size)) { - throw new ConcurrencyLimitError({ maxConcurrentRuns: options.maxConcurrentRuns }); + throw new ConcurrencyLimitError(options.maxConcurrentRuns); } active.set(runId, { cancelled: false }); } try { - if (options.beforeStart !== undefined) await options.beforeStart(); + if (options.beforeStart !== undefined) { + await options.beforeStart(); + } - // Issue #13: a scratch workspace takes its kind from the workflow and - // never evicts clone workspaces — leftover scratch dirs (kept failures) - // are excluded from retention via the run log's recorded kinds. const kind: WorkspaceKind = workflow.workspace?.kind ?? "clone"; const scratchEntries = options.workspace !== undefined && kind === "scratch" @@ -349,7 +251,10 @@ export async function startTrackedRun( ? undefined : await allocateWorkspace({ runId, - ...options.workspace, + workspaceRoot: options.workspace.workspaceRoot, + sshUrl: options.workspace.sshUrl, + identity: options.workspace.identity, + retainedWorkspaces: options.workspace.retainedWorkspaces, kind, ...(scratchEntries !== undefined ? { scratchEntries } : {}), protectedEntries: [runId, ...activeRunIds()], @@ -357,20 +262,15 @@ export async function startTrackedRun( if (dir === undefined) throw new Error("startTrackedRun needs `dir` or `workspace`"); - // Whether the runtime allocated this dir itself (explicit callers' dirs, - // incl. the legacy path-based API, are not reaped — they are caller-owned). const workspaceAllocated = options.dir === undefined; - // Issue #14: the per-run dispatch member comes from the run's dispatch - // environment, bound at this run's id — the parent's own log is where the - // depth walk and the child count are read from. const dispatch: DispatchChildFn | undefined = options.dispatchEnv === undefined ? undefined : (child, input, opts) => - dispatchChildRun(db, options.dispatchEnv!, runId, child, input, opts); + dispatchChildRun(runtime, db, options.dispatchEnv!, runId, child, input, opts); - const handle = startRun(workflow, { + const handle = await startRun(workflow, runtime, { runId, dir, ...(options.repo !== undefined ? { repo: options.repo } : {}), @@ -383,16 +283,12 @@ export async function startTrackedRun( ...(options.scheduleId !== undefined ? { scheduleId: options.scheduleId } : {}), ...(options.agentOverrides !== undefined ? { agentOverrides: options.agentOverrides } : {}), input: options.input, - adapter: options.adapter, onEvent: (event) => { appendEvent(db, event); publish(runId, event); }, }); - // A cancel that arrived while this run was only a reserved slot is - // deferred into run start (L1): the run starts, is cancelled immediately, - // and ends as a clean RunCancelled instead of orphaning the slot. const beforeStartEntry = active.get(runId); const reservedSlot = isReserved(beforeStartEntry) ? (beforeStartEntry as ReservedSlot) @@ -402,17 +298,9 @@ export async function startTrackedRun( void handle.result.finally(() => { if (active.get(runId) === handle) active.delete(runId); - // Issue #15: any terminal state — completed, failed, cancelled — - // releases the run's key. (An interrupted run — process death — is - // covered by the registry being per-process: the new process holds - // nothing.) if (options.dedupeKey !== undefined) registry.release(options.dedupeKey, runId); }); - // Issue #13: a scratch dir is reaped when the run succeeds — there is no - // tree worth keeping, and retention never applies to it. It is kept when - // the run fails (or is cancelled) so a failed precondition check remains - // inspectable. void handle.result.then((outcome) => { if (workspaceAllocated && kind === "scratch" && outcome.outcome === "completed") { void rm(dir, { recursive: true, force: true }); diff --git a/src/server/scheduler.ts b/src/server/scheduler.ts index 01bcabd..cf7f73a 100644 --- a/src/server/scheduler.ts +++ b/src/server/scheduler.ts @@ -24,7 +24,7 @@ * like any other dedupe collision; `"stack"` fires regardless. */ -import { Cron, Effect, Schedule, Schema } from "effect"; +import { Cron, Effect, Schedule, Schema, ManagedRuntime } from "effect"; import type { Database } from "bun:sqlite"; import type { FactoryConfig } from "../config"; import { DedupeKeyError, dedupeRegistry, type DedupeRegistry } from "../lib/dedupe"; @@ -37,7 +37,7 @@ import { type WorkspaceSpec, } from "./runs"; import { RunCancelledSignal, type RunRepo } from "../runtime/run"; -import type { AgentAdapter } from "../runtime/agent-adapter"; +import { AgentRuntime } from "../runtime/agent-runtime"; export class SchedulerError extends Schema.TaggedError()("SchedulerError", { cause: Schema.Defect(), @@ -252,7 +252,7 @@ export function runSchedulerLoop( */ export function makeScheduleFire(options: { readonly db: Database; - readonly adapter: AgentAdapter; + readonly runtime: ManagedRuntime.ManagedRuntime; readonly maxConcurrentRuns: number; readonly workspace: WorkspaceSpec; readonly repo: RunRepo; @@ -260,11 +260,10 @@ export function makeScheduleFire(options: { }): (schedule: RuntimeSchedule) => Promise { const env = options; return async (schedule: RuntimeSchedule): Promise => { - return startTrackedRun(env.db, schedule.workflow, { + return startTrackedRun(env.runtime, env.db, schedule.workflow, { input: schedule.input, repo: env.repo, - adapter: options.adapter, - maxConcurrentRuns: options.maxConcurrentRuns, + maxConcurrentRuns: env.maxConcurrentRuns, workspace: env.workspace, dispatchEnv: env.dispatchEnv, scheduleId: schedule.id, diff --git a/src/workflow.dispatch.test.ts b/src/workflow.dispatch.test.ts index 0f0e4f0..6f885e1 100644 --- a/src/workflow.dispatch.test.ts +++ b/src/workflow.dispatch.test.ts @@ -19,6 +19,7 @@ import { describe, expect, test } from "bun:test"; import { Schema } from "effect"; import type { RunEvent } from "./events"; import { createSlowFakeAdapter } from "./replay/adapter"; +import { makeAgentRuntime } from "./runtime/agent-runtime"; import { startRun } from "./runtime/run"; import { defineWorkflow } from "./workflow"; @@ -56,11 +57,11 @@ describe("ctx.dispatch (issue #14)", () => { }, }); - const handle = startRun(parent, { + const runtime = makeAgentRuntime(SLOW_ADAPTER); + const handle = await startRun(parent, runtime, { runId: "run-parent-1", dir: root, input: {}, - adapter: SLOW_ADAPTER, dispatch: (_child, _input): Promise => { void _child; void _input; @@ -92,27 +93,26 @@ describe("ctx.dispatch (issue #14)", () => { }); let childResult: Promise<{ outcome: string }> | undefined; + const runtime = makeAgentRuntime(SLOW_ADAPTER); - const handle = startRun(parent, { + const handle = await startRun(parent, runtime, { runId: "run-parent-x", dir: root, input: {}, - adapter: SLOW_ADAPTER, - dispatch: (childWorkflow, input) => { + dispatch: async (childWorkflow, input) => { void childWorkflow.id; childEvents = []; - const childRun = startRun(childWorkflow as never, { + const childRun = await startRun(childWorkflow as never, runtime, { runId: "run-child-x", dir: root, input, - adapter: SLOW_ADAPTER, parentRunId: "run-parent-x", onEvent: (event) => { childEvents.push(event); }, }); childResult = childRun.result; - return Promise.resolve("run-child-x"); + return "run-child-x"; }, onEvent: (event) => { parentEvents.push(event); @@ -147,11 +147,11 @@ describe("ctx.dispatch (issue #14)", () => { run: async (ctx) => ctx.dispatch(numberedChild, { n: 1 }), }); - const handle = startRun(parent, { + const runtime = makeAgentRuntime(SLOW_ADAPTER); + const handle = await startRun(parent, runtime, { runId: "run-no-daemon", dir: root, input: {}, - adapter: SLOW_ADAPTER, onEvent: () => undefined, }); From 6e0b3921fe9096db3c372d6d27d8877744d8d76f Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Sun, 20 Sep 2026 11:31:39 +0000 Subject: [PATCH 3/5] fix(runtime): reconcile replay adapter signal API with runtime naming --- src/replay/adapter.test.ts | 29 +++++++++++++++++------------ src/replay/adapter.ts | 24 +++++++++++++++++++++--- src/runtime/agent-step.test.ts | 15 ++++++++------- src/runtime/run-signals.test.ts | 6 +++--- 4 files changed, 49 insertions(+), 25 deletions(-) diff --git a/src/replay/adapter.test.ts b/src/replay/adapter.test.ts index 722db85..9ec0b9f 100644 --- a/src/replay/adapter.test.ts +++ b/src/replay/adapter.test.ts @@ -1,7 +1,11 @@ import { Effect, ManagedRuntime } from "effect"; import { describe, expect, test } from "bun:test"; import { buildAgentStepEffect } from "../runtime/agent-step"; -import type { AgentAdapter, AgentAdapterOptions, AgentStreamItem } from "../runtime/agent-adapter"; +import type { + AgentAdapter, + AgentAdapterOptions, + AgentAdapterYield, +} from "../runtime/agent-adapter"; import { AgentRuntimeLayer } from "../runtime/agent-runtime"; import { createCorpusReplayAdapter, createSlowFakeAdapter, loadCorpusBlocks } from "./adapter"; @@ -10,8 +14,8 @@ function runtimeFor(adapter: AgentAdapter) { } /** Drains an adapter stream to an array (for-of over async iterables is fine; this keeps types explicit). */ -async function drain(stream: AsyncIterable): Promise> { - const items: Array = []; +async function drain(stream: AsyncIterable): Promise> { + const items: Array = []; for await (const item of stream) items.push(item); return items; } @@ -97,11 +101,11 @@ describe("createCorpusReplayAdapter", () => { // The recorded session id surfaces as a signal, and every signal is the // normalized union — never a vendor event name. expect(signals).toContainEqual({ - kind: "session", - sessionId: "ses_f64ec04acffeJ0tjsHSkjAEqZF", + _tag: "sessionId", + value: "ses_f64ec04acffeJ0tjsHSkjAEqZF", }); for (const signal of signals) { - expect(["session", "structured-output", "error"]).toContain(signal.kind); + expect(["sessionId", "structuredOutput", "runError"]).toContain(signal._tag); } // Chunks still ride through verbatim, one item per recorded chunk. @@ -127,8 +131,8 @@ describe("createCorpusReplayAdapter", () => { ); expect(items.map((item) => item.signal)).toEqual([ - { kind: "session", sessionId: "ses_fresh" }, - { kind: "structured-output", value: { ok: true } }, + { _tag: "sessionId", value: "ses_fresh" }, + { _tag: "structuredOutput", value: { ok: true } }, ]); const runtime2 = runtimeFor(adapter); @@ -178,7 +182,7 @@ describe("createSlowFakeAdapter", () => { const adapter = createSlowFakeAdapter( [{ type: "RUN_STARTED" }, { type: "TEXT_MESSAGE_START" }], 1, - [{ index: 0, signal: { kind: "session", sessionId: "ses_slow" } }], + [{ index: 0, signal: { _tag: "sessionId", value: "ses_slow" } }], ); const runtime4 = runtimeFor(adapter); @@ -204,18 +208,19 @@ function createFakeSignalAdapter(signals: { structuredOutput?: unknown; }): AgentAdapter { return { - stream(_options: AgentAdapterOptions): AsyncIterable { + async prepareWorkspace(_dir: string): Promise {}, + stream(_options: AgentAdapterOptions): AsyncIterable { return (async function* () { if (signals.sessionId !== undefined) { yield { chunk: { type: "RUN_STARTED" }, - signal: { kind: "session", sessionId: signals.sessionId }, + signal: { _tag: "sessionId", value: signals.sessionId }, }; } if (signals.structuredOutput !== undefined) { yield { chunk: { type: "RUN_FINISHED" }, - signal: { kind: "structured-output", value: signals.structuredOutput }, + signal: { _tag: "structuredOutput", value: signals.structuredOutput }, }; } })(); diff --git a/src/replay/adapter.ts b/src/replay/adapter.ts index 5130bcc..b3713ae 100644 --- a/src/replay/adapter.ts +++ b/src/replay/adapter.ts @@ -23,6 +23,7 @@ import type { AgentAdapter, AgentAdapterOptions, AgentAdapterYield, + AgentSignal, } from "../runtime/agent-adapter"; import { extractOpencodeSignal } from "../runtime/opencode-adapter"; @@ -94,17 +95,34 @@ export function createCorpusReplayAdapter(path: string): AgentAdapter { * chunk arrives and reliably interrupt mid-stream. Used by the cancellation * regression test in place of a corpus (no recorded trace can be paused on * demand; a corpus is a fixed sequence, not a controllable one). + * + * Signals can be attached declaratively by index, so a test can exercise the + * signal path without imitating any vendor chunk shape (ADR 0012 §2). */ -export function createSlowFakeAdapter(chunks: ReadonlyArray, delayMs = 20): AgentAdapter { +export interface AttachedSignal { + /** The zero-based position of the chunk this signal attaches to. */ + readonly index: number; + readonly signal: AgentSignal; +} + +export function createSlowFakeAdapter( + chunks: ReadonlyArray, + delayMs = 20, + signals: ReadonlyArray = [], +): AgentAdapter { + const byIndex = new Map(signals.map((s) => [s.index, s.signal])); return { async prepareWorkspace(_dir: string): Promise {}, stream(_options: AgentAdapterOptions): AsyncIterable { return { async *[Symbol.asyncIterator]() { - for (const chunk of chunks) { + for (let index = 0; index < chunks.length; index++) { await new Promise((resolve) => setTimeout(resolve, delayMs)); - yield { chunk }; + const signal = byIndex.get(index); + yield signal === undefined + ? { chunk: chunks[index] } + : { chunk: chunks[index], signal }; } }, }; diff --git a/src/runtime/agent-step.test.ts b/src/runtime/agent-step.test.ts index ee01f1e..b43fc3e 100644 --- a/src/runtime/agent-step.test.ts +++ b/src/runtime/agent-step.test.ts @@ -6,13 +6,14 @@ import { Effect, ManagedRuntime } from "effect"; import { describe, expect, test } from "bun:test"; -import type { AgentAdapter, AgentAdapterOptions, AgentStreamItem } from "./agent-adapter"; +import type { AgentAdapter, AgentAdapterOptions, AgentAdapterYield } from "./agent-adapter"; import { AgentRuntimeLayer } from "./agent-runtime"; import { buildAgentStepEffect } from "./agent-step"; -function scriptedAdapter(items: ReadonlyArray): AgentAdapter { +function scriptedAdapter(items: ReadonlyArray): AgentAdapter { return { - stream(_options: AgentAdapterOptions): AsyncIterable { + async prepareWorkspace(_dir: string): Promise {}, + stream(_options: AgentAdapterOptions): AsyncIterable { return (async function* () { yield* items; })(); @@ -37,16 +38,16 @@ async function runStep(adapter: AgentAdapter, options: Parameters { test("records adapter signals; chunks reach onChunk verbatim", async () => { const seen: Array = []; - const items: ReadonlyArray = [ + const items: ReadonlyArray = [ { chunk: { type: "RUN_STARTED" } }, { chunk: { type: "CUSTOM", name: "vendor.session", value: { id: "ses_x" } }, - signal: { kind: "session", sessionId: "ses_x" }, + signal: { _tag: "sessionId", value: "ses_x" }, }, ...TEXT_CHUNKS.map((chunk) => ({ chunk })), { chunk: { type: "CUSTOM", name: "vendor.output", value: { object: { answer: 42 } } }, - signal: { kind: "structured-output", value: { answer: 42 } }, + signal: { _tag: "structuredOutput", value: { answer: 42 } }, }, { chunk: { type: "RUN_FINISHED" } }, ]; @@ -72,7 +73,7 @@ describe("buildAgentStepEffect over the signal seam (ADR 0012 §2)", () => { scriptedAdapter([ { chunk: { type: "RUN_ERROR", message: "boom" }, - signal: { kind: "error", message: "boom" }, + signal: { _tag: "runError", value: "boom" }, }, ]), { diff --git a/src/runtime/run-signals.test.ts b/src/runtime/run-signals.test.ts index 53c4fcf..c66e179 100644 --- a/src/runtime/run-signals.test.ts +++ b/src/runtime/run-signals.test.ts @@ -61,7 +61,7 @@ function stepFinished(events: ReadonlyArray) { describe("signals populate AgentStepFinished as before (issue #35)", () => { test("a session signal lands on AgentStepFinished.sessionId", async () => { const { events } = await runWith(TEXT_CHUNKS, [ - { index: 0, signal: { kind: "session", sessionId: "ses_e2e" } }, + { index: 0, signal: { _tag: "sessionId", value: "ses_e2e" } }, ]); expect(stepFinished(events).sessionId).toBe("ses_e2e"); }); @@ -70,7 +70,7 @@ describe("signals populate AgentStepFinished as before (issue #35)", () => { const events: Array = []; const runtime = makeAgentRuntime( createSlowFakeAdapter([...TEXT_CHUNKS, { type: "CUSTOM", name: "anything-at-all" }], 1, [ - { index: 3, signal: { kind: "structured-output", value: { where: "from-signal" } } }, + { index: 3, signal: { _tag: "structuredOutput", value: { where: "from-signal" } } }, ]), ); const handle = await startRun(agentWorkflow(SIGNAL_SCHEMA), runtime, { @@ -103,7 +103,7 @@ describe("signals populate AgentStepFinished as before (issue #35)", () => { test("an error signal fails the step and lands on AgentStepFinished.error", async () => { const { outcome, events } = await runWith( [...TEXT_CHUNKS, { type: "RUN_FINISHED" }], - [{ index: 3, signal: { kind: "error", message: "sandbox vanished" } }], + [{ index: 3, signal: { _tag: "runError", value: "sandbox vanished" } }], ); expect(outcome.outcome).toBe("failed"); From cadfc0fc7d6745f505f214d1530b88b6f1032012 Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Sun, 20 Sep 2026 12:46:38 +0000 Subject: [PATCH 4/5] fix(runtime): preserve usage accounting and domain-error messages on composition root --- src/runtime/agent-step.test.ts | 13 +- src/runtime/agent-step.ts | 207 ++++++++++++++------------- src/runtime/agent-step.usage.test.ts | 45 +++--- src/runtime/run.ts | 2 +- src/server/runs.ts | 27 ++-- 5 files changed, 153 insertions(+), 141 deletions(-) diff --git a/src/runtime/agent-step.test.ts b/src/runtime/agent-step.test.ts index b43fc3e..57de673 100644 --- a/src/runtime/agent-step.test.ts +++ b/src/runtime/agent-step.test.ts @@ -110,18 +110,17 @@ describe("buildAgentStepEffect over the signal seam (ADR 0012 §2)", () => { test("the service is resolved from context, not passed as an option", async () => { const fake = scriptedAdapter([{ chunk: { type: "RUN_FINISHED" } }]); - const program = Effect.gen(function* () { - const handle = yield* buildAgentStepEffect({ + const runtime = ManagedRuntime.make(AgentRuntimeLayer(fake)); + const handle = await runtime.runPromise( + buildAgentStepEffect({ threadId: "t", dir: "/tmp", model: "m", prompt: "p", onChunk: () => {}, - }); - return yield* handle.effect; - }); - const runtime = ManagedRuntime.make(AgentRuntimeLayer(fake)); - const outcome = await runtime.runPromise(program); + }), + ); + const outcome = await Effect.runPromise(handle.effect); expect(outcome.chunkCount).toBe(1); await runtime.dispose(); }); diff --git a/src/runtime/agent-step.ts b/src/runtime/agent-step.ts index e597aba..344975e 100644 --- a/src/runtime/agent-step.ts +++ b/src/runtime/agent-step.ts @@ -5,10 +5,9 @@ * §5, D17): the boundary wiring is unchanged, only the chunk source and the * bookkeeping surface (now `ctx.agent`'s granular result, ADR 0002 §2) moved. * - * ADR 0012 §2: the runtime consumes `AgentSignal`s from the adapter and never - * string-matches vendor event names. AG-UI standard types (`TEXT_MESSAGE_*`) - * are still interpreted here for `finalText` accumulation — these are part - * of the open AG-UI protocol, not vendor-specific. + * ADR 0012 §2: chunks are opaque here — the adapter interprets its own stream + * and yields each chunk alongside a normalized `AgentSignal`. This module + * records signals and forwards chunks; it never matches a vendor event name. * * WHY THE EXPLICIT `abortController.abort()` IS NEEDED (0a-1/0a-2 findings): * closing the IO stream does not terminate the opencode process; only an @@ -23,7 +22,8 @@ import { Effect, Schema, Stream } from "effect"; import type { AgentStepUsage } from "../events"; -import type { AgentAdapter, AgentAdapterYield } from "./agent-adapter"; +import type { AgentAdapterYield } from "./agent-adapter"; +import { AgentRuntime } from "./agent-runtime"; /** * Pull the four token counts out of a `RUN_FINISHED.usage` object. @@ -60,7 +60,6 @@ export interface AgentStepEffectOptions { readonly model: string; readonly prompt: string; readonly outputSchema?: unknown; - readonly adapter: AgentAdapter; /** Fired synchronously per chunk, before any bookkeeping — the runtime's `AgentChunk` emission point. */ readonly onChunk: (chunk: unknown) => void; } @@ -94,103 +93,109 @@ export interface AgentStepHandle { readonly partial: AgentStepPartial; } -export function buildAgentStepEffect(options: AgentStepEffectOptions): AgentStepHandle { - const abortController = new AbortController(); - - const iterable = options.adapter.stream({ - threadId: options.threadId, - dir: options.dir, - model: options.model, - prompt: options.prompt, - outputSchema: options.outputSchema, - abortController, - }); +export function buildAgentStepEffect( + options: AgentStepEffectOptions, +): Effect.Effect { + return Effect.gen(function* () { + const { adapter } = yield* AgentRuntime; + + const abortController = new AbortController(); + + const iterable = adapter.stream({ + threadId: options.threadId, + dir: options.dir, + model: options.model, + prompt: options.prompt, + outputSchema: options.outputSchema, + abortController, + }); + + const rawStream = Stream.fromAsyncIterable( + iterable as AsyncIterable, + (cause) => new AgentStepChunkError({ cause }), + ); + + const partial: AgentStepPartial = { + chunkCount: 0, + finalText: "", + sessionId: undefined, + usage: undefined, + }; + let currentMessageBuffer: string | undefined; + let structuredOutput: unknown; + let runError: string | undefined; + const startedAt = Date.now(); + + // Plain closure mutation (not a `Ref`) is fine: this Effect never runs + // concurrently with itself, and the callback always runs on the same + // single-threaded event loop turn (mirrors the spike's finding exactly). + const processed = Stream.mapEffect(rawStream, (yieldItem: AgentAdapterYield) => + Effect.sync(() => { + const chunk = yieldItem.chunk; + partial.chunkCount += 1; + options.onChunk(chunk); + + if (yieldItem.signal !== undefined) { + switch (yieldItem.signal._tag) { + case "sessionId": + partial.sessionId = yieldItem.signal.value; + break; + case "structuredOutput": + structuredOutput = yieldItem.signal.value; + break; + case "runError": + runError = yieldItem.signal.value; + break; + } + } - const rawStream = Stream.fromAsyncIterable( - iterable, - (cause) => new AgentStepChunkError({ cause }), - ); + // TEXT_MESSAGE_* folding stays runtime-side: `finalText` is an AG-UI + // concept (ADR 0003 §2), not a vendor event name — it feeds tier 2 of + // structured-output resolution and the cancelled-step partial. + const record = chunk as { type?: unknown; delta?: unknown; usage?: unknown }; - const partial: AgentStepPartial = { - chunkCount: 0, - finalText: "", - sessionId: undefined, - usage: undefined, - }; - let currentMessageBuffer: string | undefined; - let structuredOutput: unknown; - let runError: string | undefined; - const startedAt = Date.now(); - - // Plain closure mutation (not a `Ref`) is fine: this Effect never runs - // concurrently with itself, and the callback always runs on the same - // single-threaded event loop turn (mirrors the spike's finding exactly). - const processed = Stream.mapEffect(rawStream, (yieldItem: AgentAdapterYield) => - Effect.sync(() => { - partial.chunkCount += 1; - options.onChunk(yieldItem.chunk); - - if (yieldItem.signal !== undefined) { - switch (yieldItem.signal._tag) { - case "sessionId": - partial.sessionId = yieldItem.signal.value; - break; - case "structuredOutput": - structuredOutput = yieldItem.signal.value; - break; - case "runError": - runError = yieldItem.signal.value; - break; + if (record.type === "RUN_FINISHED") { + partial.usage = readUsage(record.usage); } - } - - const record = yieldItem.chunk as { - type?: unknown; - delta?: unknown; - usage?: unknown; - }; - - if (record.type === "RUN_FINISHED") { - partial.usage = readUsage(record.usage); - } - - if (record.type === "TEXT_MESSAGE_START") { - currentMessageBuffer = ""; - } else if (record.type === "TEXT_MESSAGE_CONTENT") { - const delta = record.delta; - if (typeof delta === "string") { - currentMessageBuffer = (currentMessageBuffer ?? "") + delta; - } - } else if (record.type === "TEXT_MESSAGE_END") { - if (currentMessageBuffer !== undefined) { - partial.finalText = currentMessageBuffer; + + if (record.type === "TEXT_MESSAGE_START") { + currentMessageBuffer = ""; + } else if (record.type === "TEXT_MESSAGE_CONTENT") { + const delta = record.delta; + if (typeof delta === "string") { + currentMessageBuffer = (currentMessageBuffer ?? "") + delta; + } + } else if (record.type === "TEXT_MESSAGE_END") { + if (currentMessageBuffer !== undefined) { + partial.finalText = currentMessageBuffer; + } + currentMessageBuffer = undefined; } - currentMessageBuffer = undefined; - } - - return yieldItem; - }), - ); - - const drain = Stream.runDrain(processed); - - // `Effect.onInterrupt`'s finalizer runs ONLY if `drain` is interrupted, not - // on normal success/failure — deliberately not `Effect.ensuring`. - const guarded = Effect.onInterrupt(drain, () => - Effect.sync(() => { - abortController.abort(); - }), - ); - - const effect = Effect.map(guarded, () => ({ - chunkCount: partial.chunkCount, - finalText: partial.finalText, - structuredOutput, - sessionId: partial.sessionId, - usage: partial.usage, - runError, - durationMs: Date.now() - startedAt, - })); - - return { effect, abortController, partial }; + + return chunk; + }), + ); + + const drain = Stream.runDrain(processed); + + // `Effect.onInterrupt`'s finalizer runs ONLY if `drain` is interrupted, not + // on normal success/failure — deliberately not `Effect.ensuring`. + const guarded = Effect.onInterrupt(drain, () => + Effect.sync(() => { + abortController.abort(); + }), + ); + + const effect = Effect.map(guarded, () => ({ + chunkCount: partial.chunkCount, + finalText: partial.finalText, + structuredOutput, + sessionId: partial.sessionId, + usage: partial.usage, + runError, + durationMs: Date.now() - startedAt, + })); + + return { effect, abortController, partial }; + }); } diff --git a/src/runtime/agent-step.usage.test.ts b/src/runtime/agent-step.usage.test.ts index 0031220..60a0671 100644 --- a/src/runtime/agent-step.usage.test.ts +++ b/src/runtime/agent-step.usage.test.ts @@ -13,31 +13,40 @@ */ import { describe, expect, test } from "bun:test"; -import { Effect } from "effect"; +import { Effect, ManagedRuntime } from "effect"; import { agentStepContextTokens } from "../events"; -import type { AgentAdapterYield } from "./agent-adapter"; +import type { AgentAdapter, AgentAdapterYield } from "./agent-adapter"; +import { AgentRuntimeLayer } from "./agent-runtime"; import { loadCorpusBlocks } from "../replay/adapter"; import { buildAgentStepEffect } from "./agent-step"; const CORPUS = "test/corpus/run-1789308170212.ndjson"; -async function runBlock(chunks: ReadonlyArray) { - const handle = buildAgentStepEffect({ - threadId: "thread", - dir: ".", - model: "model", - prompt: "prompt", - adapter: { - async prepareWorkspace(_dir: string): Promise {}, - async *stream(): AsyncGenerator { - for (const chunk of chunks) { - yield { chunk }; - } - }, +function corpusAdapter(chunks: ReadonlyArray): AgentAdapter { + return { + async prepareWorkspace(_dir: string): Promise {}, + async *stream(): AsyncGenerator { + for (const chunk of chunks) { + yield { chunk }; + } }, - onChunk: () => {}, - }); - return await Effect.runPromise(handle.effect); + }; +} + +async function runBlock(chunks: ReadonlyArray) { + const runtime = ManagedRuntime.make(AgentRuntimeLayer(corpusAdapter(chunks))); + const handle = await runtime.runPromise( + buildAgentStepEffect({ + threadId: "thread", + dir: ".", + model: "model", + prompt: "prompt", + onChunk: () => {}, + }), + ); + const outcome = await Effect.runPromise(handle.effect); + await runtime.dispose(); + return outcome; } describe("agent step usage", () => { diff --git a/src/runtime/run.ts b/src/runtime/run.ts index 03eb1cc..ddb0729 100644 --- a/src/runtime/run.ts +++ b/src/runtime/run.ts @@ -428,7 +428,7 @@ export async function startRun( emit({ _tag: "RunCancelled", durationMs }); return { outcome: "cancelled" }; } - const message = err instanceof Error ? err.message : String(err); + const message = domainErrorMessage(err); const stack = err instanceof Error ? err.stack : undefined; emit({ _tag: "RunFailed", diff --git a/src/server/runs.ts b/src/server/runs.ts index 41265f7..c7e4a0c 100644 --- a/src/server/runs.ts +++ b/src/server/runs.ts @@ -5,6 +5,7 @@ */ import type { Database } from "bun:sqlite"; +import { Schema } from "effect"; import type { ManagedRuntime } from "effect"; import { rm } from "node:fs/promises"; import { admitRun } from "./admission"; @@ -27,12 +28,9 @@ const active = new Map | ReservedSlot>(); export const DEFAULT_MAX_DISPATCH_DEPTH = 5; export const DEFAULT_MAX_CHILDREN_PER_RUN = 20; -export class DispatchCapError extends Error { - constructor(message: string) { - super(message); - this.name = "DispatchCapError"; - } -} +export class DispatchCapError extends Schema.TaggedError()("DispatchCapError", { + message: Schema.String, +}) {} function isReserved(entry: RunHandle | ReservedSlot | undefined): boolean { return entry !== undefined && !("result" in entry) && "cancelled" in entry; @@ -156,21 +154,22 @@ async function dispatchChildRun( env.maxConcurrentRuns !== undefined && !admitRun(env.maxConcurrentRuns, activeRunIds().length) ) { - throw new ConcurrencyLimitError(env.maxConcurrentRuns); + throw new ConcurrencyLimitError({ maxConcurrentRuns: env.maxConcurrentRuns }); } const depth = dispatchDepth(db, parentRunId); if (depth + 1 > maxDepth) { - throw new DispatchCapError( - `dispatch depth exceeded: run ${parentRunId} is nested ${depth} levels deep; ` + + throw new DispatchCapError({ + message: + `dispatch depth exceeded: run ${parentRunId} is nested ${depth} levels deep; ` + `max ${maxDepth} (a workflow that dispatches itself must not fill the daemon)`, - ); + }); } const childCount = countDispatchedChildren(db, parentRunId); if (childCount >= maxChildren) { - throw new DispatchCapError( - `dispatch child cap exceeded: run ${parentRunId} already dispatched ${childCount} children; max ${maxChildren}`, - ); + throw new DispatchCapError({ + message: `dispatch child cap exceeded: run ${parentRunId} already dispatched ${childCount} children; max ${maxChildren}`, + }); } const childRunId = `run-${crypto.randomUUID()}`; @@ -225,7 +224,7 @@ export async function startTrackedRun( } if (existing === undefined && options.maxConcurrentRuns !== undefined) { if (!admitRun(options.maxConcurrentRuns, active.size)) { - throw new ConcurrencyLimitError(options.maxConcurrentRuns); + throw new ConcurrencyLimitError({ maxConcurrentRuns: options.maxConcurrentRuns }); } active.set(runId, { cancelled: false }); } From 0cb44c49ae2029aa88bc4ec1b108799bf9b6812b Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Wed, 23 Sep 2026 06:49:40 +0000 Subject: [PATCH 5/5] fix(cli): restore the effect/unstable/cli entrypoint after a bad merge MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Commit a507d1e's conflict resolution against 33-parse-cli-with-effect-cli silently reverted PR #52: src/cli.ts's import.meta.main block regressed to the pre-#52 hand-rolled USAGE/parseFlags/usageError/parseArgs parser, even though src/cli-commands.ts's factoryCommand (the effect/unstable/cli command tree) and its tests kept passing in isolation — so CI stayed green while the shipped binary silently lost generated help, typed flag validation, and --wizard/--completions. Restore the base branch's entrypoint (import { factoryCommand } from "./cli-commands"; Command.run(factoryCommand, ...)) while keeping this PR's actual new work intact: the ManagedRuntime/AgentRuntimeLayer composition root in runCli, which resolves the adapter from factory.config.ts (falling back to opencodeAdapter) and threads a ManagedRuntime into startRun. Also restore `prepareWorkspace: options.clone !== undefined`, which the same bad merge had silently dropped from runCli's startRun call. cli-commands.ts's runCommand was hardcoding `adapter: opencodeAdapter` on every `factory run` invocation, which bypassed runCli's config-driven adapter resolution and defeated issue #36's "runtime selectable from factory.config.ts" criterion for the direct-run path. Drop that override so runCli's own fallback (options.adapter ?? config.agent.adapter ?? opencodeAdapter) decides. Add tests that exercise src/cli.ts's actual import.meta.main entrypoint (the same path bin/factory.js runs in production), not just factoryCommand in isolation: a static check that the file contains no hand-rolled parser, and spawned-process checks that --help renders effect/unstable/cli's generated help and that `serve --port abc` is rejected by the typed Int flag. Without these, this class of regression can pass CI again undetected. Co-Authored-By: Claude Sonnet 5 --- src/cli-argv.test.ts | 54 ++++++++++++++++++++++++++++++++++++++++++++ src/cli-commands.ts | 2 -- src/cli.ts | 9 +++----- 3 files changed, 57 insertions(+), 8 deletions(-) diff --git a/src/cli-argv.test.ts b/src/cli-argv.test.ts index 1a5594f..ed987c7 100644 --- a/src/cli-argv.test.ts +++ b/src/cli-argv.test.ts @@ -424,3 +424,57 @@ describe("factory binary: process exit codes", () => { expect(exitCode).not.toBe(0); }); }); + +/** + * PR #56 merged this branch onto `33-parse-cli-with-effect-cli` with a bad + * conflict resolution: `src/cli.ts`'s `import.meta.main` block regressed to + * the pre-#52 hand-rolled `USAGE`/`parseFlags`/`usageError` parser while + * `factoryCommand` above kept passing — the tests here only ever drove + * `factoryCommand` directly, never the actual binary entrypoint, so CI stayed + * green while the shipped CLI silently lost `effect/unstable/cli` (generated + * help, typed flag validation, `--wizard`/`--completions`, …). These tests + * exercise `src/cli.ts` itself — via `import.meta.main`, the same path + * `bin/factory.js` runs in production — so that regression can't recur + * unnoticed. + */ +describe("cli.ts entrypoint wiring", () => { + test("cli.ts contains no hand-rolled argv parser", async () => { + const source = await Bun.file(CLI).text(); + expect(source).toContain('import { factoryCommand } from "./cli-commands"'); + expect(source).toMatch(/Command\.run\(factoryCommand/); + expect(source).not.toMatch(/\bconst USAGE\b/); + expect(source).not.toMatch(/\bfunction usageError\b/); + expect(source).not.toMatch(/\bfunction parseFlags\b/); + expect(source).not.toMatch(/\bfunction parseArgs\b/); + expect(source).not.toMatch(/\bfunction parseStartArgs\b/); + expect(source).not.toMatch(/\bfunction parseServeArgs\b/); + }); + + test("--help at the real entrypoint is generated by effect/unstable/cli, not a hand-rolled banner", async () => { + const proc = Bun.spawn(["bun", CLI, "--help"], { stdout: "pipe", stderr: "pipe" }); + const [exitCode, stdout] = await Promise.all([proc.exited, new Response(proc.stdout).text()]); + expect(exitCode).toBe(0); + // effect/unstable/cli's generated help renders these section headings; + // the hand-rolled USAGE banner (a lowercase "usage:" line) never did. + expect(stdout).toContain("SUBCOMMANDS"); + expect(stdout).toContain("GLOBAL FLAGS"); + expect(stdout).not.toContain("usage:\n"); + }); + + test("factory serve --port at the real entrypoint rejects a non-numeric value before starting the daemon", async () => { + // The hand-rolled parser did `Number(portRaw)` with no validation at all — + // a bad --port silently became NaN. effect/unstable/cli's typed Int flag + // rejects it up front. + const proc = Bun.spawn(["bun", CLI, "serve", "--port", "abc"], { + stdout: "pipe", + stderr: "pipe", + }); + const [exitCode, stdout, stderr] = await Promise.all([ + proc.exited, + new Response(proc.stdout).text(), + new Response(proc.stderr).text(), + ]); + expect(exitCode).not.toBe(0); + expect(stdout + stderr).toContain("Invalid value for flag --port"); + }); +}); diff --git a/src/cli-commands.ts b/src/cli-commands.ts index d49d176..030eaf2 100644 --- a/src/cli-commands.ts +++ b/src/cli-commands.ts @@ -2,7 +2,6 @@ import { Effect, Option } from "effect"; import { Argument, Command, Flag } from "effect/unstable/cli"; import { findFactoryConfig, loadFactoryConfig } from "./config"; import { initCli } from "./init"; -import { opencodeAdapter } from "./runtime/opencode-adapter"; import { resolve } from "node:path"; import { runCli, @@ -222,7 +221,6 @@ export const runCommand = Command.make( () => `.factory/runs/run-${Date.now()}/events.ndjson`, ), dbPath: config.db, - adapter: opencodeAdapter, }; const exitCode = yield* Effect.promise(() => runCli(options)); diff --git a/src/cli.ts b/src/cli.ts index 29292b5..33c5795 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -20,8 +20,7 @@ import { Effect, FileSystem, Layer, ManagedRuntime, Path, Stdio, Terminal } from import { ChildProcessSpawner } from "effect/unstable/process"; import { CliError, Command } from "effect/unstable/cli"; import type { RunEvent } from "./events"; -import { findFactoryConfig, loadFactoryConfig } from "./config"; -import { initCli } from "./init"; +import { loadFactoryConfig } from "./config"; import type { RunRepo } from "./runtime/run"; import { resetClone, type GitIdentity } from "./lib/clone"; import { loadWorkflow } from "./lib/load-workflow"; @@ -31,10 +30,7 @@ import type { AgentAdapter } from "./runtime/agent-adapter"; import { opencodeAdapter } from "./runtime/opencode-adapter"; import { AgentRuntimeLayer } from "./runtime/agent-runtime"; import { startRun } from "./runtime/run"; -import { startDaemon, type DaemonOptions } from "./server/daemon"; - -const DEFAULT_DB_PATH = ".factory/factory.db"; -const DEFAULT_DAEMON_URL = "http://localhost:3000"; +import { factoryCommand } from "./cli-commands"; export interface CliOptions { readonly workflowPath: string; @@ -83,6 +79,7 @@ export async function runCli(options: CliOptions): Promise { runId, dir: options.dir, input: options.input, + prepareWorkspace: options.clone !== undefined, ...(repo !== undefined ? { repo } : {}), onEvent: (event) => { console.log(formatEvent(event));