From de913e9d9dfc8c854c4458e1ec47ef96e123ef30 Mon Sep 17 00:00:00 2001 From: murashit Date: Mon, 15 Jun 2026 17:44:34 +0900 Subject: [PATCH] Adopt app-server query state --- package-lock.json | 11 + package.json | 1 + src/app-server/query/cache.ts | 171 ++++++++++++++ src/app-server/query/keys.ts | 52 +++++ src/app-server/query/snapshots.ts | 90 ++++++++ src/app-server/services/shared-cache-state.ts | 215 ------------------ src/app-server/services/shared-cache.ts | 86 ------- .../chat/app-server/actions/diagnostics.ts | 8 +- .../chat/app-server/actions/metadata.ts | 18 +- .../chat/app-server/inbound/controller.ts | 8 +- .../connection/connection-controller.ts | 4 +- .../chat/application/ports/chat-host.ts | 14 +- .../chat/application/threads/composition.ts | 4 +- src/features/chat/host/connection-bundle.ts | 22 +- src/features/chat/host/runtime.ts | 4 +- src/features/chat/host/session.ts | 49 +++- src/features/chat/host/surface-handle.ts | 6 - src/features/thread-picker/modal.ts | 6 +- src/features/threads-view/session.ts | 23 +- src/features/threads-view/view.ts | 5 - src/plugin-runtime.ts | 13 +- src/settings/dynamic-data-controller.ts | 29 ++- src/settings/tab.ts | 8 + src/workspace/shared-thread-catalog.ts | 127 ++++++++--- src/workspace/thread-surface-actions.ts | 27 --- ...ared-cache.test.ts => query-cache.test.ts} | 135 ++++++----- tests/app-server/shared-cache-state.test.ts | 211 ----------------- .../server-actions/server-actions.test.ts | 54 ++--- .../chat/protocol/inbound/controller.test.ts | 16 +- tests/features/chat/view-connection.test.ts | 108 ++++++--- tests/features/thread-picker/modal.test.ts | 4 +- tests/features/threads-view/view.test.ts | 13 +- tests/main.test.ts | 36 +-- tests/mocks/obsidian.ts | 4 + tests/settings/settings-tab.test.ts | 74 +++++- tests/workspace/shared-thread-catalog.test.ts | 120 +++++++--- 36 files changed, 926 insertions(+), 850 deletions(-) create mode 100644 src/app-server/query/cache.ts create mode 100644 src/app-server/query/keys.ts create mode 100644 src/app-server/query/snapshots.ts delete mode 100644 src/app-server/services/shared-cache-state.ts delete mode 100644 src/app-server/services/shared-cache.ts rename tests/app-server/{shared-cache.test.ts => query-cache.test.ts} (55%) delete mode 100644 tests/app-server/shared-cache-state.test.ts diff --git a/package-lock.json b/package-lock.json index 957f87d7..9ba96fea 100644 --- a/package-lock.json +++ b/package-lock.json @@ -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", diff --git a/package.json b/package.json index c3c099dd..a5946dc0 100644 --- a/package.json +++ b/package.json @@ -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" } diff --git a/src/app-server/query/cache.ts b/src/app-server/query/cache.ts new file mode 100644 index 00000000..768cd449 --- /dev/null +++ b/src/app-server/query/cache.ts @@ -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(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): Promise { + 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(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(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( + queryKey: readonly unknown[], + clone: (value: T) => T, + listener: (value: T) => void, + options: { emitCurrent?: boolean }, + ): () => void { + const observer = new QueryObserver(this.client, { + queryKey, + enabled: false, + }); + const emit = (result: QueryObserverResult): 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, + }, + }, + }); +} diff --git a/src/app-server/query/keys.ts b/src/app-server/query/keys.ts new file mode 100644 index 00000000..396564c0 --- /dev/null +++ b/src/app-server/query/keys.ts @@ -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; +} diff --git a/src/app-server/query/snapshots.ts b/src/app-server/query/snapshots.ts new file mode 100644 index 00000000..05474b2a --- /dev/null +++ b/src/app-server/query/snapshots.ts @@ -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 })); +} diff --git a/src/app-server/services/shared-cache-state.ts b/src/app-server/services/shared-cache-state.ts deleted file mode 100644 index 23a35a46..00000000 --- a/src/app-server/services/shared-cache-state.ts +++ /dev/null @@ -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 = { kind: "unloaded" } | { kind: "loaded"; context: SharedAppServerCacheContext; data: T }; - -export interface SharedAppServerState { - threads: SharedCache; - appServerMetadata: SharedCache; - availableModels: SharedCache; -} - -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; -} diff --git a/src/app-server/services/shared-cache.ts b/src/app-server/services/shared-cache.ts deleted file mode 100644 index e99b949c..00000000 --- a/src/app-server/services/shared-cache.ts +++ /dev/null @@ -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 }; - -export class SharedAppServerCache { - private state: SharedAppServerState = createSharedAppServerState(); - private threadListRefreshLifecycle: ThreadListRefreshLifecycleState = { kind: "idle" }; - - refreshThreadList( - context: SharedAppServerCacheContext, - fetchThreads: () => Promise, - onSnapshot?: (threads: readonly Thread[]) => void, - ): Promise { - 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); - } -} diff --git a/src/features/chat/app-server/actions/diagnostics.ts b/src/features/chat/app-server/actions/diagnostics.ts index 1484e08c..503d8677 100644 --- a/src/features/chat/app-server/actions/diagnostics.ts +++ b/src/features/chat/app-server/actions/diagnostics.ts @@ -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[] = []; - if (!options.cachedAppServerMetadata) { + if (!options.appServerMetadataSnapshot) { probes.push( probeDiagnostic( "model/list", @@ -139,7 +139,7 @@ async function refreshPublishedDiagnosticProbes( options: RefreshDiagnosticProbesOptions = {}, ): Promise { if (!(await refreshDiagnosticProbes(host, options))) return; - host.publishAppServerMetadata(host.serverMetadataSnapshot()); + host.setAppServerMetadata(host.serverMetadataSnapshot()); } async function mcpStatusLines(host: ChatServerDiagnosticsActionsHost): Promise { diff --git a/src/features/chat/app-server/actions/metadata.ts b/src/features/chat/app-server/actions/metadata.ts index a611686b..6412e830 100644 --- a/src/features/chat/app-server/actions/metadata.ts +++ b/src/features/chat/app-server/actions/metadata.ts @@ -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; refreshAppServerMetadata: () => Promise; refreshPublishedAppServerMetadata: () => Promise; - publishAppServerMetadataSnapshot: () => void; + setAppServerMetadataSnapshot: () => void; refreshModels: () => Promise; loadModels: () => Promise; refreshSkills: (forceReload?: boolean) => Promise; @@ -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 { 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 { @@ -146,7 +146,7 @@ async function refreshSkills(host: ChatServerMetadataActionsHost, forceReload = async function refreshPublishedSkills(host: ChatServerMetadataActionsHost, forceReload = false): Promise { if (!(await refreshSkills(host, forceReload))) return; - publishAppServerMetadataSnapshot(host); + setAppServerMetadataSnapshot(host); } async function loadSkills(host: ChatServerMetadataActionsHost, forceReload = false): Promise { @@ -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 }); diff --git a/src/features/chat/app-server/inbound/controller.ts b/src/features/chat/app-server/inbound/controller.ts index 834d3272..a080b560 100644 --- a/src/features/chat/app-server/inbound/controller.ts +++ b/src/features/chat/app-server/inbound/controller.ts @@ -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); diff --git a/src/features/chat/application/connection/connection-controller.ts b/src/features/chat/application/connection/connection-controller.ts index 965a9d96..e0ce6722 100644 --- a/src/features/chat/application/connection/connection-controller.ts +++ b/src/features/chat/application/connection/connection-controller.ts @@ -89,7 +89,7 @@ export class ChatConnectionController { handleChatConnectionExit(this.host); } - async refreshThreads(): Promise { + async fetchActiveThreads(): Promise { 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 { diff --git a/src/features/chat/application/ports/chat-host.ts b/src/features/chat/application/ports/chat-host.ts index a41c3e45..ebe2159d 100644 --- a/src/features/chat/application/ports/chat-host.ts +++ b/src/features/chat/application/ports/chat-host.ts @@ -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" >; diff --git a/src/features/chat/application/threads/composition.ts b/src/features/chat/application/threads/composition.ts index 97502717..6cd23e1b 100644 --- a/src/features/chat/application/threads/composition.ts +++ b/src/features/chat/application/threads/composition.ts @@ -29,7 +29,7 @@ interface ThreadPartsContext { getClosing: () => boolean; }; thread: { - refreshThreads: () => Promise; + fetchActiveThreads: () => Promise; 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(); }, }; diff --git a/src/features/chat/host/connection-bundle.ts b/src/features/chat/host/connection-bundle.ts index b12f481f..84be2f34 100644 --- a/src/features/chat/host/connection-bundle.ts +++ b/src/features/chat/host/connection-bundle.ts @@ -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 => { - 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 }); } }); }, diff --git a/src/features/chat/host/runtime.ts b/src/features/chat/host/runtime.ts index 7fd7163b..5f2503ec 100644 --- a/src/features/chat/host/runtime.ts +++ b/src/features/chat/host/runtime.ts @@ -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(); }, diff --git a/src/features/chat/host/session.ts b/src/features/chat/host/session.ts index 5a2970b6..929860c1 100644 --- a/src/features/chat/host/session.ts +++ b/src/features/chat/host/session.ts @@ -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 { - 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); } diff --git a/src/features/chat/host/surface-handle.ts b/src/features/chat/host/surface-handle.ts index da8cf988..06e74a7c 100644 --- a/src/features/chat/host/surface-handle.ts +++ b/src/features/chat/host/surface-handle.ts @@ -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; - applyThreadListSnapshot(threads: readonly Thread[]): void; - applyAppServerMetadataSnapshot(metadata: SharedServerMetadata): void; - applyAvailableModelsSnapshot(models: readonly ModelMetadata[]): void; openPanelSnapshot(): OpenCodexPanelSnapshot; openThread(threadId: string): Promise; focusThread(threadId?: string | null): Promise; diff --git a/src/features/thread-picker/modal.ts b/src/features/thread-picker/modal.ts index 34dd1e8c..7a404e0d 100644 --- a/src/features/thread-picker/modal.ts +++ b/src/features/thread-picker/modal.ts @@ -17,7 +17,7 @@ export interface ThreadPickerHost { openThreadInAvailableView(threadId: string): Promise; } -type ThreadPickerCatalog = Pick; +type ThreadPickerCatalog = Pick; interface ThreadSuggestion { thread: Thread; @@ -78,14 +78,14 @@ function threadOpenModeFromEvent(evt: MouseEvent | KeyboardEvent): ThreadOpenMod } async function loadThreadPickerThreads(host: ThreadPickerHost): Promise { - 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); }), { diff --git a/src/features/threads-view/session.ts b/src/features/threads-view/session.ts index bf486ceb..0a1c3452 100644 --- a/src/features/threads-view/session.ts +++ b/src/features/threads-view/session.ts @@ -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(); + 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(); diff --git a/src/features/threads-view/view.ts b/src/features/threads-view/view.ts index 2b88af83..bc4ff0a7 100644 --- a/src/features/threads-view/view.ts +++ b/src/features/threads-view/view.ts @@ -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); - } } diff --git a/src/plugin-runtime.ts b/src/plugin-runtime.ts index df1439fd..48316504 100644 --- a/src/plugin-runtime.ts +++ b/src/plugin-runtime.ts @@ -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[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, diff --git a/src/settings/dynamic-data-controller.ts b/src/settings/dynamic-data-controller.ts index 50e06d81..b0414025 100644 --- a/src/settings/dynamic-data-controller.ts +++ b/src/settings/dynamic-data-controller.ts @@ -24,7 +24,10 @@ export interface SettingsDynamicDataHost { threadCatalog: SettingsThreadCatalog; } -type SettingsThreadCatalog = Pick; +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 { 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, diff --git a/src/settings/tab.ts b/src/settings/tab.ts index b03f19eb..d3cd6b1a 100644 --- a/src/settings/tab.ts +++ b/src/settings/tab.ts @@ -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 }); } }); diff --git a/src/workspace/shared-thread-catalog.ts b/src/workspace/shared-thread-catalog.ts index 1c53322b..39dd8ca8 100644 --- a/src/workspace/shared-thread-catalog.ts +++ b/src/workspace/shared-thread-catalog.ts @@ -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): Promise { - return this.options.cache.refreshThreadList(this.context(), fetchThreads, (threads) => { - this.options.surfaces.applyThreadListSnapshot(threads); - }); + async fetchActiveThreads(fetchThreads: () => Promise): Promise { + 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( + 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?.(); + }; + } } diff --git a/src/workspace/thread-surface-actions.ts b/src/workspace/thread-surface-actions.ts index d1b55028..c105e2c9 100644 --- a/src/workspace/thread-surface-actions.ts +++ b/src/workspace/thread-surface-actions.ts @@ -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(); diff --git a/tests/app-server/shared-cache.test.ts b/tests/app-server/query-cache.test.ts similarity index 55% rename from tests/app-server/shared-cache.test.ts rename to tests/app-server/query-cache.test.ts index 827672c0..24b91666 100644 --- a/tests/app-server/shared-cache.test.ts +++ b/tests/app-server/query-cache.test.ts @@ -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[]>(); 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[]>(); const newRefresh = deferred[]>(); - 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[]>(); - 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 { +function cacheContext(overrides: Partial = {}): AppServerQueryContext { return { codexPath: "codex", vaultPath: "/vault", diff --git a/tests/app-server/shared-cache-state.test.ts b/tests/app-server/shared-cache-state.test.ts deleted file mode 100644 index 60d5c0ef..00000000 --- a/tests/app-server/shared-cache-state.test.ts +++ /dev/null @@ -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 { - return { - codexPath: "codex", - vaultPath: "/vault", - ...overrides, - }; -} - -function expectPresent(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; -} diff --git a/tests/features/chat/connection/server-actions/server-actions.test.ts b/tests/features/chat/connection/server-actions/server-actions.test.ts index 16a91d8a..54fa5f38 100644 --- a/tests/features/chat/connection/server-actions/server-actions.test.ts +++ b/tests/features/chat/connection/server-actions/server-actions.test.ts @@ -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(), }); diff --git a/tests/features/chat/protocol/inbound/controller.test.ts b/tests/features/chat/protocol/inbound/controller.test.ts index 71110d74..42d79e35 100644 --- a/tests/features/chat/protocol/inbound/controller.test.ts +++ b/tests/features/chat/protocol/inbound/controller.test.ts @@ -22,10 +22,10 @@ function controllerForState( actions: Partial[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); expect(recordMcpStartupStatus).toHaveBeenCalledWith("github", "failed", "missing token"); - expect(publishAppServerMetadata).toHaveBeenCalledOnce(); + expect(setAppServerMetadata).toHaveBeenCalledOnce(); expect(chatStateMessageStreamItems(state)).toEqual([]); }); }); diff --git a/tests/features/chat/view-connection.test.ts b/tests/features/chat/view-connection.test.ts index 8f5e5468..03701a03 100644 --- a/tests/features/chat/view-connection.test.ts +++ b/tests/features/chat/view-connection.test.ts @@ -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) => fetchThreads() as Promise); + const fetchActiveThreads = vi.fn((fetchThreads: () => Promise) => fetchThreads() as Promise); 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) => fetchThreads() as Promise, - ) 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); + }; + }, }, }; } diff --git a/tests/features/thread-picker/modal.test.ts b/tests/features/thread-picker/modal.test.ts index 20022d50..1ab877e2 100644 --- a/tests/features/thread-picker/modal.test.ts +++ b/tests/features/thread-picker/modal.test.ts @@ -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); diff --git a/tests/features/threads-view/view.test.ts b/tests/features/threads-view/view.test.ts index 84f7bb5d..744322ae 100644 --- a/tests/features/threads-view/view.test.ts +++ b/tests/features/threads-view/view.test.ts @@ -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) => fetchThreads() as Promise); + const fetchActiveThreads = vi.fn((fetchThreads: () => Promise) => fetchThreads() as Promise); 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 = {}) { archiveThreadInCatalog: vi.fn(), renameThreadInCatalog: vi.fn(), refreshFromOpenSurface: vi.fn(), - refreshThreads: vi.fn((fetchThreads: () => Promise) => fetchThreads() as Promise), - cachedThreads: vi.fn(() => null), + fetchActiveThreads: vi.fn((fetchThreads: () => Promise) => fetchThreads() as Promise), + activeThreadsSnapshot: vi.fn(() => null), + observeActiveThreads: vi.fn(() => () => undefined), ...threadCatalogOverrides, }, ...hostOverrides, diff --git a/tests/main.test.ts b/tests/main.test.ts index c17dce1c..741a016c 100644 --- a/tests/main.test.ts +++ b/tests/main.test.ts @@ -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) => fetchThreads() as Promise), - cachedThreads: vi.fn(() => null), - cachedAppServerMetadata: vi.fn(() => null), + setActiveThreads: vi.fn(), + setAppServerMetadata: vi.fn(), + fetchActiveThreads: vi.fn((fetchThreads: () => Promise) => fetchThreads() as Promise), + activeThreadsSnapshot: vi.fn(() => null), + appServerMetadataSnapshot: vi.fn(() => null), + modelsSnapshot: vi.fn(() => null), + observeActiveThreads: vi.fn(() => () => undefined), + observeAppServerMetadata: vi.fn(() => () => undefined), + observeModels: vi.fn(() => () => undefined), }, }; } diff --git a/tests/mocks/obsidian.ts b/tests/mocks/obsidian.ts index 98397581..4fc39b39 100644 --- a/tests/mocks/obsidian.ts +++ b/tests/mocks/obsidian.ts @@ -253,6 +253,10 @@ export class PluginSettingTab { display(): void { // Test mock placeholder. } + + hide(): void { + this.containerEl.empty(); + } } export class Setting { diff --git a/tests/settings/settings-tab.test.ts b/tests/settings/settings-tab.test.ts index 52656878..73148122 100644 --- a/tests/settings/settings-tab.test.ts +++ b/tests/settings/settings-tab.test.ts @@ -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) => operation(oldClient)) .mockImplementationOnce((_codexPath: string, _cwd: string, operation: (client: unknown) => Promise) => 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) => 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) => 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) => 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; 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; 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(), }, }; } diff --git a/tests/workspace/shared-thread-catalog.test.ts b/tests/workspace/shared-thread-catalog.test.ts index 530feb72..db292126 100644 --- a/tests/workspace/shared-thread-catalog.test.ts +++ b/tests/workspace/shared-thread-catalog.test.ts @@ -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((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(), }; }