Adopt app-server query state

This commit is contained in:
murashit 2026-06-15 17:44:34 +09:00
parent bf4d3cc9e1
commit de913e9d9d
36 changed files with 926 additions and 850 deletions

11
package-lock.json generated
View file

@ -10,6 +10,7 @@
"license": "Apache-2.0",
"dependencies": {
"@preact/signals": "^2.9.1",
"@tanstack/query-core": "^5.101.0",
"@tanstack/virtual-core": "3.17.0",
"preact": "^10.29.2"
},
@ -2574,6 +2575,16 @@
"dev": true,
"license": "MIT"
},
"node_modules/@tanstack/query-core": {
"version": "5.101.0",
"resolved": "https://registry.npmjs.org/@tanstack/query-core/-/query-core-5.101.0.tgz",
"integrity": "sha512-cQetA74EB+seWySv1TTKr828TnP0u39m6LykwDXIo84SNortpDkp30TMEjkqtYCNP9c40uT/iwl6MLiufEt0Ow==",
"license": "MIT",
"funding": {
"type": "github",
"url": "https://github.com/sponsors/tannerlinsley"
}
},
"node_modules/@tanstack/virtual-core": {
"version": "3.17.0",
"resolved": "https://registry.npmjs.org/@tanstack/virtual-core/-/virtual-core-3.17.0.tgz",

View file

@ -62,6 +62,7 @@
},
"dependencies": {
"@preact/signals": "^2.9.1",
"@tanstack/query-core": "^5.101.0",
"@tanstack/virtual-core": "3.17.0",
"preact": "^10.29.2"
}

View file

@ -0,0 +1,171 @@
import { QueryClient, QueryObserver, type QueryObserverResult } from "@tanstack/query-core";
import type { ModelMetadata } from "../../domain/catalog/metadata";
import type { Thread } from "../../domain/threads/model";
import {
activeThreadsQueryKey,
appServerMetadataQueryKey,
appServerModelsQueryKey,
appServerQueriesFilter,
appServerQueryContextIsComplete,
cloneAppServerQueryContext,
type AppServerQueryContext,
} from "./keys";
import {
cloneModelMetadata,
cloneSharedServerMetadata,
cloneThreads,
mergeSharedServerMetadata,
type SharedServerMetadata,
} from "./snapshots";
export class AppServerQueryCache {
readonly client: QueryClient;
constructor(client: QueryClient = createAppServerQueryClient()) {
this.client = client;
}
clear(): void {
this.client.clear();
}
clearContext(context: AppServerQueryContext): void {
if (!appServerQueryContextIsComplete(context)) return;
this.client.removeQueries(appServerQueriesFilter(context));
}
activeThreadsSnapshot(context: AppServerQueryContext): readonly Thread[] | null {
if (!appServerQueryContextIsComplete(context)) return null;
const threads = this.client.getQueryData<readonly Thread[]>(activeThreadsQueryKey(context));
return threads ? cloneThreads(threads) : null;
}
observeActiveThreads(
context: AppServerQueryContext,
listener: (threads: readonly Thread[]) => void,
options: { emitCurrent?: boolean } = {},
): () => void {
return this.observeQuery(activeThreadsQueryKey(context), cloneThreads, listener, options);
}
async fetchActiveThreads(context: AppServerQueryContext, fetchThreads: () => Promise<readonly Thread[]>): Promise<readonly Thread[]> {
const refreshContext = cloneAppServerQueryContext(context);
if (!appServerQueryContextIsComplete(refreshContext)) {
return fetchThreads();
}
const key = activeThreadsQueryKey(refreshContext);
const threads = await this.client.fetchQuery({
queryKey: key,
queryFn: async () => {
const nextThreads = cloneThreads(await fetchThreads());
return nextThreads;
},
staleTime: 0,
});
return cloneThreads(threads);
}
setActiveThreads(context: AppServerQueryContext, threads: readonly Thread[]): void {
if (!appServerQueryContextIsComplete(context)) return;
this.client.setQueryData(activeThreadsQueryKey(context), cloneThreads(threads));
}
updateActiveThreads(
context: AppServerQueryContext,
updater: (threads: readonly Thread[] | null) => readonly Thread[] | null,
): readonly Thread[] | null {
const current = this.activeThreadsSnapshot(context);
const next = updater(current);
if (!next) return null;
this.setActiveThreads(context, next);
return cloneThreads(next);
}
appServerMetadataSnapshot(context: AppServerQueryContext): SharedServerMetadata | null {
if (!appServerQueryContextIsComplete(context)) return null;
const metadata = this.client.getQueryData<SharedServerMetadata>(appServerMetadataQueryKey(context));
return metadata ? cloneSharedServerMetadata(metadata) : null;
}
observeAppServerMetadata(
context: AppServerQueryContext,
listener: (metadata: SharedServerMetadata) => void,
options: { emitCurrent?: boolean } = {},
): () => void {
return this.observeQuery(appServerMetadataQueryKey(context), cloneSharedServerMetadata, listener, options);
}
setAppServerMetadata(context: AppServerQueryContext, metadata: SharedServerMetadata): SharedServerMetadata | null {
if (!appServerQueryContextIsComplete(context)) return null;
const previous = this.appServerMetadataSnapshot(context);
const next = mergeSharedServerMetadata(previous, metadata);
this.client.setQueryData(appServerMetadataQueryKey(context), cloneSharedServerMetadata(next));
if (metadata.serverDiagnostics.probes["model/list"].status === "ok") {
this.client.setQueryData(appServerModelsQueryKey(context), cloneModelMetadata(next.availableModels));
}
return cloneSharedServerMetadata(next);
}
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(
context: AppServerQueryContext,
listener: (models: readonly ModelMetadata[]) => void,
options: { emitCurrent?: boolean } = {},
): () => void {
return this.observeQuery(appServerModelsQueryKey(context), cloneModelMetadata, listener, options);
}
setModels(context: AppServerQueryContext, models: readonly ModelMetadata[]): readonly ModelMetadata[] | null {
if (!appServerQueryContextIsComplete(context)) return null;
const clonedModels = cloneModelMetadata(models);
this.client.setQueryData(appServerModelsQueryKey(context), clonedModels);
const metadata = this.appServerMetadataSnapshot(context);
if (metadata) {
this.client.setQueryData(appServerMetadataQueryKey(context), {
...metadata,
availableModels: cloneModelMetadata(clonedModels),
});
}
return cloneModelMetadata(clonedModels);
}
private observeQuery<T>(
queryKey: readonly unknown[],
clone: (value: T) => T,
listener: (value: T) => void,
options: { emitCurrent?: boolean },
): () => void {
const observer = new QueryObserver<T>(this.client, {
queryKey,
enabled: false,
});
const emit = (result: QueryObserverResult<T>): void => {
if (result.data !== undefined) listener(clone(result.data));
};
const unsubscribe = observer.subscribe(emit);
if (options.emitCurrent ?? true) emit(observer.getCurrentResult());
return unsubscribe;
}
}
function createAppServerQueryClient(): QueryClient {
return new QueryClient({
defaultOptions: {
queries: {
gcTime: Infinity,
retry: false,
refetchOnReconnect: false,
refetchOnWindowFocus: false,
},
mutations: {
retry: false,
},
},
});
}

View file

@ -0,0 +1,52 @@
import type { QueryKey } from "@tanstack/query-core";
export interface AppServerQueryContext {
codexPath: string;
vaultPath: string;
}
type AppServerQueryScope = readonly ["app-server", string, string];
export type AppServerActiveThreadsQueryKey = readonly [...AppServerQueryScope, "threads", "active"];
export type AppServerMetadataQueryKey = readonly [...AppServerQueryScope, "metadata"];
export type AppServerModelsQueryKey = readonly [...AppServerQueryScope, "models"];
export function appServerQueryContextIsComplete(context: AppServerQueryContext): boolean {
return nonEmptyString(context.codexPath) && nonEmptyString(context.vaultPath);
}
export function cloneAppServerQueryContext(context: AppServerQueryContext): AppServerQueryContext {
return { ...context };
}
export function appServerQueryContextMatches(left: AppServerQueryContext, right: AppServerQueryContext): boolean {
return (
appServerQueryContextIsComplete(left) &&
appServerQueryContextIsComplete(right) &&
left.codexPath === right.codexPath &&
left.vaultPath === right.vaultPath
);
}
function appServerQueryScope(context: AppServerQueryContext): AppServerQueryScope {
return ["app-server", context.codexPath, context.vaultPath];
}
export function activeThreadsQueryKey(context: AppServerQueryContext): AppServerActiveThreadsQueryKey {
return [...appServerQueryScope(context), "threads", "active"];
}
export function appServerMetadataQueryKey(context: AppServerQueryContext): AppServerMetadataQueryKey {
return [...appServerQueryScope(context), "metadata"];
}
export function appServerModelsQueryKey(context: AppServerQueryContext): AppServerModelsQueryKey {
return [...appServerQueryScope(context), "models"];
}
export function appServerQueriesFilter(context: AppServerQueryContext): { queryKey: QueryKey } {
return { queryKey: appServerQueryScope(context) };
}
function nonEmptyString(value: string): boolean {
return value.trim().length > 0;
}

View file

@ -0,0 +1,90 @@
import type { ModelMetadata, SkillMetadata } from "../../domain/catalog/metadata";
import { cloneRuntimeConfigSnapshot } from "../../domain/runtime/config";
import type { RateLimitSnapshot } from "../../domain/runtime/metrics";
import type { SharedServerMetadata } from "../../domain/server/metadata";
import type { Thread } from "../../domain/threads/model";
export type { SharedServerMetadata } from "../../domain/server/metadata";
export function cloneThreads(threads: readonly Thread[]): Thread[] {
return threads.map((thread) => ({ ...thread }));
}
export function cloneModelMetadata(models: readonly ModelMetadata[]): ModelMetadata[] {
return models.map((model) => ({
...model,
supportedReasoningEfforts: [...model.supportedReasoningEfforts],
inputModalities: [...model.inputModalities],
additionalSpeedTiers: [...model.additionalSpeedTiers],
serviceTiers: model.serviceTiers.map((tier) => ({ ...tier })),
}));
}
export function cloneSharedServerMetadata(metadata: SharedServerMetadata): SharedServerMetadata {
return {
...metadata,
runtimeConfig: metadata.runtimeConfig ? cloneRuntimeConfigSnapshot(metadata.runtimeConfig) : null,
rateLimit: metadata.rateLimit ? cloneRateLimitSnapshot(metadata.rateLimit) : null,
availableModels: cloneModelMetadata(metadata.availableModels),
availableSkills: cloneSkillMetadata(metadata.availableSkills),
serverDiagnostics: {
probes: { ...metadata.serverDiagnostics.probes },
mcpServers: metadata.serverDiagnostics.mcpServers.map((server) => ({ ...server })),
},
};
}
export function mergeSharedServerMetadata(previous: SharedServerMetadata | null, next: SharedServerMetadata): SharedServerMetadata {
const clonedNext = cloneSharedServerMetadata(next);
const clonedPrevious = previous ? cloneSharedServerMetadata(previous) : emptySharedServerMetadataResourceCache(clonedNext);
return {
...clonedNext,
availableModels: metadataResourceSucceeded(clonedNext, "model/list") ? clonedNext.availableModels : clonedPrevious.availableModels,
availableSkills: metadataResourceSucceeded(clonedNext, "skills/list") ? clonedNext.availableSkills : clonedPrevious.availableSkills,
rateLimit: metadataResourceSucceeded(clonedNext, "account/rateLimits/read") ? clonedNext.rateLimit : clonedPrevious.rateLimit,
serverDiagnostics: mergeServerDiagnostics(clonedPrevious, clonedNext),
};
}
function emptySharedServerMetadataResourceCache(metadata: SharedServerMetadata): SharedServerMetadata {
return {
...metadata,
availableModels: [],
availableSkills: [],
rateLimit: null,
serverDiagnostics: {
probes: { ...metadata.serverDiagnostics.probes },
mcpServers: [],
},
};
}
function metadataResourceSucceeded(
metadata: SharedServerMetadata,
method: keyof SharedServerMetadata["serverDiagnostics"]["probes"],
): boolean {
return metadata.serverDiagnostics.probes[method].status === "ok";
}
function mergeServerDiagnostics(previous: SharedServerMetadata, next: SharedServerMetadata): SharedServerMetadata["serverDiagnostics"] {
return {
probes: { ...next.serverDiagnostics.probes },
mcpServers:
next.serverDiagnostics.probes["mcpServerStatus/list"].status === "failed"
? previous.serverDiagnostics.mcpServers.map((server) => ({ ...server }))
: next.serverDiagnostics.mcpServers.map((server) => ({ ...server })),
};
}
function cloneRateLimitSnapshot(snapshot: RateLimitSnapshot): RateLimitSnapshot {
return {
...snapshot,
primary: snapshot.primary ? { ...snapshot.primary } : null,
secondary: snapshot.secondary ? { ...snapshot.secondary } : null,
individualLimit: snapshot.individualLimit ? { ...snapshot.individualLimit } : null,
};
}
function cloneSkillMetadata(skills: readonly SkillMetadata[]): SkillMetadata[] {
return skills.map((skill) => ({ ...skill }));
}

View file

