From 8decf64fa017883b9388f1ede82446d1252dd91d Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Sun, 13 Sep 2026 12:04:52 +1000 Subject: [PATCH] fix(server): restart the live session when a model change follows dead records ProviderSwitchService.plan() picked the newest providerSessions record for the current instance regardless of status, while the release filter skipped stopped/error records. A newer dead record could hide an older live session: restart_and_resume targeted the dead record, released nothing, and left the real session running with the old model. Select the newest eligible live session via a shared isLiveProviderSession predicate so stopped/error history cannot mask the session being replaced. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../ProviderSwitchService.test.ts | 202 +++++++++++++++++- .../orchestration-v2/ProviderSwitchService.ts | 25 ++- .../SelectionRestart.integration.test.ts | 178 +++++++++++++++ 3 files changed, 397 insertions(+), 8 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderSwitchService.test.ts b/apps/server/src/orchestration-v2/ProviderSwitchService.test.ts index 35bc1b803627..1e130ce33c01 100644 --- a/apps/server/src/orchestration-v2/ProviderSwitchService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSwitchService.test.ts @@ -3,6 +3,7 @@ import { ProviderDriverKind, ProviderInstanceId, ProviderSessionId, + ProviderThreadId, ThreadId, type OrchestrationV2ThreadProjection, } from "@t3tools/contracts"; @@ -13,6 +14,7 @@ import * as Layer from "effect/Layer"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts"; +import { acpSelectionTransition } from "./ProviderSelectionTransition.ts"; import * as ProviderSwitch from "./ProviderSwitchService.ts"; const driver = ProviderDriverKind.make("codex"); @@ -50,12 +52,29 @@ function projection(): OrchestrationV2ThreadProjection { } as unknown as OrchestrationV2ThreadProjection; } -function testLayer(metadata: Readonly>) { +function deadSessionRecord( + id: string, + status: "stopped" | "error", + updatedAt: DateTime.Utc = DateTime.add(now, { seconds: 1 }), +) { + return { + ...projection().providerSessions[0]!, + id: ProviderSessionId.make(id), + status, + updatedAt, + }; +} + +function testLayer( + metadata: Readonly>, + planSelectionTransition: ProviderAdapterV2Shape["planSelectionTransition"] = () => + Effect.succeed({ type: "restart_session" }), +) { const adapter = (instanceId: ProviderInstanceId): ProviderAdapterV2Shape => ({ instanceId, driver, getCapabilities: () => Effect.succeed(capabilitiesWithoutModelSwitch), - planSelectionTransition: () => Effect.succeed({ type: "restart_session" }), + planSelectionTransition, openSession: () => Effect.die("ProviderSwitchService tests do not open sessions."), }); const registry = Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({ @@ -101,6 +120,185 @@ it.effect( ), ); +for (const deadStatus of ["stopped", "error"] as const) { + it.effect( + `restarts and releases the live session when a newer ${deadStatus} session exists`, + () => + Effect.gen(function* () { + const service = yield* ProviderSwitch.ProviderSwitchServiceV2; + const thread = projection(); + const result = yield* service.plan({ + projection: { + ...thread, + providerSessions: [ + ...thread.providerSessions, + deadSessionRecord("dead_session", deadStatus), + ], + }, + targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" }, + }); + assert.equal(result.transition.type, "restart_and_resume"); + assert.deepEqual(result.releaseProviderSessionIds, [currentSessionId]); + }).pipe( + Effect.provide( + testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }), + ), + ), + ); +} + +it.effect("releases the newest live session, not the newest record overall", () => + Effect.gen(function* () { + const service = yield* ProviderSwitch.ProviderSwitchServiceV2; + const thread = projection(); + const newerLiveSessionId = ProviderSessionId.make("session_newer_live"); + const result = yield* service.plan({ + projection: { + ...thread, + providerSessions: [ + ...thread.providerSessions, + { + ...thread.providerSessions[0]!, + id: newerLiveSessionId, + updatedAt: DateTime.add(now, { seconds: 1 }), + }, + deadSessionRecord("dead_session", "stopped", DateTime.add(now, { seconds: 2 })), + ], + }, + targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" }, + }); + assert.equal(result.transition.type, "restart_and_resume"); + assert.deepEqual(result.releaseProviderSessionIds, [newerLiveSessionId]); + }).pipe( + Effect.provide( + testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }), + ), + ), +); + +it.effect("creates a fresh session with handoff when every recorded session is dead", () => + Effect.gen(function* () { + const service = yield* ProviderSwitch.ProviderSwitchServiceV2; + const thread = projection(); + const result = yield* service.plan({ + projection: { + ...thread, + providerSessions: [ + deadSessionRecord("dead_session_older", "stopped"), + deadSessionRecord("dead_session_newer", "error", DateTime.add(now, { seconds: 2 })), + ], + }, + targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" }, + }); + assert.equal(result.transition.type, "create_with_handoff"); + assert.deepEqual(result.releaseProviderSessionIds, []); + }).pipe( + Effect.provide( + testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }), + ), + ), +); + +function deadNativeThreadProjection( + status: "stopped" | "error", + capabilities = capabilitiesWithoutModelSwitch, +): OrchestrationV2ThreadProjection { + const thread = projection(); + return { + ...thread, + thread: { + ...thread.thread, + activeProviderThreadId: ProviderThreadId.make("provider-thread:native"), + }, + providerSessions: [{ ...deadSessionRecord("dead_session", status), capabilities }], + providerThreads: [ + { + id: ProviderThreadId.make("provider-thread:native"), + driver, + providerInstanceId: currentInstanceId, + providerSessionId: ProviderSessionId.make("dead_session"), + appThreadId: thread.thread.id, + ownerNodeId: null, + nativeThreadRef: { + driver, + nativeId: "native-thread:abc", + strength: "strong", + }, + nativeConversationHeadRef: null, + status: "idle", + firstRunOrdinal: null, + lastRunOrdinal: null, + handoffIds: [], + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + ], + } as OrchestrationV2ThreadProjection; +} + +it.effect("falls back to the native provider thread when every recorded session is dead", () => + Effect.gen(function* () { + const service = yield* ProviderSwitch.ProviderSwitchServiceV2; + const result = yield* service.plan({ + projection: deadNativeThreadProjection("stopped"), + targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" }, + }); + assert.equal(result.transition.type, "restart_and_resume"); + assert.deepEqual(result.releaseProviderSessionIds, []); + }).pipe( + Effect.provide( + testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }), + ), + ), +); + +for (const deadStatus of ["stopped", "error"] as const) { + it.effect( + `applies a model change on next turn when a ${deadStatus} session negotiated model switching`, + () => + Effect.gen(function* () { + const service = yield* ProviderSwitch.ProviderSwitchServiceV2; + const result = yield* service.plan({ + projection: deadNativeThreadProjection(deadStatus, CodexProviderCapabilitiesV2), + targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" }, + }); + // Static capabilities report no in-session switch, but the dead + // record's negotiated capabilities describe the provider: without + // them the ACP classification rejects the selection instead of + // reopening with the requested model on the next run. + assert.equal(result.transition.type, "switch_model_in_session"); + assert.deepEqual(result.releaseProviderSessionIds, []); + }).pipe( + Effect.provide( + testLayer( + { [currentInstanceId]: { continuationKey: "codex:account:primary" } }, + (input) => Effect.succeed(acpSelectionTransition(input)), + ), + ), + ), + ); +} + +it.effect("rejects a model change the dead record never negotiated support for", () => + Effect.gen(function* () { + const service = yield* ProviderSwitch.ProviderSwitchServiceV2; + const result = yield* service + .plan({ + projection: deadNativeThreadProjection("stopped"), + targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" }, + }) + .pipe(Effect.flip); + assert.instanceOf(result, ProviderSwitch.ProviderSwitchPlanError); + }).pipe( + Effect.provide( + testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }, (input) => + Effect.succeed(acpSelectionTransition(input)), + ), + ), + ), +); + it.effect("distinguishes compatible and incompatible instances of the same driver", () => Effect.gen(function* () { const service = yield* ProviderSwitch.ProviderSwitchServiceV2; diff --git a/apps/server/src/orchestration-v2/ProviderSwitchService.ts b/apps/server/src/orchestration-v2/ProviderSwitchService.ts index 7c00b155b3eb..d96fb8fa3680 100644 --- a/apps/server/src/orchestration-v2/ProviderSwitchService.ts +++ b/apps/server/src/orchestration-v2/ProviderSwitchService.ts @@ -51,6 +51,12 @@ export class ProviderSwitchServiceV2 extends Context.Service< ProviderSwitchServiceV2Shape >()("t3/orchestration-v2/ProviderSwitchService/ProviderSwitchServiceV2") {} +// Stopped and errored records stay in session history but can no longer be +// restarted or released; only live sessions participate in a transition. +const isLiveProviderSession = ( + session: OrchestrationV2ThreadProjection["providerSessions"][number], +) => session.status !== "stopped" && session.status !== "error"; + export const layer: Layer.Layer< ProviderSwitchServiceV2, never, @@ -83,12 +89,18 @@ export const layer: Layer.Layer< const currentInstance = yield* Effect.option(getMetadata(current.instanceId)); const targetInstance = yield* Effect.option(getMetadata(targetModelSelection.instanceId)); const targetAdapter = yield* Effect.option(adapters.get(targetModelSelection.instanceId)); - const currentSession = projection.providerSessions + const currentSessions = projection.providerSessions .filter((session) => session.providerInstanceId === current.instanceId) .toSorted( (left, right) => DateTime.toEpochMillis(right.updatedAt) - DateTime.toEpochMillis(left.updatedAt), - )[0]; + ); + const currentSession = currentSessions.find(isLiveProviderSession); + // Negotiated capabilities describe the provider, not the dead + // process; the newest record still reports what the instance + // supports after its session stops. + const negotiatedCapabilities = + currentSession?.capabilities ?? currentSessions[0]?.capabilities; // Detaching a process removes its session binding, not its native history. const currentProviderThread = projection.providerThreads.find( (thread) => @@ -105,8 +117,7 @@ export const layer: Layer.Layer< ? yield* targetAdapter.value.planSelectionTransition({ current, target: targetModelSelection, - sessionCapabilities: - currentSession?.capabilities ?? currentInstance.value.capabilities, + sessionCapabilities: negotiatedCapabilities ?? currentInstance.value.capabilities, }) : undefined; const transition = @@ -134,7 +145,7 @@ export const layer: Layer.Layer< projection.thread.worktreePath ?? "", capabilities: - currentSession?.capabilities ?? currentInstance.value.capabilities, + negotiatedCapabilities ?? currentInstance.value.capabilities, }, target: { driver: targetInstance.value.driver, @@ -174,8 +185,10 @@ export const layer: Layer.Layer< )[0]; const releaseProviderSessionIds = projection.providerSessions .filter((session) => { - if (session.status === "stopped" || session.status === "error") return false; + if (!isLiveProviderSession(session)) return false; if (transition.type === "restart_and_resume") { + // Other live records may serve pooled or delegated bindings; + // only the session being replaced is released. return session.id === currentSession?.id; } if (transition.type === "create_with_handoff") { diff --git a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts index 2b4211daebab..d62e13ba64d9 100644 --- a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts +++ b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts @@ -1,6 +1,7 @@ import { assert, it } from "@effect/vitest"; import { CommandId, + EventId, MessageId, type ModelSelection, type OrchestrationV2ProviderCapabilities, @@ -9,6 +10,7 @@ import { ProjectId, ProviderDriverKind, ProviderInstanceId, + ProviderSessionId, ProviderThreadId, ProviderTurnId, type RunId, @@ -24,6 +26,7 @@ import * as Stream from "effect/Stream"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import { OrchestrationEffectWorkerV2 } from "./EffectWorker.ts"; +import { EventSinkV2 } from "./EventSink.ts"; import { OrchestratorV2 } from "./Orchestrator.ts"; import { ProviderAdapterOpenSessionError, @@ -49,6 +52,10 @@ const replacementSelection = { instanceId: providerInstanceId, model: "restart-model-b", } satisfies ModelSelection; +const seedSelection = { + instanceId: providerInstanceId, + model: "seed-model", +} satisfies ModelSelection; const handoffDriver = ProviderDriverKind.make("claudeAgent"); const handoffProviderInstanceId = ProviderInstanceId.make("claude-handoff-test"); const handoffSelection = { @@ -534,6 +541,177 @@ it.live("restarts selection as a new attempt and retries after old-session clean ), ); +for (const deadStatus of ["stopped", "error"] as const) { + it.live( + `restarts the live session on a model change when a newer ${deadStatus} session record exists`, + () => + Effect.scoped( + Effect.gen(function* () { + const name = `selection-restart-dead-${deadStatus}`; + const cwd = yield* checkpointWorkspace(name); + const threadId = ThreadId.make(`thread:${name}`); + const state = yield* Ref.make({ + activeTurn: null, + opened: [], + started: [], + closedSessionCount: 0, + // The dead record is seeded directly, so the adapter's one-shot + // simulated replacement-open failure is skipped. + failedReplacementOpen: true, + }); + const registry = makeSingleProviderAdapterRegistryLayer( + makeRestartAdapter(state, exclusiveCapabilities), + ); + + const result = yield* Effect.gen(function* () { + const orchestrator = yield* OrchestratorV2; + const worker = yield* OrchestrationEffectWorkerV2; + const eventSink = yield* EventSinkV2; + const dispatch = (step: string, modelSelection: ModelSelection) => + Effect.gen(function* () { + const terminal = yield* orchestrator.streamDomainEvents.pipe( + Stream.filter( + (event) => event.type === "run.updated" && event.payload.status === "completed", + ), + Stream.take(1), + Stream.runDrain, + Effect.forkScoped, + ); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${name}:${step}`), + threadId, + messageId: MessageId.make(`${name}:${step}`), + text: step, + attachments: [], + modelSelection, + dispatchMode: { type: "start_immediately" }, + }); + yield* worker.drain(); + yield* Fiber.join(terminal); + yield* worker.drain(); + return yield* orchestrator.getThreadProjection(threadId); + }); + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${name}:create`), + threadId, + projectId: ProjectId.make(`project:${name}`), + title: name, + modelSelection: seedSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: cwd, + }); + const first = yield* dispatch("first", seedSelection); + const liveSession = first.providerSessions.find( + (session) => session.status !== "stopped" && session.status !== "error", + ); + assert.isDefined(liveSession); + + // Dead session records stay bound in the projection until + // detachment. A stale stopped/error record written after the + // live session attached must not hide it. + const deadAt = yield* DateTime.now; + const deadSession: OrchestrationV2ProviderSession = { + id: ProviderSessionId.make(`session:${name}:dead`), + driver, + providerInstanceId, + status: "ready", + cwd, + model: seedSelection.model, + capabilities: exclusiveCapabilities, + createdAt: deadAt, + updatedAt: deadAt, + lastError: null, + }; + yield* eventSink.write({ + events: [ + { + id: EventId.make(`event:${name}:dead-attached`), + type: "provider-session.attached", + threadId, + driver, + providerInstanceId, + occurredAt: deadAt, + payload: deadSession, + }, + { + id: EventId.make(`event:${name}:dead-updated`), + type: "provider-session.updated", + threadId, + driver, + providerInstanceId, + occurredAt: deadAt, + payload: { + ...deadSession, + status: deadStatus, + updatedAt: deadAt, + lastError: deadStatus === "error" ? "Simulated session failure." : null, + }, + }, + ], + }); + + const switchCommandId = CommandId.make(`${name}:switch`); + yield* orchestrator.dispatch({ + type: "thread.model-selection.set", + commandId: switchCommandId, + threadId, + modelSelection: replacementSelection, + }); + yield* worker.drain(); + const storedSwitchEvents = yield* eventSink + .readByCommandId({ commandId: switchCommandId }) + .pipe(Stream.runCollect); + const detachedSessionIds = [...storedSwitchEvents].flatMap((stored) => + stored.event.type === "provider-session.detached" + ? [stored.event.payload.providerSessionId] + : [], + ); + + const second = yield* dispatch("second", replacementSelection); + return { + projection: second, + captured: yield* Ref.get(state), + liveSessionId: liveSession.id, + detachedSessionIds, + }; + }).pipe(Effect.provide(makeOrchestratorV2ReplayLayerWithRegistry({ name }, registry))); + + const { projection, captured } = result; + assert.lengthOf(projection.runs, 2); + assert.equal(projection.runs[1]?.modelSelection.model, replacementSelection.model); + // The exact released session is the older live one, never the newer + // dead record. + assert.deepEqual(result.detachedSessionIds, [result.liveSessionId]); + assert.equal(captured.closedSessionCount, 1); + assert.deepEqual( + captured.opened.map((open) => open.model), + [seedSelection.model, replacementSelection.model], + ); + assert.deepEqual( + captured.started.map((turn) => turn.model), + [seedSelection.model, replacementSelection.model], + ); + const servingSession = projection.providerSessions.find( + (session) => + session.id === + projection.providerThreads.find( + (providerThread) => providerThread.id === projection.thread.activeProviderThreadId, + )?.providerSessionId, + ); + assert.equal(servingSession?.model, replacementSelection.model); + }), + ), + ); +} + it.live("detaches the old provider session after an active provider handoff", () => Effect.scoped( Effect.gen(function* () {