diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index d0475e25284..b227ff1ab66 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -248,6 +248,38 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("fails sendTurn for missing sessions through the typed error channel", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const result = yield* adapter + .sendTurn({ + threadId: asThreadId("thread-opencode-missing-send"), + input: "hello", + attachments: [], + }) + .pipe(Effect.result); + + NodeAssert.equal(result._tag, "Failure"); + NodeAssert.equal(result.failure._tag, "ProviderAdapterSessionNotFoundError"); + NodeAssert.equal(result.failure.provider, "opencode"); + NodeAssert.equal(result.failure.threadId, "thread-opencode-missing-send"); + }), + ); + + it.effect("fails stopSession for missing sessions through the typed error channel", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const result = yield* adapter + .stopSession(asThreadId("thread-opencode-missing-stop")) + .pipe(Effect.result); + + NodeAssert.equal(result._tag, "Failure"); + NodeAssert.equal(result.failure._tag, "ProviderAdapterSessionNotFoundError"); + NodeAssert.equal(result.failure.provider, "opencode"); + NodeAssert.equal(result.failure.threadId, "thread-opencode-missing-stop"); + }), + ); + it.effect("stops a configured-server session without trying to own server lifecycle", () => 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 1eb6e47bc19..956905e3a3c 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -228,28 +228,25 @@ function appendTurnItem( resolveTurnSnapshot(context, turnId).items.push(item); } -function ensureSessionContext( +const ensureSessionContext = Effect.fn("ensureSessionContext")(function* ( sessions: ReadonlyMap, threadId: ThreadId, -): OpenCodeSessionContext { +) { const session = sessions.get(threadId); if (!session) { - throw new ProviderAdapterSessionNotFoundError({ + return yield* new ProviderAdapterSessionNotFoundError({ provider: PROVIDER, threadId, }); } - // `ensureSessionContext` is a sync gate used from both sync helpers and - // Effect bodies. `Ref.getUnsafe` is an atomic read of the backing cell — - // no fiber suspension required, which keeps this callable everywhere. - if (Ref.getUnsafe(session.stopped)) { - throw new ProviderAdapterSessionClosedError({ + if (yield* Ref.get(session.stopped)) { + return yield* new ProviderAdapterSessionClosedError({ provider: PROVIDER, threadId, }); } return session; -} +}); function normalizeQuestionRequest(request: QuestionRequest): ReadonlyArray { return request.questions.map((question, index) => ({ @@ -1167,7 +1164,7 @@ export function makeOpenCodeAdapter( ); const sendTurn: OpenCodeAdapterShape["sendTurn"] = Effect.fn("sendTurn")(function* (input) { - const context = ensureSessionContext(sessions, input.threadId); + const context = yield* ensureSessionContext(sessions, input.threadId); // A sendTurn while a turn is active is a steer: OpenCode queues the // prompt into the busy session and the work continues as one turn, so // the active turn id is reused instead of opening a new turn. @@ -1291,7 +1288,7 @@ export function makeOpenCodeAdapter( const interruptTurn: OpenCodeAdapterShape["interruptTurn"] = Effect.fn("interruptTurn")( function* (threadId, turnId) { - const context = ensureSessionContext(sessions, threadId); + const context = yield* ensureSessionContext(sessions, threadId); yield* runOpenCodeSdk("session.abort", () => context.client.session.abort({ sessionID: context.openCodeSessionId }), ).pipe(Effect.mapError(toRequestError)); @@ -1313,7 +1310,7 @@ export function makeOpenCodeAdapter( const respondToRequest: OpenCodeAdapterShape["respondToRequest"] = Effect.fn( "respondToRequest", )(function* (threadId, requestId, decision) { - const context = ensureSessionContext(sessions, threadId); + const context = yield* ensureSessionContext(sessions, threadId); if (!context.pendingPermissions.has(requestId)) { return yield* new ProviderAdapterRequestError({ provider: PROVIDER, @@ -1333,7 +1330,7 @@ export function makeOpenCodeAdapter( const respondToUserInput: OpenCodeAdapterShape["respondToUserInput"] = Effect.fn( "respondToUserInput", )(function* (threadId, requestId, answers) { - const context = ensureSessionContext(sessions, threadId); + const context = yield* ensureSessionContext(sessions, threadId); const request = context.pendingQuestions.get(requestId); if (!request) { return yield* new ProviderAdapterRequestError({ @@ -1355,7 +1352,7 @@ export function makeOpenCodeAdapter( function* (threadId) { const context = sessions.get(threadId); if (!context) { - throw new ProviderAdapterSessionNotFoundError({ + return yield* new ProviderAdapterSessionNotFoundError({ provider: PROVIDER, threadId, }); @@ -1385,7 +1382,7 @@ export function makeOpenCodeAdapter( const readThread: OpenCodeAdapterShape["readThread"] = Effect.fn("readThread")( function* (threadId) { - const context = ensureSessionContext(sessions, threadId); + const context = yield* ensureSessionContext(sessions, threadId); const messages = yield* runOpenCodeSdk("session.messages", () => context.client.session.messages({ sessionID: context.openCodeSessionId, @@ -1411,7 +1408,7 @@ export function makeOpenCodeAdapter( const rollbackThread: OpenCodeAdapterShape["rollbackThread"] = Effect.fn("rollbackThread")( function* (threadId, numTurns) { - const context = ensureSessionContext(sessions, threadId); + const context = yield* ensureSessionContext(sessions, threadId); const messages = yield* runOpenCodeSdk("session.messages", () => context.client.session.messages({ sessionID: context.openCodeSessionId,