Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
202 changes: 200 additions & 2 deletions apps/server/src/orchestration-v2/ProviderSwitchService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import {
ProviderDriverKind,
ProviderInstanceId,
ProviderSessionId,
ProviderThreadId,
ThreadId,
type OrchestrationV2ThreadProjection,
} from "@t3tools/contracts";
Expand All @@ -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");
Expand Down Expand Up @@ -50,12 +52,29 @@ function projection(): OrchestrationV2ThreadProjection {
} as unknown as OrchestrationV2ThreadProjection;
}

function testLayer(metadata: Readonly<Record<string, { continuationKey: string }>>) {
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<Record<string, { continuationKey: string }>>,
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)({
Expand Down Expand Up @@ -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;
Expand Down
25 changes: 19 additions & 6 deletions apps/server/src/orchestration-v2/ProviderSwitchService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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) =>
Expand All @@ -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 =
Expand Down Expand Up @@ -134,7 +145,7 @@ export const layer: Layer.Layer<
projection.thread.worktreePath ??
"<unresolved-workspace>",
capabilities:
currentSession?.capabilities ?? currentInstance.value.capabilities,
negotiatedCapabilities ?? currentInstance.value.capabilities,
},
target: {
driver: targetInstance.value.driver,
Expand Down Expand Up @@ -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") {
Expand Down
Loading
Loading