@ -1,215 +0,0 @@
import type { SharedServerMetadata } from "../../domain/server/metadata";
import { cloneRuntimeConfigSnapshot } from "../../domain/runtime/config";
import type { RateLimitSnapshot } from "../../domain/runtime/metrics";
import type { Thread } from "../../domain/threads/model";
import type { ModelMetadata, SkillMetadata } from "../../domain/catalog/metadata";
export type { SharedServerMetadata } from "../../domain/server/metadata";
export interface SharedAppServerCacheContext {
codexPath: string;
vaultPath: string;
}
type SharedCache<T> = { kind: "unloaded" } | { kind: "loaded"; context: SharedAppServerCacheContext; data: T };
export interface SharedAppServerState {
threads: SharedCache<readonly Thread[]>;
appServerMetadata: SharedCache<SharedServerMetadata>;
availableModels: SharedCache<readonly ModelMetadata[]>;
}
export function createSharedAppServerState(): SharedAppServerState {
return {
threads: { kind: "unloaded" },
appServerMetadata: { kind: "unloaded" },
availableModels: { kind: "unloaded" },
};
}
export function applySharedThreadList(
state: SharedAppServerState,
context: SharedAppServerCacheContext,
threads: readonly Thread[],
): SharedAppServerState {
if (!sharedAppServerCacheContextIsComplete(context)) return state;
return {
...state,
threads: {
kind: "loaded",
context: cloneSharedAppServerCacheContext(context),
data: cloneThreads(threads),
},
};
}
export function applySharedServerMetadata(
state: SharedAppServerState,
context: SharedAppServerCacheContext,
metadata: SharedServerMetadata,
): SharedAppServerState {
if (!sharedAppServerCacheContextIsComplete(context)) return state;
const previous =
state.appServerMetadata.kind === "loaded" && sharedAppServerCacheContextMatches(state.appServerMetadata.context, context)
? state.appServerMetadata.data
: null;
const clonedMetadata = mergeSharedServerMetadata(previous, metadata);
return {
...state,
appServerMetadata: { kind: "loaded", context: cloneSharedAppServerCacheContext(context), data: clonedMetadata },
availableModels:
clonedMetadata.availableModels.length > 0
? {
kind: "loaded",
context: cloneSharedAppServerCacheContext(context),
data: cloneModelMetadata(clonedMetadata.availableModels),
}
: state.availableModels,
};
}
export function applySharedModels(
state: SharedAppServerState,
context: SharedAppServerCacheContext,
models: readonly ModelMetadata[],
): SharedAppServerState {
if (!sharedAppServerCacheContextIsComplete(context)) return state;
const clonedModels = cloneModelMetadata(models);
return {
...state,
appServerMetadata:
state.appServerMetadata.kind === "loaded" && sharedAppServerCacheContextMatches(state.appServerMetadata.context, context)
? {
kind: "loaded",
context: cloneSharedAppServerCacheContext(context),
data: { ...state.appServerMetadata.data, availableModels: cloneModelMetadata(clonedModels) },
}
: state.appServerMetadata,
availableModels: { kind: "loaded", context: cloneSharedAppServerCacheContext(context), data: clonedModels },
};
}
export function cachedSharedThreadList(state: SharedAppServerState, context: SharedAppServerCacheContext): readonly Thread[] | null {
if (!sharedAppServerCacheContextIsComplete(context)) return null;
if (state.threads.kind !== "loaded" || !sharedAppServerCacheContextMatches(state.threads.context, context)) return null;
return cloneThreads(state.threads.data);
}
export function cachedSharedServerMetadata(state: SharedAppServerState, context: SharedAppServerCacheContext): SharedServerMetadata | null {
if (!sharedAppServerCacheContextIsComplete(context)) return null;
if (state.appServerMetadata.kind === "loaded" && sharedAppServerCacheContextMatches(state.appServerMetadata.context, context)) {
return cloneSharedServerMetadata(state.appServerMetadata.data);
}
return null;
}
export function cachedSharedModels(state: SharedAppServerState, context: SharedAppServerCacheContext): readonly ModelMetadata[] | null {
if (!sharedAppServerCacheContextIsComplete(context)) return null;
return state.availableModels.kind === "loaded" && sharedAppServerCacheContextMatches(state.availableModels.context, context)
? cloneModelMetadata(state.availableModels.data)
: null;
}
function cloneSharedServerMetadata(metadata: SharedServerMetadata): SharedServerMetadata {
return {
...metadata,
runtimeConfig: metadata.runtimeConfig ? cloneRuntimeConfigSnapshot(metadata.runtimeConfig) : null,
rateLimit: metadata.rateLimit ? cloneRateLimitSnapshot(metadata.rateLimit) : null,
availableModels: cloneModelMetadata(metadata.availableModels),
availableSkills: cloneSkillMetadata(metadata.availableSkills),
serverDiagnostics: {
probes: { ...metadata.serverDiagnostics.probes },
mcpServers: metadata.serverDiagnostics.mcpServers.map((server) => ({ ...server })),
},
};
}
function mergeSharedServerMetadata(previous: SharedServerMetadata | null, next: SharedServerMetadata): SharedServerMetadata {
const clonedNext = cloneSharedServerMetadata(next);
const clonedPrevious = previous ? cloneSharedServerMetadata(previous) : emptySharedServerMetadataResourceCache(clonedNext);
return {
...clonedNext,
availableModels: metadataResourceSucceeded(clonedNext, "model/list") ? clonedNext.availableModels : clonedPrevious.availableModels,
availableSkills: metadataResourceSucceeded(clonedNext, "skills/list") ? clonedNext.availableSkills : clonedPrevious.availableSkills,
rateLimit: metadataResourceSucceeded(clonedNext, "account/rateLimits/read") ? clonedNext.rateLimit : clonedPrevious.rateLimit,
serverDiagnostics: mergeServerDiagnostics(clonedPrevious, clonedNext),
};
}
function emptySharedServerMetadataResourceCache(metadata: SharedServerMetadata): SharedServerMetadata {
return {
...metadata,
availableModels: [],
availableSkills: [],
rateLimit: null,
serverDiagnostics: {
probes: { ...metadata.serverDiagnostics.probes },
mcpServers: [],
},
};
}
function metadataResourceSucceeded(
metadata: SharedServerMetadata,
method: keyof SharedServerMetadata["serverDiagnostics"]["probes"],
): boolean {
if (method === "model/list") return metadata.availableModels.length > 0 && metadata.serverDiagnostics.probes[method].status === "ok";
return metadata.serverDiagnostics.probes[method].status === "ok";
}
function mergeServerDiagnostics(previous: SharedServerMetadata, next: SharedServerMetadata): SharedServerMetadata["serverDiagnostics"] {
return {
probes: { ...next.serverDiagnostics.probes },
mcpServers:
next.serverDiagnostics.probes["mcpServerStatus/list"].status === "failed"
? previous.serverDiagnostics.mcpServers.map((server) => ({ ...server }))
: next.serverDiagnostics.mcpServers.map((server) => ({ ...server })),
};
}
function cloneRateLimitSnapshot(snapshot: RateLimitSnapshot): RateLimitSnapshot {
return {
...snapshot,
primary: snapshot.primary ? { ...snapshot.primary } : null,
secondary: snapshot.secondary ? { ...snapshot.secondary } : null,
individualLimit: snapshot.individualLimit ? { ...snapshot.individualLimit } : null,
};
}
function cloneThreads(threads: readonly Thread[]): Thread[] {
return threads.map((thread) => ({ ...thread }));
}
function cloneModelMetadata(models: readonly ModelMetadata[]): ModelMetadata[] {
return models.map((model) => ({
...model,
supportedReasoningEfforts: [...model.supportedReasoningEfforts],
inputModalities: [...model.inputModalities],
additionalSpeedTiers: [...model.additionalSpeedTiers],
serviceTiers: model.serviceTiers.map((tier) => ({ ...tier })),
}));
}
function cloneSkillMetadata(skills: readonly SkillMetadata[]): SkillMetadata[] {
return skills.map((skill) => ({ ...skill }));
}
function cloneSharedAppServerCacheContext(context: SharedAppServerCacheContext): SharedAppServerCacheContext {
return { ...context };
}
export function sharedAppServerCacheContextMatches(left: SharedAppServerCacheContext, right: SharedAppServerCacheContext): boolean {
return (
sharedAppServerCacheContextIsComplete(left) &&
sharedAppServerCacheContextIsComplete(right) &&
left.codexPath === right.codexPath &&
left.vaultPath === right.vaultPath
);
}
export function sharedAppServerCacheContextIsComplete(context: SharedAppServerCacheContext): boolean {
return nonEmptyString(context.codexPath) && nonEmptyString(context.vaultPath);
}
function nonEmptyString(value: string): boolean {
return value.trim().length > 0;
}

View file

@ -1,86 +0,0 @@
import type { Thread } from "../../domain/threads/model";
import type { ModelMetadata } from "../../domain/catalog/metadata";
import {
applySharedServerMetadata,
applySharedModels,
applySharedThreadList,
cachedSharedServerMetadata,
cachedSharedModels,
cachedSharedThreadList,
createSharedAppServerState,
sharedAppServerCacheContextIsComplete,
sharedAppServerCacheContextMatches,
type SharedAppServerCacheContext,
type SharedServerMetadata,
type SharedAppServerState,
} from "./shared-cache-state";
type ThreadListRefreshLifecycleState =
| { kind: "idle" }
| { kind: "refreshing"; context: SharedAppServerCacheContext; promise: Promise<readonly Thread[]> };
export class SharedAppServerCache {
private state: SharedAppServerState = createSharedAppServerState();
private threadListRefreshLifecycle: ThreadListRefreshLifecycleState = { kind: "idle" };
refreshThreadList(
context: SharedAppServerCacheContext,
fetchThreads: () => Promise<readonly Thread[]>,
onSnapshot?: (threads: readonly Thread[]) => void,
): Promise<readonly Thread[]> {
const refreshContext = { ...context };
if (!sharedAppServerCacheContextIsComplete(refreshContext)) {
return fetchThreads();
}
if (
this.threadListRefreshLifecycle.kind === "refreshing" &&
sharedAppServerCacheContextMatches(this.threadListRefreshLifecycle.context, refreshContext)
) {
return this.threadListRefreshLifecycle.promise;
}
const promise = fetchThreads()
.then((threads) => {
if (
this.threadListRefreshLifecycle.kind === "refreshing" &&
this.threadListRefreshLifecycle.promise === promise &&
sharedAppServerCacheContextMatches(this.threadListRefreshLifecycle.context, refreshContext)
) {
this.applyThreadListSnapshot(refreshContext, threads);
onSnapshot?.(threads);
}
return threads;
})
.finally(() => {
if (this.threadListRefreshLifecycle.kind === "refreshing" && this.threadListRefreshLifecycle.promise === promise) {
this.threadListRefreshLifecycle = { kind: "idle" };
}
});
this.threadListRefreshLifecycle = { kind: "refreshing", context: refreshContext, promise };
return promise;
}
applyThreadListSnapshot(context: SharedAppServerCacheContext, threads: readonly Thread[]): void {
this.state = applySharedThreadList(this.state, context, threads);
}
cachedThreadList(context: SharedAppServerCacheContext): readonly Thread[] | null {
return cachedSharedThreadList(this.state, context);
}
applyAppServerMetadataSnapshot(context: SharedAppServerCacheContext, metadata: SharedServerMetadata): void {
this.state = applySharedServerMetadata(this.state, context, metadata);
}
cachedAppServerMetadata(context: SharedAppServerCacheContext): SharedServerMetadata | null {
return cachedSharedServerMetadata(this.state, context);
}
applyModelsSnapshot(context: SharedAppServerCacheContext, models: readonly ModelMetadata[]): void {
if (models.length === 0) return;
this.state = applySharedModels(this.state, context, models);
}
cachedModels(context: SharedAppServerCacheContext): readonly ModelMetadata[] | null {
return cachedSharedModels(this.state, context);
}
}

View file

