From 85b7d2a83dac8d30c84b472e3493a566dc8d29f8 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Tue, 29 Sep 2026 23:21:10 -0600 Subject: [PATCH] fix(opencode): settle interrupted prompt admissions Adapt T3 source e518866d28603da273a9357516dd1aef919f1353 (#12003). Capture the send-owned admission so finalization cannot settle another prompt. --- .../provider/Layers/OpenCodeAdapter.test.ts | 49 +++++++++++++++++++ .../src/provider/Layers/OpenCodeAdapter.ts | 15 +++++- 2 files changed, 63 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index c8fd8133ee..975b25f8b3 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -2,6 +2,7 @@ import * as NodeAssert from "node:assert/strict"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { it } from "@effect/vitest"; import * as Cause from "effect/Cause"; +import * as Clock from "effect/Clock"; import * as Context from "effect/Context"; import * as Crypto from "effect/Crypto"; import * as Deferred from "effect/Deferred"; @@ -833,6 +834,54 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("stops after sendTurn is interrupted before prompt submission", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-interrupted-before-submit"); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + + const baseClock = yield* Clock.clockWith(Effect.succeed); + const admissionPublished = yield* Deferred.make(); + const release = yield* Deferred.make(); + let calls = 0; + const gatedClock: Clock.Clock = { + ...baseClock, + currentTimeMillisUnsafe: () => baseClock.currentTimeMillisUnsafe(), + currentTimeMillis: Effect.suspend(() => { + calls += 1; + return calls === 2 + ? Deferred.succeed(admissionPublished, undefined).pipe( + Effect.andThen(Deferred.await(release)), + Effect.andThen(baseClock.currentTimeMillis), + ) + : baseClock.currentTimeMillis; + }), + }; + + const sendFiber = yield* adapter + .sendTurn({ + threadId, + input: "This prompt must not be submitted", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/kimi-k3", + ), + }) + .pipe(Effect.provideService(Clock.Clock, gatedClock), Effect.exit, Effect.forkChild); + yield* Deferred.await(admissionPublished); + yield* Fiber.interrupt(sendFiber); + NodeAssert.equal(runtimeMock.state.promptCalls.length, 0); + + yield* adapter.stopSession(threadId); + NodeAssert.equal(yield* adapter.hasSession(threadId), false); + NodeAssert.deepEqual(runtimeMock.state.closeCalls, ["http://127.0.0.1:9999"]); + }), + ); + it.effect("aborts a held teardown request before closing the session scope", () => 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 1397cab94d..3d38fee965 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -3597,6 +3597,7 @@ export function makeOpenCodeAdapter( }); } + let submissionSettled: Deferred.Deferred | undefined; return yield* context.promptSemaphore.withPermit( Effect.gen(function* () { yield* requireNotQuarantined(context); @@ -3658,6 +3659,7 @@ export function makeOpenCodeAdapter( }, }; context.promptGeneration = promptGeneration; + submissionSettled = promptAdmission.submissionSettled; context.promptAdmission = promptAdmission; context.activeTurnId = turnId; @@ -3937,7 +3939,18 @@ export function makeOpenCodeAdapter( ? { resumeCursor: context.session.resumeCursor } : {}), }; - }), + }).pipe( + // Stop waits for submission to settle even when the caller is + // interrupted before the prompt fiber exists. Complete only this + // send's admission, while it still owns the prompt permit. + Effect.ensuring( + Effect.suspend(() => + submissionSettled + ? Deferred.succeed(submissionSettled, undefined).pipe(Effect.ignore) + : Effect.void, + ), + ), + ), ); });