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
198 changes: 198 additions & 0 deletions apps/server/src/provider/ProviderRegistry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1105,6 +1105,204 @@ it.layer(
assert.deepStrictEqual(afterFailure.models, [authoritativeProvider.models[0]!]);
});

describe("Pi model inventories", () => {
const defaultModel = {
slug: "default",
name: "Pi default",
isCustom: false,
capabilities: null,
} as const;
const currentModel = {
slug: "openai/gpt-6.1-sol",
name: "GPT-6.1-Sol",
isCustom: false,
capabilities: null,
} as const;
const removedModel = {
slug: "openrouter/google/gemini-2.5-pro",
name: "Google: Gemini 2.5 Pro",
isCustom: false,
capabilities: null,
} as const;
const customModel = {
slug: "custom-model",
name: "Custom model",
isCustom: true,
capabilities: null,
} as const;
const cachedProvider = {
instanceId: ProviderInstanceId.make("pi-personal"),
driver: ProviderDriverKind.make("pi"),
status: "ready",
enabled: true,
installed: true,
auth: { status: "authenticated", type: "pi" },
checkedAt: "2026-10-03T00:00:00.000Z",
version: "1.0.0",
models: [
defaultModel,
currentModel,
removedModel,
{ ...customModel, slug: "removed-custom" },
],
slashCommands: [],
skills: [],
} satisfies ServerProvider;
const refreshedProvider = {
...cachedProvider,
checkedAt: "2026-10-03T00:01:00.000Z",
models: [defaultModel, currentModel, customModel],
} satisfies ServerProvider;
const pendingProvider = {
...cachedProvider,
status: "warning",
version: null,
auth: { status: "unknown" },
models: [defaultModel, customModel],
} satisfies ServerProvider;
const failedProvider = {
...pendingProvider,
checkedAt: "2026-10-03T00:02:00.000Z",
status: "ready",
version: "1.0.0",
message: "Pi is available, but T3 Code could not refresh its models and commands.",
} satisfies ServerProvider;

it("drops disconnected Pi upstream models after successful discovery", () => {
assert.deepStrictEqual(
ProviderRegistry.mergeProviderSnapshot(cachedProvider, refreshedProvider).models,
refreshedProvider.models,
);
});

it("clears discovered Pi models when successful discovery finds no usable models", () => {
const signedOutProvider = {
...refreshedProvider,
status: "warning",
auth: { status: "unauthenticated", type: "pi" },
models: [defaultModel, customModel],
} satisfies ServerProvider;

assert.deepStrictEqual(
ProviderRegistry.mergeProviderSnapshot(cachedProvider, signedOutProvider).models,
signedOutProvider.models,
);
});

it("retains Pi inventories during incomplete probes without restoring removed custom models", () => {
const incompleteProviders = [
pendingProvider,
{ ...pendingProvider, installed: false },
failedProvider,
{ ...failedProvider, status: "error" },
] satisfies ReadonlyArray<ServerProvider>;

for (const provider of incompleteProviders) {
assert.deepStrictEqual(
ProviderRegistry.mergeProviderSnapshot(cachedProvider, provider).models,
[defaultModel, customModel, currentModel, removedModel],
);
}
});

it("does not restore disconnected Pi models after a later discovery failure", () => {
const afterRemoval = ProviderRegistry.mergeProviderSnapshot(
cachedProvider,
refreshedProvider,
);
assert.deepStrictEqual(
ProviderRegistry.mergeProviderSnapshot(afterRemoval, failedProvider).models,
[defaultModel, customModel, currentModel],
);
});

it.effect(
"persists Pi inventory removals across failed refreshes and registry restarts",
() =>
Effect.gen(function* () {
const config = yield* ServerConfig.ServerConfig;
const filePath = yield* resolveProviderStatusCachePath({
cacheDir: config.providerStatusCacheDir,
instanceId: cachedProvider.instanceId,
});
yield* writeProviderStatusCache({ filePath, provider: cachedProvider });
const nextProvider = yield* Ref.make<ServerProvider>(refreshedProvider);
const instance = {
instanceId: cachedProvider.instanceId,
driverKind: cachedProvider.driver,
continuationIdentity: {
driverKind: cachedProvider.driver,
continuationKey: "pi:instance:pi-personal",
},
displayName: undefined,
enabled: true,
snapshot: {
resolveMaintenance: () =>
Effect.succeed(
makeManualOnlyProviderMaintenanceCapabilities({
provider: cachedProvider.driver,
packageName: null,
}),
),
getSnapshot: Effect.succeed(pendingProvider),
refresh: Ref.get(nextProvider),
streamChanges: Stream.empty,
applyUsageLimits: () => Effect.void,
},
orchestrationAdapter: {} as ProviderInstance["orchestrationAdapter"],
textGeneration: {} as ProviderInstance["textGeneration"],
} satisfies ProviderInstance;
const layerInstanceRegistry = Layer.succeed(
ProviderInstanceRegistry.ProviderInstanceRegistry,
{
getInstance: (id) =>
Effect.succeed(id === instance.instanceId ? instance : undefined),
listInstances: Effect.succeed([instance]),
listUnavailable: Effect.succeed([]),
streamChanges: Stream.empty,
subscribeChanges: Effect.flatMap(PubSub.unbounded<void>(), PubSub.subscribe),
},
);
const retainedModels = [defaultModel, customModel, currentModel];

for (const restarted of [false, true]) {
yield* Effect.gen(function* () {
const registry = yield* ProviderRegistry.ProviderRegistry;
assert.deepStrictEqual(
(yield* registry.getProviders)[0]?.models,
restarted
? retainedModels
: [defaultModel, customModel, currentModel, removedModel],
);

yield* registry.refreshInstance(instance.instanceId);
assert.deepStrictEqual(
(yield* readProviderStatusCache(filePath))?.models,
restarted ? retainedModels : refreshedProvider.models,
);

yield* Ref.set(nextProvider, failedProvider);
const afterFailure = yield* registry.refreshInstance(instance.instanceId);
assert.deepStrictEqual(afterFailure[0]?.models, retainedModels);
assert.deepStrictEqual(
(yield* readProviderStatusCache(filePath))?.models,
retainedModels,
);
}).pipe(
Effect.provide(ProviderRegistry.layer.pipe(Layer.provide(layerInstanceRegistry))),
Effect.scoped,
);
}
}).pipe(
Effect.provide(
ServerConfig.layerTest(process.cwd(), {
prefix: "t3-pi-model-cache-",
}).pipe(Layer.provideMerge(NodeServices.layer)),
),
),
);
});

