Skip to content
Closed
Show file tree
Hide file tree
Changes from 2 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
35 changes: 35 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -555,6 +555,41 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => {
}),
);

it.effect("renames an agent without overwriting its existing identity metadata", () =>
Effect.gen(function* () {
const { adapter, runtime } = yield* startLifecycleRuntime();
const firstEventFiber = yield* Stream.runHead(adapter.streamEvents).pipe(Effect.forkChild);

yield* runtime.emit({
id: asEventId("evt-agent-renamed"),
kind: "notification",
provider: ProviderDriverKind.make("codex"),
createdAt: "2026-01-01T00:00:00.000Z",
method: "collabAgent/renamed",
threadId: asThreadId("thread-1"),
payload: {
agentThreadId: "child-rename",
nickname: "Alpha",
},
} satisfies ProviderEvent);

const firstEvent = yield* Fiber.join(firstEventFiber);
NodeAssert.equal(firstEvent._tag, "Some");
if (firstEvent._tag !== "Some") {
return;
}
NodeAssert.equal(firstEvent.value.type, "task.updated");
if (firstEvent.value.type !== "task.updated") {
return;
}
NodeAssert.deepStrictEqual(firstEvent.value.payload, {
taskId: "child-rename",
title: "Alpha",
timelineBypass: true,
});
}),
);

