Skip to content
22 changes: 13 additions & 9 deletions e2e/implement-issue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down Expand Up @@ -118,16 +119,19 @@ describe("implement-issue workflow, replayed against the recorded round-trip cor
process.env.PATH = `${binDir}:${originalPath}`;

const events: Array<unknown> = [];
const handle = startRun(implementIssue, {
runId: "test-run-corpus-replay",
dir: workDir,
input: {
issueNumber: 1,
const handle = startRun(
implementIssue,
makeAgentRuntime(createCorpusReplayAdapter(FULL_ROUND_TRIP_CORPUS)),
{
runId: "test-run-corpus-replay",
dir: workDir,
input: {
issueNumber: 1,
},
repo: { slug: "local/fixture", baseBranch: "main" },
onEvent: (event) => events.push(event),
},
repo: { slug: "local/fixture", baseBranch: "main" },
adapter: createCorpusReplayAdapter(FULL_ROUND_TRIP_CORPUS),
onEvent: (event) => events.push(event),
});
);

const outcome = await handle.result;

Expand Down
11 changes: 5 additions & 6 deletions e2e/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -137,12 +138,11 @@ async function seedCorpusRuns(root: string, db: ReturnType<typeof openStore>): P
process.env.PATH = `${binDir}:${originalPath}`;

await awaitRun(
startRun(implementIssue, {
startRun(implementIssue, makeAgentRuntime(createCorpusReplayAdapter(CORPUS_ROUND_TRIP)), {
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),
}),
);
Expand All @@ -151,11 +151,10 @@ async function seedCorpusRuns(root: string, db: ReturnType<typeof openStore>): P

// A second, cheap run so list ordering has more than one data point.
await awaitRun(
startRun(echoWorkflow, {
startRun(echoWorkflow, makeAgentRuntime(createCorpusReplayAdapter(CORPUS_ONE_STEP)), {
runId: "run-static-echo",
dir: workDir,
input: {},
adapter: createCorpusReplayAdapter(CORPUS_ONE_STEP),
onEvent: (event) => appendEvent(db, event),
}),
);
Expand Down Expand Up @@ -253,7 +252,6 @@ async function main(): Promise<void> {
db.close();
}

const adapter = createSlowFakeAdapter(SLOW_CHUNKS, 1_000);
const config = defineConfig({
repo: {
sshUrl: remoteDir,
Expand All @@ -278,8 +276,9 @@ async function main(): Promise<void> {
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}`);
}

Expand Down
54 changes: 54 additions & 0 deletions src/cli-argv.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -419,3 +419,57 @@ describe("factory binary: process exit codes", () => {
});
});
});

/**
* 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");
});
});
11 changes: 9 additions & 2 deletions src/cli-commands.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -84,6 +83,15 @@ export const serveCommand = Command.make("serve", serveConfig, (config) =>
: "scheduler: none (no schedules in config)",
),
);
// The daemon's shutdown path: stop the scheduler and server and dispose
// the agent runtime before exiting, instead of dying with them live.
yield* Effect.sync(() => {
const shutdown = (signal: NodeJS.Signals) => {
void handle.stop().finally(() => process.exit(signal === "SIGINT" ? 130 : 143));
};
process.once("SIGINT", shutdown);
process.once("SIGTERM", shutdown);
});
}),
).pipe(Command.withDescription("Run the daemon: HTTP API, live event stream, and the web UI"));

Expand Down Expand Up @@ -216,7 +224,6 @@ export const runCommand = Command.make("run", runConfig, (config) =>
clone,
outPath: Option.getOrElse(config.out, () => `.factory/runs/run-${Date.now()}/events.ndjson`),
dbPath: config.db,
adapter: opencodeAdapter,
};

const exitCode = yield* Effect.promise(() => runCli(options));
Expand Down
22 changes: 11 additions & 11 deletions src/cli.start.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,15 +25,17 @@ async function startTestDaemon(delayMs: number, root: string): Promise<TestDaemo
const daemon = await startDaemon({
dbPath,
port: 0,
adapter: createSlowFakeAdapter(
[
{ type: "TEXT_MESSAGE_START" },
{ type: "TEXT_MESSAGE_CONTENT", delta: "hi" },
{ type: "TEXT_MESSAGE_END" },
],
delayMs,
),
config: defineConfig({
agent: {
adapter: createSlowFakeAdapter(
[
{ type: "TEXT_MESSAGE_START" },
{ type: "TEXT_MESSAGE_CONTENT", delta: "hi" },
{ type: "TEXT_MESSAGE_END" },
],
delayMs,
),
},
repo: {
sshUrl: seed,
identity: { name: "Factory", email: "factory@factory.test" },
Expand All @@ -49,9 +51,7 @@ async function startTestDaemon(delayMs: number, root: string): Promise<TestDaemo
return {
base: `http://localhost:${daemon.server.port}`,
dbPath,
stop: async () => {
await daemon.server.stop(true);
},
stop: daemon.stop,
};
}

Expand Down
16 changes: 13 additions & 3 deletions src/cli.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import { loadWorkflow } from "./lib/load-workflow";
import { streamSse } from "./lib/sse-client";
import { appendEvent, getRunEvents, listRuns, openStore } from "./persistence/store";
import type { AgentAdapter } from "./runtime/agent-adapter";
import { makeAgentRuntime } from "./runtime/agent-runtime";
import { startRun } from "./runtime/run";
import { factoryCommand } from "./cli-commands";

Expand All @@ -38,7 +39,11 @@ export interface CliOptions {
readonly clone?: { readonly sshUrl: string; readonly identity: GitIdentity };
readonly outPath: string;
readonly dbPath: string;
readonly adapter: AgentAdapter;
/**
* Injectable agent adapter for tests. When omitted, `runCli` falls back to
* `factory.config.ts`'s `agent.adapter` (or the live opencode adapter).
*/
readonly adapter?: AgentAdapter;
}

function formatEvent(event: RunEvent): string {
Expand All @@ -61,17 +66,21 @@ export async function runCli(options: CliOptions): Promise<number> {

const runId = `run-${Date.now()}`;
let repo: RunRepo | undefined;
let adapter = options.adapter;
try {
const config = await loadFactoryConfig();
repo = { slug: config.repo.slug, baseBranch: config.repo.baseBranch };
adapter ??= config.agent.adapter;
} catch {
repo = undefined;
}
const handle = startRun(workflow, {
// The direct-run path's composition root (issue #36): one runtime for this
// one run, disposed once the run settles.
const runtime = makeAgentRuntime(adapter);
const handle = startRun(workflow, runtime, {
runId,
dir: options.dir,
input: options.input,
adapter: options.adapter,
prepareWorkspace: options.clone !== undefined,
...(repo !== undefined ? { repo } : {}),
onEvent: (event) => {
Expand All @@ -89,6 +98,7 @@ export async function runCli(options: CliOptions): Promise<number> {

const outcome = await handle.result;
process.off("SIGINT", onSigint);
await runtime.dispose();
await sink.end();
db.close();

Expand Down
13 changes: 13 additions & 0 deletions src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import { Cron, Result, SchemaParser } from "effect";
import type { GitIdentity } from "./lib/clone";
import type { WorkflowDefinition } from "./workflow";
import type { AgentAdapter } from "./runtime/agent-adapter";

export const DEFAULT_WORKSPACE_ROOT = ".factory/workspaces";
export const DEFAULT_MAX_CONCURRENT_RUNS = 3;
Expand Down Expand Up @@ -69,6 +70,12 @@ export interface FactoryConfig {
*/
readonly maxDispatchDepth: number;
readonly maxChildrenPerRun: number;
/**
* Issue #36: the agent runtime. `adapter` unset means the runtime's default
* (the live opencode adapter, `runtime/agent-runtime.ts`), so existing
* configs keep working unchanged.
*/
readonly agent: { readonly adapter?: AgentAdapter };
}

/**
Expand Down Expand Up @@ -169,6 +176,11 @@ export interface FactoryConfigInput {
readonly maxDispatchDepth?: number;
readonly maxChildrenPerRun?: number;
readonly schedules?: ReadonlyArray<ScheduleConfigInput | ScheduleDefinition<any>>;
/**
* Issue #36: the agent runtime. Optional — when omitted the runtime uses
* the live opencode adapter.
*/
readonly agent?: { readonly adapter?: AgentAdapter };
}

export function defineConfig(config: FactoryConfigInput): FactoryConfig {
Expand Down Expand Up @@ -216,6 +228,7 @@ export function defineConfig(config: FactoryConfigInput): FactoryConfig {
retainedWorkspaces,
maxDispatchDepth,
maxChildrenPerRun,
agent: config.agent?.adapter !== undefined ? { adapter: config.agent.adapter } : {},
};
}

Expand Down
11 changes: 7 additions & 4 deletions src/replay/adapter.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { Effect } from "effect";
import { describe, expect, test } from "bun:test";
import { AgentRuntimeLayer } from "../runtime/agent-runtime";
import { buildAgentStepEffect } from "../runtime/agent-step";
import type { AgentAdapterYield } from "../runtime/agent-adapter";
import { createCorpusReplayAdapter, createSlowFakeAdapter, loadCorpusBlocks } from "./adapter";
Expand Down Expand Up @@ -54,11 +55,12 @@ describe("createCorpusReplayAdapter", () => {
dir: "/tmp",
model: "opencode-go/deepseek-v4.1-flash",
prompt: "irrelevant, replay ignores it",
adapter,
onChunk: (chunk) => chunks.push(chunk),
});

const outcome = await Effect.runPromise(handle.effect);
const outcome = await Effect.runPromise(
Effect.provide(handle.effect, AgentRuntimeLayer(adapter)),
);

expect(outcome.chunkCount).toBe(39);
expect(chunks.length).toBe(39);
Expand Down Expand Up @@ -107,11 +109,12 @@ describe("createSlowFakeAdapter", () => {
dir: "/tmp",
model: "m",
prompt: "p",
adapter,
onChunk: () => {},
});

const outcome = await Effect.runPromise(handle.effect);
const outcome = await Effect.runPromise(
Effect.provide(handle.effect, AgentRuntimeLayer(adapter)),
);
expect(outcome.chunkCount).toBe(3);
expect(outcome.finalText).toBe("hi");
});
Expand Down
24 changes: 21 additions & 3 deletions src/replay/adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import type {
AgentAdapter,
AgentAdapterOptions,
AgentAdapterYield,
AgentSignal,
} from "../runtime/agent-adapter";
import { extractOpencodeSignal } from "../runtime/opencode-adapter";

Expand Down Expand Up @@ -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<unknown>, 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<unknown>,
delayMs = 20,
signals: ReadonlyArray<AttachedSignal> = [],
): AgentAdapter {
const byIndex = new Map(signals.map((s) => [s.index, s.signal]));
return {
async prepareWorkspace(_dir: string): Promise<void> {},

stream(_options: AgentAdapterOptions): AsyncIterable<AgentAdapterYield> {
return {
async *[Symbol.asyncIterator]() {
for (const chunk of chunks) {
for (let index = 0; index < chunks.length; index++) {
await new Promise<void>((resolve) => setTimeout(resolve, delayMs));
yield { chunk };
const signal = byIndex.get(index);
yield signal === undefined
? { chunk: chunks[index] }
: { chunk: chunks[index], signal };
}
},
};
Expand Down
Loading
Loading