describe("Codex model inventories", () => {
const cachedProvider = {
instanceId: ProviderInstanceId.make("codex-personal"),
Expand Down
11 changes: 11 additions & 0 deletions apps/server/src/provider/ProviderRegistry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -209,6 +209,17 @@ export function upsertProviderWorkspaceSnapshot(
}

const shouldRetainMissingProviderModels = (provider: ServerProvider): boolean => {
if (provider.driver === ProviderDriverKind.make("pi")) {
// Pi reports known auth only after RPC inventory discovery completes.
// An empty inventory is warning/unauthenticated; failed or interactive
// discovery stays ready/unknown and must keep the last known models.
return !(
provider.installed &&
(provider.status === "ready" || provider.status === "warning") &&
provider.auth.status !== "unknown"
);
}

if (provider.driver === ProviderDriverKind.make("acpRegistry")) {
// ACP Registry discovery probes return the agent's complete inventory, so
// a completed probe (ready and authenticated) replaces the model list —
Expand Down
107 changes: 107 additions & 0 deletions packages/provider-pi/src/server/status.test.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,20 @@
import * as NodeServices from "@effect/platform-node/NodeServices";
import { assert, describe, it } from "@effect/vitest";
import * as Cause from "effect/Cause";
import * as Effect from "effect/Effect";
import * as Queue from "effect/Queue";
import * as Schema from "effect/Schema";
import * as Sink from "effect/Sink";
import * as Stream from "effect/Stream";
import { ChildProcess, ChildProcessSpawner } from "effect/process";

import { checkPiProviderStatus, MINIMUM_PI_VERSION } from "./status.ts";

const encoder = new TextEncoder();
const decodeRpcRequest = Schema.decodeSync(
Schema.fromJsonString(Schema.Struct({ id: Schema.String, type: Schema.String })),
);
const encodeJsonLine = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown));

function processHandle(input: {
readonly stdout?: string;
Expand Down Expand Up @@ -44,6 +51,51 @@ function piProbeSpawner(version: string) {
});
}

const piDiscoverySpawner = (inventory: unknown) =>
Effect.gen(function* () {
const stdout = yield* Queue.unbounded<Uint8Array, Cause.Done>();
let stdinBuffer = "";
const stdin = Sink.forEach((chunk: Uint8Array) =>
Effect.gen(function* () {
stdinBuffer += new TextDecoder().decode(chunk);
while (true) {
const newline = stdinBuffer.indexOf("\n");
if (newline === -1) return;
const line = stdinBuffer.slice(0, newline);
stdinBuffer = stdinBuffer.slice(newline + 1);
if (line.length === 0) continue;
const request = decodeRpcRequest(line);
const data =
request.type === "get_available_models"
? inventory
: request.type === "get_commands"
? { commands: [] }
: {};
yield* Queue.offer(
stdout,
encoder.encode(
`${encodeJsonLine({ type: "response", id: request.id, command: request.type, success: true, data })}\n`,
),
);
}
}),
);
return ChildProcessSpawner.make((command) => {
const args = ChildProcess.isStandardCommand(command) ? command.args : [];
return Effect.succeed(
args.includes("--version")
? processHandle({ stdout: "pi 0.84.3\n" })
: ChildProcessSpawner.makeHandle({
...processHandle({}),
exitCode: Effect.never,
isRunning: Effect.succeed(true),
stdin,
stdout: Stream.fromQueue(stdout),
}),
);
});
});

const settings = {
enabled: true,
binaryPath: "pi",
Expand Down Expand Up @@ -77,4 +129,59 @@ describe("PiProvider", () => {
assert.include(snapshot.message ?? "", "could not refresh its models and commands");
}).pipe(Effect.provide(NodeServices.layer)),
);

it.effect("reports only the models returned by successful Pi discovery", () =>
Effect.gen(function* () {
const spawner = yield* piDiscoverySpawner({
models: [{ provider: "anthropic", id: "claude-sonnet-4" }],
});
const snapshot = yield* checkPiProviderStatus(settings).pipe(
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
);
assert.equal(snapshot.status, "ready");
assert.equal(snapshot.auth.status, "authenticated");
assert.deepEqual(
snapshot.models.map((model) => model.slug),
["default", "anthropic/claude-sonnet-4"],
);
assert.equal(snapshot.models[1]?.subProvider, "anthropic");
}).pipe(Effect.provide(NodeServices.layer)),
);

it.effect("reports a known unauthenticated inventory when Pi returns no usable models", () =>
Effect.gen(function* () {
const spawner = yield* piDiscoverySpawner({ models: [] });
const snapshot = yield* checkPiProviderStatus(settings).pipe(
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
);
assert.equal(snapshot.status, "warning");
assert.equal(snapshot.auth.status, "unauthenticated");
assert.deepEqual(
snapshot.models.map((model) => model.slug),
["default"],
);
}).pipe(Effect.provide(NodeServices.layer)),
);

it.effect.each([
["missing models", {}],
["non-array models", { models: {} }],
["missing model provider", { models: [{ id: "claude-sonnet-4" }] }],
["missing model id", { models: [{ provider: "anthropic" }] }],
["empty model provider", { models: [{ provider: "", id: "claude-sonnet-4" }] }],
["empty model id", { models: [{ provider: "anthropic", id: "" }] }],
[
"mixed valid and invalid models",
{ models: [{ provider: "anthropic", id: "claude-sonnet-4" }, { provider: "openrouter" }] },
],
] as const)("reports unknown authentication when Pi discovery returns %s", (_label, inventory) =>
Effect.gen(function* () {
const spawner = yield* piDiscoverySpawner(inventory);
const snapshot = yield* checkPiProviderStatus(settings).pipe(
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
);
assert.equal(snapshot.status, "ready");
assert.equal(snapshot.auth.status, "unknown");
}).pipe(Effect.provide(NodeServices.layer)),
);
});
Loading
Loading