Skip to content
Open
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
131 changes: 122 additions & 9 deletions apps/server/src/orchestration-v2/Adapters/PiAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -86,8 +86,8 @@ interface FakePi {
readonly deferNextState: () => void;
/** Resolve the held `get_state` request. */
readonly resolveDeferredState: (data: unknown) => Effect.Effect<void>;
/** Reject the next `get_state` request. */
readonly failNextState: () => void;
/** Reject a `get_state` request, after `skip` successful acks. */
readonly failNextState: (skip?: number) => void;
readonly deferNextLifecycle: (type: "switch_session" | "new_session") => void;
readonly queueModels: (models: ReadonlyArray<unknown>) => void;
readonly vetoNextNewSession: () => void;
Expand Down Expand Up @@ -142,7 +142,7 @@ const makeFakePi: Effect.Effect<FakePi> = Effect.gen(function* () {
const allRequests: Array<PiRpcRecord> = [];
let deferState = false;
let deferredStateRequest: PiRpcRecord | undefined;
let failState = false;
let failStateAfter: number | null = null;
let vetoSwitch = false;
let vetoNewSession = false;
let deferredLifecycle: string | undefined;
Expand All @@ -166,8 +166,8 @@ const makeFakePi: Effect.Effect<FakePi> = Effect.gen(function* () {
};
switch (record["type"]) {
case "get_state":
if (failState) {
failState = false;
if (failStateAfter !== null && failStateAfter-- === 0) {
failStateAfter = null;
return { ...base, success: false, error: "state unavailable" };
}
// Queued data overrides fields of the recorded idle state, so a test
Expand Down Expand Up @@ -285,8 +285,8 @@ const makeFakePi: Effect.Effect<FakePi> = Effect.gen(function* () {
data,
});
}),
failNextState: () => {
failState = true;
failNextState: (skip = 0) => {
failStateAfter = skip;
},
deferNextLifecycle: (type) => {
deferredLifecycle = type;
Expand Down Expand Up @@ -338,13 +338,15 @@ const openRuntime = Effect.fnUntraced(function* (
threadId = THREAD_ID,
providerSessionId = SESSION_ID,
forkFake?: FakePi,
initialNativeThreadId?: string,
) {
const adapter = yield* makeAdapter(fake, "", forkFake);
const runtime = yield* adapter.openSession({
threadId,
providerSessionId,
modelSelection: modelSelection(model),
runtimePolicy,
...(initialNativeThreadId === undefined ? {} : { initialNativeThreadId }),
});
const emitted = yield* Queue.unbounded<ProviderAdapterV2Event>();
yield* runtime.events.pipe(
Expand All @@ -361,6 +363,28 @@ const openRuntime = Effect.fnUntraced(function* (
return { runtime, takeEvent };
});

/** A provider thread as persisted before a cold resume. */
const makePersistedProviderThread = Effect.fnUntraced(function* (nativeId: string) {
const now = yield* DateTime.now;
return {
id: ProviderThreadId.make(`thread:provider:pi:native-thread:${nativeId}`),
driver: PI_PROVIDER,
providerInstanceId: PI_INSTANCE_ID,
providerSessionId: SESSION_ID,
appThreadId: THREAD_ID,
ownerNodeId: null,
nativeThreadRef: { driver: PI_PROVIDER, nativeId, strength: "strong" },
nativeConversationHeadRef: null,
status: "not_loaded",
firstRunOrdinal: 1,
lastRunOrdinal: 1,
handoffIds: [],
forkedFrom: null,
createdAt: now,
updatedAt: now,
} satisfies OrchestrationV2ProviderThread;
});

const makeAppThread = Effect.fnUntraced(function* (model: string, threadId = THREAD_ID) {
const now = yield* DateTime.now;
return {
Expand Down Expand Up @@ -514,6 +538,43 @@ describe("PiAdapterV2", () => {
),
);

it.effect("resumes a persisted thread at spawn without replacing the session", () =>
Effect.gen(function* () {
const fake = yield* makeFakePi;
// Cold resume: the orchestrator knows the native session up front, so
// pi is spawned attached to it and no switch_session may follow. A
// switch would dispose the session and stale every extension context.
const { runtime } = yield* openRuntime(
fake,
"default",
THREAD_ID,
SESSION_ID,
undefined,
FAKE_SESSION_FILE,
);
const spawnArgs = fake.lastSpawn().args;
assert.equal(spawnArgs[spawnArgs.indexOf("--session") + 1], FAKE_SESSION_FILE);
assert.isFalse(spawnArgs.includes("--no-session"));
assert.isFalse(spawnArgs.includes("--no-extensions"));

const providerThread = yield* runtime.ensureThread({
threadId: THREAD_ID,
modelSelection: modelSelection("default"),
runtimePolicy,
existingProviderThread: yield* makePersistedProviderThread(FAKE_SESSION_FILE),
});
assert.equal(providerThread.nativeThreadRef?.nativeId, FAKE_SESSION_FILE);
assert.equal(providerThread.driver, PI_PROVIDER);
assert.isFalse(fake.allRequests().some((request) => request.type === "switch_session"));

yield* startTurn(runtime, providerThread, "anthropic/claude-sonnet");
const setModel = yield* fake.takeRequest("set_model");
assert.equal(setModel["provider"], "anthropic");
assert.equal(runtime.providerSession.model, "anthropic/claude-sonnet");
yield* fake.takeRequest("prompt");
}).pipe(Effect.scoped, Effect.provide(testLayer)),
);

it.effect("rejects a resume while a turn is active", () =>
Effect.gen(function* () {
const fake = yield* makeFakePi;
Expand All @@ -540,6 +601,8 @@ describe("PiAdapterV2", () => {
modelSelection: modelSelection("default"),
runtimePolicy,
});
// The live process holds another session file, so the resume has to switch.
fake.queueState({ sessionFile: "/fake/.pi/agent/sessions/--workspace--/0009_live.jsonl" });
fake.deferNextLifecycle("switch_session");
const resumed = yield* runtime.resumeThread({ providerThread }).pipe(Effect.forkChild);
const request = yield* fake.takeRequest("switch_session");
Expand All @@ -558,6 +621,44 @@ describe("PiAdapterV2", () => {
}).pipe(Effect.scoped, Effect.provide(testLayer)),
);

it.effect("switches a live session when it holds a different session file", () =>
Effect.gen(function* () {
const fake = yield* makeFakePi;
const { runtime, takeEvent } = yield* openRuntime(fake);
const providerThread = yield* runtime.ensureThread({
threadId: THREAD_ID,
modelSelection: modelSelection("default"),
runtimePolicy,
});
yield* startTurn(runtime, providerThread, "anthropic/claude-sonnet");
yield* fake.takeRequest("prompt");
yield* fake.emit({ type: "agent_start" });
yield* fake.emit({ type: "agent_settled" });
yield* takeEvent((event) => event.type === "turn.terminal");

// Fork adoption: the process still holds the source session, so the
// adopted file has to be switched in.
const adoptedFile = "/fake/.pi/agent/sessions/--workspace--/0002_def.jsonl";
fake.queueState({ thinkingLevel: "medium", sessionFile: FAKE_SESSION_FILE });
fake.queueState({ thinkingLevel: "medium", sessionFile: adoptedFile });
const adopted = yield* runtime.resumeThread({
providerThread: {
...providerThread,
nativeThreadRef: { driver: PI_PROVIDER, nativeId: adoptedFile, strength: "strong" },
},
});
const switchRequest = yield* fake.takeRequest("switch_session");
assert.equal(switchRequest["sessionPath"], adoptedFile);
assert.equal(adopted.nativeThreadRef?.nativeId, adoptedFile);

// The switch dropped the applied-selection cache, so the same model has
// to be re-applied to the replacement session.
yield* startTurn(runtime, adopted, "anthropic/claude-sonnet", [], "Again", undefined, 2);
yield* fake.takeRequest("prompt");
assert.equal(fake.allRequests().filter((request) => request.type === "set_model").length, 2);
}).pipe(Effect.scoped, Effect.provide(testLayer)),
);

it.effect("creates a distinct native session after a failed resume", () =>
Effect.gen(function* () {
const fake = yield* makeFakePi;
Expand All @@ -567,8 +668,11 @@ describe("PiAdapterV2", () => {
modelSelection: modelSelection("default"),
runtimePolicy,
});
fake.queueState({ sessionFile: "/fake/.pi/agent/sessions/--workspace--/0009_live.jsonl" });
fake.vetoNextSwitch();
yield* runtime.resumeThread({ providerThread }).pipe(Effect.flip);
const error = yield* runtime.resumeThread({ providerThread }).pipe(Effect.flip);
assert.equal(error._tag, "ProviderAdapterResumeThreadError");
assert.match(String(error.cause), /cancelled the session switch/);
const replacement = yield* runtime.ensureThread({
threadId: THREAD_ID,
modelSelection: modelSelection("default"),
Expand Down Expand Up @@ -629,6 +733,7 @@ describe("PiAdapterV2", () => {
modelSelection: modelSelection("default"),
runtimePolicy,
});
fake.queueState({ sessionFile: "/fake/.pi/agent/sessions/--workspace--/0009_live.jsonl" });
fake.deferNextLifecycle("switch_session");
const resumed = yield* runtime
.resumeThread({ providerThread })
Expand Down Expand Up @@ -721,8 +826,12 @@ describe("PiAdapterV2", () => {
});
const fake = yield* makeFakePi;
const { runtime } = yield* openRuntime(fake);
fake.failNextState();
// The live process holds another file: the switch succeeds and only the
// refresh after it fails.
fake.queueState({ sessionFile: "/fake/.pi/agent/sessions/--workspace--/0009_live.jsonl" });
fake.failNextState(1);
yield* runtime.resumeThread({ providerThread }).pipe(Effect.flip);
assert.isTrue(fake.allRequests().some((request) => request.type === "switch_session"));
const replacement = yield* runtime.ensureThread({
threadId: THREAD_ID,
modelSelection: modelSelection("default"),
Expand All @@ -748,6 +857,7 @@ describe("PiAdapterV2", () => {
modelSelection: modelSelection("default"),
runtimePolicy,
});
fake.queueState({ sessionFile: "/fake/.pi/agent/sessions/--workspace--/0009_live.jsonl" });
fake.deferNextLifecycle("switch_session");
const resumed = yield* runtime.resumeThread({ providerThread }).pipe(Effect.forkChild);
yield* fake.takeRequest("switch_session");
Expand Down Expand Up @@ -1015,6 +1125,9 @@ describe("PiAdapterV2", () => {
completedAt: null,
});
forkFake.queueState({ sessionFile: forkFile });
// The live process is still on the source session when the fork is
// adopted, then reports the forked file once it has switched.
fake.queueState({ sessionFile: FAKE_SESSION_FILE });
fake.queueState({ sessionFile: forkFile });
const target = ThreadId.make("fork-target");
const forked = yield* runtime.forkThread({
Expand Down
31 changes: 29 additions & 2 deletions apps/server/src/orchestration-v2/Adapters/PiAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,9 @@
* Terminal-only decoration such as status, widget, title, and editor-text
* updates has no matching T3 surface and is ignored.
*/
// @effect-diagnostics-next-line nodeBuiltinImport:off
import * as NodePath from "node:path";

import { HostProcessEnvironment } from "@t3tools/shared/hostProcess";
import { getModelSelectionStringOptionValue } from "@t3tools/shared/model";
import {
Expand Down Expand Up @@ -249,6 +252,16 @@ function providerRef(
return { driver: PI_PROVIDER, nativeId, strength };
}

/**
* Compares a live `get_state.sessionFile` with a stored `nativeThreadRef`.
* Both are absolute paths Pi echoed back, but they can be captured through
* different cwds or with redundant separators, so normalize before comparing.
*/
function samePiSessionFile(live: string | undefined, stored: string): boolean {
if (live === undefined) return false;
return NodePath.resolve(live) === NodePath.resolve(stored);
}

const PI_THINKING_LEVELS = new Set(["off", "minimal", "low", "medium", "high", "xhigh", "max"]);

// ── per-session state ─────────────────────────────────────────
Expand Down Expand Up @@ -417,6 +430,9 @@ export function makePiAdapterV2(
environment: options.environment,
mcpSession,
extensionPath,
// Resuming a persisted thread attaches its session file at spawn, so
// extensions load against the session they will actually run in.
sessionPath: input.initialNativeThreadId,
runtimeMode: input.runtimePolicy.runtimeMode,
});
const connection: PiRpcConnection = yield* makePiRpcConnection({
Expand Down Expand Up @@ -1996,7 +2012,17 @@ export function makePiAdapterV2(
const resumeId = existing?.nativeThreadRef?.nativeId;
const needsNewSession = resumeId == null && registrationAttempted;
registrationAttempted = true;
if (resumeId != null || needsNewSession) {
// Pi treats `switch_session` as a full session replacement: it disposes
// the session and re-runs every extension factory, which invalidates
// contexts already bound to it. A process spawned with `--session` is
// already on the wanted file, so read the live session first and leave
// it alone when it matches.
let stateData = resumeId != null ? yield* request({ type: "get_state" }) : undefined;
const alreadyOnResumeFile =
resumeId != null && samePiSessionFile(recordString(stateData, "sessionFile"), resumeId);
if (alreadyOnResumeFile) {
lastNativeThreadId = resumeId;
} else if (resumeId != null || needsNewSession) {
lastNativeThreadId = resumeId ?? lastNativeThreadId;
// Even a failed lifecycle operation can change Pi's native session.
// Never leave the old app binding or model defaults usable afterward.
Expand All @@ -2015,8 +2041,9 @@ export function makePiAdapterV2(
if (recordField(result, "cancelled") === true) {
return yield* protocolError("A Pi extension cancelled the session switch");
}
stateData = undefined;
}
const stateData = yield* request({ type: "get_state" });
stateData ??= yield* request({ type: "get_state" });
if (!modelsDiscovered) {
const modelsData = yield* request({ type: "get_available_models" }).pipe(
Effect.orElseSucceed(() => undefined),
Expand Down
47 changes: 47 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/piT3McpInjection.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,53 @@ describe("pi T3 MCP injection", () => {
assert.equal(permissionOnly.env[T3_PI_RUNTIME_MODE_ENV], "auto-accept-edits");
});

it("attaches a persisted session file at spawn instead of switching later", () => {
const launch = buildPiRpcLaunch({
launchArgs: ["--model", "claude-sonnet"],
environment: {},
mcpSession: undefined,
extensionPath: "/tmp/cache/pi-t3-mcp-extension.ts",
sessionPath: "/home/user/.pi/agent/sessions/--work--/0001_abc.jsonl",
});
assert.deepEqual(launch.args, [
"--mode",
"rpc",
"--session",
"/home/user/.pi/agent/sessions/--work--/0001_abc.jsonl",
"--model",
"claude-sonnet",
"--extension",
"/tmp/cache/pi-t3-mcp-extension.ts",
]);

// --session and --no-session cannot be combined.
const ephemeral = buildPiRpcLaunch({
launchArgs: [],
environment: {},
mcpSession: undefined,
extensionPath: undefined,
sessionPath: "/home/user/.pi/agent/sessions/--work--/0001_abc.jsonl",
ephemeral: true,
});
assert.deepEqual(ephemeral.args, ["--mode", "rpc", "--no-session"]);

// T3 owns session identity, so the user cannot pick the session itself.
for (const rejected of [
"--session old.jsonl",
"--session-id abc",
"--no-session",
"--resume",
"--continue",
"--fork old.jsonl",
]) {
assert.deepInclude(
resolvePiLaunchArgs(rejected),
{ ok: false },
`expected '${rejected}' to be rejected`,
);
}
});

it("falls back to Pi's first supported mode for legacy auto threads", () => {
const launch = buildPiRpcLaunch({
launchArgs: [],
Expand Down
13 changes: 12 additions & 1 deletion apps/server/src/orchestration-v2/Adapters/piT3McpInjection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -252,6 +252,13 @@ export function buildPiRpcLaunch(input: {
readonly environment: NodeJS.ProcessEnv;
readonly mcpSession: McpProviderSessionConfig | undefined;
readonly extensionPath: string | undefined;
/**
* Persisted pi session file to attach at spawn. Resuming this way avoids the
* post-spawn `switch_session`, which replaces the session and invalidates
* every extension context loaded against the default one. Ignored for
* ephemeral sessions, which never have a persisted file.
*/
readonly sessionPath?: string | undefined;
readonly ephemeral?: boolean;
readonly disableExtensions?: boolean;
readonly disableTools?: boolean;
Expand All @@ -269,10 +276,14 @@ export function buildPiRpcLaunch(input: {
: input.launchArgs;
const launchArgs =
input.disableTools === true ? withoutToolSelectionArgs(extensionSafeArgs) : extensionSafeArgs;
const ephemeral = input.ephemeral === true;
const args = [
"--mode",
"rpc",
...(input.ephemeral === true ? ["--no-session"] : []),
// `--session` and `--no-session` are mutually exclusive; an ephemeral
// session has no persisted file to attach to.
...(!ephemeral && input.sessionPath !== undefined ? ["--session", input.sessionPath] : []),
...(ephemeral ? ["--no-session"] : []),
...launchArgs,
// Restrictions follow user launch args so a configured --tools or
// --extension cannot silently re-enable unattended text-generation code.
Expand Down
Loading