diff --git a/apps/server/src/provider/Drivers/OpenCodeDriver.ts b/apps/server/src/provider/Drivers/OpenCodeDriver.ts index 0d874b9ceba8..c37e4d559dd2 100644 --- a/apps/server/src/provider/Drivers/OpenCodeDriver.ts +++ b/apps/server/src/provider/Drivers/OpenCodeDriver.ts @@ -37,6 +37,7 @@ import { import { ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; import { makeManagedServerProvider } from "../makeManagedServerProvider.ts"; import { OpenCodeRuntime, loadOpenCodeCommands } from "../opencodeRuntime.ts"; +import type * as OpenCodeChildSessionLiveness from "../Services/OpenCodeChildSessionLiveness.ts"; import * as OpenCodeServerOwner from "../OpenCodeServerOwner.ts"; import { defaultProviderContinuationIdentity, @@ -84,6 +85,7 @@ export type OpenCodeDriverEnv = | Crypto.Crypto | FileSystem.FileSystem | HttpClient.HttpClient + | OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness | OpenCodeRuntime | Path.Path | ProviderEventLoggers diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index baded8094ba6..2db8c4b52395 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,12 +30,14 @@ import { OpenCodeSettings, ProviderDriverKind, ProviderInstanceId, + type ProviderRuntimeEvent, ThreadId, } from "@t3tools/contracts"; import { createModelSelection } from "@t3tools/shared/model"; import { ServerConfig } from "../../config.ts"; import { ServerSettingsService } from "../../serverSettings.ts"; import { buildRuntimeInstructions } from "../RuntimeInstructions.ts"; +import * as OpenCodeChildSessionLiveness from "../Services/OpenCodeChildSessionLiveness.ts"; import { ProviderSessionDirectory } from "../Services/ProviderSessionDirectory.ts"; import type { OpenCodeAdapterShape } from "../Services/OpenCodeAdapter.ts"; import { @@ -613,6 +616,7 @@ const OpenCodeAdapterTestLayer = Layer.effect( }), ), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(NodeServices.layer), ); @@ -666,7 +670,323 @@ const questionRequest = (id: string, sessionID: string): QuestionRequest => ({ ], }); +const OPENCODE_PARENT_SESSION_ID = "http://127.0.0.1:9999/session"; + +/** + * One consumer for a thread's runtime events. + * + * Child `session.status` writes quiet liveness and emits nothing, so the + * following parent title is the barrier that proves those updates landed. + */ +function collectOpenCodeThreadEvents(adapter: OpenCodeAdapterShape, threadId: ThreadId) { + return Effect.gen(function* () { + 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(function* () { + 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, + ); + + const waitForTitle = (title: string) => + Stream.fromQueue(titles).pipe( + Stream.filter((candidate) => candidate === title), + Stream.take(1), + Stream.runDrain, + ); + const waitForTurnCompleted = () => Queue.take(completedTurnIds); + + return { seen, waitForTitle, waitForTurnCompleted }; + }); +} + +const openCodeParentTitle = (sessionId: string, title: string) => ({ + type: "session.updated" as const, + properties: { + info: { + id: sessionId, + title, + }, + }, +}); + it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { + it.effect( + "records related child session.status as quiet reaper liveness without task activity", + () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + 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: "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(yield* liveness.hasLive(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(openCodeParentTitle(OPENCODE_PARENT_SESSION_ID, "Nested child still running")); + yield* waitForTitle("Nested child still running"); + NodeAssert.equal(yield* liveness.hasLive(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(yield* liveness.hasLive(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(yield* liveness.hasLive(threadId), true); + NodeAssert.deepEqual( + seen.filter((event) => event.type.startsWith("task.")), + [], + ); + + yield* adapter.stopSession(threadId); + NodeAssert.equal(yield* liveness.hasLive(threadId), false); + }), + ); + + it.effect("keeps the parent turn running when a related child session goes idle", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + 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(yield* liveness.hasLive(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(yield* liveness.hasLive(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("drops quiet child liveness when the OpenCode session is rewound or stopped", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + 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(yield* liveness.hasLive(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(yield* liveness.hasLive(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(yield* liveness.hasLive(threadId), true); + + yield* adapter.stopSession(threadId); + NodeAssert.equal(yield* liveness.hasLive(threadId), false); + }), + ); + it.effect("reuses a configured OpenCode server URL instead of spawning a local server", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; @@ -1439,6 +1759,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(NodeServices.layer), ); const context = yield* Layer.buildWithScope(adapterLayer, scope); @@ -6493,6 +6814,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(NodeServices.layer), ); @@ -6559,6 +6881,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(NodeServices.layer), ); @@ -6610,6 +6933,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(NodeServices.layer), ); @@ -7852,6 +8176,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(NodeServices.layer), ); @@ -7934,6 +8259,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(NodeServices.layer), ); diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 04535edbd3af..69b301664665 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -59,6 +59,7 @@ import { toOpenCodeQuestionAnswers, type OpenCodeServerConnection, } from "../opencodeRuntime.ts"; +import * as OpenCodeChildSessionLiveness from "../Services/OpenCodeChildSessionLiveness.ts"; import * as Option from "effect/Option"; const PROVIDER = ProviderDriverKind.make("opencode"); @@ -946,6 +947,7 @@ export function makeOpenCodeAdapter( const boundInstanceId = options?.instanceId ?? ProviderInstanceId.make("opencode"); const serverConfig = yield* ServerConfig; const openCodeRuntime = yield* OpenCodeRuntime; + const childLiveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; const crypto = yield* Crypto.Crypto; const fileSystem = yield* FileSystem.FileSystem; const path = yield* Path.Path; @@ -2227,7 +2229,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 + ) { + yield* childLiveness.clearChild(context.session.threadId, deletedSessionId); + } } const payloadSessionId = openCodeEventSessionId(event); @@ -2260,6 +2268,23 @@ export function makeOpenCodeAdapter( payloadSessionId !== undefined && isOpenCodeChildRequestEvent(event) && (context.relatedSessionIds.has(payloadSessionId) || isKnownPendingTerminalEvent); + // Related child session.status is quiet liveness for the reaper. 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 = OpenCodeChildSessionLiveness.openCodeChildSessionLivenessStatus( + event.properties.status.type, + ); + if (liveness !== undefined) { + yield* childLiveness.note(context.session.threadId, payloadSessionId, liveness); + } + return; + } if (!isParentEvent && !isChildRequestEvent) { return; } @@ -2726,6 +2751,13 @@ export function makeOpenCodeAdapter( }); const startEventPump = Effect.fn("startEventPump")(function* (context: OpenCodeSessionContext) { + // Children die with the provider process. Release their hold when this + // session scope closes (stop, unexpected exit, or layer shutdown) so a + // finished session cannot keep the reaper from stopping the thread. + yield* Scope.addFinalizer( + context.sessionScope, + childLiveness.clearThread(context.session.threadId), + ); // 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 @@ -3972,6 +4004,8 @@ export function makeOpenCodeAdapter( }), ).pipe(Effect.mapError(toRequestError)); yield* clearPendingOpenCodeRequests(context, { type: "session.fork" }); + // Children of the session being replaced must not hold the fork. + yield* childLiveness.clearThread(threadId); context.openCodeSessionId = forkedSessionId; context.relatedSessionIds.clear(); context.relatedSessionIds.add(forkedSessionId); diff --git a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts index c33b1dfc690a..eebc0ee9c38d 100644 --- a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts +++ b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts @@ -56,6 +56,7 @@ import { GrokDriver } from "../Drivers/GrokDriver.ts"; import { OpenCodeDriver } from "../Drivers/OpenCodeDriver.ts"; import * as ModelManifest from "../ModelManifest.ts"; import { OpenCodeRuntimeLive } from "../opencodeRuntime.ts"; +import * as OpenCodeChildSessionLiveness from "../Services/OpenCodeChildSessionLiveness.ts"; import * as ResetCreditCoordinator from "./resetCreditCoordinator.ts"; import { NoOpProviderEventLoggers, ProviderEventLoggers } from "./ProviderEventLoggers.ts"; import { makeProviderInstanceRegistry } from "./ProviderInstanceRegistryLive.ts"; @@ -608,6 +609,7 @@ describe("ProviderInstanceRegistryLive — all drivers slice", () => { Layer.provideMerge(Layer.succeed(ProviderEventLoggers, NoOpProviderEventLoggers)), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(ResetCreditCoordinator.layerTest), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), ); it.live("boots one instance of every shipped driver from a single config map", () => diff --git a/apps/server/src/provider/Layers/ProviderRegistry.test.ts b/apps/server/src/provider/Layers/ProviderRegistry.test.ts index ee1a23573277..51644f8b6aed 100644 --- a/apps/server/src/provider/Layers/ProviderRegistry.test.ts +++ b/apps/server/src/provider/Layers/ProviderRegistry.test.ts @@ -40,6 +40,7 @@ import { AntigravityInstallation } from "../AntigravityInstallation.ts"; import * as ModelManifest from "../ModelManifest.ts"; import { applyProviderCompatibility } from "../providerCompatibility.ts"; import * as ResetCreditCoordinator from "./resetCreditCoordinator.ts"; +import * as OpenCodeChildSessionLiveness from "../Services/OpenCodeChildSessionLiveness.ts"; import * as OpenCodeRuntime from "../opencodeRuntime.ts"; import * as ProviderEventLoggers from "./ProviderEventLoggers.ts"; import { ProviderInstanceRegistryHydrationLive } from "./ProviderInstanceRegistryHydration.ts"; @@ -2323,6 +2324,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(ResetCreditCoordinator.layerTest), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), // NO spawner mock — `ChildProcessSpawner` is supplied by the // outer `NodeServices.layer` on `it.layer(...)` and will @@ -2422,6 +2424,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(ResetCreditCoordinator.layerTest), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.updateService(ChildProcessSpawner.ChildProcessSpawner, (spawner) => ChildProcessSpawner.make((command) => { if (command._tag !== "StandardCommand") return spawner.spawn(command); @@ -2538,6 +2541,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(ResetCreditCoordinator.layerTest), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(NodeServices.layer), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), ); @@ -2600,6 +2604,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(ResetCreditCoordinator.layerTest), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), Layer.provideMerge( mockCommandSpawnerLayer((command, args) => { diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts index 7b1fec90f867..b10c00628520 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts @@ -21,6 +21,7 @@ import { ProjectionSnapshotQuery } from "../../orchestration/Services/Projection import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import * as ProviderSessionRuntime from "../../persistence/ProviderSessionRuntime.ts"; import { ProviderValidationError } from "../Errors.ts"; +import * as OpenCodeChildSessionLiveness from "../Services/OpenCodeChildSessionLiveness.ts"; import { ProviderSessionReaper } from "../Services/ProviderSessionReaper.ts"; import { ProviderService, type ProviderServiceShape } from "../Services/ProviderService.ts"; import { ProviderSessionDirectoryLive } from "./ProviderSessionDirectory.ts"; @@ -62,7 +63,7 @@ function makeReadModel( 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; @@ -122,7 +123,9 @@ function makeReadModel( describe("ProviderSessionReaper", () => { let runtime: ManagedRuntime.ManagedRuntime< - ProviderSessionReaper | ProviderSessionRuntime.ProviderSessionRuntimeRepository, + | ProviderSessionReaper + | ProviderSessionRuntime.ProviderSessionRuntimeRepository + | OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness, unknown > | null = null; let scope: Scope.Closeable | null = null; @@ -150,8 +153,8 @@ describe("ProviderSessionReaper", () => { ); } - async function sweepAt(nowMs: number) { - await runtime!.runPromise( + async function sweepAt(nowMs: number, target: NonNullable = runtime!) { + await target.runPromise( Effect.gen(function* () { const reaper = yield* ProviderSessionReaper; const clock = yield* Clock.Clock; @@ -265,6 +268,7 @@ describe("ProviderSessionReaper", () => { searchThreads: () => Effect.succeed({ matches: [] }), }), ), + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(NodeServices.layer), ); @@ -734,4 +738,167 @@ describe("ProviderSessionReaper", () => { reapedThreadId, ]); }); + + it("does not reap a stale OpenCode session while a child session is still live", async () => { + 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"; + const shell = (threadId: ThreadId, providerName: "opencode" | "claudeAgent") => ({ + 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), + ); + const liveness = await runtime!.runPromise( + Effect.service(OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness), + ); + const seed = (threadId: ThreadId, providerName: "opencode" | "claudeAgent") => + 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"); + + await runtime!.runPromise( + Effect.gen(function* () { + yield* liveness.note(liveThreadId, "ses_a", "running"); + yield* liveness.note(liveThreadId, "ses_b", "running"); + yield* liveness.note(claudeThreadId, "ses_claude", "running"); + }), + ); + + await sweepAt(Date.parse(updatedAt) + 1_000); + expect(harness.stopSession.mock.calls.map(([request]) => request.threadId)).toEqual([ + claudeThreadId, + settledOpenCodeThreadId, + ]); + + const markStopped = (threadId: ThreadId, providerName: "opencode" | "claudeAgent") => + 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"); + + await runtime!.runPromise(liveness.note(liveThreadId, "ses_a", "idle")); + await sweepAt(Date.parse(updatedAt) + 1_000); + expect(harness.stopSession.mock.calls.map(([request]) => request.threadId)).toEqual([ + claudeThreadId, + settledOpenCodeThreadId, + ]); + + await runtime!.runPromise(liveness.clearThread(liveThreadId)); + await sweepAt(Date.parse(updatedAt) + 1_000); + expect(harness.stopSession.mock.calls.map(([request]) => request.threadId)).toEqual([ + claudeThreadId, + settledOpenCodeThreadId, + liveThreadId, + ]); + }); + + it("does not let one runtime's OpenCode child liveness hold another runtime", async () => { + const threadId = ThreadId.make("thread-reaper-opencode-runtime-isolation"); + const updatedAt = "2026-04-14T00:00:00.000Z"; + const shell = { + id: threadId, + session: { + threadId, + status: "ready" as const, + providerName: "opencode" as const, + runtimeMode: "full-access" as const, + activeTurnId: null, + lastError: null, + updatedAt, + }, + }; + const harnessA = await createHarness({ readModel: makeReadModel([shell]) }); + const runtimeA = runtime!; + runtime = null; + const harnessB = await createHarness({ readModel: makeReadModel([shell]) }); + const runtimeB = runtime!; + try { + const seed = ( + target: NonNullable, + repository: ProviderSessionRuntime.ProviderSessionRuntimeRepository["Service"], + ) => + target.runPromise( + repository.upsert({ + threadId, + providerName: "opencode", + providerInstanceId: null, + adapterKey: "opencode", + runtimeMode: "full-access", + status: "running", + lastSeenAt: updatedAt, + resumeCursor: { opaque: `resume-${threadId}` }, + runtimePayload: null, + }), + ); + const repositoryA = await runtimeA.runPromise( + Effect.service(ProviderSessionRuntime.ProviderSessionRuntimeRepository), + ); + const repositoryB = await runtimeB.runPromise( + Effect.service(ProviderSessionRuntime.ProviderSessionRuntimeRepository), + ); + await seed(runtimeA, repositoryA); + await seed(runtimeB, repositoryB); + const livenessA = await runtimeA.runPromise( + Effect.service(OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness), + ); + await runtimeA.runPromise(livenessA.note(threadId, "ses_background", "running")); + + await sweepAt(Date.parse(updatedAt) + 1_000, runtimeA); + await sweepAt(Date.parse(updatedAt) + 1_000, runtimeB); + expect(harnessA.stopSession).not.toHaveBeenCalled(); + expect(harnessB.stopSession.mock.calls.map(([request]) => request.threadId)).toEqual([ + threadId, + ]); + + await runtimeA.runPromise(livenessA.clearThread(threadId)); + await sweepAt(Date.parse(updatedAt) + 1_000, runtimeA); + expect(harnessA.stopSession.mock.calls.map(([request]) => request.threadId)).toEqual([ + threadId, + ]); + } finally { + await runtimeA.dispose(); + } + }); }); diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.ts index bf8199f80eac..02b40743fe68 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 * as OpenCodeChildSessionLiveness from "../Services/OpenCodeChildSessionLiveness.ts"; import { ProviderService } from "../Services/ProviderService.ts"; const DEFAULT_INACTIVITY_THRESHOLD_MS = 30 * 60 * 1000; @@ -27,6 +28,7 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = const providerService = yield* ProviderService; const directory = yield* ProviderSessionDirectory; const projectionSnapshotQuery = yield* ProjectionSnapshotQuery; + const childLiveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; const inactivityThresholdMs = Math.max( 1, @@ -92,6 +94,22 @@ 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. Leave that clock alone and skip + // the stop while a related child is still busy or retrying. Idle, + // deletion, and session teardown release the hold, so the next sweep + // can reap. Other providers are not held. + if (yield* childLiveness.holdsInactivity(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/Services/OpenCodeChildSessionLiveness.test.ts b/apps/server/src/provider/Services/OpenCodeChildSessionLiveness.test.ts new file mode 100644 index 000000000000..da2aa1ad7dcc --- /dev/null +++ b/apps/server/src/provider/Services/OpenCodeChildSessionLiveness.test.ts @@ -0,0 +1,123 @@ +import * as NodeAssert from "node:assert/strict"; +import { ThreadId } from "@t3tools/contracts"; +import { describe, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Layer from "effect/Layer"; +import * as Scope from "effect/Scope"; + +import * as OpenCodeChildSessionLiveness from "./OpenCodeChildSessionLiveness.ts"; + +const threadA = ThreadId.make("thread-opencode-child-a"); +const threadB = ThreadId.make("thread-opencode-child-b"); + +const livenessLayer = OpenCodeChildSessionLiveness.layer; + +describe("OpenCodeChildSessionLiveness", () => { + it("maps busy and retry onto running, and idle onto a release", () => { + NodeAssert.equal( + OpenCodeChildSessionLiveness.openCodeChildSessionLivenessStatus("busy"), + "running", + ); + NodeAssert.equal( + OpenCodeChildSessionLiveness.openCodeChildSessionLivenessStatus("retry"), + "running", + ); + NodeAssert.equal( + OpenCodeChildSessionLiveness.openCodeChildSessionLivenessStatus("idle"), + "idle", + ); + NodeAssert.equal( + OpenCodeChildSessionLiveness.openCodeChildSessionLivenessStatus("paused"), + undefined, + ); + }); + + it.effect( + "holds a thread while any related child is running and releases on idle or delete", + () => + Effect.gen(function* () { + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + yield* liveness.note(threadA, "ses_a", "running"); + yield* liveness.note(threadA, "ses_a", "running"); + yield* liveness.note(threadA, "ses_b", "running"); + yield* liveness.note(threadB, "ses_other", "running"); + + NodeAssert.equal(yield* liveness.hasLive(threadA), true); + NodeAssert.equal(yield* liveness.holdsInactivity("opencode", threadA), true); + NodeAssert.equal(yield* liveness.holdsInactivity("claudeAgent", threadA), false); + NodeAssert.equal(yield* liveness.holdsInactivity("codex", threadA), false); + NodeAssert.equal( + yield* liveness.holdsInactivity("opencode", ThreadId.make("thread-missing")), + false, + ); + + yield* liveness.note(threadA, "ses_a", "idle"); + yield* liveness.clearChild(threadA, "ses_missing"); + NodeAssert.equal(yield* liveness.hasLive(threadA), true); + NodeAssert.equal(yield* liveness.hasLive(threadB), true); + + yield* liveness.clearChild(threadA, "ses_b"); + NodeAssert.equal(yield* liveness.hasLive(threadA), false); + NodeAssert.equal(yield* liveness.holdsInactivity("opencode", threadA), false); + + yield* liveness.clearThread(threadB); + NodeAssert.equal(yield* liveness.hasLive(threadB), false); + }).pipe(Effect.provide(livenessLayer)), + ); + + it.effect("ignores an idle release for a child that was never held", () => + Effect.gen(function* () { + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + yield* liveness.note(threadA, "ses_missing", "idle"); + yield* liveness.clearChild(threadA, "ses_missing"); + NodeAssert.equal(yield* liveness.hasLive(threadA), false); + }).pipe(Effect.provide(livenessLayer)), + ); + + it.effect("shares liveness inside one runtime and not across runtimes", () => + Effect.gen(function* () { + const leftScope = yield* Scope.make("sequential"); + const rightScope = yield* Scope.make("sequential"); + yield* Effect.addFinalizer(() => Scope.close(leftScope, Exit.void)); + yield* Effect.addFinalizer(() => Scope.close(rightScope, Exit.void)); + const left = yield* Layer.build(livenessLayer).pipe(Scope.provide(leftScope)); + const right = yield* Layer.build(livenessLayer).pipe(Scope.provide(rightScope)); + + const hasLive = (threadId: ThreadId) => + Effect.gen(function* () { + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + return yield* liveness.hasLive(threadId); + }); + + yield* Effect.gen(function* () { + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + yield* liveness.note(threadA, "ses_a", "running"); + const again = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + NodeAssert.equal(yield* again.hasLive(threadA), true); + }).pipe(Effect.provide(left)); + + NodeAssert.equal(yield* hasLive(threadA).pipe(Effect.provide(left)), true); + NodeAssert.equal(yield* hasLive(threadA).pipe(Effect.provide(right)), false); + NodeAssert.equal( + yield* Effect.gen(function* () { + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + return yield* liveness.holdsInactivity("opencode", threadA); + }).pipe(Effect.provide(right)), + false, + ); + + yield* Effect.gen(function* () { + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + yield* liveness.note(threadA, "ses_right", "running"); + }).pipe(Effect.provide(right)); + yield* Effect.gen(function* () { + const liveness = yield* OpenCodeChildSessionLiveness.OpenCodeChildSessionLiveness; + yield* liveness.clearThread(threadA); + }).pipe(Effect.provide(left)); + + NodeAssert.equal(yield* hasLive(threadA).pipe(Effect.provide(left)), false); + NodeAssert.equal(yield* hasLive(threadA).pipe(Effect.provide(right)), true); + }), + ); +}); diff --git a/apps/server/src/provider/Services/OpenCodeChildSessionLiveness.ts b/apps/server/src/provider/Services/OpenCodeChildSessionLiveness.ts new file mode 100644 index 000000000000..aa67da08a703 --- /dev/null +++ b/apps/server/src/provider/Services/OpenCodeChildSessionLiveness.ts @@ -0,0 +1,141 @@ +/** + * Quiet liveness for OpenCode child sessions, owned by one runtime. + * + * Background subagents keep the OpenCode process busy after the parent turn + * settles. They do not emit `task.*` events, so the provider-session reaper + * would treat that silence as inactivity and stop the session. The OpenCode + * adapter records related child `session.status` here, and the reaper reads + * it before stopping. + * + * `busy` and `retry` hold a child. `idle`, deletion, and session teardown + * release it. The hold does not move the idle clock. Other providers are + * ignored. + * + * The map is created inside `make`. `layer` runs in the runtime scope, so each + * runtime build gets its own instance and closing that scope drops the map. + * The adapter and the reaper share a map only when they are built in the same + * runtime. + */ +import { ThreadId } from "@t3tools/contracts"; +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Ref from "effect/Ref"; + +export type OpenCodeChildSessionLivenessStatus = "running" | "idle"; + +/** + * Map an OpenCode child `session.status` type onto reaper liveness. + * + * `busy` and `retry` are in-flight provider work. `idle` releases that child. + * Any other type is ignored so an unknown status cannot drop a child that is + * still running. + */ +export const openCodeChildSessionLivenessStatus = ( + statusType: string, +): OpenCodeChildSessionLivenessStatus | undefined => { + switch (statusType) { + case "busy": + case "retry": + return "running"; + case "idle": + return "idle"; + default: + return undefined; + } +}; + +type ChildSessionsByThread = Map>; + +const withoutChild = ( + current: ChildSessionsByThread, + threadId: ThreadId, + sessionId: string, +): ChildSessionsByThread => { + const existing = current.get(threadId); + if (existing === undefined || !existing.has(sessionId)) { + return current; + } + const remaining = new Set(existing); + remaining.delete(sessionId); + const next = new Map(current); + if (remaining.size === 0) { + next.delete(threadId); + } else { + next.set(threadId, remaining); + } + return next; +}; + +export class OpenCodeChildSessionLiveness extends Context.Service< + OpenCodeChildSessionLiveness, + { + /** Record one related child's latest status. Idle removes only that child. */ + readonly note: ( + threadId: ThreadId, + sessionId: string, + status: OpenCodeChildSessionLivenessStatus, + ) => Effect.Effect; + /** Release one child. Siblings keep the thread held. */ + readonly clearChild: (threadId: ThreadId, sessionId: string) => Effect.Effect; + /** + * 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. + */ + readonly clearThread: (threadId: ThreadId) => Effect.Effect; + readonly hasLive: (threadId: ThreadId) => Effect.Effect; + /** + * Whether an inactivity sweep must leave this provider session running. + * + * Only OpenCode is held, and only while a related child is still live. + * The call does not change the binding's last-seen timestamp. + */ + readonly holdsInactivity: (provider: string, threadId: ThreadId) => Effect.Effect; + } +>()("t3/provider/Services/OpenCodeChildSessionLiveness") {} + +export const make = Effect.gen(function* () { + const liveChildrenByThreadId = yield* Ref.make(new Map()); + yield* Effect.addFinalizer(() => Ref.set(liveChildrenByThreadId, new Map())); + + const clearChild = (threadId: ThreadId, sessionId: string) => + Ref.update(liveChildrenByThreadId, (current) => withoutChild(current, threadId, sessionId)); + const hasLive = (threadId: ThreadId) => + Ref.get(liveChildrenByThreadId).pipe( + Effect.map((current) => (current.get(threadId)?.size ?? 0) > 0), + ); + + return OpenCodeChildSessionLiveness.of({ + note: (threadId, sessionId, status) => + status === "idle" + ? clearChild(threadId, sessionId) + : Ref.update(liveChildrenByThreadId, (current) => { + const existing = current.get(threadId); + if (existing?.has(sessionId)) { + return current; + } + const next = new Map(current); + const children = new Set(existing ?? []); + children.add(sessionId); + next.set(threadId, children); + return next; + }), + clearChild, + clearThread: (threadId) => + Ref.update(liveChildrenByThreadId, (current) => { + if (!current.has(threadId)) { + return current; + } + const next = new Map(current); + next.delete(threadId); + return next; + }), + hasLive, + holdsInactivity: (provider, threadId) => + provider === "opencode" ? hasLive(threadId) : Effect.succeed(false), + }); +}); + +export const layer = Layer.effect(OpenCodeChildSessionLiveness, make); diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 88c08fd3bdee..1d7edaad10e0 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -61,6 +61,7 @@ import { ProviderRegistry } from "./provider/Services/ProviderRegistry.ts"; import { ProviderSessionReaperLive } from "./provider/Layers/ProviderSessionReaper.ts"; import { ProviderUsageLimitsIngestionLive } from "./provider/Layers/ProviderUsageLimitsIngestion.ts"; import * as OpenCodeRuntime from "./provider/opencodeRuntime.ts"; +import * as OpenCodeChildSessionLiveness from "./provider/Services/OpenCodeChildSessionLiveness.ts"; import * as CheckpointDiffQuery from "./checkpointing/CheckpointDiffQuery.ts"; import * as CheckpointStore from "./checkpointing/CheckpointStore.ts"; import * as AzureDevOpsCli from "./sourceControl/AzureDevOpsCli.ts"; @@ -550,6 +551,9 @@ const RuntimeCoreDependenciesLive = ReactorLayerLive.pipe( // no longer transitively provides it. Exposing it at the runtime level // keeps a single Live for all opencode consumers. Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), + // One child-liveness map per runtime, shared by the OpenCode adapter and + // the provider-session reaper. A second runtime builds its own layer. + Layer.provideMerge(OpenCodeChildSessionLiveness.layer), Layer.provideMerge(WorkspaceLayerLive), Layer.provideMerge(Layer.mergeAll(NativeAppIconResolver.layer, ProjectFaviconResolverLayerLive)), Layer.provideMerge(RepositoryIdentityResolverLayerLive),