From 94002d8cccbadcc29ecbf995a3c78d3e41ae2216 Mon Sep 17 00:00:00 2001 From: Michael Johnston Date: Fri, 21 Aug 2026 15:05:41 -0700 Subject: [PATCH] fix(codex): restore pre-registration child liveness (cherry picked from commit 75b66cae0de74a7ebfa5a8d017a7444660cabe57) --- .../CodexCollabRuntime.integration.test.ts | 37 ++++++++++++++----- .../provider/Layers/CodexSessionRuntime.ts | 37 +++++++++++++++---- 2 files changed, 58 insertions(+), 16 deletions(-) diff --git a/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts b/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts index a1b46e003520..ae8eb3b6d490 100644 --- a/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts +++ b/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts @@ -181,6 +181,13 @@ describe("CodexSessionRuntime collab integration", () => { assert.isDefined(registrationA); assert.isDefined(registrationB); assert.isDefined(rootThreadStarted); + const interactedRegistrationA = { + ...registrationA, + params: { + ...registrationA.params, + item: { ...registrationA.params.item, kind: "interacted" }, + }, + }; const memoryThreadStarted = { ...rootThreadStarted, params: { @@ -207,7 +214,7 @@ describe("CodexSessionRuntime collab integration", () => { hangInterruptFor: CHILD_A, notifications: [ turnStartedA, - registrationA, + interactedRegistrationA, memoryThreadStarted, memoryTurnStarted, registrationB, @@ -233,26 +240,38 @@ describe("CodexSessionRuntime collab integration", () => { environment: { ...process.env, T3_CODEX_COLLAB_SCRIPT: scriptPath }, }); - // Wait for both children's turnStarted signals to be processed before - // stopping (B via the registered-child path; A only produces live-turn - // bookkeeping, so key on B's synthetic event). - const childBStartedFiber = yield* runtime.events.pipe( + // Wait for both children's synthetic turnStarted signals before + // stopping. B arrives through the registered-child path; A is replayed + // when its later activity registration finds the pre-registration live + // turn recorded by the foreign-notification suppressor. + const childrenStartedFiber = yield* runtime.events.pipe( Stream.filter( (event) => event.method === "collabAgent/turnStarted" && - (event.payload as { agentThreadId?: string }).agentThreadId === CHILD_B, + [CHILD_A, CHILD_B].includes( + (event.payload as { agentThreadId?: string }).agentThreadId ?? "", + ), ), - Stream.take(1), + Stream.take(2), Stream.runCollect, Effect.forkScoped, ); yield* runtime.start(); yield* runtime.sendTurn({ input: "fan out and hang" }); - const childBStarted = yield* Fiber.join(childBStartedFiber).pipe( + const childrenStarted = yield* Fiber.join(childrenStartedFiber).pipe( Effect.timeoutOption("15 seconds"), ); - assert.isTrue(childBStarted._tag === "Some", "child B turnStarted never arrived"); + assert.isTrue(childrenStarted._tag === "Some", "child turnStarted replay never arrived"); + if (childrenStarted._tag === "Some") { + const startedThreadIds = new Set( + Array.from(childrenStarted.value).map( + (event) => (event.payload as { agentThreadId?: string }).agentThreadId, + ), + ); + assert.isTrue(startedThreadIds.has(CHILD_A), "child A start must replay on registration"); + assert.isTrue(startedThreadIds.has(CHILD_B), "child B start must flow after registration"); + } // Stop everything. A's interrupt hangs forever — the bounded child // deadline must expire and the parent interrupt must still be sent. diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index fd926e43d7bf..d6ba1d7fc0bf 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -1107,8 +1107,8 @@ export const makeCodexSessionRuntime = ( return false; } const activitySpawnTurnId = (yield* Ref.get(sessionRef)).activeTurnId ?? undefined; + const existingChild = (yield* Ref.get(collabChildAgentsRef)).get(item.agentThreadId); yield* Ref.update(collabChildAgentsRef, (current) => { - const existing = current.get(item.agentThreadId); const next = new Map(current); // Merge-late semantics: when thread/started registered first, a // later subAgentActivity still carries the real agentPath (and a @@ -1120,13 +1120,13 @@ export const makeCodexSessionRuntime = ( next.set(item.agentThreadId, { agentThreadId: item.agentThreadId, nickname: - existing?.nickname ?? + existingChild?.nickname ?? item.agentPath.split("/").findLast((segment) => segment.length > 0), - role: existing?.role, - agentPath: existing?.agentPath ?? item.agentPath, - depth: existing?.depth, - parentThreadId: existing?.parentThreadId, - spawnTurnId: existing ? existing.spawnTurnId : activitySpawnTurnId, + role: existingChild?.role, + agentPath: existingChild?.agentPath ?? item.agentPath, + depth: existingChild?.depth, + parentThreadId: existingChild?.parentThreadId, + spawnTurnId: existingChild ? existingChild.spawnTurnId : activitySpawnTurnId, }); return next; }); @@ -1142,6 +1142,29 @@ export const makeCodexSessionRuntime = ( activityKind: item.kind, }, }); + // A child turn can start before this activity registers the child. + // The foreign-notification suppressor records that live turn but + // cannot emit agent lifecycle until identity is known. Replay the + // explicit start after first registration so sidebar liveness sees + // genuine work; a trailing interaction with no live turn remains + // ignored by CodexAdapter. + const preRegistrationLiveTurn = (yield* Ref.get(collabChildLiveTurnsRef)).get( + item.agentThreadId, + ); + if (!existingChild && item.kind === "interacted" && preRegistrationLiveTurn) { + yield* emitEvent({ + kind: "notification", + threadId: options.threadId, + ...(registeredChild?.spawnTurnId ? { turnId: registeredChild.spawnTurnId } : {}), + method: "collabAgent/turnStarted", + payload: { + agentThreadId: item.agentThreadId, + ...(registeredChild?.nickname ? { nickname: registeredChild.nickname } : {}), + ...(registeredChild?.role ? { role: registeredChild.role } : {}), + agentPath: item.agentPath, + }, + }); + } return true; }