diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index ee5767f9d356..1caa643adf0f 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -61,6 +61,7 @@ type MessageEntry = { info: { id: string; role: "user" | "assistant"; + parentID?: string; }; parts: Array; }; @@ -97,6 +98,7 @@ const runtimeMock = { promptEchoEvents: [] as Array, closeError: null as Error | null, messages: [] as MessageEntry[], + messagesImplementation: null as ((signal?: AbortSignal) => Promise) | null, subscribedEvents: [] as Array>, eventSubscribeObserved: null as (() => void) | null, eventStreamError: null as ((cause: unknown) => void) | null, @@ -157,6 +159,7 @@ const runtimeMock = { this.state.promptEchoEvents.length = 0; this.state.closeError = null; this.state.messages = []; + this.state.messagesImplementation = null; this.state.subscribedEvents = []; this.state.eventSubscribeObserved = null; this.state.eventStreamError = null; @@ -355,7 +358,12 @@ const OpenCodeRuntimeTestDouble: OpenCodeRuntimeShape = { runtimeMock.state.summarizeCalls.push(input); return { data: true }; }, - messages: async () => ({ data: runtimeMock.state.messages }), + messages: async (input?: { limit?: number }, options?: { signal?: AbortSignal }) => { + const messages = runtimeMock.state.messagesImplementation + ? await runtimeMock.state.messagesImplementation(options?.signal) + : runtimeMock.state.messages; + return { data: input?.limit ? messages.slice(-input.limit) : messages }; + }, message: async ({ sessionID, messageID }: { sessionID: string; messageID: string }) => { runtimeMock.state.messageCalls.push({ sessionID, messageID }); if (runtimeMock.state.messageFailures > 0) { @@ -1831,10 +1839,15 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { const threadId = asThreadId("thread-steer-reconnect-before-acceptance"); const firstUserMessageEvent = promiseWithResolvers(); const reconnectEvent = promiseWithResolvers(); + const idleEvent = promiseWithResolvers(); const steerStarted = promiseWithResolvers(); const steerRelease = promiseWithResolvers(); runtimeMock.state.autoPromptEcho = false; - runtimeMock.state.subscribedEvents = [firstUserMessageEvent.promise, reconnectEvent.promise]; + runtimeMock.state.subscribedEvents = [ + firstUserMessageEvent.promise, + reconnectEvent.promise, + idleEvent.promise, + ]; runtimeMock.state.promptAsyncImplementation = async () => { if (runtimeMock.state.promptCalls.length === 2) { steerStarted.resolve(undefined); @@ -1900,6 +1913,18 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { yield* Fiber.join(steerFiber); yield* advanceTestClock(250); + NodeAssert.equal(completedFiber.pollUnsafe(), undefined); + const session = (yield* adapter.listSessions()).find((entry) => entry.threadId === threadId); + NodeAssert.equal(session?.status, "running"); + NodeAssert.equal(session?.activeTurnId, activeTurn.turnId); + idleEvent.resolve({ + id: "evt-idle-after-recovered-steer", + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); const completed = Option.getOrUndefined( yield* Fiber.join(completedFiber).pipe(Effect.timeout("1 second")), ); @@ -1914,6 +1939,92 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("keeps a recovered prompt running until OpenCode finishes its delegated work", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-delayed-busy-with-subagent"); + const sessionID = "http://127.0.0.1:9999/session"; + const enqueue = makeOpenCodeEventQueue(); + runtimeMock.state.autoPromptEcho = false; + runtimeMock.state.promptAsyncImplementation = async () => { + const prompt = runtimeMock.state.promptCalls.at(-1) as { messageID: string }; + runtimeMock.state.messages.push({ + info: { id: prompt.messageID, role: "user" }, + parts: [], + }); + }; + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId), + Stream.takeUntil((event) => event.type === "thread.state.changed"), + Stream.runCollect, + Effect.forkChild, + ); + const turn = yield* adapter.sendTurn({ + threadId, + input: "Delegate the investigation", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + // The prompt is persisted before OpenCode starts its session loop. + // Its status map remains empty until the delayed busy event arrives. + yield* advanceTestClock(1_000); + enqueue({ type: "session.status", properties: { sessionID, status: { type: "busy" } } }); + enqueue({ + type: "session.created", + properties: { info: { id: "ses_child", parentID: sessionID, title: "Investigate" } }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_child", status: { type: "busy" } }, + }); + enqueue({ + type: "message.part.updated", + properties: { + sessionID, + part: { + id: "part-task", + sessionID, + messageID: "msg-assistant", + type: "tool", + callID: "call-task", + tool: "task", + state: { status: "pending", input: {}, raw: "" }, + }, + }, + }); + enqueue({ + type: "session.status", + properties: { sessionID: "ses_child", status: { type: "idle" } }, + }); + enqueue({ type: "session.compacted", properties: { sessionID } }); + const events = yield* Fiber.join(eventsFiber); + NodeAssert.equal( + events.some((event) => event.type === "turn.completed"), + false, + ); + NodeAssert.equal(events.find((event) => event.type === "item.started")?.turnId, turn.turnId); + const session = (yield* adapter.listSessions()).find((entry) => entry.threadId === threadId); + NodeAssert.equal(session?.status, "running"); + NodeAssert.equal(session?.activeTurnId, turn.turnId); + + const completedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.runHead, + Effect.forkChild, + ); + enqueue({ type: "session.status", properties: { sessionID, status: { type: "idle" } } }); + NodeAssert.equal(Option.getOrThrow(yield* Fiber.join(completedFiber)).turnId, turn.turnId); + yield* adapter.stopSession(threadId); + }), + ); + it.effect("resolves admission without a prompt echo when busy and idle still arrive", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; @@ -1981,6 +2092,68 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("reconciles an idle event received during the final admission status request", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-idle-during-admission-status"); + const sessionID = "ses_idle_during_admission_status"; + const enqueue = makeOpenCodeEventQueue(); + const statusRequests = Array.from({ length: 5 }, () => promiseWithResolvers()); + const finalStatus = promiseWithResolvers(); + runtimeMock.state.autoPromptEcho = false; + runtimeMock.state.createdSessionIds.push(sessionID); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + let statusRequestIndex = 0; + runtimeMock.state.sessionStatusImplementation = async () => { + const index = statusRequestIndex++; + statusRequests[index]?.resolve(undefined); + return index === 4 ? finalStatus.promise : { data: {} }; + }; + const completedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.runHead, + Effect.forkChild, + ); + const turn = yield* adapter.sendTurn({ + threadId, + input: "Run without a prompt echo", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + for (const [index, delay] of [250, 500, 1_000, 2_000].entries()) { + yield* Effect.promise(() => statusRequests[index]!.promise); + yield* advanceTestClock(delay); + } + yield* Effect.promise(() => statusRequests[4]!.promise); + const processedFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => event.threadId === threadId && event.type === "thread.state.changed", + ), + Stream.runHead, + Effect.forkChild, + ); + enqueue({ type: "session.status", properties: { sessionID, status: { type: "busy" } } }); + enqueue({ type: "session.status", properties: { sessionID, status: { type: "idle" } } }); + enqueue({ type: "session.compacted", properties: { sessionID } }); + yield* Fiber.join(processedFiber); + finalStatus.resolve({ data: {} }); + yield* advanceTestClock(2_000); + + const completed = Option.getOrThrow(yield* Fiber.join(completedFiber)); + NodeAssert.equal(completed.turnId, turn.turnId); + NodeAssert.ok(completed.type === "turn.completed"); + NodeAssert.equal(completed.payload.state, "completed"); + NodeAssert.equal(runtimeMock.state.abortCalls.includes(sessionID), false); + yield* adapter.stopSession(threadId); + }), + ); + it.effect("uses polled busy status to admit output after a stopped turn", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; @@ -7133,6 +7306,205 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect.each([ + { scenario: "missed completion", reply: "current", status: "idle", echo: false }, + { scenario: "prompt echo before acceptance", reply: "current", status: "idle", echo: true }, + { scenario: "ongoing work", reply: "current", status: "busy", echo: false }, + { scenario: "prompt not started", reply: "none", status: "idle", echo: false }, + { scenario: "reply to an earlier prompt", reply: "earlier", status: "idle", echo: false }, + { scenario: "message lookup failures", reply: "current", status: "idle", echo: false }, + { scenario: "steer during lookup", reply: "current", status: "idle", echo: false }, + { scenario: "interrupt during lookup", reply: "current", status: "idle", echo: false }, + { scenario: "stop during lookup", reply: "current", status: "idle", echo: false }, + { scenario: "new turn during lookup", reply: "current", status: "idle", echo: false }, + ] as const)( + "recovers reconnect before prompt acceptance: $scenario", + ({ scenario, reply, status, echo }) => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-reconnect-before-acceptance-completion"); + const sessionID = "http://127.0.0.1:9999/session"; + const enqueue = makeOpenCodeEventQueue(); + const promptStarted = promiseWithResolvers(); + const promptRelease = promiseWithResolvers(); + runtimeMock.state.autoPromptEcho = false; + runtimeMock.state.promptAsyncImplementation = async () => { + promptStarted.resolve(undefined); + await promptRelease.promise; + }; + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turnFiber = yield* adapter + .sendTurn({ + threadId, + input: "Work", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }) + .pipe(Effect.forkChild); + yield* Effect.promise(() => promptStarted.promise); + + const warningFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "runtime.warning"), + Stream.runHead, + Effect.forkChild, + ); + runtimeMock.state.eventStreamError?.(new Error("socket closed")); + yield* Fiber.join(warningFiber); + // The prompt echo and any completion event were lost during the outage. + const prompt = runtimeMock.state.promptCalls[0] as { messageID: string }; + runtimeMock.state.messages.push({ + info: { id: prompt.messageID, role: "user" }, + parts: [], + }); + if (reply !== "none") { + runtimeMock.state.messages.push({ + info: { + id: "msg-assistant-during-outage", + role: "assistant", + parentID: reply === "current" ? prompt.messageID : "msg-earlier-prompt", + }, + parts: [], + }); + } + runtimeMock.state.sessionStatus = status; + let failures = scenario === "message lookup failures" ? 2 : 0; + const snapshotStarted = promiseWithResolvers(); + const snapshotRelease = promiseWithResolvers(); + const holdSnapshot = scenario.endsWith("during lookup"); + let snapshotSignal: AbortSignal | undefined; + runtimeMock.state.messagesImplementation = async (signal) => { + if (failures > 0) { + failures -= 1; + throw new Error("message history temporarily unavailable"); + } + const messages = [...runtimeMock.state.messages]; + snapshotSignal = signal; + snapshotStarted.resolve(undefined); + if (holdSnapshot) await snapshotRelease.promise; + return messages; + }; + const reconnectedFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => event.threadId === threadId && event.type === "thread.state.changed", + ), + Stream.runHead, + Effect.forkChild, + ); + enqueue({ type: "server.connected", properties: {} }); + if (echo) { + enqueue({ + type: "message.updated", + properties: { sessionID, info: { id: prompt.messageID, role: "user" } }, + }); + } + enqueue({ type: "session.compacted", properties: { sessionID } }); + yield* Fiber.join(reconnectedFiber); + + const recoveryWarning = yield* Deferred.make(); + const completedFiber = yield* adapter.streamEvents.pipe( + Stream.tap((event) => + event.threadId === threadId && + event.type === "runtime.warning" && + event.payload.message === "OpenCode turn completion is waiting for message history." + ? Deferred.succeed(recoveryWarning, undefined) + : Effect.void, + ), + Stream.filter( + (event) => + event.threadId === threadId && + event.type === "turn.completed" && + event.payload.state === "completed", + ), + Stream.runHead, + Effect.forkChild, + ); + promptRelease.resolve(undefined); + const turn = yield* Fiber.join(turnFiber); + if (scenario === "message lookup failures") { + yield* Deferred.await(recoveryWarning); + yield* advanceTestClock(250); + } + yield* Effect.promise(() => snapshotStarted.promise); + + if (holdSnapshot) { + let currentTurn = turn; + if (scenario === "stop during lookup") { + yield* adapter.stopSession(threadId); + NodeAssert.equal(snapshotSignal?.aborted, true); + } else { + if (scenario !== "steer during lookup") { + yield* adapter.interruptTurn(threadId, turn.turnId); + } + if (scenario !== "interrupt during lookup") { + runtimeMock.state.promptAsyncImplementation = null; + runtimeMock.state.autoPromptEcho = true; + currentTurn = yield* adapter.sendTurn({ + threadId, + input: "More work", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + } + } + snapshotRelease.resolve(undefined); + yield* TestClock.adjust("0 millis"); + NodeAssert.equal(completedFiber.pollUnsafe(), undefined); + const session = (yield* adapter.listSessions()).find( + (entry) => entry.threadId === threadId, + ); + if (scenario === "stop during lookup") { + NodeAssert.equal(session, undefined); + } else if (scenario === "interrupt during lookup") { + NodeAssert.equal(session?.activeTurnId, undefined); + NodeAssert.equal(session?.status, "ready"); + } else { + NodeAssert.equal(session?.status, "running"); + NodeAssert.equal(session?.activeTurnId, currentTurn.turnId); + enqueue({ + type: "session.status", + properties: { sessionID, status: { type: "idle" } }, + }); + NodeAssert.equal( + Option.getOrThrow(yield* Fiber.join(completedFiber)).turnId, + currentTurn.turnId, + ); + } + return; + } + + if (reply === "current" && status === "idle") { + NodeAssert.equal( + Option.getOrThrow(yield* Fiber.join(completedFiber)).turnId, + turn.turnId, + ); + const session = (yield* adapter.listSessions()).find( + (entry) => entry.threadId === threadId, + ); + NodeAssert.equal(session?.status, "ready"); + NodeAssert.equal(session?.activeTurnId, undefined); + } else { + yield* TestClock.adjust("0 millis"); + const session = (yield* adapter.listSessions()).find( + (entry) => entry.threadId === threadId, + ); + NodeAssert.equal(session?.status, "running"); + NodeAssert.equal(session?.activeTurnId, turn.turnId); + NodeAssert.equal(completedFiber.pollUnsafe(), undefined); + runtimeMock.state.sessionStatus = "idle"; + enqueue({ type: "session.status", properties: { sessionID, status: { type: "idle" } } }); + } + NodeAssert.equal(Option.getOrThrow(yield* Fiber.join(completedFiber)).turnId, turn.turnId); + }), + ); + it.effect("warns on disconnection and recovers a completion missed during reconnect", () => 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 742ee9b86d6a..0e2ac0f3c297 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -25,6 +25,7 @@ import * as Fiber from "effect/Fiber"; import * as Path from "effect/Path"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; +import * as Schedule from "effect/Schedule"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; import * as Semaphore from "effect/Semaphore"; @@ -221,7 +222,6 @@ interface OpenCodePromptAdmission { idleObservedAfterMessage: boolean; messageObserved: boolean; busyObserved: boolean; - idleStatusConfirmations: number; accepted: boolean; cancelled: boolean; readonly acceptance: Deferred.Deferred; @@ -1391,6 +1391,18 @@ export function makeOpenCodeAdapter( } } + if ( + promptAdmission.messageObserved && + promptAdmission.idleDuringAdmission === undefined && + promptAdmission.priorIdle === undefined + ) { + // A persisted prompt proves admission, not completion. OpenCode can + // still report idle before its session loop starts processing it. + context.promptAdmission = undefined; + context.awaitingBusyAfterInterruption = false; + return; + } + const statusResponse = yield* runOpenCodeSdk("session.status", (signal) => context.client.session.status(undefined, { signal }), ).pipe(Effect.timeout("1 second"), Effect.option); @@ -1414,7 +1426,6 @@ export function makeOpenCodeAdapter( const isBusy = status?.type === "busy" || status?.type === "retry"; if (isBusy) { promptAdmission.busyObserved = true; - promptAdmission.idleStatusConfirmations = 0; context.awaitingBusyAfterInterruption = false; context.promptAdmission = undefined; return; @@ -1431,39 +1442,6 @@ export function makeOpenCodeAdapter( yield* scheduleIdleReconciliation(context, promptAdmission.turnId, idle.raw); return; } - if (isIdle && promptAdmission.messageObserved) { - promptAdmission.idleStatusConfirmations += 1; - if (promptAdmission.idleStatusConfirmations >= 2) { - context.promptAdmission = undefined; - context.awaitingBusyAfterInterruption = false; - yield* completeOpenCodeTurn( - context, - promptAdmission.turnId, - promptAdmission.generation, - { - type: "session.status.recovered", - status: statusData, - }, - ); - return; - } - } else if (!isIdle) { - promptAdmission.idleStatusConfirmations = 0; - } - if ( - isIdle && - promptAdmission.messageObserved && - promptAdmission.recoveryRaw !== undefined - ) { - context.promptAdmission = undefined; - context.awaitingBusyAfterInterruption = false; - yield* scheduleIdleReconciliation( - context, - promptAdmission.turnId, - promptAdmission.recoveryRaw, - ); - return; - } const delayMs = Math.min(250 * 2 ** retryCount, 2_000); yield* Effect.sleep(`${delayMs} millis`); @@ -2183,8 +2161,70 @@ export function makeOpenCodeAdapter( if (context.turnTokenUsage) { context.turnTokenUsage.complete = false; } + const admission = context.promptAdmission; yield* schedulePromptAdmissionRecovery(context, event); - if (context.activeTurnId !== undefined && context.promptAdmission === undefined) { + if (admission) { + const recoveryFiber = admission.recoveryFiber; + yield* Effect.gen(function* () { + if (recoveryFiber) { + yield* Fiber.await(recoveryFiber); + } + const isCurrentPrompt = () => + context.activeTurnId === admission.turnId && + context.promptGeneration === admission.generation && + context.promptAdmission === undefined; + if (!isCurrentPrompt()) { + return; + } + let warned = false; + // A user message alone can precede the session loop. An assistant + // reply proves this prompt started, so idle can recover a missed completion. + const response = yield* Effect.gen(function* () { + if (!isCurrentPrompt()) { + return yield* Effect.interrupt; + } + return yield* runOpenCodeSdk("session.messages", (signal) => + context.client.session.messages( + { sessionID: context.openCodeSessionId, limit: 1 }, + { signal }, + ), + ).pipe(Effect.timeout("1 second"), Effect.retry({ times: 1 })); + }).pipe( + Effect.tapError((cause) => + Effect.gen(function* () { + if (warned || !isCurrentPrompt()) return; + warned = true; + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId: admission.turnId, + })), + type: "runtime.warning", + payload: { + message: "OpenCode turn completion is waiting for message history.", + detail: openCodeRuntimeErrorDetail(cause), + }, + }); + }), + ), + Effect.retry({ + while: isCurrentPrompt, + schedule: Schedule.min([ + Schedule.exponential("250 millis"), + Schedule.spaced("5 seconds"), + ]), + }), + ); + const message = response.data?.at(-1)?.info; + if ( + message?.role === "assistant" && + message.parentID === admission.messageId && + isCurrentPrompt() + ) { + yield* scheduleIdleReconciliation(context, admission.turnId, event); + } + }).pipe(Effect.ignore({ log: true }), Effect.forkIn(context.sessionScope)); + } else if (context.activeTurnId !== undefined) { yield* scheduleIdleReconciliation(context, context.activeTurnId, event); } } @@ -3153,7 +3193,6 @@ export function makeOpenCodeAdapter( idleObservedAfterMessage: false, messageObserved: false, busyObserved: false, - idleStatusConfirmations: 0, accepted: false, cancelled: false, acceptance: Deferred.makeUnsafe(),