Skip to content
Closed
Show file tree
Hide file tree
Changes from 1 commit
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
11 changes: 11 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,17 @@ 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,
...statusLinkage,
},
},
];
case "collabAgent/turnCompleted": {
// Idle, not terminal: the identity is resumable via sendInput/resume.
const turn =
Expand Down
140 changes: 140 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,150 @@ 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,
},
},
],
};
}

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("replays the captured fan-out into synthetic agent events without child leaks", () =>
Effect.gen(function* () {
// @effect-diagnostics-next-line preferSchemaOverJson:off
Expand Down
29 changes: 26 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,34 @@ 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"],
);
});

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