diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index 9bc701af0837..7f85500bd2da 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -29,6 +29,7 @@ import * as Effect from "effect/Effect"; import * as Deferred from "effect/Deferred"; import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; +import * as Logger from "effect/Logger"; import * as ManagedRuntime from "effect/ManagedRuntime"; import * as Option from "effect/Option"; import * as PubSub from "effect/PubSub"; @@ -175,6 +176,8 @@ describe("ProviderCommandReactor", () => { readonly requiresNewThreadForModelChange?: boolean; readonly unreadableHistory?: boolean; readonly titleRegenerationCompletionDispatchFailures?: number; + readonly interruptRecoveryDispatchFailure?: "session" | "activity"; + readonly logger?: Logger.Logger; readonly titleRegenerationBeforeStart?: "one" | "two"; readonly serverActivation?: Effect.Effect; readonly beforeReadySessionDispatch?: () => Effect.Effect; @@ -427,6 +430,16 @@ describe("ProviderCommandReactor", () => { readThreadEvents: engine.readThreadEvents, getThreadReplayStats: engine.getThreadReplayStats, dispatch: (command) => { + if ( + (input?.interruptRecoveryDispatchFailure === "session" && + command.type === "thread.session.set" && + command.session.status === "stopped") || + (input?.interruptRecoveryDispatchFailure === "activity" && + command.type === "thread.activity.append" && + command.activity.kind === "provider.turn.interrupt.failed") + ) { + return Effect.die(new Error("Injected interrupt recovery dispatch failure")); + } if (command.type === "thread.title.regeneration.complete") { titleRegenerationCompletionDispatchAttempts += 1; if ( @@ -495,7 +508,11 @@ describe("ProviderCommandReactor", () => { Layer.provideMerge(ServerConfig.layerTest(process.cwd(), baseDir)), Layer.provideMerge(NodeServices.layer), ); - runtime = ManagedRuntime.make(layer); + runtime = ManagedRuntime.make( + input?.logger + ? layer.pipe(Layer.provide(Logger.layer([input.logger], { mergeWithExisting: false }))) + : layer, + ); const engine = await runtime.runPromise(Effect.service(OrchestrationEngineService)); const snapshotQuery = await runtime.runPromise(Effect.service(ProjectionSnapshotQuery)); @@ -3623,6 +3640,429 @@ describe("ProviderCommandReactor", () => { }); }); + // Live clock: the bound below must fire while the reactor is genuinely stuck. + effectIt.live("keeps serving other threads while a provider interrupt hangs", () => + Effect.gen(function* () { + const interruptStarted = yield* Deferred.make(); + const releaseInterrupt = yield* Deferred.make(); + const otherThreadStarted = yield* Deferred.make(); + const otherThreadId = ThreadId.make("thread-2"); + const harness = yield* Effect.promise(() => + createHarness({ + // An OpenCode cancel that never settles: the provider's problem, + // but it must not become the command worker's problem. + interruptTurnEffect: () => + Deferred.succeed(interruptStarted, undefined).pipe( + Effect.andThen(Deferred.await(releaseInterrupt)), + ), + startSessionEffect: (session) => + session.threadId === otherThreadId + ? Deferred.succeed(otherThreadStarted, undefined).pipe(Effect.as(session)) + : Effect.succeed(session), + }), + ); + const now = "2026-01-01T00:00:00.000Z"; + + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-hanging-interrupt"), + threadId: ThreadId.make("thread-1"), + session: { + threadId: ThreadId.make("thread-1"), + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-1"), + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); + yield* harness.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: CommandId.make("cmd-turn-interrupt-hanging"), + threadId: ThreadId.make("thread-1"), + turnId: asTurnId("turn-1"), + createdAt: now, + }); + yield* Deferred.await(interruptStarted); + + yield* harness.engine.dispatch({ + type: "thread.create", + commandId: CommandId.make("cmd-thread-create-other-thread"), + threadId: otherThreadId, + projectId: asProjectId("project-1"), + title: "Other thread", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + branch: null, + worktreePath: null, + createdAt: now, + }); + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-other-thread"), + threadId: otherThreadId, + message: { + messageId: MessageId.make("message-other-thread"), + role: "user", + text: "hello from another thread", + attachments: [], + }, + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: now, + }); + + // The other thread's start must not queue behind the hung interrupt: + // its session starts and its turn is sent while the interrupt is still + // pending. + yield* Deferred.await(otherThreadStarted).pipe(Effect.timeout("5 seconds")); + yield* Effect.promise(() => waitFor(() => harness.sendTurn.mock.calls.length === 1)); + const thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === otherThreadId, + ); + expect(thread?.session).toMatchObject({ providerName: "codex" }); + expect(thread?.session?.status).not.toBe("stopped"); + expect(yield* Deferred.isDone(releaseInterrupt)).toBe(false); + + yield* Deferred.succeed(releaseInterrupt, undefined); + yield* Effect.promise(() => harness.drain()); + }), + ); + + effectIt.effect( + "does not stop a newer session on the same thread when an earlier interrupt fails late", + () => + Effect.gen(function* () { + const interruptStarted = yield* Deferred.make(); + const failInterrupt = yield* Deferred.make(); + const harness = yield* Effect.promise(() => + createHarness({ + interruptTurnEffect: () => + Deferred.succeed(interruptStarted, undefined).pipe( + Effect.andThen(Deferred.await(failInterrupt)), + Effect.andThen( + Effect.fail( + new ProviderAdapterRequestError({ + provider: "codex", + method: "thread.interrupt", + detail: "provider session disappeared", + }), + ), + ), + ), + }), + ); + const threadId = ThreadId.make("thread-1"); + const now = "2026-01-01T00:00:00.000Z"; + const later = "2026-01-01T00:00:01.000Z"; + + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-before-late-failure"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-1"), + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); + yield* harness.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: CommandId.make("cmd-turn-interrupt-late-failure"), + threadId, + createdAt: now, + }); + yield* Deferred.await(interruptStarted); + + // The thread moves on to a fresh session while the interrupt is pending. + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-newer"), + threadId, + session: { + threadId, + status: "starting", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: later, + }, + createdAt: later, + }); + // drain would wait for the pending interrupt; wait for the projection instead. + yield* Effect.promise(() => + waitFor(async () => { + const current = (await harness.readModel()).threads.find( + (entry) => entry.id === threadId, + ); + return current?.session?.updatedAt === later; + }), + ); + + yield* Deferred.succeed(failInterrupt, undefined); + yield* Effect.promise(() => harness.drain()); + + const thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + expect(thread?.session).toMatchObject({ status: "starting", activeTurnId: null }); + expect(harness.stopSession).not.toHaveBeenCalled(); + }), + ); + + effectIt.effect( + "preserves a replacement session created while interrupt recovery is stopping the old session", + () => + Effect.gen(function* () { + const interruptStarted = yield* Deferred.make(); + const failInterrupt = yield* Deferred.make(); + const stopStarted = yield* Deferred.make(); + const releaseStop = yield* Deferred.make(); + const harness = yield* Effect.promise(() => + createHarness({ + stopSessionEffect: () => + Deferred.succeed(stopStarted, undefined).pipe( + Effect.andThen(Deferred.await(releaseStop)), + ), + interruptTurnEffect: () => + Deferred.succeed(interruptStarted, undefined).pipe( + Effect.andThen(Deferred.await(failInterrupt)), + Effect.andThen( + Effect.fail( + new ProviderAdapterRequestError({ + provider: "codex", + method: "thread.interrupt", + detail: "provider session disappeared", + }), + ), + ), + ), + }), + ); + const threadId = ThreadId.make("thread-1"); + const now = "2026-01-01T00:00:00.000Z"; + const later = "2026-01-01T00:00:01.000Z"; + + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-before-late-failure"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-1"), + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); + yield* harness.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: CommandId.make("cmd-turn-interrupt-late-failure"), + threadId, + createdAt: now, + }); + yield* Deferred.await(interruptStarted); + + yield* Deferred.succeed(failInterrupt, undefined); + yield* Deferred.await(stopStarted); + // A replacement starts while cleanup of the old session is awaiting the provider. + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-newer"), + threadId, + session: { + threadId, + status: "starting", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: later, + }, + createdAt: later, + }); + yield* Deferred.succeed(releaseStop, undefined); + yield* Effect.promise(() => harness.drain()); + + const thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + expect(thread?.session).toMatchObject({ status: "starting", activeTurnId: null }); + expect(harness.stopSession).toHaveBeenCalledTimes(1); + }), + ); + + effectIt.effect.each(["session", "activity"] as const)( + "logs a failure to persist the %s during interrupt recovery", + (interruptRecoveryDispatchFailure) => + Effect.gen(function* () { + const messages: Array = []; + const logger = Logger.make(({ message }) => { + messages.push(message); + }); + const harness = yield* Effect.promise(() => + createHarness({ + logger, + interruptRecoveryDispatchFailure, + interruptTurnEffect: () => + Effect.fail( + new ProviderAdapterRequestError({ + provider: "codex", + method: "thread.interrupt", + detail: "provider session disappeared", + }), + ), + }), + ); + const threadId = ThreadId.make("thread-1"); + const now = "2026-01-01T00:00:00.000Z"; + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-before-recovery-dispatch-failure"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-1"), + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); + yield* harness.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: CommandId.make("cmd-interrupt-with-recovery-dispatch-failure"), + threadId, + createdAt: now, + }); + yield* Effect.promise(() => harness.drain()); + + expect(messages).toContainEqual([ + "provider command reactor failed to recover interrupt", + { + threadId, + cause: expect.stringContaining("Injected interrupt recovery dispatch failure"), + }, + ]); + }), + ); + + effectIt.effect( + "still stops the same session when a lifecycle update rewrites it before the interrupt fails", + () => + Effect.gen(function* () { + const interruptStarted = yield* Deferred.make(); + const failInterrupt = yield* Deferred.make(); + const harness = yield* Effect.promise(() => + createHarness({ + interruptTurnEffect: () => + Deferred.succeed(interruptStarted, undefined).pipe( + Effect.andThen(Deferred.await(failInterrupt)), + Effect.andThen( + Effect.fail( + new ProviderAdapterRequestError({ + provider: "codex", + method: "thread.interrupt", + detail: "provider session disappeared", + }), + ), + ), + ), + }), + ); + const threadId = ThreadId.make("thread-1"); + const now = "2026-01-01T00:00:00.000Z"; + const later = "2026-01-01T00:00:01.000Z"; + + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-before-rewrite"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-1"), + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); + yield* harness.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: CommandId.make("cmd-turn-interrupt-rewrite"), + threadId, + turnId: asTurnId("turn-1"), + createdAt: now, + }); + yield* Deferred.await(interruptStarted); + + // Runtime ingestion rewrites the row for a lifecycle event on the same + // session and turn; only `updatedAt` moves. + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-rewrite"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-1"), + lastError: null, + updatedAt: later, + }, + createdAt: later, + }); + yield* Effect.promise(() => + waitFor(async () => { + const current = (await harness.readModel()).threads.find( + (entry) => entry.id === threadId, + ); + return current?.session?.updatedAt === later; + }), + ); + + yield* Deferred.succeed(failInterrupt, undefined); + yield* Effect.promise(() => harness.drain()); + + const thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + expect(thread?.session).toMatchObject({ + status: "stopped", + activeTurnId: null, + lastError: "provider session disappeared", + }); + expect( + thread?.activities.find((activity) => activity.kind === "provider.turn.interrupt.failed"), + ).toMatchObject({ payload: { detail: "provider session disappeared" } }); + expect(harness.stopSession).toHaveBeenCalledWith({ threadId }); + }), + ); + effectIt.effect( "stops a running session and records the failure when provider interrupt fails", () => diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 80cf2a1f4761..3e726bb22628 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -23,6 +23,7 @@ import * as DateTime from "effect/DateTime"; import * as Deferred from "effect/Deferred"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; +import * as FiberSet from "effect/FiberSet"; import * as Equal from "effect/Equal"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; @@ -112,6 +113,23 @@ function mapProviderSessionStatusToOrchestrationStatus( } } +/** + * Whether a thread's session row was replaced after `before` was read. The row + * has no identity of its own, and lifecycle events rewrite it in place with a + * new `updatedAt`, so a replacement shows up as a different provider instance, + * a different active turn, or a restart: a session only returns to `starting` + * through a new turn start. + */ +function sessionWasReplaced(before: OrchestrationSession, after: OrchestrationSession): boolean { + return ( + after.providerInstanceId !== before.providerInstanceId || + (before.activeTurnId !== null && + after.activeTurnId !== null && + after.activeTurnId !== before.activeTurnId) || + (before.status !== "starting" && after.status === "starting") + ); +} + const turnStartKeyForEvent = (event: ProviderIntentEvent): string => event.commandId !== null ? `command:${event.commandId}` : `event:${event.eventId}`; @@ -1191,6 +1209,10 @@ const make = Effect.gen(function* () { }), ), ); + // Provider interrupts run off the command worker: a provider's cancel can + // hang, and the worker is shared by every thread and provider. The set lets + // `drain` still wait for them. + const interruptFibers = yield* FiberSet.make(); const threadTitleRegenerationWorker = yield* makeDrainableWorker( processThreadTitleRegenerationSafely, ); @@ -1538,6 +1560,10 @@ const make = Effect.gen(function* () { !latestSession || latestSession.status === "stopped" || latestSession.status === "ready" || + // The interrupt runs off the command worker, so the thread may have + // moved on to a newer session or turn by the time it fails. Only the + // session this interrupt was asked about may be forced to stop. + sessionWasReplaced(session, latestSession) || (event.payload.turnId !== undefined && latestSession.activeTurnId !== null && latestSession.activeTurnId !== event.payload.turnId) @@ -1566,6 +1592,7 @@ const make = Effect.gen(function* () { !stoppedSession || stoppedSession.status === "stopped" || stoppedSession.status === "ready" || + sessionWasReplaced(session, stoppedSession) || (event.payload.turnId !== undefined && stoppedSession.activeTurnId !== null && stoppedSession.activeTurnId !== event.payload.turnId) @@ -1596,9 +1623,21 @@ const make = Effect.gen(function* () { }; // Orchestration turn ids are not provider turn ids, so interrupt by session. - yield* providerService - .interruptTurn({ threadId: event.payload.threadId }) - .pipe(Effect.catchCause(recoverInterruptFailure)); + yield* FiberSet.run( + interruptFibers, + providerService.interruptTurn({ threadId: event.payload.threadId }).pipe( + Effect.catchCause(recoverInterruptFailure), + Effect.catchCause((cause) => { + if (Cause.hasInterruptsOnly(cause)) { + return Effect.interrupt; + } + return Effect.logWarning("provider command reactor failed to recover interrupt", { + threadId: event.payload.threadId, + cause: Cause.pretty(cause), + }); + }), + ), + ); }); const processApprovalResponseRequested = Effect.fn("processApprovalResponseRequested")(function* ( @@ -1929,6 +1968,7 @@ const make = Effect.gen(function* () { start, drain: Effect.gen(function* () { yield* worker.drain; + yield* FiberSet.awaitEmpty(interruptFibers); yield* threadTitleRegenerationWorker.drain; }), } satisfies ProviderCommandReactorShape;