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(); +}