diff --git a/src/app-server/query/cache.ts b/src/app-server/query/cache.ts index be44c5bf..17955f5f 100644 --- a/src/app-server/query/cache.ts +++ b/src/app-server/query/cache.ts @@ -32,6 +32,7 @@ const THREAD_LIST_STALE_TIME_MS = 10_000; const APP_SERVER_METADATA_STALE_TIME_MS = 10_000; const MODELS_STALE_TIME_MS = 60_000; const APP_SERVER_QUERY_GC_TIME_MS = 5 * 60_000; +const FULL_ACTIVE_THREAD_FETCH_ATTEMPTS = 2; export interface AppServerQueryClientRunner { runWithClient( @@ -48,12 +49,16 @@ interface AppServerQueryOptions { } type ThreadListKind = "active" | "archived"; +export type MetadataResourceKind = "skills" | "rateLimits"; type ThreadListUpdater = (threads: readonly Thread[] | null) => readonly Thread[] | null; export class AppServerQueryCache { private readonly client: QueryClient; private readonly clientRunner: AppServerQueryClientRunner | null; private readonly activeThreadCursors = new Map(); + private readonly activeThreadRevisions = new Map(); + private readonly metadataRevisions = new Map(); + private readonly metadataWriteRevisions = new Map(); constructor(options: { client?: QueryClient; clientRunner?: AppServerQueryClientRunner } = {}) { this.client = options.client ?? createAppServerQueryClient(); @@ -62,6 +67,9 @@ export class AppServerQueryCache { clear(): void { this.activeThreadCursors.clear(); + this.activeThreadRevisions.clear(); + this.metadataRevisions.clear(); + this.metadataWriteRevisions.clear(); this.client.clear(); } @@ -107,10 +115,15 @@ export class AppServerQueryCache { const cursorKey = this.activeThreadCursorKey(refreshContext); const snapshot = this.activeThreadsSnapshot(refreshContext); if (snapshot && this.activeThreadCursors.has(cursorKey) && !this.activeThreadCursors.get(cursorKey)) return snapshot; - const threads = await this.runWithClient(refreshContext, (client) => listThreads(client, refreshContext.vaultPath)); - this.setActiveThreads(refreshContext, threads); - this.rememberActiveThreadCursor(refreshContext, null); - return cloneThreads(threads); + for (let attempt = 0; attempt < FULL_ACTIVE_THREAD_FETCH_ATTEMPTS; attempt += 1) { + const revision = this.activeThreadRevision(refreshContext); + const threads = await this.runWithClient(refreshContext, (client) => listThreads(client, refreshContext.vaultPath)); + if (this.activeThreadRevision(refreshContext) !== revision) continue; + this.setActiveThreads(refreshContext, threads); + this.rememberActiveThreadCursor(refreshContext, null); + return cloneThreads(threads); + } + throw new Error("Active thread inventory changed while it was being fetched."); } hasMoreActiveThreads(context: AppServerQueryContext): boolean { @@ -124,10 +137,12 @@ export class AppServerQueryCache { const current = this.activeThreadsSnapshot(refreshContext) ?? (await this.fetchActiveThreads(refreshContext)); const cursor = this.activeThreadCursors.get(this.activeThreadCursorKey(refreshContext)) ?? null; if (!cursor) return current; + const revision = this.activeThreadRevision(refreshContext); const page = await this.runWithClient(refreshContext, (client) => readThreadPage(client, refreshContext.vaultPath, { cursor, archived: false }), ); if (page.nextCursor === cursor) throw new Error("Codex app-server returned a repeated thread list cursor."); + if (this.activeThreadRevision(refreshContext) !== revision) return this.activeThreadsSnapshot(refreshContext) ?? current; const latest = this.activeThreadsSnapshot(refreshContext) ?? current; const existingIds = new Set(latest.map((thread) => thread.id)); const threads = [...latest, ...page.threads.filter((thread) => !existingIds.has(thread.id))]; @@ -180,6 +195,7 @@ export class AppServerQueryCache { private setThreadList(context: AppServerQueryContext, kind: ThreadListKind, threads: readonly Thread[]): void { if (!appServerQueryContextIsComplete(context)) return; this.client.setQueryData(this.threadListQueryKey(context, kind), cloneThreads(threads)); + if (kind === "active") this.bumpActiveThreadRevision(context); } private updateThreadList(context: AppServerQueryContext, kind: ThreadListKind, updater: ThreadListUpdater): readonly Thread[] | null { @@ -188,6 +204,7 @@ export class AppServerQueryCache { const next = updater(current); if (!next) return null; this.client.setQueryData(this.threadListQueryKey(context, kind), cloneThreads(next), current ? undefined : { updatedAt: 0 }); + if (kind === "active") this.bumpActiveThreadRevision(context); return cloneThreads(next); } @@ -213,6 +230,8 @@ export class AppServerQueryCache { if (!appServerQueryContextIsComplete(refreshContext)) { return null; } + this.beginMetadataResourceRefresh(refreshContext, "skills"); + this.beginMetadataResourceRefresh(refreshContext, "rateLimits"); const key = appServerMetadataQueryKey(refreshContext); await Promise.all([ this.client.invalidateQueries({ queryKey: key }), @@ -226,16 +245,43 @@ export class AppServerQueryCache { if (!appServerQueryContextIsComplete(context)) return null; const next = metadataWithLastKnownGood(metadata, this.appServerMetadataSnapshot(context)); this.client.setQueryData(appServerMetadataQueryKey(context), cloneSharedServerMetadata(next)); + this.bumpMetadataRevision(context); + this.bumpMetadataWriteRevision(context); return cloneSharedServerMetadata(next); } updateAppServerMetadata( context: AppServerQueryContext, updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null, + resource?: MetadataResourceKind, ): SharedServerMetadata | null { if (!appServerQueryContextIsComplete(context)) return null; const next = updater(this.appServerMetadataSnapshot(context)); - return next ? this.writeAppServerMetadata(context, next) : null; + if (!next) return null; + const merged = metadataWithLastKnownGood(next, this.appServerMetadataSnapshot(context)); + this.client.setQueryData(appServerMetadataQueryKey(context), cloneSharedServerMetadata(merged)); + if (resource) this.beginMetadataResourceRefresh(context, resource); + else this.bumpMetadataRevision(context); + this.bumpMetadataWriteRevision(context); + return cloneSharedServerMetadata(merged); + } + + beginMetadataResourceRefresh(context: AppServerQueryContext, resource: MetadataResourceKind): number { + const key = this.metadataResourceRevisionKey(context, resource); + const revision = (this.metadataRevisions.get(key) ?? 0) + 1; + this.metadataRevisions.delete(key); + this.metadataRevisions.set(key, revision); + while (this.metadataRevisions.size > 16) { + for (const oldestKey of this.metadataRevisions.keys()) { + this.metadataRevisions.delete(oldestKey); + break; + } + } + return revision; + } + + metadataResourceRefreshIsCurrent(context: AppServerQueryContext, resource: MetadataResourceKind, revision: number): boolean { + return this.metadataRevisions.get(this.metadataResourceRevisionKey(context, resource)) === revision; } modelsSnapshot(context: AppServerQueryContext): readonly ModelMetadata[] | null { @@ -273,9 +319,11 @@ export class AppServerQueryCache { queryKey: this.threadListQueryKey(refreshContext, kind), queryFn: async (): Promise => { if (kind === "active") { + const revision = this.activeThreadRevision(refreshContext); const page = await this.runWithClient(refreshContext, (client) => readThreadPage(client, refreshContext.vaultPath, { archived: false }), ); + if (this.activeThreadRevision(refreshContext) !== revision) return this.activeThreadsSnapshot(refreshContext) ?? []; this.rememberActiveThreadCursor(refreshContext, page.nextCursor); return cloneThreads(page.threads); } @@ -303,7 +351,8 @@ export class AppServerQueryCache { queryKey: appServerMetadataQueryKey(refreshContext), queryFn: async (): Promise => { const previous = this.appServerMetadataSnapshot(refreshContext); - return this.runWithClient(refreshContext, async (client) => { + const revision = this.metadataRevision(refreshContext); + const metadata = await this.runWithClient(refreshContext, async (client) => { const runtimeConfig = runtimeConfigSnapshotFromAppServerConfig(await readEffectiveConfig(client, refreshContext.vaultPath)); const [modelProbe, skills, permissionProfiles, rateLimit] = await Promise.all([ this.readModelMetadataProbe(refreshContext, client), @@ -326,6 +375,7 @@ export class AppServerQueryCache { previous, ); }); + return this.metadataRevision(refreshContext) === revision ? metadata : (this.appServerMetadataSnapshot(refreshContext) ?? metadata); }, staleTime: APP_SERVER_METADATA_STALE_TIME_MS, }; @@ -414,13 +464,54 @@ export class AppServerQueryCache { const key = this.activeThreadCursorKey(context); this.activeThreadCursors.delete(key); this.activeThreadCursors.set(key, cursor); + this.bumpActiveThreadRevision(context); while (this.activeThreadCursors.size > 8) { for (const oldestKey of this.activeThreadCursors.keys()) { this.activeThreadCursors.delete(oldestKey); + this.activeThreadRevisions.delete(oldestKey); break; } } } + + private activeThreadRevision(context: AppServerQueryContext): number { + return this.activeThreadRevisions.get(this.activeThreadCursorKey(context)) ?? 0; + } + + private bumpActiveThreadRevision(context: AppServerQueryContext): void { + const key = this.activeThreadCursorKey(context); + this.activeThreadRevisions.set(key, this.activeThreadRevision(context) + 1); + } + + private metadataRevision(context: AppServerQueryContext): number { + return this.metadataWriteRevisions.get(this.metadataWriteRevisionKey(context)) ?? 0; + } + + private bumpMetadataRevision(context: AppServerQueryContext): void { + this.beginMetadataResourceRefresh(context, "skills"); + this.beginMetadataResourceRefresh(context, "rateLimits"); + } + + private bumpMetadataWriteRevision(context: AppServerQueryContext): void { + const key = this.metadataWriteRevisionKey(context); + const revision = this.metadataRevision(context) + 1; + this.metadataWriteRevisions.delete(key); + this.metadataWriteRevisions.set(key, revision); + while (this.metadataWriteRevisions.size > 8) { + for (const oldestKey of this.metadataWriteRevisions.keys()) { + this.metadataWriteRevisions.delete(oldestKey); + break; + } + } + } + + private metadataWriteRevisionKey(context: AppServerQueryContext): string { + return JSON.stringify(appServerMetadataQueryKey(context)); + } + + private metadataResourceRevisionKey(context: AppServerQueryContext, resource: MetadataResourceKind): string { + return JSON.stringify([...appServerMetadataQueryKey(context), resource]); + } } function metadataWithLastKnownGood(metadata: SharedServerMetadata, previous: SharedServerMetadata | null): SharedServerMetadata { diff --git a/src/app-server/query/shared-queries.ts b/src/app-server/query/shared-queries.ts index f9476a99..1dedbff1 100644 --- a/src/app-server/query/shared-queries.ts +++ b/src/app-server/query/shared-queries.ts @@ -1,7 +1,7 @@ import type { ModelMetadata } from "../../domain/catalog/metadata"; import type { SharedServerMetadata } from "../../domain/server/metadata"; import type { Thread } from "../../domain/threads/model"; -import type { AppServerQueryCache } from "./cache"; +import type { AppServerQueryCache, MetadataResourceKind } from "./cache"; import { type AppServerQueryContext, appServerQueryContextMatches, @@ -101,8 +101,19 @@ export class AppServerSharedQueries { return this.options.cache.appServerMetadataSnapshot(this.context()); } - updateAppServerMetadata(updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null): SharedServerMetadata | null { - return this.options.cache.updateAppServerMetadata(this.context(), updater); + updateAppServerMetadata( + updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null, + resource?: MetadataResourceKind, + ): SharedServerMetadata | null { + return this.options.cache.updateAppServerMetadata(this.context(), updater, resource); + } + + beginAppServerMetadataResourceRefresh(resource: MetadataResourceKind): () => boolean { + const context = this.context(); + const revision = this.options.cache.beginMetadataResourceRefresh(context, resource); + return () => + appServerQueryContextMatches(this.context(), context) && + this.options.cache.metadataResourceRefreshIsCurrent(context, resource, revision); } refreshAppServerMetadata(options: { forceSkills?: boolean } = {}): Promise { diff --git a/src/features/chat/application/connection/server-metadata-actions.ts b/src/features/chat/application/connection/server-metadata-actions.ts index 00e7aa34..878bfcc1 100644 --- a/src/features/chat/application/connection/server-metadata-actions.ts +++ b/src/features/chat/application/connection/server-metadata-actions.ts @@ -17,7 +17,11 @@ export type AppServerResourceEvent = export interface ServerMetadataActionsHost { stateStore: ChatStateStore; metadataResourceTransport: MetadataResourceTransport; - updateAppServerMetadata: (updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null) => SharedServerMetadata | null; + beginAppServerMetadataResourceRefresh: (resource: "skills" | "rateLimits") => () => boolean; + updateAppServerMetadata: ( + updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null, + resource?: "skills" | "rateLimits", + ) => SharedServerMetadata | null; appServerMetadataSnapshot: () => SharedServerMetadata | null; refreshAppServerMetadata: (options?: { forceSkills?: boolean }) => Promise; isStaleSharedQueryError: (error: unknown) => boolean; @@ -36,6 +40,18 @@ export function createServerMetadataActions(host: ServerMetadataActionsHost): Se }, refreshAppServerMetadata: () => refreshAppServerMetadata(host), applyAppServerResourceEvent: async (event) => { + if (event.type === "skills-changed") { + await refreshSkillResource(host, event.forceReload, host.beginAppServerMetadataResourceRefresh("skills")); + return; + } + if (event.type === "rate-limits-updated") { + await refreshRateLimitResource( + host, + { preserveExistingOnFailure: event.preserveExistingOnFailure === true }, + host.beginAppServerMetadataResourceRefresh("rateLimits"), + ); + return; + } await applyAppServerResourceEvent(host, event); }, }; @@ -44,10 +60,7 @@ export function createServerMetadataActions(host: ServerMetadataActionsHost): Se async function applyAppServerResourceEvent(host: ServerMetadataActionsHost, event: AppServerResourceEvent): Promise { switch (event.type) { case "skills-changed": - await refreshSkillResource(host, event.forceReload); - return; case "rate-limits-updated": - await refreshRateLimitResource(host, { preserveExistingOnFailure: event.preserveExistingOnFailure === true }); return; case "mcp-startup-status-updated": if (event.name.length > 0) { @@ -91,9 +104,13 @@ function applyCurrentAppServerMetadataSnapshot(host: ServerMetadataActionsHost): if (metadata) applyAppServerMetadata(host, metadata); } -async function refreshSkillResource(host: ServerMetadataActionsHost, forceReload = false): Promise { +async function refreshSkillResource( + host: ServerMetadataActionsHost, + forceReload = false, + isCurrent: () => boolean = () => true, +): Promise { const skills = await host.metadataResourceTransport.readSkillMetadata(forceReload); - if (!skills) return null; + if (!skills || !isCurrent()) return null; const next = host.updateAppServerMetadata((metadata) => { if (!metadata) return null; return { @@ -101,7 +118,7 @@ async function refreshSkillResource(host: ServerMetadataActionsHost, forceReload ...(skills.probe.status === "ok" ? { availableSkills: skills.value } : {}), serverDiagnostics: diagnosticsWithProbe(cloneServerDiagnostics(metadata.serverDiagnostics), skills.probe), }; - }); + }, "skills"); if (next) { applyAppServerMetadata(host, next); return next; @@ -122,9 +139,10 @@ async function refreshSkillResource(host: ServerMetadataActionsHost, forceReload async function refreshRateLimitResource( host: ServerMetadataActionsHost, options: { preserveExistingOnFailure?: boolean } = {}, + isCurrent: () => boolean = () => true, ): Promise { const rateLimit = await host.metadataResourceTransport.readRateLimitMetadata(); - if (!rateLimit) return; + if (!rateLimit || !isCurrent()) return; const preserveExistingOnFailure = options.preserveExistingOnFailure === true; const next = updateRateLimitMetadata(host, rateLimit, { preserveRateLimitOnFailure: preserveExistingOnFailure }); if (next) { @@ -156,7 +174,7 @@ function updateRateLimitMetadata( ...(rateLimit.probe.status === "ok" || !options.preserveRateLimitOnFailure ? { rateLimit: rateLimit.value } : {}), serverDiagnostics: diagnostics, }; - }); + }, "rateLimits"); } function applyMcpStartupStatusEvent( diff --git a/src/features/chat/host/bundles/connection-bundle.ts b/src/features/chat/host/bundles/connection-bundle.ts index 1758e507..c8686c20 100644 --- a/src/features/chat/host/bundles/connection-bundle.ts +++ b/src/features/chat/host/bundles/connection-bundle.ts @@ -125,7 +125,9 @@ export function createConnectionBundle( const serverMetadata = createServerMetadataActions({ stateStore, metadataResourceTransport: appServer.metadataResource, - updateAppServerMetadata: (updater) => environment.plugin.appServerQueries.updateAppServerMetadata(updater), + beginAppServerMetadataResourceRefresh: (resource) => + environment.plugin.appServerQueries.beginAppServerMetadataResourceRefresh(resource), + updateAppServerMetadata: (updater, resource) => environment.plugin.appServerQueries.updateAppServerMetadata(updater, resource), appServerMetadataSnapshot: () => environment.plugin.appServerQueries.appServerMetadataSnapshot(), refreshAppServerMetadata: (options) => environment.plugin.appServerQueries.refreshAppServerMetadata(options), isStaleSharedQueryError: isStaleAppServerSharedQueryContextError, diff --git a/src/features/chat/host/contracts.ts b/src/features/chat/host/contracts.ts index 275187f1..199f056a 100644 --- a/src/features/chat/host/contracts.ts +++ b/src/features/chat/host/contracts.ts @@ -51,7 +51,11 @@ interface WorkspacePanels { type ChatThreadCatalog = ThreadCatalogActiveReader & ThreadCatalogEventSink; interface ChatAppServerQueries { - updateAppServerMetadata(updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null): SharedServerMetadata | null; + beginAppServerMetadataResourceRefresh(resource: "skills" | "rateLimits"): () => boolean; + updateAppServerMetadata( + updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null, + resource?: "skills" | "rateLimits", + ): SharedServerMetadata | null; appServerMetadataSnapshot(): SharedServerMetadata | null; refreshAppServerMetadata(options?: { forceSkills?: boolean }): Promise; observeAppServerMetadataResult(listener: ObservedResultListener, options?: { emitCurrent?: boolean }): () => void; diff --git a/tests/app-server/query-cache.test.ts b/tests/app-server/query-cache.test.ts index 57a8b4a2..a56aa997 100644 --- a/tests/app-server/query-cache.test.ts +++ b/tests/app-server/query-cache.test.ts @@ -263,6 +263,46 @@ describe("AppServerQueryCache", () => { }); }); + it("does not append a stale load-more page after a newer first-page refresh", async () => { + const oldPage = deferred<{ data: ReturnType[]; nextCursor: string | null }>(); + const listThreads = vi + .fn() + .mockResolvedValueOnce({ data: [thread("old-first")], nextCursor: "old-page-2" }) + .mockImplementationOnce(() => oldPage.promise) + .mockResolvedValueOnce({ data: [thread("new-first")], nextCursor: "new-page-2" }); + const cache = cacheWithRequestHandlers({ "thread/list": listThreads }); + const context = cacheContext(); + await cache.refreshActiveThreads(context); + + const loadMore = cache.loadMoreActiveThreads(context); + await Promise.resolve(); + await cache.refreshActiveThreads(context); + oldPage.resolve({ data: [thread("old-second")], nextCursor: null }); + + await expect(loadMore).resolves.toEqual([thread("new-first")]); + expect(cache.activeThreadsSnapshot(context)).toEqual([thread("new-first")]); + expect(cache.hasMoreActiveThreads(context)).toBe(true); + }); + + it("retries a full active-thread inventory when a first-page refresh wins the race", async () => { + const oldInventory = deferred<{ data: ReturnType[]; nextCursor: string | null }>(); + const listThreads = vi + .fn() + .mockImplementationOnce(() => oldInventory.promise) + .mockResolvedValueOnce({ data: [thread("new-first")], nextCursor: "new-page-2" }) + .mockResolvedValueOnce({ data: [thread("new-first"), thread("new-second")], nextCursor: null }); + const cache = cacheWithRequestHandlers({ "thread/list": listThreads }); + const context = cacheContext(); + + const inventory = cache.fetchAllActiveThreads(context); + await flushMicrotasks(); + await cache.refreshActiveThreads(context); + oldInventory.resolve({ data: [thread("old")], nextCursor: null }); + + await expect(inventory).resolves.toEqual([thread("new-first"), thread("new-second")]); + expect(cache.activeThreadsSnapshot(context)).toEqual([thread("new-first"), thread("new-second")]); + }); + it("keys thread list refresh snapshots by app-server query context", async () => { const oldContext = cacheContext({ codexPath: "codex-old" }); const newContext = cacheContext({ codexPath: "codex-new" }); @@ -407,7 +447,33 @@ describe("AppServerQueryCache", () => { expect(cache.appServerMetadataSnapshot(context)).toEqual(refreshed); }); - it("does not merge local thread list updates into in-flight app-server snapshots", async () => { + it("does not overwrite a newer sparse metadata write with an in-flight full refresh", async () => { + const skills = deferred<{ data: { skills: CatalogSkillMetadata[] }[] }>(); + const cache = cacheWithRequestHandlers({ + "config/read": vi.fn().mockResolvedValue({}), + "model/list": vi.fn().mockResolvedValue({ data: [] }), + "skills/list": vi.fn(() => skills.promise), + "permissionProfile/list": vi.fn().mockResolvedValue({ data: [], nextCursor: null }), + "account/rateLimits/read": vi.fn().mockResolvedValue({ rateLimits: appServerRateLimit(0), rateLimitsByLimitId: null }), + }); + const context = cacheContext(); + cache.writeAppServerMetadata(context, metadata({ availableSkills: [skillMetadata("initial")] })); + for (let index = 0; index < 5; index += 1) cache.beginMetadataResourceRefresh(context, "skills"); + const refresh = cache.refreshAppServerMetadata(context); + await flushMicrotasks(); + + cache.updateAppServerMetadata( + context, + () => metadata({ availableSkills: [skillMetadata("event")], rateLimit: rateLimit(42) }), + "rateLimits", + ); + skills.resolve({ data: [{ skills: [catalogSkill("old-full")] }] }); + + await expect(refresh).resolves.toMatchObject({ availableSkills: [{ name: "event" }] }); + expect(cache.appServerMetadataSnapshot(context)?.availableSkills.map((skill) => skill.name)).toEqual(["event"]); + }); + + it("does not overwrite local thread list updates with an in-flight app-server snapshot", async () => { const context = cacheContext(); const refresh = deferred[]>(); const cache = cacheWithThreads(() => refresh.promise); @@ -419,8 +485,8 @@ describe("AppServerQueryCache", () => { cache.updateActiveThreads(context, (threads) => threads?.filter((item) => item.id !== "thread") ?? null); refresh.resolve([thread("thread"), thread("other")]); - await expect(promise).resolves.toEqual([thread("thread"), thread("other")]); - expect(cache.activeThreadsSnapshot(context)).toEqual([thread("thread"), thread("other")]); + await expect(promise).resolves.toEqual([thread("other")]); + expect(cache.activeThreadsSnapshot(context)).toEqual([thread("other")]); }); }); diff --git a/tests/features/chat/application/connection/server-actions.test.ts b/tests/features/chat/application/connection/server-actions.test.ts index fc11c102..353b60f9 100644 --- a/tests/features/chat/application/connection/server-actions.test.ts +++ b/tests/features/chat/application/connection/server-actions.test.ts @@ -117,6 +117,7 @@ describe("server metadata actions", () => { const actions = createServerMetadataActions({ stateStore, metadataResourceTransport: metadataResourceTransport({ readSkillMetadata: vi.fn().mockResolvedValue(null) }), + beginAppServerMetadataResourceRefresh: () => () => true, appServerMetadataSnapshot: () => null, updateAppServerMetadata, refreshAppServerMetadata: vi.fn().mockResolvedValue(null), @@ -153,6 +154,58 @@ describe("server metadata actions", () => { expect(stateStore.getState().connection.serverDiagnostics.probes.skills).toMatchObject({ status: "failed" }); }); + it("ignores an older skill refresh that completes after a newer refresh", async () => { + const older = deferred>>(); + const newer = deferred>>(); + const stateStore = createChatStateStore(chatStateFixture()); + const actions = createServerMetadataActions({ + stateStore, + metadataResourceTransport: metadataResourceTransport({ + readSkillMetadata: vi + .fn() + .mockImplementationOnce(() => older.promise) + .mockImplementationOnce(() => newer.promise), + }), + ...metadataCacheHost(), + refreshAppServerMetadata: vi.fn().mockResolvedValue(null), + isStaleSharedQueryError: () => false, + }); + + const first = actions.applyAppServerResourceEvent({ type: "skills-changed", forceReload: true }); + const second = actions.applyAppServerResourceEvent({ type: "skills-changed", forceReload: true }); + newer.resolve({ value: [skillFixture("new")], probe: diagnosticProbeOk("skills", "new", 2) }); + await second; + older.resolve({ value: [skillFixture("old")], probe: diagnosticProbeOk("skills", "old", 1) }); + await first; + + expect(stateStore.getState().connection.availableSkills.map((skill) => skill.name)).toEqual(["new"]); + }); + + it("shares skill refresh ordering across panel action instances", async () => { + const older = deferred>>(); + const cache = { current: serverMetadataFixture() as SharedServerMetadata | null }; + const sharedCache = metadataCacheHost(cache); + const createActions = (readSkillMetadata: MetadataResourceTransport["readSkillMetadata"]) => + createServerMetadataActions({ + stateStore: createChatStateStore(chatStateFixture()), + metadataResourceTransport: metadataResourceTransport({ readSkillMetadata }), + ...sharedCache, + refreshAppServerMetadata: vi.fn().mockResolvedValue(null), + isStaleSharedQueryError: () => false, + }); + const firstPanel = createActions(vi.fn(() => older.promise)); + const secondPanel = createActions( + vi.fn().mockResolvedValue({ value: [skillFixture("new")], probe: diagnosticProbeOk("skills", "new", 2) }), + ); + + const first = firstPanel.applyAppServerResourceEvent({ type: "skills-changed", forceReload: true }); + await secondPanel.applyAppServerResourceEvent({ type: "skills-changed", forceReload: true }); + older.resolve({ value: [skillFixture("old")], probe: diagnosticProbeOk("skills", "old", 1) }); + await first; + + expect(cache.current?.availableSkills.map((skill) => skill.name)).toEqual(["new"]); + }); + it("publishes refreshed rate limits from sparse update notifications", async () => { const stateStore = createChatStateStore(chatStateFixture()); const rateLimit = rateLimitFixture({ primary: { usedPercent: 64, windowDurationMins: 300, resetsAt: null } }); @@ -490,14 +543,24 @@ function serverMetadataFixture(overrides: Partial = {}): S } function metadataCacheHost(cache: { current: SharedServerMetadata | null } = { current: null }): { + beginAppServerMetadataResourceRefresh: (resource: "skills" | "rateLimits") => () => boolean; appServerMetadataSnapshot: () => SharedServerMetadata | null; - updateAppServerMetadata: (updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null) => SharedServerMetadata | null; + updateAppServerMetadata: ( + updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null, + resource?: "skills" | "rateLimits", + ) => SharedServerMetadata | null; } { + const generations = { skills: 0, rateLimits: 0 }; return { + beginAppServerMetadataResourceRefresh: (resource) => { + const generation = ++generations[resource]; + return () => generation === generations[resource]; + }, appServerMetadataSnapshot: () => cache.current, - updateAppServerMetadata: (updater) => { + updateAppServerMetadata: (updater, resource) => { const next = updater(cache.current); cache.current = next; + if (resource) generations[resource] += 1; return next; }, }; diff --git a/tests/features/chat/host/session-graph.test.ts b/tests/features/chat/host/session-graph.test.ts index 45461294..69529c5c 100644 --- a/tests/features/chat/host/session-graph.test.ts +++ b/tests/features/chat/host/session-graph.test.ts @@ -362,6 +362,7 @@ describe("createChatPanelSessionGraph actions", () => { overrides: Partial = {}, ): ChatPanelEnvironment["plugin"]["appServerQueries"] { return { + beginAppServerMetadataResourceRefresh: vi.fn(() => () => true), updateAppServerMetadata: vi.fn(() => null), appServerMetadataSnapshot: vi.fn(() => null), refreshAppServerMetadata: vi.fn().mockResolvedValue(null), diff --git a/tests/features/chat/host/view-connection.test.ts b/tests/features/chat/host/view-connection.test.ts index bb4322ac..291770f4 100644 --- a/tests/features/chat/host/view-connection.test.ts +++ b/tests/features/chat/host/view-connection.test.ts @@ -1154,6 +1154,7 @@ function chatHost(overrides: ChatHostFixtureOverrides = {}): TestCodexChatHost { openSideChat: overrides.openSideChat ?? vi.fn().mockResolvedValue(undefined), }, appServerQueries: { + beginAppServerMetadataResourceRefresh: () => () => true, updateAppServerMetadata: overrides.updateAppServerMetadata ?? ((updater) => {