From 4e64fa9940f5a4a9bfb7b021a3c405fd76643ca5 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Mon, 5 Oct 2026 22:27:09 -0700 Subject: [PATCH] fix(native-chat): start a Codex chat without waiting on the background model probe (#23831) * fix(native-chat): start a Codex chat without waiting on the background model probe Opening a Codex chat kicks a session-less model-catalog probe (a throwaway read-only `codex app-server`) for the account, and the new chat's own options read then joined that probe's in-flight refresh through the catalog store's single-flight. When the probe's Codex hung, chat start waited out the probe's 15s deadline and failed with "codex app-server session exceeded 15000ms", even though the chat's own app-server was up. Single-flight now joins only a refresh by the same kind of lister. A live session lists over the connection it already holds; the probe still records its own failure in the store, and the picker keeps serving whichever listing succeeded. * fix(native-chat): never let a chat's model listing join another chat's listing Single-flight in the model catalog store was split by lister kind, so a new chat's acquire-time options read no longer joined the session-less probe, but it still joined any other live chat's in-flight listing for the same account. When that other chat's Codex was wedged or closing, the new chat waited out that chat's request timeout and failed to start with its error. Key in-flight listings by the lister's identity instead: a live session by its own connection, the probe by the probe itself. A lister still joins its own in-flight listing, and shouldRefresh still holds back a probe while any listing for the account is in flight. * test(native-chat): pin that the catalog store releases an account once every listing settles shouldRefresh reads 'any listing in flight for this account' from whether the per-account in-flight map exists, so that map must be dropped exactly when its last listing settles. Nothing covered this: removing the cleanup, or dropping the map on the first settle, passed every catalog test. A leftover map would stop every later background refresh and probe for the account. * fix(native-chat): name the catalog lister type instead of a bare object A live session is keyed by its per-spawn catalog handle (minted in the same acquire as its connection), the probe by itself. * test(codex): pin that a chat's own concurrent option reads share one listing * fix(native-chat): keep model picker current across parallel listings --- ...structured-session-options-catalog.test.ts | 83 ++++++++- .../codex/codex-structured-session-options.ts | 2 +- .../agent-model-catalog-service.test.ts | 149 ++++++++++++++++ .../agent-model-catalog-service.ts | 24 +-- .../agent-model-catalog-store.test.ts | 166 +++++++++++++++++- .../agent-model-catalog-store.ts | 114 +++++++++++- 6 files changed, 514 insertions(+), 24 deletions(-) diff --git a/src/main/codex/codex-structured-session-options-catalog.test.ts b/src/main/codex/codex-structured-session-options-catalog.test.ts index a5f0c1fa69d..b28c25647ff 100644 --- a/src/main/codex/codex-structured-session-options-catalog.test.ts +++ b/src/main/codex/codex-structured-session-options-catalog.test.ts @@ -12,8 +12,11 @@ import type { CodexSession } from './codex-structured-session-state' import { AGENT_MODEL_CATALOG_FRESH_MS, AGENT_MODEL_CATALOG_VALIDATION_MIN_AGE_MS, - AgentModelCatalogStore + AgentModelCatalogStore, + type AgentModelCatalogProbe } from '../native-chat/agent-model-catalog/agent-model-catalog-store' +import { createAgentModelCatalogService } from '../native-chat/agent-model-catalog/agent-model-catalog-service' +import { agentModelCatalogFingerprint } from '../native-chat/agent-model-catalog/agent-model-catalog-fingerprint' const FINGERPRINT = 'fp-session-account' @@ -102,6 +105,84 @@ describe('Codex session options through the host catalog store', () => { expect(store.get('some-other-account')).toBeNull() }) + it('restores a new chat from its own connection while a session-less probe hangs', async () => { + const store = new AgentModelCatalogStore() + // Opening the chat's picker kicked the host probe for this account; its Codex never answers. + const hungProbe: AgentModelCatalogProbe = () => new Promise(() => {}) + void store.refresh(FINGERPRINT, 'codex', hungProbe, () => hungProbe('/homes/a')) + const request = vi.fn(async () => listAnswer('gpt-live')) + const session = storeSession(request, store) + // The acquire-time restore read: joining the probe would fail the chat at the probe's deadline. + const result = await readLiveCodexSessionOptions(session, undefined) + expect(result.models.map((model) => model.id)).toEqual(['gpt-live']) + expect(modelListCalls(request)).toBe(1) + }) + + it('restores a new chat while the probe its opening picker read kicked hangs', async () => { + const store = new AgentModelCatalogStore() + const fingerprint = agentModelCatalogFingerprint({ + agent: 'codex', + accountHomeVariable: 'CODEX_HOME', + accountHomePath: '/homes/a', + wslDistro: null + }) + const hungProbe = vi.fn(() => new Promise(() => {})) + const service = createAgentModelCatalogService({ + store, + getRecord: () => undefined, + resolveAccountHome: async () => ({ variable: 'CODEX_HOME', path: '/homes/a' }), + probes: { codex: hungProbe } + }) + expect(await service.read({ agent: 'codex' })).toEqual({ + origin: 'unknown', + listingInProgress: true + }) + expect(hungProbe).toHaveBeenCalledTimes(1) + const request = vi.fn(async () => listAnswer('gpt-live')) + const session = storeSession(request, store) + session.catalogAccess = { store, fingerprint, accountHomePath: '/homes/a' } + const result = await readLiveCodexSessionOptions(session, undefined) + expect(result.models.map((model) => model.id)).toEqual(['gpt-live']) + // The picker's waiting read now answers from the chat's listing. + const picker = await service.read({ agent: 'codex', waitForListing: true }) + expect(picker.origin === 'unknown' ? null : picker.models[0]!.id).toBe('gpt-live') + }) + + it("restores a new chat from its own connection while another chat's listing hangs", async () => { + const store = new AgentModelCatalogStore() + // Another chat on the same account is mid-listing and its Codex never answers. + const wedged = storeSession( + vi.fn(() => new Promise(() => {})), + store + ) + void readLiveCodexSessionOptions(wedged, undefined) + const request = vi.fn(async () => listAnswer('gpt-live')) + const session = storeSession(request, store) + const result = await readLiveCodexSessionOptions(session, undefined) + expect(result.models.map((model) => model.id)).toEqual(['gpt-live']) + expect(modelListCalls(request)).toBe(1) + }) + + it("shares one listing between a chat's own concurrent reads", async () => { + const store = new AgentModelCatalogStore() + let answer!: () => void + const answered = new Promise((resolve) => (answer = resolve)) + const request = vi.fn(async () => { + await answered + return listAnswer('gpt-live') + }) + const session = storeSession(request, store) + const reads = [ + readLiveCodexSessionOptions(session, undefined), + readLiveCodexSessionOptions(session, undefined) + ] + answer() + for (const result of await Promise.all(reads)) { + expect(result.models.map((model) => model.id)).toEqual(['gpt-live']) + } + expect(modelListCalls(request)).toBe(1) + }) + it('answers the picker with zero provider fetches when the store is already warm', async () => { const store = new AgentModelCatalogStore() seedEntry(store, 'gpt-live', 'gpt-next') diff --git a/src/main/codex/codex-structured-session-options.ts b/src/main/codex/codex-structured-session-options.ts index d7e87412f12..3002be7b673 100644 --- a/src/main/codex/codex-structured-session-options.ts +++ b/src/main/codex/codex-structured-session-options.ts @@ -52,7 +52,7 @@ async function fetchCodexListingThroughStore( if (!access) { return fetchCodexModelCatalogListing({ connection: session.connection, timeoutMs }) } - const entry = await access.store.refresh(access.fingerprint, 'codex', async () => { + const entry = await access.store.refresh(access.fingerprint, 'codex', access, async () => { const listing = await fetchCodexModelCatalogListing({ connection: session.connection, timeoutMs diff --git a/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.test.ts b/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.test.ts index 02025485151..ee0b5550da7 100644 --- a/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.test.ts +++ b/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.test.ts @@ -210,6 +210,155 @@ describe('agent model catalog service', () => { expect(probe).toHaveBeenCalledTimes(1) }) + it("answers from a chat's listing already running instead of starting a probe", async () => { + const pending = deferredListing() + const probe = vi.fn(() => new Promise(() => {})) + const { store, service } = coldService(probe) + const chat = { + store, + fingerprint: selectedHomeFingerprint('/homes/selected'), + accountHomePath: '/homes/selected' + } + void store.refresh(chat.fingerprint, 'codex', chat, () => pending.promise) + const waited = service.read({ agent: 'codex', waitForListing: true }) + pending.resolve(listing('gpt-chat')) + const result = await waited + expect(result.origin === 'unknown' ? null : result.models[0]!.id).toBe('gpt-chat') + expect(probe).not.toHaveBeenCalled() + }) + + it('uses a chat listing that starts after the picker began waiting on a probe', async () => { + const pendingProbe = deferredListing() + const pendingChat = deferredListing() + const { store, service } = coldService(() => pendingProbe.promise) + expect(await service.read({ agent: 'codex' })).toEqual({ + origin: 'unknown', + listingInProgress: true + }) + const waited = service.read({ agent: 'codex', waitForListing: true }) + await Promise.resolve() + const fingerprint = selectedHomeFingerprint('/homes/selected') + const chat = { store, fingerprint, accountHomePath: '/homes/selected' } + const chatListing = store.refresh(fingerprint, 'codex', chat, () => pendingChat.promise) + pendingChat.resolve(listing('gpt-chat')) + await chatListing + const result = await waited + expect(result.origin === 'unknown' ? null : result.models[0]!.id).toBe('gpt-chat') + pendingProbe.reject(new Error('probe timed out')) + }) + + it('continues waiting when the probe fails before a newly started chat finishes', async () => { + const pendingProbe = deferredListing() + const pendingChat = deferredListing() + const { store, service } = coldService(() => pendingProbe.promise) + await service.read({ agent: 'codex' }) + const waited = service.read({ agent: 'codex', waitForListing: true }) + await Promise.resolve() + const fingerprint = selectedHomeFingerprint('/homes/selected') + const chat = { store, fingerprint, accountHomePath: '/homes/selected' } + const chatListing = store.refresh(fingerprint, 'codex', chat, () => pendingChat.promise) + pendingProbe.reject(new Error('probe timed out')) + await vi.waitFor(() => expect(store.hasActiveFailure(fingerprint)).toBe(true)) + let completed = false + void waited.then(() => (completed = true)) + await Promise.resolve() + expect(completed).toBe(false) + + pendingChat.resolve(listing('gpt-chat')) + await chatListing + const result = await waited + expect(result.origin === 'unknown' ? null : result.models[0]!.id).toBe('gpt-chat') + }) + + it('releases when a second chat succeeds while the first chat is still listing', async () => { + const pendingProbe = deferredListing() + const firstChat = deferredListing() + const secondChat = deferredListing() + const { store, service } = coldService(() => pendingProbe.promise) + await service.read({ agent: 'codex' }) + const waited = service.read({ agent: 'codex', waitForListing: true }) + await Promise.resolve() + const fingerprint = selectedHomeFingerprint('/homes/selected') + const first = store.refresh( + fingerprint, + 'codex', + { store, fingerprint, accountHomePath: '/homes/selected' }, + () => firstChat.promise + ) + pendingProbe.reject(new Error('probe timed out')) + await vi.waitFor(() => expect(store.hasActiveFailure(fingerprint)).toBe(true)) + const second = store.refresh( + fingerprint, + 'codex', + { store, fingerprint, accountHomePath: '/homes/selected' }, + () => secondChat.promise + ) + secondChat.resolve(listing('gpt-second')) + await second + const result = await waited + expect(result.origin === 'unknown' ? null : result.models[0]!.id).toBe('gpt-second') + firstChat.reject(new Error('first chat timed out')) + await first + }) + + it('keeps waiting for a later chat after the first chat also fails', async () => { + const pendingProbe = deferredListing() + const firstChat = deferredListing() + const secondChat = deferredListing() + const { store, service } = coldService(() => pendingProbe.promise) + await service.read({ agent: 'codex' }) + const waited = service.read({ agent: 'codex', waitForListing: true }) + await Promise.resolve() + const fingerprint = selectedHomeFingerprint('/homes/selected') + const first = store.refresh( + fingerprint, + 'codex', + { store, fingerprint, accountHomePath: '/homes/selected' }, + () => firstChat.promise + ) + pendingProbe.reject(new Error('probe timed out')) + await vi.waitFor(() => expect(store.hasActiveFailure(fingerprint)).toBe(true)) + const second = store.refresh( + fingerprint, + 'codex', + { store, fingerprint, accountHomePath: '/homes/selected' }, + () => secondChat.promise + ) + firstChat.reject(new Error('first chat timed out')) + await first + let completed = false + void waited.then(() => (completed = true)) + await Promise.resolve() + expect(completed).toBe(false) + + secondChat.resolve(listing('gpt-second')) + await second + const result = await waited + expect(result.origin === 'unknown' ? null : result.models[0]!.id).toBe('gpt-second') + }) + + it('waits for a running chat even while a failed probe is inside its TTL', async () => { + const pendingProbe = deferredListing() + const pendingChat = deferredListing() + const { store, service } = coldService(() => pendingProbe.promise) + await service.read({ agent: 'codex' }) + const fingerprint = selectedHomeFingerprint('/homes/selected') + const chat = { store, fingerprint, accountHomePath: '/homes/selected' } + const chatListing = store.refresh(fingerprint, 'codex', chat, () => pendingChat.promise) + pendingProbe.reject(new Error('probe timed out')) + await vi.waitFor(() => expect(store.hasActiveFailure(fingerprint)).toBe(true)) + + const waited = service.read({ agent: 'codex', waitForListing: true }) + let completed = false + void waited.then(() => (completed = true)) + await Promise.resolve() + expect(completed).toBe(false) + pendingChat.resolve(listing('gpt-chat')) + await chatListing + const result = await waited + expect(result.origin === 'unknown' ? null : result.models[0]!.id).toBe('gpt-chat') + }) + it('answers a plain unknown when the listing fails', async () => { const pending = deferredListing() const { service } = coldService(() => pending.promise) diff --git a/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.ts b/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.ts index a240cfa3e6d..873f652be0b 100644 --- a/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.ts +++ b/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.ts @@ -84,8 +84,8 @@ async function workspaceKeepsListedDefault( * never "whichever account listed last". `unknown` tells the client to keep * its static seed, and a missing or aged entry kicks one joined background * probe so the next read is warm. With no entry, the answer says that listing - * is running, and only a read that asks waits for it. Failures are the store's - * 30s TTL, never an answer: inside it a read answers `unknown` at once. + * is running, and only a read that asks waits for it. Failures suppress a new + * probe for 30s, but never hide another listing already running for the account. */ export function createAgentModelCatalogService( deps: AgentModelCatalogServiceDeps @@ -118,13 +118,16 @@ export function createAgentModelCatalogService( let entry = deps.store.get(fingerprint) const probe = deps.probes?.[params.agent] const home = accountHomePath - // Without an entry, join a running listing too: that is the one a waiting read answers from. - const listing = - probe && - home && - (entry ? deps.store.shouldRefresh(fingerprint) : !deps.store.hasActiveFailure(fingerprint)) - ? deps.store.refresh(fingerprint, params.agent, () => probe(home)) - : null + // Without an entry, answer from any running listing instead of starting a second one. + let listing = !entry && home ? deps.store.pendingListing(fingerprint) : null + if (probe && home) { + if (entry && deps.store.shouldRefresh(fingerprint)) { + void deps.store.refresh(fingerprint, params.agent, probe, () => probe(home)) + } else if (!entry && !listing && !deps.store.hasActiveFailure(fingerprint)) { + void deps.store.refresh(fingerprint, params.agent, probe, () => probe(home)) + listing = deps.store.pendingListing(fingerprint) + } + } if (!entry) { if (!listing) { return { origin: 'unknown' } @@ -132,7 +135,8 @@ export function createAgentModelCatalogService( if (!params.waitForListing) { return { origin: 'unknown', listingInProgress: true } } - entry = await listing + const listed = await listing + entry = deps.store.get(fingerprint) ?? listed if (!entry) { return { origin: 'unknown' } } diff --git a/src/main/native-chat/agent-model-catalog/agent-model-catalog-store.test.ts b/src/main/native-chat/agent-model-catalog/agent-model-catalog-store.test.ts index 34e0294b458..09bc3f977d9 100644 --- a/src/main/native-chat/agent-model-catalog/agent-model-catalog-store.test.ts +++ b/src/main/native-chat/agent-model-catalog/agent-model-catalog-store.test.ts @@ -13,7 +13,11 @@ import { createAgentModelCatalogFilePersistence } from './agent-model-catalog-pe import { AGENT_MODEL_CATALOG_FAILURE_TTL_MS, AGENT_MODEL_CATALOG_FRESH_MS, + AGENT_MODEL_CATALOG_MAX_ENTRIES, + AGENT_MODEL_CATALOG_PICKER_WAIT_MS, AgentModelCatalogStore, + type AgentModelCatalogProbe, + type AgentModelCatalogSessionAccess, type AgentModelCatalogSuccess } from './agent-model-catalog-store' @@ -34,6 +38,11 @@ function success(...ids: string[]): AgentModelCatalogSuccess { } } +/** A live session's per-spawn handle; each call is a distinct lister. */ +function liveLister(store: AgentModelCatalogStore): AgentModelCatalogSessionAccess { + return { store, fingerprint: 'fp-1', accountHomePath: '/homes/a' } +} + describe('agent model catalog store', () => { it('serves an entry at any age and flags staleness at the refresh threshold', () => { let at = 1_000 @@ -79,8 +88,9 @@ describe('agent model catalog store', () => { const fetch = vi.fn( () => new Promise((resolve) => (settle = resolve)) ) - const first = store.refresh('fp-1', 'codex', fetch) - const second = store.refresh('fp-1', 'codex', fetch) + const session = liveLister(store) + const first = store.refresh('fp-1', 'codex', session, fetch) + const second = store.refresh('fp-1', 'codex', session, fetch) expect(fetch).toHaveBeenCalledTimes(1) settle(success('gpt-a')) const [entryA, entryB] = await Promise.all([first, second]) @@ -88,9 +98,159 @@ describe('agent model catalog store', () => { expect(entryA!.models[0]!.id).toBe('gpt-a') }) + it('never makes a live session wait on another lister that hangs', async () => { + const store = new AgentModelCatalogStore() + let failProbe!: (error: Error) => void + const hungProbe: AgentModelCatalogProbe = () => + new Promise((_resolve, reject) => (failProbe = reject)) + const probe = store.refresh('fp-1', 'codex', hungProbe, () => hungProbe('/homes/a')) + expect(store.shouldRefresh('fp-1')).toBe(false) + + const live = await store.refresh('fp-1', 'codex', liveLister(store), async () => + success('gpt-live') + ) + expect(live!.models[0]!.id).toBe('gpt-live') + + // The probe still reports its own failure; the live listing it lost to stays served. + failProbe(new Error('codex app-server session exceeded 15000ms')) + expect(await probe).toBeNull() + expect(store.failureDetail('fp-1')).toBe('codex app-server session exceeded 15000ms') + expect(store.get('fp-1')!.models[0]!.id).toBe('gpt-live') + }) + + it('answers a pending read with the first listing that succeeds, or null once all fail', async () => { + const store = new AgentModelCatalogStore() + expect(store.pendingListing('fp-1')).toBeNull() + let failFirst!: (error: Error) => void + let settleSecond!: (success: AgentModelCatalogSuccess) => void + void store.refresh( + 'fp-1', + 'codex', + liveLister(store), + () => new Promise((_resolve, reject) => (failFirst = reject)) + ) + void store.refresh( + 'fp-1', + 'codex', + liveLister(store), + () => new Promise((resolve) => (settleSecond = resolve)) + ) + const pending = store.pendingListing('fp-1') + failFirst(new Error('stuck')) + settleSecond(success('gpt-second')) + expect((await pending)!.models[0]!.id).toBe('gpt-second') + + void store.refresh('fp-2', 'codex', liveLister(store), async () => { + throw new Error('no provider') + }) + expect(await store.pendingListing('fp-2')).toBeNull() + }) + + it('ends a picker wait at its deadline even while a listing remains active', async () => { + vi.useFakeTimers() + try { + const store = new AgentModelCatalogStore() + let settle!: (success: AgentModelCatalogSuccess) => void + const listing = store.refresh( + 'fp-1', + 'codex', + liveLister(store), + () => new Promise((resolve) => (settle = resolve)) + ) + const waited = store.pendingListing('fp-1') + await vi.advanceTimersByTimeAsync(AGENT_MODEL_CATALOG_PICKER_WAIT_MS) + expect(await waited).toBeNull() + settle(success('gpt-late')) + expect((await listing)!.models[0]!.id).toBe('gpt-late') + } finally { + vi.useRealTimers() + } + }) + + it('holds back a probe until every lister settles, then lets the account refresh again', async () => { + let at = 1_000 + const store = new AgentModelCatalogStore({ now: () => at }) + let settleSlow!: (success: AgentModelCatalogSuccess) => void + const slow = store.refresh( + 'fp-1', + 'codex', + liveLister(store), + () => new Promise((resolve) => (settleSlow = resolve)) + ) + await store.refresh('fp-1', 'codex', liveLister(store), async () => success('gpt-fast')) + at += AGENT_MODEL_CATALOG_FRESH_MS + expect(store.shouldRefresh('fp-1')).toBe(false) + + settleSlow(success('gpt-slow')) + await slow + at += AGENT_MODEL_CATALOG_FRESH_MS + // A leftover in-flight record here would suppress every later refresh for the account. + expect(store.shouldRefresh('fp-1')).toBe(true) + }) + + it('keeps the newer completed listing when an older chat finishes later', async () => { + const store = new AgentModelCatalogStore() + const save = vi.fn() + await store.attachPersistence({ load: async () => [], save, flush: async () => {} }) + let settleOlder!: (success: AgentModelCatalogSuccess) => void + const older = store.refresh( + 'fp-1', + 'codex', + liveLister(store), + () => new Promise((resolve) => (settleOlder = resolve)) + ) + const newer = await store.refresh('fp-1', 'codex', liveLister(store), async () => + success('gpt-new') + ) + settleOlder(success('gpt-old')) + const olderResult = await older + + expect(newer!.models[0]!.id).toBe('gpt-new') + expect(olderResult!.models[0]!.id).toBe('gpt-old') + expect(store.get('fp-1')!.models[0]!.id).toBe('gpt-new') + expect(save).toHaveBeenCalledTimes(1) + }) + + it('keeps a direct live update ahead of a pending older probe', async () => { + const store = new AgentModelCatalogStore() + let settleProbe!: (success: AgentModelCatalogSuccess) => void + const probe: AgentModelCatalogProbe = () => + new Promise((resolve) => (settleProbe = resolve)) + const pending = store.refresh('fp-1', 'codex', probe, () => probe('/homes/a')) + store.recordSuccess('fp-1', 'codex', success('gpt-live')) + settleProbe({ ...success('gpt-probe'), origin: 'probe' }) + expect((await pending)!.models[0]!.id).toBe('gpt-probe') + expect(store.get('fp-1')!.models[0]!.id).toBe('gpt-live') + }) + + it('keeps an older successful listing when the newer entry was evicted', async () => { + const store = new AgentModelCatalogStore() + let settleOlder!: (success: AgentModelCatalogSuccess) => void + const older = store.refresh( + 'fp-1', + 'codex', + liveLister(store), + () => new Promise((resolve) => (settleOlder = resolve)) + ) + await store.refresh('fp-1', 'codex', liveLister(store), async () => success('gpt-new')) + for (let index = 0; index < AGENT_MODEL_CATALOG_MAX_ENTRIES; index++) { + store.recordSuccess(`other-${index}`, 'codex', success('other')) + } + expect(store.get('fp-1')).toBeNull() + const save = vi.fn() + await store.attachPersistence({ load: async () => [], save, flush: async () => {} }) + + settleOlder(success('gpt-old')) + expect((await older)!.models[0]!.id).toBe('gpt-old') + expect(store.get('fp-1')!.models[0]!.id).toBe('gpt-old') + expect(save).toHaveBeenCalledWith( + expect.arrayContaining([expect.objectContaining({ fingerprint: 'fp-1' })]) + ) + }) + it('records a failed refresh as a failure and resolves null without rejecting', async () => { const store = new AgentModelCatalogStore() - const entry = await store.refresh('fp-1', 'codex', async () => { + const entry = await store.refresh('fp-1', 'codex', liveLister(store), async () => { throw new Error('no provider') }) expect(entry).toBeNull() diff --git a/src/main/native-chat/agent-model-catalog/agent-model-catalog-store.ts b/src/main/native-chat/agent-model-catalog/agent-model-catalog-store.ts index dcb6e3935ac..3e416d1cba3 100644 --- a/src/main/native-chat/agent-model-catalog/agent-model-catalog-store.ts +++ b/src/main/native-chat/agent-model-catalog/agent-model-catalog-store.ts @@ -13,6 +13,7 @@ import type { AgentModelCatalogPersistence } from './agent-model-catalog-persist export const AGENT_MODEL_CATALOG_FRESH_MS = 10 * 60_000 export const AGENT_MODEL_CATALOG_FAILURE_TTL_MS = 30_000 +export const AGENT_MODEL_CATALOG_PICKER_WAIT_MS = 30_000 /** A validation read younger than this trusts the entry even when the picked * model is missing; older, it waits for one bounded refresh before refusing. */ export const AGENT_MODEL_CATALOG_VALIDATION_MIN_AGE_MS = 60_000 @@ -38,8 +39,13 @@ export type AgentModelCatalogSuccess = { export type AgentModelCatalogProbe = (accountHomePath: string) => Promise +/** Who lists, by identity: a live session's per-spawn handle, or the session-less probe. */ +export type AgentModelCatalogLister = AgentModelCatalogSessionAccess | AgentModelCatalogProbe + type CatalogFailure = { detail: string; failedAt: number } +type InFlightListings = Map> + /** A live session's handle into the store, pinned at spawn to the account home * THAT child launched under — an account switched afterwards must never * receive or poison this session's listing. */ @@ -81,7 +87,10 @@ function listingKey(entry: AgentModelCatalogEntry): string { export class AgentModelCatalogStore { private readonly entries = new Map() private readonly failures = new Map() - private readonly refreshes = new Map>() + private readonly refreshes = new Map() + private readonly listingWaiters = new Map void>>() + private readonly latestWrittenOrder = new Map() + private nextListingOrder = 0 private persistence: AgentModelCatalogPersistence | null = null private readonly now: () => number @@ -146,13 +155,23 @@ export class AgentModelCatalogStore { fingerprint: string, agent: 'claude' | 'codex', success: AgentModelCatalogSuccess + ): AgentModelCatalogEntry | null { + const entry = this.writeSuccess(fingerprint, agent, success, ++this.nextListingOrder) + this.notifyListingWaiters(fingerprint) + return entry + } + + private entryFromSuccess( + fingerprint: string, + agent: 'claude' | 'codex', + success: AgentModelCatalogSuccess ): AgentModelCatalogEntry | null { if (success.models.length === 0) { // An empty list identifies no model; it is doubt, not a catalog. return null } const previous = this.entries.get(fingerprint) - const entry: AgentModelCatalogEntry = { + return { agent, fingerprint, models: withKnownDefaultEfforts(success.models, previous), @@ -161,8 +180,24 @@ export class AgentModelCatalogStore { origin: success.origin, fetchedAt: this.now() } + } + + private writeSuccess( + fingerprint: string, + agent: 'claude' | 'codex', + success: AgentModelCatalogSuccess, + order: number + ): AgentModelCatalogEntry | null { + const entry = this.entryFromSuccess(fingerprint, agent, success) + if (!entry) { + return null + } + const previous = this.entries.get(fingerprint) this.entries.delete(fingerprint) this.entries.set(fingerprint, entry) + if (this.refreshes.has(fingerprint)) { + this.latestWrittenOrder.set(fingerprint, order) + } this.failures.delete(fingerprint) this.evictOverCap() // Live sessions re-list every turn; an unchanged listing only refreshes the in-memory age. @@ -176,32 +211,93 @@ export class AgentModelCatalogStore { this.failures.set(fingerprint, { detail, failedAt: this.now() }) } - /** Joins an in-flight refresh for the key rather than starting a second. - * Resolves with the entry on success and null on failure — never rejects. */ + /** Joins an in-flight refresh by the same lister rather than starting a second. Never + * joins another lister's: a probe or another chat's Codex that hangs must not decide + * whether this chat starts. Resolves with the entry on success, null on failure. */ refresh( fingerprint: string, agent: 'claude' | 'codex', + lister: AgentModelCatalogLister, listModels: () => Promise ): Promise { - const inFlight = this.refreshes.get(fingerprint) + const listers: InFlightListings = this.refreshes.get(fingerprint) ?? new Map() + const inFlight = listers.get(lister) if (inFlight) { return inFlight } + const settle = (): void => { + listers.delete(lister) + if (listers.size === 0 && this.refreshes.get(fingerprint) === listers) { + this.refreshes.delete(fingerprint) + this.latestWrittenOrder.delete(fingerprint) + } + this.notifyListingWaiters(fingerprint) + } + const order = ++this.nextListingOrder const run = listModels().then( (success) => { - this.refreshes.delete(fingerprint) - return this.recordSuccess(fingerprint, agent, success) + // An older session still receives its own result, but cannot replace a newer catalog. + const entry = + (this.latestWrittenOrder.get(fingerprint) ?? 0) > order && this.entries.has(fingerprint) + ? this.entryFromSuccess(fingerprint, agent, success) + : this.writeSuccess(fingerprint, agent, success, order) + settle() + return entry }, (error: unknown) => { - this.refreshes.delete(fingerprint) + settle() this.recordFailure(fingerprint, error instanceof Error ? error.message : String(error)) return null } ) - this.refreshes.set(fingerprint, run) + listers.set(lister, run) + this.refreshes.set(fingerprint, listers) return run } + /** A picker follows the current account work until a catalog lands, all work ends, + * or its fixed deadline expires. */ + pendingListing(fingerprint: string): Promise | null { + if (!this.refreshes.has(fingerprint)) { + return null + } + return new Promise((resolve) => { + const waiters = this.listingWaiters.get(fingerprint) ?? new Set<() => void>() + let settled = false + const finish = (entry: AgentModelCatalogEntry | null): void => { + if (settled) { + return + } + settled = true + clearTimeout(deadline) + waiters.delete(check) + if (waiters.size === 0) { + this.listingWaiters.delete(fingerprint) + } + resolve(entry) + } + const check = (): void => { + const entry = this.get(fingerprint) + if (entry || !this.refreshes.has(fingerprint)) { + finish(entry) + } + } + const deadline = setTimeout( + () => finish(this.get(fingerprint)), + AGENT_MODEL_CATALOG_PICKER_WAIT_MS + ) + waiters.add(check) + this.listingWaiters.set(fingerprint, waiters) + check() + }) + } + + private notifyListingWaiters(fingerprint: string): void { + for (const check of this.listingWaiters.get(fingerprint) ?? []) { + check() + } + } + /** True when a read should kick a background refresh: nothing known or the * entry aged out, and no failure is still inside its TTL. */ shouldRefresh(fingerprint: string): boolean {