diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index a6636fcd69d0..fd1e52c6273f 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -2147,7 +2147,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); - it.effect("marks subagents when a child is proven related by ancestry lookup", () => + it.effect("replays an out-of-order child session after ancestry lookup", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; const threadId = asThreadId("thread-child-ancestry-usage"); @@ -2156,17 +2156,25 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { runtimeMock.state.sessionParentById.set("ses_child", "http://127.0.0.1:9999/session"); runtimeMock.state.sessionStatus = "busy"; const busy = promiseWithResolvers(); - const childPermission = promiseWithResolvers(); + const childSession = promiseWithResolvers(); + const childIdle = promiseWithResolvers(); const idle = promiseWithResolvers(); - runtimeMock.state.subscribedEvents = [busy.promise, childPermission.promise, idle.promise]; + runtimeMock.state.subscribedEvents = [ + busy.promise, + childSession.promise, + childIdle.promise, + idle.promise, + ]; const eventsFiber = yield* adapter.streamEvents.pipe( Stream.filter( (event) => event.threadId === threadId && - (event.type === "request.opened" || event.type === "turn.completed"), + (event.type === "task.started" || + event.type === "task.completed" || + event.type === "turn.completed"), ), - Stream.take(2), + Stream.take(3), Stream.runCollect, Effect.forkChild, ); @@ -2199,13 +2207,19 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { runtimeMock.state.sessionGetObserved = (sessionID) => { if (sessionID === "ses_child") requestOpened.resolve(undefined); }; - childPermission.resolve({ - id: "evt-child-ancestry-permission", - type: "permission.asked", - properties: permissionRequest("per_child_ancestry", "ses_child"), + childSession.resolve({ + id: "evt-child-ancestry-session", + type: "session.created", + properties: { info: { id: "ses_child", title: "Child task" } }, }); yield* Effect.promise(() => requestOpened.promise); yield* Effect.yieldNow; + childIdle.resolve({ + id: "evt-child-ancestry-idle", + type: "session.status", + properties: { sessionID: "ses_child", status: { type: "idle" } }, + }); + yield* Effect.yieldNow; runtimeMock.state.sessionStatus = "idle"; idle.resolve({ id: "evt-child-ancestry-idle", @@ -2219,9 +2233,9 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { const events = Array.from(yield* Fiber.join(eventsFiber).pipe(Effect.timeout("1 second"))); NodeAssert.deepEqual( events.map((event) => event.type), - ["request.opened", "turn.completed"], + ["task.started", "task.completed", "turn.completed"], ); - const completed = events[1]; + const completed = events[2]; if (completed?.type === "turn.completed") { NodeAssert.deepEqual(completed.payload.tokenUsage, { usageStatus: "unavailable", @@ -2233,6 +2247,100 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("completes a related child once when its idle status precedes deletion", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-child-idle-then-deleted"); + const parentSessionId = "http://127.0.0.1:9999/session"; + const busy = promiseWithResolvers(); + const childCreated = promiseWithResolvers(); + const childIdle = promiseWithResolvers(); + const childDeleted = promiseWithResolvers(); + const idle = promiseWithResolvers(); + runtimeMock.state.sessionStatus = "busy"; + runtimeMock.state.subscribedEvents = [ + busy.promise, + childCreated.promise, + childIdle.promise, + childDeleted.promise, + idle.promise, + ]; + + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => + event.threadId === threadId && + (event.type === "task.started" || + event.type === "task.completed" || + event.type === "turn.completed"), + ), + Stream.take(3), + Stream.runCollect, + Effect.forkChild, + ); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "approval-required", + }); + const sendFiber = yield* adapter + .sendTurn({ + threadId, + input: "Delegate to a child", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }) + .pipe(Effect.forkChild); + busy.resolve({ + id: "evt-child-delete-busy", + type: "session.status", + properties: { sessionID: parentSessionId, status: { type: "busy" } }, + }); + yield* Fiber.join(sendFiber); + + childCreated.resolve({ + id: "evt-child-delete-created", + type: "session.created", + properties: { + sessionID: "ses_child_delete", + info: { + id: "ses_child_delete", + parentID: parentSessionId, + title: "Child task", + }, + }, + }); + yield* Effect.yieldNow; + childIdle.resolve({ + id: "evt-child-delete-idle", + type: "session.status", + properties: { sessionID: "ses_child_delete", status: { type: "idle" } }, + }); + yield* Effect.yieldNow; + childDeleted.resolve({ + id: "evt-child-delete-deleted", + type: "session.deleted", + properties: { info: { id: "ses_child_delete", title: "Child task" } }, + }); + yield* Effect.yieldNow; + runtimeMock.state.sessionStatus = "idle"; + idle.resolve({ + id: "evt-child-delete-parent-idle", + type: "session.status", + properties: { sessionID: parentSessionId, status: { type: "idle" } }, + }); + + const events = Array.from(yield* Fiber.join(eventsFiber).pipe(Effect.timeout("1 second"))); + NodeAssert.deepEqual( + events.map((event) => event.type), + ["task.started", "task.completed", "turn.completed"], + ); + yield* adapter.stopSession(threadId); + }), + ); + it.effect("sums owned OpenCode step usage and marks unresolved usage partial", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; @@ -3889,8 +3997,8 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { ]; const openedEventsFiber = yield* adapter.streamEvents.pipe( - Stream.filter((event) => event.threadId === threadId), - Stream.take(3), + Stream.filter((event) => event.threadId === threadId && event.type === "request.opened"), + Stream.take(1), Stream.runCollect, Effect.forkChild, ); @@ -4189,8 +4297,10 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { ]; const requestedEventsFiber = yield* adapter.streamEvents.pipe( - Stream.filter((event) => event.threadId === threadId), - Stream.take(3), + Stream.filter( + (event) => event.threadId === threadId && event.type === "user-input.requested", + ), + Stream.take(1), Stream.runCollect, Effect.forkChild, ); diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 3d216bb1167b..73638c7f40ca 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -8,6 +8,7 @@ import { type ProviderSession, RuntimeItemId, RuntimeRequestId, + RuntimeTaskId, ThreadId, type ToolLifecycleItemType, type TurnTokenUsage, @@ -245,11 +246,24 @@ type OpenCodeAskedRequestEvent = Extract< type OpenCodeRoutedRequestEvent = OpenCodeAskedRequestEvent | OpenCodeTerminalRequestEvent; +type OpenCodeChildSessionEvent = Extract< + OpenCodeSubscribedEvent, + { readonly type: "session.created" | "session.updated" | "session.deleted" | "session.status" } +>; + interface OpenCodeRequestRelationRetry { warned: boolean; fiber?: Fiber.Fiber; } +interface OpenCodeSessionRelationRetry { + readonly events: Array<{ + readonly event: OpenCodeChildSessionEvent; + readonly turnId: TurnId | undefined; + }>; + fiber?: Fiber.Fiber; +} + interface OpenCodePendingRequestRecovery { warned: boolean; rerun: boolean; @@ -316,6 +330,15 @@ function isOpenCodeChildRequestEvent(event: OpenCodeSubscribedEvent): boolean { } } +function isOpenCodeChildSessionEvent(event: OpenCodeSubscribedEvent): boolean { + return ( + event.type === "session.created" || + event.type === "session.updated" || + event.type === "session.deleted" || + event.type === "session.status" + ); +} + const OPENCODE_DEFAULT_TITLE_PATTERN = /^(New session - |Child session - )\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\.\d{3}Z$/; @@ -343,7 +366,9 @@ interface OpenCodeSessionContext { readonly resolvedRequestIds: Set; readonly autoRepliedRequestIds: Set; readonly emittedTerminalRequestIds: Set; + readonly terminalChildSessionIds: Set; readonly requestRelationRetries: Map; + readonly sessionRelationRetries: Map; readonly pendingPermissions: Map; readonly pendingQuestions: Map; readonly messageRoleById: Map; @@ -1681,7 +1706,11 @@ export function makeOpenCodeAdapter( ), ); let sessionId: string | undefined = candidateSessionId; - for (let depth = 0; sessionId !== undefined && depth < 32; depth += 1) { + // `seen` bounds the walk: a session can have only one parent, so the + // chain terminates on a cycle, a root (`parentID` undefined), or a 404. + // A fixed hop cap would instead misreport a deep-but-related descendant + // as unrelated and drop its queued terminal events. + while (sessionId !== undefined) { if (context.relatedSessionIds.has(sessionId)) { addRelatedOpenCodeSession(context, candidateSessionId); return true; @@ -2034,6 +2063,112 @@ export function makeOpenCodeAdapter( retry.fiber = yield* run.pipe(Effect.forkIn(context.sessionScope)); }); + const scheduleChildSessionRelationRetry = Effect.fn("scheduleChildSessionRelationRetry")( + function* (context: OpenCodeSessionContext, event: OpenCodeChildSessionEvent) { + const sessionId = openCodeEventSessionId(event); + if (sessionId === undefined) return; + const turnId = context.activeTurnId; + const existing = context.sessionRelationRetries.get(sessionId); + if (existing) { + existing.events.push({ event, turnId }); + return; + } + const retry: OpenCodeSessionRelationRetry = { events: [{ event, turnId }] }; + context.sessionRelationRetries.set(sessionId, retry); + const run = Effect.gen(function* () { + let retryCount = 0; + while (context.sessionRelationRetries.get(sessionId) === retry) { + const relation = yield* isRelatedOpenCodeSession(context, sessionId).pipe( + Effect.match({ + onFailure: () => "unknown" as const, + onSuccess: (related): "related" | "unrelated" => + related ? "related" : "unrelated", + }), + ); + if (context.sessionRelationRetries.get(sessionId) !== retry) return; + if (relation === "unrelated") { + context.sessionRelationRetries.delete(sessionId); + return; + } + if (relation === "related") { + context.sessionRelationRetries.delete(sessionId); + for (const { event: replayEvent, turnId: replayTurnId } of retry.events) { + const base = yield* buildEventBase({ + threadId: context.session.threadId, + turnId: replayTurnId, + itemId: sessionId, + raw: replayEvent, + }); + if (replayEvent.type === "session.created") { + const replaySession = replayEvent.properties.info; + yield* emit({ + ...base, + type: "task.started", + payload: { + taskId: RuntimeTaskId.make(sessionId), + taskType: "local_agent", + title: replaySession.title, + description: replaySession.title, + }, + }); + } else if (replayEvent.type === "session.updated") { + const replaySession = replayEvent.properties.info; + yield* emit({ + ...base, + type: "task.progress", + payload: { + taskId: RuntimeTaskId.make(sessionId), + taskType: "local_agent", + title: replaySession.title, + description: replaySession.title, + summary: replaySession.title, + status: "running", + }, + }); + } else if ( + replayEvent.type === "session.status" && + replayEvent.properties.status.type !== "idle" + ) { + continue; + } else { + if (context.terminalChildSessionIds.has(sessionId)) continue; + context.terminalChildSessionIds.add(sessionId); + yield* emit({ + ...base, + type: "task.completed", + payload: { + taskId: RuntimeTaskId.make(sessionId), + taskType: "local_agent", + status: "completed", + summary: + (replayEvent.type === "session.deleted" + ? replayEvent.properties.info.title + : undefined) ?? "Completed", + }, + }); + return; + } + } + return; + } + const delayMs = Math.min(250 * 2 ** retryCount, 5_000); + retryCount += 1; + yield* Effect.sleep(`${delayMs} millis`); + } + }).pipe( + Effect.catchCause(() => Effect.void), + Effect.ensuring( + Effect.sync(() => { + if (context.sessionRelationRetries.get(sessionId) === retry) { + context.sessionRelationRetries.delete(sessionId); + } + }), + ), + ); + retry.fiber = yield* run.pipe(Effect.forkIn(context.sessionScope)); + }, + ); + const schedulePendingRequestRecovery = Effect.fn("schedulePendingRequestRecovery")(function* ( context: OpenCodeSessionContext, ) { @@ -2204,8 +2339,6 @@ export function makeOpenCodeAdapter( if (session.parentID && context.relatedSessionIds.has(session.parentID)) { addRelatedOpenCodeSession(context, session.id); } - } else if (event.type === "session.deleted") { - context.relatedSessionIds.delete(event.properties.info.id); } const payloadSessionId = openCodeEventSessionId(event); @@ -2214,9 +2347,17 @@ export function makeOpenCodeAdapter( if ( payloadSessionId !== undefined && !context.relatedSessionIds.has(payloadSessionId) && - isOpenCodeChildRequestEvent(event) + (isOpenCodeChildRequestEvent(event) || isOpenCodeChildSessionEvent(event)) ) { - if (event.type === "permission.asked") { + if ( + event.type === "session.created" || + event.type === "session.updated" || + event.type === "session.deleted" || + event.type === "session.status" + ) { + yield* scheduleChildSessionRelationRetry(context, event); + return; + } else if (event.type === "permission.asked") { yield* scheduleRequestRelationRetry(context, event); } else if (event.type === "question.asked") { yield* scheduleRequestRelationRetry(context, event); @@ -2236,7 +2377,7 @@ export function makeOpenCodeAdapter( } const isChildRequestEvent = payloadSessionId !== undefined && - isOpenCodeChildRequestEvent(event) && + (isOpenCodeChildRequestEvent(event) || isOpenCodeChildSessionEvent(event)) && (context.relatedSessionIds.has(payloadSessionId) || isKnownPendingTerminalEvent); if (!isParentEvent && !isChildRequestEvent) { return; @@ -2270,8 +2411,49 @@ export function makeOpenCodeAdapter( } switch (event.type) { + case "session.created": { + if (!isParentEvent) { + const session = event.properties.info; + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId, + itemId: session.id, + raw: event, + })), + type: "task.started", + payload: { + taskId: RuntimeTaskId.make(session.id), + taskType: "local_agent", + title: session.title, + description: session.title, + }, + }); + } + break; + } case "session.updated": { - const title = openCodeEventSessionTitle(event); + if (!isParentEvent) { + const session = event.properties.info; + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId, + itemId: session.id, + raw: event, + })), + type: "task.progress", + payload: { + taskId: RuntimeTaskId.make(session.id), + taskType: "local_agent", + title: session.title, + description: session.title, + summary: session.title, + status: "running", + }, + }); + } + const title = isParentEvent ? openCodeEventSessionTitle(event) : undefined; if (title) { yield* emit({ ...(yield* buildEventBase({ @@ -2289,6 +2471,30 @@ export function makeOpenCodeAdapter( } break; } + case "session.deleted": { + if (!isParentEvent) { + const session = event.properties.info; + context.relatedSessionIds.delete(session.id); + if (context.terminalChildSessionIds.has(session.id)) break; + context.terminalChildSessionIds.add(session.id); + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId, + itemId: session.id, + raw: event, + })), + type: "task.completed", + payload: { + taskId: RuntimeTaskId.make(session.id), + taskType: "local_agent", + status: "completed", + summary: session.title, + }, + }); + } + break; + } case "session.compacted": { yield* emit({ ...(yield* buildEventBase({ @@ -2564,6 +2770,29 @@ export function makeOpenCodeAdapter( } case "session.status": { + if (!isParentEvent) { + if (event.properties.status.type === "idle") { + const sessionId = event.properties.sessionID; + if (context.terminalChildSessionIds.has(sessionId)) break; + context.terminalChildSessionIds.add(sessionId); + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId, + itemId: sessionId, + raw: event, + })), + type: "task.completed", + payload: { + taskId: RuntimeTaskId.make(sessionId), + taskType: "local_agent", + status: "completed", + summary: "Completed", + }, + }); + } + break; + } if (event.properties.status.type === "busy" || event.properties.status.type === "retry") { if (turnId === undefined) { break; @@ -2989,7 +3218,9 @@ export function makeOpenCodeAdapter( resolvedRequestIds: new Set(), autoRepliedRequestIds: new Set(), emittedTerminalRequestIds: new Set(), + terminalChildSessionIds: new Set(), requestRelationRetries: new Map(), + sessionRelationRetries: new Map(), pendingPermissions: new Map(), pendingQuestions: new Map(), textPartsByMessageId: new Map(),