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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 7 additions & 5 deletions src/core/dev/supervisor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,17 +15,19 @@ function runtime(name: string, build: ProjectRuntime["build"] = "CodeZip"): Proj
} as ProjectRuntime;
}

/** A runner that emits `events`, then stays alive until its signal aborts (like a real server). */
/** A runner that emits `events`, stays alive until its signal aborts, then rejects with the abort reason (like the real process runner). */
function serverRunner(events: DevEvent[] = []) {
const inputs: DevServerInput[] = [];
const runner: DevRunner = {
run: async function* (input) {
inputs.push(input);
yield* events;
if (input.signal.aborted) return;
await new Promise<void>((resolve) =>
input.signal.addEventListener("abort", () => resolve(), { once: true }),
);
if (!input.signal.aborted) {
await new Promise<void>((resolve) =>
input.signal.addEventListener("abort", () => resolve(), { once: true }),
);
}
throw input.signal.reason;
},
};
return { runner, inputs };
Expand Down
9 changes: 9 additions & 0 deletions src/core/dev/supervisor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -258,6 +258,15 @@ export class DevSupervisor {
this.push(name, { type: "status", message: `Agent '${name}' stopped.` });
}
} catch (error) {
/** The runner rejects with the abort reason on teardown, which is a stop, not a crash. */
if (input.signal.aborted) {
if (entry.phase === "running") {
entry.phase = "idle";
entry.port = undefined;
this.push(name, { type: "status", message: `Agent '${name}' stopped.` });
}
return;
}
const message = error instanceof Error ? error.message : String(error);
entry.error = message;
if (entry.phase === "running") {
Expand Down
100 changes: 83 additions & 17 deletions src/handlers/project/dev/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import type { ProjectRuntime } from "../../../projectSchemas/runtime";
import {
InputValidationError,
ResourceNotFoundError,
SilentCLIError,
UserCancellationError,
} from "../../../errors";
import type { HttpRequestHandler, PortChecker } from "../../../io";
Expand Down Expand Up @@ -45,6 +46,24 @@ function captureRunner(events: DevEvent[] = []) {
return { runner, inputs };
}

/** A runner that emits `events` then stays alive until aborted, rejecting with the abort reason like the real process runner. */
function stayingRunner(events: DevEvent[] = []) {
const inputs: DevServerInput[] = [];
const runner: DevRunner = {
run: async function* (input) {
inputs.push(input);
yield* events;
if (!input.signal.aborted) {
await new Promise<void>((resolve) =>
input.signal.addEventListener("abort", () => resolve(), { once: true }),
);
}
throw input.signal.reason;
},
};
return { runner, inputs };
}

function fakeCollector() {
const starts: Parameters<DevProjectHandlerConfig["startTraceCollector"]>[0][] = [];
const state = { closed: 0 };
Expand Down Expand Up @@ -120,6 +139,9 @@ function harness(options: HarnessOptions = {}) {
resolve: async () =>
options.reloadedRuntimes ? project(...options.reloadedRuntimes) : undefined,
},
waitReady: async () => {
await Bun.sleep(5);
},
});
const ctx = ValueContext.EmptyContext()
.withValue(ProjectKey, options.project ?? project(runtime()))
Expand Down Expand Up @@ -167,16 +189,10 @@ async function inspectorStatus(subject: ReturnType<typeof harness>): Promise<{ n
describe("project dev selection and dispatch", () => {
test.each([
[project(), {}, "This project has no runtimes", InputValidationError],
[
project(runtime("orders")),
{},
"--mode headless runs a single agent in the terminal. Pass --agent <name> to choose which one. Available: orders",
InputValidationError,
],
[
project(runtime("orders"), runtime("support", "Container")),
{},
"--mode headless runs a single agent in the terminal. Pass --agent <name> to choose which one. Available: orders, support",
{ port: 4567 },
"--port applies to a single runtime. Use --agent to select one.",
InputValidationError,
],
[
Expand Down Expand Up @@ -237,6 +253,65 @@ describe("project dev selection and dispatch", () => {
});
});

describe("project dev headless multi-agent", () => {
const twoRuntimes = () => project(runtime("orders"), runtime("support", "Container"));

/** Start a headless multi-agent run and give its agents time to reach "running". */
async function supervised(subject: ReturnType<typeof harness>) {
const pending = subject.run();
pending.catch(() => undefined);
await Bun.sleep(30);
return { pending };
}

test("supervises every runtime with attributed output and per-runtime env", async () => {
const codeZip = stayingRunner([{ type: "stdout", line: "orders says hi" }]);
const container = stayingRunner();
const subject = harness({ project: twoRuntimes(), codeZip, container });
const { pending } = await supervised(subject);

expect(codeZip.inputs).toHaveLength(1);
expect(container.inputs).toHaveLength(1);
expect(codeZip.inputs[0]!.env).toMatchObject({
OTEL_EXPORTER_OTLP_ENDPOINT: "http://127.0.0.1:43180",
OTEL_SERVICE_NAME: "orders",
});
expect(container.inputs[0]!.env).toMatchObject({
OTEL_EXPORTER_OTLP_ENDPOINT: "http://host.docker.internal:43180",
OTEL_SERVICE_NAME: "support",
});
expect(subject.io.stdout()).toContain("[orders] orders says hi");
expect(subject.io.stderr()).toContain("Agent 'orders' is running on port");

process.emit("SIGINT", "SIGINT");
await expect(pending).rejects.toMatchObject({ exitCode: 130 });
expect(subject.io.stderr()).not.toContain("crashed");
expect(subject.collector.state.closed).toBe(1);
});

test("one agent failing to start leaves the others running", async () => {
const subject = harness({
project: twoRuntimes(),
codeZip: captureRunner([{ type: "status", message: "dying" }]),
container: stayingRunner(),
});
const { pending } = await supervised(subject);

expect(subject.io.stderr()).toContain("[orders] Agent 'orders' failed to start");
expect(subject.io.stderr()).toContain("Agent 'support' is running on port");

process.emit("SIGINT", "SIGINT");
await pending.catch(() => undefined);
});

test("exits non-zero when every agent fails to start", async () => {
const subject = harness({ project: twoRuntimes() });

await expect(subject.run()).rejects.toBeInstanceOf(SilentCLIError);
expect(subject.collector.state.closed).toBe(1);
});
});

describe("project dev trace collection", () => {
test("starts the collector, announces it, and points a CodeZip agent at loopback", async () => {
const subject = harness();
Expand Down Expand Up @@ -396,15 +471,6 @@ describe("project dev Inspector UI mode", () => {
"Port 9999 is already in use",
);
});

test("--port with several runtimes is rejected", async () => {
const subject = harness({
project: project(runtime("orders"), runtime("support", "Container")),
});
await expect(subject.run({ mode: "browser", port: 4567 })).rejects.toThrow(
"--port applies to a single runtime",
);
});
});

test("project dev renders attributed human and NDJSON output", async () => {
Expand Down
25 changes: 17 additions & 8 deletions src/handlers/project/dev/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import type { ProjectRuntime } from "../../../projectSchemas/runtime";
import {
InputValidationError,
ResourceNotFoundError,
SilentCLIError,
UserCancellationError,
} from "../../../errors";
import type { AppIO, BrowserOpener, FileWatcher, PortChecker, startHttpServer } from "../../../io";
Expand Down Expand Up @@ -102,7 +103,7 @@ export const createDevProjectHandler = (config: DevProjectHandlerConfig) =>
flag("traces", "disable local OTEL trace collection", z.boolean().default(true)),
flag(
"mode",
"how to run: browser (Agent Inspector web UI), headless (one agent in the terminal), or tui",
"how to run: browser (Agent Inspector web UI), headless (agents stream to the terminal), or tui",
z.enum(["browser", "headless", "tui"]).default("browser"),
),
flag(
Expand Down Expand Up @@ -132,12 +133,6 @@ export const createDevProjectHandler = (config: DevProjectHandlerConfig) =>
);
}
const runtimes = selectRuntimes(project, flags.agent);
if (flags.mode === "headless" && !flags.agent) {
const available = runtimes.map((runtime) => runtime.name).join(", ");
throw new InputValidationError(
`--mode headless runs a single agent in the terminal. Pass --agent <name> to choose which one. Available: ${available}.`,
);
}
if (runtimes.length > 1 && flags.port !== undefined) {
throw new InputValidationError(
"--port applies to a single runtime. Use --agent to select one.",
Expand Down Expand Up @@ -194,7 +189,7 @@ export const createDevProjectHandler = (config: DevProjectHandlerConfig) =>
return { ...env, ...otel };
};

if (flags.mode === "headless") {
if (flags.mode === "headless" && flags.agent) {
await runWithoutUi(
config,
runtimes[0]!,
Expand Down Expand Up @@ -227,6 +222,20 @@ export const createDevProjectHandler = (config: DevProjectHandlerConfig) =>
signal: controller.signal,
});

if (flags.mode === "headless") {
void Promise.allSettled(runtimes.map((runtime) => supervisor.start(runtime.name)));
for await (const { agentName, event } of supervisor.events()) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think normal Ctrl-C gets reported as a crash in this supervised headless path. The production process runner rejects with the child signal's UserCancellationError, so the supervisor emits Agent 'orders' crashed: Operation cancelled by user before the command exits 130. Could the supervisor treat an error from an already-aborted child signal as a normal stop? A test where the runner throws input.signal.reason on abort would cover the production behavior better than stayingRunner, which returns cleanly.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch, fixed. The pump now treats an error from an already-aborted child signal as a stop, and the runner fakes throw the abort reason like the real runner so the SIGINT tests cover the production path. Verified end to end: Ctrl-C exits 130 with no crash line.

renderAgentEvent(config.io, event, agentName, json);
const phases = supervisor.snapshot();
if (phases.every(({ phase }) => phase !== "starting" && phase !== "running")) {
if (phases.some(({ phase }) => phase === "failed")) throw new SilentCLIError();
break;
}
}
controller.signal.throwIfAborted();
return;
}

const uiPort = (
await findFreePort(UI_DEFAULT_PORT, flags["ui-port"], config.checkPort, controller.signal)
).port;
Expand Down
Loading