diff --git a/apps/server/src/provider/ProviderRegistry.test.ts b/apps/server/src/provider/ProviderRegistry.test.ts index 8653bd32ebd2..f85e273f9e66 100644 --- a/apps/server/src/provider/ProviderRegistry.test.ts +++ b/apps/server/src/provider/ProviderRegistry.test.ts @@ -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; + + 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(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(), 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"), diff --git a/apps/server/src/provider/ProviderRegistry.ts b/apps/server/src/provider/ProviderRegistry.ts index 5f5120ee28bc..9cbcba9d3f0c 100644 --- a/apps/server/src/provider/ProviderRegistry.ts +++ b/apps/server/src/provider/ProviderRegistry.ts @@ -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 — diff --git a/packages/provider-pi/src/server/status.test.ts b/packages/provider-pi/src/server/status.test.ts index 133a505612b1..08f9fa211adb 100644 --- a/packages/provider-pi/src/server/status.test.ts +++ b/packages/provider-pi/src/server/status.test.ts @@ -1,6 +1,9 @@ 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"; @@ -8,6 +11,10 @@ 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; @@ -44,6 +51,51 @@ function piProbeSpawner(version: string) { }); } +const piDiscoverySpawner = (inventory: unknown) => + Effect.gen(function* () { + const stdout = yield* Queue.unbounded(); + 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", @@ -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)), + ); }); diff --git a/packages/provider-pi/src/server/status.ts b/packages/provider-pi/src/server/status.ts index 54be47d3dfd9..ffab91fe8646 100644 --- a/packages/provider-pi/src/server/status.ts +++ b/packages/provider-pi/src/server/status.ts @@ -102,15 +102,23 @@ function piModelsFromSettings( function parseDiscoveredModels( data: unknown, defaultThinkingLevel: unknown, -): ReadonlyArray { +): Result.Result, PiRpcError> { const models = recordField(data, "models"); - if (!Array.isArray(models)) return []; + if (!Array.isArray(models)) { + return Result.fail( + new PiRpcError({ operation: "get_available_models", detail: "Invalid model inventory" }), + ); + } const seen = new Set(); const parsed: Array = []; for (const model of models) { const provider = recordString(model, "provider"); const id = recordString(model, "id"); - if (provider === undefined || id === undefined) continue; + if (provider === undefined || id === undefined || provider.trim() === "" || id.trim() === "") { + return Result.fail( + new PiRpcError({ operation: "get_available_models", detail: "Invalid model identity" }), + ); + } const slug = `${provider}/${id}`; if (seen.has(slug)) continue; seen.add(slug); @@ -122,7 +130,7 @@ function parseDiscoveredModels( capabilities: thinkingCapabilitiesForPiModel(model, defaultThinkingLevel), }); } - return parsed; + return Result.succeed(parsed); } const makePiDiscoveryConnection = Effect.fnUntraced(function* ( @@ -165,9 +173,8 @@ const discoverPiViaRpc = ( const commandsData = yield* connection .request({ type: "get_commands" }) .pipe(Effect.orElseSucceed(() => undefined)); - const discoveredModels = parseDiscoveredModels( - modelsData, - recordString(stateData, "thinkingLevel"), + const discoveredModels = yield* Effect.fromResult( + parseDiscoveredModels(modelsData, recordString(stateData, "thinkingLevel")), ); const { slashCommands, skills } = parsePiDiscoveredCommands(commandsData); return {