@ -16,7 +16,7 @@ import { mcpStatusLines as buildMcpStatusLines } from "../../application/connect
import { cloneServerDiagnostics, type ChatServerActionHost } from "./host";
interface RefreshDiagnosticProbesOptions {
cachedAppServerMetadata?: boolean;
appServerMetadataSnapshot?: boolean;
}
interface DiagnosticProbeSnapshot {
@ -26,7 +26,7 @@ interface DiagnosticProbeSnapshot {
}
export interface ChatServerDiagnosticsActionsHost extends ChatServerActionHost {
publishAppServerMetadata: (metadata: SharedServerMetadata) => void;
setAppServerMetadata: (metadata: SharedServerMetadata) => void;
serverMetadataSnapshot: () => SharedServerMetadata;
}
@ -58,7 +58,7 @@ async function refreshDiagnosticProbes(
if (!client) return false;
const probes: Promise<DiagnosticProbeSnapshot>[] = [];
if (!options.cachedAppServerMetadata) {
if (!options.appServerMetadataSnapshot) {
probes.push(
probeDiagnostic(
"model/list",
@ -139,7 +139,7 @@ async function refreshPublishedDiagnosticProbes(
options: RefreshDiagnosticProbesOptions = {},
): Promise<void> {
if (!(await refreshDiagnosticProbes(host, options))) return;
host.publishAppServerMetadata(host.serverMetadataSnapshot());
host.setAppServerMetadata(host.serverMetadataSnapshot());
}
async function mcpStatusLines(host: ChatServerDiagnosticsActionsHost): Promise<string[]> {

View file

@ -11,7 +11,7 @@ import type { SharedServerMetadata } from "../../../../domain/server/metadata";
import { cloneServerDiagnostics, type ChatServerActionHost } from "./host";
export interface ChatServerMetadataActionsHost extends ChatServerActionHost {
publishAppServerMetadata: (metadata: SharedServerMetadata) => void;
setAppServerMetadata: (metadata: SharedServerMetadata) => void;
}
export interface ChatServerMetadataActions {
@ -20,7 +20,7 @@ export interface ChatServerMetadataActions {
loadAppServerMetadata: () => Promise<SharedServerMetadata | null>;
refreshAppServerMetadata: () => Promise<SharedServerMetadata | null>;
refreshPublishedAppServerMetadata: () => Promise<SharedServerMetadata | null>;
publishAppServerMetadataSnapshot: () => void;
setAppServerMetadataSnapshot: () => void;
refreshModels: () => Promise<void>;
loadModels: () => Promise<ModelMetadataProbeResult>;
refreshSkills: (forceReload?: boolean) => Promise<void>;
@ -40,8 +40,8 @@ export function createChatServerMetadataActions(host: ChatServerMetadataActionsH
loadAppServerMetadata: () => loadAppServerMetadata(host),
refreshAppServerMetadata: () => refreshAppServerMetadata(host),
refreshPublishedAppServerMetadata: () => refreshPublishedAppServerMetadata(host),
publishAppServerMetadataSnapshot: () => {
publishAppServerMetadataSnapshot(host);
setAppServerMetadataSnapshot: () => {
setAppServerMetadataSnapshot(host);
},
refreshModels: async () => {
await refreshModels(host);
@ -104,12 +104,12 @@ async function refreshAppServerMetadata(host: ChatServerMetadataActionsHost): Pr
async function refreshPublishedAppServerMetadata(host: ChatServerMetadataActionsHost): Promise<SharedServerMetadata | null> {
const metadata = await refreshAppServerMetadata(host);
if (metadata) host.publishAppServerMetadata(metadata);
if (metadata) host.setAppServerMetadata(metadata);
return metadata;
}
function publishAppServerMetadataSnapshot(host: ChatServerMetadataActionsHost): void {
host.publishAppServerMetadata(serverMetadataSnapshot(host));
function setAppServerMetadataSnapshot(host: ChatServerMetadataActionsHost): void {
host.setAppServerMetadata(serverMetadataSnapshot(host));
}
async function refreshModels(host: ChatServerMetadataActionsHost): Promise<boolean> {
@ -146,7 +146,7 @@ async function refreshSkills(host: ChatServerMetadataActionsHost, forceReload =
async function refreshPublishedSkills(host: ChatServerMetadataActionsHost, forceReload = false): Promise<void> {
if (!(await refreshSkills(host, forceReload))) return;
publishAppServerMetadataSnapshot(host);
setAppServerMetadataSnapshot(host);
}
async function loadSkills(host: ChatServerMetadataActionsHost, forceReload = false): Promise<SkillMetadataProbeResult> {
@ -178,7 +178,7 @@ async function refreshPublishedRateLimits(host: ChatServerMetadataActionsHost):
rateLimit: rateLimit.data,
serverDiagnostics: diagnostics,
});
publishAppServerMetadataSnapshot(host);
setAppServerMetadataSnapshot(host);
return;
}
host.stateStore.dispatch({ type: "connection/metadata-applied", serverDiagnostics: diagnostics });

View file

@ -35,10 +35,10 @@ function cannotRejectServerRequestMessage(): string {
}
export interface ChatInboundControllerActions {
refreshThreads: () => void;
fetchActiveThreads: () => void;
refreshRateLimits: () => void;
refreshSkills: (forceReload?: boolean) => void;
publishAppServerMetadata: () => void;
setAppServerMetadata: () => void;
maybeNameThread: (threadId: string, turnId: string, completedSummary: ThreadConversationSummary | null) => void;
applyThreadArchived: (threadId: string) => void;
applyThreadRenamed: (threadId: string, name: string | null) => void;
@ -185,7 +185,7 @@ export class ChatInboundController {
private runNotificationEffect(effect: ChatNotificationEffect): void {
switch (effect.type) {
case "refresh-threads":
this.actions.refreshThreads();
this.actions.fetchActiveThreads();
return;
case "refresh-rate-limits":
this.actions.refreshRateLimits();
@ -194,7 +194,7 @@ export class ChatInboundController {
this.actions.refreshSkills(effect.forceReload);
return;
case "publish-app-server-metadata":
this.actions.publishAppServerMetadata();
this.actions.setAppServerMetadata();
return;
case "maybe-name-thread":
this.actions.maybeNameThread(effect.threadId, effect.turnId, effect.completedSummary);

View file

@ -89,7 +89,7 @@ export class ChatConnectionController {
handleChatConnectionExit(this.host);
}
async refreshThreads(): Promise<void> {
async fetchActiveThreads(): Promise<void> {
if (!this.host.connection.currentClient()) return;
try {
await this.host.loadSharedThreadList();
@ -114,7 +114,7 @@ export class ChatConnectionController {
} catch (error) {
this.host.addSystemMessage(error instanceof Error ? error.message : String(error));
}
await this.refreshThreads();
await this.fetchActiveThreads();
}
async refreshSkills(forceReload = false): Promise<void> {

View file

@ -25,9 +25,13 @@ export type ThreadCatalogFacade = Pick<
| "renameThreadInCatalog"
| "refreshThreadsViewLiveState"
| "refreshFromOpenSurface"
| "applyThreads"
| "publishAppServerMetadata"
| "refreshThreads"
| "cachedThreads"
| "cachedAppServerMetadata"
| "setActiveThreads"
| "setAppServerMetadata"
| "fetchActiveThreads"
| "activeThreadsSnapshot"
| "appServerMetadataSnapshot"
| "modelsSnapshot"
| "observeActiveThreads"
| "observeAppServerMetadata"
| "observeModels"
>;

View file

@ -29,7 +29,7 @@ interface ThreadPartsContext {
getClosing: () => boolean;
};
thread: {
refreshThreads: () => Promise<void>;
fetchActiveThreads: () => Promise<void>;
notifyIdentityChanged: () => void;
refreshTabHeader: () => void;
};
@ -114,7 +114,7 @@ export function createThreadParts(context: ThreadPartsContext) {
thread.notifyIdentityChanged();
},
refreshAfterThreadMutation: async () => {
await thread.refreshThreads();
await thread.fetchActiveThreads();
threadCatalog.refreshFromOpenSurface();
},
};

View file

@ -44,7 +44,7 @@ export interface ChatConnectionBundleContext {
deferredTasks: ChatViewDeferredTasks;
threadCatalog: Pick<
ThreadCatalogFacade,
"applyThreads" | "publishAppServerMetadata" | "refreshThreads" | "archiveThreadInCatalog" | "renameThreadInCatalog"
"setActiveThreads" | "setAppServerMetadata" | "fetchActiveThreads" | "archiveThreadInCatalog" | "renameThreadInCatalog"
>;
goalSync: ThreadGoalSyncActions;
autoTitle: AutoTitleController;
@ -63,16 +63,16 @@ export function createChatConnectionBundle(context: ChatConnectionBundleContext)
stateStore,
vaultPath,
currentClient,
publishAppServerMetadata: (metadata) => {
threadCatalog.publishAppServerMetadata(metadata);
setAppServerMetadata: (metadata) => {
threadCatalog.setAppServerMetadata(metadata);
},
});
const serverDiagnostics = createChatServerDiagnosticsActions({
stateStore,
vaultPath,
currentClient,
publishAppServerMetadata: (metadata) => {
threadCatalog.publishAppServerMetadata(metadata);
setAppServerMetadata: (metadata) => {
threadCatalog.setAppServerMetadata(metadata);
},
serverMetadataSnapshot: () => serverMetadata.serverMetadataSnapshot(),
});
@ -82,29 +82,29 @@ export function createChatConnectionBundle(context: ChatConnectionBundleContext)
currentClient,
runtimeSnapshotForState: runtimeSnapshotForChatState,
publishThreadList: (threads) => {
threadCatalog.applyThreads(threads);
threadCatalog.setActiveThreads(threads);
},
syncThreadGoal: (threadId) => {
void goalSync.syncThreadGoal(threadId);
},
});
const loadSharedThreadList = async (): Promise<void> => {
const threads = await threadCatalog.refreshThreads(() => serverThreads.loadThreadList());
const threads = await threadCatalog.fetchActiveThreads(() => serverThreads.loadThreadList());
serverThreads.applyThreadList(threads);
};
const serverRequestHost = {
currentClient,
};
const inboundController = new ChatInboundController(stateStore, {
refreshThreads: () => {
fetchActiveThreads: () => {
void loadSharedThreadList();
},
refreshRateLimits: () => {
void serverMetadata.refreshPublishedRateLimits();
},
refreshSkills: (forceReload) => void serverMetadata.refreshPublishedSkills(forceReload),
publishAppServerMetadata: () => {
serverMetadata.publishAppServerMetadataSnapshot();
setAppServerMetadata: () => {
serverMetadata.setAppServerMetadataSnapshot();
},
maybeNameThread: (threadId, turnId, completedSummary) => {
autoTitle.maybeAutoTitleThread(threadId, turnId, completedSummary);
@ -165,7 +165,7 @@ export function createChatConnectionBundle(context: ChatConnectionBundleContext)
scheduleDeferredDiagnostics: () => {
deferredTasks.scheduleDiagnostics(() => {
if (connection.isConnected()) {
void serverDiagnostics.refreshPublishedDiagnosticProbes({ cachedAppServerMetadata: true });
void serverDiagnostics.refreshPublishedDiagnosticProbes({ appServerMetadataSnapshot: true });
}
});
},

View file

@ -207,7 +207,7 @@ export function createChatPanelRuntime(context: ChatPanelRuntimeContext): ChatPa
connectionController = controller;
const { threads: serverThreads, diagnostics: serverDiagnostics } = serverParts.serverActions;
const ensureConnected = () => connectionController.ensureConnected();
const refreshThreads = () => connectionController.refreshThreads();
const fetchActiveThreads = () => connectionController.fetchActiveThreads();
const runtimeSettings = createChatRuntimeSettingsActions({
stateStore,
@ -251,7 +251,7 @@ export function createChatPanelRuntime(context: ChatPanelRuntimeContext): ChatPa
},
status,
thread: {
refreshThreads,
fetchActiveThreads,
notifyIdentityChanged: () => {
context.notifyActiveThreadIdentityChanged();
},

View file

@ -63,6 +63,7 @@ export class ChatPanelSession implements ChatSurfaceHandle {
private readonly resumeWork = new ChatResumeWorkTracker();
private readonly messageScrollIntent: ChatMessageScrollIntentState = createChatMessageScrollIntentState();
private readonly localItemIds: LocalChatItemIdFactory = createLocalChatItemIdFactory();
private readonly appServerStateUnsubscribers: (() => void)[] = [];
private opened = false;
private closing = false;
@ -153,16 +154,16 @@ export class ChatPanelSession implements ChatSurfaceHandle {
return this.loadSharedThreadList();
}
applyThreadListSnapshot(threads: readonly Thread[]): void {
private receiveObservedThreads(threads: readonly Thread[]): void {
this.parts.serverActions.threads.applyThreadList(threads);
this.refreshTabHeader();
}
applyAppServerMetadataSnapshot(metadata: SharedServerMetadata): void {
private receiveObservedAppServerMetadata(metadata: SharedServerMetadata): void {
this.parts.serverActions.metadata.applyAppServerMetadata(metadata);
}
applyAvailableModelsSnapshot(models: readonly ModelMetadata[]): void {
private receiveObservedModels(models: readonly ModelMetadata[]): void {
this.dispatch({ type: "connection/metadata-applied", availableModels: models });
}
@ -212,7 +213,7 @@ export class ChatPanelSession implements ChatSurfaceHandle {
this.environment.obsidian.registerPointerDown((event) => {
this.closeToolbarPanelOnOutsidePointer(event);
});
this.applyCachedAppServerState();
this.subscribeAppServerState();
this.mountOrRepairShell();
this.scheduleWarmup();
this.scheduleRestoredThreadHydration();
@ -224,6 +225,7 @@ export class ChatPanelSession implements ChatSurfaceHandle {
this.connectionWork.invalidate();
this.invalidateResumeWork();
this.deferredTasks.clearAll();
this.unsubscribeAppServerState();
const panelRoot = this.environment.view.panelRoot();
this.parts.render.messageStreamPresenter.dispose();
this.parts.composer.controller.dispose();
@ -251,10 +253,43 @@ export class ChatPanelSession implements ChatSurfaceHandle {
}
private applyCachedAppServerState(): void {
const threads = this.environment.plugin.threadCatalog.cachedThreads();
const threads = this.environment.plugin.threadCatalog.activeThreadsSnapshot();
if (threads) this.parts.serverActions.threads.applyThreadList(threads);
const metadata = this.environment.plugin.threadCatalog.cachedAppServerMetadata();
const metadata = this.environment.plugin.threadCatalog.appServerMetadataSnapshot();
if (metadata) this.parts.serverActions.metadata.applyAppServerMetadata(metadata);
const models = this.environment.plugin.threadCatalog.modelsSnapshot();
if (models) this.receiveObservedModels(models);
}
private subscribeAppServerState(): void {
this.unsubscribeAppServerState();
this.applyCachedAppServerState();
this.appServerStateUnsubscribers.push(
this.environment.plugin.threadCatalog.observeActiveThreads(
(threads) => {
this.receiveObservedThreads(threads);
},
{ emitCurrent: false },
),
this.environment.plugin.threadCatalog.observeAppServerMetadata(
(metadata) => {
this.receiveObservedAppServerMetadata(metadata);
},
{ emitCurrent: false },
),
this.environment.plugin.threadCatalog.observeModels(
(models) => {
this.receiveObservedModels(models);
},
{ emitCurrent: false },
),
);
}
private unsubscribeAppServerState(): void {
while (this.appServerStateUnsubscribers.length > 0) {
this.appServerStateUnsubscribers.pop()?.();
}
}
private mountOrRepairShell(): void {
@ -296,7 +331,7 @@ export class ChatPanelSession implements ChatSurfaceHandle {
}
private async loadSharedThreadList(): Promise<void> {
const threads = await this.environment.plugin.threadCatalog.refreshThreads(() => this.parts.serverActions.threads.loadThreadList());
const threads = await this.environment.plugin.threadCatalog.fetchActiveThreads(() => this.parts.serverActions.threads.loadThreadList());
this.parts.serverActions.threads.applyThreadList(threads);
}

View file

@ -1,6 +1,3 @@
import type { ModelMetadata } from "../../../domain/catalog/metadata";
import type { SharedServerMetadata } from "../../../domain/server/metadata";
import type { Thread } from "../../../domain/threads/model";
import type { OpenCodexPanelSnapshot } from "../../../workspace/open-panel-snapshot";
export interface ChatSurfaceHandle {
@ -11,9 +8,6 @@ export interface ChatSurfaceHandle {
close(): void;
refreshSettings(): void;
refreshSharedThreadList(): Promise<void>;
applyThreadListSnapshot(threads: readonly Thread[]): void;
applyAppServerMetadataSnapshot(metadata: SharedServerMetadata): void;
applyAvailableModelsSnapshot(models: readonly ModelMetadata[]): void;
openPanelSnapshot(): OpenCodexPanelSnapshot;
openThread(threadId: string): Promise<void>;
focusThread(threadId?: string | null): Promise<void>;

View file

@ -17,7 +17,7 @@ export interface ThreadPickerHost {
openThreadInAvailableView(threadId: string): Promise<void>;
}
type ThreadPickerCatalog = Pick<SharedThreadCatalog, "cachedThreads" | "refreshThreads">;
type ThreadPickerCatalog = Pick<SharedThreadCatalog, "activeThreadsSnapshot" | "fetchActiveThreads">;
interface ThreadSuggestion {
thread: Thread;
@ -78,14 +78,14 @@ function threadOpenModeFromEvent(evt: MouseEvent | KeyboardEvent): ThreadOpenMod
}
async function loadThreadPickerThreads(host: ThreadPickerHost): Promise<readonly Thread[]> {
const cached = host.threadCatalog.cachedThreads();
const cached = host.threadCatalog.activeThreadsSnapshot();
if (cached) return cached;
return withShortLivedAppServerClient(
host.settings.codexPath,
host.vaultPath,
async (client) =>
host.threadCatalog.refreshThreads(async () => {
host.threadCatalog.fetchActiveThreads(async () => {
return listThreads(client, host.vaultPath);
}),
{

View file

@ -41,7 +41,12 @@ export interface CodexThreadsHost {
type ThreadsThreadCatalog = Pick<
SharedThreadCatalog,
"archiveThreadInCatalog" | "renameThreadInCatalog" | "refreshFromOpenSurface" | "refreshThreads" | "cachedThreads"
| "archiveThreadInCatalog"
| "renameThreadInCatalog"
| "refreshFromOpenSurface"
| "fetchActiveThreads"
| "activeThreadsSnapshot"
| "observeActiveThreads"
>;
export interface CodexThreadsSessionEnvironment {
@ -70,6 +75,7 @@ export class CodexThreadsSession {
private status: ThreadsViewStatus = { kind: "idle" };
private threads: readonly Thread[] = [];
private readonly renameStates = new Map<string, ThreadsRenameState>();
private unsubscribeThreads: (() => void) | null = null;
private archiveConfirmThreadId: string | null = null;
constructor(private readonly environment: CodexThreadsSessionEnvironment) {
@ -125,10 +131,13 @@ export class CodexThreadsSession {
this.environment.registerPointerDown((event) => {
this.cancelArchiveConfirmOnOutsidePointer(event);
});
const cachedThreads = this.host.threadCatalog.cachedThreads();
if (cachedThreads) {
this.threads = cachedThreads;
const activeThreadsSnapshot = this.host.threadCatalog.activeThreadsSnapshot();
if (activeThreadsSnapshot) {
this.threads = activeThreadsSnapshot;
}
this.unsubscribeThreads = this.host.threadCatalog.observeActiveThreads((threads) => {
this.receiveObservedThreads(threads);
});
this.render();
void this.refresh();
}
@ -137,6 +146,8 @@ export class CodexThreadsSession {
this.connectionWork.invalidate();
this.refreshLifecycle = transitionThreadsViewRefreshLifecycle(this.refreshLifecycle, { type: "invalidated" });
this.deferredTasks.clearAll();
this.unsubscribeThreads?.();
this.unsubscribeThreads = null;
this.connection.disconnect();
this.client = null;
unmountThreadsView(this.environment.root);
@ -149,7 +160,7 @@ export class CodexThreadsSession {
try {
await this.ensureConnected();
if (this.isStaleRefresh(refresh) || !this.client) return;
const threads = await this.host.threadCatalog.refreshThreads(async () => {
const threads = await this.host.threadCatalog.fetchActiveThreads(async () => {
if (!this.client) return [];
return listThreads(this.client, this.host.vaultPath);
});
@ -168,7 +179,7 @@ export class CodexThreadsSession {
this.scheduleRender();
}
applyThreadListSnapshot(threads: readonly Thread[]): void {
private receiveObservedThreads(threads: readonly Thread[]): void {
this.threads = threads;
this.status = threads.length === 0 ? { kind: "empty", message: "No threads" } : { kind: "idle" };
this.render();

View file

@ -1,7 +1,6 @@
import { ItemView, type WorkspaceLeaf } from "obsidian";
import { VIEW_TYPE_CODEX_THREADS } from "../../constants";
import type { Thread } from "../../domain/threads/model";
import { CodexThreadsSession, type CodexThreadsHost } from "./session";
export type { CodexThreadsHost } from "./session";
@ -49,8 +48,4 @@ export class CodexThreadsView extends ItemView {
refreshLiveState(): void {
this.session.refreshLiveState();
}
applyThreadListSnapshot(threads: readonly Thread[]): void {
this.session.applyThreadListSnapshot(threads);
}
}

View file

@ -1,8 +1,8 @@
import type { App } from "obsidian";
import { VIEW_TYPE_CODEX_THREADS, VIEW_TYPE_CODEX_TURN_DIFF } from "./constants";
import { SharedAppServerCache } from "./app-server/services/shared-cache";
import type { SharedAppServerCacheContext } from "./app-server/services/shared-cache-state";
import { AppServerQueryCache } from "./app-server/query/cache";
import type { AppServerQueryContext } from "./app-server/query/keys";
import type { CodexChatHost, PluginSettingsRef } from "./features/chat/application/ports/chat-host";
import type { ChatTurnDiffViewState } from "./features/chat/domain/turn-diff";
import { persistedChatTurnDiffViewState } from "./features/chat/domain/turn-diff";
@ -21,7 +21,7 @@ export interface CodexPanelRuntimeOptions {
}
export class CodexPanelRuntime {
private readonly sharedAppServerCache = new SharedAppServerCache();
private readonly appServerQueries = new AppServerQueryCache();
readonly panels: WorkspacePanelCoordinator;
private readonly threadSurfaces: ThreadSurfaceActions;
readonly threadCatalog: SharedThreadCatalog;
@ -38,14 +38,15 @@ export class CodexPanelRuntime {
panels: this.panels,
});
this.threadCatalog = new SharedThreadCatalog({
cache: this.sharedAppServerCache,
cache: this.appServerQueries,
surfaces: this.threadSurfaces,
context: () => this.sharedAppServerCacheContext(),
context: () => this.appServerQueryContext(),
});
}
reset(): void {
this.panels.reset();
this.appServerQueries.clear();
}
recordLastFocusedPanel(leaf: Parameters<WorkspacePanelCoordinator["recordLastFocusedPanel"]>[0]): void {
@ -144,7 +145,7 @@ export class CodexPanelRuntime {
await this.options.app.workspace.revealLeaf(leaf);
}
private sharedAppServerCacheContext(): SharedAppServerCacheContext {
private appServerQueryContext(): AppServerQueryContext {
return {
codexPath: this.options.settingsRef.settings.codexPath,
vaultPath: this.options.settingsRef.vaultPath,

View file

@ -24,7 +24,10 @@ export interface SettingsDynamicDataHost {
threadCatalog: SettingsThreadCatalog;
}
type SettingsThreadCatalog = Pick<SharedThreadCatalog, "refreshFromOpenSurface" | "cachedModels" | "publishModels">;
type SettingsThreadCatalog = Pick<
SharedThreadCatalog,
"refreshFromOpenSurface" | "modelsSnapshot" | "setModels" | "observeModels" | "notifyAppServerQueryContextChanged"
>;
interface SettingsDynamicDataControllerCallbacks {
display(): void;
@ -55,12 +58,25 @@ export class SettingsDynamicDataController {
private hooksLifecycle: SettingsDynamicSectionLifecycleState = createSettingsDynamicSectionLifecycle();
private models: ModelMetadata[] = [];
private modelsLifecycle: SettingsDynamicSectionLifecycleState = createSettingsDynamicSectionLifecycle();
private unsubscribeModels: (() => void) | null = null;
constructor(
private readonly host: SettingsDynamicDataHost,
private readonly callbacks: SettingsDynamicDataControllerCallbacks,
) {
this.models = [...(host.threadCatalog.cachedModels() ?? [])];
this.activate();
}
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();
},
{ emitCurrent: false },
);
}
maybeAutoLoadSettingsData(): void {
@ -73,7 +89,7 @@ export class SettingsDynamicDataController {
this.settingsDataAutoLoadStarted = false;
this.settingsDynamicOperationId += 1;
this.settingsDataRefreshLifecycle = { kind: "idle" };
this.models = [...(this.host.threadCatalog.cachedModels() ?? [])];
this.models = [...(this.host.threadCatalog.modelsSnapshot() ?? [])];
this.modelsLifecycle = createSettingsDynamicSectionLifecycle();
this.hooks = [];
this.hookWarnings = [];
@ -83,6 +99,11 @@ export class SettingsDynamicDataController {
this.archivedThreadsLifecycle = createSettingsDynamicSectionLifecycle();
}
dispose(): void {
this.unsubscribeModels?.();
this.unsubscribeModels = null;
}
async refreshSettingsData(): Promise<void> {
this.settingsDataAutoLoadStarted = true;
const operationId = this.nextSettingsDynamicOperationId();
@ -114,7 +135,7 @@ export class SettingsDynamicDataController {
if (result.models.ok) {
this.models = result.models.data;
this.host.threadCatalog.publishModels(result.models.data);
this.host.threadCatalog.setModels(result.models.data);
this.modelsLifecycle = transitionSettingsDynamicSectionLifecycle(this.modelsLifecycle, {
type: "loaded",
status: result.models.status,

View file

@ -36,9 +36,15 @@ export class CodexPanelSettingTab extends PluginSettingTab {
}
display(): void {
this.dynamicData.activate();
this.renderSettingsTab({ autoLoadCodexData: true });
}
override hide(): void {
this.dynamicData.dispose();
super.hide();
}
private renderSettingsTab(options: { autoLoadCodexData: boolean }): void {
const { containerEl } = this;
containerEl.empty();
@ -68,6 +74,8 @@ export class CodexPanelSettingTab extends PluginSettingTab {
await this.plugin.saveSettings();
if (codexPathChanged) {
this.dynamicData.resetSettingsDataContext();
this.plugin.threadCatalog.notifyAppServerQueryContextChanged();
this.plugin.refreshOpenViews();
this.renderSettingsTab({ autoLoadCodexData: false });
}
});

View file

@ -1,50 +1,77 @@
import type { ModelMetadata } from "../domain/catalog/metadata";
import type { SharedServerMetadata } from "../domain/server/metadata";
import type { Thread } from "../domain/threads/model";
import type { SharedAppServerCache } from "../app-server/services/shared-cache";
import type { SharedAppServerCacheContext } from "../app-server/services/shared-cache-state";
import type { AppServerQueryCache } from "../app-server/query/cache";
import { appServerQueryContextMatches, cloneAppServerQueryContext, type AppServerQueryContext } from "../app-server/query/keys";
import type { ThreadSurfaceActions } from "./thread-surface-actions";
export interface SharedThreadCatalogOptions {
cache: SharedAppServerCache;
cache: AppServerQueryCache;
surfaces: ThreadSurfaceActions;
context: () => SharedAppServerCacheContext;
context: () => AppServerQueryContext;
}
export class SharedThreadCatalog {
private readonly contextChangeListeners = new Set<() => void>();
constructor(private readonly options: SharedThreadCatalogOptions) {}
cachedThreads(): readonly Thread[] | null {
return this.options.cache.cachedThreadList(this.context());
activeThreadsSnapshot(): readonly Thread[] | null {
return this.options.cache.activeThreadsSnapshot(this.context());
}
async refreshThreads(fetchThreads: () => Promise<readonly Thread[]>): Promise<readonly Thread[]> {
return this.options.cache.refreshThreadList(this.context(), fetchThreads, (threads) => {
this.options.surfaces.applyThreadListSnapshot(threads);
});
async fetchActiveThreads(fetchThreads: () => Promise<readonly Thread[]>): Promise<readonly Thread[]> {
return this.options.cache.fetchActiveThreads(this.context(), fetchThreads);
}
applyThreads(threads: readonly Thread[]): void {
this.options.cache.applyThreadListSnapshot(this.context(), threads);
this.options.surfaces.applyThreadListSnapshot(threads);
setActiveThreads(threads: readonly Thread[]): void {
this.options.cache.setActiveThreads(this.context(), threads);
}
cachedAppServerMetadata(): SharedServerMetadata | null {
return this.options.cache.cachedAppServerMetadata(this.context());
observeActiveThreads(listener: (threads: readonly Thread[]) => void, options?: { emitCurrent?: boolean }): () => void {
return this.observeCurrentContext(
(context, contextListener, observeOptions) => this.options.cache.observeActiveThreads(context, contextListener, observeOptions),
listener,
options,
);
}
publishAppServerMetadata(metadata: SharedServerMetadata): void {
this.options.cache.applyAppServerMetadataSnapshot(this.context(), metadata);
this.options.surfaces.publishAppServerMetadata(metadata);
appServerMetadataSnapshot(): SharedServerMetadata | null {
return this.options.cache.appServerMetadataSnapshot(this.context());
}
cachedModels(): readonly ModelMetadata[] | null {
return this.options.cache.cachedModels(this.context());
setAppServerMetadata(metadata: SharedServerMetadata): void {
this.options.cache.setAppServerMetadata(this.context(), metadata);
}
publishModels(models: readonly ModelMetadata[]): void {
this.options.cache.applyModelsSnapshot(this.context(), models);
this.options.surfaces.publishModels(models);
observeAppServerMetadata(listener: (metadata: SharedServerMetadata) => void, options?: { emitCurrent?: boolean }): () => void {
return this.observeCurrentContext(
(context, contextListener, observeOptions) => this.options.cache.observeAppServerMetadata(context, contextListener, observeOptions),
listener,
options,
);
}
modelsSnapshot(): readonly ModelMetadata[] | null {
return this.options.cache.modelsSnapshot(this.context());
}
setModels(models: readonly ModelMetadata[]): void {
this.options.cache.setModels(this.context(), models);
}
observeModels(listener: (models: readonly ModelMetadata[]) => void, options?: { emitCurrent?: boolean }): () => void {
return this.observeCurrentContext(
(context, contextListener, observeOptions) => this.options.cache.observeModels(context, contextListener, observeOptions),
listener,
options,
);
}
notifyAppServerQueryContextChanged(): void {
for (const listener of [...this.contextChangeListeners]) {
listener();
}
}
refreshFromOpenSurface(): void {
@ -56,18 +83,16 @@ export class SharedThreadCatalog {
}
renameThreadInCatalog(threadId: string, name: string | null): void {
const threads = this.cachedThreads();
if (threads) {
this.applyThreads(threads.map((thread) => (thread.id === threadId ? { ...thread, name } : thread)));
}
this.options.cache.updateActiveThreads(this.context(), (current) => {
return current ? current.map((thread) => (thread.id === threadId ? { ...thread, name } : thread)) : null;
});
this.options.surfaces.applyThreadRenamed(threadId, name);
}
archiveThreadInCatalog(threadId: string, options?: { closeOpenPanels?: boolean }): void {
const threads = this.cachedThreads();
if (threads) {
this.applyThreads(threads.filter((thread) => thread.id !== threadId));
}
this.options.cache.updateActiveThreads(this.context(), (current) => {
return current ? current.filter((thread) => thread.id !== threadId) : null;
});
this.options.surfaces.applyThreadArchived(threadId, options);
}
@ -75,7 +100,45 @@ export class SharedThreadCatalog {
this.options.surfaces.refreshThreadsViewLiveState();
}
private context(): SharedAppServerCacheContext {
private context(): AppServerQueryContext {
return this.options.context();
}
private observeCurrentContext<T>(
observe: (context: AppServerQueryContext, listener: (value: T) => void, options: { emitCurrent?: boolean }) => () => void,
listener: (value: T) => void,
options: { emitCurrent?: boolean } = {},
): () => void {
let observedContext: AppServerQueryContext | null = null;
let unsubscribeQuery: (() => void) | null = null;
let firstSubscribe = true;
const subscribe = (): void => {
const context = cloneAppServerQueryContext(this.context());
if (observedContext && appServerQueryContextMatches(observedContext, context)) return;
unsubscribeQuery?.();
observedContext = context;
const observeOptions: { emitCurrent?: boolean } = {};
if (firstSubscribe) {
if (options.emitCurrent !== undefined) observeOptions.emitCurrent = options.emitCurrent;
} else {
observeOptions.emitCurrent = true;
}
unsubscribeQuery = observe(
context,
(value) => {
if (appServerQueryContextMatches(this.context(), context)) listener(value);
},
observeOptions,
);
firstSubscribe = false;
};
subscribe();
this.contextChangeListeners.add(subscribe);
return () => {
this.contextChangeListeners.delete(subscribe);
unsubscribeQuery?.();
};
}
}

View file

@ -2,9 +2,6 @@ import type { App } from "obsidian";
import { VIEW_TYPE_CODEX_THREADS } from "../constants";
import { CodexThreadsView } from "../features/threads-view/view";
import type { ModelMetadata } from "../domain/catalog/metadata";
import type { Thread } from "../domain/threads/model";
import type { SharedServerMetadata } from "../domain/server/metadata";
import type { WorkspacePanelCoordinator } from "./panel-coordinator";
export interface ThreadSurfaceActionsOptions {
@ -15,11 +12,8 @@ export interface ThreadSurfaceActionsOptions {
export interface ThreadSurfaceActions {
refreshOpenViews(): void;
invalidateThreadsFromOpenSurface(): void;
applyThreadListSnapshot(threads: readonly Thread[]): void;
applyThreadArchived(threadId: string, options?: { closeOpenPanels?: boolean }): void;
applyThreadRenamed(threadId: string, name: string | null): void;
publishAppServerMetadata(metadata: SharedServerMetadata): void;
publishModels(models: readonly ModelMetadata[]): void;
refreshThreadsViewLiveState(): void;
}
@ -49,15 +43,6 @@ export function createThreadSurfaceActions(options: ThreadSurfaceActionsOptions)
invalidateThreadsFromOpenSurface,
applyThreadListSnapshot(threads: readonly Thread[]): void {
for (const view of options.panels.panelViews()) {
view.surface.applyThreadListSnapshot(threads);
}
for (const view of threadsViews()) {
view.applyThreadListSnapshot(threads);
}
},
applyThreadArchived(threadId: string, archiveOptions: { closeOpenPanels?: boolean } = {}): void {
const leavesToClose = archiveOptions.closeOpenPanels ? options.panels.panelLeavesForThread(threadId) : [];
for (const view of options.panels.panelViews()) {
@ -74,18 +59,6 @@ export function createThreadSurfaceActions(options: ThreadSurfaceActionsOptions)
}
},
publishAppServerMetadata(metadata: SharedServerMetadata): void {
for (const view of options.panels.panelViews()) {
view.surface.applyAppServerMetadataSnapshot(metadata);
}
},
publishModels(models: readonly ModelMetadata[]): void {
for (const view of options.panels.panelViews()) {
view.surface.applyAvailableModelsSnapshot(models);
}
},
refreshThreadsViewLiveState(): void {
for (const view of threadsViews()) {
view.refreshLiveState();

View file

@ -1,25 +1,26 @@
import { describe, expect, it, vi } from "vitest";
import { createServerDiagnostics, diagnosticProbeError, diagnosticProbeOk } from "../../src/domain/server/diagnostics";
import { SharedAppServerCache } from "../../src/app-server/services/shared-cache";
import { AppServerQueryCache } from "../../src/app-server/query/cache";
import type { AppServerQueryContext } from "../../src/app-server/query/keys";
import { emptyRuntimeConfigSnapshot, type RuntimeConfigSnapshot } from "../../src/app-server/protocol/runtime-config";
import type { RateLimitSnapshot } from "../../src/app-server/protocol/runtime-metrics";
import type { SharedAppServerCacheContext, SharedServerMetadata } from "../../src/app-server/services/shared-cache-state";
import type { SharedServerMetadata } from "../../src/app-server/query/snapshots";
import type { ModelMetadata, SkillMetadata } from "../../src/domain/catalog/metadata";
describe("SharedAppServerCache", () => {
describe("AppServerQueryCache", () => {
it("updates successful metadata resources while preserving failed resource cache values", () => {
const cache = new SharedAppServerCache();
const cache = new AppServerQueryCache();
const context = cacheContext();
const goodMetadata = metadata({
availableSkills: [skillMetadata("writer")],
rateLimit: rateLimit(42),
});
cache.applyAppServerMetadataSnapshot(context, goodMetadata);
expect(cache.cachedAppServerMetadata(context)?.availableModels.map((model) => model.model)).toEqual(["gpt-5.5"]);
cache.setAppServerMetadata(context, goodMetadata);
expect(cache.appServerMetadataSnapshot(context)?.availableModels.map((model) => model.model)).toEqual(["gpt-5.5"]);
cache.applyAppServerMetadataSnapshot(
cache.setAppServerMetadata(
context,
metadata({
availableModels: [modelMetadata("gpt-5.6")],
@ -30,22 +31,22 @@ describe("SharedAppServerCache", () => {
}),
);
expect(cache.cachedAppServerMetadata(context)?.availableModels.map((model) => model.model)).toEqual(["gpt-5.6"]);
expect(cache.cachedAppServerMetadata(context)?.availableSkills.map((skill) => skill.name)).toEqual(["writer"]);
expect(cache.cachedAppServerMetadata(context)?.rateLimit?.primary?.usedPercent).toBe(42);
expect(cache.cachedAppServerMetadata(context)?.serverDiagnostics.probes["skills/list"].status).toBe("failed");
expect(cache.cachedAppServerMetadata(context)?.serverDiagnostics.probes["account/rateLimits/read"].status).toBe("failed");
expect(cache.appServerMetadataSnapshot(context)?.availableModels.map((model) => model.model)).toEqual(["gpt-5.6"]);
expect(cache.appServerMetadataSnapshot(context)?.availableSkills.map((skill) => skill.name)).toEqual(["writer"]);
expect(cache.appServerMetadataSnapshot(context)?.rateLimit?.primary?.usedPercent).toBe(42);
expect(cache.appServerMetadataSnapshot(context)?.serverDiagnostics.probes["skills/list"].status).toBe("failed");
expect(cache.appServerMetadataSnapshot(context)?.serverDiagnostics.probes["account/rateLimits/read"].status).toBe("failed");
cache.applyAppServerMetadataSnapshot(context, metadata({ availableModels: [], modelProbeStatus: "failed" }));
expect(cache.cachedAppServerMetadata(context)?.availableModels.map((model) => model.model)).toEqual(["gpt-5.6"]);
expect(cache.cachedAppServerMetadata(context)?.serverDiagnostics.probes["model/list"].status).toBe("failed");
cache.setAppServerMetadata(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");
});
it("loads initial metadata snapshots without caching failed resource values", () => {
const cache = new SharedAppServerCache();
const cache = new AppServerQueryCache();
const context = cacheContext();
cache.applyAppServerMetadataSnapshot(
cache.setAppServerMetadata(
context,
metadata({
availableModels: [modelMetadata("failed-model")],
@ -56,121 +57,117 @@ describe("SharedAppServerCache", () => {
}),
);
const cached = cache.cachedAppServerMetadata(context);
const cached = cache.appServerMetadataSnapshot(context);
expect(cached?.runtimeConfig).not.toBeNull();
expect(cached?.serverDiagnostics.probes["model/list"].status).toBe("failed");
expect(cached?.availableModels).toEqual([]);
expect(cached?.availableSkills.map((skill) => skill.name)).toEqual(["writer"]);
expect(cached?.rateLimit).toBeNull();
expect(cache.cachedModels(context)).toBeNull();
expect(cache.modelsSnapshot(context)).toBeNull();
});
it("does not share or store snapshots before the cache context is complete", async () => {
const cache = new SharedAppServerCache();
const cache = new AppServerQueryCache();
const context = cacheContext({ codexPath: "" });
const firstRefresh = deferred<readonly ReturnType<typeof thread>[]>();
const firstFetch = vi.fn(() => firstRefresh.promise);
const secondFetch = vi.fn().mockResolvedValue([thread("second")]);
const onSnapshot = vi.fn();
const firstPromise = cache.refreshThreadList(context, firstFetch, onSnapshot);
await expect(cache.refreshThreadList(context, secondFetch, onSnapshot)).resolves.toEqual([thread("second")]);
const firstPromise = cache.fetchActiveThreads(context, firstFetch);
await expect(cache.fetchActiveThreads(context, secondFetch)).resolves.toEqual([thread("second")]);
firstRefresh.resolve([thread("first")]);
await expect(firstPromise).resolves.toEqual([thread("first")]);
cache.applyThreadListSnapshot(context, [thread("applied")]);
cache.applyAppServerMetadataSnapshot(context, metadata());
cache.applyModelsSnapshot(context, [modelMetadata("gpt-5.6")]);
cache.setActiveThreads(context, [thread("applied")]);
cache.setAppServerMetadata(context, metadata());
cache.setModels(context, [modelMetadata("gpt-5.6")]);
expect(firstFetch).toHaveBeenCalledOnce();
expect(secondFetch).toHaveBeenCalledOnce();
expect(onSnapshot).not.toHaveBeenCalled();
expect(cache.cachedThreadList(context)).toBeNull();
expect(cache.cachedAppServerMetadata(context)).toBeNull();
expect(cache.cachedModels(context)).toBeNull();
expect(cache.activeThreadsSnapshot(context)).toBeNull();
expect(cache.appServerMetadataSnapshot(context)).toBeNull();
expect(cache.modelsSnapshot(context)).toBeNull();
});
it("does not replace shared models with an empty model snapshot", () => {
const cache = new SharedAppServerCache();
it("stores successful empty model snapshots as shared cache truth", () => {
const cache = new AppServerQueryCache();
const context = cacheContext();
cache.applyModelsSnapshot(context, [modelMetadata("gpt-5.5")]);
cache.applyModelsSnapshot(context, []);
cache.setModels(context, [modelMetadata("gpt-5.5")]);
cache.setModels(context, []);
expect(cache.cachedModels(context)?.map((model) => model.model)).toEqual(["gpt-5.5"]);
expect(cache.modelsSnapshot(context)).toEqual([]);
cache.setAppServerMetadata(context, metadata({ availableModels: [modelMetadata("gpt-5.6")] }));
cache.setAppServerMetadata(context, metadata({ availableModels: [] }));
expect(cache.appServerMetadataSnapshot(context)?.availableModels).toEqual([]);
expect(cache.modelsSnapshot(context)).toEqual([]);
});
it("does not reuse metadata or model snapshots across app-server cache contexts", () => {
const cache = new SharedAppServerCache();
const cache = new AppServerQueryCache();
const context = cacheContext();
cache.applyAppServerMetadataSnapshot(context, metadata());
cache.applyModelsSnapshot(context, [modelMetadata("gpt-5.6")]);
cache.setAppServerMetadata(context, metadata());
cache.setModels(context, [modelMetadata("gpt-5.6")]);
expect(cache.cachedAppServerMetadata(cacheContext({ vaultPath: "/other-vault" }))).toBeNull();
expect(cache.cachedAppServerMetadata(cacheContext({ codexPath: "/opt/codex" }))).toBeNull();
expect(cache.cachedModels(cacheContext({ vaultPath: "/other-vault" }))).toBeNull();
expect(cache.cachedModels(cacheContext({ codexPath: "/opt/codex" }))).toBeNull();
expect(cache.appServerMetadataSnapshot(cacheContext({ vaultPath: "/other-vault" }))).toBeNull();
expect(cache.appServerMetadataSnapshot(cacheContext({ codexPath: "/opt/codex" }))).toBeNull();
expect(cache.modelsSnapshot(cacheContext({ vaultPath: "/other-vault" }))).toBeNull();
expect(cache.modelsSnapshot(cacheContext({ codexPath: "/opt/codex" }))).toBeNull();
});
it("stores successful empty thread list snapshots as shared cache truth", async () => {
const cache = new SharedAppServerCache();
const cache = new AppServerQueryCache();
const context = cacheContext();
const onSnapshot = vi.fn();
cache.applyThreadListSnapshot(context, [thread("cached")]);
cache.applyThreadListSnapshot(context, []);
cache.setActiveThreads(context, [thread("cached")]);
cache.setActiveThreads(context, []);
expect(cache.cachedThreadList(context)).toEqual([]);
expect(cache.activeThreadsSnapshot(context)).toEqual([]);
await expect(cache.refreshThreadList(context, () => Promise.resolve([]), onSnapshot)).resolves.toEqual([]);
expect(onSnapshot).toHaveBeenCalledWith([]);
expect(cache.cachedThreadList(context)).toEqual([]);
await expect(cache.fetchActiveThreads(context, () => Promise.resolve([]))).resolves.toEqual([]);
expect(cache.activeThreadsSnapshot(context)).toEqual([]);
});
it("ignores stale thread list refresh snapshots after the app-server cache context changes", async () => {
const cache = new SharedAppServerCache();
it("keys thread list refresh snapshots by app-server query context", async () => {
const cache = new AppServerQueryCache();
const oldContext = cacheContext({ codexPath: "codex-old" });
const newContext = cacheContext({ codexPath: "codex-new" });
const oldRefresh = deferred<readonly ReturnType<typeof thread>[]>();
const newRefresh = deferred<readonly ReturnType<typeof thread>[]>();
const oldSnapshot = vi.fn();
const newSnapshot = vi.fn();
const oldPromise = cache.refreshThreadList(oldContext, () => oldRefresh.promise, oldSnapshot);
const newPromise = cache.refreshThreadList(newContext, () => newRefresh.promise, newSnapshot);
const oldPromise = cache.fetchActiveThreads(oldContext, () => oldRefresh.promise);
const newPromise = cache.fetchActiveThreads(newContext, () => newRefresh.promise);
oldRefresh.resolve([thread("old-thread")]);
await expect(oldPromise).resolves.toEqual([thread("old-thread")]);
expect(oldSnapshot).not.toHaveBeenCalled();
expect(cache.cachedThreadList(oldContext)).toBeNull();
expect(cache.activeThreadsSnapshot(oldContext)?.map((item) => item.id)).toEqual(["old-thread"]);
newRefresh.resolve([thread("new-thread")]);
await expect(newPromise).resolves.toEqual([thread("new-thread")]);
expect(newSnapshot).toHaveBeenCalledWith([thread("new-thread")]);
expect(cache.cachedThreadList(newContext)?.map((item) => item.id)).toEqual(["new-thread"]);
expect(cache.cachedThreadList(oldContext)).toBeNull();
expect(cache.activeThreadsSnapshot(newContext)?.map((item) => item.id)).toEqual(["new-thread"]);
expect(cache.activeThreadsSnapshot(oldContext)?.map((item) => item.id)).toEqual(["old-thread"]);
});
it("keys in-flight thread list refreshes by the captured app-server cache context", async () => {
const cache = new SharedAppServerCache();
it("stores in-flight thread list refreshes under the captured app-server cache context", async () => {
const cache = new AppServerQueryCache();
const context = cacheContext({ codexPath: "codex-captured" });
const capturedContext = { ...context };
const refresh = deferred<readonly ReturnType<typeof thread>[]>();
const onSnapshot = vi.fn();
const promise = cache.refreshThreadList(context, () => refresh.promise, onSnapshot);
const promise = cache.fetchActiveThreads(context, () => refresh.promise);
context.codexPath = "codex-mutated";
refresh.resolve([thread("captured")]);
await expect(promise).resolves.toEqual([thread("captured")]);
expect(onSnapshot).toHaveBeenCalledWith([thread("captured")]);
expect(cache.cachedThreadList(capturedContext)?.map((item) => item.id)).toEqual(["captured"]);
expect(cache.cachedThreadList(context)).toBeNull();
expect(cache.activeThreadsSnapshot(capturedContext)?.map((item) => item.id)).toEqual(["captured"]);
expect(cache.activeThreadsSnapshot(context)).toBeNull();
});
});
function cacheContext(overrides: Partial<SharedAppServerCacheContext> = {}): SharedAppServerCacheContext {
function cacheContext(overrides: Partial<AppServerQueryContext> = {}): AppServerQueryContext {
return {
codexPath: "codex",
vaultPath: "/vault",

View file

@ -1,211 +0,0 @@
import { describe, expect, it } from "vitest";
import { createServerDiagnostics, diagnosticProbeError, diagnosticProbeOk } from "../../src/domain/server/diagnostics";
import type { RateLimitSnapshot } from "../../src/app-server/protocol/runtime-metrics";
import { emptyRuntimeConfigSnapshot } from "../../src/app-server/protocol/runtime-config";
import {
applySharedServerMetadata,
applySharedModels,
applySharedThreadList,
cachedSharedServerMetadata,
cachedSharedModels,
cachedSharedThreadList,
createSharedAppServerState,
sharedAppServerCacheContextIsComplete,
type SharedAppServerCacheContext,
} from "../../src/app-server/services/shared-cache-state";
import type { ModelMetadata } from "../../src/domain/catalog/metadata";
import type { Thread } from "../../src/domain/threads/model";
describe("shared app-server cache state", () => {
it("keeps snapshots detached from caller-owned arrays", () => {
const sourceThreads = [threadFixture("thread-1")];
const threadState = applySharedThreadList(createSharedAppServerState(), cacheContext(), sourceThreads);
sourceThreads.push(threadFixture("thread-2"));
const cachedThreads = cachedSharedThreadList(threadState, cacheContext());
expect(cachedThreads?.map((thread) => thread.id)).toEqual(["thread-1"]);
const mutableCachedThreads = cachedThreads as Thread[];
mutableCachedThreads.push(threadFixture("thread-3"));
expect(cachedSharedThreadList(threadState, cacheContext())?.map((thread) => thread.id)).toEqual(["thread-1"]);
const sourceModels = [modelFixture("gpt-5.5")];
const modelState = applySharedModels(createSharedAppServerState(), cacheContext(), sourceModels);
sourceModels.push(modelFixture("gpt-5.6"));
expect(expectPresent(cachedSharedModels(modelState, cacheContext())).map((model) => model.model)).toEqual(["gpt-5.5"]);
const cachedModels = expectPresent(cachedSharedModels(modelState, cacheContext()));
(cachedModels[0]?.supportedReasoningEfforts as string[] | undefined)?.push("high");
expect(expectPresent(cachedSharedModels(modelState, cacheContext()))[0]?.supportedReasoningEfforts).toEqual([]);
const sourceRateLimit = rateLimitFixture();
const metadataState = applySharedServerMetadata(createSharedAppServerState(), cacheContext(), {
runtimeConfig: emptyRuntimeConfigSnapshot(),
availableModels: sourceModels,
availableSkills: [{ name: "skill", description: "", path: "/tmp/skill", enabled: true }],
rateLimit: sourceRateLimit,
serverDiagnostics: {
...diagnostics({}),
mcpServers: [{ name: "server", startupStatus: "ready", authStatus: null, toolCount: 1, message: null }],
},
});
sourceModels.push(modelFixture("gpt-5.7"));
if (sourceRateLimit.primary) sourceRateLimit.primary.usedPercent = 99;
const cachedMetadata = cachedSharedServerMetadata(metadataState, cacheContext());
expect(cachedMetadata?.availableModels.map((model) => model.model)).toEqual(["gpt-5.5", "gpt-5.6"]);
expect(cachedMetadata?.rateLimit?.primary?.usedPercent).toBe(42);
if (cachedMetadata?.rateLimit?.primary) cachedMetadata.rateLimit.primary.usedPercent = 100;
expect(cachedSharedServerMetadata(metadataState, cacheContext())?.rateLimit?.primary?.usedPercent).toBe(42);
});
it("does not return snapshots for a different app-server cache context", () => {
const state = applySharedModels(
applySharedThreadList(createSharedAppServerState(), cacheContext(), [threadFixture("thread-1")]),
cacheContext(),
[modelFixture("gpt-5.5")],
);
expect(cachedSharedThreadList(state, cacheContext({ vaultPath: "/other-vault" }))).toBeNull();
expect(cachedSharedModels(state, cacheContext({ codexPath: "/opt/codex" }))).toBeNull();
});
it("does not load or return snapshots before the cache context is complete", () => {
const incompleteContexts = [cacheContext({ codexPath: "" }), cacheContext({ vaultPath: " " })];
expect(sharedAppServerCacheContextIsComplete(cacheContext())).toBe(true);
for (const incompleteContext of incompleteContexts) {
expect(sharedAppServerCacheContextIsComplete(incompleteContext)).toBe(false);
const state = applySharedServerMetadata(
applySharedModels(
applySharedThreadList(createSharedAppServerState(), incompleteContext, [threadFixture("thread-1")]),
incompleteContext,
[modelFixture("gpt-5.5")],
),
incompleteContext,
{
runtimeConfig: emptyRuntimeConfigSnapshot(),
availableModels: [modelFixture("gpt-5.6")],
availableSkills: [],
rateLimit: null,
serverDiagnostics: createServerDiagnostics(),
},
);
expect(cachedSharedThreadList(state, incompleteContext)).toBeNull();
expect(cachedSharedModels(state, incompleteContext)).toBeNull();
expect(cachedSharedServerMetadata(state, incompleteContext)).toBeNull();
}
});
it("merges app-server metadata by successful resource", () => {
const context = cacheContext();
const first = applySharedServerMetadata(createSharedAppServerState(), context, {
runtimeConfig: emptyRuntimeConfigSnapshot(),
availableModels: [modelFixture("gpt-5.5")],
availableSkills: [{ name: "writer", description: "", path: "/tmp/writer", enabled: true }],
rateLimit: rateLimitFixture(),
serverDiagnostics: diagnostics({}),
});
const second = applySharedServerMetadata(first, context, {
runtimeConfig: emptyRuntimeConfigSnapshot(),
availableModels: [modelFixture("gpt-5.6")],
availableSkills: [{ name: "stale", description: "", path: "/tmp/stale", enabled: true }],
rateLimit: null,
serverDiagnostics: diagnostics({ skills: "failed", rateLimit: "failed" }),
});
const cached = expectPresent(cachedSharedServerMetadata(second, context));
expect(cached.availableModels.map((model) => model.model)).toEqual(["gpt-5.6"]);
expect(cached.availableSkills.map((skill) => skill.name)).toEqual(["writer"]);
expect(cached.rateLimit?.primary?.usedPercent).toBe(42);
expect(cached.serverDiagnostics.probes["skills/list"].status).toBe("failed");
expect(cached.serverDiagnostics.probes["account/rateLimits/read"].status).toBe("failed");
expect(expectPresent(cachedSharedModels(second, context)).map((model) => model.model)).toEqual(["gpt-5.6"]);
});
it("loads initial app-server metadata while caching only successful resources", () => {
const context = cacheContext();
const state = applySharedServerMetadata(createSharedAppServerState(), context, {
runtimeConfig: emptyRuntimeConfigSnapshot(),
availableModels: [modelFixture("failed-model")],
availableSkills: [{ name: "writer", description: "", path: "/tmp/writer", enabled: true }],
rateLimit: rateLimitFixture(),
serverDiagnostics: diagnostics({ models: "failed", rateLimit: "failed" }),
});
const cached = expectPresent(cachedSharedServerMetadata(state, context));
expect(cached.runtimeConfig).not.toBeNull();
expect(cached.serverDiagnostics.probes["model/list"].status).toBe("failed");
expect(cached.availableModels).toEqual([]);
expect(cached.availableSkills.map((skill) => skill.name)).toEqual(["writer"]);
expect(cached.rateLimit).toBeNull();
expect(cachedSharedModels(state, context)).toBeNull();
});
});
function cacheContext(overrides: Partial<SharedAppServerCacheContext> = {}): SharedAppServerCacheContext {
return {
codexPath: "codex",
vaultPath: "/vault",
...overrides,
};
}
function expectPresent<T>(value: T | null): T {
if (value === null) throw new Error("Expected value to be present");
return value;
}
function threadFixture(id: string): Thread {
return {
id,
preview: "",
name: null,
archived: false,
createdAt: 1,
updatedAt: 1,
};
}
function modelFixture(model: string): ModelMetadata {
return {
id: model,
model,
displayName: model,
description: "",
hidden: false,
supportedReasoningEfforts: [],
defaultReasoningEffort: "medium",
inputModalities: [],
additionalSpeedTiers: [],
serviceTiers: [],
defaultServiceTier: null,
isDefault: false,
};
}
function rateLimitFixture(): RateLimitSnapshot {
return {
limitId: "codex",
limitName: "Codex",
primary: { usedPercent: 42, windowDurationMins: 300, resetsAt: 1 },
secondary: { usedPercent: 10, windowDurationMins: 10_080, resetsAt: 2 },
individualLimit: { limit: "100", used: "42", remainingPercent: 58, resetsAt: 3 },
rateLimitReachedType: null,
};
}
function diagnostics(overrides: { models?: "ok" | "failed"; skills?: "ok" | "failed"; rateLimit?: "ok" | "failed" }) {
const next = createServerDiagnostics();
next.probes["model/list"] =
overrides.models === "failed" ? diagnosticProbeError("model/list", new Error("offline")) : diagnosticProbeOk("model/list");
next.probes["skills/list"] =
overrides.skills === "failed" ? diagnosticProbeError("skills/list", new Error("offline")) : diagnosticProbeOk("skills/list");
next.probes["account/rateLimits/read"] =
overrides.rateLimit === "failed"
? diagnosticProbeError("account/rateLimits/read", new Error("offline"))
: diagnosticProbeOk("account/rateLimits/read");
return next;
}

View file

@ -258,13 +258,13 @@ describe("chat server actions", () => {
stateStore,
vaultPath: "/vault",
currentClient: () => client,
publishAppServerMetadata: () => undefined,
setAppServerMetadata: () => undefined,
});
const diagnostics = createChatServerDiagnosticsActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
publishAppServerMetadata: () => undefined,
setAppServerMetadata: () => undefined,
serverMetadataSnapshot: () => metadata.serverMetadataSnapshot(),
});
@ -273,7 +273,7 @@ describe("chat server actions", () => {
listSkills.mockClear();
readAccountRateLimits.mockClear();
await diagnostics.refreshDiagnosticProbes({ cachedAppServerMetadata: true });
await diagnostics.refreshDiagnosticProbes({ appServerMetadataSnapshot: true });
expect(listModels).not.toHaveBeenCalled();
expect(listSkills).not.toHaveBeenCalled();
@ -305,22 +305,22 @@ describe("chat server actions", () => {
} as unknown as AppServerClient;
const secondClient = {} as unknown as AppServerClient;
let currentClient = firstClient;
const publishAppServerMetadata = vi.fn();
const setAppServerMetadata = vi.fn();
const metadata = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
publishAppServerMetadata: () => undefined,
setAppServerMetadata: () => undefined,
});
const diagnostics = createChatServerDiagnosticsActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
publishAppServerMetadata,
setAppServerMetadata,
serverMetadataSnapshot: () => metadata.serverMetadataSnapshot(),
});
const refreshing = diagnostics.refreshPublishedDiagnosticProbes({ cachedAppServerMetadata: true });
const refreshing = diagnostics.refreshPublishedDiagnosticProbes({ appServerMetadataSnapshot: true });
currentClient = secondClient;
hooksRefresh.resolve({ data: [{ cwd: "/vault", hooks: [{}] }] });
@ -329,7 +329,7 @@ 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(publishAppServerMetadata).not.toHaveBeenCalled();
expect(setAppServerMetadata).not.toHaveBeenCalled();
});
it("loads one app-server metadata snapshot from the initially captured client", async () => {
@ -357,7 +357,7 @@ describe("chat server actions", () => {
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
publishAppServerMetadata: () => undefined,
setAppServerMetadata: () => undefined,
});
const loading = controller.loadAppServerMetadata();
@ -387,12 +387,12 @@ describe("chat server actions", () => {
} as unknown as AppServerClient;
const secondClient = {} as unknown as AppServerClient;
let currentClient = firstClient;
const publishAppServerMetadata = vi.fn();
const setAppServerMetadata = vi.fn();
const controller = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
publishAppServerMetadata,
setAppServerMetadata,
});
const refreshing = controller.refreshPublishedAppServerMetadata();
@ -402,7 +402,7 @@ describe("chat server actions", () => {
await expect(refreshing).resolves.toBeNull();
expect(stateStore.getState().connection.availableModels).toEqual([]);
expect(stateStore.getState().connection.availableSkills).toEqual([]);
expect(publishAppServerMetadata).not.toHaveBeenCalled();
expect(setAppServerMetadata).not.toHaveBeenCalled();
});
it("does not apply refreshed models after the client changes", async () => {
@ -416,7 +416,7 @@ describe("chat server actions", () => {
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
publishAppServerMetadata: () => undefined,
setAppServerMetadata: () => undefined,
});
const refreshing = controller.refreshModels();
@ -436,12 +436,12 @@ describe("chat server actions", () => {
const firstClient = { listSkills } as unknown as AppServerClient;
const secondClient = {} as unknown as AppServerClient;
let currentClient = firstClient;
const publishAppServerMetadata = vi.fn();
const setAppServerMetadata = vi.fn();
const controller = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
publishAppServerMetadata,
setAppServerMetadata,
});
const refreshing = controller.refreshPublishedSkills(true);
@ -452,14 +452,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(publishAppServerMetadata).not.toHaveBeenCalled();
expect(setAppServerMetadata).not.toHaveBeenCalled();
});
it("publishes refreshed rate limits from sparse update notifications", async () => {
const state = createChatState();
const stateStore = createChatStateStore(state);
const rateLimit = rateLimitFixture({ primary: { usedPercent: 64, windowDurationMins: 300, resetsAt: null } });
const publishAppServerMetadata = vi.fn();
const setAppServerMetadata = vi.fn();
const client = {
readAccountRateLimits: vi.fn().mockResolvedValue({ rateLimits: rateLimit, rateLimitsByLimitId: null }),
} as unknown as AppServerClient;
@ -467,13 +467,13 @@ describe("chat server actions", () => {
stateStore,
vaultPath: "/vault",
currentClient: () => client,
publishAppServerMetadata,
setAppServerMetadata,
});
await controller.refreshPublishedRateLimits();
expect(stateStore.getState().connection.rateLimit).toMatchObject({ primary: { usedPercent: 64 } });
expect(publishAppServerMetadata).toHaveBeenCalledWith(expect.objectContaining({ rateLimit }));
expect(setAppServerMetadata).toHaveBeenCalledWith(expect.objectContaining({ rateLimit }));
});
it("keeps the previous rate limit snapshot when sparse update refresh fails", async () => {
@ -484,7 +484,7 @@ describe("chat server actions", () => {
});
state.connection.rateLimit = previousRateLimit;
const stateStore = createChatStateStore(state);
const publishAppServerMetadata = vi.fn();
const setAppServerMetadata = vi.fn();
const client = {
readAccountRateLimits: vi.fn().mockRejectedValue(new Error("offline")),
} as unknown as AppServerClient;
@ -492,14 +492,14 @@ describe("chat server actions", () => {
stateStore,
vaultPath: "/vault",
currentClient: () => client,
publishAppServerMetadata,
setAppServerMetadata,
});
await controller.refreshPublishedRateLimits();
expect(stateStore.getState().connection.rateLimit).toBe(previousRateLimit);
expect(stateStore.getState().connection.serverDiagnostics.probes["account/rateLimits/read"]).toMatchObject({ status: "failed" });
expect(publishAppServerMetadata).not.toHaveBeenCalled();
expect(setAppServerMetadata).not.toHaveBeenCalled();
});
it("does not apply or publish sparse rate limit refreshes after the client changes", async () => {
@ -510,12 +510,12 @@ describe("chat server actions", () => {
} as unknown as AppServerClient;
const secondClient = {} as unknown as AppServerClient;
let currentClient = firstClient;
const publishAppServerMetadata = vi.fn();
const setAppServerMetadata = vi.fn();
const controller = createChatServerMetadataActions({
stateStore,
vaultPath: "/vault",
currentClient: () => currentClient,
publishAppServerMetadata,
setAppServerMetadata,
});
const refreshing = controller.refreshPublishedRateLimits();
@ -528,7 +528,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(publishAppServerMetadata).not.toHaveBeenCalled();
expect(setAppServerMetadata).not.toHaveBeenCalled();
});
it("loads MCP status lines with cached startup diagnostics", async () => {
@ -543,13 +543,13 @@ describe("chat server actions", () => {
stateStore,
vaultPath: "/vault",
currentClient: () => client,
publishAppServerMetadata: () => undefined,
setAppServerMetadata: () => undefined,
});
const controller = createChatServerDiagnosticsActions({
stateStore,
vaultPath: "/vault",
currentClient: () => client,
publishAppServerMetadata: () => undefined,
setAppServerMetadata: () => undefined,
serverMetadataSnapshot: () => metadata.serverMetadataSnapshot(),
});

View file

@ -22,10 +22,10 @@ function controllerForState(
actions: Partial<ConstructorParameters<typeof ChatInboundController>[1]> = {},
): ChatInboundController {
return new ChatInboundController(testStoreForState(state), {
refreshThreads: vi.fn(),
fetchActiveThreads: vi.fn(),
refreshRateLimits: vi.fn(),
refreshSkills: vi.fn(),
publishAppServerMetadata: vi.fn(),
setAppServerMetadata: vi.fn(),
maybeNameThread: vi.fn(),
applyThreadArchived: vi.fn(),
applyThreadRenamed: vi.fn(),
@ -608,8 +608,8 @@ describe("ChatInboundController", () => {
},
]);
const maybeNameThread = vi.fn();
const refreshThreads = vi.fn();
const controller = controllerForState(state, { maybeNameThread, refreshThreads });
const fetchActiveThreads = vi.fn();
const controller = controllerForState(state, { maybeNameThread, fetchActiveThreads });
controller.handleNotification({
method: "turn/completed",
@ -631,7 +631,7 @@ describe("ChatInboundController", () => {
expect(pendingTurnStart(state)).toEqual({ anchorItemId: "local-user-1", promptSubmitHookItemIds: ["hook-hook-1-1"] });
expect(chatStateMessageStreamItems(state).map((item) => item.id)).toEqual(["local-user-1", "hook-hook-1-1"]);
expect(maybeNameThread).not.toHaveBeenCalled();
expect(refreshThreads).not.toHaveBeenCalled();
expect(fetchActiveThreads).not.toHaveBeenCalled();
});
it("refreshes account rate limits after sparse update notifications", () => {
@ -662,8 +662,8 @@ describe("ChatInboundController", () => {
it("records MCP startup status for diagnostics without a chat system message", () => {
const state = createChatState();
const recordMcpStartupStatus = vi.fn();
const publishAppServerMetadata = vi.fn();
const controller = controllerForState(state, { recordMcpStartupStatus, publishAppServerMetadata });
const setAppServerMetadata = vi.fn();
const controller = controllerForState(state, { recordMcpStartupStatus, setAppServerMetadata });
controller.handleNotification({
method: "mcpServer/startupStatus/updated",
@ -676,7 +676,7 @@ describe("ChatInboundController", () => {
} satisfies Extract<ServerNotification, { method: "mcpServer/startupStatus/updated" }>);
expect(recordMcpStartupStatus).toHaveBeenCalledWith("github", "failed", "missing token");
expect(publishAppServerMetadata).toHaveBeenCalledOnce();
expect(setAppServerMetadata).toHaveBeenCalledOnce();
expect(chatStateMessageStreamItems(state)).toEqual([]);
});
});

View file

@ -6,6 +6,8 @@ import { DEFAULT_SETTINGS } from "../../../src/settings/model";
import type { CodexChatHost } from "../../../src/features/chat/application/ports/chat-host";
import { createServerDiagnostics } from "../../../src/domain/server/diagnostics";
import type { Thread } from "../../../src/domain/threads/model";
import type { ModelMetadata } from "../../../src/domain/catalog/metadata";
import type { SharedServerMetadata } from "../../../src/domain/server/metadata";
import { emptyRuntimeConfigSnapshot } from "../../../src/app-server/protocol/runtime-config";
import type { ThreadRecord } from "../../../src/app-server/protocol/thread";
import type { ServerNotification } from "../../../src/app-server/connection/rpc-messages";
@ -162,20 +164,20 @@ describe("CodexChatView connection lifecycle", () => {
});
it("refreshes shared thread lists after connecting", async () => {
const refreshThreads = vi.fn((fetchThreads: () => Promise<unknown>) => fetchThreads() as Promise<never[]>);
const fetchActiveThreads = vi.fn((fetchThreads: () => Promise<unknown>) => fetchThreads() as Promise<never[]>);
const threads = [threadFixture("thread-1")];
const client = connectedClient({
listThreads: vi.fn().mockResolvedValue({ data: threads }),
});
connectionMock.state.client = client;
const view = await chatView({
host: chatHost({ refreshThreads }),
host: chatHost({ fetchActiveThreads }),
});
await view.onOpen();
await view.surface.connect();
expect(refreshThreads).toHaveBeenCalledOnce();
expect(fetchActiveThreads).toHaveBeenCalledOnce();
expect(client.listThreads).toHaveBeenCalledWith("/vault", { archived: false, cursor: null, limit: 100 });
requiredButton(view.containerEl, '[aria-label="Show thread list"]').click();
await waitForAsyncWork(() => {
@ -184,15 +186,15 @@ describe("CodexChatView connection lifecycle", () => {
});
it("publishes app-server metadata after connecting", async () => {
const publishAppServerMetadata = vi.fn();
const setAppServerMetadata = vi.fn();
connectionMock.state.client = connectedClient();
const view = await chatView({
host: chatHost({ publishAppServerMetadata }),
host: chatHost({ setAppServerMetadata }),
});
await view.surface.connect();
expect(publishAppServerMetadata).toHaveBeenCalledWith(
expect(setAppServerMetadata).toHaveBeenCalledWith(
expect.objectContaining({
runtimeConfig: expect.any(Object),
availableModels: [],
@ -352,17 +354,19 @@ describe("CodexChatView connection lifecycle", () => {
});
it("formats the panel title from listed thread metadata", async () => {
const view = await chatView();
const host = chatHost();
const view = await chatView({ host });
await view.setState({ threadId: "thread-named" }, {} as never);
view.surface.applyThreadListSnapshot([panelThread({ id: "thread-named", name: "作業メモ" })]);
await view.onOpen();
host.threadCatalog.setActiveThreads([panelThread({ id: "thread-named", name: "作業メモ" })]);
expect(view.getDisplayText()).toBe("Codex: 作業メモ");
view.surface.applyThreadListSnapshot([panelThread({ id: "thread-named", name: null, preview: "初回依頼" })]);
host.threadCatalog.setActiveThreads([panelThread({ id: "thread-named", name: null, preview: "初回依頼" })]);
expect(view.getDisplayText()).toBe("Codex: 初回依頼");
await view.setState({ threadId: "019e061e-0000-7000-8000-000000000001" }, {} as never);
view.surface.applyThreadListSnapshot([]);
host.threadCatalog.setActiveThreads([]);
expect(view.getDisplayText()).toBe("Codex: 019e061e");
});
@ -415,8 +419,8 @@ describe("CodexChatView connection lifecycle", () => {
const cachedThread = threadFixture("thread-cached");
const view = await chatView({
host: chatHost({
cachedThreads: vi.fn(() => [cachedThread] as never[]),
cachedAppServerMetadata: vi.fn(
activeThreadsSnapshot: vi.fn(() => [cachedThread] as never[]),
appServerMetadataSnapshot: vi.fn(
() =>
({
runtimeConfig: { ...emptyRuntimeConfigSnapshot(), model: "gpt-cached" },
@ -710,14 +714,15 @@ describe("CodexChatView connection lifecycle", () => {
});
it("does not use restored thread title as a composer name before an explicit rename notification", async () => {
const view = await chatView();
const host = chatHost();
const view = await chatView({ host });
await view.setState({ threadId: "thread-1", threadTitle: "Restored title" }, {} as never);
await view.onOpen();
expect(composerPlaceholder(view)).toBe("Ask Codex to work on this task...");
view.surface.applyThreadListSnapshot([panelThread({ id: "thread-1", name: "Explicit name" })]);
host.threadCatalog.setActiveThreads([panelThread({ id: "thread-1", name: "Explicit name" })]);
view.surface.applyThreadRenamed("thread-1", "Explicit name");
await waitForAsyncWork(() => {
@ -755,18 +760,19 @@ describe("CodexChatView connection lifecycle", () => {
it("updates the composer placeholder from shared rename notifications", async () => {
const client = connectedClient();
connectionMock.state.client = client;
const view = await chatView();
const host = chatHost();
const view = await chatView({ host });
await view.onOpen();
await view.surface.openThread("thread-1");
view.surface.applyThreadListSnapshot([panelThread({ id: "thread-1", name: "Renamed thread" })]);
host.threadCatalog.setActiveThreads([panelThread({ id: "thread-1", name: "Renamed thread" })]);
view.surface.applyThreadRenamed("thread-1", "Renamed thread");
await waitForAsyncWork(() => {
expect(composerPlaceholder(view)).toBe("Ask Codex to work on “Renamed thread”...");
});
view.surface.applyThreadListSnapshot([panelThread({ id: "thread-1", name: null })]);
host.threadCatalog.setActiveThreads([panelThread({ id: "thread-1", name: null })]);
view.surface.applyThreadRenamed("thread-1", null);
await waitForAsyncWork(() => {
@ -777,7 +783,8 @@ describe("CodexChatView connection lifecycle", () => {
it("keeps composer draft and selection while updating the placeholder", async () => {
const client = connectedClient();
connectionMock.state.client = client;
const view = await chatView();
const host = chatHost();
const view = await chatView({ host });
await view.onOpen();
await view.surface.openThread("thread-1");
@ -788,7 +795,7 @@ describe("CodexChatView connection lifecycle", () => {
});
composer.setSelectionRange(5, 9);
view.surface.applyThreadListSnapshot([panelThread({ id: "thread-1", name: "Renamed thread" })]);
host.threadCatalog.setActiveThreads([panelThread({ id: "thread-1", name: "Renamed thread" })]);
view.surface.applyThreadRenamed("thread-1", "Renamed thread");
await waitForAsyncWork(() => {
@ -1290,14 +1297,21 @@ interface ChatHostFixtureOverrides {
renameThreadInCatalog?: CodexChatHost["threadCatalog"]["renameThreadInCatalog"];
refreshFromOpenSurface?: CodexChatHost["threadCatalog"]["refreshFromOpenSurface"];
refreshThreadsViewLiveState?: CodexChatHost["threadCatalog"]["refreshThreadsViewLiveState"];
applyThreads?: CodexChatHost["threadCatalog"]["applyThreads"];
publishAppServerMetadata?: CodexChatHost["threadCatalog"]["publishAppServerMetadata"];
refreshThreads?: CodexChatHost["threadCatalog"]["refreshThreads"];
cachedThreads?: CodexChatHost["threadCatalog"]["cachedThreads"];
cachedAppServerMetadata?: CodexChatHost["threadCatalog"]["cachedAppServerMetadata"];
setActiveThreads?: CodexChatHost["threadCatalog"]["setActiveThreads"];
setAppServerMetadata?: CodexChatHost["threadCatalog"]["setAppServerMetadata"];
fetchActiveThreads?: CodexChatHost["threadCatalog"]["fetchActiveThreads"];
activeThreadsSnapshot?: CodexChatHost["threadCatalog"]["activeThreadsSnapshot"];
appServerMetadataSnapshot?: CodexChatHost["threadCatalog"]["appServerMetadataSnapshot"];
modelsSnapshot?: CodexChatHost["threadCatalog"]["modelsSnapshot"];
}
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 settings = {
...DEFAULT_SETTINGS,
codexPath: "codex",
@ -1319,15 +1333,47 @@ function chatHost(overrides: ChatHostFixtureOverrides = {}): CodexChatHost {
renameThreadInCatalog: overrides.renameThreadInCatalog ?? vi.fn(),
refreshFromOpenSurface: overrides.refreshFromOpenSurface ?? vi.fn(),
refreshThreadsViewLiveState: overrides.refreshThreadsViewLiveState ?? vi.fn(),
applyThreads: overrides.applyThreads ?? vi.fn(),
publishAppServerMetadata: overrides.publishAppServerMetadata ?? vi.fn(),
refreshThreads:
overrides.refreshThreads ??
setActiveThreads:
overrides.setActiveThreads ??
((threads) => {
activeThreads = threads;
for (const listener of activeThreadListeners) listener(threads);
}),
setAppServerMetadata:
overrides.setAppServerMetadata ??
((nextMetadata) => {
metadata = nextMetadata;
for (const listener of metadataListeners) listener(nextMetadata);
}),
fetchActiveThreads:
overrides.fetchActiveThreads ??
(vi.fn(
(fetchThreads: () => Promise<unknown>) => fetchThreads() as Promise<never[]>,
) as CodexChatHost["threadCatalog"]["refreshThreads"]),
cachedThreads: overrides.cachedThreads ?? vi.fn(() => null),
cachedAppServerMetadata: overrides.cachedAppServerMetadata ?? vi.fn(() => null),
) as CodexChatHost["threadCatalog"]["fetchActiveThreads"]),
activeThreadsSnapshot: overrides.activeThreadsSnapshot ?? vi.fn(() => activeThreads),
appServerMetadataSnapshot: overrides.appServerMetadataSnapshot ?? vi.fn(() => metadata),
modelsSnapshot: overrides.modelsSnapshot ?? vi.fn(() => models),
observeActiveThreads: (listener, options = {}) => {
activeThreadListeners.add(listener);
if ((options.emitCurrent ?? true) && activeThreads) listener(activeThreads);
return () => {
activeThreadListeners.delete(listener);
};
},
observeAppServerMetadata: (listener, options = {}) => {
metadataListeners.add(listener);
if ((options.emitCurrent ?? true) && metadata) listener(metadata);
return () => {
metadataListeners.delete(listener);
};
},
observeModels: (listener, options = {}) => {
modelListeners.add(listener);
if ((options.emitCurrent ?? true) && models) listener(models);
return () => {
modelListeners.delete(listener);
};
},
},
};
}

View file

@ -110,8 +110,8 @@ function threadPickerHost(threads: readonly Thread[]): TestThreadPickerHost {
openedCurrent,
openedAvailable,
threadCatalog: {
cachedThreads: () => threads,
refreshThreads: async () => threads,
activeThreadsSnapshot: () => threads,
fetchActiveThreads: async () => threads,
},
openThreadInCurrentView: async (threadId) => {
openedCurrent.push(threadId);

View file

@ -237,20 +237,20 @@ describe("CodexThreadsView", () => {
it("refreshes thread lists through the plugin coordinator", async () => {
const threads = [threadFixture({ id: "thread", preview: "Thread preview" })];
const refreshThreads = vi.fn((fetchThreads: () => Promise<unknown>) => fetchThreads() as Promise<never[]>);
const fetchActiveThreads = vi.fn((fetchThreads: () => Promise<unknown>) => fetchThreads() as Promise<never[]>);
connectionMock.state.client = clientFixture({
listThreads: vi.fn().mockResolvedValue({ data: threads }),
});
const host = threadsHost({
threadCatalog: {
refreshThreads,
fetchActiveThreads,
},
});
const view = await threadsView(host);
await view.refresh();
expect(refreshThreads).toHaveBeenCalledOnce();
expect(fetchActiveThreads).toHaveBeenCalledOnce();
expect(view.containerEl.textContent).toContain("Thread preview");
});
@ -266,7 +266,7 @@ describe("CodexThreadsView", () => {
const view = await threadsView(
threadsHost({
threadCatalog: {
cachedThreads: vi.fn(() => [threadFixture({ id: "cached", preview: "Cached thread" })]),
activeThreadsSnapshot: vi.fn(() => [threadFixture({ id: "cached", preview: "Cached thread" })]),
},
}),
);
@ -427,8 +427,9 @@ function threadsHost(overrides: Record<string, unknown> = {}) {
archiveThreadInCatalog: vi.fn(),
renameThreadInCatalog: vi.fn(),
refreshFromOpenSurface: vi.fn(),
refreshThreads: vi.fn((fetchThreads: () => Promise<unknown>) => fetchThreads() as Promise<never[]>),
cachedThreads: vi.fn(() => null),
fetchActiveThreads: vi.fn((fetchThreads: () => Promise<unknown>) => fetchThreads() as Promise<never[]>),
activeThreadsSnapshot: vi.fn(() => null),
observeActiveThreads: vi.fn(() => () => undefined),
...threadCatalogOverrides,
},
...hostOverrides,

View file

@ -419,8 +419,8 @@ describe("CodexPanelPlugin boot restored panel loading", () => {
);
const secondFetch = vi.fn().mockResolvedValue([thread("second")]);
const first = threadCatalog(plugin).refreshThreads(fetchThreads);
const second = threadCatalog(plugin).refreshThreads(secondFetch);
const first = threadCatalog(plugin).fetchActiveThreads(fetchThreads);
const second = threadCatalog(plugin).fetchActiveThreads(secondFetch);
expect(fetchThreads).toHaveBeenCalledOnce();
expect(secondFetch).not.toHaveBeenCalled();
@ -428,7 +428,7 @@ describe("CodexPanelPlugin boot restored panel loading", () => {
await expect(first).resolves.toEqual([thread("first")]);
await expect(second).resolves.toEqual([thread("first")]);
expect(threadCatalog(plugin).cachedThreads()).toEqual([thread("first")]);
expect(threadCatalog(plugin).activeThreadsSnapshot()).toEqual([thread("first")]);
});
it("keeps shared thread list refreshes separate across app-server cache contexts", async () => {
@ -443,29 +443,29 @@ describe("CodexPanelPlugin boot restored panel loading", () => {
);
const secondFetch = vi.fn().mockResolvedValue([thread("second")]);
const first = threadCatalog(plugin).refreshThreads(firstFetch);
const first = threadCatalog(plugin).fetchActiveThreads(firstFetch);
plugin.settings.codexPath = "codex-b";
const second = threadCatalog(plugin).refreshThreads(secondFetch);
const second = threadCatalog(plugin).fetchActiveThreads(secondFetch);
expect(firstFetch).toHaveBeenCalledOnce();
expect(secondFetch).toHaveBeenCalledOnce();
await expect(second).resolves.toEqual([thread("second")]);
expect(threadCatalog(plugin).cachedThreads()).toEqual([thread("second")]);
expect(threadCatalog(plugin).activeThreadsSnapshot()).toEqual([thread("second")]);
resolveFirst([thread("first")]);
await expect(first).resolves.toEqual([thread("first")]);
expect(threadCatalog(plugin).cachedThreads()).toEqual([thread("second")]);
expect(threadCatalog(plugin).activeThreadsSnapshot()).toEqual([thread("second")]);
plugin.settings.codexPath = "codex-a";
expect(threadCatalog(plugin).cachedThreads()).toBeNull();
expect(threadCatalog(plugin).activeThreadsSnapshot()).toEqual([thread("first")]);
});
it("keeps the previous shared thread list when refresh fails", async () => {
const plugin = await pluginWithLeaves([]);
await threadCatalog(plugin).refreshThreads(() => Promise.resolve([thread("cached")]));
await threadCatalog(plugin).fetchActiveThreads(() => Promise.resolve([thread("cached")]));
await expect(threadCatalog(plugin).refreshThreads(() => Promise.reject(new Error("boom")))).rejects.toThrow("boom");
await expect(threadCatalog(plugin).fetchActiveThreads(() => Promise.reject(new Error("boom")))).rejects.toThrow("boom");
expect(threadCatalog(plugin).cachedThreads()).toEqual([thread("cached")]);
expect(threadCatalog(plugin).activeThreadsSnapshot()).toEqual([thread("cached")]);
});
it("refreshes shared thread lists from a connected chat panel", async () => {
@ -606,11 +606,15 @@ function chatHostFixture(): CodexChatHost {
renameThreadInCatalog: vi.fn(),
refreshFromOpenSurface: vi.fn(),
refreshThreadsViewLiveState: vi.fn(),
applyThreads: vi.fn(),
publishAppServerMetadata: vi.fn(),
refreshThreads: vi.fn((fetchThreads: () => Promise<unknown>) => fetchThreads() as Promise<never[]>),
cachedThreads: vi.fn(() => null),
cachedAppServerMetadata: vi.fn(() => null),
setActiveThreads: vi.fn(),
setAppServerMetadata: vi.fn(),
fetchActiveThreads: vi.fn((fetchThreads: () => Promise<unknown>) => fetchThreads() as Promise<never[]>),
activeThreadsSnapshot: vi.fn(() => null),
appServerMetadataSnapshot: vi.fn(() => null),
modelsSnapshot: vi.fn(() => null),
observeActiveThreads: vi.fn(() => () => undefined),
observeAppServerMetadata: vi.fn(() => () => undefined),
observeModels: vi.fn(() => () => undefined),
},
};
}

View file

@ -253,6 +253,10 @@ export class PluginSettingTab {
display(): void {
// Test mock placeholder.
}
hide(): void {
this.containerEl.empty();
}
}
export class Setting {

View file

@ -195,6 +195,8 @@ describe("settings tab", () => {
it("clears dynamic settings data when the Codex executable changes", async () => {
const saveSettings = vi.fn().mockResolvedValue(undefined);
const notifyAppServerQueryContextChanged = vi.fn();
const refreshOpenViews = vi.fn();
const oldClient = settingsClient({
models: [model("gpt-old")],
hooks: [hook({ key: "hook-old", command: "old hook", currentHash: "oldhash" })],
@ -208,7 +210,7 @@ describe("settings tab", () => {
withShortLivedAppServerClientMock
.mockImplementationOnce((_codexPath: string, _cwd: string, operation: (client: unknown) => Promise<unknown>) => operation(oldClient))
.mockImplementationOnce((_codexPath: string, _cwd: string, operation: (client: unknown) => Promise<unknown>) => operation(newClient));
const tab = newSettingsTab({ saveSettings });
const tab = newSettingsTab({ saveSettings, notifyAppServerQueryContextChanged, refreshOpenViews });
tab.display();
await flushPromises();
@ -223,6 +225,8 @@ describe("settings tab", () => {
await flushPromises();
expect(saveSettings).toHaveBeenCalledOnce();
expect(notifyAppServerQueryContextChanged).toHaveBeenCalledOnce();
expect(refreshOpenViews).toHaveBeenCalledOnce();
expect(withShortLivedAppServerClientMock).toHaveBeenCalledTimes(1);
expect(tab.containerEl.textContent).not.toContain("gpt-old");
expect(tab.containerEl.textContent).not.toContain("Old archived");
@ -236,6 +240,21 @@ describe("settings tab", () => {
expect(tab.containerEl.textContent).toContain("New archived");
});
it("unsubscribes model updates when the settings tab is hidden", () => {
withShortLivedAppServerClientMock.mockImplementation(
(_codexPath: string, _cwd: string, operation: (client: unknown) => Promise<unknown>) => operation(settingsClient()),
);
const unsubscribe = vi.fn();
const observeModels = vi.fn(() => unsubscribe);
const tab = newSettingsTab({ observeModels });
tab.hide();
tab.display();
expect(observeModels).toHaveBeenCalledTimes(2);
expect(unsubscribe).toHaveBeenCalledOnce();
});
it("ignores stale settings data refresh results after a newer refresh completes", async () => {
const firstModels = deferred<{ data: CatalogModel[] }>();
const firstClient = settingsClient({
@ -358,12 +377,12 @@ describe("settings tab", () => {
});
it("uses cached models initially and publishes refreshed models", async () => {
const publishModels = vi.fn();
const setModels = vi.fn();
const client = settingsClient({ models: [model("gpt-5.5")] });
withShortLivedAppServerClientMock.mockImplementation(
(_codexPath: string, _cwd: string, operation: (client: unknown) => Promise<unknown>) => operation(client),
);
const tab = newSettingsTab({ cachedModels: modelMetadataFromCatalogModels([model("gpt-cached")]), publishModels });
const tab = newSettingsTab({ modelsSnapshot: modelMetadataFromCatalogModels([model("gpt-cached")]), setModels });
tab.display();
@ -371,13 +390,40 @@ describe("settings tab", () => {
await flushPromises();
expect(publishModels).toHaveBeenCalledWith(modelMetadataFromCatalogModels([model("gpt-5.5")]));
expect(setModels).toHaveBeenCalledWith(modelMetadataFromCatalogModels([model("gpt-5.5")]));
expect(tab.containerEl.textContent).toContain("gpt-5.5");
});
it("replaces stale cached model options with an empty successful refresh while preserving saved values", async () => {
const setModels = vi.fn();
const client = settingsClient({ models: [] });
withShortLivedAppServerClientMock.mockImplementation(
(_codexPath: string, _cwd: string, operation: (client: unknown) => Promise<unknown>) => operation(client),
);
const tab = newSettingsTab({
modelsSnapshot: modelMetadataFromCatalogModels([model("gpt-cached")]),
setModels,
settings: {
threadNamingModel: "gpt-saved",
rewriteSelectionModel: "gpt-cached",
},
});
tab.display();
expect(selectOptions(tab, "Automatic thread naming")).toEqual(["Codex default", "gpt-saved (saved)", "gpt-cached"]);
expect(selectOptions(tab, "Selection rewrite")).toEqual(["Codex default", "gpt-cached"]);
await flushPromises();
expect(setModels).toHaveBeenCalledWith([]);
expect(selectOptions(tab, "Automatic thread naming")).toEqual(["Codex default", "gpt-saved (saved)"]);
expect(selectOptions(tab, "Selection rewrite")).toEqual(["Codex default", "gpt-cached (saved)"]);
});
it("uses model-provided reasoning efforts in helper settings while preserving saved unknown values", async () => {
const tab = newSettingsTab({
cachedModels: modelMetadataFromCatalogModels([model("gpt-5.5", false, false, ["extreme"])]),
modelsSnapshot: modelMetadataFromCatalogModels([model("gpt-5.5", false, false, ["extreme"])]),
settings: {
threadNamingModel: "gpt-5.5",
threadNamingEffort: "saved-custom-effort",
@ -548,8 +594,10 @@ function newSettingsTab(
options: {
saveSettings?: () => Promise<void>;
sendShortcut?: "enter" | "mod-enter";
cachedModels?: ModelMetadata[];
publishModels?: (models: ModelMetadata[]) => void;
modelsSnapshot?: ModelMetadata[];
setModels?: (models: ModelMetadata[]) => void;
observeModels?: CodexPanelSettingTabHost["threadCatalog"]["observeModels"];
notifyAppServerQueryContextChanged?: () => void;
refreshOpenViews?: () => void;
refreshFromOpenSurface?: () => void;
settings?: Partial<{
@ -567,8 +615,10 @@ function settingsTabHost(
options: {
saveSettings?: () => Promise<void>;
sendShortcut?: "enter" | "mod-enter";
cachedModels?: ModelMetadata[];
publishModels?: (models: ModelMetadata[]) => void;
modelsSnapshot?: ModelMetadata[];
setModels?: (models: ModelMetadata[]) => void;
observeModels?: CodexPanelSettingTabHost["threadCatalog"]["observeModels"];
notifyAppServerQueryContextChanged?: () => void;
refreshOpenViews?: () => void;
refreshFromOpenSurface?: () => void;
settings?: Partial<{
@ -599,8 +649,10 @@ function settingsTabHost(
refreshOpenViews: options.refreshOpenViews ?? vi.fn(),
threadCatalog: {
refreshFromOpenSurface: options.refreshFromOpenSurface ?? vi.fn(),
cachedModels: vi.fn(() => options.cachedModels ?? []),
publishModels: options.publishModels ?? vi.fn(),
modelsSnapshot: vi.fn(() => options.modelsSnapshot ?? []),
setModels: options.setModels ?? vi.fn(),
observeModels: options.observeModels ?? vi.fn(() => () => undefined),
notifyAppServerQueryContextChanged: options.notifyAppServerQueryContextChanged ?? vi.fn(),
},
};
}

View file

@ -1,6 +1,6 @@
import { describe, expect, it, vi, type Mock } from "vitest";
import { SharedAppServerCache } from "../../src/app-server/services/shared-cache";
import { AppServerQueryCache } from "../../src/app-server/query/cache";
import type { ModelMetadata } from "../../src/domain/catalog/metadata";
import { createServerDiagnostics } from "../../src/domain/server/diagnostics";
import type { SharedServerMetadata } from "../../src/domain/server/metadata";
@ -9,71 +9,128 @@ import { SharedThreadCatalog } from "../../src/workspace/shared-thread-catalog";
import type { ThreadSurfaceActions } from "../../src/workspace/thread-surface-actions";
type MockSurfaceActions = ThreadSurfaceActions & {
applyThreadListSnapshot: Mock<(threads: readonly Thread[]) => void>;
applyThreadArchived: Mock<(threadId: string, options?: { closeOpenPanels?: boolean }) => void>;
applyThreadRenamed: Mock<(threadId: string, name: string | null) => void>;
publishAppServerMetadata: Mock<(metadata: SharedServerMetadata) => void>;
publishModels: Mock<(models: readonly ModelMetadata[]) => void>;
};
describe("SharedThreadCatalog", () => {
it("applies thread snapshots to the shared cache and open surfaces", () => {
const { catalog, surfaces } = catalogFixture();
it("applies thread snapshots to the shared cache and active observers", () => {
const { catalog } = catalogFixture();
const threads = [thread("thread")];
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.applyThreads(threads);
catalog.setActiveThreads(threads);
expect(catalog.cachedThreads()).toEqual(threads);
expect(surfaces.applyThreadListSnapshot).toHaveBeenCalledWith(threads);
expect(catalog.activeThreadsSnapshot()).toEqual(threads);
expect(listener).toHaveBeenCalledWith(threads);
});
it("refreshes thread snapshots through the cache single-flight and publishes the snapshot once", async () => {
const { catalog, surfaces } = catalogFixture();
it("refreshes thread snapshots through the cache single-flight and notifies observers once", async () => {
const { catalog } = catalogFixture();
const fetchThreads = vi.fn().mockResolvedValue([thread("thread")]);
const listener = vi.fn();
catalog.observeActiveThreads(listener);
const first = catalog.refreshThreads(fetchThreads);
const second = catalog.refreshThreads(fetchThreads);
const first = catalog.fetchActiveThreads(fetchThreads);
const second = catalog.fetchActiveThreads(fetchThreads);
await expect(first).resolves.toEqual([thread("thread")]);
await expect(second).resolves.toEqual([thread("thread")]);
expect(fetchThreads).toHaveBeenCalledOnce();
expect(catalog.cachedThreads()).toEqual([thread("thread")]);
expect(surfaces.applyThreadListSnapshot).toHaveBeenCalledOnce();
expect(catalog.activeThreadsSnapshot()).toEqual([thread("thread")]);
expect(listener).toHaveBeenCalledOnce();
});
it("publishes metadata and model snapshots to cache and surfaces", () => {
const { catalog, surfaces } = catalogFixture();
it("does not notify stale thread observers after the app-server query context changes", async () => {
const context = { codexPath: "codex-a", vaultPath: "/vault" };
const surfaces = surfaceActions();
const catalog = new SharedThreadCatalog({
cache: new AppServerQueryCache(),
surfaces,
context: () => context,
});
let resolveThreads!: (threads: Thread[]) => void;
const listener = vi.fn();
catalog.observeActiveThreads(listener);
const fetch = catalog.fetchActiveThreads(
() =>
new Promise<Thread[]>((resolve) => {
resolveThreads = resolve;
}),
);
context.codexPath = "codex-b";
resolveThreads([thread("stale")]);
await expect(fetch).resolves.toEqual([thread("stale")]);
expect(listener).not.toHaveBeenCalled();
context.codexPath = "codex-a";
expect(catalog.activeThreadsSnapshot()).toEqual([thread("stale")]);
});
it("resubscribes active observers when the app-server query context changes", () => {
const context = { codexPath: "codex-a", vaultPath: "/vault" };
const catalog = new SharedThreadCatalog({
cache: new AppServerQueryCache(),
surfaces: surfaceActions(),
context: () => context,
});
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.setActiveThreads([thread("a")]);
context.codexPath = "codex-b";
catalog.notifyAppServerQueryContextChanged();
catalog.setActiveThreads([thread("b")]);
expect(listener).toHaveBeenLastCalledWith([thread("b")]);
context.codexPath = "codex-a";
catalog.notifyAppServerQueryContextChanged();
expect(listener).toHaveBeenLastCalledWith([thread("a")]);
});
it("publishes metadata and model snapshots to cache and observers", () => {
const { catalog } = catalogFixture();
const metadata = serverMetadata({ availableModels: [model("gpt-test")] });
const models = [model("gpt-other")];
const metadataListener = vi.fn();
const modelListener = vi.fn();
catalog.observeAppServerMetadata(metadataListener);
catalog.observeModels(modelListener);
catalog.publishAppServerMetadata(metadata);
catalog.publishModels(models);
catalog.setAppServerMetadata(metadata);
catalog.setModels(models);
expect(catalog.cachedAppServerMetadata()).toEqual({ ...metadata, availableModels: models });
expect(catalog.cachedModels()).toEqual(models);
expect(surfaces.publishAppServerMetadata).toHaveBeenCalledWith(metadata);
expect(surfaces.publishModels).toHaveBeenCalledWith(models);
expect(catalog.appServerMetadataSnapshot()).toEqual({ ...metadata, availableModels: models });
expect(catalog.modelsSnapshot()).toEqual(models);
expect(metadataListener).toHaveBeenLastCalledWith({ ...metadata, availableModels: models });
expect(modelListener).toHaveBeenCalledWith(models);
});
it("applies known rename mutations to cache and surfaces", () => {
const { catalog, surfaces } = catalogFixture();
catalog.applyThreads([thread("thread"), thread("other")]);
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.setActiveThreads([thread("thread"), thread("other")]);
catalog.renameThreadInCatalog("thread", "Renamed");
expect(catalog.cachedThreads()).toEqual([{ ...thread("thread"), name: "Renamed" }, thread("other")]);
expect(surfaces.applyThreadListSnapshot).toHaveBeenLastCalledWith([{ ...thread("thread"), name: "Renamed" }, thread("other")]);
expect(catalog.activeThreadsSnapshot()).toEqual([{ ...thread("thread"), name: "Renamed" }, thread("other")]);
expect(listener).toHaveBeenLastCalledWith([{ ...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();
catalog.applyThreads([thread("thread"), thread("other")]);
const listener = vi.fn();
catalog.observeActiveThreads(listener);
catalog.setActiveThreads([thread("thread"), thread("other")]);
catalog.archiveThreadInCatalog("thread", { closeOpenPanels: true });
expect(catalog.cachedThreads()).toEqual([thread("other")]);
expect(surfaces.applyThreadListSnapshot).toHaveBeenLastCalledWith([thread("other")]);
expect(catalog.activeThreadsSnapshot()).toEqual([thread("other")]);
expect(listener).toHaveBeenLastCalledWith([thread("other")]);
expect(surfaces.applyThreadArchived).toHaveBeenCalledWith("thread", { closeOpenPanels: true });
});
});
@ -81,7 +138,7 @@ describe("SharedThreadCatalog", () => {
function catalogFixture() {
const surfaces = surfaceActions();
const catalog = new SharedThreadCatalog({
cache: new SharedAppServerCache(),
cache: new AppServerQueryCache(),
surfaces,
context: () => ({ codexPath: "codex", vaultPath: "/vault" }),
});
@ -92,11 +149,8 @@ function surfaceActions(): MockSurfaceActions {
return {
refreshOpenViews: vi.fn(),
invalidateThreadsFromOpenSurface: vi.fn(),
applyThreadListSnapshot: vi.fn(),
applyThreadArchived: vi.fn(),
applyThreadRenamed: vi.fn(),
publishAppServerMetadata: vi.fn(),
publishModels: vi.fn(),
refreshThreadsViewLiveState: vi.fn(),
};
}