From f92bc4995a5aadfb68e91b7a7da9f07d9ab9124b Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Wed, 7 Oct 2026 17:37:02 +0200 Subject: [PATCH 1/7] fix(server): stop silently dropping OpenCode child requests OpenCode child permissions and questions disappeared after five routing retries, leaving the native request waiting without a visible failure. Keep the existing five-second backoff, then emit a warning and provider-session error. Reject permissions and questions only when the native parent chain proves this runtime owns them. Bound rejection delivery to ten seconds and keep failed rejections available for later routing. Ignore requests whose complete parent chain belongs to another runtime. The SDK supports both rejection paths, and OpenCode waits on these requests without a request timeout. Longer or unbounded retries only postpone the failure. Blind rejection can cancel another session on an external server, so unresolved ownership emits a failure without guessing who owns the request. Add deterministic clock and event tests for exhausted routing, failed and timed-out rejection, duplicate asks, late relations, externally settled requests, and unrelated sessions. Model: gpt-6.1-sol (Codex) --- .../Adapters/OpenCodeAdapterV2.test.ts | 200 ++++++++++++++++++ .../Adapters/OpenCodeAdapterV2.ts | 47 +++- 2 files changed, 241 insertions(+), 6 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts index a8e28a377d1a..a134b96e5ac5 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts @@ -459,6 +459,206 @@ describe("OpenCodeAdapterV2", () => { }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), ); + it.effect.each( + (["permission", "question"] as const).flatMap((kind) => + (["unknown-owner", "inactive-owner", "rejection-error", "rejection-timeout"] as const).map( + (outcome) => ({ kind, outcome }), + ), + ), + )("reports an unroutable $kind request with $outcome", ({ kind, outcome }) => + Effect.gen(function* () { + const nativeEvents = asyncEventStream(); + const failed = yield* Deferred.make(); + const rejections: unknown[] = []; + let rejectionSignal: AbortSignal | undefined; + const reject = async (input: unknown, options: { signal: AbortSignal }) => { + rejections.push(input); + rejectionSignal = options.signal; + if (outcome === "rejection-error") return { error: { message: "rejection failed" } }; + if (outcome === "rejection-timeout") return new Promise(() => {}); + return { data: true }; + }; + const harness = yield* makeOpenCodeRuntimeHarness(`unroutable-${kind}-${outcome}`, "root", { + event: { subscribe: async () => ({ stream: nativeEvents.stream }) }, + session: { + create: async () => ({ data: { id: "root", time: { created: 1, updated: 1 } } }), + get: async ({ sessionID }: { sessionID: string }) => + outcome === "unknown-owner" && sessionID === "child" + ? { error: { message: "session relation unavailable" } } + : { data: { id: sessionID, ...(sessionID === "root" ? {} : { parentID: "root" }) } }, + messages: async () => ({ data: [] }), + promptAsync: async () => ({ data: true }), + abort: async () => ({ data: true }), + children: async () => ({ data: [] }), + }, + permission: { reply: reject }, + question: { reject }, + }); + if (outcome === "unknown-owner") yield* harness.startTurn(); + const collected = yield* harness.runtime.events.pipe( + Stream.tap((event) => + event.type === "provider_session.updated" && event.providerSession.status === "error" + ? Deferred.succeed(failed, undefined) + : Effect.void, + ), + Stream.takeUntil( + (event) => + event.type === "provider_session.updated" && event.providerSession.status === "error", + ), + Stream.runCollect, + Effect.forkScoped, + ); + const asked = { + type: `${kind}.asked`, + properties: + kind === "permission" + ? { + id: "unroutable", + sessionID: "child", + permission: "bash", + patterns: ["*"], + always: [], + metadata: {}, + } + : { + id: "unroutable", + sessionID: "child", + questions: [ + { + header: "Choice", + question: "Which?", + options: [{ label: "Yes", description: "Proceed" }], + }, + ], + }, + }; + yield* Effect.promise(() => nativeEvents.push(asked)); + yield* Effect.promise(() => nativeEvents.push(asked)); + yield* TestClock.adjust("4999 millis"); + assert.isFalse(yield* Deferred.isDone(failed)); + assert.isEmpty(rejections); + yield* TestClock.adjust("1 milli"); + if (outcome === "rejection-timeout") { + yield* TestClock.adjust("10 seconds"); + assert.isTrue(rejectionSignal?.aborted); + } + if (outcome !== "unknown-owner") assert.lengthOf(rejections, 1); + yield* Deferred.await(failed); + const received = yield* Fiber.join(collected); + const failure = received.at(-1)!; + assert.equal(failure.type, "provider_session.updated"); + if (failure.type !== "provider_session.updated") return; + assert.include(failure.providerSession.lastError, `${kind} request unroutable`); + assert.equal( + failure.providerSession.lastError?.endsWith("Native rejection failed."), + outcome === "rejection-error" || outcome === "rejection-timeout", + ); + assert.isFalse(received.some((event) => event.type === "runtime_request.updated")); + assert.deepEqual( + rejections, + outcome === "unknown-owner" + ? [] + : [{ requestID: "unroutable", ...(kind === "permission" ? { reply: "reject" } : {}) }], + ); + // A failed rejection leaves the request available if its owner becomes active. + if (outcome === "inactive-owner" || outcome === "rejection-error") { + yield* harness.startTurn(); + yield* Effect.promise(() => nativeEvents.push(asked)); + const snapshot = yield* harness.runtime.readThreadSnapshot({ + providerThread: harness.providerThread, + }); + assert.lengthOf(snapshot.runtimeRequests, outcome === "inactive-owner" ? 0 : 1); + assert.lengthOf(rejections, 1); + } + }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), + ); + + it.effect.each(["late-relation", "already-settled", "foreign-owner"] as const)( + "preserves a child request with %s during routing backoff", + (outcome) => + Effect.gen(function* () { + const nativeEvents = asyncEventStream(); + const rejections: unknown[] = []; + const harness = yield* makeOpenCodeRuntimeHarness(`routing-${outcome}`, "root", { + event: { subscribe: async () => ({ stream: nativeEvents.stream }) }, + session: { + create: async () => ({ data: { id: "root", time: { created: 1, updated: 1 } } }), + get: async ({ sessionID }: { sessionID: string }) => + sessionID === "root" || outcome === "foreign-owner" + ? { + data: { + id: sessionID, + ...(sessionID === "child" ? { parentID: "foreign-root" } : {}), + }, + } + : { error: { message: "session relation unavailable" } }, + messages: async () => ({ data: [] }), + promptAsync: async () => ({ data: true }), + abort: async () => ({ data: true }), + children: async () => ({ data: [] }), + }, + permission: { + reply: async (input: unknown) => { + rejections.push(input); + return { data: true }; + }, + }, + }); + yield* harness.startTurn(); + const received: Array<{ + type: string; + providerSession?: { status: string }; + }> = []; + const routed = yield* Deferred.make(); + yield* harness.runtime.events.pipe( + Stream.runForEach((event) => { + received.push(event); + return event.type === "runtime_request.updated" + ? Deferred.succeed(routed, undefined) + : Effect.void; + }), + Effect.forkScoped, + ); + const asked = { + type: "permission.asked", + properties: { + id: "child-request", + sessionID: "child", + permission: "bash", + patterns: ["*"], + always: [], + metadata: {}, + }, + }; + yield* Effect.promise(() => nativeEvents.push(asked)); + if (outcome === "already-settled") { + yield* Effect.promise(() => + nativeEvents.push({ + type: "permission.replied", + properties: { requestID: "child-request", sessionID: "child", reply: "once" }, + }), + ); + } + if (outcome !== "foreign-owner") { + yield* Effect.promise(() => + nativeEvents.push({ + type: "session.created", + properties: { info: { id: "child", parentID: "root" } }, + }), + ); + } + yield* Effect.promise(() => nativeEvents.push(asked)); + yield* TestClock.adjust("5 seconds"); + if (outcome === "late-relation") yield* Deferred.await(routed); + const snapshot = yield* harness.runtime.readThreadSnapshot({ + providerThread: harness.providerThread, + }); + assert.lengthOf(snapshot.runtimeRequests, outcome === "late-relation" ? 1 : 0); + assert.isEmpty(rejections); + assert.isFalse(received.some((event) => event.providerSession?.status === "error")); + }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), + ); + it.effect("aborts external root and descendants before closing the event stream", () => Effect.gen(function* () { const scope = yield* Scope.make(); diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts index 510a4c859d36..74eef5bfe97f 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts @@ -1978,8 +1978,9 @@ export function makeOpenCodeAdapterV2( /** Resolve the thread state a session belongs to: its own, or for a * child not registered yet, the nearest known ancestor found by walking - * the native parent chain. Registers every hop so later requests from - * the same child resolve without another lookup. */ + * the native parent chain. A complete chain ending at an unknown root + * returns null: it belongs to another runtime. Registers every hop so + * later requests from the same child resolve without another lookup. */ const resolveSessionOwner = Effect.fnUntraced(function* (sessionId: string) { const known = threads.get(sessionId) ?? relatedSessionOwners.get(sessionId); if (known !== undefined) return known; @@ -1990,8 +1991,9 @@ export function makeOpenCodeAdapterV2( client.session.get({ sessionID: cursor }), ).pipe(Effect.option); const info = Option.getOrUndefined(response)?.data; + if (info === undefined) return undefined; const parentId = info?.parentID; - if (parentId === undefined) return undefined; + if (parentId === undefined) return null; hops.push(cursor); const owner = threads.get(parentId) ?? relatedSessionOwners.get(parentId); if (owner !== undefined) { @@ -2009,8 +2011,8 @@ export function makeOpenCodeAdapterV2( * arrive before the task part or session.created event that reveals * its relation to a thread. The first resolution attempt runs inline * (the replayable path); if the relation or the owning turn is not - * established yet, a short forked backoff keeps trying instead of - * dropping the request. */ + * established yet, a short forked backoff keeps trying. Exhaustion + * reports a session error and rejects requests known to belong here. */ const routeRuntimeRequest = Effect.fnUntraced(function* ( nativeRequestId: string, sessionId: string, @@ -2019,6 +2021,7 @@ export function makeOpenCodeAdapterV2( | { readonly type: "question"; readonly value: QuestionRequest }, ) { if (pendingChildRequestRoutes.has(nativeRequestId)) return; + let unresolvedState: OpenCodeThreadState | null | undefined; const attempt = Effect.gen(function* () { if ( settledNativeRequestIds.has(nativeRequestId) || @@ -2027,7 +2030,8 @@ export function makeOpenCodeAdapterV2( return true; } const state = yield* resolveSessionOwner(sessionId); - const owner = state === undefined ? undefined : topLevelRequestOwner(state); + unresolvedState = state; + const owner = state == null ? undefined : topLevelRequestOwner(state); if (owner === undefined) return false; yield* emitRuntimeRequest(owner, nativeRequestId, request); return true; @@ -2039,6 +2043,37 @@ export function makeOpenCodeAdapterV2( yield* Effect.sleep(Duration.millis(Math.min(200 * 2 ** retry, 2_000))); if (yield* attempt) return; } + if (unresolvedState === null) return; + const detail = `OpenCode ${request.type} request ${nativeRequestId} could not be routed to an active thread.`; + yield* Effect.logWarning(detail, { + nativeSessionId: sessionId, + providerSessionId: input.providerSessionId, + }); + // External servers broadcast other sessions' requests too. An + // unresolved relation is not proof that we own the native request. + if (unresolvedState !== undefined) { + const rejected = yield* ( + request.type === "permission" + ? sdkCall( + "permission.reply", + { requestID: nativeRequestId, reply: "reject" }, + (signal) => + client.permission.reply( + { requestID: nativeRequestId, reply: "reject" }, + { signal }, + ), + ).pipe(Effect.map((response) => unwrapData("permission.reply", response))) + : sdkCall("question.reject", { requestID: nativeRequestId }, (signal) => + client.question.reject({ requestID: nativeRequestId }, { signal }), + ).pipe(Effect.map((response) => unwrapData("question.reject", response))) + ).pipe(Effect.timeout("10 seconds"), Effect.exit); + if (Exit.isFailure(rejected)) { + yield* updateProviderSession("error", `${detail} Native rejection failed.`); + return; + } + rememberSettledRequest(nativeRequestId); + } + yield* updateProviderSession("error", detail); }).pipe( Effect.ensuring(Effect.sync(() => pendingChildRequestRoutes.delete(nativeRequestId))), Effect.forkIn(scope), From 05c1e52cfa248bf8b3a8f5ddcea21c4ef966f017 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Wed, 7 Oct 2026 18:21:09 +0200 Subject: [PATCH 2/7] fix(server): keep a proven-foreign OpenCode request from being rejected Routing retries overwrote an earlier unknown-root classification with transient lookup failures, allowing exhaustion to report a sibling request as a session error. Native rejection also treated an already-answered request as a delivery failure. Keep the observed foreign classification sticky through exhaustion. A failed lookup provides no new ownership evidence and must not overturn the earlier classification. Treat native NotFoundError as settled, so a later route cannot resurrect the answered request. Add TestClock regressions for permissions and questions, document the native permission rejection cascade, and clarify that an unknown root can also be an unregistered local root. Model: gpt-6.1-sol (Codex) --- .../Adapters/OpenCodeAdapterV2.test.ts | 100 ++++++++++++++++++ .../Adapters/OpenCodeAdapterV2.ts | 23 +++- 2 files changed, 118 insertions(+), 5 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts index a134b96e5ac5..03215f9119b1 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts @@ -659,6 +659,106 @@ describe("OpenCodeAdapterV2", () => { }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), ); + it.effect.each( + (["permission", "question"] as const).flatMap((kind) => + (["foreign-then-unavailable", "already-answered"] as const).map((outcome) => ({ + kind, + outcome, + })), + ), + )("preserves an unroutable $kind request with $outcome", ({ kind, outcome }) => + Effect.gen(function* () { + const nativeEvents = asyncEventStream(); + const rejections: unknown[] = []; + let lookups = 0; + const reject = async (input: unknown) => { + rejections.push(input); + if (outcome === "already-answered") { + // throwOnError SDK clients throw the native error response. + throw { name: "NotFoundError", data: { message: "Request already answered" } }; + } + return { data: true }; + }; + const harness = yield* makeOpenCodeRuntimeHarness(`routing-${kind}-${outcome}`, "root", { + event: { subscribe: async () => ({ stream: nativeEvents.stream }) }, + session: { + create: async () => ({ data: { id: "root", time: { created: 1, updated: 1 } } }), + get: async ({ sessionID }: { sessionID: string }) => { + if (sessionID === "root") return { data: { id: sessionID } }; + lookups += 1; + if (outcome === "foreign-then-unavailable") { + if (lookups === 1) return { data: { id: sessionID } }; + throw new Error("Session lookup temporarily unavailable"); + } + return { data: { id: sessionID, parentID: "root" } }; + }, + messages: async () => ({ data: [] }), + promptAsync: async () => ({ data: true }), + abort: async () => ({ data: true }), + children: async () => ({ data: [] }), + }, + permission: { reply: reject }, + question: { reject }, + }); + if (outcome === "foreign-then-unavailable") yield* harness.startTurn(); + const received: Array<{ + type: string; + providerSession?: { status: string }; + }> = []; + yield* harness.runtime.events.pipe( + Stream.runForEach((event) => Effect.sync(() => received.push(event))), + Effect.forkScoped, + ); + const asked = { + type: `${kind}.asked`, + properties: + kind === "permission" + ? { + id: "racing-request", + sessionID: "asking-session", + permission: "bash", + patterns: ["*"], + always: [], + metadata: {}, + } + : { + id: "racing-request", + sessionID: "asking-session", + questions: [ + { + header: "Choice", + question: "Which?", + options: [{ label: "Yes", description: "Proceed" }], + }, + ], + }, + }; + yield* Effect.promise(() => nativeEvents.push(asked)); + yield* TestClock.adjust("4999 millis"); + assert.isEmpty(rejections); + assert.isFalse(received.some((event) => event.providerSession?.status === "error")); + yield* TestClock.adjust("1 milli"); + // Drain the final retry and its queued events without advancing the clock. + yield* TestClock.adjust("0 millis"); + assert.lengthOf(rejections, outcome === "already-answered" ? 1 : 0); + assert.equal(lookups, outcome === "foreign-then-unavailable" ? 6 : 1); + assert.isFalse(received.some((event) => event.providerSession?.status === "error")); + if (outcome === "already-answered") { + yield* harness.startTurn(); + // A replay after the owner becomes active must not revive the answered request. + yield* Effect.promise(() => nativeEvents.push(asked)); + yield* TestClock.adjust("5 seconds"); + assert.lengthOf(rejections, 1); + assert.isFalse(received.some((event) => event.providerSession?.status === "error")); + } + const snapshot = yield* harness.runtime.readThreadSnapshot({ + providerThread: harness.providerThread, + }); + assert.isEmpty(snapshot.runtimeRequests); + assert.isFalse(received.some((event) => event.type === "runtime_request.updated")); + }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), + ); + it.effect("aborts external root and descendants before closing the event stream", () => Effect.gen(function* () { const scope = yield* Scope.make(); diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts index 74eef5bfe97f..c88a33c78e63 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts @@ -1979,8 +1979,10 @@ export function makeOpenCodeAdapterV2( /** Resolve the thread state a session belongs to: its own, or for a * child not registered yet, the nearest known ancestor found by walking * the native parent chain. A complete chain ending at an unknown root - * returns null: it belongs to another runtime. Registers every hop so - * later requests from the same child resolve without another lookup. */ + * returns null: no registered thread owns it. This usually means + * another runtime, but can also be a root not registered here yet. + * Registers every owned hop so later requests from the same child + * resolve without another lookup. */ const resolveSessionOwner = Effect.fnUntraced(function* (sessionId: string) { const known = threads.get(sessionId) ?? relatedSessionOwners.get(sessionId); if (known !== undefined) return known; @@ -2022,6 +2024,7 @@ export function makeOpenCodeAdapterV2( ) { if (pendingChildRequestRoutes.has(nativeRequestId)) return; let unresolvedState: OpenCodeThreadState | null | undefined; + let provedForeign = false; const attempt = Effect.gen(function* () { if ( settledNativeRequestIds.has(nativeRequestId) || @@ -2031,6 +2034,7 @@ export function makeOpenCodeAdapterV2( } const state = yield* resolveSessionOwner(sessionId); unresolvedState = state; + if (state === null) provedForeign = true; const owner = state == null ? undefined : topLevelRequestOwner(state); if (owner === undefined) return false; yield* emitRuntimeRequest(owner, nativeRequestId, request); @@ -2043,15 +2047,20 @@ export function makeOpenCodeAdapterV2( yield* Effect.sleep(Duration.millis(Math.min(200 * 2 ** retry, 2_000))); if (yield* attempt) return; } - if (unresolvedState === null) return; + // A later lookup failure must not erase an observed unknown root. + if (provedForeign) return; const detail = `OpenCode ${request.type} request ${nativeRequestId} could not be routed to an active thread.`; yield* Effect.logWarning(detail, { nativeSessionId: sessionId, providerSessionId: input.providerSessionId, }); - // External servers broadcast other sessions' requests too. An - // unresolved relation is not proof that we own the native request. + // External servers broadcast other sessions' requests too. Reject + // only when the parent chain resolves to a registered thread here. if (unresolvedState !== undefined) { + // OpenCode permission rejection cancels every pending permission + // in the asking native session, including concurrent requests. + // The SDK has no single-request denial; question.reject cancels + // only the named question. const rejected = yield* ( request.type === "permission" ? sdkCall( @@ -2068,6 +2077,10 @@ export function makeOpenCodeAdapterV2( ).pipe(Effect.map((response) => unwrapData("question.reject", response))) ).pipe(Effect.timeout("10 seconds"), Effect.exit); if (Exit.isFailure(rejected)) { + if (isOpenCodeNotFound(Cause.squash(rejected.cause))) { + rememberSettledRequest(nativeRequestId); + return; + } yield* updateProviderSession("error", `${detail} Native rejection failed.`); return; } From 5b290cf72ae3e0f3b32b11e1eb8bccac44734c8e Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Thu, 8 Oct 2026 14:14:11 +0200 Subject: [PATCH 3/7] fix(server): preserve newer OpenCode session status after routing give-up --- .../Adapters/OpenCodeAdapterV2.test.ts | 190 +++++++++++++++++- .../Adapters/OpenCodeAdapterV2.ts | 13 +- 2 files changed, 200 insertions(+), 3 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts index 03215f9119b1..7ec759c60608 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts @@ -49,7 +49,7 @@ import { OPENCODE_PROVIDER, reconcileOpenCodePromptAdmissionStatus, } from "./OpenCodeAdapterV2.ts"; -import { ProviderAdapterV2RuntimePolicy } from "../ProviderAdapter.ts"; +import { ProviderAdapterV2RuntimePolicy, type ProviderAdapterV2Event } from "../ProviderAdapter.ts"; const encodeUnknownJson = Schema.encodeUnknownSync(Schema.fromJsonString(Schema.Unknown)); const OPEN_CODE_TEST_SETTINGS = Schema.decodeSync(OpenCodeSettings)({ @@ -573,6 +573,194 @@ describe("OpenCodeAdapterV2", () => { }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), ); + it.effect.each([ + ...(["permission", "question"] as const).flatMap((kind) => + (["success", "failure", "timeout", "not-found"] as const).flatMap((outcome) => + (["unchanged", "running", "completed"] as const).map((lifecycle) => ({ + kind, + outcome, + lifecycle, + })), + ), + ), + ...(["success", "failure"] as const).map((outcome) => ({ + kind: "question" as const, + outcome, + lifecycle: "timestamp-yield" as const, + })), + { kind: "question" as const, outcome: "failure" as const, lifecycle: "other-root" as const }, + ])( + "keeps $lifecycle session state after late $kind rejection $outcome", + ({ kind, outcome, lifecycle }) => + Effect.gen(function* () { + const nativeEvents = asyncEventStream(); + const rejectionStarted = promiseGate(); + const response = promiseGate(); + const clockRead = yield* Deferred.make(); + const releaseClockRead = yield* Deferred.make(); + const baseClock = yield* Clock.Clock; + let blockNextClockRead = false; + const blockingClock: Clock.Clock = { + ...baseClock, + currentTimeMillis: Effect.suspend(() => { + if (!blockNextClockRead) return baseClock.currentTimeMillis; + blockNextClockRead = false; + return Deferred.succeed(clockRead, undefined).pipe( + Effect.andThen(Deferred.await(releaseClockRead)), + Effect.andThen(baseClock.currentTimeMillis), + ); + }), + }; + let rejectionSignal: AbortSignal | undefined; + let rejectionCalls = 0; + let createdSessions = 0; + const reject = async (_input: unknown, options: { signal: AbortSignal }) => { + rejectionCalls += 1; + rejectionSignal = options.signal; + rejectionStarted.resolve(); + await response.promise; + if (outcome === "failure") return { error: { message: "rejection failed" } }; + if (outcome === "not-found") throw { name: "NotFoundError" }; + return { data: true }; + }; + const harness = yield* makeOpenCodeRuntimeHarness( + `late-${kind}-${outcome}-${lifecycle}`, + "root", + { + event: { subscribe: async () => ({ stream: nativeEvents.stream }) }, + session: { + create: async () => ({ + data: { + id: createdSessions++ === 0 ? "root" : "other-root", + time: { created: 1, updated: 1 }, + }, + }), + get: async ({ sessionID }: { sessionID: string }) => ({ + data: { id: sessionID, ...(sessionID === "child" ? { parentID: "root" } : {}) }, + }), + messages: async () => ({ data: [] }), + promptAsync: async () => ({ data: true }), + summarize: async () => ({ data: true }), + abort: async () => ({ data: true }), + children: async () => ({ data: [] }), + }, + permission: { reply: reject }, + question: { reject }, + }, + ).pipe(Effect.provideService(Clock.Clock, blockingClock)); + const received: ProviderAdapterV2Event[] = []; + yield* harness.runtime.events.pipe( + Stream.runForEach((event) => Effect.sync(() => received.push(event))), + Effect.forkScoped, + ); + const asked = { + type: `${kind}.asked`, + properties: + kind === "permission" + ? { + id: "old-request", + sessionID: "child", + permission: "bash", + patterns: ["*"], + always: [], + metadata: {}, + } + : { + id: "old-request", + sessionID: "child", + questions: [ + { + header: "Choice", + question: "Which?", + options: [{ label: "Yes", description: "Proceed" }], + }, + ], + }, + }; + yield* Effect.promise(() => nativeEvents.push(asked)); + yield* TestClock.adjust("5 seconds"); + yield* Effect.promise(() => rejectionStarted.promise); + + if (lifecycle === "running" || lifecycle === "completed") { + yield* harness.startTurn(lifecycle === "completed" ? "/compact" : "new turn"); + } else if (lifecycle === "other-root") { + const otherThread = yield* harness.runtime.ensureThread({ + threadId: ThreadId.make("thread-other-root"), + modelSelection: { + instanceId: harness.providerThread.providerInstanceId, + model: "anthropic/claude-sonnet", + options: [], + }, + runtimePolicy: harness.policy, + }); + yield* harness.startTurn("another root's turn", otherThread); + } + // Hold the timestamp read itself to exercise the final mutation boundary. + if (lifecycle === "timestamp-yield") blockNextClockRead = true; + if (outcome === "timeout") { + yield* TestClock.adjust("10 seconds"); + assert.isTrue(rejectionSignal?.aborted); + } else { + response.resolve(); + } + if (lifecycle === "timestamp-yield") { + yield* Deferred.await(clockRead); + yield* harness.startTurn("new turn during old status timestamp read"); + yield* Deferred.succeed(releaseClockRead, undefined); + } + // Drain runnable test fibers after the SDK response or deadline. No wall-clock wait. + yield* TestClock.adjust("0 millis"); + const sessionUpdates = received.filter( + (event) => event.type === "provider_session.updated", + ); + const lastSession = + sessionUpdates.at(-1)?.providerSession ?? harness.runtime.providerSession; + if (lifecycle === "unchanged" && outcome !== "not-found") { + assert.equal(lastSession.status, "error"); + assert.include(lastSession.lastError, "old-request"); + assert.equal( + lastSession.lastError?.endsWith("Native rejection failed."), + outcome !== "success", + ); + } else { + assert.equal( + lastSession.status, + lifecycle === "unchanged" || lifecycle === "completed" ? "ready" : "running", + ); + assert.isNull(lastSession.lastError); + assert.isFalse(sessionUpdates.some((event) => event.providerSession.status === "error")); + } + if ( + lifecycle === "running" || + lifecycle === "timestamp-yield" || + lifecycle === "completed" + ) { + const snapshot = yield* harness.runtime.readThreadSnapshot({ + providerThread: harness.providerThread, + }); + assert.equal( + snapshot.providerTurns.at(-1)?.status, + lifecycle === "completed" ? "completed" : "running", + ); + } + assert.equal(rejectionCalls, 1); + + // Successful/NotFound delivery stays settled; failures can route on replay. + if (lifecycle === "unchanged" || lifecycle === "completed" || lifecycle === "other-root") { + yield* harness.startTurn(); + } + yield* Effect.promise(() => nativeEvents.push(asked)); + const replay = yield* harness.runtime.readThreadSnapshot({ + providerThread: harness.providerThread, + }); + assert.lengthOf( + replay.runtimeRequests, + outcome === "failure" || outcome === "timeout" ? 1 : 0, + ); + assert.equal(rejectionCalls, 1); + }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), + ); + it.effect.each(["late-relation", "already-settled", "foreign-owner"] as const)( "preserves a child request with %s during routing backoff", (outcome) => diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts index c88a33c78e63..85a0bedf17bb 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts @@ -1114,9 +1114,13 @@ export function makeOpenCodeAdapterV2( const updateProviderSession = ( status: OrchestrationV2ProviderSession["status"], lastError: string | null = sessionEntity.lastError, + expectedSession?: OrchestrationV2ProviderSession, ) => Effect.gen(function* () { const updatedAt = yield* DateTime.now; + // Check after the clock yield so late routing work cannot overwrite + // a newer session update, including a turn that already finished. + if (expectedSession !== undefined && sessionEntity !== expectedSession) return; sessionEntity = { ...sessionEntity, status, lastError, updatedAt }; yield* emitProviderEvent({ type: "provider_session.updated", @@ -2023,6 +2027,7 @@ export function makeOpenCodeAdapterV2( | { readonly type: "question"; readonly value: QuestionRequest }, ) { if (pendingChildRequestRoutes.has(nativeRequestId)) return; + const routingSession = sessionEntity; let unresolvedState: OpenCodeThreadState | null | undefined; let provedForeign = false; const attempt = Effect.gen(function* () { @@ -2081,12 +2086,16 @@ export function makeOpenCodeAdapterV2( rememberSettledRequest(nativeRequestId); return; } - yield* updateProviderSession("error", `${detail} Native rejection failed.`); + yield* updateProviderSession( + "error", + `${detail} Native rejection failed.`, + routingSession, + ); return; } rememberSettledRequest(nativeRequestId); } - yield* updateProviderSession("error", detail); + yield* updateProviderSession("error", detail, routingSession); }).pipe( Effect.ensuring(Effect.sync(() => pendingChildRequestRoutes.delete(nativeRequestId))), Effect.forkIn(scope), From 8cbbdf6de30d72b9a0dad31570793af79f5fb0f4 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Fri, 9 Oct 2026 01:07:20 +0200 Subject: [PATCH 4/7] test(server): verify OpenCode SDK missing request handling --- .../Adapters/OpenCodeAdapterV2.test.ts | 38 +++++++++++++++---- 1 file changed, 30 insertions(+), 8 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts index 27166d8954ec..0299c06ff56e 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts @@ -1,6 +1,6 @@ import { assert, describe, it } from "@effect/vitest"; import * as NodeServices from "@effect/platform-node/NodeServices"; -import type { OpencodeClient, ToolPart } from "@opencode-ai/sdk/v2"; +import { createOpencodeClient, type OpencodeClient, type ToolPart } from "@opencode-ai/sdk/v2"; import { CheckpointId, NodeId, @@ -852,22 +852,44 @@ describe("OpenCodeAdapterV2", () => { it.effect.each( (["permission", "question"] as const).flatMap((kind) => - (["foreign-then-unavailable", "already-answered"] as const).map((outcome) => ({ - kind, - outcome, - })), + (["foreign-then-unavailable", "already-answered", "sdk-already-answered"] as const).map( + (outcome) => ({ + kind, + outcome, + }), + ), ), )("preserves an unroutable $kind request with $outcome", ({ kind, outcome }) => Effect.gen(function* () { const nativeEvents = asyncEventStream(); const rejections: unknown[] = []; let lookups = 0; + // Exercise the installed SDK's throwOnError path, including the public + // client's error interceptor that adds HTTP status to the tagged body. + const sdk = createOpencodeClient({ + baseUrl: "http://test.invalid", + throwOnError: true, + fetch: async () => + Response.json( + { + _tag: kind === "permission" ? "PermissionNotFoundError" : "QuestionNotFoundError", + requestID: "racing-request", + message: "Request already answered", + }, + { status: 404 }, + ), + }); const reject = async (input: unknown) => { rejections.push(input); if (outcome === "already-answered") { - // throwOnError SDK clients throw the native error response. + // Older clients can throw the named native error response directly. throw { name: "NotFoundError", data: { message: "Request already answered" } }; } + if (outcome === "sdk-already-answered") { + return kind === "permission" + ? sdk.permission.reply({ requestID: "racing-request", reply: "reject" }) + : sdk.question.reject({ requestID: "racing-request" }); + } return { data: true }; }; const harness = yield* makeOpenCodeRuntimeHarness(`routing-${kind}-${outcome}`, "root", { @@ -931,10 +953,10 @@ describe("OpenCodeAdapterV2", () => { yield* TestClock.adjust("1 milli"); // Drain the final retry and its queued events without advancing the clock. yield* TestClock.adjust("0 millis"); - assert.lengthOf(rejections, outcome === "already-answered" ? 1 : 0); + assert.lengthOf(rejections, outcome === "foreign-then-unavailable" ? 0 : 1); assert.equal(lookups, outcome === "foreign-then-unavailable" ? 6 : 1); assert.isFalse(received.some((event) => event.providerSession?.status === "error")); - if (outcome === "already-answered") { + if (outcome !== "foreign-then-unavailable") { yield* harness.startTurn(); // A replay after the owner becomes active must not revive the answered request. yield* Effect.promise(() => nativeEvents.push(asked)); From 9fa09984005abb7707b253384018e862da76d223 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Fri, 9 Oct 2026 01:34:13 +0200 Subject: [PATCH 5/7] test(server): drain OpenCode retries before assertions --- .../src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts index 0299c06ff56e..eb4aeab5edb5 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts @@ -840,6 +840,8 @@ describe("OpenCodeAdapterV2", () => { } yield* Effect.promise(() => nativeEvents.push(asked)); yield* TestClock.adjust("5 seconds"); + // Drain the final retry and its queued events without advancing the clock. + yield* TestClock.adjust("0 millis"); if (outcome === "late-relation") yield* Deferred.await(routed); const snapshot = yield* harness.runtime.readThreadSnapshot({ providerThread: harness.providerThread, From a391652896d778f74bd3d5c38154998258006949 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Fri, 9 Oct 2026 13:30:52 +0200 Subject: [PATCH 6/7] fix(server): keep orphan permission rejection from cancelling newer requests --- .../Adapters/OpenCodeAdapterV2.test.ts | 217 ++++++++++++++---- .../Adapters/OpenCodeAdapterV2.ts | 43 ++-- 2 files changed, 194 insertions(+), 66 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts index eb4aeab5edb5..87513e661536 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.test.ts @@ -188,6 +188,7 @@ const makeOpenCodeRuntimeHarness = Effect.fn("makeOpenCodeRuntimeHarness")(funct const startTurn = ( text = "hello", startProviderThread: OrchestrationV2ProviderThread = providerThread, + ordinal = 1, ) => runtime.startTurn({ appThread: { @@ -214,16 +215,20 @@ const makeOpenCodeRuntimeHarness = Effect.fn("makeOpenCodeRuntimeHarness")(funct deletedAt: null, }, threadId, - runId: RunId.make(`run-opencode-${suffix}`), - runOrdinal: 1, - providerTurnOrdinal: 1, - attemptId: RunAttemptId.make(`attempt-opencode-${suffix}`), - rootNodeId: NodeId.make(`node-opencode-${suffix}`), + runId: RunId.make(`run-opencode-${suffix}${ordinal === 1 ? "" : `-${ordinal}`}`), + runOrdinal: ordinal, + providerTurnOrdinal: ordinal, + attemptId: RunAttemptId.make( + `attempt-opencode-${suffix}${ordinal === 1 ? "" : `-${ordinal}`}`, + ), + rootNodeId: NodeId.make(`node-opencode-${suffix}${ordinal === 1 ? "" : `-${ordinal}`}`), providerThread: startProviderThread, message: { createdBy: "user", creationSource: "web", - messageId: MessageId.make(`message-opencode-${suffix}`), + messageId: MessageId.make( + `message-opencode-${suffix}${ordinal === 1 ? "" : `-${ordinal}`}`, + ), text, attachments: [], }, @@ -464,9 +469,10 @@ describe("OpenCodeAdapterV2", () => { it.effect.each( (["permission", "question"] as const).flatMap((kind) => - (["unknown-owner", "inactive-owner", "rejection-error", "rejection-timeout"] as const).map( - (outcome) => ({ kind, outcome }), - ), + (kind === "permission" + ? (["unknown-owner", "inactive-owner"] as const) + : (["unknown-owner", "inactive-owner", "rejection-error", "rejection-timeout"] as const) + ).map((outcome) => ({ kind, outcome })), ), )("reports an unroutable $kind request with $outcome", ({ kind, outcome }) => Effect.gen(function* () { @@ -545,13 +551,16 @@ describe("OpenCodeAdapterV2", () => { yield* TestClock.adjust("10 seconds"); assert.isTrue(rejectionSignal?.aborted); } - if (outcome !== "unknown-owner") assert.lengthOf(rejections, 1); + if (kind === "question" && outcome !== "unknown-owner") assert.lengthOf(rejections, 1); yield* Deferred.await(failed); const received = yield* Fiber.join(collected); const failure = received.at(-1)!; assert.equal(failure.type, "provider_session.updated"); if (failure.type !== "provider_session.updated") return; assert.include(failure.providerSession.lastError, `${kind} request unroutable`); + if (kind === "permission") { + assert.include(failure.providerSession.lastError, "Native permission remains pending"); + } assert.equal( failure.providerSession.lastError?.endsWith("Native rejection failed."), outcome === "rejection-error" || outcome === "rejection-timeout", @@ -559,9 +568,7 @@ describe("OpenCodeAdapterV2", () => { assert.isFalse(received.some((event) => event.type === "runtime_request.updated")); assert.deepEqual( rejections, - outcome === "unknown-owner" - ? [] - : [{ requestID: "unroutable", ...(kind === "permission" ? { reply: "reject" } : {}) }], + outcome === "unknown-owner" || kind === "permission" ? [] : [{ requestID: "unroutable" }], ); // A failed rejection leaves the request available if its owner becomes active. if (outcome === "inactive-owner" || outcome === "rejection-error") { @@ -570,14 +577,17 @@ describe("OpenCodeAdapterV2", () => { const snapshot = yield* harness.runtime.readThreadSnapshot({ providerThread: harness.providerThread, }); - assert.lengthOf(snapshot.runtimeRequests, outcome === "inactive-owner" ? 0 : 1); - assert.lengthOf(rejections, 1); + assert.lengthOf( + snapshot.runtimeRequests, + kind === "question" && outcome === "inactive-owner" ? 0 : 1, + ); + assert.lengthOf(rejections, kind === "permission" ? 0 : 1); } }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), ); it.effect.each([ - ...(["permission", "question"] as const).flatMap((kind) => + ...(["question"] as const).flatMap((kind) => (["success", "failure", "timeout", "not-found"] as const).flatMap((outcome) => (["unchanged", "running", "completed"] as const).map((lifecycle) => ({ kind, @@ -658,27 +668,17 @@ describe("OpenCodeAdapterV2", () => { ); const asked = { type: `${kind}.asked`, - properties: - kind === "permission" - ? { - id: "old-request", - sessionID: "child", - permission: "bash", - patterns: ["*"], - always: [], - metadata: {}, - } - : { - id: "old-request", - sessionID: "child", - questions: [ - { - header: "Choice", - question: "Which?", - options: [{ label: "Yes", description: "Proceed" }], - }, - ], - }, + properties: { + id: "old-request", + sessionID: "child", + questions: [ + { + header: "Choice", + question: "Which?", + options: [{ label: "Yes", description: "Proceed" }], + }, + ], + }, }; yield* Effect.promise(() => nativeEvents.push(asked)); yield* TestClock.adjust("5 seconds"); @@ -764,6 +764,136 @@ describe("OpenCodeAdapterV2", () => { }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), ); + it.effect.each([{ askingSessionId: "root" }, { askingSessionId: "child" }])( + "keeps orphan and newer $askingSessionId permissions pending after reporting routing failure", + ({ askingSessionId }) => + Effect.gen(function* () { + const nativeEvents = asyncEventStream(); + const pending = new Set(["old-permission"]); + const cancelled: string[] = []; + const failureReported = yield* Deferred.make(); + let firstTurnId: ProviderTurnId | undefined; + let promptId = ""; + const harness = yield* makeOpenCodeRuntimeHarness("native-permission-overlap", "root", { + event: { subscribe: async () => ({ stream: nativeEvents.stream }) }, + session: { + create: async () => ({ data: { id: "root", time: { created: 1, updated: 1 } } }), + get: async ({ sessionID }: { sessionID: string }) => ({ + data: { id: sessionID, ...(sessionID === "child" ? { parentID: "root" } : {}) }, + }), + messages: async () => ({ data: [] }), + promptAsync: async (input: { messageID: string }) => { + promptId = input.messageID; + return { data: true }; + }, + abort: async () => ({ data: true }), + children: async () => ({ data: [] }), + }, + permission: { + reply: async () => { + // OpenCode 1.15.13 rejects all pending permissions in the asking session. + for (const requestID of pending) { + pending.delete(requestID); + cancelled.push(requestID); + await nativeEvents.push({ + type: "permission.replied", + properties: { sessionID: askingSessionId, requestID, reply: "reject" }, + }); + } + return { data: true }; + }, + }, + }); + yield* harness.runtime.events.pipe( + Stream.runForEach((event) => + event.type === "provider_session.updated" && event.providerSession.status === "error" + ? Deferred.succeed(failureReported, undefined) + : Effect.void, + ), + Effect.forkScoped, + ); + const asked = (id: string) => ({ + type: "permission.asked", + properties: { + id, + sessionID: askingSessionId, + permission: "bash", + patterns: ["*"], + always: [], + metadata: {}, + }, + }); + if (askingSessionId === "child") { + // A native background child stays busy after its parent root settles. + yield* harness.startTurn("first root turn"); + yield* Effect.promise(() => + nativeEvents.push({ + type: "message.updated", + properties: { + sessionID: "root", + info: { id: promptId, sessionID: "root", role: "user", time: { created: 1 } }, + }, + }), + ); + yield* Effect.promise(() => + nativeEvents.push({ + type: "session.created", + properties: { + info: { id: "child", parentID: "root", time: { created: 2, updated: 2 } }, + }, + }), + ); + yield* Effect.promise(() => + nativeEvents.push({ + type: "session.status", + properties: { sessionID: "child", status: { type: "busy" } }, + }), + ); + yield* Effect.promise(() => + nativeEvents.push({ + type: "session.status", + properties: { sessionID: "root", status: { type: "idle" } }, + }), + ); + const settled = yield* harness.runtime.readThreadSnapshot({ + providerThread: harness.providerThread, + }); + assert.equal(settled.providerTurns.at(-1)?.status, "completed"); + firstTurnId = settled.providerTurns.at(-1)?.id; + assert.isTrue(yield* harness.runtime.hasPendingBackgroundWork!); + } + yield* Effect.promise(() => nativeEvents.push(asked("old-permission"))); + yield* TestClock.adjust("5 seconds"); + yield* Deferred.await(failureReported); + assert.deepEqual([...pending], ["old-permission"]); + yield* harness.startTurn("new turn", harness.providerThread, 2); + pending.add("new-permission"); + yield* Effect.promise(() => nativeEvents.push(asked("new-permission"))); + const beforeReply = yield* harness.runtime.readThreadSnapshot({ + providerThread: harness.providerThread, + }); + assert.equal(beforeReply.providerTurns.at(-1)?.status, "running"); + assert.notEqual(beforeReply.providerTurns.at(-1)?.id, firstTurnId); + const newerRequest = beforeReply.runtimeRequests.find( + (request) => request.nativeRequestRef?.nativeId === "new-permission", + ); + assert.equal(newerRequest?.status, "pending"); + assert.equal(newerRequest?.providerTurnId, beforeReply.providerTurns.at(-1)?.id); + yield* TestClock.adjust("0 millis"); + const snapshot = yield* harness.runtime.readThreadSnapshot({ + providerThread: harness.providerThread, + }); + assert.isEmpty(cancelled); + assert.deepEqual([...pending], ["old-permission", "new-permission"]); + assert.equal( + snapshot.runtimeRequests.find( + (request) => request.nativeRequestRef?.nativeId === "new-permission", + )?.status, + "pending", + ); + }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), + ); + it.effect.each(["late-relation", "already-settled", "foreign-owner"] as const)( "preserves a child request with %s during routing backoff", (outcome) => @@ -854,12 +984,13 @@ describe("OpenCodeAdapterV2", () => { it.effect.each( (["permission", "question"] as const).flatMap((kind) => - (["foreign-then-unavailable", "already-answered", "sdk-already-answered"] as const).map( - (outcome) => ({ - kind, - outcome, - }), - ), + (kind === "permission" + ? (["foreign-then-unavailable"] as const) + : (["foreign-then-unavailable", "already-answered", "sdk-already-answered"] as const) + ).map((outcome) => ({ + kind, + outcome, + })), ), )("preserves an unroutable $kind request with $outcome", ({ kind, outcome }) => Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts index 47b1d852bcc0..0df89ef94d58 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts @@ -2018,7 +2018,8 @@ export function makeOpenCodeAdapterV2( * its relation to a thread. The first resolution attempt runs inline * (the replayable path); if the relation or the owning turn is not * established yet, a short forked backoff keeps trying. Exhaustion - * reports a session error and rejects requests known to belong here. */ + * reports a session error. Known local questions can be rejected by ID; + * permissions stay pending because native rejection cancels siblings. */ const routeRuntimeRequest = Effect.fnUntraced(function* ( nativeRequestId: string, sessionId: string, @@ -2054,33 +2055,29 @@ export function makeOpenCodeAdapterV2( } // A later lookup failure must not erase an observed unknown root. if (provedForeign) return; - const detail = `OpenCode ${request.type} request ${nativeRequestId} could not be routed to an active thread.`; + const detail = + `OpenCode ${request.type} request ${nativeRequestId} could not be routed to an active thread.` + + (request.type === "permission" + ? " Native permission remains pending because rejection can cancel other permissions in the same native session." + : ""); yield* Effect.logWarning(detail, { nativeSessionId: sessionId, providerSessionId: input.providerSessionId, }); // External servers broadcast other sessions' requests too. Reject - // only when the parent chain resolves to a registered thread here. - if (unresolvedState !== undefined) { - // OpenCode permission rejection cancels every pending permission - // in the asking native session, including concurrent requests. - // The SDK has no single-request denial; question.reject cancels - // only the named question. - const rejected = yield* ( - request.type === "permission" - ? sdkCall( - "permission.reply", - { requestID: nativeRequestId, reply: "reject" }, - (signal) => - client.permission.reply( - { requestID: nativeRequestId, reply: "reject" }, - { signal }, - ), - ).pipe(Effect.map((response) => unwrapData("permission.reply", response))) - : sdkCall("question.reject", { requestID: nativeRequestId }, (signal) => - client.question.reject({ requestID: nativeRequestId }, { signal }), - ).pipe(Effect.map((response) => unwrapData("question.reject", response))) - ).pipe(Effect.timeout("10 seconds"), Effect.exit); + // only named questions whose parent chain belongs to this runtime. + // Permission rejection cancels every pending permission in the + // asking session, including requests from a newer owning turn. + if (unresolvedState !== undefined && request.type === "question") { + const rejected = yield* sdkCall( + "question.reject", + { requestID: nativeRequestId }, + (signal) => client.question.reject({ requestID: nativeRequestId }, { signal }), + ).pipe( + Effect.map((response) => unwrapData("question.reject", response)), + Effect.timeout("10 seconds"), + Effect.exit, + ); if (Exit.isFailure(rejected)) { if (isOpenCodeNotFound(Cause.squash(rejected.cause))) { rememberSettledRequest(nativeRequestId); From 929dcc4553a70727c02086f7131cd75cd3e7fa3e Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Fri, 9 Oct 2026 13:44:46 +0200 Subject: [PATCH 7/7] fix(opencode): refresh request ownership after root registration --- .../src/server/adapter.test.ts | 98 +++++++++++++++++++ .../provider-opencode/src/server/adapter.ts | 3 +- 2 files changed, 100 insertions(+), 1 deletion(-) diff --git a/packages/provider-opencode/src/server/adapter.test.ts b/packages/provider-opencode/src/server/adapter.test.ts index 5aa191d68e88..0b844c8d5130 100644 --- a/packages/provider-opencode/src/server/adapter.test.ts +++ b/packages/provider-opencode/src/server/adapter.test.ts @@ -891,6 +891,104 @@ describe("OpenCodeAdapterV2", () => { }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), ); + it.effect.each(["permission", "question"] as const)( + "handles an initially unknown root registered during %s routing backoff", + (kind) => + Effect.gen(function* () { + const nativeEvents = asyncEventStream(); + const rejections: unknown[] = []; + const lookups: string[] = []; + const reject = async (input: unknown) => { + rejections.push(input); + return { data: true }; + }; + const harness = yield* makeOpenCodeRuntimeHarness(`registered-root-${kind}`, "root", { + event: { subscribe: async () => ({ stream: nativeEvents.stream }) }, + session: { + create: async () => ({ data: { id: "root", time: { created: 1, updated: 1 } } }), + get: async ({ sessionID }: { sessionID: string }) => { + lookups.push(sessionID); + return { + data: { + id: sessionID, + ...(sessionID === "child" ? { parentID: "late-root" } : {}), + time: { created: 1, updated: 1 }, + }, + }; + }, + messages: async () => ({ data: [] }), + abort: async () => ({ data: true }), + children: async () => ({ data: [] }), + }, + permission: { reply: reject }, + question: { reject }, + }); + const received: ProviderAdapter.ProviderAdapterV2Event[] = []; + yield* harness.runtime.events.pipe( + Stream.runForEach((event) => Effect.sync(() => received.push(event))), + Effect.forkScoped, + ); + yield* Effect.promise(() => + nativeEvents.push({ + type: `${kind}.asked`, + properties: + kind === "permission" + ? { + id: "late-owner-request", + sessionID: "child", + permission: "bash", + patterns: ["*"], + always: [], + metadata: {}, + } + : { + id: "late-owner-request", + sessionID: "child", + questions: [ + { + header: "Choice", + question: "Which?", + options: [{ label: "Yes", description: "Proceed" }], + }, + ], + }, + }), + ); + assert.deepEqual(lookups, ["child", "late-root"]); + assert.isEmpty(rejections); + // Attach the previously unknown native root without starting a turn. + const resumed = yield* harness.runtime.resumeThread({ + providerThread: { + ...harness.providerThread, + nativeThreadRef: { + driver: OPENCODE_PROVIDER, + nativeId: "late-root", + strength: "weak", + }, + }, + }); + yield* TestClock.adjust("5 seconds"); + yield* TestClock.adjust("0 millis"); + assert.deepEqual( + rejections, + kind === "question" ? [{ requestID: "late-owner-request" }] : [], + ); + const failure = received.find( + (event) => + event.type === "provider_session.updated" && event.providerSession.status === "error", + ); + assert.isDefined(failure); + if (failure?.type !== "provider_session.updated") return; + assert.include(failure.providerSession.lastError, `${kind} request late-owner-request`); + if (kind === "permission") { + assert.include(failure.providerSession.lastError, "Native permission remains pending"); + } + const snapshot = yield* harness.runtime.readThreadSnapshot({ providerThread: resumed }); + assert.isEmpty(snapshot.providerTurns); + assert.isEmpty(snapshot.runtimeRequests); + }).pipe(Effect.provide(IdAllocator.layer), Effect.scoped), + ); + it.effect.each(["late-relation", "already-settled", "foreign-owner"] as const)( "preserves a child request with %s during routing backoff", (outcome) => diff --git a/packages/provider-opencode/src/server/adapter.ts b/packages/provider-opencode/src/server/adapter.ts index bcf2fc5d1072..7ae559694b0d 100644 --- a/packages/provider-opencode/src/server/adapter.ts +++ b/packages/provider-opencode/src/server/adapter.ts @@ -2039,7 +2039,8 @@ export const makeOpenCodeAdapterV2 = Effect.fn("makeOpenCodeAdapterV2")(function } const state = yield* resolveSessionOwner(sessionId); unresolvedState = state; - if (state === null) provedForeign = true; + // Failed lookups preserve an observed unknown root; local ownership replaces it. + if (state !== undefined) provedForeign = state === null; const owner = state == null ? undefined : topLevelRequestOwner(state); if (owner === undefined) return false; yield* emitRuntimeRequest(owner, nativeRequestId, request);