From 7e2bf1c366b5f0daf1e0532aa857df08f4464825 Mon Sep 17 00:00:00 2001 From: axisrow Date: Tue, 29 Sep 2026 18:29:20 +0800 Subject: [PATCH 1/3] feat(daemon): async providers, scoped caches, refresh dedup, cancellation MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implements the daemon-epic step for issue #18 on top of the IPC transport (#44): provider work now runs in the daemon's scope through async twins, and widget output in one-shot mode is unchanged. - RefreshGroup (src/daemon/provider-scope.ts): single-flight refresh per key with prompt cancellation when the last consumer leaves; capMap FIFO bound for shared maps. - Daemon render path: no more process.env/cwd swapping — requests carry an env/cwd snapshot (RenderInvocation); identical display contexts join one in-flight render (deduped counter); bounded in-flight renders; idle cache sweep reclaims provider caches after quiet periods. - Prefetch (src/daemon/prefetch.ts): transcript analysis, usage (per credential fingerprint), service status, block metrics and git/jj/review/ custom-command warmups run async before sync formatting. - Async provider twins: git/jj execFile, custom-command capture (bounded output, process-group kill), usage fetch with signal/env scope, git review cache, block metrics — sync paths unchanged for one-shot mode. - Scoped caches: usage memory cache gated by credential fingerprint BEFORE the fast return (missing credentials never take a logged-in entry); transcript whole-analysis reuse validated by file identity with invalidation on truncation, append, replacement and subagent updates. - Bounded memory: capped caches/maps, idle sweep, no env or token logging. Tests: RefreshGroup dedup/cancellation, daemon render join semantics, transcript reuse/invalidation, usage identity gating (child-process probes, same style as usage-fetch.test.ts). Co-Authored-By: Claude Code --- src/daemon/__tests__/provider-scope.test.ts | 120 ++++++ src/daemon/__tests__/server-dedup.test.ts | 90 +++++ src/daemon/__tests__/test-daemon.ts | 13 +- src/daemon/prefetch.ts | 370 +++++++++++++++++++ src/daemon/provider-scope.ts | 73 ++++ src/daemon/server.ts | 286 ++++++++------ src/render.ts | 70 +++- src/types/RenderContext.ts | 10 + src/utils/__tests__/transcript-reuse.test.ts | 108 ++++++ src/utils/__tests__/usage-identity.test.ts | 185 ++++++++++ src/utils/claude-service-status.ts | 45 ++- src/utils/claude-settings.ts | 10 +- src/utils/context-percentage.ts | 4 +- src/utils/custom-command-capture.ts | 142 ++++++- src/utils/custom-command.ts | 124 ++++++- src/utils/git-review-cache.ts | 319 ++++++++++++++++ src/utils/git.ts | 108 +++++- src/utils/jj.ts | 75 +++- src/utils/jsonl-blocks.ts | 201 ++++++++-- src/utils/jsonl-cache.ts | 46 ++- src/utils/jsonl-metrics.ts | 166 ++++++++- src/utils/jsonl.ts | 1 + src/utils/model-context.ts | 15 +- src/utils/terminal.ts | 22 ++ src/utils/usage-fetch.ts | 224 +++++++++-- src/utils/usage-prefetch.ts | 41 +- src/utils/usage-windows.ts | 5 +- src/widgets/ContextBar.ts | 2 +- src/widgets/ContextPercentageUsable.ts | 2 +- src/widgets/ContextWindow.ts | 2 +- src/widgets/CustomCommand.tsx | 22 +- src/widgets/GitCiStatus.ts | 2 +- src/widgets/GitPr.ts | 2 +- 33 files changed, 2643 insertions(+), 262 deletions(-) create mode 100644 src/daemon/__tests__/provider-scope.test.ts create mode 100644 src/daemon/__tests__/server-dedup.test.ts create mode 100644 src/daemon/prefetch.ts create mode 100644 src/daemon/provider-scope.ts create mode 100644 src/utils/__tests__/transcript-reuse.test.ts create mode 100644 src/utils/__tests__/usage-identity.test.ts diff --git a/src/daemon/__tests__/provider-scope.test.ts b/src/daemon/__tests__/provider-scope.test.ts new file mode 100644 index 000000000..c0f92e308 --- /dev/null +++ b/src/daemon/__tests__/provider-scope.test.ts @@ -0,0 +1,120 @@ +import { + describe, + expect, + it, + vi +} from 'vitest'; + +import { + RefreshGroup, + capMap +} from '../provider-scope'; + +describe('capMap', () => { + it('evicts oldest insertions beyond the cap', () => { + const map = new Map(); + map.set('a', 1); + map.set('b', 2); + map.set('c', 3); + capMap(map, 2); + expect([...map.keys()]).toEqual(['b', 'c']); + }); + + it('keeps a map at or below the cap untouched', () => { + const map = new Map([['a', 1]]); + capMap(map, 4); + expect(map.size).toBe(1); + }); +}); + +describe('RefreshGroup', () => { + it('dedups concurrent refreshes of the same key into one work run', async () => { + const group = new RefreshGroup(); + const work = vi.fn((): Promise => Promise.resolve('value')); + + const first = group.refresh('git', work); + const second = group.refresh('git', work); + + expect(work).toHaveBeenCalledTimes(1); + await expect(first.promise).resolves.toBe('value'); + await expect(second.promise).resolves.toBe('value'); + first.release(); + second.release(); + }); + + it('runs different keys independently', async () => { + const group = new RefreshGroup(); + const work = vi.fn((): Promise => Promise.resolve('ok')); + + const a = group.refresh('cwd-a', work); + const b = group.refresh('cwd-b', work); + expect(group.size).toBe(2); + + a.release(); + await expect(a.promise).resolves.toBeDefined(); + b.release(); + await expect(b.promise).resolves.toBeDefined(); + }); + + it('aborts the work when the last consumer releases before completion', async () => { + const group = new RefreshGroup(); + let observedAbort = false; + const work = (signal: AbortSignal): Promise => new Promise((resolve, reject) => { + signal.addEventListener('abort', () => { + observedAbort = true; + reject(new Error('aborted')); + }); + }); + + const only = group.refresh('usage', work); + const outcome = only.promise.then(() => 'resolved', () => 'rejected'); + only.release(); + expect(observedAbort).toBe(true); + await expect(outcome).resolves.toBe('rejected'); + // The settled job is unregistered on its own microtask. + await Promise.resolve(); + expect(group.size).toBe(0); + }); + + it('does not abort while another consumer is still joined', async () => { + const group = new RefreshGroup(); + const work = vi.fn((): Promise => Promise.resolve('shared')); + + const first = group.refresh('git', work); + const second = group.refresh('git', work); + first.release(); + + await expect(second.promise).resolves.toBe('shared'); + second.release(); + }); + + it('does not join a cancelled job: a later refresh starts fresh work', async () => { + const group = new RefreshGroup(); + const resolvers: ((value: string) => void)[] = []; + const work = vi.fn((): Promise => new Promise((resolve) => { resolvers.push(resolve); })); + + const first = group.refresh('review', work); + first.release(); // Still pending -> aborts the job. + expect(work).toHaveBeenCalledTimes(1); + + const second = group.refresh('review', work); + expect(work).toHaveBeenCalledTimes(2); + resolvers[1]?.('fresh'); + await expect(second.promise).resolves.toBe('fresh'); + second.release(); + }); + + it('ignores a release after completion', async () => { + const group = new RefreshGroup(); + const work = vi.fn((): Promise => Promise.resolve('done')); + + const only = group.refresh('git', work); + await only.promise; + only.release(); + only.release(); + + const next = group.refresh('git', work); + await expect(next.promise).resolves.toBe('done'); + next.release(); + }); +}); diff --git a/src/daemon/__tests__/server-dedup.test.ts b/src/daemon/__tests__/server-dedup.test.ts new file mode 100644 index 000000000..a9de17bdd --- /dev/null +++ b/src/daemon/__tests__/server-dedup.test.ts @@ -0,0 +1,90 @@ +import { + afterEach, + describe, + expect, + it +} from 'vitest'; + +import type { LoadedSettings } from '../../utils/config'; + +import type { StartedTestDaemon } from './test-daemon'; +import { + MODEL_ONLY_SETTINGS, + renderRequest, + startTestDaemon, + stopTestDaemon, + waitFor +} from './test-daemon'; + +// Render dedup at the daemon boundary (#18): identical display contexts join +// one in-flight render. Consumer-release cancellation semantics are covered +// deterministically in provider-scope.test.ts (RefreshGroup); the socket-level +// "last client left" wiring cannot be exercised under bun, which never +// surfaces client disconnects to node:http servers (verified against 1.3.13). + +interface LoadSettingsGate { + loadSettings: (configPath: string) => Promise; + readonly pending: number; + openAll: () => void; +} + +function makeLoadSettingsGate(): LoadSettingsGate { + const resolvers: (() => void)[] = []; + return { + loadSettings: () => new Promise((resolve) => { + resolvers.push(() => { resolve({ settings: { ...MODEL_ONLY_SETTINGS }, loadError: null }); }); + }), + get pending() { + return resolvers.length; + }, + openAll: () => { + while (resolvers.length > 0) { + resolvers.shift()?.(); + } + } + }; +} + +const BODY_A = JSON.stringify({ model: { id: 'claude-dedup-a' }, cwd: '/tmp' }); +const BODY_B = JSON.stringify({ model: { id: 'claude-dedup-b' }, cwd: '/tmp' }); + +describe('daemon render dedup (#18)', () => { + let started: StartedTestDaemon | undefined; + + afterEach(async () => { + if (started) { + await stopTestDaemon(started); + started = undefined; + } + }); + + it('joins identical in-flight renders into one job', async () => { + const gate = makeLoadSettingsGate(); + started = await startTestDaemon({ loadSettings: gate.loadSettings }); + + const first = renderRequest(started.daemon, started.daemon.token, BODY_A); + const second = renderRequest(started.daemon, started.daemon.token, BODY_A); + await waitFor(() => gate.pending === 1); + gate.openAll(); + + const [firstResponse, secondResponse] = await Promise.all([first, second]); + expect(firstResponse.status).toBe(200); + expect(secondResponse.text).toBe(firstResponse.text); + expect(started.daemon.counters.deduped).toBe(1); + // One render ran: one invocation snapshot built. + expect(started.dependencies.invocations).toHaveLength(1); + }); + + it('does not join renders with different payloads', async () => { + started = await startTestDaemon({}); + + const [firstResponse, secondResponse] = await Promise.all([ + renderRequest(started.daemon, started.daemon.token, BODY_A), + renderRequest(started.daemon, started.daemon.token, BODY_B) + ]); + expect(firstResponse.status).toBe(200); + expect(secondResponse.status).toBe(200); + expect(started.daemon.counters.deduped ?? 0).toBe(0); + expect(started.dependencies.invocations).toHaveLength(2); + }); +}); diff --git a/src/daemon/__tests__/test-daemon.ts b/src/daemon/__tests__/test-daemon.ts index f6b4542d5..1607d3605 100644 --- a/src/daemon/__tests__/test-daemon.ts +++ b/src/daemon/__tests__/test-daemon.ts @@ -36,12 +36,19 @@ export function hermeticDependencies(overrides: Partial = {} }); return { loadSettings, - resolveTerminalWidth: () => 120, - buildInvocation: (context, terminalWidth) => { + resolveTerminalWidth: (_sessionId, _ttlSeconds, env) => { + // Honor the request env the way the production dependency does: + // an explicit override wins, otherwise a fixed width. + if (env.CCSTATUSLINE_WIDTH !== undefined) { + return Number.parseInt(env.CCSTATUSLINE_WIDTH, 10) || 120; + } + return 120; + }, + buildInvocation: (context, terminalWidth, env) => { const invocation: RenderInvocation = { configPath: '/tmp/hermetic-settings.json', cwd: context.cwd ?? '/tmp', - env: { ...process.env }, + env, terminalWidth }; invocations.push(invocation); diff --git a/src/daemon/prefetch.ts b/src/daemon/prefetch.ts new file mode 100644 index 000000000..556bd2c33 --- /dev/null +++ b/src/daemon/prefetch.ts @@ -0,0 +1,370 @@ +import type { RenderPrefetch } from '../render'; +import type { BlockMetrics } from '../types'; +import type { RenderContext } from '../types/RenderContext'; +import type { Settings } from '../types/Settings'; +import type { StatusJSON } from '../types/StatusJSON'; +import type { WidgetItem } from '../types/Widget'; +import { prefetchClaudeStatusIfNeeded } from '../utils/claude-service-status'; +import { getClaudeConfigDir } from '../utils/claude-settings'; +import { + buildCustomCommandRequest, + runCustomCommandAsync +} from '../utils/custom-command'; +import { + getExecutedGitCommands, + resolveGitCwd, + runGitArgsAsync +} from '../utils/git'; +import { fetchGitReviewDataAsync } from '../utils/git-review-cache'; +import { + getExecutedJjCommands, + runJjArgsAsync +} from '../utils/jj'; +import { getCachedBlockMetricsAsync } from '../utils/jsonl-cache'; +import { getTranscriptAnalysis } from '../utils/jsonl-metrics'; +import { + getWidgetSpeedWindowSeconds, + isWidgetSpeedWindowEnabled +} from '../utils/speed-window'; +import type { UsageMemoryCache } from '../utils/usage-fetch'; +import { + getUsageCredentialsAsync, + getUsageScopeKey +} from '../utils/usage-fetch'; +import { + hasUsageDependentWidgets, + prefetchUsageDataIfNeeded +} from '../utils/usage-prefetch'; + +import { + RefreshGroup, + capMap +} from './provider-scope'; + +// Daemon render prefetch (#18): everything below runs before the synchronous +// formatting section, off the render loop, deduped through one RefreshGroup so +// concurrent sessions that need the same git command / usage account / review +// lookup share a single execution, and released as soon as the requesting +// render goes away. Widget output never changes: prefetch warms exactly the +// caches and computes exactly the values the sync formatter reads. + +/** Per-daemon state shared across requests. */ +export interface PrefetchState { + refresh: RefreshGroup; + /** One usage memory cache per credential fingerprint (account scope). */ + usageCaches: Map; +} + +export function createPrefetchState(): PrefetchState { + return { + refresh: new RefreshGroup(), + usageCaches: new Map() + }; +} + +export interface PrefetchRequestScope { + /** Merged request environment (daemon env + allowlisted request values). */ + env: NodeJS.ProcessEnv; + /** Absolute request working directory. */ + cwd: string; + /** Terminal width resolved for this request (feeds custom command input). */ + terminalWidth: number | null; + signal: AbortSignal; + state: PrefetchState; +} + +const USAGE_CACHES_MAX_ENTRIES = 8; + +// git commands each widget family runs, keyed by widget type. Every git widget +// also runs the work-tree check. ponytail: this table must list new git widget +// commands or their first render per repo pays one sync spawn; the executed- +// command log (recorded on every sync run) covers everything after that. +const GIT_PREFETCH_COMMANDS: Record = { + 'git-branch': ['symbolic-ref --short HEAD'], + 'git-root-dir': ['rev-parse --show-toplevel'], + 'git-worktree': ['rev-parse --git-dir'], + 'git-changes': ['diff --shortstat', 'diff --cached --shortstat'], + 'git-insertions': ['diff --shortstat', 'diff --cached --shortstat'], + 'git-deletions': ['diff --shortstat', 'diff --cached --shortstat'], + 'git-status': ['status --porcelain -z'], + 'git-staged': ['status --porcelain -z'], + 'git-unstaged': ['status --porcelain -z'], + 'git-untracked': ['status --porcelain -z'], + 'git-clean-status': ['status --porcelain -z'], + 'git-staged-files': ['status --porcelain -z'], + 'git-unstaged-files': ['status --porcelain -z'], + 'git-untracked-files': ['status --porcelain -z'], + 'git-ahead-behind': ['rev-list --left-right --count HEAD...@{upstream}'], + 'git-conflicts': ['ls-files --unmerged'], + 'git-sha': ['rev-parse --short HEAD'], + 'git-is-fork': ['remote', 'remote get-url -- origin'], + 'git-origin-owner': ['remote', 'remote get-url -- origin'], + 'git-origin-repo': ['remote', 'remote get-url -- origin'], + 'git-origin-owner-repo': ['remote', 'remote get-url -- origin'], + 'git-upstream-owner': ['remote', 'remote get-url -- upstream', 'rev-parse --abbrev-ref --symbolic-full-name @{upstream}'], + 'git-upstream-repo': ['remote', 'remote get-url -- upstream', 'rev-parse --abbrev-ref --symbolic-full-name @{upstream}'], + 'git-upstream-owner-repo': ['remote', 'remote get-url -- upstream', 'rev-parse --abbrev-ref --symbolic-full-name @{upstream}'] +}; + +const GIT_BASE_COMMAND = 'rev-parse --is-inside-work-tree'; +const JJ_BASE_COMMAND = 'root'; + +function collectGitCommandTokens(lineItems: WidgetItem[]): Set { + const tokens = new Set(); + for (const item of lineItems) { + if (!item.type.startsWith('git-')) { + continue; + } + tokens.add(GIT_BASE_COMMAND); + for (const command of GIT_PREFETCH_COMMANDS[item.type] ?? []) { + tokens.add(command); + } + } + return tokens; +} + +function hasGitReviewWidgets(lineItems: WidgetItem[][]): boolean { + return lineItems.some(line => line.some(item => item.type === 'git-pr' || item.type === 'git-ci-status')); +} + +function hasBlockTimerWidgets(lineItems: WidgetItem[][]): boolean { + return lineItems.some(line => line.some(item => item.type === 'block-timer' || item.type === 'block-reset-timer')); +} + +/** + * Mirror of the transcript options renderStatusLines computes for itself. Must + * stay in sync with render.ts — the analysis cache keys on these options, so a + * drift means two different option sets thrash one cache entry per path. + */ +function transcriptOptionsFor(data: StatusJSON, settings: Settings): Parameters[1] { + const lines = settings.lines; + const speedWidgetTypes = new Set(['output-speed', 'input-speed', 'total-speed']); + const hasSessionClock = lines.some(line => line.some(item => item.type === 'session-clock')); + const hasSpeedItems = lines.some(line => line.some(item => speedWidgetTypes.has(item.type))); + const hasCompactionWidget = lines.some(line => line.some(item => item.type === 'compaction-counter')); + const hasThinkingEffortWidget = lines.some(line => line.some(item => item.type === 'thinking-effort')); + const hasSessionNameWidget = lines.some(line => line.some(item => item.type === 'session-name')); + const hasLastTurnTokensWidget = lines.some(line => line.some(item => item.type === 'tokens-last-turn')); + const needsTranscriptThinkingEffort = hasThinkingEffortWidget + && (!data.effort || !('level' in data.effort)); + const hasSessionDurationInStatusJson = typeof data.cost?.total_duration_ms === 'number' + && Number.isFinite(data.cost.total_duration_ms) && data.cost.total_duration_ms >= 0; + const requestedSpeedWindows = new Set(); + for (const line of lines) { + for (const item of line) { + if (speedWidgetTypes.has(item.type) && isWidgetSpeedWindowEnabled(item)) { + requestedSpeedWindows.add(getWidgetSpeedWindowSeconds(item)); + } + } + } + + return { + includeSessionDuration: hasSessionClock && !hasSessionDurationInStatusJson, + includeSpeedMetrics: hasSpeedItems, + includeSubagents: true, + speedWindowSeconds: Array.from(requestedSpeedWindows), + includeCompactionStats: hasCompactionWidget, + includeThinkingEffort: needsTranscriptThinkingEffort, + includeSessionName: hasSessionNameWidget, + includeLastTurnTokens: hasLastTurnTokensWidget + }; +} + +/** + * Build the minimal render context the provider twins resolve cwd/env from. + * resolveGitCwd reads data.cwd/workspace; the request cwd is the same fallback + * the widgets see (context.cwd). + */ +function prefetchRenderContext(scope: PrefetchRequestScope, data: StatusJSON, settings: Settings): RenderContext { + return { + data, + env: scope.env, + cwd: scope.cwd, + terminalWidth: scope.terminalWidth, + gitCacheTtlSeconds: settings.gitCacheTtlSeconds, + customCommandCacheTtlSeconds: settings.customCommandCacheTtlSeconds + }; +} + +/** Track one shared refresh for this request: aborting the request releases it. */ +function trackRefresh( + scope: PrefetchRequestScope, + key: string, + work: (signal: AbortSignal) => Promise +): Promise { + const { promise, release } = scope.state.refresh.refresh(key, work); + if (scope.signal.aborted) { + release(); + return promise; + } + scope.signal.addEventListener('abort', release, { once: true }); + return promise; +} + +async function prefetchGitCommands( + lineItems: WidgetItem[][], + scope: PrefetchRequestScope, + data: StatusJSON, + settings: Settings +): Promise { + const gitContext: RenderContext = prefetchRenderContext(scope, data, settings); + const gitCwd = resolveGitCwd(gitContext) ?? scope.cwd; + + const tokens = collectGitCommandTokens(lineItems.flat()); + for (const args of getExecutedGitCommands(gitCwd)) { + tokens.add(args.join(' ')); + } + for (const token of tokens) { + const args = token.split(/\s+/).filter(Boolean); + if (args.length === 0) { + continue; + } + // ponytail: the git twins enforce their own 5s timeout, so the group + // signal is not wired into the spawn here; cancellation bounds joins, + // not local git. Wire execFile's signal through if that ever matters. + await trackRefresh(scope, `git:${gitCwd}\0${token}`, () => runGitArgsAsync(args, gitContext, token)) + .catch(() => undefined); + } + + if (lineItems.some(line => line.some(item => item.type.startsWith('jj-')))) { + const jjTokens = new Set([JJ_BASE_COMMAND]); + for (const args of getExecutedJjCommands(gitCwd)) { + jjTokens.add(args.join('\0')); + } + for (const token of jjTokens) { + const args = token === JJ_BASE_COMMAND ? [JJ_BASE_COMMAND] : token.split('\0'); + await trackRefresh(scope, `jj:${gitCwd}\0${token}`, () => runJjArgsAsync(args, gitContext)) + .catch(() => undefined); + } + } +} + +async function prefetchGitReview( + lineItems: WidgetItem[][], + scope: PrefetchRequestScope, + data: StatusJSON, + settings: Settings +): Promise { + if (!hasGitReviewWidgets(lineItems)) { + return; + } + const gitContext: RenderContext = prefetchRenderContext(scope, data, settings); + const gitCwd = resolveGitCwd(gitContext) ?? scope.cwd; + const includeChecks = lineItems.some(line => line.some(item => item.type === 'git-ci-status')); + + await trackRefresh( + scope, + `git-review:${gitCwd}\0${includeChecks ? 'checks' : 'metadata'}`, + signal => fetchGitReviewDataAsync(gitCwd, { includeChecks }, signal) + ).catch(() => undefined); +} + +async function prefetchCustomCommands( + lineItems: WidgetItem[][], + scope: PrefetchRequestScope, + data: StatusJSON, + settings: Settings +): Promise { + const context: RenderContext = prefetchRenderContext(scope, data, settings); + for (const item of lineItems.flat()) { + if (item.type !== 'custom-command') { + continue; + } + const request = buildCustomCommandRequest(item, context); + if (request === null) { + continue; + } + const key = `custom:${request.cwd ?? scope.cwd}\0${request.command}\0${request.sessionId ?? ''}\0${request.terminalWidth ?? ''}`; + await trackRefresh(scope, key, signal => runCustomCommandAsync(request, signal)) + .catch(() => undefined); + } +} + +async function prefetchUsage( + lineItems: WidgetItem[][], + scope: PrefetchRequestScope, + data: StatusJSON +): Promise>> { + if (!hasUsageDependentWidgets(lineItems)) { + return null; + } + + // Account scope before anything else (#18): credentials resolved once per + // render (deduped), and the memory cache picked by credential fingerprint + // so no fast return can cross accounts. + const credentials = await trackRefresh( + scope, + 'usage-credentials', + signal => getUsageCredentialsAsync(scope.env, signal) + ); + const scopeKey = credentials ? getUsageScopeKey(credentials) : 'none'; + let cache = scope.state.usageCaches.get(scopeKey); + if (!cache) { + cache = { data: null, time: 0, errorMaxAge: 30 }; + scope.state.usageCaches.set(scopeKey, cache); + capMap(scope.state.usageCaches, USAGE_CACHES_MAX_ENTRIES); + } + + return trackRefresh(scope, `usage:${scopeKey}`, (groupSignal) => { + return prefetchUsageDataIfNeeded(lineItems, data, { + cache: cache, + credentials, + signal: groupSignal, + env: scope.env + }); + }); +} + +/** + * Gather everything the render needs before the synchronous formatting + * section (#18). Individual prefetch failures never fail the render: the + * formatter falls back to computing (or the cached fallback) exactly as the + * one-shot path does. + */ +export async function prefetchRenderData( + data: StatusJSON, + settings: Settings, + scope: PrefetchRequestScope +): Promise { + const lineItems = settings.lines; + + const transcriptAnalysis: Promise = data.transcript_path + ? getTranscriptAnalysis(data.transcript_path, transcriptOptionsFor(data, settings)) + .catch(() => null) + : Promise.resolve(null); + + const usageData = prefetchUsage(lineItems, scope, data).catch(() => null); + + const claudeStatusData = trackRefresh( + scope, + 'claude-status', + signal => prefetchClaudeStatusIfNeeded(lineItems, { signal, env: scope.env }) + ).catch(() => null); + + // Warm-ups land in shared provider caches the sync section reads. + const warmups = [ + prefetchGitCommands(lineItems, scope, data, settings), + prefetchGitReview(lineItems, scope, data, settings), + prefetchCustomCommands(lineItems, scope, data, settings) + ]; + + const blockMetrics: Promise = hasBlockTimerWidgets(lineItems) + ? trackRefresh(scope, `block:${getClaudeConfigDir(scope.env)}`, () => getCachedBlockMetricsAsync(scope.env)) + .catch(() => undefined) + : Promise.resolve(undefined); + + const [transcript, usage, claudeStatus, block] = await Promise.all([ + transcriptAnalysis, + usageData, + claudeStatusData, + blockMetrics + ]); + await Promise.all(warmups); + + return { + transcriptAnalysis: transcript, + usageData: usage, + claudeStatusData: claudeStatus, + ...(block !== undefined ? { blockMetrics: block } : {}) + }; +} diff --git a/src/daemon/provider-scope.ts b/src/daemon/provider-scope.ts new file mode 100644 index 000000000..d5d7b41c7 --- /dev/null +++ b/src/daemon/provider-scope.ts @@ -0,0 +1,73 @@ +// Shared provider-refresh machinery for the daemon (#18): a bounded Map cap +// and a single-flight group. The group keys in-flight provider work (a git +// spawn, a usage fetch, a review lookup) so concurrent renders that need the +// same data join one execution, and the underlying subprocess/HTTP work is +// aborted only when the last consumer goes away — one session disconnecting +// must never kill a refresh another session is still waiting on. + +/** Truncate a Map to the newest `maxEntries` insertions (FIFO eviction). */ +export function capMap(map: Map, maxEntries: number): void { + while (map.size > maxEntries) { + const oldest = map.keys().next(); + if (oldest.done) { + return; + } + map.delete(oldest.value); + } +} + +interface RefreshJob { + controller: AbortController; + consumers: number; + promise: Promise; +} + +export class RefreshGroup { + private readonly jobs = new Map(); + + /** + * Run `work` once per key while consumers remain. Every caller gets a + * release function; when the last consumer releases before completion the + * job's signal fires and the work is expected to settle on its own. A + * caller releasing after completion is a no-op. Joining an already + * cancelled job starts a fresh one. + */ + refresh(key: string, work: (signal: AbortSignal) => Promise): { promise: Promise; release: () => void } { + const existing = this.jobs.get(key); + if (existing && !existing.controller.signal.aborted) { + existing.consumers++; + return { + promise: existing.promise as Promise, + release: () => { this.release(key, existing); } + }; + } + + const controller = new AbortController(); + const job: RefreshJob = { + controller, + consumers: 1, + promise: work(controller.signal).finally(() => { + this.jobs.delete(key); + }) + }; + this.jobs.set(key, job); + return { + promise: job.promise as Promise, + release: () => { this.release(key, job); } + }; + } + + private release(key: string, job: RefreshJob): void { + if (job.consumers > 0) { + job.consumers--; + } + if (job.consumers === 0 && this.jobs.get(key) === job) { + job.controller.abort(); + } + } + + /** In-flight job count; diagnostics and tests only. */ + get size(): number { + return this.jobs.size; + } +} diff --git a/src/daemon/server.ts b/src/daemon/server.ts index a7cb219d5..93bbd53bc 100644 --- a/src/daemon/server.ts +++ b/src/daemon/server.ts @@ -10,10 +10,13 @@ import { getConfigPath, loadSettingsFrom } from '../utils/config'; +import { clearCustomCommandCache } from '../utils/custom-command'; +import { clearGitCache } from '../utils/git'; +import { clearJjCommandLog } from '../utils/jj'; +import { clearTranscriptAnalysisCache } from '../utils/jsonl-metrics'; import { getPackageVersion, - getTerminalWidth, - resetTerminalWidthCache + getTerminalWidth } from '../utils/terminal'; import { @@ -24,6 +27,11 @@ import { prepareSocketPath, sweepStaleSockets } from './paths'; +import type { PrefetchState } from './prefetch'; +import { + createPrefetchState, + prefetchRenderData +} from './prefetch'; import type { InvocationContext } from './protocol'; import { AUTH_SCHEME, @@ -43,15 +51,18 @@ import { // through the 0600 discovery file in the 0700 runtime directory, so another // user can neither connect to the socket nor read the token. // -// Renders are strictly serialized inside the process: the render path still -// reads process.env/process.cwd in a few places (claude-settings, terminal -// width, custom-command caches), so each request applies its invocation -// context there, one at a time, and restores it afterwards. +// Renders run concurrently (#18): the render path resolves everything through +// the per-request invocation snapshot (env/cwd) and prefetched provider data, +// so no process-global state is swapped. Requests with the exact same display +// context (config path + env + cwd + width + payload) join one in-flight +// render instead of repeating the work, and each client connection holds a +// consumer slot — the last one leaving cancels the provider work that only it +// still needed. export interface DaemonDependencies { loadSettings: (configPath: string) => Promise; - resolveTerminalWidth: (sessionId: string | undefined, ttlSeconds: number) => number | null; - buildInvocation: (context: InvocationContext, terminalWidth: number | null) => RenderInvocation; + resolveTerminalWidth: (sessionId: string | undefined, ttlSeconds: number, env: NodeJS.ProcessEnv) => number | null; + buildInvocation: (context: InvocationContext, terminalWidth: number | null, env: NodeJS.ProcessEnv) => RenderInvocation; } export interface DaemonServerOptions { @@ -81,64 +92,62 @@ class BodyTooLargeError extends Error { } } -/** Production dependencies: same request-scoped wiring as the --serve loop. */ +/** + * Apply the allowlisted request env on top of the daemon's own environment: + * names present in the snapshot are set, names absent from it are cleared — + * that is what preserves absent-vs-empty across the IPC boundary. Everything + * outside the allowlist stays the daemon's, so spawned providers always have + * PATH/HOME. Pure: returns a fresh object, never touches process.env (#18). + */ +export function mergeRequestEnvironment(context: InvocationContext): NodeJS.ProcessEnv { + const merged: NodeJS.ProcessEnv = { ...process.env }; + for (const name of ENV_ALLOWLIST) { + const value = context.env[name]; + if (value === undefined) { + Reflect.deleteProperty(merged, name); + } else { + merged[name] = value; + } + } + return merged; +} + +/** Production dependencies: request-scoped wiring for the --serve loop. */ export function createProcessDaemonDependencies(): DaemonDependencies { return { loadSettings: configPath => loadSettingsFrom(configPath), - resolveTerminalWidth: (sessionId, ttlSeconds) => { - // Same rationale as serve.ts: the width memo is process-global, - // so reset before each probe or a resize would serve stale widths. - resetTerminalWidthCache(); - return getTerminalWidth({ sessionId, ttlSeconds }); + resolveTerminalWidth: (sessionId, ttlSeconds, env) => { + // The env snapshot carries CCSTATUSLINE_WIDTH/COLUMNS per request; + // without an explicit width the shared memoized probe answers + // (the daemon's own ancestry is stable for the process lifetime). + return getTerminalWidth({ sessionId, ttlSeconds, env }); }, - buildInvocation: (context, terminalWidth) => ({ + buildInvocation: (context, terminalWidth, env) => ({ configPath: getConfigPath(), cwd: context.cwd ?? process.cwd(), - // The context is applied to process.env for the duration of the - // render, so this snapshot is exactly what render-path reads see. - env: { ...process.env }, + env, terminalWidth }) }; } -/** - * Unset an environment variable. `delete process.env[name]` is the only - * correct unset (assigning undefined would stringify); the Reflect form is - * the same operation, kept here so the dynamic-delete lint rule holds. - */ -function unsetEnv(name: string): void { - Reflect.deleteProperty(process.env, name); -} - -/** - * Swap the allowlisted env slice for the request snapshot: names present in - * the snapshot are set, names absent from it are cleared — that is what - * preserves absent-vs-empty across the IPC boundary. Only ever called inside - * the serialized render section. - */ -function applyContextEnvironment(env: InvocationContext['env']): { name: string; value: string | undefined }[] { - const saved: { name: string; value: string | undefined }[] = []; - for (const name of ENV_ALLOWLIST) { - saved.push({ name, value: process.env[name] }); - const next = env[name]; - if (next === undefined) { - unsetEnv(name); - } else { - process.env[name] = next; - } - } - return saved; +interface RenderJob { + controller: AbortController; + consumers: number; + promise: Promise; } -function restoreEnvironment(saved: { name: string; value: string | undefined }[]): void { - for (const { name, value } of saved) { - if (value === undefined) { - unsetEnv(name); - } else { - process.env[name] = value; - } - } +// How long the daemon must see no requests before provider caches are +// reclaimed (#18). Bounded caches cap steady-state growth; the idle sweep +// gives the memory back after quiet periods. +const IDLE_SWEEP_INTERVAL_MS = 60_000; +const IDLE_SWEEP_AFTER_MS = 5 * 60_000; + +function sweepProviderCaches(): void { + clearGitCache(); + clearJjCommandLog(); + clearCustomCommandCache(); + clearTranscriptAnalysisCache(); } export function createDaemonServer(options: DaemonServerOptions): DaemonServerHandle { @@ -153,57 +162,85 @@ export function createDaemonServer(options: DaemonServerOptions): DaemonServerHa const counters: DaemonCounters = { requests: 0, ok: 0 }; let lastRenderMs: number | null = null; + let lastActivityAt = Date.now(); - // Strict serialization: one render runs at a time; inFlight counts - // renders that are running or queued, and overflow answers 503. - let inFlight = 0; - let renderTail: Promise = Promise.resolve(); - let savedCwd = process.cwd(); - - function serializeRender(render: () => Promise): Promise { - const run = renderTail.then(render, render); - // Keep the tail resolved regardless of outcome; `run` itself still - // propagates errors to the requesting handler. - renderTail = run.catch(() => undefined); - return run; - } + const prefetchState: PrefetchState = createPrefetchState(); + const renderJobs = new Map(); - function renderOne(data: StatusJSON, context: InvocationContext): Promise { - return serializeRender(async () => { - const savedEnv = applyContextEnvironment(context.env); - let chdirApplied = false; - if (context.cwd !== null) { - try { - process.chdir(context.cwd); - chdirApplied = true; - } catch { - // cwd vanished between validation and render: proceed in - // the daemon cwd rather than failing the repaint. - } + const idleSweeper = setInterval(() => { + if (Date.now() - lastActivityAt > IDLE_SWEEP_AFTER_MS && renderJobs.size === 0) { + sweepProviderCaches(); + for (const cache of prefetchState.usageCaches.values()) { + cache.data = null; + cache.identity = undefined; } - try { - const loaded = await dependencies.loadSettings(getConfigPath()); - const terminalWidth = dependencies.resolveTerminalWidth( - data.session_id, - loaded.settings.terminalWidthCacheTtlSeconds - ); - const invocation = dependencies.buildInvocation(context, terminalWidth); - const { text } = await renderStatusLines(data, loaded, invocation); - return text; - } finally { - restoreEnvironment(savedEnv); - if (chdirApplied) { - try { - process.chdir(savedCwd); - } catch { - // The daemon's own cwd was removed; fall back to the - // runtime directory, which this process created. - process.chdir(runtimeDir); - savedCwd = runtimeDir; - } + } + }, IDLE_SWEEP_INTERVAL_MS); + idleSweeper.unref(); + + /** Join an in-flight render for `key`, or start one. */ + function scheduleRender(key: string, run: (signal: AbortSignal) => Promise): { promise: Promise; release: () => void } { + const existing = renderJobs.get(key); + if (existing && !existing.controller.signal.aborted) { + existing.consumers++; + counters.deduped = (counters.deduped ?? 0) + 1; + return { + promise: existing.promise, + release: () => { releaseRender(key, existing); } + }; + } + + const controller = new AbortController(); + const basePromise = run(controller.signal); + const job: RenderJob = { + controller, + consumers: 1, + promise: basePromise.finally(() => { + if (renderJobs.get(key) === job) { + renderJobs.delete(key); } - } + }) + }; + renderJobs.set(key, job); + return { + promise: job.promise, + release: () => { releaseRender(key, job); } + }; + } + + function releaseRender(key: string, job: RenderJob): void { + if (job.consumers > 0) { + job.consumers--; + } + // Zero consumers before completion: nobody will read the result, so + // the provider work backing this render is cancelled. The job stays + // registered until its promise settles (it still occupies an + // in-flight slot while finishing). + if (job.consumers === 0 && renderJobs.get(key) === job) { + job.controller.abort(); + } + } + + async function renderOne(data: StatusJSON, context: InvocationContext, signal: AbortSignal): Promise { + const env = mergeRequestEnvironment(context); + const loaded = await dependencies.loadSettings(getConfigPath()); + const terminalWidth = dependencies.resolveTerminalWidth( + data.session_id, + loaded.settings.terminalWidthCacheTtlSeconds, + env + ); + const invocation = dependencies.buildInvocation(context, terminalWidth, env); + + const prefetch = await prefetchRenderData(data, loaded.settings, { + env, + cwd: invocation.cwd, + terminalWidth, + signal, + state: prefetchState }); + + const { text } = await renderStatusLines(data, loaded, invocation, prefetch); + return text; } const server = http.createServer((request, response) => { @@ -269,6 +306,7 @@ export function createDaemonServer(options: DaemonServerOptions): DaemonServerHa async function handleRequest(request: http.IncomingMessage, response: http.ServerResponse): Promise { counters.requests++; + lastActivityAt = Date.now(); if (!isAuthorized(request)) { fail(response, 401, 'unauthorized'); @@ -290,6 +328,7 @@ export function createDaemonServer(options: DaemonServerOptions): DaemonServerHa startedAt: startedAt.toISOString(), uptimeSeconds: Math.floor((Date.now() - startedAt.getTime()) / 1000), lastRenderMs, + activeRenders: renderJobs.size, counters }); return; @@ -364,7 +403,7 @@ export function createDaemonServer(options: DaemonServerOptions): DaemonServerHa Object.assign(context, decoded.context); } if (context.cwd !== null) { - // Boundary validation before any work: the render chdirs here. + // Boundary validation before any work: the render reads it. try { if (!fs.statSync(context.cwd).isDirectory()) { throw new Error('not a directory'); @@ -375,24 +414,60 @@ export function createDaemonServer(options: DaemonServerOptions): DaemonServerHa } } - if (inFlight >= maxInFlight) { + if (renderJobs.size >= maxInFlight) { fail(response, 503, 'busy', 'render queue is full'); return; } - inFlight++; + + // Exact display context (#18): identical (config, env, cwd, width, + // payload) requests join the in-flight render and receive its text. + const env = mergeRequestEnvironment(context); + const terminalWidth = getTerminalWidth({ env }); + const renderKey = JSON.stringify([ + getConfigPath(), + Object.entries(env).filter(([name]) => (ENV_ALLOWLIST as readonly string[]).includes(name)), + context.cwd, + terminalWidth, + JSON.stringify(statusResult.data) + ]); + + let release: (() => void) | undefined; + let settled = false; + const onClientGone = () => { + // The client stopped waiting (repaint superseded, Claude Code + // restarted): release the consumer slot; the last one out cancels + // provider work nobody else needs. Node fires request 'close' on + // premature client disconnects; bun (1.3) surfaces nothing until + // the first write, so there the release happens only at settle — + // acceptable, the work was already done by then. + release?.(); + }; + request.on('close', () => { + if (!settled) { + onClientGone(); + } + }); + try { const startedAtRender = Date.now(); - const text = await renderOne(statusResult.data, context); + const scheduled = scheduleRender(renderKey, signal => renderOneWithSignal(statusResult.data, context, signal)); + release = scheduled.release; + const text = await scheduled.promise; + settled = true; lastRenderMs = Date.now() - startedAtRender; bump('ok'); sendText(response, 200, text); } catch (error) { + settled = true; fail(response, 500, 'render_failed', error instanceof Error ? error.message : String(error)); - } finally { - inFlight--; } } + /** renderOne with the request's cancellation signal wired into prefetch. */ + function renderOneWithSignal(data: StatusJSON, context: InvocationContext, signal: AbortSignal): Promise { + return renderOne(data, context, signal); + } + function writeDiscoveryFile(): void { const lines = [ '# ccstatusline daemon discovery v1 — one KEY=VALUE per line, safe to parse with a POSIX shell', @@ -436,6 +511,7 @@ export function createDaemonServer(options: DaemonServerOptions): DaemonServerHa writeDiscoveryFile(); }, async stop(): Promise { + clearInterval(idleSweeper); await new Promise((resolve) => { server.closeIdleConnections(); server.close(() => { resolve(); }); diff --git a/src/render.ts b/src/render.ts index fa255b409..ac3ad331d 100644 --- a/src/render.ts +++ b/src/render.ts @@ -1,7 +1,11 @@ import chalk from 'chalk'; -import type { SkillsMetrics } from './types'; import type { + BlockMetrics, + SkillsMetrics +} from './types'; +import type { + ClaudeStatusRenderData, RenderContext, RenderInvocation } from './types/RenderContext'; @@ -13,6 +17,7 @@ import { updateColorMap } from './utils/colors'; import { ZERO_COMPACTION_STATS } from './utils/compaction'; import type { LoadedSettings } from './utils/config'; import { saveSettingsTo } from './utils/config'; +import type { TranscriptAnalysis } from './utils/jsonl'; import { getTranscriptAnalysis } from './utils/jsonl'; import { advanceGlobalPowerlineThemeIndex, @@ -40,6 +45,21 @@ export interface RenderedStatusLines { loadError: string | null; } +/** + * Provider data gathered before the render (#18). The shared daemon prefetches + * all of it concurrently (deduped per account/repo key, cancellable) and hands + * the bundle in; one-shot and serve mode leave it unset and the render + * computes each piece itself, exactly as before. A field left `undefined` + * means "not prefetched — compute here"; an explicit null is a computed empty + * result and suppresses the fallback (e.g. the block-metrics directory walk). + */ +export interface RenderPrefetch { + transcriptAnalysis?: TranscriptAnalysis | null; + usageData?: Awaited>; + claudeStatusData?: ClaudeStatusRenderData | null; + blockMetrics?: BlockMetrics | null; +} + function hasSessionDurationInStatusJson(data: StatusJSON): boolean { const durationMs = data.cost?.total_duration_ms; return typeof durationMs === 'number' && Number.isFinite(durationMs) && durationMs >= 0; @@ -59,7 +79,8 @@ function hasSessionDurationInStatusJson(data: StatusJSON): boolean { export async function renderStatusLines( data: StatusJSON, loaded: LoadedSettings, - invocation: RenderInvocation + invocation: RenderInvocation, + prefetch: RenderPrefetch = {} ): Promise { const { settings, loadError: configError } = loaded; @@ -86,22 +107,28 @@ export async function renderStatusLines( } } - const transcriptAnalysisPromise = data.transcript_path - ? getTranscriptAnalysis(data.transcript_path, { - includeSessionDuration: hasSessionClock && !hasSessionDurationInStatusJson(data), - includeSpeedMetrics: hasSpeedItems, - includeSubagents: true, - speedWindowSeconds: Array.from(requestedSpeedWindows), - includeCompactionStats: hasCompactionWidget, - includeThinkingEffort: needsTranscriptThinkingEffort, - includeSessionName: hasSessionNameWidget, - includeLastTurnTokens: hasLastTurnTokensWidget - }) - : Promise.resolve(null); + const transcriptAnalysisPromise = prefetch.transcriptAnalysis !== undefined + ? Promise.resolve(prefetch.transcriptAnalysis) + : data.transcript_path + ? getTranscriptAnalysis(data.transcript_path, { + includeSessionDuration: hasSessionClock && !hasSessionDurationInStatusJson(data), + includeSpeedMetrics: hasSpeedItems, + includeSubagents: true, + speedWindowSeconds: Array.from(requestedSpeedWindows), + includeCompactionStats: hasCompactionWidget, + includeThinkingEffort: needsTranscriptThinkingEffort, + includeSessionName: hasSessionNameWidget, + includeLastTurnTokens: hasLastTurnTokensWidget + }) + : Promise.resolve(null); const [transcriptAnalysis, usageData, claudeStatusData] = await Promise.all([ transcriptAnalysisPromise, - prefetchUsageDataIfNeeded(lines, data), - prefetchClaudeStatusIfNeeded(lines) + prefetch.usageData !== undefined + ? Promise.resolve(prefetch.usageData) + : prefetchUsageDataIfNeeded(lines, data), + prefetch.claudeStatusData !== undefined + ? Promise.resolve(prefetch.claudeStatusData) + : prefetchClaudeStatusIfNeeded(lines) ]); // --- Start of the synchronous formatting section (chalk setup through @@ -150,7 +177,16 @@ export async function renderStatusLines( minimalist: settings.minimalistMode, gitCacheTtlSeconds: settings.gitCacheTtlSeconds, customCommandCacheTtlSeconds: settings.customCommandCacheTtlSeconds, - gitReviewNeedsChecks: lines.some(line => line.some(item => item.type === 'git-ci-status')) + gitReviewNeedsChecks: lines.some(line => line.some(item => item.type === 'git-ci-status')), + // Request env/cwd snapshots (#18): provider calls in the formatting + // section resolve through these instead of process state, so the + // daemon never swaps process.env/cwd between concurrent renders. + env: invocation.env, + cwd: invocation.cwd, + // Only meaningful when the daemon prefetch computed it: undefined + // (one-shot) lets usage widgets run the directory walk themselves, + // an explicit null suppresses it. + ...(prefetch.blockMetrics !== undefined ? { blockMetrics: prefetch.blockMetrics } : {}) }; const outputLines: string[] = []; diff --git a/src/types/RenderContext.ts b/src/types/RenderContext.ts index 94ff909d7..e4781deea 100644 --- a/src/types/RenderContext.ts +++ b/src/types/RenderContext.ts @@ -81,6 +81,16 @@ export interface RenderContext { lineIndex?: number; // Index of the current line being rendered (for theme cycling) globalSeparatorIndex?: number; // Global separator index that continues across lines + /** + * Environment and working directory snapshots for this render (#18). + * Provider calls that spawn child processes or read env-dependent config + * resolve through these instead of process.env/process.cwd, so the daemon + * can serve concurrent requests for different sessions without swapping + * process-global state. Unset in one-shot mode (process state is used). + */ + env?: NodeJS.ProcessEnv; + cwd?: string; + // For git widget thresholds gitData?: { changedFiles?: number; diff --git a/src/utils/__tests__/transcript-reuse.test.ts b/src/utils/__tests__/transcript-reuse.test.ts new file mode 100644 index 000000000..c55a91338 --- /dev/null +++ b/src/utils/__tests__/transcript-reuse.test.ts @@ -0,0 +1,108 @@ +import * as fs from 'node:fs'; +import * as os from 'node:os'; +import * as path from 'node:path'; +import { + afterEach, + beforeEach, + describe, + expect, + it +} from 'vitest'; + +import { + clearTranscriptAnalysisCache, + getTranscriptAnalysis +} from '../jsonl-metrics'; + +// Whole-analysis reuse (#18): unchanged transcripts are returned from cache by +// reference; any write to the main file or a subagent file forces a rescan. + +let transcriptDir: string; +let transcriptPath: string; + +function writeTranscript(lines: string[]): void { + fs.writeFileSync(transcriptPath, `${lines.join('\n')}\n`); +} + +function usageLine(inputTokens: number): string { + return JSON.stringify({ + timestamp: '2026-01-01T10:00:00.000Z', + message: { usage: { input_tokens: inputTokens, output_tokens: 50 } } + }); +} + +beforeEach(() => { + transcriptDir = fs.mkdtempSync(path.join(os.tmpdir(), 'ccsd-transcript-')); + transcriptPath = path.join(transcriptDir, 'session.jsonl'); + clearTranscriptAnalysisCache(); +}); + +afterEach(() => { + clearTranscriptAnalysisCache(); + fs.rmSync(transcriptDir, { recursive: true, force: true }); +}); + +describe('transcript analysis reuse (#18)', () => { + it('returns the same analysis object for an unchanged transcript', async () => { + writeTranscript([usageLine(100)]); + const first = await getTranscriptAnalysis(transcriptPath); + const second = await getTranscriptAnalysis(transcriptPath); + expect(second).toBe(first); + }); + + it('dedups concurrent scans of the same transcript into one', async () => { + writeTranscript([usageLine(100)]); + const [a, b] = await Promise.all([ + getTranscriptAnalysis(transcriptPath), + getTranscriptAnalysis(transcriptPath) + ]); + expect(b).toBe(a); + }); + + it('rescans after an append to the transcript', async () => { + writeTranscript([usageLine(100)]); + const first = await getTranscriptAnalysis(transcriptPath); + + writeTranscript([usageLine(100), usageLine(200)]); + const second = await getTranscriptAnalysis(transcriptPath); + + expect(second).not.toBe(first); + }); + + it('rescans after the transcript is truncated', async () => { + writeTranscript([usageLine(100), usageLine(200), usageLine(300)]); + const first = await getTranscriptAnalysis(transcriptPath); + + writeTranscript([usageLine(100)]); + const second = await getTranscriptAnalysis(transcriptPath); + + expect(second).not.toBe(first); + }); + + it('rescans when a subagent transcript is appended to', async () => { + writeTranscript([usageLine(100)]); + const subagentsDir = path.join(transcriptDir, 'subagents'); + fs.mkdirSync(subagentsDir); + const subagentPath = path.join(subagentsDir, 'agent-1.jsonl'); + fs.writeFileSync(subagentPath, `${usageLine(10)}\n`); + const first = await getTranscriptAnalysis(transcriptPath); + + fs.appendFileSync(subagentPath, `${usageLine(20)}\n`); + const second = await getTranscriptAnalysis(transcriptPath); + + expect(second).not.toBe(first); + }); + + it('rescans when a new subagent transcript appears', async () => { + const subagentsDir = path.join(transcriptDir, 'subagents'); + fs.mkdirSync(subagentsDir); + writeTranscript([usageLine(100)]); + fs.writeFileSync(path.join(subagentsDir, 'agent-1.jsonl'), `${usageLine(10)}\n`); + const first = await getTranscriptAnalysis(transcriptPath); + + fs.writeFileSync(path.join(subagentsDir, 'agent-2.jsonl'), `${usageLine(20)}\n`); + const second = await getTranscriptAnalysis(transcriptPath); + + expect(second).not.toBe(first); + }); +}); diff --git a/src/utils/__tests__/usage-identity.test.ts b/src/utils/__tests__/usage-identity.test.ts new file mode 100644 index 000000000..ee32adece --- /dev/null +++ b/src/utils/__tests__/usage-identity.test.ts @@ -0,0 +1,185 @@ +import * as fs from 'node:fs'; +import { createRequire } from 'node:module'; +import * as os from 'node:os'; +import * as path from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { + afterAll, + describe, + expect, + it +} from 'vitest'; + +// Identity gating for the shared daemon usage cache (#18), probed in child +// processes (same harness style as usage-fetch.test.ts): each probe is a +// fresh module registry with its own fake HOME, so the on-disk usage cache +// under ~/.cache/ccstatusline is sandboxed and the account switch scenarios +// share one in-process memory cache. + +interface IdentityProbeResult { + homedir: string; + scopeKeys: { a: string; b: string }; + requestCounts: number[]; + lastData: { error?: string; sessionUsage?: number }; +} + +const usageModulePath = fileURLToPath(new URL('../usage-fetch.ts', import.meta.url)); + +// The CJS binding: bun's ESM execFileSync binding drops its return value here. +const realExecFileSync = ( + createRequire(import.meta.url)('child_process') as { execFileSync: (file: string, args: readonly string[], options: { encoding: string; env: NodeJS.ProcessEnv }) => string } +).execFileSync; + +function createIdentityHarness() { + const tempRoot = fs.mkdtempSync(path.join(os.tmpdir(), 'ccstatusline-usage-identity-')); + const probeScriptPath = path.join(tempRoot, 'probe-usage-identity.mjs'); + let probeCounter = 0; + + const probeScript = ` +import * as fs from 'fs'; +import * as os from 'os'; +import * as path from 'path'; +import { createRequire } from 'module'; + +const require = createRequire(import.meta.url); +const https = require('https'); +const scenario = process.env.IDENTITY_SCENARIO; +let requestCount = 0; + +https.request = (...args) => { + requestCount += 1; + const callback = args.find(value => typeof value === 'function'); + const responseHandlers = new Map(); + const response = { + statusCode: 200, + headers: {}, + setEncoding() {}, + on(event, handler) { + const existing = responseHandlers.get(event) || []; + existing.push(handler); + responseHandlers.set(event, existing); + return response; + } + }; + const request = { + on() { return request; }, + destroy() {}, + end() { + if (callback) { + callback(response); + } + const body = JSON.stringify({ + five_hour: { utilization: 10, resets_at: '2030-01-01T00:00:00.000Z' }, + seven_day: { utilization: 20, resets_at: '2030-01-02T00:00:00.000Z' } + }); + for (const handler of responseHandlers.get('data') || []) { + handler(body); + } + for (const handler of responseHandlers.get('end') || []) { + handler(); + } + } + }; + return request; +}; + +const { fetchUsageData, createUsageMemoryCache, getUsageScopeKey } = await import(${JSON.stringify(usageModulePath)}); + +const ACCOUNT_A = { accessToken: 'a-access', refreshToken: 'a-refresh' }; +const ACCOUNT_B = { accessToken: 'b-access', refreshToken: 'b-refresh' }; +const emptyConfigDir = path.join(os.homedir(), 'empty-config'); +fs.mkdirSync(emptyConfigDir, { recursive: true }); +const emptyEnv = { CLAUDE_CONFIG_DIR: emptyConfigDir }; + +const sharedCache = createUsageMemoryCache(); +const run = (account) => fetchUsageData({ + requiredFields: ['sessionUsage'], + cache: sharedCache, + resolveCredentials: () => Promise.resolve(account), + env: emptyEnv +}); + +const requestCounts = []; +let lastData = null; +if (scenario === 'memory') { + const sequence = [ACCOUNT_A, ACCOUNT_A, ACCOUNT_B, ACCOUNT_A, null, null]; + for (const account of sequence) { + lastData = await run(account); + requestCounts.push(requestCount); + } +} else { + // file scenario: prime under B, then read as A with a fresh memory cache + lastData = await run(ACCOUNT_B); + requestCounts.push(requestCount); + lastData = await run(ACCOUNT_A); + requestCounts.push(requestCount); +} + +process.stdout.write(JSON.stringify({ + homedir: os.homedir(), + scopeKeys: { a: getUsageScopeKey(ACCOUNT_A), b: getUsageScopeKey(ACCOUNT_B) }, + requestCounts, + lastData: { ...(lastData.error ? { error: lastData.error } : {}), ...(lastData.sessionUsage !== undefined ? { sessionUsage: lastData.sessionUsage } : {}) } +})); +`; + + fs.writeFileSync(probeScriptPath, probeScript); + + function runProbe(scenario: 'memory' | 'file'): IdentityProbeResult { + probeCounter += 1; + const home = path.join(tempRoot, `home-${scenario}-${probeCounter}`); + fs.mkdirSync(home, { recursive: true }); + const env = Object.fromEntries(Object.entries(process.env).filter(([key]) => { + const normalizedKey = key.toUpperCase(); + return normalizedKey !== 'CLAUDE_CONFIG_DIR' && normalizedKey !== 'HTTPS_PROXY'; + })); + Object.assign(env, { HOME: home, USERPROFILE: home, PATH: '/nonexistent', IDENTITY_SCENARIO: scenario }); + + const output = realExecFileSync(process.execPath, [probeScriptPath], { encoding: 'utf8', env }); + let result: IdentityProbeResult; + try { + result = JSON.parse(output) as IdentityProbeResult; + } catch { + throw new Error(`probe output was not JSON: ${output.slice(0, 500)}`); + } + + // A probe resolving a different home has escaped its sandbox and would + // read or write the real user's ~/.cache/ccstatusline + expect(result.homedir).toBe(home); + + return result; + } + + return { + runProbe, + cleanup: (): void => { + fs.rmSync(tempRoot, { recursive: true, force: true }); + } + }; +} + +const harness = createIdentityHarness(); + +afterAll(() => { + harness.cleanup(); +}); + +describe('usage cache identity gating (#18)', () => { + it('fingerprints differ per account', () => { + const probe = harness.runProbe('memory'); + expect(probe.scopeKeys.a).not.toBe(probe.scopeKeys.b); + }); + + it('never serves another account\'s entry and refetches on account switch', () => { + const probe = harness.runProbe('memory'); + // Sequence A, A, B, A, none, none: + expect(probe.requestCounts).toEqual([1, 1, 2, 3, 3, 3]); + }); + + it('does not take the file cache of another account', () => { + const probe = harness.runProbe('file'); + // First fetch (B) hits the API; the A fetch must not read B's entry. + expect(probe.requestCounts).toEqual([1, 2]); + expect(probe.lastData.sessionUsage).toBe(10); + }); +}); diff --git a/src/utils/claude-service-status.ts b/src/utils/claude-service-status.ts index 174837ead..d4e4dfa50 100644 --- a/src/utils/claude-service-status.ts +++ b/src/utils/claude-service-status.ts @@ -256,8 +256,8 @@ function clearFailureLock(): void { } } -function getStatusPageProxyUrl(): string | null { - const proxyUrl = process.env.HTTPS_PROXY?.trim(); +function getStatusPageProxyUrl(env: NodeJS.ProcessEnv = process.env): string | null { + const proxyUrl = env.HTTPS_PROXY?.trim(); return proxyUrl?.length ? proxyUrl : null; } @@ -281,8 +281,8 @@ type StatusPageRequestFn = ( const requestStatusPage: StatusPageRequestFn = (options, onResponse) => https.request(options, onResponse); -async function getStatusPageRequestOptions(): Promise { - const proxyUrl = getStatusPageProxyUrl(); +async function getStatusPageRequestOptions(env: NodeJS.ProcessEnv): Promise { + const proxyUrl = getStatusPageProxyUrl(env); try { let agent: https.RequestOptions['agent'] | undefined; @@ -308,9 +308,11 @@ async function getStatusPageRequestOptions(): Promise { - return getStatusPageRequestOptions().then((baseOptions) => { + return getStatusPageRequestOptions(env).then((baseOptions) => { if (!baseOptions) { return null; } @@ -323,9 +325,17 @@ function fetchStatusPagePath( return; } settled = true; + if (signal) { + signal.removeEventListener('abort', onAbort); + } resolve(value); }; + const onAbort = () => { + request.destroy(); + finish(null); + }; + const requestOptions: https.RequestOptions = { ...baseOptions, path: pathName @@ -349,6 +359,13 @@ function fetchStatusPagePath( request.destroy(); finish(null); }); + if (signal) { + if (signal.aborted) { + finish(null); + return; + } + signal.addEventListener('abort', onAbort, { once: true }); + } request.end(); }); }); @@ -368,7 +385,10 @@ function isCacheUsable(cache: CachedClaudeStatus, includeIncidents: boolean): bo return !includeIncidents || cache.incidentsQueried; } -async function fetchClaudeServiceStatus(includeIncidents: boolean): Promise { +async function fetchClaudeServiceStatus( + includeIncidents: boolean, + scope: { signal?: AbortSignal; env?: NodeJS.ProcessEnv } = {} +): Promise { const nowMs = Date.now(); const nowSeconds = Math.floor(nowMs / 1000); @@ -398,8 +418,8 @@ async function fetchClaudeServiceStatus(includeIncidents: boolean): Promise { +export async function prefetchClaudeStatusIfNeeded( + lines: WidgetItem[][], + scope: { signal?: AbortSignal; env?: NodeJS.ProcessEnv } = {} +): Promise { if (!hasClaudeStatusWidgets(lines)) { return null; } - return fetchClaudeServiceStatus(claudeStatusNeedsIncidents(lines)); + return fetchClaudeServiceStatus(claudeStatusNeedsIncidents(lines), scope); } diff --git a/src/utils/claude-settings.ts b/src/utils/claude-settings.ts index 6b9ea69c1..8811b3ec7 100644 --- a/src/utils/claude-settings.ts +++ b/src/utils/claude-settings.ts @@ -88,8 +88,8 @@ function quotePathIfNeeded(filePath: string): string { * Determines the Claude config directory, checking CLAUDE_CONFIG_DIR environment variable first, * then falling back to the default ~/.claude directory. */ -export function getClaudeConfigDir(): string { - const envConfigDir = process.env.CLAUDE_CONFIG_DIR; +export function getClaudeConfigDir(env: NodeJS.ProcessEnv = process.env): string { + const envConfigDir = env.CLAUDE_CONFIG_DIR; if (envConfigDir) { try { @@ -122,9 +122,9 @@ export function getClaudeConfigDir(): string { * Claude Code stores this as a sibling of the default ~/.claude directory, but * inside CLAUDE_CONFIG_DIR when a valid config directory override is active. */ -export function getClaudeJsonPath(): string { - const configDir = getClaudeConfigDir(); - const envConfigDir = process.env.CLAUDE_CONFIG_DIR; +export function getClaudeJsonPath(env: NodeJS.ProcessEnv = process.env): string { + const configDir = getClaudeConfigDir(env); + const envConfigDir = env.CLAUDE_CONFIG_DIR; if (envConfigDir && configDir === path.resolve(envConfigDir)) { return path.join(configDir, '.claude.json'); diff --git a/src/utils/context-percentage.ts b/src/utils/context-percentage.ts index ea1b565b0..b1e71775c 100644 --- a/src/utils/context-percentage.ts +++ b/src/utils/context-percentage.ts @@ -16,10 +16,10 @@ export interface ContextPercentageMetrics { * percentage. Returns null when neither status JSON nor transcript metrics can * provide context usage. */ -export function calculateContextPercentageMetrics(context: Pick): ContextPercentageMetrics | null { +export function calculateContextPercentageMetrics(context: Pick): ContextPercentageMetrics | null { const contextWindowMetrics = getContextWindowMetrics(context.data); const modelIdentifier = getModelContextIdentifier(context.data?.model); - const contextConfig = getContextConfig(modelIdentifier, contextWindowMetrics.windowSize); + const contextConfig = getContextConfig(modelIdentifier, contextWindowMetrics.windowSize, context.env); if (contextWindowMetrics.usedPercentage !== null) { return { diff --git a/src/utils/custom-command-capture.ts b/src/utils/custom-command-capture.ts index 5d372e0ab..8da8ab25c 100644 --- a/src/utils/custom-command-capture.ts +++ b/src/utils/custom-command-capture.ts @@ -1,4 +1,7 @@ -import type { spawn } from 'child_process'; +import type { + ChildProcessWithoutNullStreams, + spawn +} from 'child_process'; import type { CustomCommandRequest, @@ -101,3 +104,140 @@ export function captureCustomCommand( } child.stdin.end(request.input); } + +export interface CaptureCustomCommandAsyncOptions { + /** Environment the command runs with; defaults to the current process env. */ + env?: NodeJS.ProcessEnv; + /** Working directory the command runs in. */ + cwd?: string; + /** Cooperative cancellation from the daemon scheduler (#18). */ + signal?: AbortSignal; +} + +/** + * In-process asynchronous twin of captureCustomCommand for the daemon (#18): + * the command is spawned directly (no helper runtime hop) with the same + * guarantees — bounded stdout, deadline kill of the whole process group, + * EPIPE tolerance, stdin payload delivery — but resolves a Promise instead of + * writing the result to stdout, so the awaiting render is never blocked. + */ +export function captureCustomCommandAsync( + spawnCommand: typeof spawn, + request: CustomCommandRequest, + maxBytes: number, + maxChars: number, + options: CaptureCustomCommandAsyncOptions = {} +): Promise { + return new Promise((resolve) => { + let child: ChildProcessWithoutNullStreams; + try { + // stdio is all-pipe below, so the streams are always present; the + // typed cast carries that through the spread-built options. + child = spawnCommand(request.command, { + shell: true, + stdio: ['pipe', 'pipe', 'ignore'], + windowsHide: true, + detached: process.platform !== 'win32', + ...(options.env !== undefined ? { env: options.env } : {}), + ...(options.cwd !== undefined ? { cwd: options.cwd } : {}) + }) as unknown as ChildProcessWithoutNullStreams; + } catch { + resolve({ status: 'failed', marker: '[Error]' }); + return; + } + + const output = Buffer.alloc(maxBytes); + let length = 0; + let finished = false; + let exited = false; + let exitMarker: string | null = null; + let timer: ReturnType | undefined; + + const killTree = (): void => { + try { + if (process.platform !== 'win32' && child.pid !== undefined) { + process.kill(-child.pid, 'SIGKILL'); + } else { + child.kill('SIGKILL'); + } + } catch { + // The command may already have exited. + } + }; + + const finish = (marker: string | null, terminate = false): void => { + if (finished) { + return; + } + finished = true; + clearTimeout(timer); + options.signal?.removeEventListener('abort', onAbort); + + if (terminate) { + killTree(); + } + + // Descendants must not keep the daemon alive via inherited pipes + // after the deadline, an overflow, or a spawn failure. + child.stdin.destroy(); + child.stdout.destroy(); + child.unref(); + resolve(marker === null + ? { status: 'ok', stdout: output.toString('utf8', 0, length).slice(0, maxChars).trim() } + : { status: 'failed', marker }); + }; + + function onAbort(): void { + finish(exited ? exitMarker : '[Error]', !exited); + } + + child.stdout.on('data', (chunk: Buffer) => { + if (finished) { + return; + } + if (chunk.length > maxBytes - length) { + finish('[Error]', true); + return; + } + chunk.copy(output, length); + length += chunk.length; + }); + child.stdout.on('error', () => { + finish('[Error]', true); + }); + child.stdin.on('error', (error: NodeJS.ErrnoException) => { + // Commands need not consume their stdin payload. + if (error.code !== 'EPIPE') { + finish('[Error]', true); + } + }); + child.on('error', (error: NodeJS.ErrnoException) => { + const marker = error.code === 'ENOENT' ? '[Cmd not found]' + : error.code === 'EACCES' ? '[Permission denied]' : '[Error]'; + finish(marker, true); + }); + child.on('exit', (code, signal) => { + exited = true; + exitMarker = signal ? `[Signal: ${signal}]` + : code === 0 ? null : typeof code === 'number' ? `[Exit: ${code}]` : '[Error]'; + }); + // 'exit' can precede the last stdout data. Drain the pipe until 'close', but + // never wait beyond the deadline for a background descendant to close it. + child.on('close', () => { + finish(exitMarker); + }); + if (request.timeoutMs > 0) { + timer = setTimeout(() => { + finish(exited ? exitMarker : '[Timeout]', !exited); + }, request.timeoutMs); + } + if (options.signal) { + if (options.signal.aborted) { + onAbort(); + return; + } + options.signal.addEventListener('abort', onAbort, { once: true }); + } + child.stdin.end(request.input); + }); +} diff --git a/src/utils/custom-command.ts b/src/utils/custom-command.ts index 41151c863..22ed70a16 100644 --- a/src/utils/custom-command.ts +++ b/src/utils/custom-command.ts @@ -1,11 +1,21 @@ import type { SpawnSyncReturns } from 'child_process'; -import { spawnSync } from 'child_process'; +import { + spawn, + spawnSync +} from 'child_process'; import { createHash } from 'node:crypto'; import * as fs from 'node:fs'; import * as os from 'node:os'; import * as path from 'node:path'; -import { captureCustomCommand } from './custom-command-capture'; +import { capMap } from '../daemon/provider-scope'; +import type { RenderContext } from '../types/RenderContext'; +import type { WidgetItem } from '../types/Widget'; + +import { + captureCustomCommand, + captureCustomCommandAsync +} from './custom-command-capture'; /** Outcome of one custom command invocation. */ export type CustomCommandResult @@ -29,6 +39,39 @@ export interface CustomCommandRequest { * waiting out the TTL. */ terminalWidth?: number | null; + /** + * Environment and working directory for the child process. The daemon + * passes the request's snapshot so concurrent renders need no + * process-global env swap (#18); one-shot mode leaves them unset and the + * child inherits the current process. + */ + env?: NodeJS.ProcessEnv; + cwd?: string; +} + +/** + * The exact request the CustomCommand widget renders with, shared with the + * daemon prefetch (#18) so both paths build identical keys and payloads. + */ +export function buildCustomCommandRequest(item: WidgetItem, context: RenderContext): CustomCommandRequest | null { + if (!item.commandPath || !context.data) { + return null; + } + const jsonInput = JSON.stringify( + typeof context.terminalWidth === 'number' + ? { ...context.data, terminal_width: context.terminalWidth } + : context.data + ); + return { + command: item.commandPath, + input: jsonInput, + timeoutMs: item.timeout ?? 1000, + ttlSeconds: context.customCommandCacheTtlSeconds, + sessionId: context.data.session_id, + terminalWidth: context.terminalWidth, + env: context.env, + cwd: context.cwd + }; } interface CustomCommandCacheEntry { @@ -57,6 +100,8 @@ const MAX_STDOUT_BYTES = 1024 * 1024; // In-process cache keeps cwd in the key. The persistent cache stores cwd once at // the file level and keys entries by command, session and terminal width. const customCommandCache = new Map(); +// Bounded for the long-lived daemon (#18); one-shot processes die anyway. +const CUSTOM_COMMAND_CACHE_MAX_ENTRIES = 256; function getCacheDir(): string { return path.join(os.homedir(), '.cache', 'ccstatusline'); @@ -267,7 +312,8 @@ function executeCommand(request: CustomCommandRequest): CustomCommandResult { // result delivery here, with a backstop if the helper fails to reply. timeout: request.timeoutMs > 0 ? request.timeoutMs + 1000 : 0, killSignal: 'SIGKILL', - env: process.env, + env: request.env ?? process.env, + ...(request.cwd !== undefined ? { cwd: request.cwd } : {}), windowsHide: true }); const marker = getFailureMarker(result); @@ -280,6 +326,29 @@ function executeCommand(request: CustomCommandRequest): CustomCommandResult { } } +/** + * Async twin of executeCommand for the daemon (#18): the capture runs + * in-process (see captureCustomCommandAsync), so no helper runtime is spawned + * and the awaiting render is never blocked by the child. + */ +async function executeCommandAsync(request: CustomCommandRequest, signal?: AbortSignal): Promise { + try { + return await captureCustomCommandAsync( + spawn, + request, + MAX_STDOUT_BYTES, + MAX_CACHED_OUTPUT_CHARS, + { + env: request.env ?? process.env, + ...(request.cwd !== undefined ? { cwd: request.cwd } : {}), + signal + } + ); + } catch { + return { status: 'failed', marker: '[Error]' }; + } +} + /** * Run a custom command, reusing a recent result when one is still within the TTL. * @@ -302,7 +371,7 @@ export function runCustomCommand(request: CustomCommandRequest): CustomCommandRe return executeCommand(request); } - const cwd = process.cwd(); + const cwd = request.cwd ?? process.cwd(); const entryKey = getEntryKey(request); const memoryCacheKey = `${entryKey}\0${cwd}`; const canShareAcrossProcesses = typeof request.sessionId === 'string' && request.sessionId.length > 0; @@ -329,6 +398,53 @@ export function runCustomCommand(request: CustomCommandRequest): CustomCommandRe createdAt: Date.now() }; customCommandCache.set(memoryCacheKey, entry); + capMap(customCommandCache, CUSTOM_COMMAND_CACHE_MAX_ENTRIES); + if (canShareAcrossProcesses) { + writePersistentCacheEntry(cwd, entryKey, entry, entry.createdAt); + } + + return result; +} + +/** + * Async twin of runCustomCommand for the daemon (#18): same cache keys and + * TTL semantics (TTL 0 still executes per request — only concurrent callers + * share one in-flight run via the daemon's refresh group), non-blocking + * execution. The cache is shared with the sync path, so a prefetched result + * is picked up by the synchronous formatter without re-running anything. + */ +export async function runCustomCommandAsync(request: CustomCommandRequest, signal?: AbortSignal): Promise { + const ttlMs = getCacheTtlMs(request.ttlSeconds); + if (ttlMs === 0) { + return executeCommandAsync(request, signal); + } + + const cwd = request.cwd ?? process.cwd(); + const entryKey = getEntryKey(request); + const memoryCacheKey = `${entryKey}\0${cwd}`; + const canShareAcrossProcesses = typeof request.sessionId === 'string' && request.sessionId.length > 0; + const now = Date.now(); + + const memoryEntry = customCommandCache.get(memoryCacheKey); + if (memoryEntry && isCacheEntryFresh(memoryEntry, ttlMs, now)) { + return memoryEntry.result; + } + + if (canShareAcrossProcesses) { + const persistentEntry = readPersistentCacheEntry(cwd, entryKey, ttlMs, now); + if (persistentEntry) { + customCommandCache.set(memoryCacheKey, persistentEntry); + return persistentEntry.result; + } + } + + const result = await executeCommandAsync(request, signal); + const entry: CustomCommandCacheEntry = { + result, + createdAt: Date.now() + }; + customCommandCache.set(memoryCacheKey, entry); + capMap(customCommandCache, CUSTOM_COMMAND_CACHE_MAX_ENTRIES); if (canShareAcrossProcesses) { writePersistentCacheEntry(cwd, entryKey, entry, entry.createdAt); } diff --git a/src/utils/git-review-cache.ts b/src/utils/git-review-cache.ts index a846b3c0f..4b71997c4 100644 --- a/src/utils/git-review-cache.ts +++ b/src/utils/git-review-cache.ts @@ -1,4 +1,6 @@ +import type { ExecFileOptionsWithStringEncoding } from 'child_process'; import { + execFile, execFileSync, spawn } from 'child_process'; @@ -15,6 +17,14 @@ import { import { createHash } from 'node:crypto'; import os from 'node:os'; import path from 'node:path'; +import { promisify } from 'util'; + +// Promise-based execFile with string decoding (see usage-fetch.ts for the shape). +const execFileAsync = promisify(execFile) as ( + file: string, + args: readonly string[], + options: ExecFileOptionsWithStringEncoding +) => Promise<{ stdout: string; stderr: string }>; import { parseRemoteUrl } from './git-remote'; @@ -755,6 +765,315 @@ export function refreshGitReviewCacheFromCli( } } +// --- Async refresh path (#18) --------------------------------------------------- +// +// The daemon prefetch refreshes the review cache directly (non-blocking child +// processes, cancellable), so the sync formatter reads a fresh cache file and +// the detached `node