diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 9b5b56309749..d4330d19641c 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -4784,6 +4784,66 @@ describe("ProviderRuntimeIngestion", () => { ).toBe("# Plan title"); }); + it("treats OpenCode child task status as background liveness without moving the idle clock", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const provider = ProviderDriverKind.make("opencode"); + const threadId = asThreadId("thread-1"); + + const childStatus = ( + eventId: string, + taskId: string, + status: "running" | "idle", + createdAt: string, + ) => ({ + type: "task.updated" as const, + eventId: asEventId(eventId), + provider, + createdAt, + threadId, + payload: { + taskId, + status, + taskType: "subagent", + description: taskId, + title: taskId, + }, + }); + + await harness.emitAndDrain([ + childStatus("evt-child-a-running", "ses_a", "running", now), + childStatus("evt-child-b-running", "ses_b", "running", now), + ]); + const live = await harness.readThreadShell(); + expect(live.backgroundLiveness).toBe("working"); + expect(live.session?.activeTurnId ?? null).toBeNull(); + expect(live.session?.updatedAt).toBe(now); + + await harness.emitAndDrain([childStatus("evt-child-a-idle", "ses_a", "idle", now)]); + expect((await harness.readThreadShell()).backgroundLiveness).toBe("working"); + + await harness.emitAndDrain([childStatus("evt-child-b-idle", "ses_b", "idle", now)]); + expect((await harness.readThreadShell()).backgroundLiveness).toBeNull(); + + await harness.emitAndDrain([ + childStatus("evt-child-a-resumed", "ses_a", "running", "2026-01-01T00:00:01.000Z"), + ]); + expect((await harness.readThreadShell()).backgroundLiveness).toBe("working"); + expect((await harness.readThreadShell()).session?.updatedAt).toBe(now); + + await harness.emitAndDrain([ + { + type: "session.exited", + eventId: asEventId("evt-opencode-child-session-exited"), + provider, + createdAt: "2026-01-01T00:00:02.000Z", + threadId, + payload: {}, + }, + ]); + expect((await harness.readThreadShell()).backgroundLiveness).toBeNull(); + }); + it("titles task activities with the task description, including on completion", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index baded8094ba6..0bccf319e5e9 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -47,6 +47,9 @@ import { isSameOpenCodeDirectory, makeOpenCodeAdapter, mergeOpenCodeAssistantText, + nextOpenCodeChildLivenessStatus, + openCodeChildSessionLivenessStatus, + openCodeChildSessionTaskUpdate, } from "./OpenCodeAdapter.ts"; import { symlinksSupported } from "@t3tools/shared/testing/symlinks"; @@ -666,6 +669,41 @@ const questionRequest = (id: string, sessionID: string): QuestionRequest => ({ ], }); +it.effect("maps related OpenCode child session status onto background liveness", () => + Effect.sync(() => { + NodeAssert.equal(openCodeChildSessionLivenessStatus("busy"), "running"); + NodeAssert.equal(openCodeChildSessionLivenessStatus("retry"), "running"); + NodeAssert.equal(openCodeChildSessionLivenessStatus("idle"), "idle"); + NodeAssert.equal(openCodeChildSessionLivenessStatus("paused"), undefined); + + NodeAssert.equal(nextOpenCodeChildLivenessStatus(undefined, "running"), "running"); + NodeAssert.equal(nextOpenCodeChildLivenessStatus("running", "running"), undefined); + NodeAssert.equal(nextOpenCodeChildLivenessStatus("running", "idle"), "idle"); + NodeAssert.equal(nextOpenCodeChildLivenessStatus(undefined, "idle"), undefined); + NodeAssert.equal(nextOpenCodeChildLivenessStatus("idle", "idle"), undefined); + NodeAssert.equal(nextOpenCodeChildLivenessStatus("idle", "running"), "running"); + + NodeAssert.deepEqual( + openCodeChildSessionTaskUpdate({ + sessionId: "ses_child", + status: "running", + title: " Research agent ", + }), + { + taskId: "ses_child", + status: "running", + taskType: "subagent", + description: "Research agent", + title: "Research agent", + }, + ); + NodeAssert.equal( + openCodeChildSessionTaskUpdate({ sessionId: "ses_child", status: "idle" }).description, + "Background agent", + ); + }), +); + it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { it.effect("reuses a configured OpenCode server URL instead of spawning a local server", () => Effect.gen(function* () { @@ -4219,6 +4257,275 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect( + "records related child session status as background liveness without touching the parent turn", + () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-child-liveness"); + const parentSessionId = "http://127.0.0.1:9999/session"; + const enqueue = makeOpenCodeEventQueue(); + const updatesFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "task.updated"), + Stream.take(5), + Stream.runCollect, + Effect.forkChild, + ); + const completedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.runHead, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + + enqueue({ + type: "session.status", + properties: { sessionID: "ses_research", status: { type: "busy" } }, + }); + enqueue({ + type: "session.created", + properties: { + info: { + id: "ses_research", + parentID: parentSessionId, + title: "Research agent", + }, + }, + }); + enqueue({ + type: "session.created", + properties: { + info: { + id: "ses_unrelated", + parentID: "ses_other_parent", + title: "Someone else", + }, + }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_unrelated", status: { type: "busy" } }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_research", status: { type: "busy" } }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_research", status: { type: "busy" } }, + }); + enqueue({ + type: "session.status", + properties: { + sessionID: "ses_research", + status: { type: "retry", attempt: 2, message: "rate limit", next: 10 }, + }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_research", status: { type: "paused" } }, + }); + enqueue({ + type: "session.created", + properties: { + info: { + id: "ses_nested", + parentID: "ses_research", + title: "Child session - 2026-09-24T21:48:38.700Z", + }, + }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_nested", status: { type: "busy" } }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_research", status: { type: "idle" } }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_research", status: { type: "idle" } }, + }); + enqueue({ + type: "session.deleted", + properties: { info: { id: "ses_nested" } }, + }); + enqueue({ + type: "session.status", + properties: { + sessionID: "ses_research", + status: { type: "retry", attempt: 3, message: "again", next: 20 }, + }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: parentSessionId, status: { type: "idle" } }, + }); + + const updates = Array.from( + yield* Fiber.join(updatesFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.deepEqual( + updates.map((event) => + event.type === "task.updated" + ? { + taskId: event.payload.taskId, + status: event.payload.status, + description: event.payload.description, + title: event.payload.title, + taskType: event.payload.taskType, + turnId: event.turnId, + } + : null, + ), + [ + { + taskId: "ses_research", + status: "running", + description: "Research agent", + title: "Research agent", + taskType: "subagent", + turnId: undefined, + }, + { + taskId: "ses_nested", + status: "running", + description: "Background agent", + title: "Background agent", + taskType: "subagent", + turnId: undefined, + }, + { + taskId: "ses_research", + status: "idle", + description: "Research agent", + title: "Research agent", + taskType: "subagent", + turnId: undefined, + }, + { + taskId: "ses_nested", + status: "idle", + description: "Background agent", + title: "Background agent", + taskType: "subagent", + turnId: undefined, + }, + { + taskId: "ses_research", + status: "running", + description: "Research agent", + title: "Research agent", + taskType: "subagent", + turnId: undefined, + }, + ], + ); + NodeAssert.equal(completedFiber.pollUnsafe(), undefined); + const session = (yield* adapter.listSessions()).find( + (candidate) => candidate.threadId === threadId, + ); + NodeAssert.equal(session?.activeTurnId, undefined); + NodeAssert.notEqual(session?.status, "running"); + + yield* Fiber.interrupt(completedFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("keeps the parent turn running when a related child session goes idle", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-child-liveness-turn"); + const parentSessionId = "http://127.0.0.1:9999/session"; + runtimeMock.state.sessionStatus = "busy"; + const enqueue = makeOpenCodeEventQueue(); + const completedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.runHead, + Effect.forkChild, + ); + const childUpdatesFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "task.updated"), + Stream.take(2), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turn = yield* adapter.sendTurn({ + threadId, + input: "Run the suite in the background", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + enqueue({ + type: "session.created", + properties: { + info: { + id: "ses_live", + parentID: parentSessionId, + title: "Browser run", + }, + }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_live", status: { type: "busy" } }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_live", status: { type: "idle" } }, + }); + + const childUpdates = Array.from( + yield* Fiber.join(childUpdatesFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.deepEqual( + childUpdates.map((event) => + event.type === "task.updated" ? event.payload.status : undefined, + ), + ["running", "idle"], + ); + NodeAssert.equal( + childUpdates.every((event) => event.turnId === undefined), + true, + ); + NodeAssert.equal(completedFiber.pollUnsafe(), undefined); + NodeAssert.equal( + (yield* adapter.listSessions()).find((candidate) => candidate.threadId === threadId) + ?.activeTurnId, + turn.turnId, + ); + + enqueue({ + type: "session.status", + properties: { sessionID: parentSessionId, status: { type: "idle" } }, + }); + const completed = Option.getOrThrow( + yield* Fiber.join(completedFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.equal( + completed.type === "turn.completed" ? completed.turnId : undefined, + turn.turnId, + ); + + yield* adapter.stopSession(threadId); + }), + ); + it.effect.each([ { name: "a doom-loop ask on the parent session", diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 04535edbd3af..33468462a970 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, @@ -319,6 +320,94 @@ function isOpenCodeChildRequestEvent(event: OpenCodeSubscribedEvent): boolean { } } +/** + * Liveness status a related OpenCode child session contributes to the + * thread's background-work registry. + * + * `running` keeps the provider session alive. `idle` drops that child. + */ +export type OpenCodeChildLiveness = "running" | "idle"; + +/** + * Label used when a related child has no real session title. + * + * OpenCode's placeholder titles (`Child session - `) are not + * names, so the task row falls back to this instead of locking onto them. + */ +const OPENCODE_CHILD_SESSION_DESCRIPTION = "Background agent"; + +/** + * Map an OpenCode child `session.status` type onto the task status the + * background-liveness registry already understands. + * + * `busy` and `retry` are in-flight provider work and count as `running`. + * `idle` clears that child. Any other type is ignored so a future OpenCode + * status cannot drop a child that is still running. + */ +export function openCodeChildSessionLivenessStatus( + statusType: string, +): OpenCodeChildLiveness | undefined { + switch (statusType) { + case "busy": + case "retry": + return "running"; + case "idle": + return "idle"; + default: + return undefined; + } +} + +/** + * The next liveness status to forward, or `undefined` when this observation + * does not change the registry. + * + * Repeated `busy`/`retry` while a child is already live must not emit another + * task row. `idle` only clears a child that was live. A later `busy` or + * `retry` after `idle` starts it again. + */ +export function nextOpenCodeChildLivenessStatus( + previous: OpenCodeChildLiveness | undefined, + next: OpenCodeChildLiveness, +): OpenCodeChildLiveness | undefined { + if (previous === next) { + return undefined; + } + if (next === "idle" && previous !== "running") { + return undefined; + } + return next; +} + +/** + * `task.updated` payload for one related OpenCode child session. + * + * The task id is the child session id, so one child's idle does not clear + * its siblings. The registry treats `running` as live work and `idle` as + * not live, which is what keeps the session reaper from stopping a settled + * thread while background subagents are still streaming. + */ +export function openCodeChildSessionTaskUpdate(input: { + readonly sessionId: string; + readonly status: OpenCodeChildLiveness; + readonly title?: string | undefined; +}): { + readonly taskId: RuntimeTaskId; + readonly status: OpenCodeChildLiveness; + readonly taskType: "subagent"; + readonly description: string; + readonly title: string; +} { + const label = input.title?.trim() || OPENCODE_CHILD_SESSION_DESCRIPTION; + return { + taskId: RuntimeTaskId.make(input.sessionId), + status: input.status, + taskType: "subagent", + description: label, + title: label, + }; +} + 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,6 +432,13 @@ interface OpenCodeSessionContext { readonly directory: string; openCodeSessionId: string; readonly relatedSessionIds: Set; + /** + * Last liveness status forwarded for a related child. Repeated busy/retry + * while a child is already live does not emit another task row. + */ + readonly childSessionLiveness: Map; + /** Real titles for related children. Placeholder OpenCode titles are omitted. */ + readonly childSessionTitles: Map; readonly resolvedRequestIds: Set; readonly autoRepliedRequestIds: Set; readonly emittedTerminalRequestIds: Set; @@ -1674,6 +1770,78 @@ export function makeOpenCodeAdapter( } }; + /** + * Remember a related child's display title when OpenCode assigned a real + * one. Placeholder titles are ignored so a later liveness row does not + * lock onto `Child session - `. + */ + const rememberOpenCodeChildSessionTitle = ( + context: OpenCodeSessionContext, + sessionId: string, + title: string, + ) => { + const trimmed = trimText(title); + if (!trimmed || isOpenCodeDefaultTitle(trimmed)) { + return; + } + context.childSessionTitles.set(sessionId, trimmed); + }; + + /** + * Forward one child-session liveness transition as `task.updated`. + * + * OpenCode background subagents keep the provider process busy after the + * parent turn settles, but they never emit `task.*`. The reaper already + * skips a thread while `backgroundLiveness` is set; this is the signal + * that fills it. Child status does not touch the parent turn. + */ + const recordOpenCodeChildSessionLiveness = Effect.fn("recordOpenCodeChildSessionLiveness")( + function* ( + context: OpenCodeSessionContext, + sessionId: string, + next: OpenCodeChildLiveness, + raw: unknown, + ) { + const status = nextOpenCodeChildLivenessStatus( + context.childSessionLiveness.get(sessionId), + next, + ); + if (status === undefined) { + return; + } + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + raw, + })), + type: "task.updated", + payload: openCodeChildSessionTaskUpdate({ + sessionId, + status, + title: context.childSessionTitles.get(sessionId), + }), + }); + context.childSessionLiveness.set(sessionId, status); + }, + ); + + /** + * Drop every tracked child from the liveness registry. + * + * Used when the parent session is replaced (rewind/fork) and the previous + * children no longer belong to this thread. Only children that were live + * emit `idle`; the rest were already clear. + */ + const clearOpenCodeChildSessionLiveness = Effect.fn("clearOpenCodeChildSessionLiveness")( + function* (context: OpenCodeSessionContext, raw: unknown) { + for (const sessionId of context.childSessionLiveness.keys()) { + yield* recordOpenCodeChildSessionLiveness(context, sessionId, "idle", raw); + } + context.childSessionLiveness.clear(); + context.childSessionTitles.clear(); + }, + ); + const isRelatedOpenCodeSession = Effect.fn("isRelatedOpenCodeSession")(function* ( context: OpenCodeSessionContext, candidateSessionId: string, @@ -2175,6 +2343,13 @@ export function makeOpenCodeAdapter( yield* run.pipe(Effect.forkIn(context.sessionScope)); }); + /** + * Route one OpenCode subscription event into runtime events. + * + * Parent `session.status` still owns turn admission and completion. + * Related child `session.status` is forwarded as background-task liveness + * and does not settle the parent turn. + */ const handleSubscribedEvent = Effect.fn("handleSubscribedEvent")(function* ( context: OpenCodeSessionContext, event: OpenCodeSubscribedEvent, @@ -2225,9 +2400,16 @@ export function makeOpenCodeAdapter( const session = event.properties.info; if (session.parentID && context.relatedSessionIds.has(session.parentID)) { addRelatedOpenCodeSession(context, session.id); + rememberOpenCodeChildSessionTitle(context, session.id, session.title); } } else if (event.type === "session.deleted") { - context.relatedSessionIds.delete(event.properties.info.id); + const deletedId = event.properties.info.id; + if (deletedId !== context.openCodeSessionId && context.relatedSessionIds.has(deletedId)) { + yield* recordOpenCodeChildSessionLiveness(context, deletedId, "idle", event); + } + context.relatedSessionIds.delete(deletedId); + context.childSessionLiveness.delete(deletedId); + context.childSessionTitles.delete(deletedId); } const payloadSessionId = openCodeEventSessionId(event); @@ -2260,7 +2442,14 @@ export function makeOpenCodeAdapter( payloadSessionId !== undefined && isOpenCodeChildRequestEvent(event) && (context.relatedSessionIds.has(payloadSessionId) || isKnownPendingTerminalEvent); - if (!isParentEvent && !isChildRequestEvent) { + // Related child `session.status` is the liveness signal for OpenCode + // background subagents. Permission/question events are not. + const isRelatedChildStatusEvent = + event.type === "session.status" && + payloadSessionId !== undefined && + !isParentEvent && + context.relatedSessionIds.has(payloadSessionId); + if (!isParentEvent && !isChildRequestEvent && !isRelatedChildStatusEvent) { return; } @@ -2587,6 +2776,13 @@ export function makeOpenCodeAdapter( } case "session.status": { + if (!isParentEvent && payloadSessionId !== undefined) { + const liveness = openCodeChildSessionLivenessStatus(event.properties.status.type); + if (liveness !== undefined) { + yield* recordOpenCodeChildSessionLiveness(context, payloadSessionId, liveness, event); + } + break; + } if (event.properties.status.type === "busy" || event.properties.status.type === "retry") { if (turnId === undefined) { break; @@ -3010,6 +3206,8 @@ export function makeOpenCodeAdapter( directory, openCodeSessionId: started.openCodeSession.id, relatedSessionIds: new Set([started.openCodeSession.id]), + childSessionLiveness: new Map(), + childSessionTitles: new Map(), resolvedRequestIds: new Set(), autoRepliedRequestIds: new Set(), emittedTerminalRequestIds: new Set(), @@ -3912,6 +4110,11 @@ export function makeOpenCodeAdapter( }, ); + /** + * Rewind the thread by forking a new OpenCode session at the retained + * boundary. Children of the session being replaced are cleared from + * background liveness because they no longer belong to this thread. + */ const rollbackThread: OpenCodeAdapterShape["rollbackThread"] = Effect.fn("rollbackThread")( function* (threadId, numTurns) { const context = yield* ensureSessionContext(sessions, threadId); @@ -3972,6 +4175,7 @@ export function makeOpenCodeAdapter( }), ).pipe(Effect.mapError(toRequestError)); yield* clearPendingOpenCodeRequests(context, { type: "session.fork" }); + yield* clearOpenCodeChildSessionLiveness(context, { type: "session.fork" }); context.openCodeSessionId = forkedSessionId; context.relatedSessionIds.clear(); context.relatedSessionIds.add(forkedSessionId);