it.effect("labels MCP lifecycle entries with server and tool names", () =>
Effect.gen(function* () {
const { adapter, runtime } = yield* startLifecycleRuntime();
Expand Down
12 changes: 12 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -605,6 +605,18 @@ function mapCollabAgentEvent(
payload: { taskId, status: "running", ...statusLinkage },
},
];
case "collabAgent/renamed":
return [
{
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
...base,
type: "task.updated",
payload: {
taskId,
...(nickname ? { title: nickname } : {}),
timelineBypass: true,
},
},
];
case "collabAgent/turnCompleted": {
// Idle, not terminal: the identity is resumable via sendInput/resume.
const turn =
Expand Down
244 changes: 244 additions & 0 deletions apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -70,10 +70,254 @@ function buildScript() {
};
}

function buildDirectChildScript() {
const rootTurnId = "019fcfd6-1806-7de1-8564-de69fd55bffb";
const childTurnId = `${CHILD_A}-direct-turn`;
const childItemId = `${CHILD_A}-direct-message`;
return {
rootThreadId: ROOT,
notifications: [
{
method: "item/completed",
params: {
threadId: ROOT,
turnId: rootTurnId,
item: {
type: "collabAgentToolCall",
id: "call_direct_spawn",
tool: "spawnAgent",
status: "completed",
senderThreadId: ROOT,
receiverThreadIds: [CHILD_A],
prompt: "Return one concise result.",
agentsStates: {
[CHILD_A]: { status: "pendingInit", message: null },
},
},
completedAtMs: 1785898350000,
},
},
{
method: "item/started",
params: {
threadId: CHILD_A,
turnId: childTurnId,
item: { type: "agentMessage", id: childItemId, text: "" },
},
},
{
method: "item/agentMessage/delta",
params: {
threadId: CHILD_A,
turnId: childTurnId,
itemId: childItemId,
delta: "child narration must not enter the parent transcript",
},
},
{
method: "item/completed",
params: {
threadId: CHILD_A,
turnId: childTurnId,
completedAtMs: 1785898350000,
item: {
type: "agentMessage",
id: childItemId,
phase: "final_answer",
text: "child result is consumed by the parent model",
},
},
},
{
method: "item/completed",
params: {
threadId: ROOT,
turnId: rootTurnId,
item: {
type: "agentMessage",
id: "root-summary",
phase: "final_answer",
text: "The agent completed:\n- Direct Child Researcher: returned one concise result.",
},
completedAtMs: 1785898350000,
},
},
],
};
}

function buildOutOfOrderNamingScript() {
const rootTurnId = "019fcfd6-1806-7de1-8564-de69fd55bffb";
return {
rootThreadId: ROOT,
notifications: [
// Child B speaks before the parent-side fan-out arrives. This makes
// provisional registration order differ from receiverThreadIds order.
{
method: "turn/started",
params: {
threadId: CHILD_B,
turn: {
id: `${CHILD_B}-early-turn`,
items: [],
itemsView: "notLoaded",
status: "inProgress",
error: null,
startedAt: 1785898342,
completedAt: null,
durationMs: null,
},
},
},
{
method: "item/completed",
params: {
threadId: ROOT,
turnId: rootTurnId,
item: {
type: "collabAgentToolCall",
id: "call_direct_spawn_ordered",
tool: "spawnAgent",
status: "completed",
senderThreadId: ROOT,
receiverThreadIds: [CHILD_A, CHILD_B],
prompt: "Return one concise result.",
agentsStates: {
[CHILD_A]: { status: "pendingInit", message: null },
[CHILD_B]: { status: "pendingInit", message: null },
},
},
completedAtMs: 1785898350000,
},
},
{
method: "item/completed",
params: {
threadId: ROOT,
turnId: rootTurnId,
item: {
type: "agentMessage",
id: "root-order-summary",
phase: "final_answer",
text: "Agents:\n- Alpha\n- Beta",
},
completedAtMs: 1785898350000,
},
},
],
};
}

const scriptPath = NodePath.join(import.meta.dirname, "../testFixtures/.collab-script.json");
const peerPath = NodePath.join(import.meta.dirname, "../testFixtures/codexCollabMockPeer.sh");

describe("CodexSessionRuntime collab integration", () => {
it.effect("registers receiver ids and keeps child narration out of the parent stream", () =>
Effect.gen(function* () {
// @effect-diagnostics-next-line preferSchemaOverJson:off
NodeFS.writeFileSync(scriptPath, JSON.stringify(buildDirectChildScript()), "utf8");
yield* Effect.addFinalizer(() =>
Effect.sync(() => NodeFS.rmSync(scriptPath, { force: true })),
);

const runtime = yield* makeCodexSessionRuntime({
threadId: ThreadId.make("thread-collab-direct-child"),
binaryPath: peerPath,
cwd: "/tmp",
runtimeMode: "full-access",
environment: { ...process.env, T3_CODEX_COLLAB_SCRIPT: scriptPath },
});

const eventsFiber = yield* runtime.events.pipe(
Stream.takeUntil((event) => event.method === "turn/completed"),
Stream.runCollect,
Effect.forkScoped,
);

yield* runtime.start();
yield* runtime.sendTurn({ input: "direct child" });

const events = Array.from(yield* Fiber.join(eventsFiber));
const methods = events.map((event) => event.method);
assert.include(methods, "collabAgent/started");
assert.include(methods, "collabAgent/item");
assert.include(methods, "collabAgent/statusChanged");
assert.include(methods, "collabAgent/renamed");
assert.include(methods, "turn/completed");
assert.notInclude(methods, "item/agentMessage/delta");
const started = events.find((event) => event.method === "collabAgent/started");
assert.isUndefined((started?.payload as { nickname?: string } | undefined)?.nickname);
const renamed = events.find((event) => event.method === "collabAgent/renamed");
assert.equal(
(renamed?.payload as { nickname?: string } | undefined)?.nickname,
"Direct Child Researcher",
);
const statusChanged = events.find(
(event) =>
event.method === "collabAgent/statusChanged" &&
(event.payload as { status?: { type?: string } } | undefined)?.status?.type === "idle",
);
assert.deepEqual((statusChanged?.payload as { status?: unknown } | undefined)?.status, {
type: "idle",
});

const leakedChildEvents = events.filter((event) => {
if (event.method === "item/agentMessage/delta") return true;
const payload = event.payload as { threadId?: string } | undefined;
return payload?.threadId === CHILD_A;
});
assert.deepEqual(
leakedChildEvents.map((event) => event.method),
[],
"child notifications must not be emitted as parent-timeline events",
);

yield* runtime.close;
}).pipe(Effect.scoped, Effect.provide(NodeServices.layer)),
);

it.effect("assigns coordinator names in receiver order after child-first registration", () =>
Effect.gen(function* () {
// @effect-diagnostics-next-line preferSchemaOverJson:off
NodeFS.writeFileSync(scriptPath, JSON.stringify(buildOutOfOrderNamingScript()), "utf8");
yield* Effect.addFinalizer(() =>
Effect.sync(() => NodeFS.rmSync(scriptPath, { force: true })),
);

const runtime = yield* makeCodexSessionRuntime({
threadId: ThreadId.make("thread-collab-name-order"),
binaryPath: peerPath,
cwd: "/tmp",
runtimeMode: "full-access",
environment: { ...process.env, T3_CODEX_COLLAB_SCRIPT: scriptPath },
});

const eventsFiber = yield* runtime.events.pipe(
Stream.takeUntil((event) => event.method === "turn/completed"),
Stream.runCollect,
Effect.forkScoped,
);

yield* runtime.start();
yield* runtime.sendTurn({ input: "name the children" });

const events = Array.from(yield* Fiber.join(eventsFiber));
const renamed = events.filter((event) => event.method === "collabAgent/renamed");
assert.deepEqual(
renamed.map((event) => {
const payload = event.payload as { agentThreadId?: string; nickname?: string };
return [payload.agentThreadId, payload.nickname];
}),
[
[CHILD_A, "Alpha"],
[CHILD_B, "Beta"],
],
);

yield* runtime.close;
}).pipe(Effect.scoped, Effect.provide(NodeServices.layer)),
);

it.effect("replays the captured fan-out into synthetic agent events without child leaks", () =>
Effect.gen(function* () {
// @effect-diagnostics-next-line preferSchemaOverJson:off
Expand Down
30 changes: 27 additions & 3 deletions apps/server/src/provider/Layers/CodexCollabWire.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
import { assert, describe, it } from "vite-plus/test";

import fixture from "../testFixtures/codexMultiAgentWire.json" with { type: "json" };
import { routeCodexChildNotification } from "./CodexSessionRuntime.ts";
import { readCoordinatorAgentNames, routeCodexChildNotification } from "./CodexSessionRuntime.ts";

interface WireNotification {
readonly method: string;
Expand Down Expand Up @@ -61,8 +61,8 @@ describe("codex multi-agent wire capture", () => {
it("emits child traffic BEFORE the item that registers the child", () => {
// Ordering hazard: the child's own thread/status/changed arrives before
// the parent-side subAgentActivity naming it. Registration must tolerate
// child-first arrival, so unregistered child traffic passes through
// rather than being eaten (no regression vs. pre-feature behavior).
// child-first arrival without leaking the child's conversation into the
// parent timeline.
const firstChildTraffic = notifications.findIndex((entry) => {
const threadId = notificationThreadId(entry);
return threadId !== undefined && childThreadIds.has(threadId);
Expand Down Expand Up @@ -128,11 +128,35 @@ describe("routeCodexChildNotification", () => {
}
});

it("reads coordinator-assigned names from fleet summaries", () => {
assert.deepEqual(
readCoordinatorAgentNames(
"Three agents are running in parallel: Halley, Banach, and Parfit.",
),
["Halley", "Banach", "Parfit"],
);
assert.deepEqual(
readCoordinatorAgentNames(
"All three agents completed successfully:\n\n- Halley: generated names\n- Banach: wrote a riddle\n- Parfit: sorted terms",
),
["Halley", "Banach", "Parfit"],
);
assert.deepEqual(
readCoordinatorAgentNames(
"Started 3 Luna medium sub-agents:\n\n- Planck\n- Parfit\n- Avicenna",
),
["Planck", "Parfit", "Avicenna"],
);
assert.deepEqual(readCoordinatorAgentNames("Summary:\n\n- Tests: passed\n- Build: passed"), []);
});

it("drops only enumerated child chatter", () => {
for (const method of [
"item/agentMessage/delta",
"item/reasoning/textDelta",
"item/commandExecution/outputDelta",
"item/commandExecution/terminalInteraction",
"item/mcpToolCall/progress",
"turn/plan/updated",
"thread/name/updated",
]) {
Expand Down
Loading
Loading