mirror of
https://github.com/murashit/codex-panel.git
synced 2026-07-22 06:57:10 +00:00
Prevent stale shared app-server snapshots
This commit is contained in:
parent
26adfa4d46
commit
eb56414cd7
9 changed files with 282 additions and 25 deletions
|
|
@ -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<T>(
|
||||
|
|
@ -48,12 +49,16 @@ interface AppServerQueryOptions<T> {
|
|||
}
|
||||
|
||||
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<string, string | null>();
|
||||
private readonly activeThreadRevisions = new Map<string, number>();
|
||||
private readonly metadataRevisions = new Map<string, number>();
|
||||
private readonly metadataWriteRevisions = new Map<string, number>();
|
||||
|
||||
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<readonly Thread[]> => {
|
||||
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<SharedServerMetadata> => {
|
||||
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 {
|
||||
|
|
|
|||
|
|
@ -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<SharedServerMetadata | null> {
|
||||
|
|
|
|||
|
|
@ -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<SharedServerMetadata | null>;
|
||||
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<void> {
|
||||
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<SharedServerMetadata | null> {
|
||||
async function refreshSkillResource(
|
||||
host: ServerMetadataActionsHost,
|
||||
forceReload = false,
|
||||
isCurrent: () => boolean = () => true,
|
||||
): Promise<SharedServerMetadata | null> {
|
||||
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<void> {
|
||||
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(
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<SharedServerMetadata | null>;
|
||||
observeAppServerMetadataResult(listener: ObservedResultListener<SharedServerMetadata>, options?: { emitCurrent?: boolean }): () => void;
|
||||
|
|
|
|||
|
|
@ -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<typeof thread>[]; 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<typeof thread>[]; 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<readonly ReturnType<typeof thread>[]>();
|
||||
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")]);
|
||||
});
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -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<Awaited<ReturnType<MetadataResourceTransport["readSkillMetadata"]>>>();
|
||||
const newer = deferred<Awaited<ReturnType<MetadataResourceTransport["readSkillMetadata"]>>>();
|
||||
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<Awaited<ReturnType<MetadataResourceTransport["readSkillMetadata"]>>>();
|
||||
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<SharedServerMetadata> = {}): 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;
|
||||
},
|
||||
};
|
||||
|
|
|
|||
|
|
@ -362,6 +362,7 @@ describe("createChatPanelSessionGraph actions", () => {
|
|||
overrides: Partial<ChatPanelEnvironment["plugin"]["appServerQueries"]> = {},
|
||||
): ChatPanelEnvironment["plugin"]["appServerQueries"] {
|
||||
return {
|
||||
beginAppServerMetadataResourceRefresh: vi.fn(() => () => true),
|
||||
updateAppServerMetadata: vi.fn(() => null),
|
||||
appServerMetadataSnapshot: vi.fn(() => null),
|
||||
refreshAppServerMetadata: vi.fn().mockResolvedValue(null),
|
||||
|
|
|
|||
|
|
@ -1154,6 +1154,7 @@ function chatHost(overrides: ChatHostFixtureOverrides = {}): TestCodexChatHost {
|
|||
openSideChat: overrides.openSideChat ?? vi.fn().mockResolvedValue(undefined),
|
||||
},
|
||||
appServerQueries: {
|
||||
beginAppServerMetadataResourceRefresh: () => () => true,
|
||||
updateAppServerMetadata:
|
||||
overrides.updateAppServerMetadata ??
|
||||
((updater) => {
|
||||
|
|
|
|||
Loading…
Reference in a new issue