Consolidate app-server query cache ownership

This commit is contained in:
murashit 2026-06-16 06:46:29 +09:00
parent 6afd94b1de
commit 85897bcb1d
22 changed files with 877 additions and 403 deletions

View file

@ -1,9 +1,11 @@
import { QueryClient, QueryObserver, type QueryObserverResult } from "@tanstack/query-core";
import { MutationObserver, QueryClient, QueryObserver, type QueryObserverResult } from "@tanstack/query-core";
import type { AppServerClient } from "../connection/client";
import { listModelMetadata } from "../services/catalog";
import { readRateLimitMetadataProbe, readRuntimeConfigSnapshot, readSkillMetadataProbe } from "../services/metadata";
import { listThreads } from "../services/threads";
import type { ModelMetadata } from "../../domain/catalog/metadata";
import { createServerDiagnostics, diagnosticProbeError, diagnosticProbeOk, diagnosticsWithProbe } from "../../domain/server/diagnostics";
import type { Thread } from "../../domain/threads/model";
import {
activeThreadsQueryKey,
@ -11,12 +13,14 @@ import {
appServerModelsQueryKey,
appServerQueriesFilter,
appServerQueryContextIsComplete,
appServerQueryScope,
cloneAppServerQueryContext,
type AppServerQueryContext,
} from "./keys";
import { cloneModelMetadata, cloneSharedServerMetadata, cloneThreads, type SharedServerMetadata } from "./snapshots";
const ACTIVE_THREADS_STALE_TIME_MS = 10_000;
const APP_SERVER_METADATA_STALE_TIME_MS = 10_000;
const MODELS_STALE_TIME_MS = 60_000;
export interface AppServerQueryClientRunner {
@ -27,6 +31,17 @@ export interface AppServerQueryClientRunner {
): Promise<T>;
}
export type AppServerObservedQueryResult<T> = Omit<QueryObserverResult<T>, "data" | "error"> & {
readonly data: T | null;
readonly error: Error | null;
};
interface AppServerQueryOptions<T> {
readonly queryKey: readonly unknown[];
readonly queryFn: () => Promise<T>;
readonly staleTime: number;
}
export class AppServerQueryCache {
readonly client: QueryClient;
private readonly clientRunner: AppServerQueryClientRunner | null;
@ -42,7 +57,10 @@ export class AppServerQueryCache {
clearContext(context: AppServerQueryContext): void {
if (!appServerQueryContextIsComplete(context)) return;
this.client.removeQueries(appServerQueriesFilter(context));
const filter = appServerQueriesFilter(context);
void this.client.cancelQueries(filter);
this.client.removeQueries(filter);
this.removeContextMutations(context);
}
activeThreadsSnapshot(context: AppServerQueryContext): readonly Thread[] | null {
@ -51,12 +69,12 @@ export class AppServerQueryCache {
return threads ? cloneThreads(threads) : null;
}
observeActiveThreads(
observeActiveThreadsResult(
context: AppServerQueryContext,
listener: (threads: readonly Thread[]) => void,
listener: (result: AppServerObservedQueryResult<readonly Thread[]>) => void,
options: { emitCurrent?: boolean } = {},
): () => void {
return this.observeQuery(activeThreadsQueryKey(context), cloneThreads, listener, options);
return this.observeQueryResult(this.activeThreadsQueryOptions(context), cloneThreads, listener, options);
}
async fetchActiveThreads(context: AppServerQueryContext, options: { force?: boolean } = {}): Promise<readonly Thread[]> {
@ -66,16 +84,7 @@ export class AppServerQueryCache {
}
const key = activeThreadsQueryKey(refreshContext);
if (options.force) await this.client.invalidateQueries({ queryKey: key });
const threads = await this.client.fetchQuery({
queryKey: key,
queryFn: async () => {
const nextThreads = cloneThreads(
await this.runWithClient(refreshContext, (client) => listThreads(client, refreshContext.vaultPath)),
);
return nextThreads;
},
staleTime: ACTIVE_THREADS_STALE_TIME_MS,
});
const threads = await this.client.fetchQuery(this.activeThreadsQueryOptions(refreshContext));
return cloneThreads(threads);
}
@ -95,7 +104,7 @@ export class AppServerQueryCache {
const current = this.activeThreadsSnapshot(context);
const next = updater(current);
if (!next) return null;
this.setActiveThreads(context, next);
this.applyActiveThreadsMutation(context, current, next);
return cloneThreads(next);
}
@ -105,15 +114,41 @@ export class AppServerQueryCache {
return metadata ? cloneSharedServerMetadata(metadata) : null;
}
observeAppServerMetadata(
observeAppServerMetadataResult(
context: AppServerQueryContext,
listener: (metadata: SharedServerMetadata) => void,
listener: (result: AppServerObservedQueryResult<SharedServerMetadata>) => void,
options: { emitCurrent?: boolean } = {},
): () => void {
return this.observeQuery(appServerMetadataQueryKey(context), cloneSharedServerMetadata, listener, options);
return this.observeQueryResult(this.appServerMetadataQueryOptions(context), cloneSharedServerMetadata, listener, options);
}
setAppServerMetadata(context: AppServerQueryContext, metadata: SharedServerMetadata): SharedServerMetadata | null {
async fetchAppServerMetadata(
context: AppServerQueryContext,
options: { force?: boolean; forceSkills?: boolean } = {},
): Promise<SharedServerMetadata | null> {
const refreshContext = cloneAppServerQueryContext(context);
if (!appServerQueryContextIsComplete(refreshContext)) {
return null;
}
const key = appServerMetadataQueryKey(refreshContext);
if (options.force) {
await Promise.all([
this.client.invalidateQueries({ queryKey: key }),
this.client.invalidateQueries({ queryKey: appServerModelsQueryKey(refreshContext) }),
]);
}
const metadata = await this.client.fetchQuery(this.appServerMetadataQueryOptions(refreshContext, options));
return this.writeAppServerMetadata(refreshContext, metadata);
}
async refreshAppServerMetadata(
context: AppServerQueryContext,
options: { forceSkills?: boolean } = {},
): Promise<SharedServerMetadata | null> {
return this.fetchAppServerMetadata(context, { ...options, force: true });
}
writeAppServerMetadata(context: AppServerQueryContext, metadata: SharedServerMetadata): SharedServerMetadata | null {
if (!appServerQueryContextIsComplete(context)) return null;
const next = cloneSharedServerMetadata({
...metadata,
@ -127,18 +162,27 @@ export class AppServerQueryCache {
return cloneSharedServerMetadata(next);
}
updateAppServerMetadata(
context: AppServerQueryContext,
updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null,
): SharedServerMetadata | null {
if (!appServerQueryContextIsComplete(context)) return null;
const next = updater(this.appServerMetadataSnapshot(context));
return next ? this.writeAppServerMetadata(context, next) : null;
}
modelsSnapshot(context: AppServerQueryContext): readonly ModelMetadata[] | null {
if (!appServerQueryContextIsComplete(context)) return null;
const models = this.client.getQueryData<readonly ModelMetadata[]>(appServerModelsQueryKey(context));
return models ? cloneModelMetadata(models) : null;
}
observeModels(
observeModelsResult(
context: AppServerQueryContext,
listener: (models: readonly ModelMetadata[]) => void,
listener: (result: AppServerObservedQueryResult<readonly ModelMetadata[]>) => void,
options: { emitCurrent?: boolean } = {},
): () => void {
return this.observeQuery(appServerModelsQueryKey(context), cloneModelMetadata, listener, options);
return this.observeQueryResult(this.modelsQueryOptions(context), cloneModelMetadata, listener, options);
}
async fetchModels(context: AppServerQueryContext, options: { force?: boolean } = {}): Promise<readonly ModelMetadata[]> {
@ -148,17 +192,7 @@ export class AppServerQueryCache {
}
const key = appServerModelsQueryKey(refreshContext);
if (options.force) await this.client.invalidateQueries({ queryKey: key });
const models = await this.client.fetchQuery({
queryKey: key,
queryFn: async () => {
return cloneModelMetadata(
await this.runWithClient(refreshContext, (client) => listModelMetadata(client), {
unhandledServerRequestMessage: "Codex model list refresh does not handle server requests.",
}),
);
},
staleTime: MODELS_STALE_TIME_MS,
});
const models = await this.client.fetchQuery(this.modelsQueryOptions(refreshContext));
return cloneModelMetadata(models);
}
@ -166,24 +200,157 @@ export class AppServerQueryCache {
return this.fetchModels(context, { force: true });
}
private observeQuery<T>(
queryKey: readonly unknown[],
private activeThreadsQueryOptions(context: AppServerQueryContext): AppServerQueryOptions<readonly Thread[]> {
const refreshContext = cloneAppServerQueryContext(context);
return {
queryKey: activeThreadsQueryKey(refreshContext),
queryFn: async (): Promise<readonly Thread[]> => {
return cloneThreads(await this.runWithClient(refreshContext, (client) => listThreads(client, refreshContext.vaultPath)));
},
staleTime: ACTIVE_THREADS_STALE_TIME_MS,
};
}
private appServerMetadataQueryOptions(
context: AppServerQueryContext,
options: { forceSkills?: boolean } = {},
): AppServerQueryOptions<SharedServerMetadata> {
const refreshContext = cloneAppServerQueryContext(context);
return {
queryKey: appServerMetadataQueryKey(refreshContext),
queryFn: async (): Promise<SharedServerMetadata> => {
return this.runWithClient(refreshContext, async (client) => {
const runtimeConfig = await readRuntimeConfigSnapshot(client, refreshContext.vaultPath);
const [models, skills, rateLimit] = await Promise.all([
this.readModelMetadataProbe(refreshContext, client),
readSkillMetadataProbe(client, refreshContext.vaultPath, options.forceSkills ?? false),
readRateLimitMetadataProbe(client),
]);
const diagnostics = [models.probe, skills.probe, rateLimit.probe].reduce(
(current, probe) => diagnosticsWithProbe(current, probe),
this.appServerMetadataSnapshot(refreshContext)?.serverDiagnostics ?? createServerDiagnostics(),
);
return {
runtimeConfig,
availableModels: models.data,
availableSkills: skills.data,
rateLimit: rateLimit.data,
serverDiagnostics: diagnostics,
};
});
},
staleTime: APP_SERVER_METADATA_STALE_TIME_MS,
};
}
private modelsQueryOptions(context: AppServerQueryContext): AppServerQueryOptions<readonly ModelMetadata[]> {
const refreshContext = cloneAppServerQueryContext(context);
return {
queryKey: appServerModelsQueryKey(refreshContext),
queryFn: async (): Promise<readonly ModelMetadata[]> => {
return cloneModelMetadata(
await this.runWithClient(refreshContext, (client) => listModelMetadata(client), {
unhandledServerRequestMessage: "Codex model list refresh does not handle server requests.",
}),
);
},
staleTime: MODELS_STALE_TIME_MS,
};
}
private async readModelMetadataProbe(
context: AppServerQueryContext,
client: AppServerClient,
): Promise<{
data: readonly ModelMetadata[];
probe: SharedServerMetadata["serverDiagnostics"]["probes"]["model/list"];
}> {
try {
const data = cloneModelMetadata(await this.client.fetchQuery(this.modelsQueryOptionsWithClient(context, client)));
return { data, probe: diagnosticProbeOk("model/list", `${String(data.length)} models`) };
} catch (error) {
return {
data: this.modelsSnapshot(context) ?? [],
probe: diagnosticProbeError("model/list", error),
};
}
}
private modelsQueryOptionsWithClient(
context: AppServerQueryContext,
client: AppServerClient,
): AppServerQueryOptions<readonly ModelMetadata[]> {
const refreshContext = cloneAppServerQueryContext(context);
return {
...this.modelsQueryOptions(refreshContext),
queryFn: async (): Promise<readonly ModelMetadata[]> => cloneModelMetadata(await listModelMetadata(client)),
};
}
private observeQueryResult<T>(
queryOptions: AppServerQueryOptions<T>,
clone: (value: T) => T,
listener: (value: T) => void,
listener: (result: AppServerObservedQueryResult<T>) => void,
options: { emitCurrent?: boolean },
): () => void {
const observer = new QueryObserver<T>(this.client, {
queryKey,
...queryOptions,
enabled: false,
});
const emit = (result: QueryObserverResult<T>): void => {
if (result.data !== undefined) listener(clone(result.data));
listener(this.cloneObservedResult(result, clone));
};
const unsubscribe = observer.subscribe(emit);
if (options.emitCurrent ?? true) emit(observer.getCurrentResult());
return unsubscribe;
}
private cloneObservedResult<T>(result: QueryObserverResult<T>, clone: (value: T) => T): AppServerObservedQueryResult<T> {
return {
...result,
data: result.data === undefined ? null : clone(result.data),
error: result.error instanceof Error ? result.error : null,
};
}
private applyActiveThreadsMutation(context: AppServerQueryContext, previous: readonly Thread[] | null, next: readonly Thread[]): void {
if (!appServerQueryContextIsComplete(context)) return;
const mutationKey = [...activeThreadsQueryKey(context), "optimistic-update"] as const;
const observer = new MutationObserver<
readonly Thread[],
Error,
{ readonly previous: readonly Thread[] | null; readonly next: readonly Thread[] },
{ readonly previous: readonly Thread[] | null }
>(this.client, {
mutationKey,
mutationFn: async (variables) => cloneThreads(variables.next),
onMutate: (variables) => {
void this.client.cancelQueries({ queryKey: activeThreadsQueryKey(context) });
this.client.setQueryData(activeThreadsQueryKey(context), cloneThreads(variables.next));
return { previous: variables.previous ? cloneThreads(variables.previous) : null };
},
onError: (_error, _variables, mutationContext) => {
if (mutationContext?.previous) {
this.client.setQueryData(activeThreadsQueryKey(context), cloneThreads(mutationContext.previous));
}
},
});
void observer.mutate({ previous, next }).finally(() => {
observer.reset();
});
}
private removeContextMutations(context: AppServerQueryContext): void {
const scope = appServerQueryScope(context);
for (const mutation of this.client.getMutationCache().getAll()) {
const mutationKey = mutation.options.mutationKey;
if (!Array.isArray(mutationKey)) continue;
if (scope.every((part, index) => mutationKey[index] === part)) {
this.client.getMutationCache().remove(mutation);
}
}
}
private runWithClient<T>(
context: AppServerQueryContext,
operation: (client: AppServerClient) => Promise<T>,

View file

@ -27,7 +27,7 @@ export function appServerQueryContextMatches(left: AppServerQueryContext, right:
);
}
function appServerQueryScope(context: AppServerQueryContext): AppServerQueryScope {
export function appServerQueryScope(context: AppServerQueryContext): AppServerQueryScope {
return ["app-server", context.codexPath, context.vaultPath];
}

View file

@ -18,6 +18,7 @@ import { cloneServerDiagnostics, type ChatServerActionHost } from "./host";
interface RefreshDiagnosticProbesOptions {
appServerMetadataSnapshot?: boolean;
forceResourceProbes?: boolean;
}
interface DiagnosticProbeSnapshot {
@ -27,8 +28,8 @@ interface DiagnosticProbeSnapshot {
}
export interface ChatServerDiagnosticsActionsHost extends ChatServerActionHost {
setAppServerMetadata: (metadata: SharedServerMetadata) => void;
serverMetadataSnapshot: () => SharedServerMetadata;
updateAppServerMetadata: (updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null) => SharedServerMetadata | null;
appServerMetadataSnapshot: () => SharedServerMetadata | null;
}
export interface ChatServerDiagnosticsActions {
@ -59,7 +60,7 @@ async function refreshDiagnosticProbes(
if (!client) return false;
const probes: Promise<DiagnosticProbeSnapshot>[] = [];
if (!options.appServerMetadataSnapshot) {
if (options.forceResourceProbes === true && options.appServerMetadataSnapshot !== true) {
probes.push(
probeDiagnostic(
"model/list",
@ -121,11 +122,12 @@ async function refreshDiagnosticProbes(
const results = await Promise.all(probes);
if (host.currentClient() !== client) return false;
let diagnostics = cloneServerDiagnostics(host.stateStore.getState().connection.serverDiagnostics);
let diagnostics = currentMetadataDiagnostics(host);
for (const result of results) {
diagnostics = diagnosticsWithProbe(diagnostics, result.probe);
if (result.mcpServerStatuses) diagnostics = upsertMcpServerStatusDiagnostics(diagnostics, result.mcpServerStatuses);
}
host.updateAppServerMetadata((metadata) => (metadata ? { ...metadata, serverDiagnostics: diagnostics } : null));
host.stateStore.dispatch({ type: "connection/metadata-applied", serverDiagnostics: diagnostics });
return true;
}
@ -139,8 +141,7 @@ async function refreshPublishedDiagnosticProbes(
host: ChatServerDiagnosticsActionsHost,
options: RefreshDiagnosticProbesOptions = {},
): Promise<void> {
if (!(await refreshDiagnosticProbes(host, options))) return;
host.setAppServerMetadata(host.serverMetadataSnapshot());
await refreshDiagnosticProbes(host, options);
}
async function mcpStatusLines(host: ChatServerDiagnosticsActionsHost): Promise<string[]> {
@ -163,18 +164,26 @@ function recordMcpStartupStatus(
startupStatus: McpServerStartupStatus,
message: string | null,
): void {
const diagnostics = upsertMcpServerDiagnostic(currentMetadataDiagnostics(host), {
name,
startupStatus,
authStatus: null,
toolCount: null,
message,
});
host.updateAppServerMetadata((metadata) => (metadata ? { ...metadata, serverDiagnostics: diagnostics } : null));
host.stateStore.dispatch({
type: "connection/metadata-applied",
serverDiagnostics: upsertMcpServerDiagnostic(host.stateStore.getState().connection.serverDiagnostics, {
name,
startupStatus,
authStatus: null,
toolCount: null,
message,
}),
serverDiagnostics: diagnostics,
});
}
function currentMetadataDiagnostics(host: ChatServerDiagnosticsActionsHost): SharedServerMetadata["serverDiagnostics"] {
return (
host.appServerMetadataSnapshot()?.serverDiagnostics ?? cloneServerDiagnostics(host.stateStore.getState().connection.serverDiagnostics)
);
}
async function probeDiagnostic<T>(
method: DiagnosticProbeMethod,
request: () => Promise<T>,

View file

@ -1,29 +1,26 @@
import {
readRateLimitMetadataProbe,
readRuntimeConfigSnapshot,
readSkillMetadataProbe,
type RateLimitMetadataProbeResult,
type SkillMetadataProbeResult,
} from "../../../../app-server/services/metadata";
import type { ModelMetadata } from "../../../../domain/catalog/metadata";
import { diagnosticsWithProbe, diagnosticProbeError, diagnosticProbeOk } from "../../../../domain/server/diagnostics";
import { diagnosticsWithProbe } from "../../../../domain/server/diagnostics";
import type { SharedServerMetadata } from "../../../../domain/server/metadata";
import { cloneServerDiagnostics, type ChatServerActionHost } from "./host";
export interface ChatServerMetadataActionsHost extends ChatServerActionHost {
setAppServerMetadata: (metadata: SharedServerMetadata) => void;
modelsSnapshot: () => readonly ModelMetadata[] | null;
fetchModels: () => Promise<readonly ModelMetadata[]>;
refreshModels: () => Promise<readonly ModelMetadata[]>;
updateAppServerMetadata: (updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null) => SharedServerMetadata | null;
appServerMetadataSnapshot: () => SharedServerMetadata | null;
fetchAppServerMetadata: () => Promise<SharedServerMetadata | null>;
refreshAppServerMetadata: (options?: { forceSkills?: boolean }) => Promise<SharedServerMetadata | null>;
}
export interface ChatServerMetadataActions {
serverMetadataSnapshot: () => SharedServerMetadata;
applyAppServerMetadata: (metadata: SharedServerMetadata) => void;
loadAppServerMetadata: () => Promise<SharedServerMetadata | null>;
refreshAppServerMetadata: () => Promise<SharedServerMetadata | null>;
refreshPublishedAppServerMetadata: () => Promise<SharedServerMetadata | null>;
setAppServerMetadataSnapshot: () => void;
applyAppServerMetadataSnapshot: () => void;
refreshSkills: (forceReload?: boolean) => Promise<void>;
refreshPublishedSkills: (forceReload?: boolean) => Promise<void>;
loadSkills: (forceReload?: boolean) => Promise<SkillMetadataProbeResult>;
@ -34,15 +31,14 @@ export interface ChatServerMetadataActions {
export function createChatServerMetadataActions(host: ChatServerMetadataActionsHost): ChatServerMetadataActions {
return {
serverMetadataSnapshot: () => serverMetadataSnapshot(host),
applyAppServerMetadata: (metadata) => {
applyAppServerMetadata(host, metadata);
},
loadAppServerMetadata: () => loadAppServerMetadata(host),
refreshAppServerMetadata: () => refreshAppServerMetadata(host),
refreshPublishedAppServerMetadata: () => refreshPublishedAppServerMetadata(host),
setAppServerMetadataSnapshot: () => {
setAppServerMetadataSnapshot(host);
applyAppServerMetadataSnapshot: () => {
applyAppServerMetadataSnapshot(host);
},
refreshSkills: async (forceReload) => {
await refreshSkills(host, forceReload);
@ -55,17 +51,6 @@ export function createChatServerMetadataActions(host: ChatServerMetadataActionsH
};
}
function serverMetadataSnapshot(host: ChatServerMetadataActionsHost): SharedServerMetadata {
const state = host.stateStore.getState();
return {
runtimeConfig: state.connection.runtimeConfig,
availableModels: state.connection.availableModels,
availableSkills: state.connection.availableSkills,
rateLimit: state.connection.rateLimit,
serverDiagnostics: state.connection.serverDiagnostics,
};
}
function applyAppServerMetadata(host: ChatServerMetadataActionsHost, metadata: SharedServerMetadata): void {
host.stateStore.dispatch({
type: "connection/metadata-applied",
@ -78,84 +63,52 @@ function applyAppServerMetadata(host: ChatServerMetadataActionsHost, metadata: S
}
async function loadAppServerMetadata(host: ChatServerMetadataActionsHost): Promise<SharedServerMetadata | null> {
const client = host.currentClient();
if (!client) return null;
return loadAppServerMetadataFromClient(host, client);
}
async function loadAppServerMetadataFromClient(
host: ChatServerMetadataActionsHost,
client: NonNullable<ReturnType<ChatServerMetadataActionsHost["currentClient"]>>,
): Promise<SharedServerMetadata> {
const runtimeConfig = await readRuntimeConfigSnapshot(client, host.vaultPath);
const [models, skills, rateLimit] = await Promise.all([
loadModelMetadataFromQuery(host),
readSkillMetadataProbe(client, host.vaultPath),
readRateLimitMetadataProbe(client),
]);
const diagnostics = [models.probe, skills.probe, rateLimit.probe].reduce(
(current, probe) => diagnosticsWithProbe(current, probe),
cloneServerDiagnostics(host.stateStore.getState().connection.serverDiagnostics),
);
return {
runtimeConfig,
availableModels: models.data,
availableSkills: skills.data,
rateLimit: rateLimit.data,
serverDiagnostics: diagnostics,
};
return host.fetchAppServerMetadata();
}
async function refreshAppServerMetadata(host: ChatServerMetadataActionsHost): Promise<SharedServerMetadata | null> {
const client = host.currentClient();
if (!client) return null;
const metadata = await loadAppServerMetadataFromClient(host, client);
if (host.currentClient() !== client) return null;
const metadata = await host.refreshAppServerMetadata();
if (!metadata) return null;
applyAppServerMetadata(host, metadata);
return metadata;
}
async function refreshPublishedAppServerMetadata(host: ChatServerMetadataActionsHost): Promise<SharedServerMetadata | null> {
const metadata = await refreshAppServerMetadata(host);
if (metadata) host.setAppServerMetadata(metadata);
return metadata;
return refreshAppServerMetadata(host);
}
async function loadModelMetadataFromQuery(host: ChatServerMetadataActionsHost): Promise<{
data: readonly ModelMetadata[];
probe: SharedServerMetadata["serverDiagnostics"]["probes"]["model/list"];
}> {
try {
const data = await host.fetchModels();
return { data, probe: diagnosticProbeOk("model/list", `${String(data.length)} models`) };
} catch (error) {
return {
data: host.modelsSnapshot() ?? [],
probe: diagnosticProbeError("model/list", error),
};
}
function applyAppServerMetadataSnapshot(host: ChatServerMetadataActionsHost): void {
const metadata = host.appServerMetadataSnapshot();
if (metadata) applyAppServerMetadata(host, metadata);
}
function setAppServerMetadataSnapshot(host: ChatServerMetadataActionsHost): void {
host.setAppServerMetadata(serverMetadataSnapshot(host));
}
async function refreshSkills(host: ChatServerMetadataActionsHost, forceReload = false): Promise<boolean> {
async function refreshSkills(host: ChatServerMetadataActionsHost, forceReload = false): Promise<SharedServerMetadata | null> {
const client = host.currentClient();
const skills = await readSkillMetadataProbe(client, host.vaultPath, forceReload);
if (client && host.currentClient() !== client) return false;
const diagnostics = diagnosticsWithProbe(cloneServerDiagnostics(host.stateStore.getState().connection.serverDiagnostics), skills.probe);
if (client && host.currentClient() !== client) return null;
const next = host.updateAppServerMetadata((metadata) => {
if (!metadata) return null;
return {
...metadata,
availableSkills: skills.data,
serverDiagnostics: diagnosticsWithProbe(cloneServerDiagnostics(metadata.serverDiagnostics), skills.probe),
};
});
if (next) {
applyAppServerMetadata(host, next);
return next;
}
const diagnostics = diagnosticsWithProbe(currentMetadataDiagnostics(host), skills.probe);
host.stateStore.dispatch({
type: "connection/metadata-applied",
availableSkills: skills.data,
serverDiagnostics: diagnostics,
});
return true;
return null;
}
async function refreshPublishedSkills(host: ChatServerMetadataActionsHost, forceReload = false): Promise<void> {
if (!(await refreshSkills(host, forceReload))) return;
setAppServerMetadataSnapshot(host);
await refreshSkills(host, forceReload);
}
async function loadSkills(host: ChatServerMetadataActionsHost, forceReload = false): Promise<SkillMetadataProbeResult> {
@ -166,10 +119,12 @@ async function refreshRateLimits(host: ChatServerMetadataActionsHost): Promise<v
const client = host.currentClient();
const rateLimit = await readRateLimitMetadataProbe(client);
if (client && host.currentClient() !== client) return;
const diagnostics = diagnosticsWithProbe(
cloneServerDiagnostics(host.stateStore.getState().connection.serverDiagnostics),
rateLimit.probe,
);
const next = updateRateLimitMetadata(host, rateLimit, { preserveRateLimitOnFailure: false });
if (next) {
applyAppServerMetadata(host, next);
return;
}
const diagnostics = diagnosticsWithProbe(currentMetadataDiagnostics(host), rateLimit.probe);
host.stateStore.dispatch({
type: "connection/metadata-applied",
rateLimit: rateLimit.data,
@ -181,17 +136,18 @@ async function refreshPublishedRateLimits(host: ChatServerMetadataActionsHost):
const client = host.currentClient();
const rateLimit = await readRateLimitMetadataProbe(client);
if (client && host.currentClient() !== client) return;
const diagnostics = diagnosticsWithProbe(
cloneServerDiagnostics(host.stateStore.getState().connection.serverDiagnostics),
rateLimit.probe,
);
const next = updateRateLimitMetadata(host, rateLimit, { preserveRateLimitOnFailure: true });
if (next) {
applyAppServerMetadata(host, next);
return;
}
const diagnostics = diagnosticsWithProbe(currentMetadataDiagnostics(host), rateLimit.probe);
if (rateLimit.probe.status === "ok") {
host.stateStore.dispatch({
type: "connection/metadata-applied",
rateLimit: rateLimit.data,
serverDiagnostics: diagnostics,
});
setAppServerMetadataSnapshot(host);
return;
}
host.stateStore.dispatch({ type: "connection/metadata-applied", serverDiagnostics: diagnostics });
@ -200,3 +156,25 @@ async function refreshPublishedRateLimits(host: ChatServerMetadataActionsHost):
async function loadRateLimit(host: ChatServerMetadataActionsHost): Promise<RateLimitMetadataProbeResult> {
return readRateLimitMetadataProbe(host.currentClient());
}
function updateRateLimitMetadata(
host: ChatServerMetadataActionsHost,
rateLimit: RateLimitMetadataProbeResult,
options: { preserveRateLimitOnFailure: boolean },
): SharedServerMetadata | null {
return host.updateAppServerMetadata((metadata) => {
if (!metadata) return null;
const diagnostics = diagnosticsWithProbe(cloneServerDiagnostics(metadata.serverDiagnostics), rateLimit.probe);
return {
...metadata,
...(rateLimit.probe.status === "ok" || !options.preserveRateLimitOnFailure ? { rateLimit: rateLimit.data } : {}),
serverDiagnostics: diagnostics,
};
});
}
function currentMetadataDiagnostics(host: ChatServerMetadataActionsHost): SharedServerMetadata["serverDiagnostics"] {
return (
host.appServerMetadataSnapshot()?.serverDiagnostics ?? cloneServerDiagnostics(host.stateStore.getState().connection.serverDiagnostics)
);
}

View file

@ -38,7 +38,7 @@ export interface ChatInboundControllerActions {
fetchActiveThreads: () => void;
refreshRateLimits: () => void;
refreshSkills: (forceReload?: boolean) => void;
setAppServerMetadata: () => void;
applyAppServerMetadataSnapshot: () => void;
maybeNameThread: (threadId: string, turnId: string, completedSummary: ThreadConversationSummary | null) => void;
applyThreadArchived: (threadId: string) => void;
applyThreadRenamed: (threadId: string, name: string | null) => void;
@ -193,8 +193,8 @@ export class ChatInboundController {
case "refresh-skills":
this.actions.refreshSkills(effect.forceReload);
return;
case "publish-app-server-metadata":
this.actions.setAppServerMetadata();
case "apply-app-server-metadata-snapshot":
this.actions.applyAppServerMetadataSnapshot();
return;
case "maybe-name-thread":
this.actions.maybeNameThread(effect.threadId, effect.turnId, effect.completedSummary);

View file

@ -52,7 +52,7 @@ export type ChatNotificationEffect =
| { type: "refresh-threads" }
| { type: "refresh-rate-limits" }
| { type: "refresh-skills"; forceReload: boolean }
| { type: "publish-app-server-metadata" }
| { type: "apply-app-server-metadata-snapshot" }
| { type: "maybe-name-thread"; threadId: string; turnId: string; completedSummary: ThreadConversationSummary | null }
| { type: "apply-thread-archived"; threadId: string }
| { type: "apply-thread-renamed"; threadId: string; name: string | null }
@ -105,7 +105,7 @@ const DIAGNOSTIC_STATUS_PLANNERS = {
actions: [],
effects:
notification.params.name.length === 0
? [{ type: "publish-app-server-metadata" }]
? [{ type: "apply-app-server-metadata-snapshot" }]
: [
{
type: "record-mcp-startup-status",
@ -113,7 +113,7 @@ const DIAGNOSTIC_STATUS_PLANNERS = {
status: notification.params.status,
message: notification.params.error,
},
{ type: "publish-app-server-metadata" },
{ type: "apply-app-server-metadata-snapshot" },
],
}),
} satisfies ServerNotificationPlannerMap<DiagnosticStatusNotificationMethod>;

View file

@ -24,7 +24,7 @@ export interface ChatConnectionMetadataActions {
}
export interface ChatConnectionDiagnosticsActions {
refreshPublishedDiagnosticProbes: () => Promise<void>;
refreshPublishedDiagnosticProbes: (options?: { appServerMetadataSnapshot?: boolean; forceResourceProbes?: boolean }) => Promise<void>;
}
export interface ChatConnectionControllerHost {
@ -93,7 +93,6 @@ export class ChatConnectionController {
if (!this.host.connection.currentClient()) return;
try {
await this.host.loadSharedThreadList();
await this.host.metadata.refreshPublishedAppServerMetadata();
this.host.refreshTabHeader();
} catch (error) {
this.host.addSystemMessage(error instanceof Error ? error.message : String(error));
@ -105,7 +104,8 @@ export class ChatConnectionController {
await this.ensureConnected();
if (!this.host.connection.currentClient()) return;
this.host.clearDeferredDiagnostics();
await this.host.diagnostics.refreshPublishedDiagnosticProbes();
await this.host.metadata.refreshPublishedAppServerMetadata();
await this.host.diagnostics.refreshPublishedDiagnosticProbes({ appServerMetadataSnapshot: true });
}
async refreshStatusPanel(): Promise<void> {

View file

@ -46,7 +46,10 @@ export interface ChatConnectionBundleContext {
deferredTasks: ChatViewDeferredTasks;
threadCatalog: {
setActiveThreads(threads: readonly Thread[]): void;
setAppServerMetadata(metadata: SharedServerMetadata): void;
updateAppServerMetadata(updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null): SharedServerMetadata | null;
appServerMetadataSnapshot(): SharedServerMetadata | null;
fetchAppServerMetadata(): Promise<SharedServerMetadata | null>;
refreshAppServerMetadata(options?: { forceSkills?: boolean }): Promise<SharedServerMetadata | null>;
refreshActiveThreads(): Promise<readonly Thread[]>;
modelsSnapshot(): readonly ModelMetadata[] | null;
fetchModels(): Promise<readonly ModelMetadata[]>;
@ -71,21 +74,17 @@ export function createChatConnectionBundle(context: ChatConnectionBundleContext)
stateStore,
vaultPath,
currentClient,
setAppServerMetadata: (metadata) => {
threadCatalog.setAppServerMetadata(metadata);
},
modelsSnapshot: () => threadCatalog.modelsSnapshot(),
fetchModels: () => threadCatalog.fetchModels(),
refreshModels: () => threadCatalog.refreshModels(),
updateAppServerMetadata: (updater) => threadCatalog.updateAppServerMetadata(updater),
appServerMetadataSnapshot: () => threadCatalog.appServerMetadataSnapshot(),
fetchAppServerMetadata: () => threadCatalog.fetchAppServerMetadata(),
refreshAppServerMetadata: (options) => threadCatalog.refreshAppServerMetadata(options),
});
const serverDiagnostics = createChatServerDiagnosticsActions({
stateStore,
vaultPath,
currentClient,
setAppServerMetadata: (metadata) => {
threadCatalog.setAppServerMetadata(metadata);
},
serverMetadataSnapshot: () => serverMetadata.serverMetadataSnapshot(),
updateAppServerMetadata: (updater) => threadCatalog.updateAppServerMetadata(updater),
appServerMetadataSnapshot: () => threadCatalog.appServerMetadataSnapshot(),
});
const serverThreads = createChatServerThreadActions({
stateStore,
@ -114,8 +113,8 @@ export function createChatConnectionBundle(context: ChatConnectionBundleContext)
void serverMetadata.refreshPublishedRateLimits();
},
refreshSkills: (forceReload) => void serverMetadata.refreshPublishedSkills(forceReload),
setAppServerMetadata: () => {
serverMetadata.setAppServerMetadataSnapshot();
applyAppServerMetadataSnapshot: () => {
serverMetadata.applyAppServerMetadataSnapshot();
},
maybeNameThread: (threadId, turnId, completedSummary) => {
autoTitle.maybeAutoTitleThread(threadId, turnId, completedSummary);
@ -170,7 +169,7 @@ export function createChatConnectionBundle(context: ChatConnectionBundleContext)
refreshPublishedSkills: (forceReload) => serverMetadata.refreshPublishedSkills(forceReload),
},
diagnostics: {
refreshPublishedDiagnosticProbes: () => serverDiagnostics.refreshPublishedDiagnosticProbes(),
refreshPublishedDiagnosticProbes: (options) => serverDiagnostics.refreshPublishedDiagnosticProbes(options),
},
loadSharedThreadList,
scheduleDeferredDiagnostics: () => {

View file

@ -73,16 +73,18 @@ type ChatThreadCatalog = Pick<
| "refreshThreadsViewLiveState"
| "refreshFromOpenSurface"
| "setActiveThreads"
| "setAppServerMetadata"
| "updateAppServerMetadata"
| "refreshActiveThreads"
| "activeThreadsSnapshot"
| "appServerMetadataSnapshot"
| "fetchAppServerMetadata"
| "refreshAppServerMetadata"
| "modelsSnapshot"
| "fetchModels"
| "refreshModels"
| "observeActiveThreads"
| "observeAppServerMetadata"
| "observeModels"
| "observeActiveThreadsResult"
| "observeAppServerMetadataResult"
| "observeModelsResult"
>;
export interface ChatPanelEnvironment {

View file

@ -1,4 +1,5 @@
import type { AppServerClient } from "../../../app-server/connection/client";
import type { AppServerObservedQueryResult } from "../../../app-server/query/cache";
import type { ModelMetadata } from "../../../domain/catalog/metadata";
import type { Thread } from "../../../domain/threads/model";
@ -170,14 +171,26 @@ export class ChatPanelSession implements ChatSurfaceHandle {
this.refreshTabHeader();
}
private receiveObservedThreadResult(result: AppServerObservedQueryResult<readonly Thread[]>): void {
if (result.data) this.receiveObservedThreads(result.data);
}
private receiveObservedAppServerMetadata(metadata: SharedServerMetadata): void {
this.parts.serverActions.metadata.applyAppServerMetadata(metadata);
}
private receiveObservedAppServerMetadataResult(result: AppServerObservedQueryResult<SharedServerMetadata>): void {
if (result.data) this.receiveObservedAppServerMetadata(result.data);
}
private receiveObservedModels(models: readonly ModelMetadata[]): void {
this.dispatch({ type: "connection/metadata-applied", availableModels: models });
}
private receiveObservedModelsResult(result: AppServerObservedQueryResult<readonly ModelMetadata[]>): void {
if (result.data) this.receiveObservedModels(result.data);
}
openPanelSnapshot(): OpenCodexPanelSnapshot {
return {
viewId: this.environment.obsidian.viewId,
@ -276,21 +289,21 @@ export class ChatPanelSession implements ChatSurfaceHandle {
this.unsubscribeAppServerState();
this.applyCachedAppServerState();
this.appServerStateUnsubscribers.push(
this.environment.plugin.threadCatalog.observeActiveThreads(
(threads) => {
this.receiveObservedThreads(threads);
this.environment.plugin.threadCatalog.observeActiveThreadsResult(
(result) => {
this.receiveObservedThreadResult(result);
},
{ emitCurrent: false },
),
this.environment.plugin.threadCatalog.observeAppServerMetadata(
(metadata) => {
this.receiveObservedAppServerMetadata(metadata);
this.environment.plugin.threadCatalog.observeAppServerMetadataResult(
(result) => {
this.receiveObservedAppServerMetadataResult(result);
},
{ emitCurrent: false },
),
this.environment.plugin.threadCatalog.observeModels(
(models) => {
this.receiveObservedModels(models);
this.environment.plugin.threadCatalog.observeModelsResult(
(result) => {
this.receiveObservedModelsResult(result);
},
{ emitCurrent: false },
),

View file

@ -1,6 +1,7 @@
import { Notice } from "obsidian";
import type { AppServerClient } from "../../app-server/connection/client";
import type { AppServerObservedQueryResult } from "../../app-server/query/cache";
import { ConnectionManager, type ConnectionManagerHandlers, StaleConnectionError } from "../../app-server/connection/connection-manager";
import type { Thread } from "../../domain/threads/model";
import type { CodexPanelSettings } from "../../settings/model";
@ -45,7 +46,7 @@ type ThreadsThreadCatalog = Pick<
| "refreshFromOpenSurface"
| "refreshActiveThreads"
| "activeThreadsSnapshot"
| "observeActiveThreads"
| "observeActiveThreadsResult"
>;
export interface CodexThreadsSessionEnvironment {
@ -134,8 +135,8 @@ export class CodexThreadsSession {
if (activeThreadsSnapshot) {
this.threads = activeThreadsSnapshot;
}
this.unsubscribeThreads = this.host.threadCatalog.observeActiveThreads((threads) => {
this.receiveObservedThreads(threads);
this.unsubscribeThreads = this.host.threadCatalog.observeActiveThreadsResult((result) => {
this.receiveObservedThreadsResult(result);
});
this.render();
void this.refresh();
@ -181,6 +182,22 @@ export class CodexThreadsSession {
this.render();
}
private receiveObservedThreadsResult(result: AppServerObservedQueryResult<readonly Thread[]>): void {
if (result.data) {
this.receiveObservedThreads(result.data);
return;
}
if (result.isFetching && this.threads.length === 0) {
this.status = { kind: "loading", message: "Loading threads..." };
this.render();
return;
}
if (result.error && this.threads.length === 0) {
this.status = { kind: "error", message: result.error.message };
this.render();
}
}
private get host(): CodexThreadsHost {
return this.environment.host;
}

View file

@ -1,4 +1,5 @@
import type { AppServerClient } from "../app-server/connection/client";
import type { AppServerObservedQueryResult } from "../app-server/query/cache";
import { withShortLivedAppServerClient } from "../app-server/connection/short-lived-client";
import { setHookItemEnabled, trustHookItem } from "../app-server/services/catalog";
import { restoreArchivedThread as restoreArchivedThreadOnAppServer } from "../app-server/services/threads";
@ -26,7 +27,12 @@ export interface SettingsDynamicDataHost {
type SettingsThreadCatalog = Pick<
SharedThreadCatalog,
"refreshFromOpenSurface" | "modelsSnapshot" | "observeModels" | "fetchModels" | "refreshModels" | "notifyAppServerQueryContextChanged"
| "refreshFromOpenSurface"
| "modelsSnapshot"
| "observeModelsResult"
| "fetchModels"
| "refreshModels"
| "notifyAppServerQueryContextChanged"
>;
interface SettingsDynamicDataControllerCallbacks {
@ -70,10 +76,9 @@ export class SettingsDynamicDataController {
activate(): void {
if (this.unsubscribeModels) return;
this.models = [...(this.host.threadCatalog.modelsSnapshot() ?? [])];
this.unsubscribeModels = this.host.threadCatalog.observeModels(
(models) => {
this.models = [...models];
this.callbacks.display();
this.unsubscribeModels = this.host.threadCatalog.observeModelsResult(
(result) => {
this.receiveObservedModelsResult(result);
},
{ emitCurrent: false },
);
@ -104,6 +109,12 @@ export class SettingsDynamicDataController {
this.unsubscribeModels = null;
}
private receiveObservedModelsResult(result: AppServerObservedQueryResult<readonly ModelMetadata[]>): void {
if (!result.data) return;
this.models = [...result.data];
this.callbacks.display();
}
async refreshSettingsData(options: { forceModels?: boolean } = {}): Promise<void> {
this.settingsDataAutoLoadStarted = true;
const operationId = this.nextSettingsDynamicOperationId();

View file

@ -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 "../app-server/query/cache";
import type { AppServerObservedQueryResult, AppServerQueryCache } from "../app-server/query/cache";
import { appServerQueryContextMatches, cloneAppServerQueryContext, type AppServerQueryContext } from "../app-server/query/keys";
interface ThreadSurfaceActions {
@ -38,9 +38,12 @@ export class SharedThreadCatalog {
this.options.cache.setActiveThreads(this.context(), threads);
}
observeActiveThreads(listener: (threads: readonly Thread[]) => void, options?: { emitCurrent?: boolean }): () => void {
observeActiveThreadsResult(
listener: (result: AppServerObservedQueryResult<readonly Thread[]>) => void,
options?: { emitCurrent?: boolean },
): () => void {
return this.observeCurrentContext(
(context, contextListener, observeOptions) => this.options.cache.observeActiveThreads(context, contextListener, observeOptions),
(context, contextListener, observeOptions) => this.options.cache.observeActiveThreadsResult(context, contextListener, observeOptions),
listener,
options,
);
@ -50,13 +53,25 @@ export class SharedThreadCatalog {
return this.options.cache.appServerMetadataSnapshot(this.context());
}
setAppServerMetadata(metadata: SharedServerMetadata): void {
this.options.cache.setAppServerMetadata(this.context(), metadata);
updateAppServerMetadata(updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null): SharedServerMetadata | null {
return this.options.cache.updateAppServerMetadata(this.context(), updater);
}
observeAppServerMetadata(listener: (metadata: SharedServerMetadata) => void, options?: { emitCurrent?: boolean }): () => void {
async fetchAppServerMetadata(): Promise<SharedServerMetadata | null> {
return this.options.cache.fetchAppServerMetadata(this.context());
}
async refreshAppServerMetadata(options: { forceSkills?: boolean } = {}): Promise<SharedServerMetadata | null> {
return this.options.cache.refreshAppServerMetadata(this.context(), options);
}
observeAppServerMetadataResult(
listener: (result: AppServerObservedQueryResult<SharedServerMetadata>) => void,
options?: { emitCurrent?: boolean },
): () => void {
return this.observeCurrentContext(
(context, contextListener, observeOptions) => this.options.cache.observeAppServerMetadata(context, contextListener, observeOptions),
(context, contextListener, observeOptions) =>
this.options.cache.observeAppServerMetadataResult(context, contextListener, observeOptions),
listener,
options,
);
@ -74,9 +89,12 @@ export class SharedThreadCatalog {
return this.options.cache.refreshModels(this.context());
}
observeModels(listener: (models: readonly ModelMetadata[]) => void, options?: { emitCurrent?: boolean }): () => void {
observeModelsResult(
listener: (result: AppServerObservedQueryResult<readonly ModelMetadata[]>) => void,
options?: { emitCurrent?: boolean },
): () => void {
return this.observeCurrentContext(
(context, contextListener, observeOptions) => this.options.cache.observeModels(context, contextListener, observeOptions),
(context, contextListener, observeOptions) => this.options.cache.observeModelsResult(context, contextListener, observeOptions),
listener,
options,
);

View file

@ -12,6 +12,7 @@ import { emptyRuntimeConfigSnapshot, type RuntimeConfigSnapshot } from "../../sr
import type { RateLimitSnapshot } from "../../src/app-server/protocol/runtime-metrics";
import type { SharedServerMetadata } from "../../src/app-server/query/snapshots";
import type { ModelMetadata, SkillMetadata } from "../../src/domain/catalog/metadata";
import type { CatalogModel, CatalogSkillMetadata } from "../../src/app-server/protocol/catalog";
describe("AppServerQueryCache", () => {
it("stores metadata snapshots without replacing failed resource values with stale data", () => {
@ -22,10 +23,10 @@ describe("AppServerQueryCache", () => {
rateLimit: rateLimit(42),
});
cache.setAppServerMetadata(context, goodMetadata);
cache.writeAppServerMetadata(context, goodMetadata);
expect(cache.appServerMetadataSnapshot(context)?.availableModels.map((model) => model.model)).toEqual(["gpt-5.5"]);
cache.setAppServerMetadata(
cache.writeAppServerMetadata(
context,
metadata({
availableModels: [modelMetadata("gpt-5.6")],
@ -42,7 +43,7 @@ describe("AppServerQueryCache", () => {
expect(cache.appServerMetadataSnapshot(context)?.serverDiagnostics.probes["skills/list"].status).toBe("failed");
expect(cache.appServerMetadataSnapshot(context)?.serverDiagnostics.probes["account/rateLimits/read"].status).toBe("failed");
cache.setAppServerMetadata(context, metadata({ availableModels: [], modelProbeStatus: "failed" }));
cache.writeAppServerMetadata(context, metadata({ availableModels: [], modelProbeStatus: "failed" }));
expect(cache.appServerMetadataSnapshot(context)?.availableModels.map((model) => model.model)).toEqual(["gpt-5.6"]);
expect(cache.appServerMetadataSnapshot(context)?.serverDiagnostics.probes["model/list"].status).toBe("failed");
});
@ -51,7 +52,7 @@ describe("AppServerQueryCache", () => {
const cache = new AppServerQueryCache();
const context = cacheContext();
cache.setAppServerMetadata(
cache.writeAppServerMetadata(
context,
metadata({
availableModels: [modelMetadata("failed-model")],
@ -77,7 +78,7 @@ describe("AppServerQueryCache", () => {
await expect(cache.fetchActiveThreads(context)).resolves.toEqual([]);
cache.setActiveThreads(context, [thread("applied")]);
cache.setAppServerMetadata(context, metadata());
cache.writeAppServerMetadata(context, metadata());
expect(cache.activeThreadsSnapshot(context)).toBeNull();
expect(cache.appServerMetadataSnapshot(context)).toBeNull();
@ -88,8 +89,8 @@ describe("AppServerQueryCache", () => {
const cache = new AppServerQueryCache();
const context = cacheContext();
cache.setAppServerMetadata(context, metadata({ availableModels: [modelMetadata("gpt-5.6")] }));
cache.setAppServerMetadata(context, metadata({ availableModels: [] }));
cache.writeAppServerMetadata(context, metadata({ availableModels: [modelMetadata("gpt-5.6")] }));
cache.writeAppServerMetadata(context, metadata({ availableModels: [] }));
expect(cache.appServerMetadataSnapshot(context)?.availableModels).toEqual([]);
expect(cache.modelsSnapshot(context)).toEqual([]);
@ -99,7 +100,7 @@ describe("AppServerQueryCache", () => {
const cache = new AppServerQueryCache();
const context = cacheContext();
cache.setAppServerMetadata(context, metadata());
cache.writeAppServerMetadata(context, metadata());
expect(cache.appServerMetadataSnapshot(cacheContext({ vaultPath: "/other-vault" }))).toBeNull();
expect(cache.appServerMetadataSnapshot(cacheContext({ codexPath: "/opt/codex" }))).toBeNull();
@ -160,6 +161,85 @@ describe("AppServerQueryCache", () => {
expect(cache.activeThreadsSnapshot(capturedContext)?.map((item) => item.id)).toEqual(["captured"]);
expect(cache.activeThreadsSnapshot(context)).toBeNull();
});
it("fetches app-server metadata through the query cache and syncs model snapshots", async () => {
const context = cacheContext();
const cache = cacheWithClient({
readEffectiveConfig: vi.fn().mockResolvedValue({}),
listModels: vi.fn().mockResolvedValue({ data: [catalogModel("gpt-meta")] }),
listSkills: vi.fn().mockResolvedValue({ data: [{ skills: [catalogSkill("writer")] }] }),
readAccountRateLimits: vi.fn().mockResolvedValue({ rateLimits: appServerRateLimit(64), rateLimitsByLimitId: null }),
});
const metadata = await cache.refreshAppServerMetadata(context);
expect(metadata?.availableModels.map((model) => model.model)).toEqual(["gpt-meta"]);
expect(metadata?.availableSkills.map((skill) => skill.name)).toEqual(["writer"]);
expect(metadata?.rateLimit?.primary?.usedPercent).toBe(64);
expect(metadata?.serverDiagnostics.probes["model/list"].status).toBe("ok");
expect(cache.modelsSnapshot(context)?.map((model) => model.model)).toEqual(["gpt-meta"]);
});
it("shares in-flight model fetches between metadata and models queries", async () => {
const context = cacheContext();
const modelRefresh = deferred<{ data: CatalogModel[] }>();
const listModels = vi.fn(() => modelRefresh.promise);
const cache = cacheWithClient({
readEffectiveConfig: vi.fn().mockResolvedValue({}),
listModels,
listSkills: vi.fn().mockResolvedValue({ data: [{ skills: [] }] }),
readAccountRateLimits: vi.fn().mockResolvedValue({ rateLimits: appServerRateLimit(0), rateLimitsByLimitId: null }),
});
const metadataPromise = cache.refreshAppServerMetadata(context);
await flushMicrotasks();
const modelsPromise = cache.fetchModels(context);
await flushMicrotasks();
expect(listModels).toHaveBeenCalledOnce();
modelRefresh.resolve({ data: [catalogModel("gpt-shared")] });
await expect(modelsPromise).resolves.toMatchObject([{ model: "gpt-shared" }]);
await expect(metadataPromise).resolves.toMatchObject({
availableModels: [{ model: "gpt-shared" }],
});
expect(listModels).toHaveBeenCalledOnce();
expect(cache.modelsSnapshot(context)?.map((model) => model.model)).toEqual(["gpt-shared"]);
});
it("keeps query-cached models when app-server metadata model refresh fails", async () => {
const context = cacheContext();
const cache = cacheWithClient({
readEffectiveConfig: vi.fn().mockResolvedValue({}),
listModels: vi.fn().mockRejectedValue(new Error("offline")),
listSkills: vi.fn().mockResolvedValue({ data: [{ skills: [] }] }),
readAccountRateLimits: vi.fn().mockResolvedValue({ rateLimits: appServerRateLimit(0), rateLimitsByLimitId: null }),
});
cache.writeAppServerMetadata(context, metadata({ availableModels: [modelMetadata("gpt-cached")] }));
const metadataSnapshot = await cache.refreshAppServerMetadata(context);
expect(metadataSnapshot?.availableModels.map((model) => model.model)).toEqual(["gpt-cached"]);
expect(metadataSnapshot?.serverDiagnostics.probes["model/list"].status).toBe("failed");
expect(cache.modelsSnapshot(context)?.map((model) => model.model)).toEqual(["gpt-cached"]);
});
it("tracks optimistic active thread updates in the mutation cache and clears them by context", () => {
const cache = new AppServerQueryCache();
const context = cacheContext();
cache.setActiveThreads(context, [thread("thread")]);
cache.updateActiveThreads(context, (threads) => threads?.map((item) => ({ ...item, name: "Renamed" })) ?? null);
expect(cache.activeThreadsSnapshot(context)).toEqual([{ ...thread("thread"), name: "Renamed" }]);
expect(cache.client.getMutationCache().getAll()).toHaveLength(1);
cache.clearContext(context);
expect(cache.client.getMutationCache().getAll()).toHaveLength(0);
expect(cache.activeThreadsSnapshot(context)).toBeNull();
});
});
function cacheContext(overrides: Partial<AppServerQueryContext> = {}): AppServerQueryContext {
@ -187,6 +267,14 @@ function cacheWithThreads(
});
}
function cacheWithClient(client: Record<string, unknown>): AppServerQueryCache {
return new AppServerQueryCache({
clientRunner: {
runWithClient: async (_context, operation) => operation(client as never),
},
});
}
function metadata(
overrides: {
availableModels?: readonly ModelMetadata[];
@ -258,6 +346,43 @@ function modelMetadata(model: string): ModelMetadata {
};
}
function catalogModel(model: string): CatalogModel {
return {
id: model,
model,
displayName: model,
description: "",
hidden: false,
supportedReasoningEfforts: [],
defaultReasoningEffort: "medium",
inputModalities: ["text"],
additionalSpeedTiers: [],
serviceTiers: [],
defaultServiceTier: null,
isDefault: false,
};
}
function catalogSkill(name: string): CatalogSkillMetadata {
return {
name,
description: "",
path: `/tmp/${name}`,
enabled: true,
};
}
function appServerRateLimit(usedPercent: number): RateLimitSnapshot {
return {
limitId: "codex",
limitName: "Codex",
primary: { usedPercent, windowDurationMins: 300, resetsAt: null },
secondary: null,
individualLimit: null,
rateLimitReachedType: null,
};
}
function thread(id: string) {
return {
id,
@ -281,3 +406,9 @@ function deferred<T>(): Deferred<T> {
});
return { promise, resolve };
}
async function flushMicrotasks(): Promise<void> {
for (let index = 0; index < 10; index += 1) {
await Promise.resolve();
}
}

View file

@ -81,13 +81,23 @@ describe("ChatConnectionController", () => {
expect(host.setStatus).toHaveBeenCalledWith("Connected.", { kind: "connected" });
});
it("refreshes diagnostics after clearing deferred diagnostics", async () => {
const { controller, host, refreshPublishedDiagnosticProbes } = createController({ connected: true });
it("refreshes metadata before metadata-backed diagnostics", async () => {
const { controller, host, refreshPublishedAppServerMetadata, refreshPublishedDiagnosticProbes } = createController({ connected: true });
await controller.refreshDiagnostics();
expect(host.clearDeferredDiagnostics).toHaveBeenCalledTimes(2);
expect(refreshPublishedDiagnosticProbes).toHaveBeenCalledOnce();
expect(refreshPublishedAppServerMetadata).toHaveBeenCalledOnce();
expect(refreshPublishedDiagnosticProbes).toHaveBeenCalledWith({ appServerMetadataSnapshot: true });
});
it("refreshes active threads without refreshing metadata", async () => {
const { controller, host, refreshPublishedAppServerMetadata } = createController({ connected: true });
await controller.fetchActiveThreads();
expect(host.loadSharedThreadList).toHaveBeenCalledOnce();
expect(refreshPublishedAppServerMetadata).not.toHaveBeenCalled();
});
it("clears disconnected connection state on server exit while keeping last startup metadata", () => {

View file

@ -1,7 +1,14 @@
import { describe, expect, it, vi } from "vitest";
import type { AppServerClient } from "../../../../../src/app-server/connection/client";
import type { McpServerStatus } from "../../../../../src/domain/server/diagnostics";
import {
createServerDiagnostics,
diagnosticProbeError,
diagnosticProbeOk,
diagnosticsWithProbe,
type McpServerStatus,
} from "../../../../../src/domain/server/diagnostics";
import type { SharedServerMetadata } from "../../../../../src/domain/server/metadata";
import { emptyRuntimeConfigSnapshot } from "../../../../../src/app-server/protocol/runtime-config";
import type { RateLimitSnapshot } from "../../../../../src/app-server/protocol/runtime-metrics";
import { threadFromThreadRecord } from "../../../../../src/app-server/protocol/thread";
@ -243,47 +250,52 @@ describe("chat server actions", () => {
const state = chatStateFixture();
const stateStore = createChatStateStore(state);
const fetchModels = vi.fn().mockResolvedValue(modelMetadataFromCatalogModels([modelFixture("gpt-5.1")]));
const listSkills = vi.fn().mockResolvedValue({ data: [{ skills: [skillFixture("writer")] }] });
const readAccountRateLimits = vi.fn().mockResolvedValue({ rateLimits: {} as RateLimitSnapshot });
const listHooks = vi.fn().mockResolvedValue({ data: [{ cwd: "/vault", hooks: [] }] });
const refreshedMetadata = serverMetadataFixture({
availableModels: modelMetadataFromCatalogModels([modelFixture("gpt-5.1")]),
availableSkills: [{ name: "writer", description: "", path: "/tmp/writer", enabled: true }],
rateLimit: rateLimitFixture(),
serverDiagnostics: diagnosticsWithProbe(
diagnosticsWithProbe(
diagnosticsWithProbe(createServerDiagnostics(), diagnosticProbeOk("model/list", "1 models")),
diagnosticProbeOk("skills/list", "1 skills"),
),
diagnosticProbeOk("account/rateLimits/read", "available"),
),
});
const refreshAppServerMetadata = vi.fn<() => Promise<SharedServerMetadata | null>>().mockResolvedValue(refreshedMetadata);
const client = {
readEffectiveConfig: vi.fn().mockResolvedValue({}),
listSkills,
readAccountRateLimits,
listHooks,
listMcpServerStatus: vi.fn().mockResolvedValue({ data: [] }),
listCollaborationModes: vi.fn().mockResolvedValue({ data: [] }),
readModelProviderCapabilities: vi.fn().mockResolvedValue({}),
} as unknown as AppServerClient;
const metadataCache = metadataCacheHost({ current: null });
const metadata = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
setAppServerMetadata: () => undefined,
modelsSnapshot: () => null,
fetchModels,
refreshModels: async () => [],
...metadataCache,
fetchAppServerMetadata: async () => refreshedMetadata,
refreshAppServerMetadata: async () => {
const next = await refreshAppServerMetadata();
metadataCache.updateAppServerMetadata(() => next);
return next;
},
});
const diagnostics = createChatServerDiagnosticsActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
setAppServerMetadata: () => undefined,
serverMetadataSnapshot: () => metadata.serverMetadataSnapshot(),
...metadataCache,
});
await metadata.refreshAppServerMetadata();
fetchModels.mockClear();
listSkills.mockClear();
readAccountRateLimits.mockClear();
await diagnostics.refreshDiagnosticProbes({ appServerMetadataSnapshot: true });
expect(fetchModels).not.toHaveBeenCalled();
expect(listSkills).not.toHaveBeenCalled();
expect(readAccountRateLimits).not.toHaveBeenCalled();
expect(refreshAppServerMetadata).toHaveBeenCalledOnce();
expect(listHooks).toHaveBeenCalledWith("/vault");
expect(stateStore.getState().connection.serverDiagnostics.probes["model/list"]).toMatchObject({
status: "ok",
@ -299,6 +311,77 @@ describe("chat server actions", () => {
});
});
it("uses metadata diagnostics as the default resource probe source", async () => {
const stateStore = createChatStateStore(chatStateFixture());
const metadataCache = metadataCacheHost({
current: serverMetadataFixture({
serverDiagnostics: diagnosticsWithProbe(createServerDiagnostics(), diagnosticProbeOk("model/list", "cached models")),
}),
});
const listModels = vi.fn().mockResolvedValue({ data: [modelFixture("gpt-direct")] });
const listSkills = vi.fn().mockResolvedValue({ data: [{ skills: [skillFixture("direct-skill")] }] });
const readAccountRateLimits = vi.fn().mockResolvedValue({ rateLimits: rateLimitFixture(), rateLimitsByLimitId: null });
const listHooks = vi.fn().mockResolvedValue({ data: [{ cwd: "/vault", hooks: [] }] });
const client = {
listModels,
listSkills,
readAccountRateLimits,
listHooks,
listMcpServerStatus: vi.fn().mockResolvedValue({ data: [] }),
listCollaborationModes: vi.fn().mockResolvedValue({ data: [] }),
readModelProviderCapabilities: vi.fn().mockResolvedValue({}),
} as unknown as AppServerClient;
const diagnostics = createChatServerDiagnosticsActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
...metadataCache,
});
await diagnostics.refreshDiagnosticProbes();
expect(listModels).not.toHaveBeenCalled();
expect(listSkills).not.toHaveBeenCalled();
expect(readAccountRateLimits).not.toHaveBeenCalled();
expect(stateStore.getState().connection.serverDiagnostics.probes["model/list"]).toMatchObject({
status: "ok",
summary: "cached models",
});
expect(listHooks).toHaveBeenCalledWith("/vault");
});
it("can force resource probes for explicit health checks", async () => {
const stateStore = createChatStateStore(chatStateFixture());
const listModels = vi.fn().mockResolvedValue({ data: [modelFixture("gpt-direct")] });
const listSkills = vi.fn().mockResolvedValue({ data: [{ skills: [skillFixture("direct-skill")] }] });
const readAccountRateLimits = vi.fn().mockResolvedValue({ rateLimits: rateLimitFixture(), rateLimitsByLimitId: null });
const client = {
listModels,
listSkills,
readAccountRateLimits,
listHooks: vi.fn().mockResolvedValue({ data: [{ cwd: "/vault", hooks: [] }] }),
listMcpServerStatus: vi.fn().mockResolvedValue({ data: [] }),
listCollaborationModes: vi.fn().mockResolvedValue({ data: [] }),
readModelProviderCapabilities: vi.fn().mockResolvedValue({}),
} as unknown as AppServerClient;
const diagnostics = createChatServerDiagnosticsActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
...metadataCacheHost(),
});
await diagnostics.refreshDiagnosticProbes({ forceResourceProbes: true });
expect(listModels).toHaveBeenCalledWith(false);
expect(listSkills).toHaveBeenCalledWith("/vault");
expect(readAccountRateLimits).toHaveBeenCalledOnce();
expect(stateStore.getState().connection.serverDiagnostics.probes["model/list"]).toMatchObject({
status: "ok",
summary: "1 models",
});
});
it("does not apply or publish diagnostic probes after the client changes", async () => {
const stateStore = createChatStateStore(chatStateFixture());
const hooksRefresh = deferred<{ data: { cwd: string; hooks: unknown[] }[] }>();
@ -311,22 +394,13 @@ describe("chat server actions", () => {
} as unknown as AppServerClient;
const secondClient = {} as unknown as AppServerClient;
let currentClient = firstClient;
const setAppServerMetadata = vi.fn();
const metadata = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
setAppServerMetadata: () => undefined,
modelsSnapshot: () => null,
fetchModels: async () => [],
refreshModels: async () => [],
});
const updateAppServerMetadata = vi.fn(() => null);
const diagnostics = createChatServerDiagnosticsActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
setAppServerMetadata,
serverMetadataSnapshot: () => metadata.serverMetadataSnapshot(),
appServerMetadataSnapshot: () => null,
updateAppServerMetadata,
});
const refreshing = diagnostics.refreshPublishedDiagnosticProbes({ appServerMetadataSnapshot: true });
@ -338,100 +412,68 @@ describe("chat server actions", () => {
expect(stateStore.getState().connection.serverDiagnostics.probes["hooks/list"].status).toBe("unknown");
expect(stateStore.getState().connection.serverDiagnostics.probes["mcpServerStatus/list"].status).toBe("unknown");
expect(stateStore.getState().connection.serverDiagnostics.mcpServers).toEqual([]);
expect(setAppServerMetadata).not.toHaveBeenCalled();
expect(updateAppServerMetadata).not.toHaveBeenCalled();
});
it("loads one app-server metadata snapshot from the initially captured client", async () => {
const stateStore = createChatStateStore(chatStateFixture());
const readEffectiveConfig = deferred<Record<string, never>>();
const fetchModels = vi.fn().mockResolvedValue(modelMetadataFromCatalogModels([modelFixture("gpt-first")]));
const firstListSkills = vi.fn().mockResolvedValue({ data: [{ skills: [skillFixture("first-skill")] }] });
const firstReadAccountRateLimits = vi.fn().mockResolvedValue({ rateLimits: rateLimitFixture() });
const secondListSkills = vi.fn().mockResolvedValue({ data: [{ skills: [skillFixture("second-skill")] }] });
const secondReadAccountRateLimits = vi.fn().mockResolvedValue({ rateLimits: rateLimitFixture({ limitName: "Second" }) });
const firstClient = {
readEffectiveConfig: vi.fn().mockReturnValue(readEffectiveConfig.promise),
listSkills: firstListSkills,
readAccountRateLimits: firstReadAccountRateLimits,
} as unknown as AppServerClient;
const secondClient = {
listSkills: secondListSkills,
readAccountRateLimits: secondReadAccountRateLimits,
} as unknown as AppServerClient;
let currentClient = firstClient;
const fetchAppServerMetadata = vi.fn().mockResolvedValue(
serverMetadataFixture({
availableModels: modelMetadataFromCatalogModels([modelFixture("gpt-first")]),
availableSkills: [{ name: "first-skill", description: "", path: "/tmp/first-skill", enabled: true }],
}),
);
const controller = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
setAppServerMetadata: () => undefined,
modelsSnapshot: () => null,
fetchModels,
refreshModels: async () => [],
currentClient: () => ({}) as AppServerClient,
...metadataCacheHost(),
fetchAppServerMetadata,
refreshAppServerMetadata: async () => null,
});
const loading = controller.loadAppServerMetadata();
currentClient = secondClient;
readEffectiveConfig.resolve({});
await expect(loading).resolves.toMatchObject({
availableModels: [{ model: "gpt-first" }],
availableSkills: [{ name: "first-skill" }],
});
expect(fetchModels).toHaveBeenCalledOnce();
expect(firstListSkills).toHaveBeenCalledOnce();
expect(firstReadAccountRateLimits).toHaveBeenCalledOnce();
expect(secondListSkills).not.toHaveBeenCalled();
expect(secondReadAccountRateLimits).not.toHaveBeenCalled();
expect(fetchAppServerMetadata).toHaveBeenCalledOnce();
});
it("does not apply or publish app-server metadata when the client changes before refresh completes", async () => {
const stateStore = createChatStateStore(chatStateFixture());
const readEffectiveConfig = deferred<Record<string, never>>();
const firstClient = {
readEffectiveConfig: vi.fn().mockReturnValue(readEffectiveConfig.promise),
listModels: vi.fn().mockResolvedValue({ data: [modelFixture("gpt-stale")] }),
listSkills: vi.fn().mockResolvedValue({ data: [{ skills: [skillFixture("stale-skill")] }] }),
readAccountRateLimits: vi.fn().mockResolvedValue({ rateLimits: rateLimitFixture() }),
} as unknown as AppServerClient;
const secondClient = {} as unknown as AppServerClient;
let currentClient = firstClient;
const setAppServerMetadata = vi.fn();
const refreshAppServerMetadata = vi.fn().mockResolvedValue(null);
const controller = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
setAppServerMetadata,
modelsSnapshot: () => null,
fetchModels: async () => modelMetadataFromCatalogModels([modelFixture("gpt-stale")]),
refreshModels: async () => modelMetadataFromCatalogModels([modelFixture("gpt-stale")]),
currentClient: () => ({}) as AppServerClient,
...metadataCacheHost(),
fetchAppServerMetadata: async () => null,
refreshAppServerMetadata,
});
const refreshing = controller.refreshPublishedAppServerMetadata();
currentClient = secondClient;
readEffectiveConfig.resolve({});
await expect(refreshing).resolves.toBeNull();
expect(stateStore.getState().connection.availableModels).toEqual([]);
expect(stateStore.getState().connection.availableSkills).toEqual([]);
expect(setAppServerMetadata).not.toHaveBeenCalled();
});
it("keeps query-cached models visible when metadata model refresh fails", async () => {
const state = chatStateFixture();
const stateStore = createChatStateStore(state);
const client = {
readEffectiveConfig: vi.fn().mockResolvedValue({}),
listSkills: vi.fn().mockResolvedValue({ data: [{ skills: [] }] }),
readAccountRateLimits: vi.fn().mockResolvedValue({ rateLimits: rateLimitFixture() }),
} as unknown as AppServerClient;
const metadata = serverMetadataFixture({
availableModels: modelMetadataFromCatalogModels([modelFixture("gpt-cached")]),
serverDiagnostics: diagnosticsWithProbe(createServerDiagnostics(), diagnosticProbeError("model/list", new Error("offline"))),
});
const controller = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
setAppServerMetadata: () => undefined,
modelsSnapshot: () => modelMetadataFromCatalogModels([modelFixture("gpt-cached")]),
fetchModels: () => Promise.reject(new Error("offline")),
refreshModels: () => Promise.reject(new Error("offline")),
currentClient: () => ({}) as AppServerClient,
...metadataCacheHost({ current: metadata }),
fetchAppServerMetadata: async () => metadata,
refreshAppServerMetadata: async () => metadata,
});
await controller.refreshAppServerMetadata();
@ -444,19 +486,17 @@ describe("chat server actions", () => {
let state = chatStateFixture();
state = chatStateWith(state, { connection: { availableModels: modelMetadataFromCatalogModels([modelFixture("gpt-state-only")]) } });
const stateStore = createChatStateStore(state);
const client = {
readEffectiveConfig: vi.fn().mockResolvedValue({}),
listSkills: vi.fn().mockResolvedValue({ data: [{ skills: [] }] }),
readAccountRateLimits: vi.fn().mockResolvedValue({ rateLimits: rateLimitFixture() }),
} as unknown as AppServerClient;
const metadata = serverMetadataFixture({
availableModels: [],
serverDiagnostics: diagnosticsWithProbe(createServerDiagnostics(), diagnosticProbeError("model/list", new Error("offline"))),
});
const controller = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
setAppServerMetadata: () => undefined,
modelsSnapshot: () => null,
fetchModels: () => Promise.reject(new Error("offline")),
refreshModels: () => Promise.reject(new Error("offline")),
currentClient: () => ({}) as AppServerClient,
...metadataCacheHost({ current: metadata }),
fetchAppServerMetadata: async () => metadata,
refreshAppServerMetadata: async () => metadata,
});
await controller.refreshAppServerMetadata();
@ -472,15 +512,15 @@ describe("chat server actions", () => {
const firstClient = { listSkills } as unknown as AppServerClient;
const secondClient = {} as unknown as AppServerClient;
let currentClient = firstClient;
const setAppServerMetadata = vi.fn();
const updateAppServerMetadata = vi.fn(() => null);
const controller = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
setAppServerMetadata,
modelsSnapshot: () => null,
fetchModels: async () => [],
refreshModels: async () => [],
appServerMetadataSnapshot: () => null,
updateAppServerMetadata,
fetchAppServerMetadata: async () => null,
refreshAppServerMetadata: async () => null,
});
const refreshing = controller.refreshPublishedSkills(true);
@ -491,14 +531,14 @@ describe("chat server actions", () => {
expect(listSkills).toHaveBeenCalledWith("/vault", true);
expect(stateStore.getState().connection.availableSkills).toEqual([]);
expect(stateStore.getState().connection.serverDiagnostics.probes["skills/list"].status).toBe("unknown");
expect(setAppServerMetadata).not.toHaveBeenCalled();
expect(updateAppServerMetadata).not.toHaveBeenCalled();
});
it("publishes refreshed rate limits from sparse update notifications", async () => {
const state = chatStateFixture();
const stateStore = createChatStateStore(state);
const rateLimit = rateLimitFixture({ primary: { usedPercent: 64, windowDurationMins: 300, resetsAt: null } });
const setAppServerMetadata = vi.fn();
const cachedMetadata = { current: serverMetadataFixture() as SharedServerMetadata | null };
const client = {
readAccountRateLimits: vi.fn().mockResolvedValue({ rateLimits: rateLimit, rateLimitsByLimitId: null }),
} as unknown as AppServerClient;
@ -506,16 +546,15 @@ describe("chat server actions", () => {
stateStore,
vaultPath: "/vault",
currentClient: () => client,
setAppServerMetadata,
modelsSnapshot: () => null,
fetchModels: async () => [],
refreshModels: async () => [],
...metadataCacheHost(cachedMetadata),
fetchAppServerMetadata: async () => null,
refreshAppServerMetadata: async () => null,
});
await controller.refreshPublishedRateLimits();
expect(stateStore.getState().connection.rateLimit).toMatchObject({ primary: { usedPercent: 64 } });
expect(setAppServerMetadata).toHaveBeenCalledWith(expect.objectContaining({ rateLimit }));
expect(cachedMetadata.current?.rateLimit).toStrictEqual(rateLimit);
});
it("keeps the previous rate limit snapshot when sparse update refresh fails", async () => {
@ -526,7 +565,6 @@ describe("chat server actions", () => {
});
state = chatStateWith(state, { connection: { rateLimit: previousRateLimit } });
const stateStore = createChatStateStore(state);
const setAppServerMetadata = vi.fn();
const client = {
readAccountRateLimits: vi.fn().mockRejectedValue(new Error("offline")),
} as unknown as AppServerClient;
@ -534,17 +572,15 @@ describe("chat server actions", () => {
stateStore,
vaultPath: "/vault",
currentClient: () => client,
setAppServerMetadata,
modelsSnapshot: () => null,
fetchModels: async () => [],
refreshModels: async () => [],
...metadataCacheHost(),
fetchAppServerMetadata: async () => null,
refreshAppServerMetadata: async () => null,
});
await controller.refreshPublishedRateLimits();
expect(stateStore.getState().connection.rateLimit).toBe(previousRateLimit);
expect(stateStore.getState().connection.serverDiagnostics.probes["account/rateLimits/read"]).toMatchObject({ status: "failed" });
expect(setAppServerMetadata).not.toHaveBeenCalled();
});
it("does not apply or publish sparse rate limit refreshes after the client changes", async () => {
@ -555,15 +591,15 @@ describe("chat server actions", () => {
} as unknown as AppServerClient;
const secondClient = {} as unknown as AppServerClient;
let currentClient = firstClient;
const setAppServerMetadata = vi.fn();
const updateAppServerMetadata = vi.fn(() => null);
const controller = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
setAppServerMetadata,
modelsSnapshot: () => null,
fetchModels: async () => [],
refreshModels: async () => [],
appServerMetadataSnapshot: () => null,
updateAppServerMetadata,
fetchAppServerMetadata: async () => null,
refreshAppServerMetadata: async () => null,
});
const refreshing = controller.refreshPublishedRateLimits();
@ -576,7 +612,7 @@ describe("chat server actions", () => {
await refreshing;
expect(stateStore.getState().connection.rateLimit).toBeNull();
expect(stateStore.getState().connection.serverDiagnostics.probes["account/rateLimits/read"].status).toBe("unknown");
expect(setAppServerMetadata).not.toHaveBeenCalled();
expect(updateAppServerMetadata).not.toHaveBeenCalled();
});
it("loads MCP status lines with cached startup diagnostics", async () => {
@ -587,21 +623,12 @@ describe("chat server actions", () => {
const client = {
listMcpServerStatus,
} as unknown as AppServerClient;
const metadata = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
setAppServerMetadata: () => undefined,
modelsSnapshot: () => null,
fetchModels: async () => [],
refreshModels: async () => [],
});
const metadataCache = metadataCacheHost({ current: serverMetadataFixture() });
const controller = createChatServerDiagnosticsActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
setAppServerMetadata: () => undefined,
serverMetadataSnapshot: () => metadata.serverMetadataSnapshot(),
...metadataCache,
});
controller.recordMcpStartupStatus("github", "ready", null);
@ -684,6 +711,31 @@ function rateLimitFixture(overrides: Partial<RateLimitSnapshot> = {}): RateLimit
};
}
function serverMetadataFixture(overrides: Partial<SharedServerMetadata> = {}): SharedServerMetadata {
return {
runtimeConfig: emptyRuntimeConfigSnapshot(),
availableModels: [],
availableSkills: [],
rateLimit: null,
serverDiagnostics: createServerDiagnostics(),
...overrides,
};
}
function metadataCacheHost(cache: { current: SharedServerMetadata | null } = { current: null }): {
appServerMetadataSnapshot: () => SharedServerMetadata | null;
updateAppServerMetadata: (updater: (metadata: SharedServerMetadata | null) => SharedServerMetadata | null) => SharedServerMetadata | null;
} {
return {
appServerMetadataSnapshot: () => cache.current,
updateAppServerMetadata: (updater) => {
const next = updater(cache.current);
cache.current = next;
return next;
},
};
}
function mcpServerStatus(): McpServerStatus {
return {
name: "github",

View file

@ -27,7 +27,7 @@ function controllerForState(
fetchActiveThreads: vi.fn(),
refreshRateLimits: vi.fn(),
refreshSkills: vi.fn(),
setAppServerMetadata: vi.fn(),
applyAppServerMetadataSnapshot: vi.fn(),
maybeNameThread: vi.fn(),
applyThreadArchived: vi.fn(),
applyThreadRenamed: vi.fn(),
@ -694,8 +694,8 @@ describe("ChatInboundController", () => {
it("records MCP startup status for diagnostics without a chat system message", () => {
const state = chatStateFixture();
const recordMcpStartupStatus = vi.fn();
const setAppServerMetadata = vi.fn();
const controller = controllerForState(state, { recordMcpStartupStatus, setAppServerMetadata });
const applyAppServerMetadataSnapshot = vi.fn();
const controller = controllerForState(state, { recordMcpStartupStatus, applyAppServerMetadataSnapshot });
controller.handleNotification({
method: "mcpServer/startupStatus/updated",
@ -708,7 +708,7 @@ describe("ChatInboundController", () => {
} satisfies Extract<ServerNotification, { method: "mcpServer/startupStatus/updated" }>);
expect(recordMcpStartupStatus).toHaveBeenCalledWith("github", "failed", "missing token");
expect(setAppServerMetadata).toHaveBeenCalledOnce();
expect(applyAppServerMetadataSnapshot).toHaveBeenCalledOnce();
expect(chatStateMessageStreamItems(controller.currentState())).toEqual([]);
});
});

View file

@ -4,6 +4,8 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { DEFAULT_SETTINGS } from "../../../src/settings/model";
import type { CodexChatHost } from "../../../src/features/chat/host/runtime";
import type { AppServerObservedQueryResult } from "../../../src/app-server/query/cache";
import { modelMetadataFromCatalogModels } from "../../../src/app-server/protocol/catalog";
import { createServerDiagnostics } from "../../../src/domain/server/diagnostics";
import type { Thread } from "../../../src/domain/threads/model";
import type { ModelMetadata } from "../../../src/domain/catalog/metadata";
@ -187,22 +189,14 @@ describe("CodexChatView connection lifecycle", () => {
});
});
it("publishes app-server metadata after connecting", async () => {
const setAppServerMetadata = vi.fn();
it("loads app-server metadata after connecting", async () => {
connectionMock.state.client = connectedClient();
const view = await chatView({
host: chatHost({ setAppServerMetadata }),
});
const view = await chatView();
await view.surface.connect();
expect(setAppServerMetadata).toHaveBeenCalledWith(
expect.objectContaining({
runtimeConfig: expect.any(Object),
availableModels: [],
availableSkills: [],
}),
);
expect(connectionMock.state.client["readEffectiveConfig"]).toHaveBeenCalledOnce();
expect(view.surface.openPanelSnapshot()).toMatchObject({ connected: true });
});
it("renders the chat shell on the view content root", async () => {
@ -1312,32 +1306,71 @@ interface ChatHostFixtureOverrides {
refreshFromOpenSurface?: CodexChatHost["threadCatalog"]["refreshFromOpenSurface"];
refreshThreadsViewLiveState?: CodexChatHost["threadCatalog"]["refreshThreadsViewLiveState"];
setActiveThreads?: CodexChatHost["threadCatalog"]["setActiveThreads"];
setAppServerMetadata?: CodexChatHost["threadCatalog"]["setAppServerMetadata"];
updateAppServerMetadata?: CodexChatHost["threadCatalog"]["updateAppServerMetadata"];
refreshActiveThreads?: CodexChatHost["threadCatalog"]["refreshActiveThreads"];
activeThreadsSnapshot?: CodexChatHost["threadCatalog"]["activeThreadsSnapshot"];
appServerMetadataSnapshot?: CodexChatHost["threadCatalog"]["appServerMetadataSnapshot"];
modelsSnapshot?: CodexChatHost["threadCatalog"]["modelsSnapshot"];
fetchModels?: CodexChatHost["threadCatalog"]["fetchModels"];
refreshModels?: CodexChatHost["threadCatalog"]["refreshModels"];
fetchAppServerMetadata?: CodexChatHost["threadCatalog"]["fetchAppServerMetadata"];
refreshAppServerMetadata?: CodexChatHost["threadCatalog"]["refreshAppServerMetadata"];
}
function chatHost(overrides: ChatHostFixtureOverrides = {}): CodexChatHost {
let activeThreads = overrides.activeThreadsSnapshot?.() ?? null;
let metadata = overrides.appServerMetadataSnapshot?.() ?? null;
const models = overrides.modelsSnapshot?.() ?? null;
const activeThreadListeners = new Set<(threads: readonly Thread[]) => void>();
const metadataListeners = new Set<(metadata: SharedServerMetadata) => void>();
const modelListeners = new Set<(models: readonly ModelMetadata[]) => void>();
const activeThreadResultListeners = new Set<(result: AppServerObservedQueryResult<readonly Thread[]>) => void>();
const metadataResultListeners = new Set<(result: AppServerObservedQueryResult<SharedServerMetadata>) => void>();
const modelResultListeners = new Set<(result: AppServerObservedQueryResult<readonly ModelMetadata[]>) => void>();
const settings = {
...DEFAULT_SETTINGS,
codexPath: "codex",
sendShortcut: "enter" as const,
...overrides.settings,
};
const vaultPath = overrides.vaultPath ?? "/vault";
const applyMetadataToCache = (nextMetadata: SharedServerMetadata): SharedServerMetadata => {
metadata = nextMetadata;
for (const listener of metadataResultListeners) listener(queryResult(nextMetadata));
for (const listener of modelResultListeners) listener(queryResult(nextMetadata.availableModels));
return nextMetadata;
};
const loadAppServerMetadata = async (): Promise<SharedServerMetadata | null> => {
const client = connectionMock.state.client as ReturnType<typeof baseClient> | null;
if (!client || !connectionMock.state.connected) return null;
const connectionStillCurrent = () => connectionMock.state.client === client && connectionMock.state.connected;
await client.readEffectiveConfig(vaultPath);
if (!connectionStillCurrent()) return null;
const availableModels = overrides.fetchModels
? await overrides.fetchModels()
: modelMetadataFromCatalogModels((await client.listModels(false)).data as Parameters<typeof modelMetadataFromCatalogModels>[0]);
if (!connectionStillCurrent()) return null;
const skillsResponse = await client.listSkills(vaultPath);
if (!connectionStillCurrent()) return null;
await client.readAccountRateLimits();
if (!connectionStillCurrent()) return null;
return {
runtimeConfig: emptyRuntimeConfigSnapshot(),
availableModels,
availableSkills: skillsResponse.data.flatMap(
(entry: { skills: { name: string; description?: string; path?: string; enabled?: boolean }[] }) =>
entry.skills.map((skill) => ({
name: skill.name,
description: skill.description ?? "",
path: skill.path ?? "",
enabled: skill.enabled ?? true,
})),
),
rateLimit: null,
serverDiagnostics: createServerDiagnostics(),
};
};
return {
settingsRef: {
settings,
vaultPath: overrides.vaultPath ?? "/vault",
vaultPath,
},
workspace: {
openThreadInNewView: overrides.openThreadInNewView ?? vi.fn(),
@ -1353,13 +1386,14 @@ function chatHost(overrides: ChatHostFixtureOverrides = {}): CodexChatHost {
overrides.setActiveThreads ??
((threads) => {
activeThreads = threads;
for (const listener of activeThreadListeners) listener(threads);
for (const listener of activeThreadResultListeners) listener(queryResult(threads));
}),
setAppServerMetadata:
overrides.setAppServerMetadata ??
((nextMetadata) => {
metadata = nextMetadata;
for (const listener of metadataListeners) listener(nextMetadata);
updateAppServerMetadata:
overrides.updateAppServerMetadata ??
((updater) => {
const nextMetadata = updater(metadata);
if (!nextMetadata) return null;
return applyMetadataToCache(nextMetadata);
}),
refreshActiveThreads:
overrides.refreshActiveThreads ??
@ -1369,39 +1403,66 @@ function chatHost(overrides: ChatHostFixtureOverrides = {}): CodexChatHost {
const listThreads = client["listThreads"] as (cwd: string, options: Record<string, unknown>) => Promise<{ data: ThreadRecord[] }>;
const response = await listThreads("/vault", { archived: false, cursor: null, limit: 100 });
activeThreads = response.data.map(threadFromRecord);
for (const listener of activeThreadListeners) listener(activeThreads);
for (const listener of activeThreadResultListeners) listener(queryResult(activeThreads));
return activeThreads;
}) as CodexChatHost["threadCatalog"]["refreshActiveThreads"]),
activeThreadsSnapshot: overrides.activeThreadsSnapshot ?? vi.fn(() => activeThreads),
appServerMetadataSnapshot: overrides.appServerMetadataSnapshot ?? vi.fn(() => metadata),
fetchAppServerMetadata:
overrides.fetchAppServerMetadata ??
vi.fn(async () => {
const nextMetadata = await loadAppServerMetadata();
return nextMetadata ? applyMetadataToCache(nextMetadata) : null;
}),
refreshAppServerMetadata:
overrides.refreshAppServerMetadata ??
vi.fn(async () => {
const nextMetadata = await loadAppServerMetadata();
return nextMetadata ? applyMetadataToCache(nextMetadata) : null;
}),
modelsSnapshot: overrides.modelsSnapshot ?? vi.fn(() => models),
fetchModels: overrides.fetchModels ?? vi.fn(async () => models ?? []),
refreshModels: overrides.refreshModels ?? vi.fn(async () => models ?? []),
observeActiveThreads: (listener, options = {}) => {
activeThreadListeners.add(listener);
if ((options.emitCurrent ?? true) && activeThreads) listener(activeThreads);
observeActiveThreadsResult: (listener, options = {}) => {
activeThreadResultListeners.add(listener);
if ((options.emitCurrent ?? true) && activeThreads) listener(queryResult(activeThreads));
return () => {
activeThreadListeners.delete(listener);
activeThreadResultListeners.delete(listener);
};
},
observeAppServerMetadata: (listener, options = {}) => {
metadataListeners.add(listener);
if ((options.emitCurrent ?? true) && metadata) listener(metadata);
observeAppServerMetadataResult: (listener, options = {}) => {
metadataResultListeners.add(listener);
if ((options.emitCurrent ?? true) && metadata) listener(queryResult(metadata));
return () => {
metadataListeners.delete(listener);
metadataResultListeners.delete(listener);
};
},
observeModels: (listener, options = {}) => {
modelListeners.add(listener);
if ((options.emitCurrent ?? true) && models) listener(models);
observeModelsResult: (listener, options = {}) => {
modelResultListeners.add(listener);
if ((options.emitCurrent ?? true) && models) listener(queryResult(models));
return () => {
modelListeners.delete(listener);
modelResultListeners.delete(listener);
};
},
},
};
}
function queryResult<T>(data: T | null): AppServerObservedQueryResult<T> {
return {
data,
error: null,
isFetching: false,
isLoading: false,
isPending: data === null,
isSuccess: data !== null,
isError: false,
isStale: false,
status: data === null ? "pending" : "success",
fetchStatus: "idle",
} as AppServerObservedQueryResult<T>;
}
async function chatView(
options: { activeLeafChangeListeners?: ((leaf: unknown) => void)[]; host?: CodexChatHost; requestSaveLayout?: () => void } = {},
) {

View file

@ -441,7 +441,7 @@ function threadsHost(overrides: Record<string, unknown> = {}) {
return response.data.map(threadFromRecord);
}),
activeThreadsSnapshot: vi.fn(() => null),
observeActiveThreads: vi.fn(() => () => undefined),
observeActiveThreadsResult: vi.fn(() => () => undefined),
...threadCatalogOverrides,
},
...hostOverrides,

View file

@ -651,16 +651,18 @@ function chatHostFixture(): CodexChatHost {
refreshFromOpenSurface: vi.fn(),
refreshThreadsViewLiveState: vi.fn(),
setActiveThreads: vi.fn(),
setAppServerMetadata: vi.fn(),
updateAppServerMetadata: vi.fn(() => null),
refreshActiveThreads: vi.fn(() => Promise.resolve([])),
activeThreadsSnapshot: vi.fn(() => null),
appServerMetadataSnapshot: vi.fn(() => null),
fetchAppServerMetadata: vi.fn(() => Promise.resolve(null)),
refreshAppServerMetadata: vi.fn(() => Promise.resolve(null)),
modelsSnapshot: vi.fn(() => null),
fetchModels: vi.fn(() => Promise.resolve([])),
refreshModels: vi.fn(() => Promise.resolve([])),
observeActiveThreads: vi.fn(() => () => undefined),
observeAppServerMetadata: vi.fn(() => () => undefined),
observeModels: vi.fn(() => () => undefined),
observeActiveThreadsResult: vi.fn(() => () => undefined),
observeAppServerMetadataResult: vi.fn(() => () => undefined),
observeModelsResult: vi.fn(() => () => undefined),
},
};
}

View file

@ -608,7 +608,7 @@ function newSettingsTab(
modelsSnapshot?: ModelMetadata[];
fetchModels?: () => Promise<readonly ModelMetadata[]>;
refreshModels?: () => Promise<readonly ModelMetadata[]>;
observeModels?: CodexPanelSettingTabHost["threadCatalog"]["observeModels"];
observeModels?: CodexPanelSettingTabHost["threadCatalog"]["observeModelsResult"];
notifyAppServerQueryContextChanged?: () => void;
refreshOpenViews?: () => void;
refreshFromOpenSurface?: () => void;
@ -630,7 +630,7 @@ function settingsTabHost(
modelsSnapshot?: ModelMetadata[];
fetchModels?: () => Promise<readonly ModelMetadata[]>;
refreshModels?: () => Promise<readonly ModelMetadata[]>;
observeModels?: CodexPanelSettingTabHost["threadCatalog"]["observeModels"];
observeModels?: CodexPanelSettingTabHost["threadCatalog"]["observeModelsResult"];
notifyAppServerQueryContextChanged?: () => void;
refreshOpenViews?: () => void;
refreshFromOpenSurface?: () => void;
@ -665,7 +665,7 @@ function settingsTabHost(
modelsSnapshot: vi.fn(() => options.modelsSnapshot ?? []),
fetchModels: options.fetchModels ?? vi.fn().mockResolvedValue(options.modelsSnapshot ?? []),
refreshModels: options.refreshModels ?? vi.fn().mockResolvedValue(options.modelsSnapshot ?? []),
observeModels: options.observeModels ?? vi.fn(() => () => undefined),
observeModelsResult: options.observeModels ?? vi.fn(() => () => undefined),
notifyAppServerQueryContextChanged: options.notifyAppServerQueryContextChanged ?? vi.fn(),
},
};

View file

@ -20,19 +20,19 @@ describe("SharedThreadCatalog", () => {
const { catalog } = catalogFixture();
const threads = [thread("thread")];
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.observeActiveThreadsResult(listener);
catalog.setActiveThreads(threads);
expect(catalog.activeThreadsSnapshot()).toEqual(threads);
expect(listener).toHaveBeenCalledWith(threads);
expect(listener).toHaveBeenCalledWith(expect.objectContaining({ data: threads }));
});
it("refreshes thread snapshots through the cache single-flight and notifies observers once", async () => {
const fetchThreads = vi.fn().mockResolvedValue([thread("thread")]);
const { catalog } = catalogFixture({ fetchThreads });
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.observeActiveThreadsResult(listener);
const first = catalog.refreshActiveThreads();
const second = catalog.refreshActiveThreads();
@ -41,7 +41,8 @@ describe("SharedThreadCatalog", () => {
await expect(second).resolves.toEqual([thread("thread")]);
expect(fetchThreads).toHaveBeenCalledOnce();
expect(catalog.activeThreadsSnapshot()).toEqual([thread("thread")]);
expect(listener).toHaveBeenCalledOnce();
expect(listener.mock.calls.filter(([result]) => result.data !== null)).toHaveLength(1);
expect(listener).toHaveBeenLastCalledWith(expect.objectContaining({ data: [thread("thread")] }));
});
it("does not notify stale thread observers after the app-server query context changes", async () => {
@ -59,11 +60,12 @@ describe("SharedThreadCatalog", () => {
});
let resolveThreads!: (threads: Thread[]) => void;
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.observeActiveThreadsResult(listener);
const fetch = catalog.refreshActiveThreads();
await flushMicrotasks();
context.codexPath = "codex-b";
listener.mockClear();
resolveThreads([thread("stale")]);
await expect(fetch).resolves.toEqual([thread("stale")]);
@ -80,17 +82,17 @@ describe("SharedThreadCatalog", () => {
context: () => context,
});
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.observeActiveThreadsResult(listener);
catalog.setActiveThreads([thread("a")]);
context.codexPath = "codex-b";
catalog.notifyAppServerQueryContextChanged();
catalog.setActiveThreads([thread("b")]);
expect(listener).toHaveBeenLastCalledWith([thread("b")]);
expect(listener).toHaveBeenLastCalledWith(expect.objectContaining({ data: [thread("b")] }));
context.codexPath = "codex-a";
catalog.notifyAppServerQueryContextChanged();
expect(listener).toHaveBeenLastCalledWith([thread("a")]);
expect(listener).toHaveBeenLastCalledWith(expect.objectContaining({ data: [thread("a")] }));
});
it("publishes metadata and model snapshots to cache and observers", () => {
@ -98,40 +100,42 @@ describe("SharedThreadCatalog", () => {
const metadata = serverMetadata({ availableModels: [model("gpt-test")] });
const metadataListener = vi.fn();
const modelListener = vi.fn();
catalog.observeAppServerMetadata(metadataListener);
catalog.observeModels(modelListener);
catalog.observeAppServerMetadataResult(metadataListener);
catalog.observeModelsResult(modelListener);
catalog.setAppServerMetadata(metadata);
catalog.updateAppServerMetadata(() => metadata);
expect(catalog.appServerMetadataSnapshot()).toEqual(metadata);
expect(catalog.modelsSnapshot()).toEqual(metadata.availableModels);
expect(metadataListener).toHaveBeenLastCalledWith(metadata);
expect(modelListener).toHaveBeenCalledWith(metadata.availableModels);
expect(metadataListener).toHaveBeenLastCalledWith(expect.objectContaining({ data: metadata }));
expect(modelListener).toHaveBeenCalledWith(expect.objectContaining({ data: metadata.availableModels }));
});
it("applies known rename mutations to cache and surfaces", () => {
const { catalog, surfaces } = catalogFixture();
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.observeActiveThreadsResult(listener);
catalog.setActiveThreads([thread("thread"), thread("other")]);
catalog.renameThreadInCatalog("thread", "Renamed");
expect(catalog.activeThreadsSnapshot()).toEqual([{ ...thread("thread"), name: "Renamed" }, thread("other")]);
expect(listener).toHaveBeenLastCalledWith([{ ...thread("thread"), name: "Renamed" }, thread("other")]);
expect(listener).toHaveBeenLastCalledWith(
expect.objectContaining({ data: [{ ...thread("thread"), name: "Renamed" }, thread("other")] }),
);
expect(surfaces.applyThreadRenamed).toHaveBeenCalledWith("thread", "Renamed");
});
it("applies known archive mutations to cache and surfaces", () => {
const { catalog, surfaces } = catalogFixture();
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.observeActiveThreadsResult(listener);
catalog.setActiveThreads([thread("thread"), thread("other")]);
catalog.archiveThreadInCatalog("thread", { closeOpenPanels: true });
expect(catalog.activeThreadsSnapshot()).toEqual([thread("other")]);
expect(listener).toHaveBeenLastCalledWith([thread("other")]);
expect(listener).toHaveBeenLastCalledWith(expect.objectContaining({ data: [thread("other")] }));
expect(surfaces.applyThreadArchived).toHaveBeenCalledWith("thread", { closeOpenPanels: true });
});
});