diff --git a/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts b/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts index 02b7a45f33ad..f400e7f9d103 100644 --- a/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts +++ b/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts @@ -29,6 +29,10 @@ import { HostProcessPlatform } from "@t3tools/shared/hostProcess"; const ROOT = wireFixture.rootThreadId; const [CHILD_A, CHILD_B] = wireFixture.childThreadIds as [string, string]; const MEMORY = "memory-consolidation-thread"; +const FOREIGN = "unregistered-foreign-thread"; +const FOREIGN_RECEIVER = "unregistered-foreign-receiver"; +const FOREIGN_ACTIVITY_CHILD = "unregistered-foreign-activity-child"; +const NESTED_ACTIVITY_CHILD = "owned-nested-activity-child"; const decodeMcpElicitationResponse = Schema.decodeUnknownEffect( Schema.fromJsonString( Schema.Struct({ @@ -62,6 +66,85 @@ function buildScript() { }, }, }, + // An already-owned child can establish ownership for descendants. + { + method: "item/completed", + params: { + completedAtMs: 5, + threadId: CHILD_A, + turnId: `${CHILD_A}-turn-1`, + item: { + id: "nested-sub-agent-activity", + type: "subAgentActivity", + kind: "started", + agentThreadId: NESTED_ACTIVITY_CHILD, + agentPath: "/root/alpha/nested", + }, + }, + }, + // Foreign collaboration announcements cannot establish ownership for + // receiver IDs or subAgentActivity children in this runtime. + { + method: "item/completed", + params: { + completedAtMs: 2, + threadId: FOREIGN, + turnId: "foreign-turn", + item: { + id: "foreign-collab-call", + type: "collabAgentToolCall", + tool: "wait", + status: "completed", + senderThreadId: FOREIGN, + receiverThreadIds: [FOREIGN_RECEIVER], + }, + }, + }, + { + method: "item/completed", + params: { + completedAtMs: 3, + threadId: FOREIGN_RECEIVER, + turnId: "foreign-receiver-turn", + item: { + id: "foreign-receiver-message", + type: "agentMessage", + phase: "final_answer", + text: "foreign receiver report", + }, + }, + }, + { + method: "item/completed", + params: { + completedAtMs: 4, + threadId: FOREIGN, + turnId: "foreign-turn", + item: { + id: "foreign-sub-agent-activity", + type: "subAgentActivity", + kind: "started", + agentThreadId: FOREIGN_ACTIVITY_CHILD, + agentPath: "/root/foreign", + }, + }, + }, + // A different app-server session can emit an assistant item without a + // preceding thread/started notification. It must not become parent chat. + { + method: "item/completed", + params: { + completedAtMs: 1, + threadId: FOREIGN, + turnId: "foreign-turn", + item: { + id: "foreign-message", + type: "agentMessage", + phase: "final_answer", + text: "unrelated background report", + }, + }, + }, // Child terminal lifecycle AFTER the receiver map knows the children — // pre-fix, the legacy suppressor dropped these before interception saw // them, so no synthetic agent events were emitted. @@ -166,6 +249,90 @@ const peerPath = NodePath.join( ); describe("CodexSessionRuntime collab integration", () => { + it.effect("keeps startup ownership on the authoritative thread response", () => + Effect.gen(function* () { + const script = { + rootThreadId: ROOT, + startupResponseDelayMs: 25, + startupNotifications: [ + capturedSpawnedThread(FOREIGN), + { + method: "item/agentMessage/delta", + params: { + delta: "foreign startup content", + itemId: "foreign-startup-message", + threadId: FOREIGN, + turnId: "foreign-startup-turn", + }, + }, + ], + notifications: [ + { + method: "item/agentMessage/delta", + params: { + delta: "owned root content", + itemId: "root-message", + threadId: ROOT, + turnId: `${ROOT}-turn`, + }, + }, + { method: "warning", params: { message: "foreign warning", threadId: FOREIGN } }, + { method: "warning", params: { message: "root warning", threadId: ROOT } }, + ], + }; + // @effect-diagnostics-next-line preferSchemaOverJson:off + NodeFS.writeFileSync(scriptPath, JSON.stringify(script), "utf8"); + yield* Effect.addFinalizer(() => + Effect.sync(() => NodeFS.rmSync(scriptPath, { force: true })), + ); + + const runtime = yield* makeCodexSessionRuntime({ + threadId: ThreadId.make("thread-startup-ownership"), + binaryPath: peerPath, + cwd: NodeOS.tmpdir(), + 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, + ); + + const session = yield* runtime.start(); + assert.deepEqual(session.resumeCursor, { threadId: ROOT }); + yield* runtime.sendTurn({ input: "continue on the owned root" }); + const events = Array.from(yield* Fiber.join(eventsFiber)); + + assert.isFalse( + events.some((event) => event.textDelta === "foreign startup content"), + "foreign content emitted before the root response must stay out of chat", + ); + assert.isTrue( + events.some((event) => event.textDelta === "owned root content"), + "root content remains visible after authoritative ownership is established", + ); + assert.isFalse( + events.some( + (event) => + event.method === "warning" && + (event.payload as { message?: string } | undefined)?.message === "foreign warning", + ), + "foreign thread warnings must stay out of the root conversation", + ); + assert.isTrue( + events.some( + (event) => + event.method === "warning" && + (event.payload as { message?: string } | undefined)?.message === "root warning", + ), + "root thread warnings remain visible", + ); + + yield* runtime.close; + }).pipe(Effect.scoped, Effect.provide(NodeServices.layer)), + ); + it.effect("looks up child model metadata once after activity registration", () => Effect.gen(function* () { const script = { @@ -468,6 +635,32 @@ describe("CodexSessionRuntime collab integration", () => { [], "child thread/* lifecycle must not appear as parent events", ); + assert.isFalse( + events.some( + (event) => (event.payload as { threadId?: string } | undefined)?.threadId === FOREIGN, + ), + "unregistered foreign assistant items must not appear as parent events", + ); + assert.isFalse( + events.some((event) => { + const payload = event.payload as + | { threadId?: string; agentThreadId?: string } + | undefined; + return ( + payload?.threadId === FOREIGN_RECEIVER || + payload?.agentThreadId === FOREIGN_ACTIVITY_CHILD + ); + }), + "foreign collaboration announcements must not establish child ownership", + ); + assert.isTrue( + events.some( + (event) => + (event.payload as { agentThreadId?: string } | undefined)?.agentThreadId === + NESTED_ACTIVITY_CHILD, + ), + "owned children can register nested subAgentActivity descendants", + ); yield* runtime.close; }).pipe(Effect.scoped, Effect.provide(NodeServices.layer)), diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts index ec113ab7c521..fa36e553f0e6 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts @@ -18,6 +18,7 @@ import { isRecoverableThreadResumeError, makeMemoryConsolidationNotificationFilter, openCodexThread, + shouldSuppressUnownedCodexNotification, readCodexThread, rollbackCodexThread, toMcpElicitationResponse, @@ -802,6 +803,121 @@ describe("makeMemoryConsolidationNotificationFilter", () => { }); }); +describe("shouldSuppressUnownedCodexNotification", () => { + const agentMessageDelta = (threadId: string, delta: string) => ({ + method: "item/agentMessage/delta" as const, + params: { + delta, + itemId: `${threadId}-message`, + threadId, + turnId: `${threadId}-turn`, + }, + }); + + it("suppresses unregistered foreign assistant text without hiding root or child replies", () => { + const foreignNotifications: ReadonlyArray< + Parameters[0] + > = [ + agentMessageDelta("foreign-thread", "unrelated background report"), + { + method: "item/started" as const, + params: { + startedAtMs: 1, + threadId: "foreign-thread", + turnId: "foreign-turn", + item: { + id: "foreign-message", + type: "agentMessage" as const, + text: "", + }, + }, + }, + { + method: "item/completed" as const, + params: { + completedAtMs: 2, + threadId: "foreign-thread", + turnId: "foreign-turn", + item: { + id: "foreign-message", + type: "agentMessage" as const, + text: "unrelated background report", + }, + }, + }, + ]; + for (const notification of foreignNotifications) { + NodeAssert.equal( + shouldSuppressUnownedCodexNotification(notification, "root-thread", false), + true, + ); + } + NodeAssert.equal( + shouldSuppressUnownedCodexNotification( + agentMessageDelta("root-thread", "root commentary"), + "root-thread", + false, + ), + false, + ); + NodeAssert.equal( + shouldSuppressUnownedCodexNotification( + agentMessageDelta("registered-child", "child reply"), + "root-thread", + true, + ), + false, + ); + NodeAssert.equal( + shouldSuppressUnownedCodexNotification( + { method: "warning", params: { message: "foreign warning", threadId: "foreign-thread" } }, + "root-thread", + false, + ), + true, + ); + NodeAssert.equal( + shouldSuppressUnownedCodexNotification( + { method: "warning", params: { message: "root warning", threadId: "root-thread" } }, + "root-thread", + false, + ), + false, + ); + }); + + it("keeps foreign request resolution on the parent correlation path", () => { + NodeAssert.equal( + shouldSuppressUnownedCodexNotification( + { + method: "serverRequest/resolved", + params: { requestId: "request-1", threadId: "foreign-thread" }, + }, + "root-thread", + false, + ), + false, + ); + }); + + it("suppresses thread-addressed startup traffic until the root response establishes ownership", () => { + const threadStarted = makeThreadStartedNotification("foreign-thread", "appServer"); + NodeAssert.equal(shouldSuppressUnownedCodexNotification(threadStarted, undefined, false), true); + NodeAssert.equal( + shouldSuppressUnownedCodexNotification(threadStarted, "root-thread", false), + true, + ); + NodeAssert.equal( + shouldSuppressUnownedCodexNotification( + { method: "warning", params: { message: "global warning" } }, + undefined, + false, + ), + false, + ); + }); +}); + describe("codexSessionAppServerArgs", () => { it("keeps the app-server subcommand when explicit args are provided", () => { NodeAssert.deepStrictEqual(codexSessionAppServerArgs(["-c", "model=gpt-5"], undefined), [ diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index c8aa069cb2b7..a32693efea33 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -779,6 +779,7 @@ function readNotificationThreadId(notification: CodexServerNotification): string case "thread/started": return notification.params.thread.id; case "error": + case "warning": case "thread/status/changed": case "thread/archived": case "thread/unarchived": @@ -818,12 +819,34 @@ function readNotificationThreadId(notification: CodexServerNotification): string case "thread/realtime/sdp": case "thread/realtime/error": case "thread/realtime/closed": - return notification.params.threadId; + return notification.params.threadId ?? undefined; default: return undefined; } } +export function shouldSuppressUnownedCodexNotification( + notification: CodexServerNotification, + rootProviderThreadId: string | undefined, + isRegisteredChild: boolean, +): boolean { + const providerThreadId = readNotificationThreadId(notification); + if (providerThreadId !== undefined && rootProviderThreadId === undefined) { + return true; + } + if ( + providerThreadId === undefined || + providerThreadId === rootProviderThreadId || + isRegisteredChild + ) { + return false; + } + + // Resolution notifications complete requests owned by the parent runtime, + // even when Codex addresses them to a child provider thread. + return notification.method !== "serverRequest/resolved"; +} + export function makeMemoryConsolidationNotificationFilter(): ( notification: CodexServerNotification, ) => boolean { @@ -936,12 +959,10 @@ function readRouteFields(notification: CodexServerNotification): { * synthetic `collabAgent/*` provider events the adapter turns into task.* * runtime events (timelineBypass keeps them out of the parent chat). * - * WIP, probe-gated: registration is deliberately explicit-signals-only. The - * spec's "provisionally treat unknown foreign thread ids as v2 children" rule - * needs a live wire capture of the packaged binary before it lands — blind - * capture risks eating unrelated traffic. Until then a child whose first - * notification precedes registration passes through as today (no regression - * vs main, which passes everything through). + * Registration is deliberately explicit-signals-only. Notifications from an + * unknown foreign thread stay unregistered and are suppressed at the parent + * boundary until one of those signals arrives, preventing unrelated provider + * sessions from leaking into this conversation. */ interface CollabChildAgentState { readonly agentThreadId: string; @@ -1024,8 +1045,9 @@ function rememberCollabReceiverTurns( collabReceiverTurns: Map, notification: CodexServerNotification, parentTurnId: TurnId | undefined, + sourceOwnedBySession: boolean, ): void { - if (!parentTurnId) { + if (!parentTurnId || !sourceOwnedBySession) { return; } @@ -1042,6 +1064,21 @@ function rememberCollabReceiverTurns( } } +function isOwnedCollabThreadId( + providerThreadId: string | undefined, + rootProviderThreadId: string | undefined, + collabReceiverTurns: ReadonlyMap, + collabChildAgents: ReadonlyMap, +): boolean { + return ( + providerThreadId !== undefined && + rootProviderThreadId !== undefined && + (providerThreadId === rootProviderThreadId || + collabReceiverTurns.has(providerThreadId) || + collabChildAgents.has(providerThreadId)) + ); +} + function shouldSuppressChildConversationNotification( method: CodexRpc.ServerNotificationMethod, ): boolean { @@ -1557,7 +1594,10 @@ export const makeCodexSessionRuntime = ( * Returns true when the notification was fully handled (must not reach * parent-timeline mapping). */ - const interceptCollabChildNotification = (notification: CodexServerNotification) => + const interceptCollabChildNotification = ( + notification: CodexServerNotification, + sourceOwnedBySession: boolean, + ) => Effect.gen(function* () { // Registration path 1: child thread announces itself with a // subAgent thread_spawn source. @@ -1571,6 +1611,14 @@ export const makeCodexSessionRuntime = ( if (thread.id === rootProviderThreadId) { return false; } + const parentThreadId = spawn.parentThreadId ?? thread.parentThreadId ?? undefined; + const receiverTurns = yield* Ref.get(collabReceiverTurnsRef); + const childAgents = yield* Ref.get(collabChildAgentsRef); + if ( + !isOwnedCollabThreadId(parentThreadId, rootProviderThreadId, receiverTurns, childAgents) + ) { + return false; + } // Merge with any subAgentActivity registration that got here // first. spawnTurnId is REGISTRATION-time-only on both paths: for // an already-known child we keep its value (set or unset) — a @@ -1619,6 +1667,9 @@ export const makeCodexSessionRuntime = ( (notification.method === "item/started" || notification.method === "item/completed") && notification.params.item.type === "subAgentActivity" ) { + if (!sourceOwnedBySession) { + return false; + } const item = notification.params.item; // Never register the session's ROOT thread as its own child. The // wire emits subAgentActivity {agentPath: "/root", interacted} @@ -1868,20 +1919,33 @@ export const makeCodexSessionRuntime = ( const payload = notification.params; const route = readRouteFields(notification); const collabReceiverTurns = yield* Ref.get(collabReceiverTurnsRef); + const suppressRootId = currentProviderThreadId(yield* Ref.get(sessionRef)); + const providerConversationId = readNotificationThreadId(notification); + const collabChildAgents = yield* Ref.get(collabChildAgentsRef); + const sourceOwnedBySession = isOwnedCollabThreadId( + providerConversationId, + suppressRootId, + collabReceiverTurns, + collabChildAgents, + ); const childParentTurnId = (() => { - const providerConversationId = readNotificationThreadId(notification); return providerConversationId ? collabReceiverTurns.get(providerConversationId) : undefined; })(); - rememberCollabReceiverTurns(collabReceiverTurns, notification, route.turnId); + rememberCollabReceiverTurns( + collabReceiverTurns, + notification, + route.turnId, + sourceOwnedBySession, + ); // Interception FIRST: a registered v2 child is usually also in the // receiver-turn map (collabAgentToolCall.receiverThreadIds), and the // legacy suppressor below would drop its lifecycle before it could // become synthetic collabAgent events (review finding). The // suppressor still covers UNREGISTERED children. - if (yield* interceptCollabChildNotification(notification)) { + if (yield* interceptCollabChildNotification(notification, sourceOwnedBySession)) { yield* Ref.set(collabReceiverTurnsRef, collabReceiverTurns); return; } @@ -1891,11 +1955,9 @@ export const makeCodexSessionRuntime = ( // (codexMultiAgentWire.json) shows a child's thread/status/changed // arriving BEFORE anything registers the child — pre-registration // lifecycle must not reach the parent path, where the adapter maps - // thread/* onto parent session state. Root-id-known guard keeps the - // root's own early notifications flowing during session open. - const suppressRootId = currentProviderThreadId(yield* Ref.get(sessionRef)); + // thread/* onto parent session state. The unowned guard below enforces + // the authoritative root boundary for every other thread method. const foreignConversation = (() => { - const providerConversationId = readNotificationThreadId(notification); return ( providerConversationId !== undefined && suppressRootId !== undefined && @@ -1942,6 +2004,20 @@ export const makeCodexSessionRuntime = ( return; } + // Codex app-server can emit another session's thread traffic without a + // preceding thread/started notification. Registered collaboration + // children were handled above; all other foreign thread traffic must + // stay out of this canonical T3 conversation. + if ( + shouldSuppressUnownedCodexNotification( + notification, + suppressRootId, + childParentTurnId !== undefined, + ) + ) { + return; + } + if (isMemoryConsolidationNotification) { return; } @@ -1993,7 +2069,10 @@ export const makeCodexSessionRuntime = ( yield* client.handleServerNotification("thread/started", (payload) => currentSessionProviderThreadId.pipe( Effect.flatMap((providerThreadId) => { - if (providerThreadId && payload.thread.id !== providerThreadId) { + // thread/start and thread/resume responses authoritatively establish + // ownership. An unsolicited startup notification must not claim an + // uninitialized runtime for another provider thread. + if (!providerThreadId || payload.thread.id !== providerThreadId) { return Effect.void; } return updateSession(sessionRef, { diff --git a/apps/server/src/provider/testFixtures/codexCollabMockPeer.mjs b/apps/server/src/provider/testFixtures/codexCollabMockPeer.mjs index fa567d75cf8a..1a7f91da8921 100644 --- a/apps/server/src/provider/testFixtures/codexCollabMockPeer.mjs +++ b/apps/server/src/provider/testFixtures/codexCollabMockPeer.mjs @@ -66,7 +66,17 @@ rl.on("line", (line) => { return; } if (method === "thread/start") { - write({ id, result: fixture.responses.threadStart }); + for (const notification of script.startupNotifications ?? []) { + write({ jsonrpc: "2.0", method: notification.method, params: notification.params }); + } + if (script.startupResponseDelayMs) { + setTimeout( + () => write({ id, result: fixture.responses.threadStart }), + script.startupResponseDelayMs, + ); + } else { + write({ id, result: fixture.responses.threadStart }); + } return; } if (method === "thread/resume") {