From 0c1d18c3acf2135b27ed824d2d4085d9319f611f Mon Sep 17 00:00:00 2001 From: gubin-dev Date: Tue, 29 Sep 2026 00:03:03 +0300 Subject: [PATCH 1/6] refactor(code-index): isolate watcher sessions and preserve batch outcomes --- .../code-index-watcher-session.spec.ts | 164 ++++++++++++++++++ .../code-index/__tests__/orchestrator.spec.ts | 33 ++++ .../code-index/code-index-watcher-session.ts | 111 ++++++++++++ src/services/code-index/orchestrator.ts | 57 ++---- 4 files changed, 319 insertions(+), 46 deletions(-) create mode 100644 src/services/code-index/__tests__/code-index-watcher-session.spec.ts create mode 100644 src/services/code-index/code-index-watcher-session.ts diff --git a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts new file mode 100644 index 0000000000..fa0b27001c --- /dev/null +++ b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts @@ -0,0 +1,164 @@ +import type { Event } from "vscode" +import type { BatchProcessingSummary, IFileWatcher } from "../interfaces" +import { CodeIndexStateManager } from "../state-manager" +import { WatcherSession } from "../code-index-watcher-session" + +vi.mock("vscode", async () => { + const { makeEventEmitter } = await import("../../../test-utils/vscode") + return { + EventEmitter: vi.fn().mockImplementation(function () { + return makeEventEmitter() + }), + } +}) + +function eventSource() { + const listeners = new Set<(value: T) => void>() + const dispose = vi.fn(() => listeners.clear()) + const event: Event = (listener) => { + listeners.add(listener) + return { dispose } + } + return { event: vi.fn(event), fire: (value: T) => listeners.forEach((listener) => listener(value)), dispose } +} + +function setup() { + const start = eventSource() + const progress = eventSource<{ processedInBatch: number; totalInBatch: number; currentFile?: string }>() + const finish = eventSource() + const watcher = { + initialize: vi.fn<() => Promise>().mockResolvedValue(undefined), + dispose: vi.fn(), + processFile: vi.fn(), + onDidStartBatchProcessing: start.event, + onBatchProgressUpdate: progress.event, + onDidFinishBatchProcessing: finish.event, + } satisfies IFileWatcher + const state = new CodeIndexStateManager() + const session = new WatcherSession(watcher, state) + return { start, progress, finish, watcher, state, session } +} + +describe("WatcherSession", () => { + it("does not initialize after being stopped before startup", async () => { + const { session, watcher } = setup() + session.stop() + await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) + expect(watcher.initialize).not.toHaveBeenCalled() + expect(watcher.dispose).toHaveBeenCalledTimes(1) + }) + + it("reuses pending and active sessions without duplicating subscriptions", async () => { + const { session, watcher, start } = setup() + const pending = session.start() + expect(session.start()).toBe(pending) + await pending + await session.start() + expect(watcher.initialize).toHaveBeenCalledTimes(1) + expect(start.event).toHaveBeenCalledTimes(1) + }) + + it.each(["initialize", "progress", "finish"])("cleans partial startup when %s fails", async (stage) => { + const { session, watcher, start, progress, finish } = setup() + const error = new Error("startup failed") + if (stage === "initialize") watcher.initialize.mockRejectedValue(error) + else + (stage === "progress" ? progress : finish).event.mockImplementation(() => { + throw error + }) + await expect(session.start()).rejects.toBe(error) + expect(watcher.dispose).toHaveBeenCalledTimes(1) + if (stage !== "initialize") expect(start.dispose).toHaveBeenCalledTimes(1) + if (stage === "finish") expect(progress.dispose).toHaveBeenCalledTimes(1) + session.stop() + expect(watcher.dispose).toHaveBeenCalledTimes(1) + }) + + it("does not revive initialization stopped while pending", async () => { + const { session, watcher, start } = setup() + let resolve!: () => void + watcher.initialize.mockReturnValue( + new Promise((done) => { + resolve = done + }), + ) + const pending = session.start() + session.stop() + resolve() + await expect(pending).rejects.toMatchObject({ name: "AbortError" }) + expect(start.event).not.toHaveBeenCalled() + expect(watcher.dispose).toHaveBeenCalledTimes(2) + await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) + expect(watcher.initialize).toHaveBeenCalledTimes(1) + }) + + it("unsubscribes once and ignores retained callbacks after stop", async () => { + const { session, watcher, start, progress, finish, state } = setup() + await session.start() + const callback = finish.event.mock.calls[0][0] + session.stop() + session.stop() + callback({ processedFiles: [], batchError: new Error("late") }) + expect(state.state).toBe("Standby") + for (const source of [start, progress, finish]) expect(source.dispose).toHaveBeenCalledTimes(1) + expect(watcher.dispose).toHaveBeenCalledTimes(1) + await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) + }) + + describe("real state manager integration", () => { + it.each(["success", "skipped", "error", "local_error"] as const)( + "preserves the %s outcome through terminal and empty progress", + async (status) => { + const { session, start, progress, finish, state } = setup() + await session.start() + start.fire(["/workspace/file.ts"]) + progress.fire({ processedInBatch: 0, totalInBatch: 1, currentFile: "/workspace/file.ts" }) + expect(state.getCurrentStatus()).toMatchObject({ systemStatus: "Indexing", currentItemUnit: "files" }) + expect(state.getCurrentStatus().message).toContain("Current: file.ts") + progress.fire({ processedInBatch: 1, totalInBatch: 1 }) + expect(state.state).toBe("Indexing") + finish.fire({ processedFiles: [{ path: "/workspace/file.ts", status }] }) + const outcome = state.getCurrentStatus() + expect(outcome.systemStatus).toBe(status === "error" || status === "local_error" ? "Error" : "Indexed") + progress.fire({ processedInBatch: 1, totalInBatch: 1 }) + progress.fire({ processedInBatch: 0, totalInBatch: 0 }) + expect(state.getCurrentStatus()).toEqual(outcome) + }, + ) + + it("reports batch errors, then allows a subsequent successful batch to recover", async () => { + const { session, start, finish, state } = setup() + await session.start() + finish.fire({ processedFiles: [], batchError: new Error("database unavailable") }) + expect(state.state).toBe("Error") + expect(state.getCurrentStatus().message).toContain("database unavailable") + start.fire(["next.ts"]) + expect(state.state).toBe("Indexing") + finish.fire({ processedFiles: [{ path: "next.ts", status: "success" }] }) + expect(state.state).toBe("Indexed") + }) + + it("reports file errors even in a mixed successful batch", async () => { + const { session, finish, state } = setup() + await session.start() + finish.fire({ + processedFiles: [ + { path: "good.ts", status: "success" }, + { path: "bad.ts", status: "local_error", error: new Error("parse failed") }, + ], + }) + expect(state.state).toBe("Error") + expect(state.getCurrentStatus().message).toContain("parse failed") + }) + + it("does not override Stopping with any watcher event", async () => { + const { session, start, progress, finish, state } = setup() + await session.start() + state.setSystemState("Stopping", "Stopped") + start.fire(["file.ts"]) + progress.fire({ processedInBatch: 0, totalInBatch: 1 }) + finish.fire({ processedFiles: [] }) + expect(state.state).toBe("Stopping") + }) + }) +}) diff --git a/src/services/code-index/__tests__/orchestrator.spec.ts b/src/services/code-index/__tests__/orchestrator.spec.ts index bad3d69be3..58f292ec82 100644 --- a/src/services/code-index/__tests__/orchestrator.spec.ts +++ b/src/services/code-index/__tests__/orchestrator.spec.ts @@ -346,6 +346,39 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { } }) + it("does not publish Indexed when stopped during watcher initialization", async () => { + let finishInitialization!: () => void + let enteredInitialization!: () => void + const entered = new Promise((resolve) => { + enteredInitialization = resolve + }) + fileWatcher.initialize.mockImplementation(() => { + enteredInitialization() + return new Promise((resolve) => { + finishInitialization = resolve + }) + }) + scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + const indexing = orchestrator.startIndexing() + await entered + orchestrator.stopIndexing() + finishInitialization() + await indexing + expect(stateManager.state).toBe("Standby") + expect(stateManager.setSystemState).not.toHaveBeenCalledWith("Indexed", expect.anything()) + expect(fileWatcher.onDidStartBatchProcessing).not.toHaveBeenCalled() + expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() + }) + it("should abort indexing when stopIndexing() is called", async () => { // Make scanner hang until aborted scanner.scanDirectory.mockImplementation( diff --git a/src/services/code-index/code-index-watcher-session.ts b/src/services/code-index/code-index-watcher-session.ts new file mode 100644 index 0000000000..18b1a27072 --- /dev/null +++ b/src/services/code-index/code-index-watcher-session.ts @@ -0,0 +1,111 @@ +import * as path from "path" +import type { Disposable, Event } from "vscode" +import type { BatchProcessingSummary, IFileWatcher } from "./interfaces" +import type { CodeIndexStateManager } from "./state-manager" + +interface Session { + stopped: boolean + subscriptions: Disposable[] + ready: Promise +} + +/** Owns watcher startup and subscriptions; only batch summaries publish final outcomes. */ +export class WatcherSession { + private session?: Session + private stopped = false + + constructor( + private readonly watcher: IFileWatcher, + private readonly stateManager: Pick< + CodeIndexStateManager, + "state" | "setSystemState" | "reportFileQueueProgress" + >, + ) {} + + start(): Promise { + // A disposed watcher owns disposed event emitters and cannot be restarted. + if (this.stopped || this.session?.stopped) + return Promise.reject(new DOMException("Watcher session stopped", "AbortError")) + if (this.session) return this.session.ready + + const session: Session = { stopped: false, subscriptions: [], ready: Promise.resolve() } + this.session = session + session.ready = this.initialize(session) + return session.ready + } + + stop(): void { + if (this.stopped || this.session?.stopped) return + this.stopped = true + if (this.session) { + this.session.stopped = true + this.disposeSubscriptions(this.session) + } + this.watcher.dispose() + } + + private disposeSubscriptions(session: Session): void { + for (const subscription of session.subscriptions.splice(0)) subscription.dispose() + } + + private async initialize(session: Session): Promise { + try { + await this.watcher.initialize() + if (session.stopped) throw new DOMException("Watcher startup stopped", "AbortError") + + this.subscribeToWatcher(session) + } catch (error) { + session.stopped = true + this.disposeSubscriptions(session) + // Initialization may have allocated resources after stop() disposed the watcher. + this.watcher.dispose() + throw error + } + } + + private subscribeToWatcher(session: Session): void { + this.subscribe(session, this.watcher.onDidStartBatchProcessing, (files) => this.handleBatchStarted(files)) + this.subscribe(session, this.watcher.onBatchProgressUpdate, ({ processedInBatch, totalInBatch, currentFile }) => + this.handleBatchProgress(processedInBatch, totalInBatch, currentFile), + ) + this.subscribe(session, this.watcher.onDidFinishBatchProcessing, (summary) => this.handleBatchFinished(summary)) + } + + private subscribe(session: Session, event: Event, handler: (value: T) => void): void { + // Record each subscription immediately so a later registration failure cannot leak it. + const subscription = event((value) => { + if (session.stopped || this.stateManager.state === "Stopping") return + handler(value) + }) + session.subscriptions.push(subscription) + } + + private handleBatchStarted(files: string[]): void { + if (files.length === 0) return + this.stateManager.setSystemState("Indexing", "Processing file changes...") + } + + private handleBatchProgress(processed: number, total: number, currentFile?: string): void { + // Terminal progress can follow the summary; reporting it would erase the final outcome. + if (total === 0 || processed >= total) return + this.stateManager.reportFileQueueProgress( + processed, + total, + currentFile ? path.basename(currentFile) : undefined, + ) + } + + private handleBatchFinished(summary: BatchProcessingSummary): void { + const errors = summary.processedFiles.filter((file) => file.status === "error" || file.status === "local_error") + if (!summary.batchError && errors.length === 0) { + this.stateManager.setSystemState("Indexed", "File changes processed. Index up-to-date.") + return + } + + const detail = summary.batchError?.message ?? errors.find((file) => file.error)?.error?.message + this.stateManager.setSystemState( + "Error", + `File changes failed (${errors.length} file errors).${detail ? ` ${detail}` : ""}`, + ) + } +} diff --git a/src/services/code-index/orchestrator.ts b/src/services/code-index/orchestrator.ts index c055fff370..bf25ac1f1d 100644 --- a/src/services/code-index/orchestrator.ts +++ b/src/services/code-index/orchestrator.ts @@ -1,8 +1,8 @@ import * as vscode from "vscode" -import * as path from "path" import { CodeIndexConfigManager } from "./config-manager" import { CodeIndexStateManager, IndexingState } from "./state-manager" -import { IFileWatcher, IVectorStore, BatchProcessingSummary } from "./interfaces" +import { IFileWatcher, IVectorStore } from "./interfaces" +import { WatcherSession } from "./code-index-watcher-session" import { DirectoryScanner } from "./processors" import { CacheManager } from "./cache-manager" import { CodeIndexScanExecutor } from "./code-index-scan-executor" @@ -14,7 +14,7 @@ import { t } from "../../i18n" * Manages the code indexing workflow, coordinating between different services and managers. */ export class CodeIndexOrchestrator { - private _fileWatcherSubscriptions: vscode.Disposable[] = [] + private readonly watcherSession: WatcherSession private _isProcessing: boolean = false private _abortController: AbortController | null = null private readonly scanExecutor: CodeIndexScanExecutor @@ -26,9 +26,10 @@ export class CodeIndexOrchestrator { private readonly cacheManager: CacheManager, private readonly vectorStore: IVectorStore, scanner: DirectoryScanner, - private readonly fileWatcher: IFileWatcher, + fileWatcher: IFileWatcher, ) { this.scanExecutor = new CodeIndexScanExecutor(workspacePath, scanner, vectorStore, stateManager) + this.watcherSession = new WatcherSession(fileWatcher, stateManager) } /** @@ -42,45 +43,7 @@ export class CodeIndexOrchestrator { this.stateManager.setSystemState("Indexing", "Initializing file watcher...") try { - await this.fileWatcher.initialize() - - this._fileWatcherSubscriptions = [ - this.fileWatcher.onDidStartBatchProcessing((filePaths: string[]) => {}), - this.fileWatcher.onBatchProgressUpdate(({ processedInBatch, totalInBatch, currentFile }) => { - if (totalInBatch > 0 && this.stateManager.state !== "Indexing") { - this.stateManager.setSystemState("Indexing", "Processing file changes...") - } - this.stateManager.reportFileQueueProgress( - processedInBatch, - totalInBatch, - currentFile ? path.basename(currentFile) : undefined, - ) - if (processedInBatch === totalInBatch) { - // Covers (N/N) and (0/0) - if (totalInBatch > 0) { - // Batch with items completed - this.stateManager.setSystemState("Indexed", "File changes processed. Index up-to-date.") - } else { - if (this.stateManager.state === "Indexing") { - // Only transition if it was "Indexing" - this.stateManager.setSystemState("Indexed", "Index up-to-date. File queue empty.") - } - } - } - }), - this.fileWatcher.onDidFinishBatchProcessing((summary: BatchProcessingSummary) => { - if (summary.batchError) { - console.error(`[CodeIndexOrchestrator] Batch processing failed:`, summary.batchError) - } else { - const successCount = summary.processedFiles.filter( - (f: { status: string }) => f.status === "success", - ).length - const errorCount = summary.processedFiles.filter( - (f: { status: string }) => f.status === "error" || f.status === "local_error", - ).length - } - }), - ] + await this.watcherSession.start() } catch (error) { console.error("[CodeIndexOrchestrator] Failed to start file watcher:", error) TelemetryService.instance.captureEvent(TelemetryEventName.CODE_INDEX_ERROR, { @@ -157,9 +120,11 @@ export class CodeIndexOrchestrator { } await this._startWatcher() + signal.throwIfAborted() // Mark indexing as complete after successful incremental scan await this.vectorStore.markIndexingComplete() + signal.throwIfAborted() this.stateManager.setSystemState("Indexed", t("embeddings:orchestrator.fileWatcherStarted")) } else { @@ -171,9 +136,11 @@ export class CodeIndexOrchestrator { } await this._startWatcher() + signal.throwIfAborted() // Mark indexing as complete after successful full scan await this.vectorStore.markIndexingComplete() + signal.throwIfAborted() this.stateManager.setSystemState("Indexed", t("embeddings:orchestrator.fileWatcherStarted")) } @@ -250,9 +217,7 @@ export class CodeIndexOrchestrator { * Stops the file watcher and cleans up resources. */ public stopWatcher(): void { - this.fileWatcher.dispose() - this._fileWatcherSubscriptions.forEach((sub) => sub.dispose()) - this._fileWatcherSubscriptions = [] + this.watcherSession.stop() if (this.stateManager.state !== "Error" && this.stateManager.state !== "Stopping") { this.stateManager.setSystemState("Standby", t("embeddings:orchestrator.fileWatcherStopped")) From 93a37c181878bd1b39656289a77a8e2f7cf4f422 Mon Sep 17 00:00:00 2001 From: gubin-dev Date: Tue, 29 Sep 2026 00:24:40 +0300 Subject: [PATCH 2/6] fix(code-index): recreate watcher sessions on restart --- .../code-index-watcher-session.spec.ts | 66 ++++++++++++++++--- .../code-index/__tests__/manager.spec.ts | 14 ++-- .../code-index/__tests__/orchestrator.spec.ts | 59 ++++++++++++++--- .../code-index/code-index-watcher-session.ts | 45 +++++++------ src/services/code-index/manager.ts | 4 +- src/services/code-index/orchestrator.ts | 4 +- .../__tests__/processor-factories.spec.ts | 7 +- src/services/code-index/service-factory.ts | 23 +++---- 8 files changed, 162 insertions(+), 60 deletions(-) diff --git a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts index fa0b27001c..b22108d65d 100644 --- a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts +++ b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts @@ -35,17 +35,18 @@ function setup() { onDidFinishBatchProcessing: finish.event, } satisfies IFileWatcher const state = new CodeIndexStateManager() - const session = new WatcherSession(watcher, state) - return { start, progress, finish, watcher, state, session } + const createWatcher = vi.fn(() => watcher) + const session = new WatcherSession(createWatcher, state) + return { start, progress, finish, watcher, state, session, createWatcher } } describe("WatcherSession", () => { - it("does not initialize after being stopped before startup", async () => { + it("allows startup after stopping an idle owner without allocating resources", async () => { const { session, watcher } = setup() session.stop() - await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) - expect(watcher.initialize).not.toHaveBeenCalled() - expect(watcher.dispose).toHaveBeenCalledTimes(1) + expect(watcher.dispose).not.toHaveBeenCalled() + await session.start() + expect(watcher.initialize).toHaveBeenCalledTimes(1) }) it("reuses pending and active sessions without duplicating subscriptions", async () => { @@ -88,7 +89,6 @@ describe("WatcherSession", () => { await expect(pending).rejects.toMatchObject({ name: "AbortError" }) expect(start.event).not.toHaveBeenCalled() expect(watcher.dispose).toHaveBeenCalledTimes(2) - await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) expect(watcher.initialize).toHaveBeenCalledTimes(1) }) @@ -102,7 +102,57 @@ describe("WatcherSession", () => { expect(state.state).toBe("Standby") for (const source of [start, progress, finish]) expect(source.dispose).toHaveBeenCalledTimes(1) expect(watcher.dispose).toHaveBeenCalledTimes(1) - await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) + }) + + it.each(["resolve", "reject"])("keeps the replacement session when stopped startup later %ss", async (outcome) => { + const first = setup() + const next = setup() + let resolve!: () => void + let reject!: (error: Error) => void + first.watcher.initialize.mockReturnValue( + new Promise((done, fail) => { + resolve = done + reject = fail + }), + ) + first.createWatcher.mockReturnValueOnce(first.watcher).mockReturnValue(next.watcher) + const pending = first.session.start() + const rejected = expect(pending).rejects.toBeInstanceOf(Error) + first.session.stop() + await first.session.start() + if (outcome === "resolve") resolve() + else reject(new Error("late initialization failure")) + await rejected + await first.session.start() + expect(first.createWatcher).toHaveBeenCalledTimes(2) + expect(next.watcher.initialize).toHaveBeenCalledTimes(1) + expect(next.watcher.dispose).not.toHaveBeenCalled() + next.progress.fire({ processedInBatch: 0, totalInBatch: 1 }) + expect(first.state.state).toBe("Indexing") + next.finish.fire({ processedFiles: [{ path: "next.ts", status: "success" }] }) + expect(first.state.state).toBe("Indexed") + }) + + it.each(["stop", "failure"])("creates a working replacement after %s", async (reason) => { + const first = setup() + const next = setup() + first.createWatcher.mockReturnValueOnce(first.watcher).mockReturnValue(next.watcher) + if (reason === "failure") { + first.watcher.initialize.mockRejectedValue(new Error("startup failed")) + await expect(first.session.start()).rejects.toThrow("startup failed") + } else { + await first.session.start() + first.session.stop() + } + await first.session.start() + expect(first.watcher.dispose).toHaveBeenCalledTimes(1) + next.start.fire(["next.ts"]) + next.progress.fire({ processedInBatch: 0, totalInBatch: 1 }) + expect(first.state.state).toBe("Indexing") + first.finish.fire({ processedFiles: [], batchError: new Error("stale") }) + expect(first.state.state).toBe("Indexing") + next.finish.fire({ processedFiles: [], batchError: new Error("new failure") }) + expect(first.state.getCurrentStatus().message).toContain("new failure") }) describe("real state manager integration", () => { diff --git a/src/services/code-index/__tests__/manager.spec.ts b/src/services/code-index/__tests__/manager.spec.ts index 952e205fdd..5cf9e8c365 100644 --- a/src/services/code-index/__tests__/manager.spec.ts +++ b/src/services/code-index/__tests__/manager.spec.ts @@ -304,13 +304,13 @@ describe("CodeIndexManager - handleSettingsChange regression", () => { embedder: { embedderInfo: { name: "openai" } }, vectorStore: {}, scanner: {}, - fileWatcher: { + createWatcher: () => ({ onDidStartBatchProcessing: vi.fn(), onBatchProgressUpdate: vi.fn(), watch: vi.fn(), stopWatcher: vi.fn(), dispose: vi.fn(), - }, + }), }), validateEmbedder: vi.fn().mockResolvedValue({ valid: true }), } @@ -380,13 +380,13 @@ describe("CodeIndexManager - handleSettingsChange regression", () => { embedder: { embedderInfo: { name: "openai" } }, vectorStore: {}, scanner: {}, - fileWatcher: { + createWatcher: () => ({ onDidStartBatchProcessing: vi.fn(), onBatchProgressUpdate: vi.fn(), watch: vi.fn(), stopWatcher: vi.fn(), dispose: vi.fn(), - }, + }), }), validateEmbedder: vi.fn().mockResolvedValue({ valid: true }), } @@ -442,7 +442,7 @@ describe("CodeIndexManager - handleSettingsChange regression", () => { embedder: mockEmbedder, vectorStore: mockVectorStore, scanner: mockScanner, - fileWatcher: mockFileWatcher, + createWatcher: () => mockFileWatcher, }), validateEmbedder: vi.fn(), } @@ -637,13 +637,13 @@ describe("CodeIndexManager - handleSettingsChange regression", () => { embedder: { embedderInfo: { name: "openai" } }, vectorStore: {}, scanner: {}, - fileWatcher: { + createWatcher: () => ({ onDidStartBatchProcessing: vi.fn(), onBatchProgressUpdate: vi.fn(), watch: vi.fn(), stopWatcher: vi.fn(), dispose: vi.fn(), - }, + }), }), validateEmbedder: vi.fn().mockResolvedValue({ valid: true }), } diff --git a/src/services/code-index/__tests__/orchestrator.spec.ts b/src/services/code-index/__tests__/orchestrator.spec.ts index 58f292ec82..3b80b14c7b 100644 --- a/src/services/code-index/__tests__/orchestrator.spec.ts +++ b/src/services/code-index/__tests__/orchestrator.spec.ts @@ -133,7 +133,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) await orchestrator.startIndexing() expect(events).toEqual(["incomplete", "scan", "watcher", "complete"]) @@ -155,7 +155,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) scanner.scanDirectory.mockImplementation(async () => { orchestrator.stopIndexing() @@ -184,7 +184,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) // Act @@ -213,7 +213,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) // Act @@ -249,7 +249,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) await orchestrator.startIndexing() @@ -279,7 +279,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) await orchestrator.startIndexing() @@ -346,6 +346,45 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { } }) + it.each(["stop", "clear", "error"])("restarts the same orchestrator after %s", async (reason) => { + const createWatcher = vi.fn(() => ({ + initialize: vi.fn<() => Promise>().mockResolvedValue(undefined), + dispose: vi.fn(), + processFile: vi.fn(), + onDidStartBatchProcessing: vi.fn().mockReturnValue({ dispose: vi.fn() }), + onBatchProgressUpdate: vi.fn().mockReturnValue({ dispose: vi.fn() }), + onDidFinishBatchProcessing: vi.fn().mockReturnValue({ dispose: vi.fn() }), + })) + scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) + vectorStore.deleteCollection = vi.fn().mockResolvedValue(undefined) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + createWatcher, + ) + if (reason === "error") vectorStore.initialize.mockRejectedValueOnce(new Error("Qdrant unavailable")) + await orchestrator.startIndexing() + if (reason === "error") { + expect(orchestrator.state).toBe("Error") + expect(createWatcher).not.toHaveBeenCalled() + } else { + expect(orchestrator.state).toBe("Indexed") + if (reason === "clear") await orchestrator.clearIndexData() + else orchestrator.stopIndexing() + expect(createWatcher.mock.results[0].value.dispose).toHaveBeenCalledTimes(1) + } + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Indexed") + expect(createWatcher).toHaveBeenCalledTimes(reason === "error" ? 1 : 2) + for (const { value: watcher } of createWatcher.mock.results) { + expect(watcher.initialize).toHaveBeenCalledTimes(1) + } + }) + it("does not publish Indexed when stopped during watcher initialization", async () => { let finishInitialization!: () => void let enteredInitialization!: () => void @@ -366,7 +405,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) const indexing = orchestrator.startIndexing() await entered @@ -402,7 +441,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) // Start indexing (async, don't await) @@ -445,7 +484,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) const indexingPromise = orchestrator.startIndexing() @@ -483,7 +522,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + () => fileWatcher, ) const indexingPromise = orchestrator.startIndexing() diff --git a/src/services/code-index/code-index-watcher-session.ts b/src/services/code-index/code-index-watcher-session.ts index 18b1a27072..78392678ed 100644 --- a/src/services/code-index/code-index-watcher-session.ts +++ b/src/services/code-index/code-index-watcher-session.ts @@ -4,6 +4,7 @@ import type { BatchProcessingSummary, IFileWatcher } from "./interfaces" import type { CodeIndexStateManager } from "./state-manager" interface Session { + watcher: IFileWatcher stopped: boolean subscriptions: Disposable[] ready: Promise @@ -12,10 +13,9 @@ interface Session { /** Owns watcher startup and subscriptions; only batch summaries publish final outcomes. */ export class WatcherSession { private session?: Session - private stopped = false constructor( - private readonly watcher: IFileWatcher, + private readonly createWatcher: () => IFileWatcher, private readonly stateManager: Pick< CodeIndexStateManager, "state" | "setSystemState" | "reportFileQueueProgress" @@ -23,25 +23,26 @@ export class WatcherSession { ) {} start(): Promise { - // A disposed watcher owns disposed event emitters and cannot be restarted. - if (this.stopped || this.session?.stopped) - return Promise.reject(new DOMException("Watcher session stopped", "AbortError")) if (this.session) return this.session.ready - const session: Session = { stopped: false, subscriptions: [], ready: Promise.resolve() } + const session: Session = { + watcher: this.createWatcher(), + stopped: false, + subscriptions: [], + ready: Promise.resolve(), + } this.session = session session.ready = this.initialize(session) return session.ready } stop(): void { - if (this.stopped || this.session?.stopped) return - this.stopped = true - if (this.session) { - this.session.stopped = true - this.disposeSubscriptions(this.session) - } - this.watcher.dispose() + const session = this.session + if (!session) return + this.session = undefined + session.stopped = true + this.disposeSubscriptions(session) + session.watcher.dispose() } private disposeSubscriptions(session: Session): void { @@ -50,25 +51,31 @@ export class WatcherSession { private async initialize(session: Session): Promise { try { - await this.watcher.initialize() + await session.watcher.initialize() if (session.stopped) throw new DOMException("Watcher startup stopped", "AbortError") this.subscribeToWatcher(session) } catch (error) { + if (this.session === session) this.session = undefined session.stopped = true this.disposeSubscriptions(session) // Initialization may have allocated resources after stop() disposed the watcher. - this.watcher.dispose() + session.watcher.dispose() throw error } } private subscribeToWatcher(session: Session): void { - this.subscribe(session, this.watcher.onDidStartBatchProcessing, (files) => this.handleBatchStarted(files)) - this.subscribe(session, this.watcher.onBatchProgressUpdate, ({ processedInBatch, totalInBatch, currentFile }) => - this.handleBatchProgress(processedInBatch, totalInBatch, currentFile), + this.subscribe(session, session.watcher.onDidStartBatchProcessing, (files) => this.handleBatchStarted(files)) + this.subscribe( + session, + session.watcher.onBatchProgressUpdate, + ({ processedInBatch, totalInBatch, currentFile }) => + this.handleBatchProgress(processedInBatch, totalInBatch, currentFile), + ) + this.subscribe(session, session.watcher.onDidFinishBatchProcessing, (summary) => + this.handleBatchFinished(summary), ) - this.subscribe(session, this.watcher.onDidFinishBatchProcessing, (summary) => this.handleBatchFinished(summary)) } private subscribe(session: Session, event: Event, handler: (value: T) => void): void { diff --git a/src/services/code-index/manager.ts b/src/services/code-index/manager.ts index ee181b17d6..57e1561f38 100644 --- a/src/services/code-index/manager.ts +++ b/src/services/code-index/manager.ts @@ -404,7 +404,7 @@ export class CodeIndexManager { await rooIgnoreController.initialize() // (Re)Create shared service instances - const { embedder, vectorStore, scanner, fileWatcher } = this._serviceFactory.createServices( + const { embedder, vectorStore, scanner, createWatcher } = this._serviceFactory.createServices( this.context, this._cacheManager!, ignoreInstance, @@ -427,7 +427,7 @@ export class CodeIndexManager { this._cacheManager!, vectorStore, scanner, - fileWatcher, + createWatcher, ) // (Re)Initialize search service diff --git a/src/services/code-index/orchestrator.ts b/src/services/code-index/orchestrator.ts index bf25ac1f1d..f66241e10d 100644 --- a/src/services/code-index/orchestrator.ts +++ b/src/services/code-index/orchestrator.ts @@ -26,10 +26,10 @@ export class CodeIndexOrchestrator { private readonly cacheManager: CacheManager, private readonly vectorStore: IVectorStore, scanner: DirectoryScanner, - fileWatcher: IFileWatcher, + createWatcher: () => IFileWatcher, ) { this.scanExecutor = new CodeIndexScanExecutor(workspacePath, scanner, vectorStore, stateManager) - this.watcherSession = new WatcherSession(fileWatcher, stateManager) + this.watcherSession = new WatcherSession(createWatcher, stateManager) } /** diff --git a/src/services/code-index/processors/__tests__/processor-factories.spec.ts b/src/services/code-index/processors/__tests__/processor-factories.spec.ts index 4af92c14cd..c3632aebf3 100644 --- a/src/services/code-index/processors/__tests__/processor-factories.spec.ts +++ b/src/services/code-index/processors/__tests__/processor-factories.spec.ts @@ -159,8 +159,13 @@ describe("processor factories", () => { vectorStore: options.vectorStore, parser: codeParser, scanner: vi.mocked(DirectoryScanner).mock.instances[0], - fileWatcher: vi.mocked(FileWatcher).mock.instances[0], + createWatcher: expect.any(Function), }) + expect(FileWatcher).not.toHaveBeenCalled() + const first = services.createWatcher() + const second = services.createWatcher() + expect(first).not.toBe(second) + expect(FileWatcher).toHaveBeenCalledTimes(2) expect(vi.mocked(DirectoryScanner).mock.calls[0][3]).toBe(options.cacheManager) expect(vi.mocked(FileWatcher).mock.calls[0][2]).toBe(watcherCache) expect(vi.mocked(FileWatcher).mock.calls[0][0]).toBe(options.workspacePath) diff --git a/src/services/code-index/service-factory.ts b/src/services/code-index/service-factory.ts index d2bbe43f2a..e4b81348c0 100644 --- a/src/services/code-index/service-factory.ts +++ b/src/services/code-index/service-factory.ts @@ -59,7 +59,7 @@ export class CodeIndexServiceFactory { vectorStore: IVectorStore parser: ICodeParser scanner: DirectoryScanner - fileWatcher: IFileWatcher + createWatcher: () => IFileWatcher } { if (!this.configManager.isFeatureConfigured) { throw new Error(t("embeddings:serviceFactory.codeIndexingNotConfigured")) @@ -76,22 +76,23 @@ export class CodeIndexServiceFactory { cacheManager: this.cacheManager, ignoreInstance, }) - const fileWatcher = this.fileWatcherFactory.create({ - workspacePath: this.workspacePath, - context, - embedder, - vectorStore, - cacheManager, - ignoreInstance, - rooIgnoreController, - }) + const createWatcher = () => + this.fileWatcherFactory.create({ + workspacePath: this.workspacePath, + context, + embedder, + vectorStore, + cacheManager, + ignoreInstance, + rooIgnoreController, + }) return { embedder, vectorStore, parser, scanner, - fileWatcher, + createWatcher, } } } From d2a34b0066d7cd155c031b4480799116d64ba6a8 Mon Sep 17 00:00:00 2001 From: gubin-dev Date: Tue, 29 Sep 2026 00:32:09 +0300 Subject: [PATCH 3/6] Revert "fix(code-index): recreate watcher sessions on restart" This reverts commit 93a37c181878bd1b39656289a77a8e2f7cf4f422. --- .../code-index-watcher-session.spec.ts | 66 +++---------------- .../code-index/__tests__/manager.spec.ts | 14 ++-- .../code-index/__tests__/orchestrator.spec.ts | 59 +++-------------- .../code-index/code-index-watcher-session.ts | 45 ++++++------- src/services/code-index/manager.ts | 4 +- src/services/code-index/orchestrator.ts | 4 +- .../__tests__/processor-factories.spec.ts | 7 +- src/services/code-index/service-factory.ts | 23 ++++--- 8 files changed, 60 insertions(+), 162 deletions(-) diff --git a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts index b22108d65d..fa0b27001c 100644 --- a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts +++ b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts @@ -35,18 +35,17 @@ function setup() { onDidFinishBatchProcessing: finish.event, } satisfies IFileWatcher const state = new CodeIndexStateManager() - const createWatcher = vi.fn(() => watcher) - const session = new WatcherSession(createWatcher, state) - return { start, progress, finish, watcher, state, session, createWatcher } + const session = new WatcherSession(watcher, state) + return { start, progress, finish, watcher, state, session } } describe("WatcherSession", () => { - it("allows startup after stopping an idle owner without allocating resources", async () => { + it("does not initialize after being stopped before startup", async () => { const { session, watcher } = setup() session.stop() - expect(watcher.dispose).not.toHaveBeenCalled() - await session.start() - expect(watcher.initialize).toHaveBeenCalledTimes(1) + await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) + expect(watcher.initialize).not.toHaveBeenCalled() + expect(watcher.dispose).toHaveBeenCalledTimes(1) }) it("reuses pending and active sessions without duplicating subscriptions", async () => { @@ -89,6 +88,7 @@ describe("WatcherSession", () => { await expect(pending).rejects.toMatchObject({ name: "AbortError" }) expect(start.event).not.toHaveBeenCalled() expect(watcher.dispose).toHaveBeenCalledTimes(2) + await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) expect(watcher.initialize).toHaveBeenCalledTimes(1) }) @@ -102,57 +102,7 @@ describe("WatcherSession", () => { expect(state.state).toBe("Standby") for (const source of [start, progress, finish]) expect(source.dispose).toHaveBeenCalledTimes(1) expect(watcher.dispose).toHaveBeenCalledTimes(1) - }) - - it.each(["resolve", "reject"])("keeps the replacement session when stopped startup later %ss", async (outcome) => { - const first = setup() - const next = setup() - let resolve!: () => void - let reject!: (error: Error) => void - first.watcher.initialize.mockReturnValue( - new Promise((done, fail) => { - resolve = done - reject = fail - }), - ) - first.createWatcher.mockReturnValueOnce(first.watcher).mockReturnValue(next.watcher) - const pending = first.session.start() - const rejected = expect(pending).rejects.toBeInstanceOf(Error) - first.session.stop() - await first.session.start() - if (outcome === "resolve") resolve() - else reject(new Error("late initialization failure")) - await rejected - await first.session.start() - expect(first.createWatcher).toHaveBeenCalledTimes(2) - expect(next.watcher.initialize).toHaveBeenCalledTimes(1) - expect(next.watcher.dispose).not.toHaveBeenCalled() - next.progress.fire({ processedInBatch: 0, totalInBatch: 1 }) - expect(first.state.state).toBe("Indexing") - next.finish.fire({ processedFiles: [{ path: "next.ts", status: "success" }] }) - expect(first.state.state).toBe("Indexed") - }) - - it.each(["stop", "failure"])("creates a working replacement after %s", async (reason) => { - const first = setup() - const next = setup() - first.createWatcher.mockReturnValueOnce(first.watcher).mockReturnValue(next.watcher) - if (reason === "failure") { - first.watcher.initialize.mockRejectedValue(new Error("startup failed")) - await expect(first.session.start()).rejects.toThrow("startup failed") - } else { - await first.session.start() - first.session.stop() - } - await first.session.start() - expect(first.watcher.dispose).toHaveBeenCalledTimes(1) - next.start.fire(["next.ts"]) - next.progress.fire({ processedInBatch: 0, totalInBatch: 1 }) - expect(first.state.state).toBe("Indexing") - first.finish.fire({ processedFiles: [], batchError: new Error("stale") }) - expect(first.state.state).toBe("Indexing") - next.finish.fire({ processedFiles: [], batchError: new Error("new failure") }) - expect(first.state.getCurrentStatus().message).toContain("new failure") + await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) }) describe("real state manager integration", () => { diff --git a/src/services/code-index/__tests__/manager.spec.ts b/src/services/code-index/__tests__/manager.spec.ts index 5cf9e8c365..952e205fdd 100644 --- a/src/services/code-index/__tests__/manager.spec.ts +++ b/src/services/code-index/__tests__/manager.spec.ts @@ -304,13 +304,13 @@ describe("CodeIndexManager - handleSettingsChange regression", () => { embedder: { embedderInfo: { name: "openai" } }, vectorStore: {}, scanner: {}, - createWatcher: () => ({ + fileWatcher: { onDidStartBatchProcessing: vi.fn(), onBatchProgressUpdate: vi.fn(), watch: vi.fn(), stopWatcher: vi.fn(), dispose: vi.fn(), - }), + }, }), validateEmbedder: vi.fn().mockResolvedValue({ valid: true }), } @@ -380,13 +380,13 @@ describe("CodeIndexManager - handleSettingsChange regression", () => { embedder: { embedderInfo: { name: "openai" } }, vectorStore: {}, scanner: {}, - createWatcher: () => ({ + fileWatcher: { onDidStartBatchProcessing: vi.fn(), onBatchProgressUpdate: vi.fn(), watch: vi.fn(), stopWatcher: vi.fn(), dispose: vi.fn(), - }), + }, }), validateEmbedder: vi.fn().mockResolvedValue({ valid: true }), } @@ -442,7 +442,7 @@ describe("CodeIndexManager - handleSettingsChange regression", () => { embedder: mockEmbedder, vectorStore: mockVectorStore, scanner: mockScanner, - createWatcher: () => mockFileWatcher, + fileWatcher: mockFileWatcher, }), validateEmbedder: vi.fn(), } @@ -637,13 +637,13 @@ describe("CodeIndexManager - handleSettingsChange regression", () => { embedder: { embedderInfo: { name: "openai" } }, vectorStore: {}, scanner: {}, - createWatcher: () => ({ + fileWatcher: { onDidStartBatchProcessing: vi.fn(), onBatchProgressUpdate: vi.fn(), watch: vi.fn(), stopWatcher: vi.fn(), dispose: vi.fn(), - }), + }, }), validateEmbedder: vi.fn().mockResolvedValue({ valid: true }), } diff --git a/src/services/code-index/__tests__/orchestrator.spec.ts b/src/services/code-index/__tests__/orchestrator.spec.ts index 3b80b14c7b..58f292ec82 100644 --- a/src/services/code-index/__tests__/orchestrator.spec.ts +++ b/src/services/code-index/__tests__/orchestrator.spec.ts @@ -133,7 +133,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) await orchestrator.startIndexing() expect(events).toEqual(["incomplete", "scan", "watcher", "complete"]) @@ -155,7 +155,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) scanner.scanDirectory.mockImplementation(async () => { orchestrator.stopIndexing() @@ -184,7 +184,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) // Act @@ -213,7 +213,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) // Act @@ -249,7 +249,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) await orchestrator.startIndexing() @@ -279,7 +279,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) await orchestrator.startIndexing() @@ -346,45 +346,6 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { } }) - it.each(["stop", "clear", "error"])("restarts the same orchestrator after %s", async (reason) => { - const createWatcher = vi.fn(() => ({ - initialize: vi.fn<() => Promise>().mockResolvedValue(undefined), - dispose: vi.fn(), - processFile: vi.fn(), - onDidStartBatchProcessing: vi.fn().mockReturnValue({ dispose: vi.fn() }), - onBatchProgressUpdate: vi.fn().mockReturnValue({ dispose: vi.fn() }), - onDidFinishBatchProcessing: vi.fn().mockReturnValue({ dispose: vi.fn() }), - })) - scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) - vectorStore.deleteCollection = vi.fn().mockResolvedValue(undefined) - const orchestrator = new CodeIndexOrchestrator( - configManager, - stateManager, - workspacePath, - cacheManager, - vectorStore, - scanner, - createWatcher, - ) - if (reason === "error") vectorStore.initialize.mockRejectedValueOnce(new Error("Qdrant unavailable")) - await orchestrator.startIndexing() - if (reason === "error") { - expect(orchestrator.state).toBe("Error") - expect(createWatcher).not.toHaveBeenCalled() - } else { - expect(orchestrator.state).toBe("Indexed") - if (reason === "clear") await orchestrator.clearIndexData() - else orchestrator.stopIndexing() - expect(createWatcher.mock.results[0].value.dispose).toHaveBeenCalledTimes(1) - } - await orchestrator.startIndexing() - expect(orchestrator.state).toBe("Indexed") - expect(createWatcher).toHaveBeenCalledTimes(reason === "error" ? 1 : 2) - for (const { value: watcher } of createWatcher.mock.results) { - expect(watcher.initialize).toHaveBeenCalledTimes(1) - } - }) - it("does not publish Indexed when stopped during watcher initialization", async () => { let finishInitialization!: () => void let enteredInitialization!: () => void @@ -405,7 +366,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) const indexing = orchestrator.startIndexing() await entered @@ -441,7 +402,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) // Start indexing (async, don't await) @@ -484,7 +445,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) const indexingPromise = orchestrator.startIndexing() @@ -522,7 +483,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - () => fileWatcher, + fileWatcher, ) const indexingPromise = orchestrator.startIndexing() diff --git a/src/services/code-index/code-index-watcher-session.ts b/src/services/code-index/code-index-watcher-session.ts index 78392678ed..18b1a27072 100644 --- a/src/services/code-index/code-index-watcher-session.ts +++ b/src/services/code-index/code-index-watcher-session.ts @@ -4,7 +4,6 @@ import type { BatchProcessingSummary, IFileWatcher } from "./interfaces" import type { CodeIndexStateManager } from "./state-manager" interface Session { - watcher: IFileWatcher stopped: boolean subscriptions: Disposable[] ready: Promise @@ -13,9 +12,10 @@ interface Session { /** Owns watcher startup and subscriptions; only batch summaries publish final outcomes. */ export class WatcherSession { private session?: Session + private stopped = false constructor( - private readonly createWatcher: () => IFileWatcher, + private readonly watcher: IFileWatcher, private readonly stateManager: Pick< CodeIndexStateManager, "state" | "setSystemState" | "reportFileQueueProgress" @@ -23,26 +23,25 @@ export class WatcherSession { ) {} start(): Promise { + // A disposed watcher owns disposed event emitters and cannot be restarted. + if (this.stopped || this.session?.stopped) + return Promise.reject(new DOMException("Watcher session stopped", "AbortError")) if (this.session) return this.session.ready - const session: Session = { - watcher: this.createWatcher(), - stopped: false, - subscriptions: [], - ready: Promise.resolve(), - } + const session: Session = { stopped: false, subscriptions: [], ready: Promise.resolve() } this.session = session session.ready = this.initialize(session) return session.ready } stop(): void { - const session = this.session - if (!session) return - this.session = undefined - session.stopped = true - this.disposeSubscriptions(session) - session.watcher.dispose() + if (this.stopped || this.session?.stopped) return + this.stopped = true + if (this.session) { + this.session.stopped = true + this.disposeSubscriptions(this.session) + } + this.watcher.dispose() } private disposeSubscriptions(session: Session): void { @@ -51,31 +50,25 @@ export class WatcherSession { private async initialize(session: Session): Promise { try { - await session.watcher.initialize() + await this.watcher.initialize() if (session.stopped) throw new DOMException("Watcher startup stopped", "AbortError") this.subscribeToWatcher(session) } catch (error) { - if (this.session === session) this.session = undefined session.stopped = true this.disposeSubscriptions(session) // Initialization may have allocated resources after stop() disposed the watcher. - session.watcher.dispose() + this.watcher.dispose() throw error } } private subscribeToWatcher(session: Session): void { - this.subscribe(session, session.watcher.onDidStartBatchProcessing, (files) => this.handleBatchStarted(files)) - this.subscribe( - session, - session.watcher.onBatchProgressUpdate, - ({ processedInBatch, totalInBatch, currentFile }) => - this.handleBatchProgress(processedInBatch, totalInBatch, currentFile), - ) - this.subscribe(session, session.watcher.onDidFinishBatchProcessing, (summary) => - this.handleBatchFinished(summary), + this.subscribe(session, this.watcher.onDidStartBatchProcessing, (files) => this.handleBatchStarted(files)) + this.subscribe(session, this.watcher.onBatchProgressUpdate, ({ processedInBatch, totalInBatch, currentFile }) => + this.handleBatchProgress(processedInBatch, totalInBatch, currentFile), ) + this.subscribe(session, this.watcher.onDidFinishBatchProcessing, (summary) => this.handleBatchFinished(summary)) } private subscribe(session: Session, event: Event, handler: (value: T) => void): void { diff --git a/src/services/code-index/manager.ts b/src/services/code-index/manager.ts index 57e1561f38..ee181b17d6 100644 --- a/src/services/code-index/manager.ts +++ b/src/services/code-index/manager.ts @@ -404,7 +404,7 @@ export class CodeIndexManager { await rooIgnoreController.initialize() // (Re)Create shared service instances - const { embedder, vectorStore, scanner, createWatcher } = this._serviceFactory.createServices( + const { embedder, vectorStore, scanner, fileWatcher } = this._serviceFactory.createServices( this.context, this._cacheManager!, ignoreInstance, @@ -427,7 +427,7 @@ export class CodeIndexManager { this._cacheManager!, vectorStore, scanner, - createWatcher, + fileWatcher, ) // (Re)Initialize search service diff --git a/src/services/code-index/orchestrator.ts b/src/services/code-index/orchestrator.ts index f66241e10d..bf25ac1f1d 100644 --- a/src/services/code-index/orchestrator.ts +++ b/src/services/code-index/orchestrator.ts @@ -26,10 +26,10 @@ export class CodeIndexOrchestrator { private readonly cacheManager: CacheManager, private readonly vectorStore: IVectorStore, scanner: DirectoryScanner, - createWatcher: () => IFileWatcher, + fileWatcher: IFileWatcher, ) { this.scanExecutor = new CodeIndexScanExecutor(workspacePath, scanner, vectorStore, stateManager) - this.watcherSession = new WatcherSession(createWatcher, stateManager) + this.watcherSession = new WatcherSession(fileWatcher, stateManager) } /** diff --git a/src/services/code-index/processors/__tests__/processor-factories.spec.ts b/src/services/code-index/processors/__tests__/processor-factories.spec.ts index c3632aebf3..4af92c14cd 100644 --- a/src/services/code-index/processors/__tests__/processor-factories.spec.ts +++ b/src/services/code-index/processors/__tests__/processor-factories.spec.ts @@ -159,13 +159,8 @@ describe("processor factories", () => { vectorStore: options.vectorStore, parser: codeParser, scanner: vi.mocked(DirectoryScanner).mock.instances[0], - createWatcher: expect.any(Function), + fileWatcher: vi.mocked(FileWatcher).mock.instances[0], }) - expect(FileWatcher).not.toHaveBeenCalled() - const first = services.createWatcher() - const second = services.createWatcher() - expect(first).not.toBe(second) - expect(FileWatcher).toHaveBeenCalledTimes(2) expect(vi.mocked(DirectoryScanner).mock.calls[0][3]).toBe(options.cacheManager) expect(vi.mocked(FileWatcher).mock.calls[0][2]).toBe(watcherCache) expect(vi.mocked(FileWatcher).mock.calls[0][0]).toBe(options.workspacePath) diff --git a/src/services/code-index/service-factory.ts b/src/services/code-index/service-factory.ts index e4b81348c0..d2bbe43f2a 100644 --- a/src/services/code-index/service-factory.ts +++ b/src/services/code-index/service-factory.ts @@ -59,7 +59,7 @@ export class CodeIndexServiceFactory { vectorStore: IVectorStore parser: ICodeParser scanner: DirectoryScanner - createWatcher: () => IFileWatcher + fileWatcher: IFileWatcher } { if (!this.configManager.isFeatureConfigured) { throw new Error(t("embeddings:serviceFactory.codeIndexingNotConfigured")) @@ -76,23 +76,22 @@ export class CodeIndexServiceFactory { cacheManager: this.cacheManager, ignoreInstance, }) - const createWatcher = () => - this.fileWatcherFactory.create({ - workspacePath: this.workspacePath, - context, - embedder, - vectorStore, - cacheManager, - ignoreInstance, - rooIgnoreController, - }) + const fileWatcher = this.fileWatcherFactory.create({ + workspacePath: this.workspacePath, + context, + embedder, + vectorStore, + cacheManager, + ignoreInstance, + rooIgnoreController, + }) return { embedder, vectorStore, parser, scanner, - createWatcher, + fileWatcher, } } } From d07446c1563b40de331274c2d9d155a01ffd748d Mon Sep 17 00:00:00 2001 From: gubin-dev Date: Tue, 29 Sep 2026 00:40:12 +0300 Subject: [PATCH 4/6] fix(code-index): recreate watchers through injected factory on restart --- .../code-index-watcher-session.spec.ts | 67 ++++++++++++++++--- .../code-index/__tests__/orchestrator.spec.ts | 59 +++++++++++++--- .../code-index/code-index-watcher-session.ts | 46 +++++++------ .../interfaces/file-watcher-factory.ts | 2 +- src/services/code-index/manager.ts | 4 +- src/services/code-index/orchestrator.ts | 7 +- .../__tests__/processor-factories.spec.ts | 10 +-- .../processors/file-watcher-factory.ts | 14 ++-- src/services/code-index/service-factory.ts | 10 +-- 9 files changed, 157 insertions(+), 62 deletions(-) diff --git a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts index fa0b27001c..266ffa6973 100644 --- a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts +++ b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts @@ -35,17 +35,19 @@ function setup() { onDidFinishBatchProcessing: finish.event, } satisfies IFileWatcher const state = new CodeIndexStateManager() - const session = new WatcherSession(watcher, state) - return { start, progress, finish, watcher, state, session } + const factory = { create: vi.fn(() => watcher) } + const session = new WatcherSession(factory, state) + return { start, progress, finish, watcher, state, session, factory } } describe("WatcherSession", () => { - it("does not initialize after being stopped before startup", async () => { - const { session, watcher } = setup() + it("allows startup after stopping an idle owner without allocating resources", async () => { + const { session, watcher, factory } = setup() session.stop() - await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) - expect(watcher.initialize).not.toHaveBeenCalled() - expect(watcher.dispose).toHaveBeenCalledTimes(1) + expect(factory.create).not.toHaveBeenCalled() + expect(watcher.dispose).not.toHaveBeenCalled() + await session.start() + expect(watcher.initialize).toHaveBeenCalledTimes(1) }) it("reuses pending and active sessions without duplicating subscriptions", async () => { @@ -88,7 +90,6 @@ describe("WatcherSession", () => { await expect(pending).rejects.toMatchObject({ name: "AbortError" }) expect(start.event).not.toHaveBeenCalled() expect(watcher.dispose).toHaveBeenCalledTimes(2) - await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) expect(watcher.initialize).toHaveBeenCalledTimes(1) }) @@ -102,7 +103,55 @@ describe("WatcherSession", () => { expect(state.state).toBe("Standby") for (const source of [start, progress, finish]) expect(source.dispose).toHaveBeenCalledTimes(1) expect(watcher.dispose).toHaveBeenCalledTimes(1) - await expect(session.start()).rejects.toMatchObject({ name: "AbortError" }) + }) + + it.each(["stop", "failure"])("creates a working replacement after %s", async (reason) => { + const first = setup() + const next = setup() + first.factory.create.mockReturnValueOnce(first.watcher).mockReturnValue(next.watcher) + if (reason === "failure") { + first.watcher.initialize.mockRejectedValueOnce(new Error("startup failed")) + await expect(first.session.start()).rejects.toThrow("startup failed") + } else { + await first.session.start() + first.session.stop() + } + await first.session.start() + expect(first.factory.create).toHaveBeenCalledTimes(2) + expect(first.watcher.dispose).toHaveBeenCalledTimes(1) + next.progress.fire({ processedInBatch: 0, totalInBatch: 1 }) + expect(first.state.state).toBe("Indexing") + next.finish.fire({ processedFiles: [] }) + expect(first.state.state).toBe("Indexed") + }) + + it.each(["resolve", "reject"])("preserves replacement when old startup later %ss", async (outcome) => { + const first = setup() + const next = setup() + let resolve!: () => void + let reject!: (error: Error) => void + first.watcher.initialize.mockReturnValue( + new Promise((done, fail) => { + resolve = done + reject = fail + }), + ) + first.factory.create.mockReturnValueOnce(first.watcher).mockReturnValue(next.watcher) + const pending = first.session.start() + const rejected = expect(pending).rejects.toBeInstanceOf(Error) + first.session.stop() + await first.session.start() + if (outcome === "resolve") resolve() + else reject(new Error("late failure")) + await rejected + await first.session.start() + expect(first.factory.create).toHaveBeenCalledTimes(2) + expect(next.watcher.initialize).toHaveBeenCalledTimes(1) + expect(next.watcher.dispose).not.toHaveBeenCalled() + next.progress.fire({ processedInBatch: 0, totalInBatch: 1 }) + expect(first.state.state).toBe("Indexing") + next.finish.fire({ processedFiles: [] }) + expect(first.state.state).toBe("Indexed") }) describe("real state manager integration", () => { diff --git a/src/services/code-index/__tests__/orchestrator.spec.ts b/src/services/code-index/__tests__/orchestrator.spec.ts index 58f292ec82..8fa48be304 100644 --- a/src/services/code-index/__tests__/orchestrator.spec.ts +++ b/src/services/code-index/__tests__/orchestrator.spec.ts @@ -133,7 +133,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) await orchestrator.startIndexing() expect(events).toEqual(["incomplete", "scan", "watcher", "complete"]) @@ -155,7 +155,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) scanner.scanDirectory.mockImplementation(async () => { orchestrator.stopIndexing() @@ -184,7 +184,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) // Act @@ -213,7 +213,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) // Act @@ -249,7 +249,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) await orchestrator.startIndexing() @@ -279,7 +279,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) await orchestrator.startIndexing() @@ -366,7 +366,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) const indexing = orchestrator.startIndexing() await entered @@ -379,6 +379,45 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() }) + it.each(["stop", "clear", "error"])("restarts the same orchestrator after %s", async (reason) => { + scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) + vectorStore.deleteCollection = vi.fn().mockResolvedValue(undefined) + const nextWatcher = { + initialize: vi.fn().mockResolvedValue(undefined), + processFile: vi.fn(), + onDidStartBatchProcessing: vi.fn().mockReturnValue({ dispose: vi.fn() }), + onBatchProgressUpdate: vi.fn().mockReturnValue({ dispose: vi.fn() }), + onDidFinishBatchProcessing: vi.fn().mockReturnValue({ dispose: vi.fn() }), + dispose: vi.fn(), + } + const factory = { create: vi.fn().mockReturnValueOnce(fileWatcher).mockReturnValue(nextWatcher) } + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + factory, + ) + if (reason === "error") vectorStore.initialize.mockRejectedValueOnce(new Error("Qdrant unavailable")) + await orchestrator.startIndexing() + if (reason === "error") { + expect(stateManager.state).toBe("Error") + expect(factory.create).not.toHaveBeenCalled() + } else { + expect(stateManager.state).toBe("Indexed") + if (reason === "stop") orchestrator.stopIndexing() + else await orchestrator.clearIndexData() + expect(fileWatcher.dispose).toHaveBeenCalledTimes(1) + } + await orchestrator.startIndexing() + expect(stateManager.state).toBe("Indexed") + expect(factory.create).toHaveBeenCalledTimes(reason === "error" ? 1 : 2) + expect(fileWatcher.initialize).toHaveBeenCalledTimes(1) + expect(nextWatcher.initialize).toHaveBeenCalledTimes(reason === "error" ? 0 : 1) + }) + it("should abort indexing when stopIndexing() is called", async () => { // Make scanner hang until aborted scanner.scanDirectory.mockImplementation( @@ -402,7 +441,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) // Start indexing (async, don't await) @@ -445,7 +484,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) const indexingPromise = orchestrator.startIndexing() @@ -483,7 +522,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) const indexingPromise = orchestrator.startIndexing() diff --git a/src/services/code-index/code-index-watcher-session.ts b/src/services/code-index/code-index-watcher-session.ts index 18b1a27072..f9a960e94a 100644 --- a/src/services/code-index/code-index-watcher-session.ts +++ b/src/services/code-index/code-index-watcher-session.ts @@ -2,8 +2,10 @@ import * as path from "path" import type { Disposable, Event } from "vscode" import type { BatchProcessingSummary, IFileWatcher } from "./interfaces" import type { CodeIndexStateManager } from "./state-manager" +import type { IFileWatcherFactory } from "./interfaces/file-watcher-factory" interface Session { + watcher: IFileWatcher stopped: boolean subscriptions: Disposable[] ready: Promise @@ -12,10 +14,9 @@ interface Session { /** Owns watcher startup and subscriptions; only batch summaries publish final outcomes. */ export class WatcherSession { private session?: Session - private stopped = false constructor( - private readonly watcher: IFileWatcher, + private readonly watcherFactory: IFileWatcherFactory, private readonly stateManager: Pick< CodeIndexStateManager, "state" | "setSystemState" | "reportFileQueueProgress" @@ -23,25 +24,26 @@ export class WatcherSession { ) {} start(): Promise { - // A disposed watcher owns disposed event emitters and cannot be restarted. - if (this.stopped || this.session?.stopped) - return Promise.reject(new DOMException("Watcher session stopped", "AbortError")) if (this.session) return this.session.ready - const session: Session = { stopped: false, subscriptions: [], ready: Promise.resolve() } + const session: Session = { + watcher: this.watcherFactory.create(), + stopped: false, + subscriptions: [], + ready: Promise.resolve(), + } this.session = session session.ready = this.initialize(session) return session.ready } stop(): void { - if (this.stopped || this.session?.stopped) return - this.stopped = true - if (this.session) { - this.session.stopped = true - this.disposeSubscriptions(this.session) - } - this.watcher.dispose() + const session = this.session + if (!session) return + this.session = undefined + session.stopped = true + this.disposeSubscriptions(session) + session.watcher.dispose() } private disposeSubscriptions(session: Session): void { @@ -50,7 +52,7 @@ export class WatcherSession { private async initialize(session: Session): Promise { try { - await this.watcher.initialize() + await session.watcher.initialize() if (session.stopped) throw new DOMException("Watcher startup stopped", "AbortError") this.subscribeToWatcher(session) @@ -58,17 +60,23 @@ export class WatcherSession { session.stopped = true this.disposeSubscriptions(session) // Initialization may have allocated resources after stop() disposed the watcher. - this.watcher.dispose() + session.watcher.dispose() + if (this.session === session) this.session = undefined throw error } } private subscribeToWatcher(session: Session): void { - this.subscribe(session, this.watcher.onDidStartBatchProcessing, (files) => this.handleBatchStarted(files)) - this.subscribe(session, this.watcher.onBatchProgressUpdate, ({ processedInBatch, totalInBatch, currentFile }) => - this.handleBatchProgress(processedInBatch, totalInBatch, currentFile), + this.subscribe(session, session.watcher.onDidStartBatchProcessing, (files) => this.handleBatchStarted(files)) + this.subscribe( + session, + session.watcher.onBatchProgressUpdate, + ({ processedInBatch, totalInBatch, currentFile }) => + this.handleBatchProgress(processedInBatch, totalInBatch, currentFile), + ) + this.subscribe(session, session.watcher.onDidFinishBatchProcessing, (summary) => + this.handleBatchFinished(summary), ) - this.subscribe(session, this.watcher.onDidFinishBatchProcessing, (summary) => this.handleBatchFinished(summary)) } private subscribe(session: Session, event: Event, handler: (value: T) => void): void { diff --git a/src/services/code-index/interfaces/file-watcher-factory.ts b/src/services/code-index/interfaces/file-watcher-factory.ts index 94fc8882fc..9be70ec642 100644 --- a/src/services/code-index/interfaces/file-watcher-factory.ts +++ b/src/services/code-index/interfaces/file-watcher-factory.ts @@ -15,5 +15,5 @@ export interface FileWatcherFactoryOptions { } export interface IFileWatcherFactory { - create(options: FileWatcherFactoryOptions): IFileWatcher + create(): IFileWatcher } diff --git a/src/services/code-index/manager.ts b/src/services/code-index/manager.ts index ee181b17d6..8a32104314 100644 --- a/src/services/code-index/manager.ts +++ b/src/services/code-index/manager.ts @@ -404,7 +404,7 @@ export class CodeIndexManager { await rooIgnoreController.initialize() // (Re)Create shared service instances - const { embedder, vectorStore, scanner, fileWatcher } = this._serviceFactory.createServices( + const { embedder, vectorStore, scanner, fileWatcherFactory } = this._serviceFactory.createServices( this.context, this._cacheManager!, ignoreInstance, @@ -427,7 +427,7 @@ export class CodeIndexManager { this._cacheManager!, vectorStore, scanner, - fileWatcher, + fileWatcherFactory, ) // (Re)Initialize search service diff --git a/src/services/code-index/orchestrator.ts b/src/services/code-index/orchestrator.ts index bf25ac1f1d..ddb4dfe86c 100644 --- a/src/services/code-index/orchestrator.ts +++ b/src/services/code-index/orchestrator.ts @@ -1,7 +1,8 @@ import * as vscode from "vscode" import { CodeIndexConfigManager } from "./config-manager" import { CodeIndexStateManager, IndexingState } from "./state-manager" -import { IFileWatcher, IVectorStore } from "./interfaces" +import { IVectorStore } from "./interfaces" +import type { IFileWatcherFactory } from "./interfaces/file-watcher-factory" import { WatcherSession } from "./code-index-watcher-session" import { DirectoryScanner } from "./processors" import { CacheManager } from "./cache-manager" @@ -26,10 +27,10 @@ export class CodeIndexOrchestrator { private readonly cacheManager: CacheManager, private readonly vectorStore: IVectorStore, scanner: DirectoryScanner, - fileWatcher: IFileWatcher, + fileWatcherFactory: IFileWatcherFactory, ) { this.scanExecutor = new CodeIndexScanExecutor(workspacePath, scanner, vectorStore, stateManager) - this.watcherSession = new WatcherSession(fileWatcher, stateManager) + this.watcherSession = new WatcherSession(fileWatcherFactory, stateManager) } /** diff --git a/src/services/code-index/processors/__tests__/processor-factories.spec.ts b/src/services/code-index/processors/__tests__/processor-factories.spec.ts index 4af92c14cd..8226352e8e 100644 --- a/src/services/code-index/processors/__tests__/processor-factories.spec.ts +++ b/src/services/code-index/processors/__tests__/processor-factories.spec.ts @@ -55,7 +55,7 @@ describe("processor factories", () => { const options = dependencies() const scanner = new DirectoryScannerFactory().create(options) - const watcher = new FileWatcherFactory().create(options) + const watcher = new FileWatcherFactory(options).create() expect(scanner).toBeInstanceOf(DirectoryScanner) expect(watcher).toBeInstanceOf(FileWatcher) @@ -92,10 +92,10 @@ describe("processor factories", () => { vi.mocked(vscode.workspace.getConfiguration).mockReturnValue(configuration) const options = dependencies() const scanners = new DirectoryScannerFactory() - const watchers = new FileWatcherFactory() + const watchers = new FileWatcherFactory(options) expect(scanners.create(options)).not.toBe(scanners.create(options)) - expect(watchers.create(options)).not.toBe(watchers.create(options)) + expect(watchers.create()).not.toBe(watchers.create()) expect(vi.mocked(DirectoryScanner).mock.calls.map((args) => args[5])).toEqual([32, 64]) expect(vi.mocked(FileWatcher).mock.calls.map((args) => args[7])).toEqual([128, 256]) }) @@ -159,8 +159,10 @@ describe("processor factories", () => { vectorStore: options.vectorStore, parser: codeParser, scanner: vi.mocked(DirectoryScanner).mock.instances[0], - fileWatcher: vi.mocked(FileWatcher).mock.instances[0], + fileWatcherFactory: expect.any(FileWatcherFactory), }) + expect(FileWatcher).not.toHaveBeenCalled() + expect(services.fileWatcherFactory.create()).not.toBe(services.fileWatcherFactory.create()) expect(vi.mocked(DirectoryScanner).mock.calls[0][3]).toBe(options.cacheManager) expect(vi.mocked(FileWatcher).mock.calls[0][2]).toBe(watcherCache) expect(vi.mocked(FileWatcher).mock.calls[0][0]).toBe(options.workspacePath) diff --git a/src/services/code-index/processors/file-watcher-factory.ts b/src/services/code-index/processors/file-watcher-factory.ts index 2f4cc227b6..044e5715f3 100644 --- a/src/services/code-index/processors/file-watcher-factory.ts +++ b/src/services/code-index/processors/file-watcher-factory.ts @@ -4,15 +4,11 @@ import { FileWatcher } from "./file-watcher" import { getEmbeddingBatchSize } from "./get-embedding-batch-size" export class FileWatcherFactory implements IFileWatcherFactory { - public create({ - workspacePath, - context, - cacheManager, - embedder, - vectorStore, - ignoreInstance, - rooIgnoreController, - }: FileWatcherFactoryOptions): IFileWatcher { + constructor(private readonly options: FileWatcherFactoryOptions) {} + + public create(): IFileWatcher { + const { workspacePath, context, cacheManager, embedder, vectorStore, ignoreInstance, rooIgnoreController } = + this.options return new FileWatcher( workspacePath, context, diff --git a/src/services/code-index/service-factory.ts b/src/services/code-index/service-factory.ts index d2bbe43f2a..0344eac305 100644 --- a/src/services/code-index/service-factory.ts +++ b/src/services/code-index/service-factory.ts @@ -11,7 +11,8 @@ import { VectorStoreFactory } from "./vector-store/vector-store-factory" import { codeParser, DirectoryScanner } from "./processors" import { DirectoryScannerFactory } from "./processors/directory-scanner-factory" import { FileWatcherFactory } from "./processors/file-watcher-factory" -import { ICodeParser, IEmbedder, IFileWatcher, IVectorStore } from "./interfaces" +import { ICodeParser, IEmbedder, IVectorStore } from "./interfaces" +import type { IFileWatcherFactory } from "./interfaces/file-watcher-factory" import { CodeIndexConfigManager } from "./config-manager" import { CacheManager } from "./cache-manager" @@ -25,7 +26,6 @@ export class CodeIndexServiceFactory { private readonly embedderFactory = new EmbedderFactory() private readonly vectorStoreFactory = new VectorStoreFactory() private readonly directoryScannerFactory = new DirectoryScannerFactory() - private readonly fileWatcherFactory = new FileWatcherFactory() private readonly embedderValidationManager = new EmbedderValidationManager() constructor( @@ -59,7 +59,7 @@ export class CodeIndexServiceFactory { vectorStore: IVectorStore parser: ICodeParser scanner: DirectoryScanner - fileWatcher: IFileWatcher + fileWatcherFactory: IFileWatcherFactory } { if (!this.configManager.isFeatureConfigured) { throw new Error(t("embeddings:serviceFactory.codeIndexingNotConfigured")) @@ -76,7 +76,7 @@ export class CodeIndexServiceFactory { cacheManager: this.cacheManager, ignoreInstance, }) - const fileWatcher = this.fileWatcherFactory.create({ + const fileWatcherFactory = new FileWatcherFactory({ workspacePath: this.workspacePath, context, embedder, @@ -91,7 +91,7 @@ export class CodeIndexServiceFactory { vectorStore, parser, scanner, - fileWatcher, + fileWatcherFactory, } } } From 5d02a51a75637ed4318beeedcfac94b314bc1a73 Mon Sep 17 00:00:00 2001 From: gubin-dev Date: Tue, 29 Sep 2026 00:50:07 +0300 Subject: [PATCH 5/6] test(code-index): cover empty watcher batch preserving error state --- .../__tests__/code-index-watcher-session.spec.ts | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts index 266ffa6973..f37698f626 100644 --- a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts +++ b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts @@ -155,6 +155,16 @@ describe("WatcherSession", () => { }) describe("real state manager integration", () => { + it("preserves the current outcome when an empty batch starts", async () => { + const { session, start, finish, state } = setup() + await session.start() + finish.fire({ processedFiles: [], batchError: new Error("database unavailable") }) + const outcome = state.getCurrentStatus() + expect(outcome.systemStatus).toBe("Error") + start.fire([]) + expect(state.getCurrentStatus()).toEqual(outcome) + }) + it.each(["success", "skipped", "error", "local_error"] as const)( "preserves the %s outcome through terminal and empty progress", async (status) => { From f33afe370a873aaf758333ddb3900243e677cdec Mon Sep 17 00:00:00 2001 From: gubin-dev Date: Tue, 29 Sep 2026 19:26:19 +0300 Subject: [PATCH 6/6] refactor(code-index): clarify watcher session naming and types --- .../code-index-watcher-session.spec.ts | 6 +++--- .../code-index/code-index-watcher-session.ts | 19 +++++-------------- .../code-index/interfaces/watcher-session.ts | 9 +++++++++ src/services/code-index/orchestrator.ts | 6 +++--- 4 files changed, 20 insertions(+), 20 deletions(-) create mode 100644 src/services/code-index/interfaces/watcher-session.ts diff --git a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts index f37698f626..02ff48b390 100644 --- a/src/services/code-index/__tests__/code-index-watcher-session.spec.ts +++ b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts @@ -1,7 +1,7 @@ import type { Event } from "vscode" import type { BatchProcessingSummary, IFileWatcher } from "../interfaces" import { CodeIndexStateManager } from "../state-manager" -import { WatcherSession } from "../code-index-watcher-session" +import { CodeIndexWatcherSession } from "../code-index-watcher-session" vi.mock("vscode", async () => { const { makeEventEmitter } = await import("../../../test-utils/vscode") @@ -36,11 +36,11 @@ function setup() { } satisfies IFileWatcher const state = new CodeIndexStateManager() const factory = { create: vi.fn(() => watcher) } - const session = new WatcherSession(factory, state) + const session = new CodeIndexWatcherSession(factory, state) return { start, progress, finish, watcher, state, session, factory } } -describe("WatcherSession", () => { +describe("CodeIndexWatcherSession", () => { it("allows startup after stopping an idle owner without allocating resources", async () => { const { session, watcher, factory } = setup() session.stop() diff --git a/src/services/code-index/code-index-watcher-session.ts b/src/services/code-index/code-index-watcher-session.ts index f9a960e94a..3a86262aca 100644 --- a/src/services/code-index/code-index-watcher-session.ts +++ b/src/services/code-index/code-index-watcher-session.ts @@ -1,26 +1,17 @@ import * as path from "path" -import type { Disposable, Event } from "vscode" -import type { BatchProcessingSummary, IFileWatcher } from "./interfaces" +import type { Event } from "vscode" +import type { BatchProcessingSummary } from "./interfaces" import type { CodeIndexStateManager } from "./state-manager" import type { IFileWatcherFactory } from "./interfaces/file-watcher-factory" - -interface Session { - watcher: IFileWatcher - stopped: boolean - subscriptions: Disposable[] - ready: Promise -} +import type { Session } from "./interfaces/watcher-session" /** Owns watcher startup and subscriptions; only batch summaries publish final outcomes. */ -export class WatcherSession { +export class CodeIndexWatcherSession { private session?: Session constructor( private readonly watcherFactory: IFileWatcherFactory, - private readonly stateManager: Pick< - CodeIndexStateManager, - "state" | "setSystemState" | "reportFileQueueProgress" - >, + private readonly stateManager: CodeIndexStateManager, ) {} start(): Promise { diff --git a/src/services/code-index/interfaces/watcher-session.ts b/src/services/code-index/interfaces/watcher-session.ts new file mode 100644 index 0000000000..56bcd184e1 --- /dev/null +++ b/src/services/code-index/interfaces/watcher-session.ts @@ -0,0 +1,9 @@ +import type { Disposable } from "vscode" +import type { IFileWatcher } from "../interfaces" + +export interface Session { + watcher: IFileWatcher + stopped: boolean + subscriptions: Disposable[] + ready: Promise +} diff --git a/src/services/code-index/orchestrator.ts b/src/services/code-index/orchestrator.ts index ddb4dfe86c..632b936dd9 100644 --- a/src/services/code-index/orchestrator.ts +++ b/src/services/code-index/orchestrator.ts @@ -3,7 +3,7 @@ import { CodeIndexConfigManager } from "./config-manager" import { CodeIndexStateManager, IndexingState } from "./state-manager" import { IVectorStore } from "./interfaces" import type { IFileWatcherFactory } from "./interfaces/file-watcher-factory" -import { WatcherSession } from "./code-index-watcher-session" +import { CodeIndexWatcherSession } from "./code-index-watcher-session" import { DirectoryScanner } from "./processors" import { CacheManager } from "./cache-manager" import { CodeIndexScanExecutor } from "./code-index-scan-executor" @@ -15,7 +15,7 @@ import { t } from "../../i18n" * Manages the code indexing workflow, coordinating between different services and managers. */ export class CodeIndexOrchestrator { - private readonly watcherSession: WatcherSession + private readonly watcherSession: CodeIndexWatcherSession private _isProcessing: boolean = false private _abortController: AbortController | null = null private readonly scanExecutor: CodeIndexScanExecutor @@ -30,7 +30,7 @@ export class CodeIndexOrchestrator { fileWatcherFactory: IFileWatcherFactory, ) { this.scanExecutor = new CodeIndexScanExecutor(workspacePath, scanner, vectorStore, stateManager) - this.watcherSession = new WatcherSession(fileWatcherFactory, stateManager) + this.watcherSession = new CodeIndexWatcherSession(fileWatcherFactory, stateManager) } /**