From 5e7ee990cf4d9cbca6729c422055973ec37a4cdb Mon Sep 17 00:00:00 2001 From: Jan Heinsohn Date: Tue, 29 Sep 2026 13:01:45 +0200 Subject: [PATCH 1/2] feat(hoosage): cache session indexes across windows Every new window re-read all Copilot CLI session files (including WSL over 9P) and re-parsed every Chat transcript, so "Indexing saved usage" showed for seconds on each start. Keep a local index cache under globalStorageUri/cache/: - CLI/JetBrains/WSL session-state: per-file byte offset, inode, reader state and the usage read so far; later starts only read appended bytes. Line splitting is now byte-based so the resume offset is exact. - Chat transcripts: size, mtime and allowlisted usage per file; unchanged transcripts are not parsed again. - Previously imported Chat history is shown before the rescan. The cache is keyed by extension version, validated per file and falls back to a full read when missing, corrupt, truncated or replaced. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- README.md | 2 + src/core/chat-history-import.ts | 192 +++++++++++++++++++---- src/core/cli.ts | 262 +++++++++++++++++++++++++++++-- src/core/scan-cache.ts | 68 ++++++++ src/core/wsl-sessions.ts | 4 +- src/extension.ts | 54 +++++-- test/scan-cache.test.ts | 268 ++++++++++++++++++++++++++++++++ 7 files changed, 789 insertions(+), 61 deletions(-) create mode 100644 src/core/scan-cache.ts create mode 100644 test/scan-cache.test.ts diff --git a/README.md b/README.md index b0ebe0a..582181d 100644 --- a/README.md +++ b/README.md @@ -137,6 +137,8 @@ History lives under `globalStorageUri/projects//`; shared collec Saved usage is indexed in the background when the extension starts. The dashboard and status bar show an indexing state until totals are complete; large local histories are read in bounded batches without delaying extension activation. +To avoid reading every session file again in each new window, hoosage keeps a local index cache under `globalStorageUri/cache/`. For each Copilot CLI/JetBrains `session-state` folder (including WSL folders read from Windows), it stores each session file's read position, the reader state and the usage entries read so far, including the session's working directory for local project matching. For VS Code Chat transcripts, it stores each file's size, modification time and the allowlisted usage entries. Later starts only read what was appended or changed, so indexing normally finishes within a second. Previously imported Chat history is shown immediately, before transcripts are rescanned. The cache never replaces the source files: sessions that were removed, truncated or replaced are read again or dropped. The cache is rebuilt once after each hoosage update. You can delete the folder at any time; the next start then reads everything again. + VS Code 1.138 restricts these Copilot settings to application scope. Setup therefore uses user settings, without writing secrets into project files. When the configured Hoosage endpoint and `collector.json` differ, startup rebinds the configured local endpoint; foreign destinations are never adopted. Local desktop VS Code is verified; separate profiles and remote hosts must not share another host's endpoint. Telemetry environment overrides and enterprise policy conflicts are reported instead of silently redirecting them. VS Code's global telemetry-off preference is respected; it also disables Copilot's local OTel exporter upstream. ### WSL and Dev Containers diff --git a/src/core/chat-history-import.ts b/src/core/chat-history-import.ts index f2381b4..907204f 100644 --- a/src/core/chat-history-import.ts +++ b/src/core/chat-history-import.ts @@ -1,8 +1,15 @@ import { createHash } from "node:crypto"; +import type { Stats } from "node:fs"; import { readdir, readFile, stat } from "node:fs/promises"; -import { basename, dirname, join } from "node:path"; +import { basename, dirname, join, relative } from "node:path"; import { fileURLToPath } from "node:url"; import { folderPathHash } from "./cli"; +import { + cacheFile, + readCache, + writeCache, + type ScanCache, +} from "./scan-cache"; import type { Project, UsageCall } from "./types"; /** Attribution bucket for Chat sessions VS Code stored without a folder. */ @@ -231,14 +238,48 @@ export function requestUsage( }; } +type RecoveredCall = UsageCall & { dedupeKey: string }; + +interface TranscriptEntry { + size: number; + mtimeMs: number; + projectId: string; + calls: RecoveredCall[]; +} + +/** Results of transcripts parsed earlier, keyed by path below the VS Code + * user folder; unchanged files (same size and modification time) are not + * parsed again. */ +interface TranscriptIndex { + userDir: string; + previous: Map; + next: Map; + changed: boolean; +} + +const TRANSCRIPT_CACHE_FORMAT = 1; + async function sessionCalls( file: string, projectId: string, -): Promise<(UsageCall & { dedupeKey: string })[]> { + index?: TranscriptIndex, +): Promise { let session: ChatSession | undefined; + let info: Stats; + const key = index ? relative(index.userDir, file) : ""; try { - const info = await stat(file); + info = await stat(file); if (!info.isFile() || info.size > MAX_SESSION_BYTES) return []; + const cached = index?.previous.get(key); + if ( + cached && + cached.size === info.size && + cached.mtimeMs === info.mtimeMs && + cached.projectId === projectId + ) { + index!.next.set(key, cached); + return cached.calls; + } const raw = await readFile(file, "utf8"); session = file.endsWith(".jsonl") ? replayMutationLog(raw) @@ -246,6 +287,24 @@ async function sessionCalls( } catch { return []; } + const calls = requestCalls(file, session, projectId); + if (index) { + index.next.set(key, { + size: info.size, + mtimeMs: info.mtimeMs, + projectId, + calls, + }); + index.changed = true; + } + return calls; +} + +function requestCalls( + file: string, + session: ChatSession | undefined, + projectId: string, +): RecoveredCall[] { if (!session || !Array.isArray(session.requests)) return []; const sessionId = typeof session.sessionId === "string" && session.sessionId @@ -261,20 +320,68 @@ async function sessionCalls( async function directoryCalls( dir: string, projectId: string, -): Promise<(UsageCall & { dedupeKey: string })[]> { + index?: TranscriptIndex, +): Promise { let files: string[]; try { files = await readdir(dir); } catch { return []; } - const calls: (UsageCall & { dedupeKey: string })[] = []; + const calls: RecoveredCall[] = []; for (const file of files.filter((f) => /\.jsonl?$/.test(f))) - for (const call of await sessionCalls(join(dir, file), projectId)) + for (const call of await sessionCalls(join(dir, file), projectId, index)) calls.push(call); return calls; } +async function loadTranscriptIndex( + cache: ScanCache, + userDir: string, +): Promise { + const index: TranscriptIndex = { + userDir, + previous: new Map(), + next: new Map(), + changed: false, + }; + const payload = await readCache( + cacheFile(cache, "chat", userDir), + TRANSCRIPT_CACHE_FORMAT, + cache.key, + userDir, + ); + const files = payload?.files; + if (!files || typeof files !== "object" || Array.isArray(files)) return index; + for (const [key, value] of Object.entries(files)) { + const entry = value as Record | null; + if (!entry || typeof entry !== "object" || !Array.isArray(entry.calls)) + continue; + const size = count(entry.size); + const mtimeMs = entry.mtimeMs; + const projectId = label(entry.projectId); + if ( + size === undefined || + typeof mtimeMs !== "number" || + !Number.isFinite(mtimeMs) || + !projectId + ) + continue; + const calls: RecoveredCall[] = []; + for (const item of entry.calls) { + const call = restoreCall(item, projectId); + const dedupeKey = (item as { dedupeKey?: unknown } | null)?.dedupeKey; + if (!call || typeof dedupeKey !== "string" || dedupeKey.length > 2048) + break; + calls.push({ ...call, dedupeKey }); + } + // A partly invalid entry is parsed again from the transcript. + if (calls.length === entry.calls.length) + index.previous.set(key, { size, mtimeMs, projectId, calls }); + } + return index; +} + const projectId = (identity: string) => createHash("sha256").update(identity).digest("hex").slice(0, 24); @@ -334,31 +441,34 @@ export function restoreHistory(raw: unknown, projectId: string): UsageCall[] { const list = (raw as { calls?: unknown } | undefined)?.calls; if (!Array.isArray(list)) return []; return list.flatMap((item): UsageCall[] => { - if (!item || typeof item !== "object") return []; - const c = item as Record; - const id = label(c.id); - const timestamp = count(c.timestamp); - if (!id || !/^[a-f0-9]{24}$/.test(id) || timestamp === undefined) - return []; - return [ - { - id, - projectId, - timestamp, - model: label(c.model) ?? "unknown", - sessionId: label(c.sessionId), - source: "chat-history", - input: count(c.input), - output: count(c.output), - nanoAiu: count(c.nanoAiu), - requests: count(c.requests), - durationMs: count(c.durationMs), - failed: c.failed === true, - }, - ]; + const call = restoreCall(item, projectId); + return call ? [call] : []; }); } +function restoreCall(item: unknown, projectId: string): UsageCall | undefined { + if (!item || typeof item !== "object") return undefined; + const c = item as Record; + const id = label(c.id); + const timestamp = count(c.timestamp); + if (!id || !/^[a-f0-9]{24}$/.test(id) || timestamp === undefined) + return undefined; + return { + id, + projectId, + timestamp, + model: label(c.model) ?? "unknown", + sessionId: label(c.sessionId), + source: "chat-history", + input: count(c.input), + output: count(c.output), + nanoAiu: count(c.nanoAiu), + requests: count(c.requests), + durationMs: count(c.durationMs), + failed: c.failed === true, + }; +} + /** Adds newly recovered entries to the stored ones. Returns the merged list * and whether anything new was added, so callers only rewrite on change. */ export function mergeStoredHistory( @@ -398,13 +508,16 @@ export function withoutLiveOverlap( /** Scans every Chat transcript of this VS Code profile. The earliest copy of * a request wins when VS Code duplicated a session (for example after - * continuing a chat in a new session), so copied requests count once. */ + * continuing a chat in a new session), so copied requests count once. With a + * cache, only transcripts that changed since the last scan are parsed. */ export async function scanChatHistory( globalStorageFsPath: string, + cache?: ScanCache, ): Promise { const userDir = dirname(dirname(globalStorageFsPath)); const storageRoot = join(userDir, "workspaceStorage"); - const found: (UsageCall & { dedupeKey: string })[] = []; + const index = cache ? await loadTranscriptIndex(cache, userDir) : undefined; + const found: RecoveredCall[] = []; const projects = new Map(); let entries: string[] = []; try { @@ -424,7 +537,11 @@ export async function scanChatHistory( continue; } if (!project) continue; - const calls = await directoryCalls(join(dir, "chatSessions"), project.id); + const calls = await directoryCalls( + join(dir, "chatSessions"), + project.id, + index, + ); if (!calls.length) continue; let createdAt = Infinity; for (const call of calls) { @@ -440,8 +557,21 @@ export async function scanChatHistory( const noFolder = await directoryCalls( join(userDir, "globalStorage", "emptyWindowChatSessions"), NO_FOLDER_CHAT_PROJECT_ID, + index, ); for (const call of noFolder) found.push(call); + if ( + cache && + index && + (index.changed || index.next.size !== index.previous.size) + ) + await writeCache( + cacheFile(cache, "chat", userDir), + TRANSCRIPT_CACHE_FORMAT, + cache.key, + userDir, + { files: Object.fromEntries(index.next) }, + ); found.sort((a, b) => a.timestamp - b.timestamp); const seen = new Set(); const calls: UsageCall[] = []; diff --git a/src/core/cli.ts b/src/core/cli.ts index b54eba3..ad49366 100644 --- a/src/core/cli.ts +++ b/src/core/cli.ts @@ -2,12 +2,21 @@ import { createHash } from "node:crypto"; import { open, readFile, readdir, stat } from "node:fs/promises"; import { homedir } from "node:os"; import { join, resolve as resolvePath } from "node:path"; -import { StringDecoder } from "node:string_decoder"; +import { + cacheFile, + readCache, + writeCache, + type ScanCache, +} from "./scan-cache"; import type { UsageCall } from "./types"; const MAX_LINE = 2 * 1024 * 1024; const BATCH = 4 * 1024 * 1024; const MAX_SKEW = 24 * 60 * 60 * 1000; +const EMPTY = Buffer.alloc(0); +/** Bump when the cached reader state changes shape. */ +const CACHE_FORMAT = 1; +const SAVE_INTERVAL_MS = 30_000; /** Attribution bucket for CLI sessions whose cwd matches no known project. */ export const CLI_PROJECT_ID = "copilot-cli"; @@ -41,8 +50,8 @@ interface ModelSnapshot { interface FileState { offset: number; inode?: number; - decoder: StringDecoder; - pending: string; + /** Bytes after the last newline; always a partial line. */ + pending: Buffer; dropping: boolean; cwd?: string; /** Event ordinal within this file; disambiguates events lacking an id. */ @@ -62,8 +71,7 @@ interface FileState { } const freshState = (): FileState => ({ offset: 0, - decoder: new StringDecoder("utf8"), - pending: "", + pending: EMPTY, dropping: false, seq: 0, emitted: new Set(), @@ -125,6 +133,100 @@ const modelName = (value: unknown): string | undefined => ? value.replace(/[\u0000-\u001f\u007f]/g, "").slice(0, 160) || undefined : undefined; +const text = (value: unknown): string | undefined => + typeof value === "string" && + value.length > 0 && + value.length <= 4096 && + !value.includes("\0") + ? value + : undefined; + +type CachedCall = Omit< + UsageCall, + "projectId" | "sessionId" | "source" | "failed" +> & { cwd?: string }; + +/** Validates one cached file entry; undefined means "read the file again". */ +function restoreFile( + sessionId: string, + value: unknown, +): { st: FileState; calls: CachedCall[] } | undefined { + const saved = record(value); + if ( + !saved || + !/^[^\\/\0]{1,255}$/.test(sessionId) || + sessionId === "." || + sessionId === ".." + ) + return undefined; + const offset = num(saved.offset); + const seq = num(saved.seq); + if (offset === undefined || seq === undefined) return undefined; + const workspace = record(saved.workspace); + const st: FileState = { + ...freshState(), + offset, + inode: + typeof saved.inode === "number" && Number.isFinite(saved.inode) + ? saved.inode + : undefined, + dropping: saved.dropping === true, + cwd: text(saved.cwd), + seq, + workspace: workspace + ? { clientName: text(workspace.clientName), cwd: text(workspace.cwd) } + : undefined, + model: modelName(saved.model), + totalRaw: num(saved.totalRaw), + totalOffset: num(saved.totalOffset) ?? 0, + accounted: num(saved.accounted) ?? 0, + }; + for (const [model, raw] of Object.entries(record(saved.baselines) ?? {})) { + const baseline = record(raw); + if (!model || !baseline) return undefined; + st.baselines.set(model, { + input: num(baseline.input), + output: num(baseline.output), + cacheRead: num(baseline.cacheRead), + cacheWrite: num(baseline.cacheWrite), + requests: num(baseline.requests), + nanoAiu: num(baseline.nanoAiu), + }); + } + if (!Array.isArray(saved.calls)) return undefined; + const calls: CachedCall[] = []; + for (const item of saved.calls) { + const call = record(item); + const id = call?.id; + const timestamp = num(call?.timestamp); + const model = modelName(call?.model); + // Entries missing here would never be emitted again from the saved offset. + if ( + !call || + typeof id !== "string" || + id.length > 2048 || + !id.startsWith(`cli:${sessionId}:`) || + timestamp === undefined || + !model + ) + return undefined; + calls.push({ + id, + timestamp, + model, + input: num(call.input), + output: num(call.output), + cacheRead: num(call.cacheRead), + cacheWrite: num(call.cacheWrite), + requests: num(call.requests), + nanoAiu: num(call.nanoAiu), + cwd: text(call.cwd), + }); + st.emitted.add(id); + } + return { st, calls }; +} + type Resolve = (cwd: string) => string | undefined; /** Incremental reader for Copilot CLI session-state event streams @@ -138,7 +240,11 @@ type Resolve = (cwd: string) => string | undefined; * totalNanoAiu (usage checkpoints and shutdowns) is present it is * authoritative for cost: checkpoints emit cost-only increments (tokens and * requests are counted at shutdown) and each shutdown adds what is still - * missing to the current model. */ + * missing to the current model. + * + * With a cache, the reader state of each file (read position, cumulative + * baselines and emitted entries) is saved locally, so a new window only + * reads bytes appended since then. */ export class CliUsageScanner { readonly calls = new Map(); skippedLines = 0; @@ -149,20 +255,32 @@ export class CliUsageScanner { /** Local-only cwd for reattributing a session when new workspaces appear. */ private readonly callCwds = new Map(); private busy = false; + private readonly cacheFile?: string; + private restored = false; + private dirty = false; + private savedAt = -Infinity; - constructor(root?: string) { + constructor( + root?: string, + private readonly cache?: ScanCache, + ) { this.root = root ?? join( process.env.COPILOT_HOME || join(homedir(), ".copilot"), "session-state", ); + if (cache) this.cacheFile = cacheFile(cache, "cli", this.root); } async poll(resolve: Resolve): Promise { if (this.busy) return; this.busy = true; try { + if (!this.restored) { + this.restored = true; + await this.restore(); + } const rootInfo = await stat(this.root).catch((error) => { if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; @@ -176,6 +294,11 @@ export class CliUsageScanner { const entries = await readdir(this.root, { withFileTypes: true }); let caughtUp = true; const sessions = entries.filter((entry) => entry.isDirectory()); + // Usage of removed session folders disappears, as it would after a + // restart without cache. + const listed = new Set(sessions.map((entry) => entry.name)); + for (const sessionId of [...this.files.keys()]) + if (!listed.has(sessionId)) this.forget(sessionId); for (let start = 0; start < sessions.length; start += 16) { const batch = await Promise.allSettled( sessions @@ -205,12 +328,114 @@ export class CliUsageScanner { : CLI_PROJECT_ID); } } + const wasCaughtUp = this.caughtUp; this.caughtUp = caughtUp; + if ( + this.dirty && + (Date.now() - this.savedAt >= SAVE_INTERVAL_MS || + (caughtUp && !wasCaughtUp)) + ) + await this.persist(); } finally { this.busy = false; } } + private forget(sessionId: string): void { + const st = this.files.get(sessionId); + if (!st) return; + for (const id of st.emitted) { + this.calls.delete(id); + this.callCwds.delete(id); + } + this.files.delete(sessionId); + this.dirty = true; + } + + /** Restores the reader state an earlier window saved for this folder. A + * file entry that fails validation is dropped and read again. */ + private async restore(): Promise { + if (!this.cacheFile || !this.cache) return; + const payload = await readCache( + this.cacheFile, + CACHE_FORMAT, + this.cache.key, + this.root, + ); + const files = record(payload?.files); + if (!payload || !files) return; + this.skippedLines = num(payload.skippedLines) ?? 0; + for (const [sessionId, value] of Object.entries(files)) { + const restored = restoreFile(sessionId, value); + if (!restored) continue; + const { st, calls } = restored; + const source = + st.workspace?.clientName === JETBRAINS_CLIENT ? "jetbrains" : "cli"; + for (const { cwd, ...call } of calls) { + // Project attribution is resolved again at the end of each poll. + this.calls.set(call.id, { + ...call, + projectId: + source === "jetbrains" ? JETBRAINS_PROJECT_ID : CLI_PROJECT_ID, + sessionId, + source, + failed: false, + }); + if (cwd) this.callCwds.set(call.id, cwd); + } + this.files.set(sessionId, st); + } + } + + private async persist(): Promise { + if (!this.cacheFile || !this.cache) return; + const files: Record = {}; + for (const [sessionId, st] of this.files) + files[sessionId] = { + // Resume at the start of the unfinished line. + offset: st.offset - st.pending.length, + inode: st.inode, + dropping: st.dropping, + seq: st.seq, + cwd: st.cwd, + workspace: st.workspace ?? undefined, + model: st.model, + totalRaw: st.totalRaw, + totalOffset: st.totalOffset, + accounted: st.accounted, + baselines: Object.fromEntries(st.baselines), + calls: [...st.emitted].flatMap((id) => { + const call = this.calls.get(id); + return call + ? [ + { + id, + timestamp: call.timestamp, + model: call.model, + input: call.input, + output: call.output, + cacheRead: call.cacheRead, + cacheWrite: call.cacheWrite, + requests: call.requests, + nanoAiu: call.nanoAiu, + cwd: this.callCwds.get(id), + }, + ] + : []; + }), + }; + this.dirty = false; + this.savedAt = Date.now(); + const saved = await writeCache( + this.cacheFile, + CACHE_FORMAT, + this.cache.key, + this.root, + { skippedLines: this.skippedLines, files }, + ); + if (!saved) this.dirty = true; + } + /** Returns false while the file still has unread bytes. */ private async pollFile( sessionId: string, @@ -234,6 +459,7 @@ export class CliUsageScanner { this.callCwds.delete(id); } st = freshState(); + this.dirty = true; } st.inode = info.ino; this.files.set(sessionId, st); @@ -250,6 +476,7 @@ export class CliUsageScanner { }); st.workspace = content === undefined ? null : parseWorkspace(content); if (st.workspace) { + this.dirty = true; if (st.cwd === undefined && st.workspace.cwd) st.cwd = st.workspace.cwd; // Re-tag calls emitted before the file appeared. @@ -272,12 +499,8 @@ export class CliUsageScanner { const buffer = Buffer.alloc(length); const { bytesRead } = await file.read(buffer, 0, length, st.offset); st.offset += bytesRead; - this.consume( - sessionId, - st, - st.decoder.write(buffer.subarray(0, bytesRead)), - resolve, - ); + if (bytesRead) this.dirty = true; + this.consume(sessionId, st, buffer.subarray(0, bytesRead), resolve); } catch { return false; } @@ -287,14 +510,19 @@ export class CliUsageScanner { } } + /** Splits on newline bytes, so a resumable position is always the start of + * `pending` (UTF-8 never contains 0x0a inside a multi-byte character). */ private consume( sessionId: string, st: FileState, - chunk: string, + chunk: Buffer, resolve: Resolve, ): void { - const lines = (st.pending + chunk).split("\n"); - st.pending = lines.pop() ?? ""; + const data = st.pending.length ? Buffer.concat([st.pending, chunk]) : chunk; + const end = data.lastIndexOf(0x0a); + // Copy, so the pending tail does not retain the whole read buffer. + st.pending = Buffer.from(data.subarray(end + 1)); + const lines = end < 0 ? [] : data.toString("utf8", 0, end).split("\n"); for (const line of lines) { if (st.dropping) { st.dropping = false; @@ -312,7 +540,7 @@ export class CliUsageScanner { } } if (st.pending.length > MAX_LINE) { - st.pending = ""; + st.pending = EMPTY; st.dropping = true; this.skippedLines++; } diff --git a/src/core/scan-cache.ts b/src/core/scan-cache.ts new file mode 100644 index 0000000..7cd0336 --- /dev/null +++ b/src/core/scan-cache.ts @@ -0,0 +1,68 @@ +import { createHash, randomBytes } from "node:crypto"; +import { mkdir, readFile, rename, rm, writeFile } from "node:fs/promises"; +import { dirname, join } from "node:path"; + +/** Where a reader keeps its local index cache, and the value (the extension + * version) a stored cache must match to be reused. A cache only speeds up + * start-up: the source files stay authoritative, and a missing, corrupt or + * outdated cache just means reading them again. */ +export interface ScanCache { + directory: string; + key: string; +} + +/** One cache file per reader kind and source folder. */ +export const cacheFile = (cache: ScanCache, kind: string, source: string) => + join( + cache.directory, + `${kind}-${createHash("sha256").update(source).digest("hex").slice(0, 16)}.json`, + ); + +/** The stored payload, or undefined when the cache is missing, unreadable or + * was written for another format, extension version or source folder. */ +export async function readCache( + file: string, + format: number, + key: string, + source: string, +): Promise | undefined> { + try { + const data: unknown = JSON.parse(await readFile(file, "utf8")); + if (!data || typeof data !== "object" || Array.isArray(data)) + return undefined; + const stored = data as Record; + const payload = stored.payload; + return stored.format === format && + stored.key === key && + stored.source === source && + payload && + typeof payload === "object" && + !Array.isArray(payload) + ? (payload as Record) + : undefined; + } catch { + return undefined; + } +} + +/** Best-effort atomic replace, so windows sharing a cache never read a + * partial file. Returns false when the cache could not be written. */ +export async function writeCache( + file: string, + format: number, + key: string, + source: string, + payload: unknown, +): Promise { + const content = JSON.stringify({ format, key, source, payload }); + const temporary = `${file}.${randomBytes(6).toString("hex")}.tmp`; + try { + await mkdir(dirname(file), { recursive: true, mode: 0o700 }); + await writeFile(temporary, content, { mode: 0o600 }); + await rename(temporary, file); + return true; + } catch { + await rm(temporary, { force: true }).catch(() => undefined); + return false; + } +} diff --git a/src/core/wsl-sessions.ts b/src/core/wsl-sessions.ts index c72ca30..9ef5825 100644 --- a/src/core/wsl-sessions.ts +++ b/src/core/wsl-sessions.ts @@ -2,6 +2,7 @@ import { execFile } from "node:child_process"; import { readdir, stat } from "node:fs/promises"; import { join, win32 } from "node:path"; import { CliUsageScanner } from "./cli"; +import type { ScanCache } from "./scan-cache"; import type { UsageCall } from "./types"; // Reading \\wsl.localhost\ starts a stopped distribution, so the @@ -82,6 +83,7 @@ export class WslCliSessions { constructor( private readonly listDistros: () => Promise = runningWslDistros, private readonly shareRoot = "\\\\wsl.localhost", + private readonly cache?: ScanCache, ) {} get calls(): UsageCall[] { @@ -164,7 +166,7 @@ export class WslCliSessions { if (info?.isDirectory()) this.scanners.set(root, { distro, - scanner: new CliUsageScanner(root), + scanner: new CliUsageScanner(root, this.cache), }); } } diff --git a/src/extension.ts b/src/extension.ts index e50bf1d..4622609 100644 --- a/src/extension.ts +++ b/src/extension.ts @@ -54,6 +54,7 @@ import { registerWindow, routeWindow } from "./core/routing"; import { hasTelemetryEnvironmentConflict } from "./core/environment"; import { SharedUsageSync } from "./core/shared-usage-sync"; import { WslCliSessions } from "./core/wsl-sessions"; +import type { ScanCache } from "./core/scan-cache"; const KEYS = [ "enabled", @@ -161,9 +162,17 @@ export async function activate(context: vscode.ExtensionContext) { const capture = (id: string) => join(root, id, "copilot.jsonl"); const tailers = new Map(); const discovering = new Set(); - const cliScanner = new CliUsageScanner(); + // Local index caches let a new window resume where the last one stopped + // instead of re-reading every session file. Rebuilt after each update. + const scanCache: ScanCache = { + directory: join(storage, "cache"), + key: String(context.extension.packageJSON?.version ?? ""), + }; + const cliScanner = new CliUsageScanner(undefined, scanCache); const wslSessions = - process.platform === "win32" ? new WslCliSessions() : undefined; + process.platform === "win32" + ? new WslCliSessions(undefined, undefined, scanCache) + : undefined; const readWsl = () => wslSessions !== undefined && vscode.workspace @@ -312,13 +321,40 @@ export async function activate(context: vscode.ExtensionContext) { // what an open window would derive, so later live usage joins the same card. let chatHistory = new Map(); let chatHistoryComplete = Boolean(remoteName); + let chatHistoryLoaded = Boolean(remoteName); const importedHistoryFile = (id: string) => join(root, id, "chat-history.json"); + const isHistoryId = (id: string) => + /^[a-f0-9]{24}$/.test(id) || id === NO_FOLDER_CHAT_PROJECT_ID; + const storedHistory = async (id: string): Promise => { + try { + return restoreHistory( + JSON.parse(await readFile(importedHistoryFile(id), "utf8")), + id, + ); + } catch { + return []; + } + }; if (!remoteName) void (async () => { + // Entries imported by earlier windows are shown right away; the + // transcript scan below only adds what is new since then. + try { + const loaded = new Map(); + for (const id of (await readdir(root).catch(() => [])).filter( + isHistoryId, + )) { + const calls = await storedHistory(id); + if (calls.length) loaded.set(id, calls); + } + chatHistory = loaded; + } finally { + chatHistoryLoaded = true; + } const recovered = new Map(); try { - const scanned = await scanChatHistory(storage); + const scanned = await scanChatHistory(storage, scanCache); for (const project of scanned.projects) { if (project.id === current?.id) continue; discovering.add(project.id); @@ -348,15 +384,8 @@ export async function activate(context: vscode.ExtensionContext) { ids = await readdir(root); } catch {} for (const id of new Set([...ids, ...recovered.keys()])) { - if (!/^[a-f0-9]{24}$/.test(id) && id !== NO_FOLDER_CHAT_PROJECT_ID) - continue; - let stored: UsageCall[] = []; - try { - stored = restoreHistory( - JSON.parse(await readFile(importedHistoryFile(id), "utf8")), - id, - ); - } catch {} + if (!isHistoryId(id)) continue; + const stored = await storedHistory(id); const merged = mergeStoredHistory(stored, recovered.get(id) ?? []); if (merged.changed) try { @@ -617,6 +646,7 @@ export async function activate(context: vscode.ExtensionContext) { ? `Collector ready on the ${remoteLabel} host. Run Copilot Chat, then Diagnose Tracking to confirm a Chat span arrives here.` : "This project is registered automatically. Use Copilot Chat to record usage; reload if you just enabled tracking."); const indexing = + !chatHistoryLoaded || [...tailers.values()].some((t) => !t.caughtUp) || !cliScanner.caughtUp || (readWsl() && !wslSessions!.caughtUp); diff --git a/test/scan-cache.test.ts b/test/scan-cache.test.ts new file mode 100644 index 0000000..fa40bf1 --- /dev/null +++ b/test/scan-cache.test.ts @@ -0,0 +1,268 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + appendFile, + mkdir, + mkdtemp, + open, + readdir, + rm, + stat, + truncate, + utimes, + writeFile, +} from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { pathToFileURL } from "node:url"; +import { CliUsageScanner } from "../src/core/cli"; +import { scanChatHistory } from "../src/core/chat-history-import"; +import type { ScanCache } from "../src/core/scan-cache"; +import type { UsageCall } from "../src/core/types"; + +const line = (value: Record) => JSON.stringify(value) + "\n"; +const start = (cwd: string) => + line({ + type: "session.start", + id: "start", + timestamp: "2026-09-22T10:00:00.000Z", + data: { context: { cwd } }, + }); +const shutdown = (id: string, minute: number, input: number, total: number) => + line({ + type: "session.shutdown", + id, + timestamp: `2026-09-22T10:${String(minute).padStart(2, "0")}:00.000Z`, + data: { + currentModel: "gpt-5.5", + totalNanoAiu: total, + modelMetrics: { + "gpt-5.5": { + usage: { inputTokens: input, outputTokens: input / 10 }, + requests: { count: input / 100 }, + totalNanoAiu: total, + }, + }, + }, + }); + +// Stable comparison independent of key order and undefined fields. +const view = (calls: Iterable) => + [...calls] + .map((call) => + JSON.stringify( + Object.fromEntries( + Object.entries(call) + .filter(([, v]) => v !== undefined) + .sort(([a], [b]) => a.localeCompare(b)), + ), + ), + ) + .sort(); + +async function fixture() { + const dir = await mkdtemp(join(tmpdir(), "hoosage-cache-")); + const root = join(dir, "session-state"); + const cache: ScanCache = { directory: join(dir, "cache"), key: "1.0.0" }; + const project = join(dir, "repo"); + const resolve = (cwd: string) => (cwd === project ? "proj" : undefined); + const session = async (id: string, content: string) => { + await mkdir(join(root, id), { recursive: true }); + const file = join(root, id, "events.jsonl"); + await writeFile(file, content); + return file; + }; + return { dir, root, cache, project, resolve, session }; +} + +/** Overwrites bytes in place: same file, same size, so only a reader that + * resumes from its saved position keeps the original entries. */ +async function scramble(file: string, length: number) { + const handle = await open(file, "r+"); + try { + await handle.write(Buffer.alloc(length, 0x78), 0, length, 0); + } finally { + await handle.close(); + } +} + +test("a new CLI reader resumes from the cache instead of re-reading", async () => { + const f = await fixture(); + try { + const file = await f.session( + "s1", + start(f.project) + shutdown("e1", 5, 1000, 300_000_000), + ); + const first = new CliUsageScanner(f.root, f.cache); + await first.poll(f.resolve); + assert.equal(first.calls.size, 1); + assert.equal((await readdir(f.cache.directory)).length, 1); + + const size = (await stat(file)).size; + await scramble(file, size - 1); + await appendFile(file, shutdown("e2", 9, 3000, 700_000_000)); + + const second = new CliUsageScanner(f.root, f.cache); + await second.poll(f.resolve); + assert.equal(second.caughtUp, true); + assert.equal(second.skippedLines, 0); + const calls = [...second.calls.values()]; + assert.deepEqual( + calls.map((c) => [c.id, c.projectId, c.input, c.nanoAiu]), + [ + ["cli:s1:e1:gpt-5.5", "proj", 1000, 300_000_000], + ["cli:s1:e2:gpt-5.5", "proj", 2000, 400_000_000], + ], + ); + } finally { + await rm(f.dir, { recursive: true, force: true }); + } +}); + +test("cached CLI results match a full read, including split UTF-8 lines", async () => { + const f = await fixture(); + try { + const renamed = line({ + type: "session.model_change", + id: "m1", + timestamp: "2026-09-22T10:01:00.000Z", + data: { newModel: "gpt-5.5 “größer”" }, + }); + const bytes = Buffer.from(renamed); + // Cut inside a multi-byte character of an unfinished line. + const cut = bytes.indexOf(Buffer.from("“")) + 1; + const file = await f.session( + "s1", + start(f.project) + shutdown("e1", 5, 1000, 300_000_000), + ); + await appendFile(file, bytes.subarray(0, cut)); + const first = new CliUsageScanner(f.root, f.cache); + await first.poll(f.resolve); + + await appendFile(file, bytes.subarray(cut)); + await appendFile( + file, + line({ + type: "session.usage_checkpoint", + id: "c1", + timestamp: "2026-09-22T10:06:00.000Z", + data: { totalNanoAiu: 500_000_000 }, + }), + ); + const resumed = new CliUsageScanner(f.root, f.cache); + await resumed.poll(f.resolve); + const full = new CliUsageScanner(f.root); + await full.poll(f.resolve); + assert.deepEqual(view(resumed.calls.values()), view(full.calls.values())); + assert.equal( + resumed.calls.get("cli:s1:c1:gpt-5.5 “größer”")?.nanoAiu, + 200_000_000, + ); + } finally { + await rm(f.dir, { recursive: true, force: true }); + } +}); + +test("the CLI cache is ignored for another version and after truncation", async () => { + const f = await fixture(); + try { + const file = await f.session( + "s1", + start(f.project) + shutdown("e1", 5, 1000, 300_000_000), + ); + const first = new CliUsageScanner(f.root, f.cache); + await first.poll(f.resolve); + + await scramble(file, 20); + const updated = new CliUsageScanner(f.root, { ...f.cache, key: "1.0.1" }); + await updated.poll(f.resolve); + assert.equal(updated.skippedLines, 1, "re-read the scrambled file"); + + await truncate(file, 0); + await appendFile(file, start(f.project) + shutdown("e9", 7, 500, 1)); + const truncated = new CliUsageScanner(f.root, { + ...f.cache, + key: "1.0.1", + }); + await truncated.poll(f.resolve); + assert.deepEqual([...truncated.calls.keys()], ["cli:s1:e9:gpt-5.5"]); + } finally { + await rm(f.dir, { recursive: true, force: true }); + } +}); + +test("removed sessions and corrupt caches do not leave stale usage", async () => { + const f = await fixture(); + try { + await f.session("s1", start(f.project) + shutdown("e1", 5, 1000, 1)); + await f.session("s2", start(f.project) + shutdown("e2", 6, 1000, 1)); + const scanner = new CliUsageScanner(f.root, f.cache); + await scanner.poll(f.resolve); + assert.equal(scanner.calls.size, 2); + + await rm(join(f.root, "s2"), { recursive: true }); + const restored = new CliUsageScanner(f.root, f.cache); + await restored.poll(f.resolve); + assert.deepEqual([...restored.calls.keys()], ["cli:s1:e1:gpt-5.5"]); + await scanner.poll(f.resolve); + assert.deepEqual([...scanner.calls.keys()], ["cli:s1:e1:gpt-5.5"]); + + const [name] = await readdir(f.cache.directory); + await writeFile(join(f.cache.directory, name!), "{not json"); + const fresh = new CliUsageScanner(f.root, f.cache); + await fresh.poll(f.resolve); + assert.deepEqual([...fresh.calls.keys()], ["cli:s1:e1:gpt-5.5"]); + } finally { + await rm(f.dir, { recursive: true, force: true }); + } +}); + +test("unchanged Chat transcripts are served from the cache", async () => { + const dir = await mkdtemp(join(tmpdir(), "hoosage-cache-chat-")); + try { + const user = join(dir, "User"); + const globalStorage = join(user, "globalStorage", "openhoo.hoosage"); + const sessions = join(user, "workspaceStorage", "w1", "chatSessions"); + await mkdir(sessions, { recursive: true }); + await mkdir(globalStorage, { recursive: true }); + await writeFile( + join(user, "workspaceStorage", "w1", "workspace.json"), + JSON.stringify({ folder: pathToFileURL(join(dir, "repo")).toString() }), + ); + const transcript = (output: number) => + JSON.stringify({ + sessionId: "s1", + requests: [ + { + requestId: "r1", + timestamp: 1000, + modelId: "copilot/gpt-5.5", + completionTokens: output, + result: { metadata: { toolCallRounds: [{}] } }, + }, + ], + }); + const file = join(sessions, "s1.json"); + await writeFile(file, transcript(100)); + await utimes(file, 1_700_000_000, 1_700_000_000); + const cache: ScanCache = { + directory: join(dir, "cache"), + key: "1.0.0", + }; + const first = await scanChatHistory(globalStorage, cache); + assert.equal(first.calls[0]!.output, 100); + + // Same size and modification time: the cached result is used. + await writeFile(file, transcript(200)); + await utimes(file, 1_700_000_000, 1_700_000_000); + const cached = await scanChatHistory(globalStorage, cache); + assert.deepEqual(cached, first); + + await writeFile(file, transcript(3000)); + const changed = await scanChatHistory(globalStorage, cache); + assert.equal(changed.calls[0]!.output, 3000); + assert.equal((await scanChatHistory(globalStorage)).calls[0]!.output, 3000); + } finally { + await rm(dir, { recursive: true, force: true }); + } +}); From f84e91678d57cd045b63448db9224034d2b72006 Mon Sep 17 00:00:00 2001 From: Jan Heinsohn Date: Tue, 29 Sep 2026 13:06:50 +0200 Subject: [PATCH 2/2] fix(hoosage): check and read transcripts through one file handle Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- src/core/chat-history-import.ts | 11 ++++++++--- test/scan-cache.test.ts | 8 ++++---- 2 files changed, 12 insertions(+), 7 deletions(-) diff --git a/src/core/chat-history-import.ts b/src/core/chat-history-import.ts index 907204f..ab4da8d 100644 --- a/src/core/chat-history-import.ts +++ b/src/core/chat-history-import.ts @@ -1,6 +1,6 @@ import { createHash } from "node:crypto"; import type { Stats } from "node:fs"; -import { readdir, readFile, stat } from "node:fs/promises"; +import { open, readdir, readFile, type FileHandle } from "node:fs/promises"; import { basename, dirname, join, relative } from "node:path"; import { fileURLToPath } from "node:url"; import { folderPathHash } from "./cli"; @@ -266,9 +266,12 @@ async function sessionCalls( ): Promise { let session: ChatSession | undefined; let info: Stats; + let handle: FileHandle | undefined; const key = index ? relative(index.userDir, file) : ""; try { - info = await stat(file); + // Checks and reads the same open file, so a replaced path cannot slip in. + handle = await open(file, "r"); + info = await handle.stat(); if (!info.isFile() || info.size > MAX_SESSION_BYTES) return []; const cached = index?.previous.get(key); if ( @@ -280,12 +283,14 @@ async function sessionCalls( index!.next.set(key, cached); return cached.calls; } - const raw = await readFile(file, "utf8"); + const raw = await handle.readFile("utf8"); session = file.endsWith(".jsonl") ? replayMutationLog(raw) : (JSON.parse(raw) as ChatSession); } catch { return []; + } finally { + await handle?.close().catch(() => {}); } const calls = requestCalls(file, session, projectId); if (index) { diff --git a/test/scan-cache.test.ts b/test/scan-cache.test.ts index fa40bf1..b035539 100644 --- a/test/scan-cache.test.ts +++ b/test/scan-cache.test.ts @@ -7,7 +7,6 @@ import { open, readdir, rm, - stat, truncate, utimes, writeFile, @@ -76,10 +75,12 @@ async function fixture() { } /** Overwrites bytes in place: same file, same size, so only a reader that - * resumes from its saved position keeps the original entries. */ + * resumes from its saved position keeps the original entries. A negative + * length keeps that many bytes at the end. */ async function scramble(file: string, length: number) { const handle = await open(file, "r+"); try { + if (length < 0) length += (await handle.stat()).size; await handle.write(Buffer.alloc(length, 0x78), 0, length, 0); } finally { await handle.close(); @@ -98,8 +99,7 @@ test("a new CLI reader resumes from the cache instead of re-reading", async () = assert.equal(first.calls.size, 1); assert.equal((await readdir(f.cache.directory)).length, 1); - const size = (await stat(file)).size; - await scramble(file, size - 1); + await scramble(file, -1); await appendFile(file, shutdown("e2", 9, 3000, 700_000_000)); const second = new CliUsageScanner(f.root, f.cache);