From ed9c2f65bbffc476832d4a3e0c6ff8324f4f9f33 Mon Sep 17 00:00:00 2001 From: Federico Jaramillo Martinez Date: Fri, 24 Jul 2026 23:44:35 +0200 Subject: [PATCH] fix(sessions): move provider catalog network refreshes off request paths MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The shared ModelRuntime was constructed with network refreshes enabled, so reloadConfig()/login()/logout() — called on the model picker, session model changes, and auth dialogs — performed unbounded provider-catalog fetches. A single stalled fetch blocked those requests for minutes and, through pi's coalesced per-provider refresh, dragged session creation along with it. Construct the runtime with PI_OFFLINE forced so every runtime-driven refresh stays local, and add ModelCatalogRefresher as the single deliberate network path: bounded by an abort timeout, serialized through one in-flight run, scheduled in the background, and triggered after provider auth changes. --- .changeset/bounded-model-catalog-refresh.md | 5 + src/server/sessiond.ts | 12 +- src/server/sessions/authService.test.ts | 37 +++- src/server/sessions/authService.ts | 32 +++- .../sessions/modelCatalogRefresher.test.ts | 173 ++++++++++++++++++ src/server/sessions/modelCatalogRefresher.ts | 108 +++++++++++ 6 files changed, 358 insertions(+), 9 deletions(-) create mode 100644 .changeset/bounded-model-catalog-refresh.md create mode 100644 src/server/sessions/modelCatalogRefresher.test.ts create mode 100644 src/server/sessions/modelCatalogRefresher.ts diff --git a/.changeset/bounded-model-catalog-refresh.md b/.changeset/bounded-model-catalog-refresh.md new file mode 100644 index 0000000..d6c9e45 --- /dev/null +++ b/.changeset/bounded-model-catalog-refresh.md @@ -0,0 +1,5 @@ +--- +"@jmfederico/pi-web": patch +--- + +Fix multi-minute stalls when opening the model selector, starting sessions, or using the auth dialogs. Provider catalog refreshes no longer run unbounded network fetches on request paths: the shared model runtime now operates offline and pi-web refreshes provider catalogs itself on a bounded, background schedule and after provider logins. diff --git a/src/server/sessiond.ts b/src/server/sessiond.ts index 8703dba..949eeb0 100644 --- a/src/server/sessiond.ts +++ b/src/server/sessiond.ts @@ -9,6 +9,7 @@ import { SessionEventHub } from "./realtime/sessionEventHub.js"; import { AuthService } from "./sessions/authService.js"; import { bootstrapAndFreezeGlobalExtensionProviders } from "./sessions/globalProviderPolicy.js"; import { registerAuthRoutes } from "./sessions/authRoutes.js"; +import { ModelCatalogRefresher } from "./sessions/modelCatalogRefresher.js"; import { PiSessionService } from "./sessions/piSessionService.js"; import { createPiSessionManagerGateway } from "./sessions/piSessionManagerGateway.js"; import { registerSessionRoutes } from "./sessions/sessionRoutes.js"; @@ -55,6 +56,12 @@ await runSessionDaemonStartup({ // still mutable, then freeze every later extension-provider mutation before // any real session can load project resources. await bootstrapAndFreezeGlobalExtensionProviders(auth.runtime, activeAgentProfile.dir, app.log); + // The shared model runtime is constructed offline so request paths never + // wait on provider-catalog fetches; this is the single bounded network + // refresher, and auth changes (login/logout) ask it for a prompt run. + const catalogRefresher = new ModelCatalogRefresher({ runtime: auth.runtime, logger: app.log }); + catalogRefresher.start(); + auth.subscribe(() => { catalogRefresher.requestRefresh(); }); const spawnTargets = config.spawnSessions ? new ProjectScopedSpawnTargetResolver({ projects: new ProjectService(new ProjectStore()), workspaces: new WorkspaceService() }) : undefined; @@ -79,7 +86,7 @@ await runSessionDaemonStartup({ ...getPiWebRuntimeComponent("sessiond", SESSIOND_RUNTIME_CAPABILITIES), activeAgentProfile, }); - return { eventHub, workspaceActivity, auth, sessions, terminals, unreadStore, activeAgentProfile, runtimeComponent }; + return { eventHub, workspaceActivity, auth, sessions, terminals, unreadStore, activeAgentProfile, runtimeComponent, catalogRefresher }; }, registerRoutes({ eventHub, workspaceActivity, auth, sessions, terminals, runtimeComponent }) { registerWorkspaceActivityRoutes(app, workspaceActivity); @@ -102,7 +109,7 @@ await runSessionDaemonStartup({ app.get("/runtime", () => runtimeComponent); }, - async listen({ auth, sessions, terminals, unreadStore }) { + async listen({ auth, sessions, terminals, unreadStore, catalogRefresher }) { let shuttingDown = false; async function shutdown(signal: NodeJS.Signals): Promise { if (shuttingDown) return; @@ -117,6 +124,7 @@ await runSessionDaemonStartup({ } }; await attempt("dispose terminals", () => { terminals.dispose(); }); + await attempt("dispose catalog refresher", () => { catalogRefresher.dispose(); }); await attempt("dispose auth", () => { auth.dispose(); }); await attempt("dispose sessions", () => sessions.dispose()); await attempt("flush session unread state", () => unreadStore.flush()); diff --git a/src/server/sessions/authService.test.ts b/src/server/sessions/authService.test.ts index 8a07d0a..7779966 100644 --- a/src/server/sessions/authService.test.ts +++ b/src/server/sessions/authService.test.ts @@ -362,7 +362,7 @@ describe("AuthService", () => { it("stores credentials in the configured agent directory", async () => { const agentDir = await tempAgentDir(); - const runtime = await createModelRuntimeForAgentDir(agentDir, false); + const runtime = await createModelRuntimeForAgentDir(agentDir); const auth = await AuthService.create({ runtime }); await auth.saveApiKey("anthropic", "sk-test"); @@ -470,6 +470,39 @@ describe("AuthService", () => { }); }); +describe("createModelRuntimeForAgentDir", () => { + it("disables runtime-owned network refreshes so request paths stay local", async () => { + // The stall fix relies on this construction-time flag: reloadConfig(), + // login(), and logout() refresh with allowNetwork = modelNetworkEnabled. + // Asserting the private field is proportionate here because the flag is + // exactly the contract this change depends on. + const agentDir = await tempAgentDir(); + const runtime = await createModelRuntimeForAgentDir(agentDir); + expect(Reflect.get(runtime, "modelNetworkEnabled")).toBe(false); + }); + + it("restores a previously set PI_OFFLINE after runtime creation", async () => { + // The file-level beforeEach stubs PI_OFFLINE=1. + const agentDir = await tempAgentDir(); + await createModelRuntimeForAgentDir(agentDir); + expect(process.env["PI_OFFLINE"]).toBe("1"); + }); + + it("restores a previously unset PI_OFFLINE after runtime creation", async () => { + vi.unstubAllEnvs(); + const previous = process.env["PI_OFFLINE"]; + delete process.env["PI_OFFLINE"]; + try { + const agentDir = await tempAgentDir(); + const runtime = await createModelRuntimeForAgentDir(agentDir); + expect(Reflect.get(runtime, "modelNetworkEnabled")).toBe(false); + expect(process.env["PI_OFFLINE"]).toBeUndefined(); + } finally { + if (previous !== undefined) process.env["PI_OFFLINE"] = previous; + } + }); +}); + async function createAuthService(seed: Record = {}, logger?: AuthServiceLogger) { const credentials = new InMemoryCredentialStore(); for (const [providerId, credential] of Object.entries(seed)) { @@ -486,7 +519,7 @@ async function createFileBackedAuthService(seed: Record) { const agentDir = await tempAgentDir(); const authPath = join(agentDir, "auth.json"); await writeFile(authPath, JSON.stringify(seed, null, 2)); - const runtime = await createModelRuntimeForAgentDir(agentDir, false); + const runtime = await createModelRuntimeForAgentDir(agentDir); const auth = await AuthService.create({ runtime }); const changes: AuthChange[] = []; auth.subscribe((change) => { changes.push(change); }); diff --git a/src/server/sessions/authService.ts b/src/server/sessions/authService.ts index 5a237ff..6e7eaca 100644 --- a/src/server/sessions/authService.ts +++ b/src/server/sessions/authService.ts @@ -1,5 +1,5 @@ import { join } from "node:path"; -import { ModelRuntime } from "@earendil-works/pi-coding-agent"; +import { ModelRuntime, type CreateModelRuntimeOptions } from "@earendil-works/pi-coding-agent"; import type { AuthInteraction } from "@earendil-works/pi-ai"; import type { AuthProvidersResponse, AuthType, OAuthFlowState } from "../../shared/apiTypes.js"; import { getLoginProviderOptions, getLogoutProviderOptions } from "./authProviderOptions.js"; @@ -31,11 +31,33 @@ interface AuthChangeContext { const noopLogger: AuthServiceLogger = { error() { /* no-op */ } }; -export function createModelRuntimeForAgentDir(agentDir: string, allowModelNetwork?: boolean): Promise { - return ModelRuntime.create({ +/** + * Create the shared model runtime with runtime-owned network refreshes disabled. + * + * Upstream `ModelRuntime.reloadConfig()`, `login()`, and `logout()` always refresh + * with `allowNetwork: modelNetworkEnabled` and accept no abort signal. With the + * default (`PI_OFFLINE` unset) a single stalled provider-catalog fetch can block + * those call paths for minutes — and, because pi-web shares one runtime and pi + * coalesces per-provider refreshes, session creation joins the same stalled + * fetch. Forcing `PI_OFFLINE` during construction makes every runtime-driven + * refresh local-only; pi-web performs its own bounded catalog refreshes in the + * background instead (see modelCatalogRefresher.ts). + */ +async function createOfflineModelRuntime(options: CreateModelRuntimeOptions): Promise { + const previous = process.env["PI_OFFLINE"]; + process.env["PI_OFFLINE"] = "1"; + try { + return await ModelRuntime.create(options); + } finally { + if (previous === undefined) delete process.env["PI_OFFLINE"]; + else process.env["PI_OFFLINE"] = previous; + } +} + +export function createModelRuntimeForAgentDir(agentDir: string): Promise { + return createOfflineModelRuntime({ authPath: join(agentDir, "auth.json"), modelsPath: join(agentDir, "models.json"), - ...(allowModelNetwork === undefined ? {} : { allowModelNetwork }), }); } @@ -52,7 +74,7 @@ export class AuthService { } static async create(deps: AuthServiceDependencies = {}): Promise { - const runtime = deps.runtime ?? (deps.agentDir === undefined ? await ModelRuntime.create({}) : await createModelRuntimeForAgentDir(deps.agentDir)); + const runtime = deps.runtime ?? (deps.agentDir === undefined ? await createOfflineModelRuntime({}) : await createModelRuntimeForAgentDir(deps.agentDir)); const logger = deps.logger ?? noopLogger; const authFlows = deps.authFlows ?? new OAuthLoginFlowService({ logger }); return new AuthService(runtime, authFlows, logger); diff --git a/src/server/sessions/modelCatalogRefresher.test.ts b/src/server/sessions/modelCatalogRefresher.test.ts new file mode 100644 index 0000000..c311bb8 --- /dev/null +++ b/src/server/sessions/modelCatalogRefresher.test.ts @@ -0,0 +1,173 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { ModelCatalogRefresher } from "./modelCatalogRefresher.js"; + +interface RefreshCall { + allowNetwork?: boolean; + signal?: AbortSignal; +} + +interface RefreshResult { + aborted: boolean; + errors: Map; +} + +const okResult = (): RefreshResult => ({ aborted: false, errors: new Map() }); + +function deferred() { + let resolveValue: (value: T) => void = () => undefined; + const promise = new Promise((resolve) => { + resolveValue = resolve; + }); + return { promise, resolve: resolveValue }; +} + +function createRuntime() { + const calls: RefreshCall[] = []; + const refresh = vi.fn((options?: RefreshCall) => { + calls.push(options ?? {}); + return Promise.resolve(okResult()); + }); + return { refresh, calls }; +} + +function createLogger() { + const warn = vi.fn(); + const error = vi.fn(); + return { logger: { warn, error }, warn, error }; +} + +/** Let a started refresh run to completion, including the finally-queue bookkeeping. */ +async function flushMicrotasks(rounds = 5): Promise { + for (let index = 0; index < rounds; index++) await Promise.resolve(); +} + +describe("ModelCatalogRefresher", () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it("runs a bounded network refresh when one is requested", async () => { + const runtime = createRuntime(); + const refresher = new ModelCatalogRefresher({ runtime }); + + refresher.requestRefresh(); + await flushMicrotasks(); + + expect(runtime.refresh).toHaveBeenCalledOnce(); + const call = runtime.calls.at(0); + expect(call?.allowNetwork).toBe(true); + expect(call?.signal).toBeInstanceOf(AbortSignal); + refresher.dispose(); + }); + + it("coalesces overlapping requests into a single follow-up run", async () => { + const gate = deferred(); + const refresh = vi.fn() + .mockImplementationOnce(() => gate.promise) + .mockResolvedValue(okResult()); + const refresher = new ModelCatalogRefresher({ runtime: { refresh } }); + + refresher.requestRefresh(); + refresher.requestRefresh(); + refresher.requestRefresh(); + expect(refresh).toHaveBeenCalledOnce(); + + gate.resolve(okResult()); + await flushMicrotasks(); + + expect(refresh).toHaveBeenCalledTimes(2); + refresher.dispose(); + }); + + it("refreshes after the initial delay and then on the interval", async () => { + const runtime = createRuntime(); + const refresher = new ModelCatalogRefresher({ runtime, initialDelayMs: 1_000, intervalMs: 60_000 }); + refresher.start(); + + await vi.advanceTimersByTimeAsync(999); + expect(runtime.refresh).not.toHaveBeenCalled(); + + await vi.advanceTimersByTimeAsync(1); + expect(runtime.refresh).toHaveBeenCalledOnce(); + + await vi.advanceTimersByTimeAsync(60_000); + expect(runtime.refresh).toHaveBeenCalledTimes(2); + refresher.dispose(); + }); + + it("stops scheduling refreshes after dispose", async () => { + const runtime = createRuntime(); + const refresher = new ModelCatalogRefresher({ runtime, initialDelayMs: 1_000, intervalMs: 60_000 }); + refresher.start(); + refresher.dispose(); + + await vi.advanceTimersByTimeAsync(120_000); + expect(runtime.refresh).not.toHaveBeenCalled(); + }); + + it("does not run a queued follow-up after dispose", async () => { + const gate = deferred(); + const refresh = vi.fn().mockImplementation(() => gate.promise); + const refresher = new ModelCatalogRefresher({ runtime: { refresh } }); + + refresher.requestRefresh(); + refresher.requestRefresh(); + refresher.dispose(); + gate.resolve(okResult()); + await flushMicrotasks(); + + expect(refresh).toHaveBeenCalledOnce(); + }); + + it("warns and keeps going when a refresh is aborted by its timeout", async () => { + const { logger, warn, error } = createLogger(); + const refresh = vi.fn(() => Promise.resolve({ aborted: true, errors: new Map() })); + const refresher = new ModelCatalogRefresher({ runtime: { refresh }, logger }); + + refresher.requestRefresh(); + await flushMicrotasks(); + + expect(warn).toHaveBeenCalledOnce(); + expect(error).not.toHaveBeenCalled(); + refresher.dispose(); + }); + + it("warns with provider details when a refresh reports provider errors", async () => { + const { logger, warn } = createLogger(); + const errors = new Map([["openrouter", new Error("boom")]]); + const refresh = vi.fn(() => Promise.resolve({ aborted: false, errors })); + const refresher = new ModelCatalogRefresher({ runtime: { refresh }, logger }); + + refresher.requestRefresh(); + await flushMicrotasks(); + + expect(warn).toHaveBeenCalledWith( + { providers: ["openrouter: boom"] }, + "model catalog refresh failed for some providers; keeping cached catalogs", + ); + refresher.dispose(); + }); + + it("logs and swallows a rejecting refresh so timers stay alive", async () => { + const { logger, error } = createLogger(); + const failure = new Error("refresh exploded"); + const refresh = vi.fn() + .mockRejectedValueOnce(failure) + .mockResolvedValue(okResult()); + const refresher = new ModelCatalogRefresher({ runtime: { refresh }, logger, initialDelayMs: 1_000, intervalMs: 60_000 }); + + refresher.requestRefresh(); + await flushMicrotasks(); + + expect(error).toHaveBeenCalledWith({ err: failure }, "model catalog refresh failed; keeping cached catalogs"); + + refresher.start(); + await vi.advanceTimersByTimeAsync(61_000); + expect(refresh).toHaveBeenCalledTimes(3); + refresher.dispose(); + }); +}); diff --git a/src/server/sessions/modelCatalogRefresher.ts b/src/server/sessions/modelCatalogRefresher.ts new file mode 100644 index 0000000..7d4ad54 --- /dev/null +++ b/src/server/sessions/modelCatalogRefresher.ts @@ -0,0 +1,108 @@ +import type { ModelRuntime } from "@earendil-works/pi-coding-agent"; + +/** + * Matches pi's REMOTE_CATALOG_REFRESH_INTERVAL_MS: provider catalog entries in + * models-store.json are treated as fresh for four hours. + */ +const DEFAULT_INTERVAL_MS = 4 * 60 * 60 * 1000; +/** Give the daemon a moment to finish startup before the first network refresh. */ +const DEFAULT_INITIAL_DELAY_MS = 15_000; +/** Bound every catalog refresh so a stalled provider fetch can never block for minutes. */ +const DEFAULT_TIMEOUT_MS = 15_000; + +/** Minimal structured-logging seam for non-fatal refresh failures. */ +export interface ModelCatalogRefresherLogger { + warn(details: Record, message: string): void; + error(details: Record, message: string): void; +} + +export interface ModelCatalogRefresherOptions { + runtime: Pick; + logger?: ModelCatalogRefresherLogger; + intervalMs?: number; + initialDelayMs?: number; + timeoutMs?: number; +} + +const noopLogger: ModelCatalogRefresherLogger = { + warn() { /* no-op */ }, + error() { /* no-op */ }, +}; + +/** + * Refreshes provider model catalogs over the network on a background schedule. + * + * The shared ModelRuntime is constructed offline (see authService.ts), so its + * own refreshes never touch the network and stay fast on request paths. This + * refresher is the single place that deliberately performs network refreshes — + * bounded by an abort timeout, serialized through one in-flight run, and off + * any request path. `requestRefresh()` additionally asks for a prompt refresh + * after events that change what should be listed, such as provider logins. + */ +export class ModelCatalogRefresher { + private readonly runtime: Pick; + private readonly logger: ModelCatalogRefresherLogger; + private readonly intervalMs: number; + private readonly initialDelayMs: number; + private readonly timeoutMs: number; + private initialTimer?: NodeJS.Timeout; + private intervalTimer?: NodeJS.Timeout; + private inflight: Promise | undefined; + private queued = false; + private disposed = false; + + constructor(options: ModelCatalogRefresherOptions) { + this.runtime = options.runtime; + this.logger = options.logger ?? noopLogger; + this.intervalMs = options.intervalMs ?? DEFAULT_INTERVAL_MS; + this.initialDelayMs = options.initialDelayMs ?? DEFAULT_INITIAL_DELAY_MS; + this.timeoutMs = options.timeoutMs ?? DEFAULT_TIMEOUT_MS; + } + + start(): void { + if (this.disposed) return; + this.initialTimer = setTimeout(() => { this.requestRefresh(); }, this.initialDelayMs); + this.initialTimer.unref(); + this.intervalTimer = setInterval(() => { this.requestRefresh(); }, this.intervalMs); + this.intervalTimer.unref(); + } + + /** Ask for a refresh, coalescing concurrent and overlapping requests. */ + requestRefresh(): void { + if (this.disposed) return; + if (this.inflight !== undefined) { + this.queued = true; + return; + } + const run = this.run(); + this.inflight = run; + this.inflight.finally(() => { + this.inflight = undefined; + if (this.queued && !this.disposed) { + this.queued = false; + this.requestRefresh(); + } + }).catch(() => { /* run() never rejects; finally() re-throws otherwise */ }); + } + + dispose(): void { + this.disposed = true; + if (this.initialTimer !== undefined) clearTimeout(this.initialTimer); + if (this.intervalTimer !== undefined) clearInterval(this.intervalTimer); + } + + private async run(): Promise { + try { + const result = await this.runtime.refresh({ allowNetwork: true, signal: AbortSignal.timeout(this.timeoutMs) }); + if (result.aborted) { + this.logger.warn({ timeoutMs: this.timeoutMs }, "model catalog refresh timed out; keeping cached catalogs"); + } + if (result.errors.size > 0) { + const providers = [...result.errors.entries()].map(([providerId, error]) => `${providerId}: ${error.message}`); + this.logger.warn({ providers }, "model catalog refresh failed for some providers; keeping cached catalogs"); + } + } catch (error: unknown) { + this.logger.error({ err: error }, "model catalog refresh failed; keeping cached catalogs"); + } + } +}