Skip to content

Commit 57cd8a1

Browse files
saphidgithub-actions[bot]devin-ai-integration[bot]
authored andcommitted
fix(server): restart the live session on model changes after dead records (#11505)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
1 parent 1a79960 commit 57cd8a1

3 files changed

Lines changed: 397 additions & 8 deletions

File tree

‎apps/server/src/orchestration-v2/ProviderSwitchService.test.ts‎

Lines changed: 200 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ import {
33
ProviderDriverKind,
44
ProviderInstanceId,
55
ProviderSessionId,
6+
ProviderThreadId,
67
ThreadId,
78
type OrchestrationV2ThreadProjection,
89
} from "@t3tools/contracts";
@@ -13,6 +14,7 @@ import * as Layer from "effect/Layer";
1314
import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts";
1415
import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts";
1516
import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts";
17+
import { acpSelectionTransition } from "./ProviderSelectionTransition.ts";
1618
import * as ProviderSwitch from "./ProviderSwitchService.ts";
1719

1820
const driver = ProviderDriverKind.make("codex");
@@ -50,12 +52,29 @@ function projection(): OrchestrationV2ThreadProjection {
5052
} as unknown as OrchestrationV2ThreadProjection;
5153
}
5254

53-
function testLayer(metadata: Readonly<Record<string, { continuationKey: string }>>) {
55+
function deadSessionRecord(
56+
id: string,
57+
status: "stopped" | "error",
58+
updatedAt: DateTime.Utc = DateTime.add(now, { seconds: 1 }),
59+
) {
60+
return {
61+
...projection().providerSessions[0]!,
62+
id: ProviderSessionId.make(id),
63+
status,
64+
updatedAt,
65+
};
66+
}
67+
68+
function testLayer(
69+
metadata: Readonly<Record<string, { continuationKey: string }>>,
70+
planSelectionTransition: ProviderAdapterV2Shape["planSelectionTransition"] = () =>
71+
Effect.succeed({ type: "restart_session" }),
72+
) {
5473
const adapter = (instanceId: ProviderInstanceId): ProviderAdapterV2Shape => ({
5574
instanceId,
5675
driver,
5776
getCapabilities: () => Effect.succeed(capabilitiesWithoutModelSwitch),
58-
planSelectionTransition: () => Effect.succeed({ type: "restart_session" }),
77+
planSelectionTransition,
5978
openSession: () => Effect.die("ProviderSwitchService tests do not open sessions."),
6079
});
6180
const registry = Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({
@@ -101,6 +120,185 @@ it.effect(
101120
),
102121
);
103122

123+
for (const deadStatus of ["stopped", "error"] as const) {
124+
it.effect(
125+
`restarts and releases the live session when a newer ${deadStatus} session exists`,
126+
() =>
127+
Effect.gen(function* () {
128+
const service = yield* ProviderSwitch.ProviderSwitchServiceV2;
129+
const thread = projection();
130+
const result = yield* service.plan({
131+
projection: {
132+
...thread,
133+
providerSessions: [
134+
...thread.providerSessions,
135+
deadSessionRecord("dead_session", deadStatus),
136+
],
137+
},
138+
targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" },
139+
});
140+
assert.equal(result.transition.type, "restart_and_resume");
141+
assert.deepEqual(result.releaseProviderSessionIds, [currentSessionId]);
142+
}).pipe(
143+
Effect.provide(
144+
testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }),
145+
),
146+
),
147+
);
148+
}
149+
150+
it.effect("releases the newest live session, not the newest record overall", () =>
151+
Effect.gen(function* () {
152+
const service = yield* ProviderSwitch.ProviderSwitchServiceV2;
153+
const thread = projection();
154+
const newerLiveSessionId = ProviderSessionId.make("session_newer_live");
155+
const result = yield* service.plan({
156+
projection: {
157+
...thread,
158+
providerSessions: [
159+
...thread.providerSessions,
160+
{
161+
...thread.providerSessions[0]!,
162+
id: newerLiveSessionId,
163+
updatedAt: DateTime.add(now, { seconds: 1 }),
164+
},
165+
deadSessionRecord("dead_session", "stopped", DateTime.add(now, { seconds: 2 })),
166+
],
167+
},
168+
targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" },
169+
});
170+
assert.equal(result.transition.type, "restart_and_resume");
171+
assert.deepEqual(result.releaseProviderSessionIds, [newerLiveSessionId]);
172+
}).pipe(
173+
Effect.provide(
174+
testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }),
175+
),
176+
),
177+
);
178+
179+
it.effect("creates a fresh session with handoff when every recorded session is dead", () =>
180+
Effect.gen(function* () {
181+
const service = yield* ProviderSwitch.ProviderSwitchServiceV2;
182+
const thread = projection();
183+
const result = yield* service.plan({
184+
projection: {
185+
...thread,
186+
providerSessions: [
187+
deadSessionRecord("dead_session_older", "stopped"),
188+
deadSessionRecord("dead_session_newer", "error", DateTime.add(now, { seconds: 2 })),
189+
],
190+
},
191+
targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" },
192+
});
193+
assert.equal(result.transition.type, "create_with_handoff");
194+
assert.deepEqual(result.releaseProviderSessionIds, []);
195+
}).pipe(
196+
Effect.provide(
197+
testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }),
198+
),
199+
),
200+
);
201+
202+
function deadNativeThreadProjection(
203+
status: "stopped" | "error",
204+
capabilities = capabilitiesWithoutModelSwitch,
205+
): OrchestrationV2ThreadProjection {
206+
const thread = projection();
207+
return {
208+
...thread,
209+
thread: {
210+
...thread.thread,
211+
activeProviderThreadId: ProviderThreadId.make("provider-thread:native"),
212+
},
213+
providerSessions: [{ ...deadSessionRecord("dead_session", status), capabilities }],
214+
providerThreads: [
215+
{
216+
id: ProviderThreadId.make("provider-thread:native"),
217+
driver,
218+
providerInstanceId: currentInstanceId,
219+
providerSessionId: ProviderSessionId.make("dead_session"),
220+
appThreadId: thread.thread.id,
221+
ownerNodeId: null,
222+
nativeThreadRef: {
223+
driver,
224+
nativeId: "native-thread:abc",
225+
strength: "strong",
226+
},
227+
nativeConversationHeadRef: null,
228+
status: "idle",
229+
firstRunOrdinal: null,
230+
lastRunOrdinal: null,
231+
handoffIds: [],
232+
forkedFrom: null,
233+
createdAt: now,
234+
updatedAt: now,
235+
},
236+
],
237+
} as OrchestrationV2ThreadProjection;
238+
}
239+
240+
it.effect("falls back to the native provider thread when every recorded session is dead", () =>
241+
Effect.gen(function* () {
242+
const service = yield* ProviderSwitch.ProviderSwitchServiceV2;
243+
const result = yield* service.plan({
244+
projection: deadNativeThreadProjection("stopped"),
245+
targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" },
246+
});
247+
assert.equal(result.transition.type, "restart_and_resume");
248+
assert.deepEqual(result.releaseProviderSessionIds, []);
249+
}).pipe(
250+
Effect.provide(
251+
testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }),
252+
),
253+
),
254+
);
255+
256+
for (const deadStatus of ["stopped", "error"] as const) {
257+
it.effect(
258+
`applies a model change on next turn when a ${deadStatus} session negotiated model switching`,
259+
() =>
260+
Effect.gen(function* () {
261+
const service = yield* ProviderSwitch.ProviderSwitchServiceV2;
262+
const result = yield* service.plan({
263+
projection: deadNativeThreadProjection(deadStatus, CodexProviderCapabilitiesV2),
264+
targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" },
265+
});
266+
// Static capabilities report no in-session switch, but the dead
267+
// record's negotiated capabilities describe the provider: without
268+
// them the ACP classification rejects the selection instead of
269+
// reopening with the requested model on the next run.
270+
assert.equal(result.transition.type, "switch_model_in_session");
271+
assert.deepEqual(result.releaseProviderSessionIds, []);
272+
}).pipe(
273+
Effect.provide(
274+
testLayer(
275+
{ [currentInstanceId]: { continuationKey: "codex:account:primary" } },
276+
(input) => Effect.succeed(acpSelectionTransition(input)),
277+
),
278+
),
279+
),
280+
);
281+
}
282+
283+
it.effect("rejects a model change the dead record never negotiated support for", () =>
284+
Effect.gen(function* () {
285+
const service = yield* ProviderSwitch.ProviderSwitchServiceV2;
286+
const result = yield* service
287+
.plan({
288+
projection: deadNativeThreadProjection("stopped"),
289+
targetModelSelection: { instanceId: currentInstanceId, model: "gpt-5.2-codex" },
290+
})
291+
.pipe(Effect.flip);
292+
assert.instanceOf(result, ProviderSwitch.ProviderSwitchPlanError);
293+
}).pipe(
294+
Effect.provide(
295+
testLayer({ [currentInstanceId]: { continuationKey: "codex:account:primary" } }, (input) =>
296+
Effect.succeed(acpSelectionTransition(input)),
297+
),
298+
),
299+
),
300+
);
301+
104302
it.effect("distinguishes compatible and incompatible instances of the same driver", () =>
105303
Effect.gen(function* () {
106304
const service = yield* ProviderSwitch.ProviderSwitchServiceV2;

‎apps/server/src/orchestration-v2/ProviderSwitchService.ts‎

Lines changed: 19 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,12 @@ export class ProviderSwitchServiceV2 extends Context.Service<
5151
ProviderSwitchServiceV2Shape
5252
>()("t3/orchestration-v2/ProviderSwitchService/ProviderSwitchServiceV2") {}
5353

54+
// Stopped and errored records stay in session history but can no longer be
55+
// restarted or released; only live sessions participate in a transition.
56+
const isLiveProviderSession = (
57+
session: OrchestrationV2ThreadProjection["providerSessions"][number],
58+
) => session.status !== "stopped" && session.status !== "error";
59+
5460
export const layer: Layer.Layer<
5561
ProviderSwitchServiceV2,
5662
never,
@@ -83,12 +89,18 @@ export const layer: Layer.Layer<
8389
const currentInstance = yield* Effect.option(getMetadata(current.instanceId));
8490
const targetInstance = yield* Effect.option(getMetadata(targetModelSelection.instanceId));
8591
const targetAdapter = yield* Effect.option(adapters.get(targetModelSelection.instanceId));
86-
const currentSession = projection.providerSessions
92+
const currentSessions = projection.providerSessions
8793
.filter((session) => session.providerInstanceId === current.instanceId)
8894
.toSorted(
8995
(left, right) =>
9096
DateTime.toEpochMillis(right.updatedAt) - DateTime.toEpochMillis(left.updatedAt),
91-
)[0];
97+
);
98+
const currentSession = currentSessions.find(isLiveProviderSession);
99+
// Negotiated capabilities describe the provider, not the dead
100+
// process; the newest record still reports what the instance
101+
// supports after its session stops.
102+
const negotiatedCapabilities =
103+
currentSession?.capabilities ?? currentSessions[0]?.capabilities;
92104
// Detaching a process removes its session binding, not its native history.
93105
const currentProviderThread = projection.providerThreads.find(
94106
(thread) =>
@@ -105,8 +117,7 @@ export const layer: Layer.Layer<
105117
? yield* targetAdapter.value.planSelectionTransition({
106118
current,
107119
target: targetModelSelection,
108-
sessionCapabilities:
109-
currentSession?.capabilities ?? currentInstance.value.capabilities,
120+
sessionCapabilities: negotiatedCapabilities ?? currentInstance.value.capabilities,
110121
})
111122
: undefined;
112123
const transition =
@@ -134,7 +145,7 @@ export const layer: Layer.Layer<
134145
projection.thread.worktreePath ??
135146
"<unresolved-workspace>",
136147
capabilities:
137-
currentSession?.capabilities ?? currentInstance.value.capabilities,
148+
negotiatedCapabilities ?? currentInstance.value.capabilities,
138149
},
139150
target: {
140151
driver: targetInstance.value.driver,
@@ -174,8 +185,10 @@ export const layer: Layer.Layer<
174185
)[0];
175186
const releaseProviderSessionIds = projection.providerSessions
176187
.filter((session) => {
177-
if (session.status === "stopped" || session.status === "error") return false;
188+
if (!isLiveProviderSession(session)) return false;
178189
if (transition.type === "restart_and_resume") {
190+
// Other live records may serve pooled or delegated bindings;
191+
// only the session being replaced is released.
179192
return session.id === currentSession?.id;
180193
}
181194
if (transition.type === "create_with_handoff") {

0 commit comments

Comments
 (0)