From 7db7fd048be749ea1aab87ba4409e812f22b08e7 Mon Sep 17 00:00:00 2001 From: Simon van Laak <32648751+simonvanlaak@users.noreply.github.com> Date: Sun, 13 Sep 2026 20:17:03 +0200 Subject: [PATCH 1/2] fix(server): reconcile delayed OpenCode prompt admission --- .../provider/Layers/OpenCodeAdapter.test.ts | 314 ++++++++++++++++++ .../src/provider/Layers/OpenCodeAdapter.ts | 42 ++- 2 files changed, 355 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index a6636fcd69d0..0b3bb18ab558 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -29,6 +29,7 @@ import { OpenCodeSettings, ProviderDriverKind, ProviderInstanceId, + type ProviderRuntimeEvent, ThreadId, } from "@t3tools/contracts"; import { createModelSelection } from "@t3tools/shared/model"; @@ -2021,6 +2022,319 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("admits a busy event as the only evidence just after the probe deadline", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-busy-event-after-probe-deadline"); + const sessionId = "ses_busy_event_after_probe_deadline"; + const pushEvent = makeOpenCodeEventQueue(); + const fifthProbeObserved = promiseWithResolvers(); + const terminalObserved = promiseWithResolvers(); + const terminalEvents: Array = []; + runtimeMock.state.autoPromptEcho = false; + runtimeMock.state.createdSessionIds.push(sessionId); + runtimeMock.state.sessionStatusImplementation = async () => { + if (runtimeMock.state.sessionStatusCalls === 5) { + fifthProbeObserved.resolve(undefined); + } + return { data: {} }; + }; + + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => + event.threadId === threadId && + (event.type === "turn.completed" || event.type === "runtime.error"), + ), + Stream.runForEach((event) => + Effect.sync(() => { + terminalEvents.push(event); + if (event.type === "turn.completed") { + terminalObserved.resolve(undefined); + } + }), + ), + Effect.forkChild, + ); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turn = yield* adapter.sendTurn({ + threadId, + input: "Keep this admitted turn alive", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + + for (const delayMs of [250, 500, 1_000, 2_000]) { + yield* advanceTestClock(delayMs); + } + yield* Effect.promise(() => fifthProbeObserved.promise); + yield* advanceTestClock(2_000); + const abortCallsBeforeAdmissionEvidence = runtimeMock.state.abortCalls.filter( + (candidate) => candidate === sessionId, + ).length; + + pushEvent({ + id: "evt-busy-after-probe-deadline", + type: "session.status", + properties: { + sessionID: sessionId, + status: { type: "busy" }, + }, + }); + pushEvent({ + id: "evt-idle-after-probe-deadline", + type: "session.status", + properties: { + sessionID: sessionId, + status: { type: "idle" }, + }, + }); + yield* Effect.promise(() => terminalObserved.promise); + + const turnCompleted = terminalEvents.find((event) => event.type === "turn.completed"); + NodeAssert.equal(turnCompleted?.payload.state, "completed"); + NodeAssert.equal( + terminalEvents.some((event) => event.type === "runtime.error"), + false, + ); + NodeAssert.equal( + runtimeMock.state.abortCalls.filter((candidate) => candidate === sessionId).length, + abortCallsBeforeAdmissionEvidence, + ); + NodeAssert.equal(turnCompleted?.turnId, turn.turnId); + + yield* Fiber.interrupt(eventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("fails admission after bounded probes and grace produce no evidence", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-admission-without-evidence"); + const sessionId = "ses_admission_without_evidence"; + const fifthProbeObserved = promiseWithResolvers(); + runtimeMock.state.autoPromptEcho = false; + runtimeMock.state.createdSessionIds.push(sessionId); + runtimeMock.state.sessionStatusImplementation = async () => { + if (runtimeMock.state.sessionStatusCalls === 5) { + fifthProbeObserved.resolve(undefined); + } + return { data: {} }; + }; + + const terminalFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => + event.threadId === threadId && + (event.type === "turn.completed" || event.type === "runtime.error"), + ), + Stream.take(2), + Stream.runCollect, + Effect.forkChild, + ); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId, + input: "Fail when admission has no evidence", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + + for (const delayMs of [250, 500, 1_000, 2_000]) { + yield* advanceTestClock(delayMs); + } + yield* Effect.promise(() => fifthProbeObserved.promise); + yield* advanceTestClock(2_000); + const abortCallsBeforeTerminalFailure = runtimeMock.state.abortCalls.filter( + (candidate) => candidate === sessionId, + ).length; + yield* advanceTestClock(1_000); + + const terminalEvents = Array.from(yield* Fiber.join(terminalFiber)); + NodeAssert.equal(runtimeMock.state.sessionStatusCalls, 5); + NodeAssert.equal( + runtimeMock.state.abortCalls.filter((candidate) => candidate === sessionId).length, + abortCallsBeforeTerminalFailure + 1, + ); + NodeAssert.deepEqual( + terminalEvents.map((event) => event.type), + ["turn.completed", "runtime.error"], + ); + const completed = terminalEvents[0]; + NodeAssert.equal(completed?.type, "turn.completed"); + if (completed?.type === "turn.completed") { + NodeAssert.equal(completed.payload.state, "failed"); + } + + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("does not fail admission when the turn is interrupted during grace", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-interrupt-during-admission-grace"); + const sessionId = "ses_interrupt_during_admission_grace"; + const fifthProbeObserved = promiseWithResolvers(); + const terminalEvents: Array = []; + runtimeMock.state.autoPromptEcho = false; + runtimeMock.state.createdSessionIds.push(sessionId); + runtimeMock.state.sessionStatusImplementation = async () => { + if (runtimeMock.state.sessionStatusCalls === 5) { + fifthProbeObserved.resolve(undefined); + } + return { data: {} }; + }; + + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => + event.threadId === threadId && + (event.type === "turn.aborted" || + event.type === "turn.completed" || + event.type === "runtime.error"), + ), + Stream.runForEach((event) => Effect.sync(() => terminalEvents.push(event))), + Effect.forkChild, + ); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turn = yield* adapter.sendTurn({ + threadId, + input: "Interrupt while admission waits for late evidence", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + + for (const delayMs of [250, 500, 1_000, 2_000]) { + yield* advanceTestClock(delayMs); + } + yield* Effect.promise(() => fifthProbeObserved.promise); + yield* advanceTestClock(2_000); + const abortCallsBeforeInterrupt = runtimeMock.state.abortCalls.filter( + (candidate) => candidate === sessionId, + ).length; + yield* adapter.interruptTurn(threadId, turn.turnId); + yield* advanceTestClock(1_000); + + NodeAssert.equal( + runtimeMock.state.abortCalls.filter((candidate) => candidate === sessionId).length, + abortCallsBeforeInterrupt + 1, + ); + NodeAssert.deepEqual( + terminalEvents.map((event) => event.type), + ["turn.aborted"], + ); + + yield* Fiber.interrupt(eventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("does not fail admission after cleanup abort loses ownership to interruption", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-interrupt-during-admission-cleanup"); + const sessionId = "ses_interrupt_during_admission_cleanup"; + const fifthProbeObserved = promiseWithResolvers(); + const admissionAbortStarted = promiseWithResolvers(); + const interruptAbortStarted = promiseWithResolvers(); + const releaseAdmissionAbort = promiseWithResolvers(); + const terminalEvents: Array = []; + let ownedAbortCalls = 0; + runtimeMock.state.autoPromptEcho = false; + runtimeMock.state.createdSessionIds.push(sessionId); + runtimeMock.state.sessionStatusImplementation = async () => { + if (runtimeMock.state.sessionStatusCalls === 5) { + fifthProbeObserved.resolve(undefined); + } + return { data: {} }; + }; + runtimeMock.state.abortImplementation = async (candidate) => { + if (candidate !== sessionId) return; + ownedAbortCalls += 1; + if (ownedAbortCalls === 1) { + admissionAbortStarted.resolve(undefined); + await releaseAdmissionAbort.promise; + } else if (ownedAbortCalls === 2) { + interruptAbortStarted.resolve(undefined); + } + }; + + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => + event.threadId === threadId && + (event.type === "turn.aborted" || + event.type === "turn.completed" || + event.type === "runtime.error"), + ), + Stream.runForEach((event) => Effect.sync(() => terminalEvents.push(event))), + Effect.forkChild, + ); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turn = yield* adapter.sendTurn({ + threadId, + input: "Interrupt while admission cleanup is aborting", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + + for (const delayMs of [250, 500, 1_000, 2_000]) { + yield* advanceTestClock(delayMs); + } + yield* Effect.promise(() => fifthProbeObserved.promise); + yield* advanceTestClock(3_000); + yield* Effect.promise(() => admissionAbortStarted.promise); + + const interruptFiber = yield* adapter + .interruptTurn(threadId, turn.turnId) + .pipe(Effect.forkChild); + yield* Effect.promise(() => interruptAbortStarted.promise); + releaseAdmissionAbort.resolve(undefined); + yield* Fiber.join(interruptFiber); + yield* Effect.yieldNow; + + NodeAssert.equal(ownedAbortCalls, 2); + NodeAssert.deepEqual( + terminalEvents.map((event) => event.type), + ["turn.aborted"], + ); + const sessions = yield* adapter.listSessions(); + const session = sessions.find((candidate) => candidate.threadId === threadId); + NodeAssert.equal(session?.status, "ready"); + NodeAssert.equal(session?.activeTurnId, undefined); + + runtimeMock.state.abortImplementation = null; + yield* Fiber.interrupt(eventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + it.effect("uses polled busy status to admit output after a stopped turn", () => 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 3d216bb1167b..cc4c70888cdb 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -225,6 +225,7 @@ interface OpenCodePromptAdmission { accepted: boolean; cancelled: boolean; readonly acceptance: Deferred.Deferred; + readonly admissionObserved: Deferred.Deferred; readonly submissionSettled: Deferred.Deferred; promptFiber?: Fiber.Fiber; recoveryFiber?: Fiber.Fiber; @@ -1275,9 +1276,12 @@ export function makeOpenCodeAdapter( promptAdmission: OpenCodePromptAdmission, ) { if ( + (yield* Ref.get(context.stopped)) || + sessions.get(context.session.threadId) !== context || context.promptAdmission !== promptAdmission || context.activeTurnId !== promptAdmission.turnId || - context.promptGeneration !== promptAdmission.generation + context.promptGeneration !== promptAdmission.generation || + promptAdmission.cancelled ) { return; } @@ -1288,6 +1292,16 @@ export function makeOpenCodeAdapter( context.client.session.abort({ sessionID: context.openCodeSessionId }, { signal }), ).pipe(Effect.timeout("1 second")), ); + if ( + (yield* Ref.get(context.stopped)) || + sessions.get(context.session.threadId) !== context || + context.promptAdmission !== promptAdmission || + context.activeTurnId !== promptAdmission.turnId || + context.promptGeneration !== promptAdmission.generation || + promptAdmission.cancelled + ) { + return; + } if (Exit.isFailure(abortExit)) { yield* emitUnexpectedExit( context, @@ -1468,6 +1482,25 @@ export function makeOpenCodeAdapter( const delayMs = Math.min(250 * 2 ** retryCount, 2_000); yield* Effect.sleep(`${delayMs} millis`); } + const eventObserved = yield* Deferred.await(promptAdmission.admissionObserved).pipe( + Effect.timeout("1 second"), + Effect.option, + ); + if ( + Option.isSome(eventObserved) && + context.promptAdmission === promptAdmission && + context.activeTurnId === promptAdmission.turnId && + context.promptGeneration === promptAdmission.generation && + !promptAdmission.cancelled + ) { + context.promptAdmission = undefined; + context.awaitingBusyAfterInterruption = false; + const idle = promptAdmission.idleDuringAdmission ?? promptAdmission.priorIdle; + if (idle) { + yield* scheduleIdleReconciliation(context, idle.turnId, idle.raw); + } + return; + } yield* failPromptAdmissionRecovery(context, promptAdmission); }).pipe( Effect.catchCause(() => Effect.void), @@ -2312,6 +2345,9 @@ export function makeOpenCodeAdapter( promptAdmission?.messageId === event.properties.info.id ) { promptAdmission.messageObserved = true; + yield* Deferred.succeed(promptAdmission.admissionObserved, undefined).pipe( + Effect.ignore, + ); if (promptAdmission.accepted) { const idle = promptAdmission.idleDuringAdmission; context.awaitingBusyAfterInterruption = false; @@ -2572,6 +2608,9 @@ export function makeOpenCodeAdapter( context.awaitingBusyAfterInterruption = false; if (context.promptAdmission?.turnId === turnId) { context.promptAdmission.busyObserved = true; + yield* Deferred.succeed(context.promptAdmission.admissionObserved, undefined).pipe( + Effect.ignore, + ); yield* schedulePromptAdmissionRecovery(context, event); } yield* updateProviderSession(context, { @@ -3160,6 +3199,7 @@ export function makeOpenCodeAdapter( accepted: false, cancelled: false, acceptance: Deferred.makeUnsafe(), + admissionObserved: Deferred.makeUnsafe(), submissionSettled: Deferred.makeUnsafe(), recoveryRaw: undefined, }; From c95d93dbcd007807f97892215d889610f38ffacd Mon Sep 17 00:00:00 2001 From: Simon van Laak <32648751+simonvanlaak@users.noreply.github.com> Date: Mon, 14 Sep 2026 09:31:24 +0200 Subject: [PATCH 2/2] fix(server): retain OpenCode admission probe evidence --- .../provider/Layers/OpenCodeAdapter.test.ts | 100 ++++++++++++++++++ .../src/provider/Layers/OpenCodeAdapter.ts | 4 +- 2 files changed, 103 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index 0b3bb18ab558..b878eced4d6f 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -2114,6 +2114,106 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("keeps the turn admitted when the fifth probe finds the exact user message", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-message-on-fifth-admission-probe"); + const sessionId = "ses_message_on_fifth_admission_probe"; + const pushEvent = makeOpenCodeEventQueue(); + const fifthProbeObserved = promiseWithResolvers(); + const terminalObserved = promiseWithResolvers(); + const terminalEvents: Array = []; + runtimeMock.state.autoPromptEcho = false; + runtimeMock.state.createdSessionIds.push(sessionId); + runtimeMock.state.messageFailures = 4; + runtimeMock.state.promptAsyncImplementation = async () => { + const prompt = runtimeMock.state.promptCalls.at(-1) as { messageID?: string } | undefined; + if (prompt?.messageID) { + runtimeMock.state.messages.push({ + info: { id: prompt.messageID, role: "user" }, + parts: [], + }); + } + }; + runtimeMock.state.sessionStatusImplementation = async () => { + if (runtimeMock.state.sessionStatusCalls === 5) { + fifthProbeObserved.resolve(undefined); + } + throw new Error("status unavailable"); + }; + + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => + event.threadId === threadId && + (event.type === "turn.completed" || event.type === "runtime.error"), + ), + Stream.runForEach((event) => + Effect.sync(() => { + terminalEvents.push(event); + if (event.type === "turn.completed") { + terminalObserved.resolve(undefined); + } + }), + ), + Effect.forkChild, + ); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turn = yield* adapter.sendTurn({ + threadId, + input: "Recover from the exact fifth-probe message", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }); + + for (const delayMs of [250, 500, 1_000, 2_000]) { + yield* advanceTestClock(delayMs); + } + yield* Effect.promise(() => fifthProbeObserved.promise); + yield* advanceTestClock(3_000); + + NodeAssert.equal(runtimeMock.state.messageCalls.length, 5); + NodeAssert.equal( + runtimeMock.state.abortCalls.filter((candidate) => candidate === sessionId).length, + 0, + ); + pushEvent({ + id: "evt-busy-after-fifth-message-probe", + type: "session.status", + properties: { + sessionID: sessionId, + status: { type: "busy" }, + }, + }); + pushEvent({ + id: "evt-idle-after-fifth-message-probe", + type: "session.status", + properties: { + sessionID: sessionId, + status: { type: "idle" }, + }, + }); + yield* Effect.promise(() => terminalObserved.promise); + + const turnCompleted = terminalEvents.find((event) => event.type === "turn.completed"); + NodeAssert.equal(turnCompleted?.payload.state, "completed"); + NodeAssert.equal(turnCompleted?.turnId, turn.turnId); + NodeAssert.equal( + terminalEvents.some((event) => event.type === "runtime.error"), + false, + ); + + yield* Fiber.interrupt(eventsFiber); + yield* adapter.stopSession(threadId); + }), + ); + it.effect("fails admission after bounded probes and grace produce no evidence", () => 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 cc4c70888cdb..6f1420461f9d 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -1487,7 +1487,9 @@ export function makeOpenCodeAdapter( Effect.option, ); if ( - Option.isSome(eventObserved) && + (Option.isSome(eventObserved) || + promptAdmission.messageObserved || + promptAdmission.busyObserved) && context.promptAdmission === promptAdmission && context.activeTurnId === promptAdmission.turnId && context.promptGeneration === promptAdmission.generation &&