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 {