diff --git a/packages/server/src/server/bootstrap.ts b/packages/server/src/server/bootstrap.ts index 37403e9bd..35fa32f4b 100644 --- a/packages/server/src/server/bootstrap.ts +++ b/packages/server/src/server/bootstrap.ts @@ -100,6 +100,7 @@ import { createAllClients, shutdownProviders } from "./agent/provider-registry.j import { bootstrapWorkspaceRegistries } from "./workspace-registry-bootstrap.js"; import { FileBackedProjectRegistry, FileBackedWorkspaceRegistry } from "./workspace-registry.js"; import { FileBackedChatService } from "./chat/chat-service.js"; +import { CheckoutDiffManager } from "./checkout-diff-manager.js"; import { LoopService } from "./loop-service.js"; import { ScheduleService } from "./schedule/service.js"; import { createTerminalManager, type TerminalManager } from "../terminal/terminal-manager.js"; @@ -410,6 +411,10 @@ export async function createPaseoDaemon( logger.info({ elapsed: elapsed() }, "Workspace registries bootstrapped"); await chatService.initialize(); logger.info({ elapsed: elapsed() }, "Chat service initialized"); + const checkoutDiffManager = new CheckoutDiffManager({ + logger, + paseoHome: config.paseoHome, + }); const loopService = new LoopService({ paseoHome: config.paseoHome, logger, @@ -662,6 +667,7 @@ export async function createPaseoDaemon( chatService, loopService, scheduleService, + checkoutDiffManager, ); unsubscribeSpeechReadiness = subscribeSpeechReadiness((snapshot) => { wsServer?.publishSpeechReadiness(snapshot); diff --git a/packages/server/src/server/checkout-diff-manager.ts b/packages/server/src/server/checkout-diff-manager.ts new file mode 100644 index 000000000..db7a9096b --- /dev/null +++ b/packages/server/src/server/checkout-diff-manager.ts @@ -0,0 +1,361 @@ +import { watch, type FSWatcher } from "node:fs"; +import { exec } from "child_process"; +import { promisify } from "util"; +import type pino from "pino"; +import type { SubscribeCheckoutDiffRequest, SessionOutboundMessage } from "./messages.js"; +import { getCheckoutDiff } from "../utils/checkout-git.js"; +import { expandTilde } from "../utils/path.js"; +import { READ_ONLY_GIT_ENV, resolveCheckoutGitDir, toCheckoutError } from "./checkout-git-utils.js"; + +const execAsync = promisify(exec); + +const CHECKOUT_DIFF_WATCH_DEBOUNCE_MS = 150; +const CHECKOUT_DIFF_FALLBACK_REFRESH_MS = 5_000; + +export type CheckoutDiffCompareInput = SubscribeCheckoutDiffRequest["compare"]; + +export type CheckoutDiffSnapshotPayload = Omit< + Extract["payload"], + "subscriptionId" +>; + +export type CheckoutDiffMetrics = { + checkoutDiffTargetCount: number; + checkoutDiffSubscriptionCount: number; + checkoutDiffWatcherCount: number; + checkoutDiffFallbackRefreshTargetCount: number; +}; + +type CheckoutDiffWatchTarget = { + key: string; + cwd: string; + diffCwd: string; + compare: CheckoutDiffCompareInput; + listeners: Set<(snapshot: CheckoutDiffSnapshotPayload) => void>; + watchers: FSWatcher[]; + fallbackRefreshInterval: NodeJS.Timeout | null; + debounceTimer: NodeJS.Timeout | null; + refreshPromise: Promise | null; + refreshQueued: boolean; + latestPayload: CheckoutDiffSnapshotPayload | null; + latestFingerprint: string | null; +}; + +export class CheckoutDiffManager { + private readonly logger: pino.Logger; + private readonly paseoHome: string; + private readonly targets = new Map(); + + constructor(options: { logger: pino.Logger; paseoHome: string }) { + this.logger = options.logger.child({ module: "checkout-diff-manager" }); + this.paseoHome = options.paseoHome; + } + + async subscribe( + params: { + cwd: string; + compare: CheckoutDiffCompareInput; + }, + listener: (snapshot: CheckoutDiffSnapshotPayload) => void, + ): Promise<{ initial: CheckoutDiffSnapshotPayload; unsubscribe: () => void }> { + const cwd = params.cwd; + const compare = this.normalizeCompare(params.compare); + const target = await this.ensureTarget(cwd, compare); + target.listeners.add(listener); + + const initial = + target.latestPayload ?? + (await this.computeCheckoutDiffSnapshot(target.cwd, target.compare, { + diffCwd: target.diffCwd, + })); + target.latestPayload = initial; + target.latestFingerprint = JSON.stringify(initial); + return { + initial, + unsubscribe: () => { + this.removeListener(target.key, listener); + }, + }; + } + + scheduleRefreshForCwd(cwd: string): void { + const resolvedCwd = expandTilde(cwd); + for (const target of this.targets.values()) { + if (target.cwd !== resolvedCwd && target.diffCwd !== resolvedCwd) { + continue; + } + this.scheduleTargetRefresh(target); + } + } + + getMetrics(): CheckoutDiffMetrics { + let checkoutDiffSubscriptionCount = 0; + let checkoutDiffWatcherCount = 0; + let checkoutDiffFallbackRefreshTargetCount = 0; + + for (const target of this.targets.values()) { + checkoutDiffSubscriptionCount += target.listeners.size; + checkoutDiffWatcherCount += target.watchers.length; + if (target.fallbackRefreshInterval) { + checkoutDiffFallbackRefreshTargetCount += 1; + } + } + + return { + checkoutDiffTargetCount: this.targets.size, + checkoutDiffSubscriptionCount, + checkoutDiffWatcherCount, + checkoutDiffFallbackRefreshTargetCount, + }; + } + + dispose(): void { + for (const target of this.targets.values()) { + this.closeTarget(target); + } + this.targets.clear(); + } + + private normalizeCompare(compare: CheckoutDiffCompareInput): CheckoutDiffCompareInput { + if (compare.mode === "uncommitted") { + return { mode: "uncommitted" }; + } + const trimmedBaseRef = compare.baseRef?.trim(); + return trimmedBaseRef ? { mode: "base", baseRef: trimmedBaseRef } : { mode: "base" }; + } + + private buildTargetKey(cwd: string, compare: CheckoutDiffCompareInput): string { + return JSON.stringify([ + cwd, + compare.mode, + compare.mode === "base" ? (compare.baseRef ?? "") : "", + ]); + } + + private closeTarget(target: CheckoutDiffWatchTarget): void { + if (target.debounceTimer) { + clearTimeout(target.debounceTimer); + target.debounceTimer = null; + } + if (target.fallbackRefreshInterval) { + clearInterval(target.fallbackRefreshInterval); + target.fallbackRefreshInterval = null; + } + for (const watcher of target.watchers) { + watcher.close(); + } + target.watchers = []; + target.listeners.clear(); + } + + private removeListener( + targetKey: string, + listener: (snapshot: CheckoutDiffSnapshotPayload) => void, + ): void { + const target = this.targets.get(targetKey); + if (!target) { + return; + } + target.listeners.delete(listener); + if (target.listeners.size > 0) { + return; + } + this.closeTarget(target); + this.targets.delete(targetKey); + } + + private async resolveCheckoutWatchRoot(cwd: string): Promise { + try { + const { stdout } = await execAsync("git rev-parse --path-format=absolute --show-toplevel", { + cwd, + env: READ_ONLY_GIT_ENV, + }); + const root = stdout.trim(); + return root.length > 0 ? root : null; + } catch { + return null; + } + } + + private scheduleTargetRefresh(target: CheckoutDiffWatchTarget): void { + if (target.debounceTimer) { + clearTimeout(target.debounceTimer); + } + target.debounceTimer = setTimeout(() => { + target.debounceTimer = null; + void this.refreshTarget(target); + }, CHECKOUT_DIFF_WATCH_DEBOUNCE_MS); + } + + private async computeCheckoutDiffSnapshot( + cwd: string, + compare: CheckoutDiffCompareInput, + options?: { diffCwd?: string }, + ): Promise { + const diffCwd = options?.diffCwd ?? cwd; + try { + const diffResult = await getCheckoutDiff( + diffCwd, + { + mode: compare.mode, + baseRef: compare.baseRef, + includeStructured: true, + }, + { paseoHome: this.paseoHome }, + ); + const files = [...(diffResult.structured ?? [])]; + files.sort((a, b) => { + if (a.path === b.path) return 0; + return a.path < b.path ? -1 : 1; + }); + return { + cwd, + files, + error: null, + }; + } catch (error) { + return { + cwd, + files: [], + error: toCheckoutError(error), + }; + } + } + + private async refreshTarget(target: CheckoutDiffWatchTarget): Promise { + if (target.refreshPromise) { + target.refreshQueued = true; + return; + } + + target.refreshPromise = (async () => { + do { + target.refreshQueued = false; + const snapshot = await this.computeCheckoutDiffSnapshot(target.cwd, target.compare, { + diffCwd: target.diffCwd, + }); + target.latestPayload = snapshot; + const fingerprint = JSON.stringify(snapshot); + if (fingerprint !== target.latestFingerprint) { + target.latestFingerprint = fingerprint; + for (const listener of target.listeners) { + listener(snapshot); + } + } + } while (target.refreshQueued); + })(); + + try { + await target.refreshPromise; + } finally { + target.refreshPromise = null; + } + } + + private async ensureTarget( + cwd: string, + compare: CheckoutDiffCompareInput, + ): Promise { + const targetKey = this.buildTargetKey(cwd, compare); + const existing = this.targets.get(targetKey); + if (existing) { + return existing; + } + + const watchRoot = await this.resolveCheckoutWatchRoot(cwd); + const target: CheckoutDiffWatchTarget = { + key: targetKey, + cwd, + diffCwd: watchRoot ?? cwd, + compare, + listeners: new Set(), + watchers: [], + fallbackRefreshInterval: null, + debounceTimer: null, + refreshPromise: null, + refreshQueued: false, + latestPayload: null, + latestFingerprint: null, + }; + + const repoWatchPath = watchRoot ?? cwd; + const watchPaths = new Set([repoWatchPath]); + const gitDir = await resolveCheckoutGitDir(cwd); + if (gitDir) { + watchPaths.add(gitDir); + } + + let hasRecursiveRepoCoverage = false; + const allowRecursiveRepoWatch = process.platform !== "linux"; + for (const watchPath of watchPaths) { + const shouldTryRecursive = watchPath === repoWatchPath && allowRecursiveRepoWatch; + const createWatcher = (recursive: boolean): FSWatcher => + watch(watchPath, { recursive }, () => { + this.scheduleTargetRefresh(target); + }); + + let watcher: FSWatcher | null = null; + let watcherIsRecursive = false; + try { + if (shouldTryRecursive) { + watcher = createWatcher(true); + watcherIsRecursive = true; + } else { + watcher = createWatcher(false); + } + } catch (error) { + if (shouldTryRecursive) { + try { + watcher = createWatcher(false); + this.logger.warn( + { err: error, watchPath, cwd, compare }, + "Checkout diff recursive watch unavailable; using non-recursive fallback", + ); + } catch (fallbackError) { + this.logger.warn( + { err: fallbackError, watchPath, cwd, compare }, + "Failed to start checkout diff watcher", + ); + } + } else { + this.logger.warn( + { err: error, watchPath, cwd, compare }, + "Failed to start checkout diff watcher", + ); + } + } + + if (!watcher) { + continue; + } + + watcher.on("error", (error) => { + this.logger.warn({ err: error, watchPath, cwd, compare }, "Checkout diff watcher error"); + }); + target.watchers.push(watcher); + if (watchPath === repoWatchPath && watcherIsRecursive) { + hasRecursiveRepoCoverage = true; + } + } + + const missingRepoCoverage = !hasRecursiveRepoCoverage; + if (target.watchers.length === 0 || missingRepoCoverage) { + target.fallbackRefreshInterval = setInterval(() => { + this.scheduleTargetRefresh(target); + }, CHECKOUT_DIFF_FALLBACK_REFRESH_MS); + this.logger.warn( + { + cwd, + compare, + intervalMs: CHECKOUT_DIFF_FALLBACK_REFRESH_MS, + reason: + target.watchers.length === 0 ? "no_watchers" : "missing_recursive_repo_root_coverage", + }, + "Checkout diff watchers unavailable; using timed refresh fallback", + ); + } + + this.targets.set(targetKey, target); + return target; + } +} diff --git a/packages/server/src/server/checkout-git-utils.ts b/packages/server/src/server/checkout-git-utils.ts new file mode 100644 index 000000000..a27460e1b --- /dev/null +++ b/packages/server/src/server/checkout-git-utils.ts @@ -0,0 +1,50 @@ +import { exec } from "child_process"; +import { promisify } from "util"; +import { + MergeConflictError, + MergeFromBaseConflictError, + NotGitRepoError, +} from "../utils/checkout-git.js"; + +const execAsync = promisify(exec); + +export const READ_ONLY_GIT_ENV: NodeJS.ProcessEnv = { + ...process.env, + GIT_OPTIONAL_LOCKS: "0", +}; + +export type CheckoutErrorCode = "NOT_GIT_REPO" | "NOT_ALLOWED" | "MERGE_CONFLICT" | "UNKNOWN"; + +export type CheckoutErrorPayload = { + code: CheckoutErrorCode; + message: string; +}; + +export async function resolveCheckoutGitDir(cwd: string): Promise { + try { + const { stdout } = await execAsync("git rev-parse --absolute-git-dir", { + cwd, + env: READ_ONLY_GIT_ENV, + }); + const gitDir = stdout.trim(); + return gitDir.length > 0 ? gitDir : null; + } catch { + return null; + } +} + +export function toCheckoutError(error: unknown): CheckoutErrorPayload { + if (error instanceof NotGitRepoError) { + return { code: "NOT_GIT_REPO", message: error.message }; + } + if (error instanceof MergeConflictError) { + return { code: "MERGE_CONFLICT", message: error.message }; + } + if (error instanceof MergeFromBaseConflictError) { + return { code: "MERGE_CONFLICT", message: error.message }; + } + if (error instanceof Error) { + return { code: "UNKNOWN", message: error.message }; + } + return { code: "UNKNOWN", message: String(error) }; +} diff --git a/packages/server/src/server/session.ts b/packages/server/src/server/session.ts index 63da9ac78..2d1d32146 100644 --- a/packages/server/src/server/session.ts +++ b/packages/server/src/server/session.ts @@ -153,9 +153,6 @@ import { getCheckoutStatus, getCheckoutStatusLite, listBranchSuggestions, - NotGitRepoError, - MergeConflictError, - MergeFromBaseConflictError, commitChanges, mergeToBase, mergeFromBase, @@ -167,6 +164,12 @@ import { import { getProjectIcon } from "../utils/project-icon.js"; import { expandTilde } from "../utils/path.js"; import { searchHomeDirectories, searchWorkspaceEntries } from "../utils/directory-suggestions.js"; +import { + READ_ONLY_GIT_ENV, + resolveCheckoutGitDir, + toCheckoutError, +} from "./checkout-git-utils.js"; +import { CheckoutDiffManager } from "./checkout-diff-manager.js"; import { ensureLocalSpeechModels, getLocalSpeechModelDir, @@ -187,14 +190,8 @@ import { ScheduleService } from "./schedule/service.js"; const execAsync = promisify(exec); const MAX_INITIAL_AGENT_TITLE_CHARS = Math.min(60, MAX_EXPLICIT_AGENT_TITLE_CHARS); -const READ_ONLY_GIT_ENV: NodeJS.ProcessEnv = { - ...process.env, - GIT_OPTIONAL_LOCKS: "0", -}; const pendingAgentInitializations = new Map>(); const DEFAULT_AGENT_PROVIDER = AGENT_PROVIDER_IDS[0]; -const CHECKOUT_DIFF_WATCH_DEBOUNCE_MS = 150; -const CHECKOUT_DIFF_FALLBACK_REFRESH_MS = 5_000; const WORKSPACE_GIT_WATCH_DEBOUNCE_MS = 500; const WORKSPACE_GIT_WATCH_REMOVED_FINGERPRINT = "__removed__"; const TERMINAL_STREAM_HIGH_WATER_BYTES = 256 * 1024; @@ -237,28 +234,6 @@ export function resolveCreateAgentTitles(options: { type ProcessingPhase = "idle" | "transcribing"; -type CheckoutDiffCompareInput = SubscribeCheckoutDiffRequest["compare"]; - -type CheckoutDiffSnapshotPayload = Omit< - Extract["payload"], - "subscriptionId" ->; - -type CheckoutDiffWatchTarget = { - key: string; - cwd: string; - diffCwd: string; - compare: CheckoutDiffCompareInput; - subscriptions: Set; - watchers: FSWatcher[]; - fallbackRefreshInterval: NodeJS.Timeout | null; - debounceTimer: NodeJS.Timeout | null; - refreshPromise: Promise | null; - refreshQueued: boolean; - latestPayload: CheckoutDiffSnapshotPayload | null; - latestFingerprint: string | null; -}; - type WorkspaceGitWatchTarget = { cwd: string; watchers: FSWatcher[]; @@ -276,13 +251,6 @@ type NormalizedGitOptions = { worktreeSlug?: string; }; -type CheckoutErrorCode = "NOT_GIT_REPO" | "NOT_ALLOWED" | "MERGE_CONFLICT" | "UNKNOWN"; - -type CheckoutErrorPayload = { - code: CheckoutErrorCode; - message: string; -}; - type ActiveTerminalStream = { terminalId: string; slot: number; @@ -292,10 +260,6 @@ type ActiveTerminalStream = { }; export type SessionRuntimeMetrics = { - checkoutDiffTargetCount: number; - checkoutDiffSubscriptionCount: number; - checkoutDiffWatcherCount: number; - checkoutDiffFallbackRefreshTargetCount: number; terminalDirectorySubscriptionCount: number; terminalSubscriptionCount: number; inflightRequests: number; @@ -417,6 +381,7 @@ export type SessionOptions = { chatService: FileBackedChatService; scheduleService: ScheduleService; loopService: LoopService; + checkoutDiffManager: CheckoutDiffManager; createAgentMcpTransport: AgentMcpTransportFactory; stt: Resolvable; tts: Resolvable; @@ -603,6 +568,7 @@ export class Session { private readonly chatService: FileBackedChatService; private readonly scheduleService: ScheduleService; private readonly loopService: LoopService; + private readonly checkoutDiffManager: CheckoutDiffManager; private readonly createAgentMcpTransport: AgentMcpTransportFactory; private readonly downloadTokenStore: DownloadTokenStore; private readonly pushTokenStore: PushTokenStore; @@ -627,8 +593,7 @@ export class Session { private nextTerminalSlot = 0; private inflightRequests = 0; private peakInflightRequests = 0; - private readonly checkoutDiffSubscriptions = new Map(); - private readonly checkoutDiffTargets = new Map(); + private readonly checkoutDiffSubscriptions = new Map void>(); private readonly workspaceGitWatchTargets = new Map(); private readonly voiceAgentMcpStdio: VoiceMcpStdioConfig | null; private readonly localSpeechModelsDir: string; @@ -668,6 +633,7 @@ export class Session { chatService, scheduleService, loopService, + checkoutDiffManager, createAgentMcpTransport, stt, tts, @@ -693,6 +659,7 @@ export class Session { this.chatService = chatService; this.scheduleService = scheduleService; this.loopService = loopService; + this.checkoutDiffManager = checkoutDiffManager; this.createAgentMcpTransport = createAgentMcpTransport; this.terminalManager = terminalManager; if (this.terminalManager) { @@ -761,20 +728,7 @@ export class Session { } public getRuntimeMetrics(): SessionRuntimeMetrics { - let checkoutDiffWatcherCount = 0; - let checkoutDiffFallbackRefreshTargetCount = 0; - for (const target of this.checkoutDiffTargets.values()) { - checkoutDiffWatcherCount += target.watchers.length; - if (target.fallbackRefreshInterval) { - checkoutDiffFallbackRefreshTargetCount += 1; - } - } - return { - checkoutDiffTargetCount: this.checkoutDiffTargets.size, - checkoutDiffSubscriptionCount: this.checkoutDiffSubscriptions.size, - checkoutDiffWatcherCount, - checkoutDiffFallbackRefreshTargetCount, terminalDirectorySubscriptionCount: this.subscribedTerminalDirectories.size, terminalSubscriptionCount: this.activeTerminalStreams.size, inflightRequests: this.inflightRequests, @@ -3283,22 +3237,6 @@ export class Session { return baseBranch; } - private toCheckoutError(error: unknown): CheckoutErrorPayload { - if (error instanceof NotGitRepoError) { - return { code: "NOT_GIT_REPO", message: error.message }; - } - if (error instanceof MergeConflictError) { - return { code: "MERGE_CONFLICT", message: error.message }; - } - if (error instanceof MergeFromBaseConflictError) { - return { code: "MERGE_CONFLICT", message: error.message }; - } - if (error instanceof Error) { - return { code: "UNKNOWN", message: error.message }; - } - return { code: "UNKNOWN", message: String(error) }; - } - private isPathWithinRoot(rootPath: string, candidatePath: string): boolean { const resolvedRoot = resolve(rootPath); const resolvedCandidate = resolve(candidatePath); @@ -3900,7 +3838,7 @@ export class Session { hasRemote: false, remoteUrl: null, isPaseoOwnedWorktree: false, - error: this.toCheckoutError(error), + error: toCheckoutError(error), requestId, }, }); @@ -4055,70 +3993,6 @@ export class Session { } } - private normalizeCheckoutDiffCompare( - compare: CheckoutDiffCompareInput, - ): CheckoutDiffCompareInput { - if (compare.mode === "uncommitted") { - return { mode: "uncommitted" }; - } - const trimmedBaseRef = compare.baseRef?.trim(); - return trimmedBaseRef ? { mode: "base", baseRef: trimmedBaseRef } : { mode: "base" }; - } - - private buildCheckoutDiffTargetKey(cwd: string, compare: CheckoutDiffCompareInput): string { - return JSON.stringify([ - cwd, - compare.mode, - compare.mode === "base" ? (compare.baseRef ?? "") : "", - ]); - } - - private closeCheckoutDiffWatchTarget(target: CheckoutDiffWatchTarget): void { - if (target.debounceTimer) { - clearTimeout(target.debounceTimer); - target.debounceTimer = null; - } - if (target.fallbackRefreshInterval) { - clearInterval(target.fallbackRefreshInterval); - target.fallbackRefreshInterval = null; - } - for (const watcher of target.watchers) { - watcher.close(); - } - target.watchers = []; - } - - private removeCheckoutDiffSubscription(subscriptionId: string): void { - const subscription = this.checkoutDiffSubscriptions.get(subscriptionId); - if (!subscription) { - return; - } - this.checkoutDiffSubscriptions.delete(subscriptionId); - - const target = this.checkoutDiffTargets.get(subscription.targetKey); - if (!target) { - return; - } - target.subscriptions.delete(subscriptionId); - if (target.subscriptions.size === 0) { - this.closeCheckoutDiffWatchTarget(target); - this.checkoutDiffTargets.delete(subscription.targetKey); - } - } - - private async resolveCheckoutGitDir(cwd: string): Promise { - try { - const { stdout } = await execAsync("git rev-parse --absolute-git-dir", { - cwd, - env: READ_ONLY_GIT_ENV, - }); - const gitDir = stdout.trim(); - return gitDir.length > 0 ? gitDir : null; - } catch { - return null; - } - } - private async resolveWorkspaceGitRefsRoot(gitDir: string): Promise { try { const commonDir = (await readFile(join(gitDir, "commondir"), "utf8")).trim(); @@ -4235,7 +4109,7 @@ export class Session { return; } - const gitDir = await this.resolveCheckoutGitDir(cwd); + const gitDir = await resolveCheckoutGitDir(cwd); if (!gitDir) { return; } @@ -4292,267 +4166,39 @@ export class Session { await this.ensureWorkspaceGitWatchTarget(cwd); } - private async resolveCheckoutWatchRoot(cwd: string): Promise { - try { - const { stdout } = await execAsync("git rev-parse --path-format=absolute --show-toplevel", { - cwd, - env: READ_ONLY_GIT_ENV, - }); - const root = stdout.trim(); - return root.length > 0 ? root : null; - } catch { - return null; - } - } - - private scheduleCheckoutDiffTargetRefresh(target: CheckoutDiffWatchTarget): void { - if (target.debounceTimer) { - clearTimeout(target.debounceTimer); - } - target.debounceTimer = setTimeout(() => { - target.debounceTimer = null; - void this.refreshCheckoutDiffTarget(target); - }, CHECKOUT_DIFF_WATCH_DEBOUNCE_MS); - } - - private emitCheckoutDiffUpdate( - target: CheckoutDiffWatchTarget, - snapshot: CheckoutDiffSnapshotPayload, - ): void { - if (target.subscriptions.size === 0) { - return; - } - for (const subscriptionId of target.subscriptions) { - this.emit({ - type: "checkout_diff_update", - payload: { - subscriptionId, - ...snapshot, - }, - }); - } - } - - private checkoutDiffSnapshotFingerprint(snapshot: CheckoutDiffSnapshotPayload): string { - return JSON.stringify(snapshot); - } - - private async computeCheckoutDiffSnapshot( - cwd: string, - compare: CheckoutDiffCompareInput, - options?: { diffCwd?: string }, - ): Promise { - const diffCwd = options?.diffCwd ?? cwd; - try { - const diffResult = await getCheckoutDiff( - diffCwd, - { - mode: compare.mode, - baseRef: compare.baseRef, - includeStructured: true, - }, - { paseoHome: this.paseoHome }, - ); - const files = [...(diffResult.structured ?? [])]; - files.sort((a, b) => { - if (a.path === b.path) return 0; - return a.path < b.path ? -1 : 1; - }); - return { - cwd, - files, - error: null, - }; - } catch (error) { - return { - cwd, - files: [], - error: this.toCheckoutError(error), - }; - } - } - - private async refreshCheckoutDiffTarget(target: CheckoutDiffWatchTarget): Promise { - if (target.refreshPromise) { - target.refreshQueued = true; - return; - } - - target.refreshPromise = (async () => { - do { - target.refreshQueued = false; - const snapshot = await this.computeCheckoutDiffSnapshot(target.cwd, target.compare, { - diffCwd: target.diffCwd, - }); - target.latestPayload = snapshot; - const fingerprint = this.checkoutDiffSnapshotFingerprint(snapshot); - if (fingerprint !== target.latestFingerprint) { - target.latestFingerprint = fingerprint; - this.emitCheckoutDiffUpdate(target, snapshot); - } - } while (target.refreshQueued); - })(); - - try { - await target.refreshPromise; - } finally { - target.refreshPromise = null; - } - } - - private async ensureCheckoutDiffWatchTarget( - cwd: string, - compare: CheckoutDiffCompareInput, - ): Promise { - const targetKey = this.buildCheckoutDiffTargetKey(cwd, compare); - const existing = this.checkoutDiffTargets.get(targetKey); - if (existing) { - return existing; - } - - const watchRoot = await this.resolveCheckoutWatchRoot(cwd); - const target: CheckoutDiffWatchTarget = { - key: targetKey, - cwd, - diffCwd: watchRoot ?? cwd, - compare, - subscriptions: new Set(), - watchers: [], - fallbackRefreshInterval: null, - debounceTimer: null, - refreshPromise: null, - refreshQueued: false, - latestPayload: null, - latestFingerprint: null, - }; - - const repoWatchPath = watchRoot ?? cwd; - const watchPaths = new Set([repoWatchPath]); - const gitDir = await this.resolveCheckoutGitDir(cwd); - if (gitDir) { - watchPaths.add(gitDir); - } - - let hasRecursiveRepoCoverage = false; - const allowRecursiveRepoWatch = process.platform !== "linux"; - for (const watchPath of watchPaths) { - const shouldTryRecursive = watchPath === repoWatchPath && allowRecursiveRepoWatch; - const createWatcher = (recursive: boolean): FSWatcher => - watch(watchPath, { recursive }, () => { - this.scheduleCheckoutDiffTargetRefresh(target); - }); - - let watcher: FSWatcher | null = null; - let watcherIsRecursive = false; - try { - if (shouldTryRecursive) { - watcher = createWatcher(true); - watcherIsRecursive = true; - } else { - watcher = createWatcher(false); - } - } catch (error) { - if (shouldTryRecursive) { - try { - watcher = createWatcher(false); - this.sessionLogger.warn( - { err: error, watchPath, cwd, compare }, - "Checkout diff recursive watch unavailable; using non-recursive fallback", - ); - } catch (fallbackError) { - this.sessionLogger.warn( - { err: fallbackError, watchPath, cwd, compare }, - "Failed to start checkout diff watcher", - ); - } - } else { - this.sessionLogger.warn( - { err: error, watchPath, cwd, compare }, - "Failed to start checkout diff watcher", - ); - } - } - - if (!watcher) { - continue; - } - - watcher.on("error", (error) => { - this.sessionLogger.warn( - { err: error, watchPath, cwd, compare }, - "Checkout diff watcher error", - ); - }); - target.watchers.push(watcher); - if (watchPath === repoWatchPath && watcherIsRecursive) { - hasRecursiveRepoCoverage = true; - } - } - - const missingRepoCoverage = !hasRecursiveRepoCoverage; - if (target.watchers.length === 0 || missingRepoCoverage) { - target.fallbackRefreshInterval = setInterval(() => { - this.scheduleCheckoutDiffTargetRefresh(target); - }, CHECKOUT_DIFF_FALLBACK_REFRESH_MS); - this.sessionLogger.warn( - { - cwd, - compare, - intervalMs: CHECKOUT_DIFF_FALLBACK_REFRESH_MS, - reason: - target.watchers.length === 0 ? "no_watchers" : "missing_recursive_repo_root_coverage", - }, - "Checkout diff watchers unavailable; using timed refresh fallback", - ); - } - - this.checkoutDiffTargets.set(targetKey, target); - return target; - } - private async handleSubscribeCheckoutDiffRequest( msg: SubscribeCheckoutDiffRequest, ): Promise { const cwd = expandTilde(msg.cwd); - const compare = this.normalizeCheckoutDiffCompare(msg.compare); - - this.removeCheckoutDiffSubscription(msg.subscriptionId); - const target = await this.ensureCheckoutDiffWatchTarget(cwd, compare); - target.subscriptions.add(msg.subscriptionId); - this.checkoutDiffSubscriptions.set(msg.subscriptionId, { - targetKey: target.key, - }); - - const snapshot = - target.latestPayload ?? - (await this.computeCheckoutDiffSnapshot(cwd, compare, { - diffCwd: target.diffCwd, - })); - target.latestPayload = snapshot; - target.latestFingerprint = this.checkoutDiffSnapshotFingerprint(snapshot); + this.checkoutDiffSubscriptions.get(msg.subscriptionId)?.(); + this.checkoutDiffSubscriptions.delete(msg.subscriptionId); + const subscription = await this.checkoutDiffManager.subscribe( + { cwd, compare: msg.compare }, + (snapshot) => { + this.emit({ + type: "checkout_diff_update", + payload: { + subscriptionId: msg.subscriptionId, + ...snapshot, + }, + }); + }, + ); + this.checkoutDiffSubscriptions.set(msg.subscriptionId, subscription.unsubscribe); this.emit({ type: "subscribe_checkout_diff_response", payload: { subscriptionId: msg.subscriptionId, - ...snapshot, + ...subscription.initial, requestId: msg.requestId, }, }); } private handleUnsubscribeCheckoutDiffRequest(msg: UnsubscribeCheckoutDiffRequest): void { - this.removeCheckoutDiffSubscription(msg.subscriptionId); - } - - private scheduleCheckoutDiffRefreshForCwd(cwd: string): void { - const resolvedCwd = expandTilde(cwd); - for (const target of this.checkoutDiffTargets.values()) { - if (target.cwd !== resolvedCwd && target.diffCwd !== resolvedCwd) { - continue; - } - this.scheduleCheckoutDiffTargetRefresh(target); - } + this.checkoutDiffSubscriptions.get(msg.subscriptionId)?.(); + this.checkoutDiffSubscriptions.delete(msg.subscriptionId); } private async handleCheckoutCommitRequest( @@ -4573,7 +4219,7 @@ export class Session { message, addAll: msg.addAll ?? true, }); - this.scheduleCheckoutDiffRefreshForCwd(cwd); + this.checkoutDiffManager.scheduleRefreshForCwd(cwd); this.emit({ type: "checkout_commit_response", @@ -4590,7 +4236,7 @@ export class Session { payload: { cwd, success: false, - error: this.toCheckoutError(error), + error: toCheckoutError(error), requestId, }, }); @@ -4647,7 +4293,7 @@ export class Session { }, { paseoHome: this.paseoHome }, ); - this.scheduleCheckoutDiffRefreshForCwd(cwd); + this.checkoutDiffManager.scheduleRefreshForCwd(cwd); this.emit({ type: "checkout_merge_response", @@ -4664,7 +4310,7 @@ export class Session { payload: { cwd, success: false, - error: this.toCheckoutError(error), + error: toCheckoutError(error), requestId, }, }); @@ -4691,7 +4337,7 @@ export class Session { baseRef: msg.baseRef, requireCleanTarget: msg.requireCleanTarget ?? true, }); - this.scheduleCheckoutDiffRefreshForCwd(cwd); + this.checkoutDiffManager.scheduleRefreshForCwd(cwd); this.emit({ type: "checkout_merge_from_base_response", @@ -4708,7 +4354,7 @@ export class Session { payload: { cwd, success: false, - error: this.toCheckoutError(error), + error: toCheckoutError(error), requestId, }, }); @@ -4737,7 +4383,7 @@ export class Session { payload: { cwd, success: false, - error: this.toCheckoutError(error), + error: toCheckoutError(error), requestId, }, }); @@ -4782,7 +4428,7 @@ export class Session { cwd, url: null, number: null, - error: this.toCheckoutError(error), + error: toCheckoutError(error), requestId, }, }); @@ -4813,7 +4459,7 @@ export class Session { cwd, status: null, githubFeaturesEnabled: true, - error: this.toCheckoutError(error), + error: toCheckoutError(error), requestId, }, }); @@ -4857,7 +4503,7 @@ export class Session { type: "paseo_worktree_list_response", payload: { worktrees: [], - error: this.toCheckoutError(error), + error: toCheckoutError(error), requestId, }, }); @@ -5005,7 +4651,7 @@ export class Session { payload: { success: false, removedAgents: [], - error: this.toCheckoutError(error), + error: toCheckoutError(error), requestId, }, }); @@ -7590,10 +7236,9 @@ export class Session { this.terminalExitSubscriptions.clear(); this.disposeTerminalSubscriptions(); - for (const target of this.checkoutDiffTargets.values()) { - this.closeCheckoutDiffWatchTarget(target); + for (const unsubscribe of this.checkoutDiffSubscriptions.values()) { + unsubscribe(); } - this.checkoutDiffTargets.clear(); this.checkoutDiffSubscriptions.clear(); for (const target of this.workspaceGitWatchTargets.values()) { diff --git a/packages/server/src/server/session.workspace-git-watch.test.ts b/packages/server/src/server/session.workspace-git-watch.test.ts index 4a13cfd92..77818a928 100644 --- a/packages/server/src/server/session.workspace-git-watch.test.ts +++ b/packages/server/src/server/session.workspace-git-watch.test.ts @@ -32,6 +32,8 @@ const { watchCalls, watchMock } = vi.hoisted(() => { }; }); +const resolveCheckoutGitDirMock = vi.hoisted(() => vi.fn(async () => null)); + vi.mock("node:fs", async () => { const actual = await vi.importActual("node:fs"); return { @@ -40,6 +42,10 @@ vi.mock("node:fs", async () => { }; }); +vi.mock("./checkout-git-utils.js", () => ({ + resolveCheckoutGitDir: resolveCheckoutGitDirMock, +})); + import { Session } from "./session.js"; function createSessionForWorkspaceGitWatchTests(): { @@ -120,6 +126,20 @@ function createSessionForWorkspaceGitWatchTests(): { workspaces.delete(workspaceId); }, } as any, + checkoutDiffManager: { + subscribe: async () => ({ + initial: { cwd: "/tmp", files: [], error: null }, + unsubscribe: () => {}, + }), + scheduleRefreshForCwd: () => {}, + getMetrics: () => ({ + checkoutDiffTargetCount: 0, + checkoutDiffSubscriptionCount: 0, + checkoutDiffWatcherCount: 0, + checkoutDiffFallbackRefreshTargetCount: 0, + }), + dispose: () => {}, + } as any, createAgentMcpTransport: async () => { throw new Error("not used"); }, @@ -140,6 +160,8 @@ describe("workspace git watch targets", () => { beforeEach(() => { watchCalls.length = 0; watchMock.mockClear(); + resolveCheckoutGitDirMock.mockReset(); + resolveCheckoutGitDirMock.mockResolvedValue(null); vi.useFakeTimers(); }); @@ -163,7 +185,7 @@ describe("workspace git watch targets", () => { mainRepoRoot: null, }, }); - sessionAny.resolveCheckoutGitDir = async () => "/tmp/repo/.git"; + resolveCheckoutGitDirMock.mockResolvedValue("/tmp/repo/.git"); sessionAny.workspaceUpdatesSubscription = { subscriptionId: "sub-1", filter: undefined, @@ -250,7 +272,7 @@ describe("workspace git watch targets", () => { }, }); - sessionAny.resolveCheckoutGitDir = async (cwd: string) => path.join(cwd, ".git"); + resolveCheckoutGitDirMock.mockImplementation(async (cwd: string) => path.join(cwd, ".git")); await sessionAny.ensureWorkspaceRegistered("/tmp/repo-one"); expect(sessionAny.workspaceGitWatchTargets.size).toBe(1); diff --git a/packages/server/src/server/session.workspaces.test.ts b/packages/server/src/server/session.workspaces.test.ts index 4c2b4fadb..7fd90d892 100644 --- a/packages/server/src/server/session.workspaces.test.ts +++ b/packages/server/src/server/session.workspaces.test.ts @@ -105,6 +105,20 @@ function createSessionForWorkspaceTests(): Session { archive: async () => {}, remove: async () => {}, } as any, + checkoutDiffManager: { + subscribe: async () => ({ + initial: { cwd: "/tmp", files: [], error: null }, + unsubscribe: () => {}, + }), + scheduleRefreshForCwd: () => {}, + getMetrics: () => ({ + checkoutDiffTargetCount: 0, + checkoutDiffSubscriptionCount: 0, + checkoutDiffWatcherCount: 0, + checkoutDiffFallbackRefreshTargetCount: 0, + }), + dispose: () => {}, + } as any, createAgentMcpTransport: async () => { throw new Error("not used"); }, @@ -278,6 +292,20 @@ describe("workspace aggregation", () => { archive: async () => {}, remove: async () => {}, } as any, + checkoutDiffManager: { + subscribe: async () => ({ + initial: { cwd: "/tmp", files: [], error: null }, + unsubscribe: () => {}, + }), + scheduleRefreshForCwd: () => {}, + getMetrics: () => ({ + checkoutDiffTargetCount: 0, + checkoutDiffSubscriptionCount: 0, + checkoutDiffWatcherCount: 0, + checkoutDiffFallbackRefreshTargetCount: 0, + }), + dispose: () => {}, + } as any, createAgentMcpTransport: async () => { throw new Error("not used"); }, diff --git a/packages/server/src/server/websocket-server.notifications.test.ts b/packages/server/src/server/websocket-server.notifications.test.ts index 5a27d4eed..5fe0502ee 100644 --- a/packages/server/src/server/websocket-server.notifications.test.ts +++ b/packages/server/src/server/websocket-server.notifications.test.ts @@ -88,6 +88,17 @@ function createServer(agentManagerOverrides?: Record) { {} as any, {} as any, {} as any, + { + subscribe: vi.fn(), + scheduleRefreshForCwd: vi.fn(), + getMetrics: vi.fn(() => ({ + checkoutDiffTargetCount: 0, + checkoutDiffSubscriptionCount: 0, + checkoutDiffWatcherCount: 0, + checkoutDiffFallbackRefreshTargetCount: 0, + })), + dispose: vi.fn(), + } as any, ); return { server, agentManager }; diff --git a/packages/server/src/server/websocket-server.relay-reconnect.test.ts b/packages/server/src/server/websocket-server.relay-reconnect.test.ts index 25b2f3210..e19dd7d90 100644 --- a/packages/server/src/server/websocket-server.relay-reconnect.test.ts +++ b/packages/server/src/server/websocket-server.relay-reconnect.test.ts @@ -183,6 +183,17 @@ function createServer(options?: { speechReadiness?: SpeechReadinessSnapshot | nu {} as any, {} as any, {} as any, + { + subscribe: vi.fn(), + scheduleRefreshForCwd: vi.fn(), + getMetrics: vi.fn(() => ({ + checkoutDiffTargetCount: 0, + checkoutDiffSubscriptionCount: 0, + checkoutDiffWatcherCount: 0, + checkoutDiffFallbackRefreshTargetCount: 0, + })), + dispose: vi.fn(), + } as any, ); } diff --git a/packages/server/src/server/websocket-server.ts b/packages/server/src/server/websocket-server.ts index 6e57731e5..ddd6a29ba 100644 --- a/packages/server/src/server/websocket-server.ts +++ b/packages/server/src/server/websocket-server.ts @@ -12,6 +12,7 @@ import type { ProjectRegistry, WorkspaceRegistry } from "./workspace-registry.js import type { FileBackedChatService } from "./chat/chat-service.js"; import type { LoopService } from "./loop-service.js"; import type { ScheduleService } from "./schedule/service.js"; +import type { CheckoutDiffManager, CheckoutDiffMetrics } from "./checkout-diff-manager.js"; import { type ServerInfoStatusPayload, type WSHelloMessage, @@ -65,6 +66,8 @@ type WebSocketServerConfig = { allowedHosts?: AllowedHostsConfig; }; +type WebSocketRuntimeMetrics = SessionRuntimeMetrics & CheckoutDiffMetrics; + function createNoopProjectRegistry(): ProjectRegistry { return { initialize: async () => {}, @@ -234,6 +237,7 @@ export class VoiceAssistantWebSocketServer { private readonly chatService: FileBackedChatService; private readonly loopService: LoopService; private readonly scheduleService: ScheduleService; + private readonly checkoutDiffManager: CheckoutDiffManager; private readonly downloadTokenStore: DownloadTokenStore; private readonly paseoHome: string; private readonly pushTokenStore: PushTokenStore; @@ -323,6 +327,7 @@ export class VoiceAssistantWebSocketServer { chatService?: FileBackedChatService, loopService?: LoopService, scheduleService?: ScheduleService, + checkoutDiffManager?: CheckoutDiffManager, ) { this.logger = logger.child({ module: "websocket-server" }); this.serverId = serverId; @@ -346,6 +351,10 @@ export class VoiceAssistantWebSocketServer { throw new Error("VoiceAssistantWebSocketServer requires a schedule service."); } this.scheduleService = scheduleService; + if (!checkoutDiffManager) { + throw new Error("VoiceAssistantWebSocketServer requires a checkout diff manager."); + } + this.checkoutDiffManager = checkoutDiffManager; this.downloadTokenStore = downloadTokenStore; this.paseoHome = paseoHome; this.createAgentMcpTransport = createAgentMcpTransport; @@ -504,6 +513,7 @@ export class VoiceAssistantWebSocketServer { } await Promise.all(cleanupPromises); + this.checkoutDiffManager.dispose(); this.pendingConnections.clear(); this.sessions.clear(); this.externalSessionsByKey.clear(); @@ -646,6 +656,7 @@ export class VoiceAssistantWebSocketServer { chatService: this.chatService, loopService: this.loopService, scheduleService: this.scheduleService, + checkoutDiffManager: this.checkoutDiffManager, createAgentMcpTransport: this.createAgentMcpTransport, stt: this.stt, tts: this.tts, @@ -1230,12 +1241,8 @@ export class VoiceAssistantWebSocketServer { return stats.slice(0, 15); } - private collectSessionRuntimeMetrics(): SessionRuntimeMetrics { + private collectSessionRuntimeMetrics(): WebSocketRuntimeMetrics { const uniqueConnections = new Set(this.externalSessionsByKey.values()); - let checkoutDiffTargetCount = 0; - let checkoutDiffSubscriptionCount = 0; - let checkoutDiffWatcherCount = 0; - let checkoutDiffFallbackRefreshTargetCount = 0; let terminalDirectorySubscriptionCount = 0; let terminalSubscriptionCount = 0; let inflightRequests = 0; @@ -1243,11 +1250,6 @@ export class VoiceAssistantWebSocketServer { for (const connection of uniqueConnections) { const sessionMetrics = connection.session.getRuntimeMetrics(); - checkoutDiffTargetCount += sessionMetrics.checkoutDiffTargetCount; - checkoutDiffSubscriptionCount += sessionMetrics.checkoutDiffSubscriptionCount; - checkoutDiffWatcherCount += sessionMetrics.checkoutDiffWatcherCount; - checkoutDiffFallbackRefreshTargetCount += - sessionMetrics.checkoutDiffFallbackRefreshTargetCount; terminalDirectorySubscriptionCount += sessionMetrics.terminalDirectorySubscriptionCount; terminalSubscriptionCount += sessionMetrics.terminalSubscriptionCount; inflightRequests += sessionMetrics.inflightRequests; @@ -1256,10 +1258,7 @@ export class VoiceAssistantWebSocketServer { } return { - checkoutDiffTargetCount, - checkoutDiffSubscriptionCount, - checkoutDiffWatcherCount, - checkoutDiffFallbackRefreshTargetCount, + ...this.checkoutDiffManager.getMetrics(), terminalDirectorySubscriptionCount, terminalSubscriptionCount, inflightRequests,