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..02ff48b390 --- /dev/null +++ b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts @@ -0,0 +1,223 @@ +import type { Event } from "vscode" +import type { BatchProcessingSummary, IFileWatcher } from "../interfaces" +import { CodeIndexStateManager } from "../state-manager" +import { CodeIndexWatcherSession } 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 factory = { create: vi.fn(() => watcher) } + const session = new CodeIndexWatcherSession(factory, state) + return { start, progress, finish, watcher, state, session, factory } +} + +describe("CodeIndexWatcherSession", () => { + it("allows startup after stopping an idle owner without allocating resources", async () => { + const { session, watcher, factory } = setup() + session.stop() + 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 () => { + 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) + 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) + }) + + 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", () => { + 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) => { + 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..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() @@ -346,6 +346,78 @@ 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, + { create: () => 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.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( @@ -369,7 +441,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) // Start indexing (async, don't await) @@ -412,7 +484,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { cacheManager, vectorStore, scanner, - fileWatcher, + { create: () => fileWatcher }, ) const indexingPromise = orchestrator.startIndexing() @@ -450,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 new file mode 100644 index 0000000000..3a86262aca --- /dev/null +++ b/src/services/code-index/code-index-watcher-session.ts @@ -0,0 +1,110 @@ +import * as path from "path" +import type { Event } from "vscode" +import type { BatchProcessingSummary } from "./interfaces" +import type { CodeIndexStateManager } from "./state-manager" +import type { IFileWatcherFactory } from "./interfaces/file-watcher-factory" +import type { Session } from "./interfaces/watcher-session" + +/** Owns watcher startup and subscriptions; only batch summaries publish final outcomes. */ +export class CodeIndexWatcherSession { + private session?: Session + + constructor( + private readonly watcherFactory: IFileWatcherFactory, + private readonly stateManager: CodeIndexStateManager, + ) {} + + start(): Promise { + if (this.session) return this.session.ready + + 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 { + const session = this.session + if (!session) return + this.session = undefined + session.stopped = true + this.disposeSubscriptions(session) + session.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 session.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. + session.watcher.dispose() + if (this.session === session) this.session = undefined + 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), + ) + } + + 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/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/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/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 c055fff370..632b936dd9 100644 --- a/src/services/code-index/orchestrator.ts +++ b/src/services/code-index/orchestrator.ts @@ -1,8 +1,9 @@ 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 { IVectorStore } from "./interfaces" +import type { IFileWatcherFactory } from "./interfaces/file-watcher-factory" +import { CodeIndexWatcherSession } from "./code-index-watcher-session" import { DirectoryScanner } from "./processors" import { CacheManager } from "./cache-manager" import { CodeIndexScanExecutor } from "./code-index-scan-executor" @@ -14,7 +15,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: CodeIndexWatcherSession private _isProcessing: boolean = false private _abortController: AbortController | null = null private readonly scanExecutor: CodeIndexScanExecutor @@ -26,9 +27,10 @@ export class CodeIndexOrchestrator { private readonly cacheManager: CacheManager, private readonly vectorStore: IVectorStore, scanner: DirectoryScanner, - private readonly fileWatcher: IFileWatcher, + fileWatcherFactory: IFileWatcherFactory, ) { this.scanExecutor = new CodeIndexScanExecutor(workspacePath, scanner, vectorStore, stateManager) + this.watcherSession = new CodeIndexWatcherSession(fileWatcherFactory, stateManager) } /** @@ -42,45 +44,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 +121,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 +137,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 +218,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")) 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, } } }