From 40f57f202b584f0359ffb7b3a7b3730d943c4dc4 Mon Sep 17 00:00:00 2001 From: macodev00 <273427913+macodev00@users.noreply.github.com> Date: Mon, 28 Sep 2026 07:18:29 +0000 Subject: [PATCH] fix(server): keep OpenCode sessions alive while child sessions run OpenCode background subagents stay busy after the parent turn settles, but they never emit task events, so the reaper treated the thread as idle and stopped the provider session. Related child session.status is now quiet liveness for that reaper check. busy and retry hold the session; idle and deletion release it. The signal stays inside the reaper and leaves task activity unchanged. Fixes #13514 --- .../provider/Layers/OpenCodeAdapter.test.ts | 387 ++++++++++++++++++ .../src/provider/Layers/OpenCodeAdapter.ts | 61 ++- .../Layers/ProviderSessionReaper.test.ts | 123 +++++- .../provider/Layers/ProviderSessionReaper.ts | 34 ++ .../OpenCodeChildSessionLiveness.test.ts | 54 +++ .../provider/OpenCodeChildSessionLiveness.ts | 127 ++++++ 6 files changed, 784 insertions(+), 2 deletions(-) create mode 100644 apps/server/src/provider/OpenCodeChildSessionLiveness.test.ts create mode 100644 apps/server/src/provider/OpenCodeChildSessionLiveness.ts diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index baded8094ba6..ec8338073f27 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -12,6 +12,7 @@ import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Path from "effect/Path"; +import * as Queue from "effect/Queue"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; @@ -29,6 +30,7 @@ import { OpenCodeSettings, ProviderDriverKind, ProviderInstanceId, + type ProviderRuntimeEvent, ThreadId, } from "@t3tools/contracts"; import { createModelSelection } from "@t3tools/shared/model"; @@ -42,6 +44,10 @@ import { OpenCodeRuntimeError, type OpenCodeRuntimeShape, } from "../opencodeRuntime.ts"; +import { + hasLiveOpenCodeChildSessions, + resetOpenCodeChildSessionLiveness, +} from "../OpenCodeChildSessionLiveness.ts"; import { isOpenCodeNotFound, isSameOpenCodeDirectory, @@ -618,6 +624,7 @@ const OpenCodeAdapterTestLayer = Layer.effect( beforeEach(() => { runtimeMock.reset(); + resetOpenCodeChildSessionLiveness(); }); const advanceTestClock = (ms: number) => @@ -654,6 +661,84 @@ const permissionRequest = (id: string, sessionID: string): PermissionRequest => always: [], }); +const OPENCODE_PARENT_SESSION_ID = "http://127.0.0.1:9999/session"; + +/** + * Single consumer for one thread's runtime events. + * + * Child `session.status` writes quiet liveness and emits no runtime event, so + * a following parent title is the barrier that proves those updates landed. + * One consumer owns the adapter queue; a second subscriber would drop events. + */ +function collectOpenCodeThreadEvents(adapter: OpenCodeAdapterShape, threadId: ThreadId) { + /** + * Subscribe once and expose title and turn-completion barriers. + * + * The parent title is the barrier because quiet child status emits nothing. + */ + function* subscribe() { + const seen: Array = []; + const titles = yield* Queue.unbounded(); + const completedTurnIds = yield* Queue.unbounded(); + yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId), + Stream.runForEach((event) => + Effect.gen( + /** + * Record one runtime event and publish the barriers tests wait on. + */ + function* recordRuntimeEvent() { + seen.push(event); + if (event.type === "thread.metadata.updated" && event.payload.name) { + yield* Queue.offer(titles, event.payload.name); + } + if (event.type === "turn.completed") { + yield* Queue.offer(completedTurnIds, event.turnId); + } + }, + ), + ), + Effect.forkChild, + ); + + /** + * Wait until a parent session title has been mirrored onto the thread. + */ + const waitForTitle = (title: string) => + Stream.fromQueue(titles).pipe( + Stream.filter((candidate) => candidate === title), + Stream.take(1), + Stream.runDrain, + Effect.timeout("2 seconds"), + ); + + /** + * Wait until the parent turn completes and return that turn id. + */ + const waitForTurnCompleted = () => + Queue.take(completedTurnIds).pipe(Effect.timeout("2 seconds")); + + return { seen, waitForTitle, waitForTurnCompleted }; + } + + return Effect.gen(subscribe); +} + +/** + * Parent `session.updated` used as a processing barrier after quiet child events. + */ +function openCodeParentTitle(sessionId: string, title: string) { + return { + type: "session.updated" as const, + properties: { + info: { + id: sessionId, + title, + }, + }, + }; +} + const questionRequest = (id: string, sessionID: string): QuestionRequest => ({ id, sessionID, @@ -667,6 +752,308 @@ const questionRequest = (id: string, sessionID: string): QuestionRequest => ({ }); it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { + /** + * Assert related child `session.status` updates quiet reaper liveness only. + * + * Unrelated sessions are ignored. Busy, retry, and paused must not emit + * task activity or complete a parent turn. Idle and deletion release only + * that child, and stopping the session releases the thread. + */ + function* recordRelatedChildSessionLiveness() { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-child-liveness"); + const enqueue = makeOpenCodeEventQueue(); + const { seen, waitForTitle } = yield* collectOpenCodeThreadEvents(adapter, threadId); + + 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: OPENCODE_PARENT_SESSION_ID, + 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(openCodeParentTitle(OPENCODE_PARENT_SESSION_ID, "Children are running")); + + yield* waitForTitle("Children are running"); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), true); + const session = (yield* adapter.listSessions()).find( + (candidate) => candidate.threadId === threadId, + ); + NodeAssert.equal(session?.activeTurnId, undefined); + NodeAssert.notEqual(session?.status, "running"); + + enqueue({ + type: "session.status", + properties: { sessionID: "ses_research", status: { type: "idle" } }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_research", status: { type: "idle" } }, + }); + enqueue(openCodeParentTitle(OPENCODE_PARENT_SESSION_ID, "Nested child still running")); + yield* waitForTitle("Nested child still running"); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), true); + + enqueue({ + type: "session.deleted", + properties: { info: { id: "ses_nested" } }, + }); + enqueue(openCodeParentTitle(OPENCODE_PARENT_SESSION_ID, "Nested child deleted")); + yield* waitForTitle("Nested child deleted"); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), false); + + enqueue({ + type: "session.status", + properties: { + sessionID: "ses_research", + status: { type: "retry", attempt: 3, message: "again", next: 20 }, + }, + }); + enqueue(openCodeParentTitle(OPENCODE_PARENT_SESSION_ID, "Research resumed")); + yield* waitForTitle("Research resumed"); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), true); + NodeAssert.deepEqual( + seen.filter((event) => event.type.startsWith("task.")), + [], + ); + + yield* adapter.stopSession(threadId); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), false); + } + + it.effect( + "records related child session.status as quiet reaper liveness without task activity", + /** + * Run the quiet-liveness regression against the OpenCode adapter layer. + */ + () => Effect.gen(recordRelatedChildSessionLiveness), + ); + + /** + * Assert a child going idle does not complete the parent turn. + * + * The parent stays running until its own `session.status` is idle, and the + * child transition emits no task activity. + */ + function* keepParentTurnRunningWhileChildGoesIdle() { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-child-liveness-turn"); + runtimeMock.state.sessionStatus = "busy"; + const enqueue = makeOpenCodeEventQueue(); + const { seen, waitForTitle, waitForTurnCompleted } = yield* collectOpenCodeThreadEvents( + adapter, + threadId, + ); + + 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: OPENCODE_PARENT_SESSION_ID, + title: "Browser run", + }, + }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_live", status: { type: "busy" } }, + }); + enqueue(openCodeParentTitle(OPENCODE_PARENT_SESSION_ID, "Child is busy")); + yield* waitForTitle("Child is busy"); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), true); + NodeAssert.equal( + seen.some((event) => event.type === "turn.completed"), + false, + ); + + enqueue({ + type: "session.status", + properties: { sessionID: "ses_live", status: { type: "idle" } }, + }); + enqueue(openCodeParentTitle(OPENCODE_PARENT_SESSION_ID, "Child is idle")); + yield* waitForTitle("Child is idle"); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), false); + NodeAssert.equal( + seen.some((event) => event.type === "turn.completed"), + false, + ); + NodeAssert.equal( + (yield* adapter.listSessions()).find((candidate) => candidate.threadId === threadId) + ?.activeTurnId, + turn.turnId, + ); + + enqueue({ + type: "session.status", + properties: { sessionID: OPENCODE_PARENT_SESSION_ID, status: { type: "idle" } }, + }); + NodeAssert.equal(yield* waitForTurnCompleted(), turn.turnId); + NodeAssert.deepEqual( + seen.filter((event) => event.type.startsWith("task.")), + [], + ); + + yield* adapter.stopSession(threadId); + } + + it.effect( + "keeps the parent turn running when a related child session goes idle", + /** + * Run the parent-turn regression against the OpenCode adapter layer. + */ + () => Effect.gen(keepParentTurnRunningWhileChildGoesIdle), + ); + + /** + * Assert rewind and stop release quiet child liveness. + * + * Children of the replaced session must not hold the forked session, and + * a child of the replacement can hold it again until stop. + */ + function* dropQuietChildLivenessOnRewindOrStop() { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-child-liveness-teardown"); + const enqueue = makeOpenCodeEventQueue(); + const { waitForTitle } = yield* collectOpenCodeThreadEvents(adapter, threadId); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + enqueue({ + type: "session.created", + properties: { + info: { + id: "ses_old", + parentID: OPENCODE_PARENT_SESSION_ID, + title: "Old child", + }, + }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_old", status: { type: "busy" } }, + }); + enqueue(openCodeParentTitle(OPENCODE_PARENT_SESSION_ID, "Old child is live")); + yield* waitForTitle("Old child is live"); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), true); + + runtimeMock.state.messages = [ + { info: { id: "user-1", role: "user" }, parts: [] }, + { + info: { id: "assistant-1", role: "assistant" }, + parts: [{ id: "part-1", type: "text", text: "answer" }], + }, + ]; + yield* adapter.rollbackThread(threadId, 1); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), false); + + const forkedSessionId = `${OPENCODE_PARENT_SESSION_ID}_fork`; + enqueue({ + type: "session.created", + properties: { + info: { + id: "ses_new", + parentID: forkedSessionId, + title: "Replacement child", + }, + }, + }); + enqueue({ + type: "session.status", + properties: { + sessionID: "ses_new", + status: { type: "retry", attempt: 1, message: "again", next: 5 }, + }, + }); + enqueue(openCodeParentTitle(forkedSessionId, "Replacement child is live")); + yield* waitForTitle("Replacement child is live"); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), true); + + yield* adapter.stopSession(threadId); + NodeAssert.equal(hasLiveOpenCodeChildSessions(threadId), false); + } + + it.effect( + "drops quiet child liveness when the OpenCode session is rewound or stopped", + /** + * Run the rewind and stop regression against the OpenCode adapter layer. + */ + () => Effect.gen(dropQuietChildLivenessOnRewindOrStop), + ); + it.effect("reuses a configured OpenCode server URL instead of spawning a local server", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 04535edbd3af..0bf7ee397942 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -59,6 +59,12 @@ import { toOpenCodeQuestionAnswers, type OpenCodeServerConnection, } from "../opencodeRuntime.ts"; +import { + clearOpenCodeChildSession, + clearOpenCodeThreadChildSessions, + noteOpenCodeChildSessionLiveness, + openCodeChildSessionLivenessStatus, +} from "../OpenCodeChildSessionLiveness.ts"; import * as Option from "effect/Option"; const PROVIDER = ProviderDriverKind.make("opencode"); @@ -2175,6 +2181,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` updates quiet reaper liveness only. + * That update leaves the parent turn untouched and emits no task activity. + */ const handleSubscribedEvent = Effect.fn("handleSubscribedEvent")(function* ( context: OpenCodeSessionContext, event: OpenCodeSubscribedEvent, @@ -2227,7 +2240,13 @@ export function makeOpenCodeAdapter( addRelatedOpenCodeSession(context, session.id); } } else if (event.type === "session.deleted") { - context.relatedSessionIds.delete(event.properties.info.id); + const deletedSessionId = event.properties.info.id; + if ( + context.relatedSessionIds.delete(deletedSessionId) && + deletedSessionId !== context.openCodeSessionId + ) { + clearOpenCodeChildSession(context.session.threadId, deletedSessionId); + } } const payloadSessionId = openCodeEventSessionId(event); @@ -2260,6 +2279,21 @@ export function makeOpenCodeAdapter( payloadSessionId !== undefined && isOpenCodeChildRequestEvent(event) && (context.relatedSessionIds.has(payloadSessionId) || isKnownPendingTerminalEvent); + // Related child session.status is the quiet liveness signal for OpenCode + // background subagents. It must not fall through into parent turn + // admission or completion, and it must not emit task activity. + if ( + event.type === "session.status" && + payloadSessionId !== undefined && + !isParentEvent && + context.relatedSessionIds.has(payloadSessionId) + ) { + const liveness = openCodeChildSessionLivenessStatus(event.properties.status.type); + if (liveness !== undefined) { + noteOpenCodeChildSessionLiveness(context.session.threadId, payloadSessionId, liveness); + } + return; + } if (!isParentEvent && !isChildRequestEvent) { return; } @@ -2725,7 +2759,25 @@ export function makeOpenCodeAdapter( } }); + /** + * Subscribe to the OpenCode event stream for one session. + * + * The session scope owns quiet child-session liveness. Closing it on + * stop, unexpected exit, or layer shutdown releases every child held + * for the thread, because those children die with the provider process. + */ const startEventPump = Effect.fn("startEventPump")(function* (context: OpenCodeSessionContext) { + /** + * Release quiet child liveness when this session scope closes. + * + * Stop, unexpected exit, and layer shutdown close the scope. Those + * children die with the provider process, so they must not keep the + * inactivity reaper holding the thread. + */ + function releaseOpenCodeChildLiveness() { + clearOpenCodeThreadChildSessions(context.session.threadId); + } + yield* Scope.addFinalizer(context.sessionScope, Effect.sync(releaseOpenCodeChildLiveness)); // One AbortController per session scope. The finalizer fires when // the scope closes (explicit stop, unexpected exit, or layer // shutdown) and cancels the in-flight `event.subscribe` fetch so @@ -3912,6 +3964,12 @@ export function makeOpenCodeAdapter( }, ); + /** + * Rewind the thread by forking a new OpenCode session at the retained boundary. + * + * Children of the session being replaced are released from quiet liveness. + * They belong to the discarded session and must not keep the replacement alive. + */ const rollbackThread: OpenCodeAdapterShape["rollbackThread"] = Effect.fn("rollbackThread")( function* (threadId, numTurns) { const context = yield* ensureSessionContext(sessions, threadId); @@ -3972,6 +4030,7 @@ export function makeOpenCodeAdapter( }), ).pipe(Effect.mapError(toRequestError)); yield* clearPendingOpenCodeRequests(context, { type: "session.fork" }); + clearOpenCodeThreadChildSessions(threadId); context.openCodeSessionId = forkedSessionId; context.relatedSessionIds.clear(); context.relatedSessionIds.add(forkedSessionId); diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts index 7b1fec90f867..bbdc3cc4268c 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts @@ -25,6 +25,11 @@ import { ProviderSessionReaper } from "../Services/ProviderSessionReaper.ts"; import { ProviderService, type ProviderServiceShape } from "../Services/ProviderService.ts"; import { ProviderSessionDirectoryLive } from "./ProviderSessionDirectory.ts"; import { makeProviderSessionReaperLive } from "./ProviderSessionReaper.ts"; +import { + clearOpenCodeThreadChildSessions, + noteOpenCodeChildSessionLiveness, + resetOpenCodeChildSessionLiveness, +} from "../OpenCodeChildSessionLiveness.ts"; const defaultModelSelection = { instanceId: ProviderInstanceId.make("codex"), @@ -56,13 +61,20 @@ const drainFibers = Effect.forEach(Array.from({ length: 10 }), () => Effect.yiel const unsupported = () => Effect.die(new Error("Unsupported provider call in test")) as never; +/** + * Build a shell snapshot for the reaper tests. + * + * Session timestamps and provider names are the inputs the inactivity sweep + * reads. `opencode` is included so a live child can be distinguished from + * Claude and Codex bindings. + */ function makeReadModel( threads: ReadonlyArray<{ readonly id: ThreadId; readonly session: { readonly threadId: ThreadId; readonly status: "starting" | "running" | "ready" | "interrupted" | "stopped" | "error"; - readonly providerName: "codex" | "claudeAgent"; + readonly providerName: "codex" | "claudeAgent" | "opencode"; readonly runtimeMode: "approval-required" | "full-access" | "auto-accept-edits"; readonly activeTurnId: TurnId | null; readonly lastError: string | null; @@ -128,6 +140,7 @@ describe("ProviderSessionReaper", () => { let scope: Scope.Closeable | null = null; afterEach(async () => { + resetOpenCodeChildSessionLiveness(); if (scope) { await Effect.runPromise(Scope.close(scope, Exit.void)); } @@ -734,4 +747,112 @@ describe("ProviderSessionReaper", () => { reapedThreadId, ]); }); + + it("does not reap a stale OpenCode session while a child session is still live" /** + * Hold a stale OpenCode binding while a child is live, and still reap the others. + * + * Claude is never held by this signal. Idle of one child leaves its sibling + * holding the binding. Clearing the thread lets the next sweep reap it. + */, async function reapSkipsLiveOpenCodeChildSession() { + const liveThreadId = ThreadId.make("thread-reaper-opencode-child-live"); + const settledOpenCodeThreadId = ThreadId.make("thread-reaper-opencode-settled"); + const claudeThreadId = ThreadId.make("thread-reaper-claude-not-held"); + const updatedAt = "2026-04-14T00:00:00.000Z"; + /** + * Settled thread shell whose session timestamp is already past the idle window. + */ + function shell(threadId: ThreadId, providerName: "opencode" | "claudeAgent") { + return { + id: threadId, + session: { + threadId, + status: "ready" as const, + providerName, + runtimeMode: "full-access" as const, + activeTurnId: null, + lastError: null, + updatedAt, + }, + }; + } + const harness = await createHarness({ + readModel: makeReadModel([ + shell(claudeThreadId, "claudeAgent"), + shell(liveThreadId, "opencode"), + shell(settledOpenCodeThreadId, "opencode"), + ]), + }); + const repository = await runtime!.runPromise( + Effect.service(ProviderSessionRuntime.ProviderSessionRuntimeRepository), + ); + /** + * Persist a live binding old enough for the next sweep to consider it idle. + */ + function seed(threadId: ThreadId, providerName: "opencode" | "claudeAgent") { + return runtime!.runPromise( + repository.upsert({ + threadId, + providerName, + providerInstanceId: null, + adapterKey: providerName, + runtimeMode: "full-access", + status: "running", + lastSeenAt: updatedAt, + resumeCursor: { opaque: `resume-${threadId}` }, + runtimePayload: null, + }), + ); + } + await seed(claudeThreadId, "claudeAgent"); + await seed(liveThreadId, "opencode"); + await seed(settledOpenCodeThreadId, "opencode"); + + // busy and retry both record "running". A second child stays held after + // its sibling goes idle. Claude is not held by this OpenCode-only signal. + noteOpenCodeChildSessionLiveness(liveThreadId, "ses_a", "running"); + noteOpenCodeChildSessionLiveness(liveThreadId, "ses_b", "running"); + noteOpenCodeChildSessionLiveness(claudeThreadId, "ses_claude", "running"); + + await sweepAt(Date.parse(updatedAt) + 1_000); + expect(harness.stopSession.mock.calls.map(([request]) => request.threadId)).toEqual([ + claudeThreadId, + settledOpenCodeThreadId, + ]); + + /** + * Mark a binding stopped so a later sweep does not stop it again. + */ + function markStopped(threadId: ThreadId, providerName: "opencode" | "claudeAgent") { + return runtime!.runPromise( + repository.upsert({ + threadId, + providerName, + providerInstanceId: null, + adapterKey: providerName, + runtimeMode: "full-access", + status: "stopped", + lastSeenAt: updatedAt, + resumeCursor: { opaque: `resume-${threadId}` }, + runtimePayload: null, + }), + ); + } + await markStopped(claudeThreadId, "claudeAgent"); + await markStopped(settledOpenCodeThreadId, "opencode"); + + noteOpenCodeChildSessionLiveness(liveThreadId, "ses_a", "idle"); + await sweepAt(Date.parse(updatedAt) + 1_000); + expect(harness.stopSession.mock.calls.map(([request]) => request.threadId)).toEqual([ + claudeThreadId, + settledOpenCodeThreadId, + ]); + + clearOpenCodeThreadChildSessions(liveThreadId); + await sweepAt(Date.parse(updatedAt) + 1_000); + expect(harness.stopSession.mock.calls.map(([request]) => request.threadId)).toEqual([ + claudeThreadId, + settledOpenCodeThreadId, + liveThreadId, + ]); + }); }); diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.ts index bf8199f80eac..f417ad4aa5bb 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.ts @@ -12,6 +12,7 @@ import { type ProviderSessionReaperShape, } from "../Services/ProviderSessionReaper.ts"; import { forkParked } from "../../serverActivation.ts"; +import { openCodeInactivityHeldByChildSession } from "../OpenCodeChildSessionLiveness.ts"; import { ProviderService } from "../Services/ProviderService.ts"; const DEFAULT_INACTIVITY_THRESHOLD_MS = 30 * 60 * 1000; @@ -22,6 +23,15 @@ export interface ProviderSessionReaperLiveOptions { readonly sweepIntervalMs?: number; } +/** + * Build the provider-session reaper. + * + * A sweep stops a live binding once user-facing activity is older than the + * idle window, unless the thread still has an active turn, background-task + * liveness, or a live OpenCode child session. Child `session.status` holds + * the binding quietly: inactivity has to ignore it, and it does not add task + * activity. + */ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) => Effect.gen(function* () { const providerService = yield* ProviderService; @@ -34,6 +44,13 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = ); const sweepIntervalMs = Math.max(1, options?.sweepIntervalMs ?? DEFAULT_SWEEP_INTERVAL_MS); + /** + * Stop live bindings whose user-facing activity is past the idle window. + * + * A binding stays up while its turn is active, while background-task + * liveness is set, or while an OpenCode child session bound to the + * thread is still alive. + */ const sweep = Effect.gen(function* () { // Stopped rows stay for their resume cursors and far outnumber live // ones, so the query skips them. @@ -92,6 +109,23 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = continue; } + // OpenCode child sessions keep the provider process working after the + // parent turn settles. They never emit task lifecycle events, so + // backgroundLiveness stays empty and the idle clock keeps advancing + // from the last user-facing activity. Inactivity has to ignore that + // clock while a related child session.status is still busy or retry: + // stopping the session aborts those children. Idle, deletion, and + // session teardown release the hold. The signal is quiet reaper state, + // separate from task rows and the sidebar liveness pill. + if (openCodeInactivityHeldByChildSession(binding.provider, binding.threadId)) { + yield* Effect.logDebug("provider.session.reaper.skipped-opencode-child-session", { + threadId: binding.threadId, + provider: binding.provider, + idleDurationMs, + }); + continue; + } + const reaped = yield* providerService.stopSession({ threadId: binding.threadId }).pipe( Effect.tap(() => Effect.logInfo("provider.session.reaped", { diff --git a/apps/server/src/provider/OpenCodeChildSessionLiveness.test.ts b/apps/server/src/provider/OpenCodeChildSessionLiveness.test.ts new file mode 100644 index 000000000000..fc55e48633da --- /dev/null +++ b/apps/server/src/provider/OpenCodeChildSessionLiveness.test.ts @@ -0,0 +1,54 @@ +import { describe, expect, it, beforeEach } from "vite-plus/test"; + +import { + clearOpenCodeChildSession, + clearOpenCodeThreadChildSessions, + hasLiveOpenCodeChildSessions, + noteOpenCodeChildSessionLiveness, + openCodeChildSessionLivenessStatus, + openCodeInactivityHeldByChildSession, + resetOpenCodeChildSessionLiveness, +} from "./OpenCodeChildSessionLiveness.ts"; + +describe("OpenCodeChildSessionLiveness", () => { + beforeEach(() => { + resetOpenCodeChildSessionLiveness(); + }); + + it("maps busy and retry onto running, and idle onto a release", () => { + expect(openCodeChildSessionLivenessStatus("busy")).toBe("running"); + expect(openCodeChildSessionLivenessStatus("retry")).toBe("running"); + expect(openCodeChildSessionLivenessStatus("idle")).toBe("idle"); + expect(openCodeChildSessionLivenessStatus("paused")).toBeUndefined(); + }); + + it("holds a thread while any related child is running", () => { + noteOpenCodeChildSessionLiveness("thread-a", "ses_a", "running"); + noteOpenCodeChildSessionLiveness("thread-a", "ses_a", "running"); + noteOpenCodeChildSessionLiveness("thread-a", "ses_b", "running"); + noteOpenCodeChildSessionLiveness("thread-b", "ses_other", "running"); + + expect(hasLiveOpenCodeChildSessions("thread-a")).toBe(true); + expect(openCodeInactivityHeldByChildSession("opencode", "thread-a")).toBe(true); + expect(openCodeInactivityHeldByChildSession("claudeAgent", "thread-a")).toBe(false); + expect(openCodeInactivityHeldByChildSession("opencode", "thread-missing")).toBe(false); + + noteOpenCodeChildSessionLiveness("thread-a", "ses_a", "idle"); + clearOpenCodeChildSession("thread-a", "ses_missing"); + expect(hasLiveOpenCodeChildSessions("thread-a")).toBe(true); + expect(hasLiveOpenCodeChildSessions("thread-b")).toBe(true); + + noteOpenCodeChildSessionLiveness("thread-a", "ses_b", "idle"); + expect(hasLiveOpenCodeChildSessions("thread-a")).toBe(false); + expect(openCodeInactivityHeldByChildSession("opencode", "thread-a")).toBe(false); + + clearOpenCodeThreadChildSessions("thread-b"); + expect(hasLiveOpenCodeChildSessions("thread-b")).toBe(false); + }); + + it("ignores an idle release for a child that was never held", () => { + noteOpenCodeChildSessionLiveness("thread-a", "ses_missing", "idle"); + clearOpenCodeChildSession("thread-a", "ses_missing"); + expect(hasLiveOpenCodeChildSessions("thread-a")).toBe(false); + }); +}); diff --git a/apps/server/src/provider/OpenCodeChildSessionLiveness.ts b/apps/server/src/provider/OpenCodeChildSessionLiveness.ts new file mode 100644 index 000000000000..4f37e6c76a4e --- /dev/null +++ b/apps/server/src/provider/OpenCodeChildSessionLiveness.ts @@ -0,0 +1,127 @@ +/** + * Quiet liveness for OpenCode child sessions. + * + * Background subagents keep the OpenCode process busy after the parent turn + * settles. They do not emit `task.*` events, so the sidebar liveness pill and + * the task activity stream never hear about them. The provider-session reaper + * would otherwise treat that silence as inactivity and stop the session, + * which aborts the children. + * + * This map is the reaper's signal only. `busy` and `retry` hold a child; + * `idle`, deletion, and session teardown release it. The hold is empty after + * a restart, which matches a provider process that did not survive it. + * + * @module provider/OpenCodeChildSessionLiveness + */ + +export type OpenCodeChildSessionLiveness = "running" | "idle"; + +const liveChildrenByThreadId = new Map>(); + +/** + * Map an OpenCode child `session.status` type onto quiet reaper liveness. + * + * `busy` and `retry` are in-flight provider work and count as running. + * `idle` releases that child. Any other type is ignored so an unknown + * OpenCode status cannot drop a child that is still running. + */ +export function openCodeChildSessionLivenessStatus( + statusType: string, +): OpenCodeChildSessionLiveness | undefined { + switch (statusType) { + case "busy": + case "retry": + return "running"; + case "idle": + return "idle"; + default: + return undefined; + } +} + +/** + * Record one related child's latest `session.status` for the reaper. + * + * Running adds the child. Idle removes that child and leaves its siblings + * held. Repeating the same status is safe: the map only cares which children + * are still live. + */ +export function noteOpenCodeChildSessionLiveness( + threadId: string, + sessionId: string, + status: OpenCodeChildSessionLiveness, +): void { + if (status === "idle") { + clearOpenCodeChildSession(threadId, sessionId); + return; + } + + const existing = liveChildrenByThreadId.get(threadId); + if (existing) { + existing.add(sessionId); + return; + } + + liveChildrenByThreadId.set(threadId, new Set([sessionId])); +} + +/** + * Release one child. + * + * The thread stays held while any sibling is still running. Clearing a child + * that was never held is a no-op. + */ +export function clearOpenCodeChildSession(threadId: string, sessionId: string): void { + const live = liveChildrenByThreadId.get(threadId); + if (!live) { + return; + } + live.delete(sessionId); + if (live.size === 0) { + liveChildrenByThreadId.delete(threadId); + } +} + +/** + * Release every child held for a thread. + * + * Session stop, unexpected exit, and rewind use this. Those children belong + * to a provider session that is gone or has been replaced. + */ +export function clearOpenCodeThreadChildSessions(threadId: string): void { + liveChildrenByThreadId.delete(threadId); +} + +/** + * True while any related OpenCode child session is still busy or retrying. + */ +export function hasLiveOpenCodeChildSessions(threadId: string): boolean { + return (liveChildrenByThreadId.get(threadId)?.size ?? 0) > 0; +} + +/** + * Whether an inactivity sweep must leave this provider session running. + * + * OpenCode background subagents keep working after the parent turn settles. + * They do not emit task lifecycle events, so background-task liveness stays + * empty and the idle clock keeps moving from the last user-facing activity. + * Related child `session.status` of `busy` or `retry` is recorded here + * instead. Stopping the session aborts those children, so inactivity has to + * ignore that clock while any of them are still live. Idle and deleted + * children do not hold the session. Other providers are left alone, and this + * hold does not surface task rows. + */ +export function openCodeInactivityHeldByChildSession(provider: string, threadId: string): boolean { + return provider === "opencode" && hasLiveOpenCodeChildSessions(threadId); +} + +/** + * Drop every tracked child. + * + * Tests use this so one case cannot hold a thread created by another. The + * production reaper never needs a global reset; session teardown clears the + * thread it owns. + */ +export function resetOpenCodeChildSessionLiveness(): void { + liveChildrenByThreadId.clear(); +}