diff --git a/src/eslint-suppressions.json b/src/eslint-suppressions.json index 24db0bf433..e2d4344669 100644 --- a/src/eslint-suppressions.json +++ b/src/eslint-suppressions.json @@ -1306,7 +1306,7 @@ }, "services/code-index/__tests__/orchestrator.spec.ts": { "@typescript-eslint/no-explicit-any": { - "count": 25 + "count": 23 } }, "services/code-index/__tests__/service-factory.spec.ts": { @@ -1389,11 +1389,6 @@ "count": 2 } }, - "services/code-index/orchestrator.ts": { - "@typescript-eslint/no-explicit-any": { - "count": 2 - } - }, "services/code-index/processors/__tests__/file-watcher.spec.ts": { "@typescript-eslint/no-explicit-any": { "count": 25 diff --git a/src/services/code-index/__tests__/code-index-recovery.spec.ts b/src/services/code-index/__tests__/code-index-recovery.spec.ts new file mode 100644 index 0000000000..c58c2ce193 --- /dev/null +++ b/src/services/code-index/__tests__/code-index-recovery.spec.ts @@ -0,0 +1,88 @@ +import { CodeIndexRecovery } from "../code-index-recovery" +import { CodeIndexRun, type CodeIndexRunState } from "../code-index-run" +import { StateHolder } from "../../../utils/StateHolder" + +vi.mock("@roo-code/telemetry", () => ({ + TelemetryService: { instance: { captureEvent: vi.fn() } }, +})) +vi.mock("../../../i18n", () => ({ + t: (key: string, options?: { errorMessage?: unknown }) => `${key}:${options?.errorMessage ?? ""}`, +})) + +function setup() { + const cache = { flush: vi.fn().mockResolvedValue(undefined), clearCacheFile: vi.fn().mockResolvedValue(undefined) } + const store = { clearCollection: vi.fn().mockResolvedValue(undefined) } + const state = { setSystemState: vi.fn() } + const watcher = { stop: vi.fn() } + const run = new CodeIndexRun(new AbortController(), new StateHolder("running")) + return { cache, store, state, watcher, run, recovery: new CodeIndexRecovery(cache, store, state, watcher) } +} + +describe("CodeIndexRecovery", () => { + it.each([true, undefined])( + "preserves preexisting or unknown data during full-scan recovery (%s)", + async (presence) => { + const { recovery, cache, store, run } = setup() + run.preexistingCodePoints = presence + run.markScanStarted("full") + await recovery.handle(new Error("scan failed"), run) + expect(store.clearCollection).not.toHaveBeenCalled() + expect(cache.clearCacheFile).not.toHaveBeenCalled() + }, + ) + + it.each(["preparation", "incremental", "full"] as const)("applies cleanup policy for %s failures", async (mode) => { + const { recovery, cache, store, state, watcher, run } = setup() + run.preexistingCodePoints = false + if (mode !== "preparation") run.markScanStarted(mode) + await recovery.handle(new Error("scan failed"), run) + expect(store.clearCollection).toHaveBeenCalledTimes(mode === "full" ? 1 : 0) + expect(cache.clearCacheFile).toHaveBeenCalledTimes(mode === "full" ? 1 : 0) + expect(state.setSystemState).toHaveBeenLastCalledWith("Error", expect.stringContaining("scan failed")) + expect(watcher.stop).toHaveBeenCalledOnce() + expect(run.state.value).toBe("running") + }) + + it.each(["signal", "AbortError"])("preserves data when cancellation is identified by %s", async (kind) => { + const { recovery, cache, store, state, watcher, run } = setup() + run.markScanStarted("full") + if (kind === "signal") run.cancel() + const error = kind === "signal" ? new Error("interrupted") : new DOMException("Stopped", "AbortError") + await recovery.handle(error, run) + expect(cache.flush).toHaveBeenCalledOnce() + expect(store.clearCollection).not.toHaveBeenCalled() + expect(cache.clearCacheFile).not.toHaveBeenCalled() + expect(watcher.stop).toHaveBeenCalledOnce() + expect(state.setSystemState).toHaveBeenLastCalledWith("Standby", expect.any(String)) + expect(run.state.value).not.toBe("finished") + }) + + it("finishes cancellation recovery even if flushing fails", async () => { + const { recovery, cache, state, watcher, run } = setup() + run.cancel() + cache.flush.mockRejectedValue(new Error("flush failed")) + await recovery.handle(new Error("stopped"), run) + expect(watcher.stop).toHaveBeenCalledOnce() + expect(state.setSystemState).toHaveBeenLastCalledWith("Standby", expect.any(String)) + }) + + it("attempts cache cleanup after collection cleanup fails and preserves the original failure", async () => { + const { recovery, cache, store, state, watcher, run } = setup() + run.preexistingCodePoints = false + run.markScanStarted("full") + store.clearCollection.mockRejectedValue(new Error("collection cleanup failed")) + cache.clearCacheFile.mockRejectedValue(new Error("cache cleanup failed")) + await recovery.handle(new Error("original failure"), run) + expect(cache.clearCacheFile).toHaveBeenCalledOnce() + expect(watcher.stop).toHaveBeenCalledOnce() + expect(state.setSystemState).toHaveBeenLastCalledWith("Error", expect.stringContaining("original failure")) + }) + + it("reports clearing errors without initiating additional destructive cleanup", () => { + const { recovery, cache, store, state } = setup() + recovery.handleClearError("delete failed") + expect(state.setSystemState).toHaveBeenCalledWith("Error", "Failed to clear index data: delete failed") + expect(store.clearCollection).not.toHaveBeenCalled() + expect(cache.clearCacheFile).not.toHaveBeenCalled() + }) +}) diff --git a/src/services/code-index/__tests__/code-index-run.spec.ts b/src/services/code-index/__tests__/code-index-run.spec.ts new file mode 100644 index 0000000000..724aa1c7f6 --- /dev/null +++ b/src/services/code-index/__tests__/code-index-run.spec.ts @@ -0,0 +1,110 @@ +import { CodeIndexRun, CodeIndexRunState } from "../code-index-run" +import { StateHolder } from "../../../utils/StateHolder" + +function createRun(): CodeIndexRun { + return new CodeIndexRun(new AbortController(), new StateHolder("running")) +} + +describe("CodeIndexRun", () => { + it("uses the injected controller and state holder", async () => { + const controller = new AbortController() + const stateHolder = new StateHolder("running") + const run = new CodeIndexRun(controller, stateHolder) + expect(run.signal).toBe(controller.signal) + expect(run.state).toBe(stateHolder) + run.cancel() + expect(controller.signal.aborted).toBe(true) + expect(stateHolder.value).toBe("cancelling") + run.finish() + expect(stateHolder.value).toBe("finished") + await expect(run.waitUntilFinished()).resolves.toBeUndefined() + }) + + it("starts without cancellation or destructive cleanup eligibility", () => { + const run = createRun() + expect(run.signal.aborted).toBe(false) + expect(run.fullScanStarted).toBe(false) + }) + + it.each(["full", "incremental"] as const)("records when a %s scan starts", (mode) => { + const run = createRun() + run.markScanStarted(mode) + expect(run.fullScanStarted).toBe(mode === "full") + }) + + it("keeps completion pending after cancellation until cleanup finishes", async () => { + const run = createRun() + const finished = vi.fn() + const completion = run.waitUntilFinished().then(finished) + const aborted = vi.fn() + run.signal.addEventListener("abort", aborted) + + run.cancel() + run.cancel() + await Promise.resolve() + expect(run.signal.aborted).toBe(true) + expect(aborted).toHaveBeenCalledOnce() + expect(finished).not.toHaveBeenCalled() + + run.finish() + run.finish() + await completion + expect(finished).toHaveBeenCalledOnce() + }) + + it("finishes a successful run without cancelling it", async () => { + const run = createRun() + run.finish() + await expect(run.waitUntilFinished()).resolves.toBeUndefined() + expect(run.signal.aborted).toBe(false) + }) + + it("isolates cancellation, scan mode and completion between runs", async () => { + const first = createRun() + const second = createRun() + const secondFinished = vi.fn() + const completion = second.waitUntilFinished().then(secondFinished) + first.markScanStarted("full") + first.cancel() + first.finish() + await first.waitUntilFinished() + + expect(second.signal.aborted).toBe(false) + expect(second.fullScanStarted).toBe(false) + expect(secondFinished).not.toHaveBeenCalled() + second.finish() + await completion + }) + + it("notifies all waiters and releases their subscriptions", async () => { + const run = createRun() + const firstNotified = vi.fn() + const secondNotified = vi.fn() + const first = run.waitUntilFinished().then(firstNotified) + const second = run.waitUntilFinished().then(secondNotified) + await Promise.resolve() + expect(firstNotified).not.toHaveBeenCalled() + expect(secondNotified).not.toHaveBeenCalled() + + run.finish() + await Promise.all([first, second]) + expect(firstNotified).toHaveBeenCalledOnce() + expect(secondNotified).toHaveBeenCalledOnce() + + await expect(run.waitUntilFinished()).resolves.toBeUndefined() + expect(run.state.value).toBe("finished") + }) + + it("replays state and publishes cancellation followed by completion", () => { + const run = createRun() + const states: string[] = [] + const subscription = run.state.subscribe((state) => states.push(state)) + run.cancel() + run.cancel() + run.finish() + run.finish() + run.cancel() + expect(states).toEqual(["running", "cancelling", "finished"]) + subscription.unsubscribe() + }) +}) diff --git a/src/services/code-index/__tests__/code-index-scan-executor.spec.ts b/src/services/code-index/__tests__/code-index-scan-executor.spec.ts new file mode 100644 index 0000000000..e835744a5e --- /dev/null +++ b/src/services/code-index/__tests__/code-index-scan-executor.spec.ts @@ -0,0 +1,79 @@ +import { CodeIndexScanExecutor } from "../code-index-scan-executor" +import type { IDirectoryScanner } from "../interfaces" + +vi.mock("../../../i18n", () => ({ t: (key: string) => key })) + +describe("CodeIndexScanExecutor", () => { + function setup() { + const scanner = { scanDirectory: vi.fn() } + const vectorStore = { markIndexingIncomplete: vi.fn().mockResolvedValue(undefined) } + const stateManager = { setSystemState: vi.fn(), reportBlockIndexingProgress: vi.fn() } + const executor = new CodeIndexScanExecutor("/workspace", scanner, vectorStore, stateManager) + return { scanner, vectorStore, stateManager, executor } + } + + it.each(["runFullScan", "runIncrementalScan"] as const)("%s rejects a missing scanner result", async (method) => { + const { executor } = setup() + // An unconfigured mock returns undefined, simulating a broken scanner contract. + await expect(executor[method](new AbortController().signal)).rejects.toThrow( + method === "runFullScan" + ? "Scan failed, is scanner initialized?" + : "Incremental scan failed, is scanner initialized?", + ) + }) + + it.each(["runFullScan", "runIncrementalScan"] as const)( + "%s reports progress without completing the operation", + async (method) => { + const { scanner, vectorStore, stateManager, executor } = setup() + const signal = new AbortController().signal + scanner.scanDirectory.mockImplementation(async (_path, _onError, onIndexed, onParsed, receivedSignal) => { + expect(receivedSignal).toBe(signal) + expect(vectorStore.markIndexingIncomplete).toHaveBeenCalledOnce() + onParsed?.(3) + onIndexed?.(1) + onIndexed?.(2) + return { stats: { processed: 1, skipped: 0 }, totalBlockCount: 3 } + }) + await executor[method](signal) + expect(stateManager.reportBlockIndexingProgress.mock.calls).toEqual([ + [0, 3], + [1, 3], + [3, 3], + ]) + expect(stateManager.setSystemState).not.toHaveBeenCalledWith("Indexed", expect.any(String)) + }, + ) + + it.each(["runFullScan", "runIncrementalScan"] as const)( + "%s propagates cancellation to its owner", + async (method) => { + const { scanner, executor } = setup() + const controller = new AbortController() + scanner.scanDirectory.mockImplementation(async () => { + controller.abort() + return { stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 } + }) + await expect(executor[method](controller.signal)).rejects.toMatchObject({ name: "AbortError" }) + }, + ) + + it.each([ + { method: "runFullScan", rejects: false }, + { method: "runIncrementalScan", rejects: true }, + ] as const)("$method preserves its partial-failure policy", async ({ method, rejects }) => { + const { scanner, executor } = setup() + scanner.scanDirectory.mockImplementation(async (_path, onError, onIndexed, onParsed) => { + onParsed?.(10) + onIndexed?.(9) + onError?.(new Error("batch failed")) + return { stats: { processed: 2, skipped: 0 }, totalBlockCount: 10 } + }) + const result = executor[method](new AbortController().signal) + if (rejects) { + await expect(result).rejects.toThrow("batch failed") + } else { + await expect(result).resolves.toBeUndefined() + } + }) +}) 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..cab4c6022a --- /dev/null +++ b/src/services/code-index/__tests__/code-index-watcher-session.spec.ts @@ -0,0 +1,179 @@ +import type { IFileWatcher } from "../interfaces" +import { CodeIndexWatcherSession } from "../code-index-watcher-session" + +function setup() { + const progress = { dispose: vi.fn() } + const finished = { dispose: vi.fn() } + const watcher = { + initialize: vi.fn().mockResolvedValue(undefined), + dispose: vi.fn(), + onDidStartBatchProcessing: vi.fn(), + onBatchProgressUpdate: vi.fn().mockReturnValue(progress), + onDidFinishBatchProcessing: vi.fn().mockReturnValue(finished), + processFile: vi.fn(), + } satisfies IFileWatcher + const stateManager = { + state: "Standby" as const, + setSystemState: vi.fn(), + reportFileQueueProgress: vi.fn(), + } + const session = new CodeIndexWatcherSession(watcher, stateManager) + return { session, watcher, progress, finished, stateManager } +} + +describe("CodeIndexWatcherSession", () => { + it.each(["error", "local_error"] as const)("preserves the specific message of a %s file result", async (status) => { + const { session, watcher, stateManager } = setup() + await session.start(new AbortController().signal) + watcher.onDidFinishBatchProcessing.mock.calls[0][0]({ + processedFiles: [{ path: "file.ts", status, error: new Error("specific failure") }], + }) + expect(stateManager.setSystemState).toHaveBeenLastCalledWith("Error", "specific failure") + session.stop() + }) + + it("unsubscribes the registered listeners so later events cannot change state", async () => { + const { session, watcher, stateManager } = setup() + const progressListeners = new Set[0]>() + const finishedListeners = new Set[0]>() + watcher.onBatchProgressUpdate.mockImplementation((listener) => { + progressListeners.add(listener) + return { + dispose: () => { + progressListeners.delete(listener) + }, + } + }) + watcher.onDidFinishBatchProcessing.mockImplementation((listener) => { + finishedListeners.add(listener) + return { + dispose: () => { + finishedListeners.delete(listener) + }, + } + }) + await session.start(new AbortController().signal) + expect(progressListeners.has(watcher.onBatchProgressUpdate.mock.calls[0][0])).toBe(true) + expect(finishedListeners.has(watcher.onDidFinishBatchProcessing.mock.calls[0][0])).toBe(true) + session.stop() + for (const listener of progressListeners) listener({ processedInBatch: 0, totalInBatch: 1 }) + for (const listener of finishedListeners) listener({ processedFiles: [] }) + expect(progressListeners.size).toBe(0) + expect(finishedListeners.size).toBe(0) + expect(stateManager.setSystemState).not.toHaveBeenCalled() + expect(stateManager.reportFileQueueProgress).not.toHaveBeenCalled() + }) + + it("reports batch progress with the file name and ignores terminal progress", async () => { + const { session, watcher, stateManager } = setup() + await session.start(new AbortController().signal) + const onProgress = watcher.onBatchProgressUpdate.mock.calls[0][0] + onProgress({ processedInBatch: 1, totalInBatch: 2, currentFile: "/workspace/file.ts" }) + expect(stateManager.setSystemState).toHaveBeenCalledWith("Indexing", "Processing file changes...") + expect(stateManager.reportFileQueueProgress).toHaveBeenCalledWith(1, 2, "file.ts") + stateManager.setSystemState.mockClear() + onProgress({ processedInBatch: 2, totalInBatch: 2 }) + onProgress({ processedInBatch: 0, totalInBatch: 0 }) + expect(stateManager.setSystemState).not.toHaveBeenCalled() + expect(stateManager.reportFileQueueProgress).toHaveBeenCalledOnce() + session.stop() + }) + + it("reports batch success and individual file errors", async () => { + const { session, watcher, stateManager } = setup() + await session.start(new AbortController().signal) + const onFinished = watcher.onDidFinishBatchProcessing.mock.calls[0][0] + onFinished({ processedFiles: [] }) + expect(stateManager.setSystemState).toHaveBeenLastCalledWith( + "Indexed", + "File changes processed. Index up-to-date.", + ) + onFinished({ processedFiles: [{ path: "file.ts", status: "error" }] }) + expect(stateManager.setSystemState).toHaveBeenLastCalledWith("Error", "Failed to index file: file.ts") + session.stop() + }) + + it("starts once, registers callbacks and releases subscriptions on stop", async () => { + const { session, watcher, progress, finished } = setup() + const signal = new AbortController().signal + await session.start(signal) + await session.start(signal) + expect(session.isRunning).toBe(true) + expect(watcher.initialize).toHaveBeenCalledOnce() + expect(watcher.onBatchProgressUpdate).toHaveBeenCalledOnce() + expect(watcher.onDidFinishBatchProcessing).toHaveBeenCalledOnce() + session.stop() + session.stop() + expect(session.isRunning).toBe(false) + expect(progress.dispose).toHaveBeenCalledOnce() + expect(finished.dispose).toHaveBeenCalledOnce() + }) + + it("does not initialize after cancellation", async () => { + const { session, watcher } = setup() + const controller = new AbortController() + controller.abort() + await expect(session.start(controller.signal)).rejects.toMatchObject({ name: "AbortError" }) + expect(watcher.initialize).not.toHaveBeenCalled() + }) + + it("releases partially registered subscriptions and preserves the startup error", async () => { + const { session, watcher, progress } = setup() + const error = new Error("registration failed") + watcher.onDidFinishBatchProcessing.mockImplementation(() => { + throw error + }) + await expect(session.start(new AbortController().signal)).rejects.toBe(error) + expect(progress.dispose).toHaveBeenCalledOnce() + expect(watcher.dispose).toHaveBeenCalledOnce() + expect(session.isRunning).toBe(false) + }) + + it("cleans up a failed initialization", async () => { + const { session, watcher } = setup() + const error = new Error("initialization failed") + watcher.initialize.mockRejectedValue(error) + await expect(session.start(new AbortController().signal)).rejects.toBe(error) + expect(watcher.dispose).toHaveBeenCalledOnce() + expect(watcher.onBatchProgressUpdate).not.toHaveBeenCalled() + }) + + it.each(["stop", "abort"] as const)("does not revive after %s during initialization", async (action) => { + const { session, watcher } = setup() + let release!: () => void + watcher.initialize.mockImplementation( + () => + new Promise((resolve) => { + release = resolve + }), + ) + const controller = new AbortController() + const starting = session.start(controller.signal) + const rejected = expect(starting).rejects.toMatchObject({ name: "AbortError" }) + if (action === "stop") session.stop() + else controller.abort() + await expect(session.start(new AbortController().signal)).rejects.toThrow("already in progress") + release() + await rejected + expect(session.isRunning).toBe(false) + expect(watcher.onBatchProgressUpdate).not.toHaveBeenCalled() + expect(watcher.dispose).toHaveBeenCalled() + }) + + it("attempts all disposals even if one subscription throws", async () => { + const { session, watcher, progress, finished } = setup() + const log = vi.spyOn(console, "error").mockImplementation(() => undefined) + try { + await session.start(new AbortController().signal) + progress.dispose.mockImplementation(() => { + throw new Error("dispose failed") + }) + session.stop() + expect(finished.dispose).toHaveBeenCalledOnce() + expect(watcher.dispose).toHaveBeenCalledOnce() + expect(log).toHaveBeenCalledOnce() + } finally { + log.mockRestore() + } + }) +}) diff --git a/src/services/code-index/__tests__/orchestrator.spec.ts b/src/services/code-index/__tests__/orchestrator.spec.ts index 86b0f94808..323ad5fdeb 100644 --- a/src/services/code-index/__tests__/orchestrator.spec.ts +++ b/src/services/code-index/__tests__/orchestrator.spec.ts @@ -1,5 +1,8 @@ import { describe, it, expect, beforeEach, vi } from "vitest" import { CodeIndexOrchestrator } from "../orchestrator" +import { TelemetryService } from "@roo-code/telemetry" +import * as vscode from "vscode" +import { CodeIndexStateManager } from "../state-manager" import { clearAllMocks } from "../../../test-utils/reset" @@ -8,6 +11,11 @@ vi.mock("vscode", () => { const path = require("path") const testWorkspacePath = path.join(path.sep, "test", "workspace") return { + EventEmitter: class { + event = vi.fn().mockReturnValue({ dispose: vi.fn() }) + fire = vi.fn() + dispose = vi.fn() + }, window: { activeTextEditor: null, }, @@ -42,8 +50,8 @@ vi.mock("@roo-code/telemetry", () => ({ })) // Mock i18n translator used in orchestrator messages -vi.mock("../../i18n", () => ({ - t: (key: string, params?: any) => { +vi.mock("../../../i18n", () => ({ + t: (key: string, params?: { errorMessage?: string }) => { if (key === "embeddings:orchestrator.failedDuringInitialScan" && params?.errorMessage) { return `Failed during initial scan: ${params.errorMessage}` } @@ -88,6 +96,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { vectorStore = { initialize: vi.fn(), + hasCodePoints: vi.fn().mockResolvedValue(false), hasIndexedData: vi.fn(), markIndexingIncomplete: vi.fn(), markIndexingComplete: vi.fn(), @@ -107,6 +116,519 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { } }) + it.each([ + { collectionCreated: false, hasExistingData: true, message: "Checking for new or modified files..." }, + { collectionCreated: false, hasExistingData: false, message: "Services ready. Starting workspace scan..." }, + { collectionCreated: true, hasExistingData: true, message: "Services ready. Starting workspace scan..." }, + ])( + "completes scanning with $collectionCreated / $hasExistingData", + async ({ collectionCreated, hasExistingData, message }) => { + vectorStore.initialize.mockResolvedValue(collectionCreated) + vectorStore.hasIndexedData.mockResolvedValue(hasExistingData) + scanner.scanDirectory.mockImplementation( + async ( + _dir: string, + _onBatchError: (error: Error) => void, + onBlocksIndexed: (count: number) => void, + onFileParsed: (count: number) => void, + ) => { + onFileParsed(3) + onBlocksIndexed(2) + onBlocksIndexed(1) + return { stats: { processed: 1, skipped: 0 }, totalBlockCount: 3 } + }, + ) + + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + + await orchestrator.startIndexing() + + expect(stateManager.setSystemState).toHaveBeenCalledWith("Indexing", message) + expect(stateManager.reportBlockIndexingProgress.mock.calls).toEqual([ + [0, 3], + [2, 3], + [3, 3], + ]) + expect(scanner.scanDirectory).toHaveBeenCalledWith( + workspacePath, + expect.any(Function), + expect.any(Function), + expect.any(Function), + expect.any(AbortSignal), + ) + expect(vectorStore.markIndexingIncomplete).toHaveBeenCalledTimes(1) + expect(fileWatcher.initialize).toHaveBeenCalledTimes(1) + expect(vectorStore.markIndexingComplete).toHaveBeenCalledTimes(1) + expect(stateManager.state).toBe("Indexed") + expect(vectorStore.clearCollection).not.toHaveBeenCalled() + expect(cacheManager.clearCacheFile).toHaveBeenCalledTimes(collectionCreated ? 1 : 0) + }, + ) + + it.each([ + { found: 0, indexed: 0, batchError: false, expected: "Indexed" }, + { found: 3, indexed: 0, batchError: false, expected: "Error" }, + { found: 3, indexed: 0, batchError: true, expected: "Error" }, + { found: 0, indexed: 0, batchError: true, expected: "Error" }, + { found: 10, indexed: 9, batchError: true, expected: "Indexed" }, + { found: 10, indexed: 8, batchError: true, expected: "Error" }, + { found: 10, indexed: 8, batchError: false, expected: "Indexed" }, + ])( + "preserves full-scan validation for $indexed/$found blocks, batch error: $batchError", + async ({ found, indexed, batchError, expected }) => { + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockResolvedValue(false) + scanner.scanDirectory.mockImplementation( + async ( + _dir: string, + onError: (error: Error) => void, + onBlocksIndexed: (count: number) => void, + onFileParsed: (count: number) => void, + ) => { + onFileParsed(found) + onBlocksIndexed(indexed) + if (batchError) onError(new Error("batch failure")) + return { stats: { processed: 1, skipped: 0 }, totalBlockCount: found } + }, + ) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + + await orchestrator.startIndexing() + + expect(orchestrator.state).toBe(expected) + expect(vectorStore.markIndexingComplete).toHaveBeenCalledTimes(expected === "Indexed" ? 1 : 0) + }, + ) + + it.each(["success", "batch error", "file error"])( + "preserves %s after trailing progress with the real state manager", + async (outcome) => { + const realState = new CodeIndexStateManager() + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockResolvedValue(true) + scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) + const orchestrator = new CodeIndexOrchestrator( + configManager, + realState, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + try { + await orchestrator.startIndexing() + const onProgress = fileWatcher.onBatchProgressUpdate.mock.calls[0][0] + const onFinished = fileWatcher.onDidFinishBatchProcessing.mock.calls[0][0] + onProgress({ processedInBatch: 0, totalInBatch: 1, currentFile: "/workspace/test.ts" }) + expect(realState.getCurrentStatus()).toMatchObject({ + systemStatus: "Indexing", + totalItems: 1, + currentItemUnit: "files", + }) + onFinished({ + processedFiles: [{ path: "test.ts", status: outcome === "file error" ? "error" : "success" }], + batchError: outcome === "batch error" ? new Error("batch failed") : undefined, + }) + const terminalStatus = realState.getCurrentStatus() + expect(terminalStatus.systemStatus).toBe(outcome === "success" ? "Indexed" : "Error") + onProgress({ processedInBatch: 1, totalInBatch: 1 }) + expect(realState.getCurrentStatus()).toEqual(terminalStatus) + onProgress({ processedInBatch: 0, totalInBatch: 0 }) + expect(realState.getCurrentStatus()).toEqual(terminalStatus) + await orchestrator.startIndexing() + expect(scanner.scanDirectory).toHaveBeenCalledTimes(2) + } finally { + orchestrator.stopWatcher() + realState.dispose() + } + }, + ) + + it("should handle watcher progress and completion through registered callbacks", async () => { + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockResolvedValue(false) + scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + const onProgress = fileWatcher.onBatchProgressUpdate.mock.calls[0][0] + const onFinished = fileWatcher.onDidFinishBatchProcessing.mock.calls[0][0] + + onProgress({ processedInBatch: 1, totalInBatch: 2, currentFile: "/test/workspace/example.ts" }) + expect(orchestrator.state).toBe("Indexing") + expect(stateManager.reportFileQueueProgress).toHaveBeenLastCalledWith(1, 2, "example.ts") + onProgress({ processedInBatch: 1, totalInBatch: 2 }) + expect(stateManager.reportFileQueueProgress).toHaveBeenLastCalledWith(1, 2, undefined) + onProgress({ processedInBatch: 2, totalInBatch: 2 }) + expect(orchestrator.state).toBe("Indexing") + onFinished({ processedFiles: [{ path: "example.ts", status: "success" }] }) + expect(orchestrator.state).toBe("Indexed") + onProgress({ processedInBatch: 2, totalInBatch: 2 }) + expect(orchestrator.state).toBe("Indexed") + onProgress({ processedInBatch: 0, totalInBatch: 1 }) + onFinished({ processedFiles: [] }) + onProgress({ processedInBatch: 0, totalInBatch: 0 }) + expect(orchestrator.state).toBe("Indexed") + + const error = new Error("watcher batch failed") + const log = vi.spyOn(console, "error").mockImplementation(() => {}) + try { + onFinished({ processedFiles: [], batchError: error }) + expect(log).toHaveBeenCalledWith("[CodeIndexWatcherSession] Batch processing failed:", error) + onProgress({ processedInBatch: 1, totalInBatch: 1 }) + onProgress({ processedInBatch: 0, totalInBatch: 0 }) + expect(orchestrator.state).toBe("Error") + onFinished({ processedFiles: [] }) + } finally { + log.mockRestore() + } + }) + + it.each(["error", "local_error"])( + "should report watcher file failures without a batch error (%s)", + async (status) => { + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockResolvedValue(true) + scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + const onFinished = fileWatcher.onDidFinishBatchProcessing.mock.calls[0][0] + const onProgress = fileWatcher.onBatchProgressUpdate.mock.calls[0][0] + onFinished({ processedFiles: [{ path: "bad.ts", status }] }) + onProgress({ processedInBatch: 1, totalInBatch: 1 }) + onProgress({ processedInBatch: 0, totalInBatch: 0 }) + expect(orchestrator.state).toBe("Error") + onProgress({ processedInBatch: 0, totalInBatch: 1 }) + onFinished({ processedFiles: [{ path: "bad.ts", status: "success" }] }) + expect(orchestrator.state).toBe("Indexed") + }, + ) + + it("should reuse active watcher subscriptions on repeated indexing", async () => { + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockResolvedValue(true) + scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) + const disposeProgress = vi.fn() + const disposeFinished = vi.fn() + fileWatcher.onBatchProgressUpdate.mockReturnValue({ dispose: disposeProgress }) + fileWatcher.onDidFinishBatchProcessing.mockReturnValue({ dispose: disposeFinished }) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + await orchestrator.startIndexing() + expect(fileWatcher.initialize).toHaveBeenCalledTimes(1) + expect(fileWatcher.onBatchProgressUpdate).toHaveBeenCalledTimes(1) + expect(fileWatcher.onDidFinishBatchProcessing).toHaveBeenCalledTimes(1) + orchestrator.stopWatcher() + expect(disposeProgress).toHaveBeenCalledTimes(1) + expect(disposeFinished).toHaveBeenCalledTimes(1) + }) + + it.each([undefined, []])("rejects indexing without workspace folders (%s)", async (folders) => { + const original = vscode.workspace.workspaceFolders + Object.defineProperty(vscode.workspace, "workspaceFolders", { value: folders, configurable: true }) + try { + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Error") + expect(vectorStore.initialize).not.toHaveBeenCalled() + } finally { + Object.defineProperty(vscode.workspace, "workspaceFolders", { value: original, configurable: true }) + } + }) + + it("rejects indexing without configuration", async () => { + configManager.isFeatureConfigured = false + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Standby") + expect(vectorStore.initialize).not.toHaveBeenCalled() + }) + + it.each(["Indexing", "Stopping"])("rejects indexing in %s state", async (state) => { + stateManager.setSystemState(state, "busy") + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + expect(orchestrator.state).toBe(state) + expect(vectorStore.initialize).not.toHaveBeenCalled() + }) + + it.each(["configuration lost", "watcher failed"])("handles watcher startup failure: %s", async (failure) => { + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockResolvedValue(true) + scanner.scanDirectory.mockImplementation(async () => { + if (failure === "configuration lost") configManager.isFeatureConfigured = false + return { stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 } + }) + fileWatcher.initialize.mockRejectedValue(new Error("watcher failed")) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Error") + expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() + expect(fileWatcher.dispose).toHaveBeenCalled() + expect(TelemetryService.instance.captureEvent).toHaveBeenCalledWith( + expect.any(String), + expect.objectContaining({ location: "startIndexing" }), + ) + expect(TelemetryService.instance.captureEvent).toHaveBeenCalledTimes(1) + if (failure === "configuration lost") { + expect(fileWatcher.initialize).not.toHaveBeenCalled() + } + }) + + it("preserves the original error when collection cleanup rejects a non-Error value", async () => { + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockResolvedValue(false) + scanner.scanDirectory.mockRejectedValue(new Error("original failure")) + vectorStore.clearCollection.mockRejectedValue("cleanup failed") + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + expect(cacheManager.clearCacheFile).toHaveBeenCalledOnce() + expect(stateManager.setSystemState).toHaveBeenLastCalledWith( + "Error", + expect.stringContaining("original failure"), + ) + expect(TelemetryService.instance.captureEvent).toHaveBeenCalledWith( + expect.any(String), + expect.objectContaining({ location: "startIndexing.cleanup", error: "cleanup failed", stack: undefined }), + ) + }) + + it.each([null, "connection failed", { message: "" }])("handles unexpected rejection values (%s)", async (error) => { + vectorStore.initialize.mockRejectedValue(error) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await expect(orchestrator.startIndexing()).resolves.toBeUndefined() + expect(stateManager.setSystemState).toHaveBeenLastCalledWith("Error", expect.stringContaining("unknownError")) + expect(fileWatcher.dispose).toHaveBeenCalled() + }) + + it("clears cache without deleting storage when configuration is missing", async () => { + configManager.isFeatureConfigured = false + vectorStore.deleteCollection = vi.fn() + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + const stopWatcher = vi.spyOn(orchestrator, "stopWatcher") + await orchestrator.clearIndexData() + expect(stopWatcher).not.toHaveBeenCalled() + expect(vectorStore.deleteCollection).not.toHaveBeenCalled() + expect(cacheManager.clearCacheFile).toHaveBeenCalledOnce() + expect(orchestrator.state).toBe("Standby") + }) + + it("reports non-Error deletion failures and releases the clearing guard", async () => { + vectorStore.deleteCollection = vi.fn().mockRejectedValueOnce("delete failed").mockResolvedValue(undefined) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.clearIndexData() + expect(stateManager.setSystemState).toHaveBeenLastCalledWith( + "Error", + "Failed to clear index data: delete failed", + ) + expect(cacheManager.clearCacheFile).not.toHaveBeenCalled() + await orchestrator.clearIndexData() + expect(orchestrator.state).toBe("Standby") + }) + + it("should preserve existing data when checking collection contents fails", async () => { + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockRejectedValue(new Error("query failed")) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Error") + expect(vectorStore.clearCollection).not.toHaveBeenCalled() + expect(cacheManager.clearCacheFile).not.toHaveBeenCalled() + expect(scanner.scanDirectory).not.toHaveBeenCalled() + }) + + it.each([false, true])( + "never clears data when point presence is unknown (created/recreated: %s)", + async (created) => { + vectorStore.initialize.mockResolvedValue(created) + vectorStore.hasCodePoints.mockRejectedValue(new Error("point query failed")) + vectorStore.hasIndexedData.mockResolvedValue(false) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Error") + expect(vectorStore.clearCollection).not.toHaveBeenCalled() + expect(cacheManager.clearCacheFile).not.toHaveBeenCalled() + expect(scanner.scanDirectory).not.toHaveBeenCalled() + expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() + }, + ) + + it.each([false, true])("cleans a failed full scan that started empty (created/recreated: %s)", async (created) => { + const points = new Set() + const cache = new Map() + vectorStore.initialize.mockResolvedValue(created) + vectorStore.hasCodePoints.mockImplementation(async () => points.size > 0) + vectorStore.hasIndexedData.mockResolvedValue(false) + vectorStore.clearCollection.mockImplementation(async () => { + points.clear() + }) + cacheManager.clearCacheFile.mockImplementation(async () => { + cache.clear() + }) + scanner.scanDirectory.mockImplementation(async () => { + points.add("partial-block") + cache.set("partial.ts", "partial-hash") + throw new Error("scan failed") + }) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Error") + expect(points.size).toBe(0) + expect(cache.size).toBe(0) + expect(vectorStore.clearCollection).toHaveBeenCalledOnce() + expect(cacheManager.clearCacheFile).toHaveBeenCalledTimes(created ? 2 : 1) + expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() + }) + + it("should finish error handling when cache cleanup fails", async () => { + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockResolvedValue(false) + scanner.scanDirectory.mockRejectedValue(new Error("scan failed")) + cacheManager.clearCacheFile.mockRejectedValue(new Error("cache cleanup failed")) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + + await expect(orchestrator.startIndexing()).resolves.toBeUndefined() + expect(orchestrator.state).toBe("Error") + expect(fileWatcher.dispose).toHaveBeenCalled() + expect(stateManager.setSystemState).toHaveBeenLastCalledWith("Error", expect.stringContaining("scan failed")) + }) + it("should not call clearCollection() or clear cache when initialize() fails (indexing not started)", async () => { // Arrange: fail at initialize() vectorStore.initialize.mockRejectedValue(new Error("Qdrant unreachable")) @@ -193,36 +715,108 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => { expect(calls[calls.length - 1]).toBe("Error") }) - it("collects batch errors from incremental scan and still completes indexing", async () => { - const batchError = new Error("incremental batch failure") - vectorStore.initialize.mockResolvedValue(false) // existing collection - vectorStore.hasIndexedData.mockResolvedValue(true) // force incremental scan path - vectorStore.markIndexingIncomplete.mockResolvedValue(undefined) - vectorStore.markIndexingComplete.mockResolvedValue(undefined) + it.each([false, true])( + "preserves stored points and cache after an incremental failure and a failed retry (new orchestrator: %s)", + async (recreate) => { + const points = new Set(["existing-code-block"]) + const cache = new Map([["existing.ts", "unchanged-hash"]]) + let complete = true + vectorStore.initialize.mockResolvedValue(false) + vectorStore.hasIndexedData.mockImplementation(async () => points.size > 0 && complete) + vectorStore.hasCodePoints = vi.fn(async () => points.size > 0) + vectorStore.markIndexingIncomplete.mockImplementation(async () => { + complete = false + }) + vectorStore.clearCollection.mockImplementation(async () => { + points.clear() + }) + cacheManager.clearCacheFile.mockImplementation(async () => { + cache.clear() + }) + scanner.scanDirectory.mockImplementation(async (_dir: string, onError: (error: Error) => void) => { + onError(new Error("embedding unavailable")) + return { stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 } + }) + let orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Error") + expect(complete).toBe(false) + expect(points.size).toBe(1) + if (recreate) { + orchestrator = new CodeIndexOrchestrator( + configManager, + new CodeIndexStateManager(), + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + } + await orchestrator.startIndexing() + expect(scanner.scanDirectory).toHaveBeenCalledTimes(2) + expect(orchestrator.state).toBe("Error") + expect(points).toEqual(new Set(["existing-code-block"])) + expect(cache).toEqual(new Map([["existing.ts", "unchanged-hash"]])) + expect(vectorStore.clearCollection).not.toHaveBeenCalled() + expect(cacheManager.clearCacheFile).not.toHaveBeenCalled() + expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() + }, + ) - // Incremental scan reports a batch error but returns a result — orchestrator completes normally - scanner.scanDirectory.mockImplementation(async (_dir: string, onBatchError: (e: Error) => void) => { - onBatchError(batchError) - return { stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 } - }) + it.each([0, 2])( + "should report incremental batch failure without deleting existing data (%s blocks indexed)", + async (indexedCount) => { + const batchError = new Error("incremental batch failure") + vectorStore.initialize.mockResolvedValue(false) // existing collection + vectorStore.hasIndexedData.mockResolvedValue(true) // force incremental scan path + vectorStore.markIndexingIncomplete.mockResolvedValue(undefined) + vectorStore.markIndexingComplete.mockResolvedValue(undefined) - const orchestrator = new CodeIndexOrchestrator( - configManager, - stateManager, - workspacePath, - cacheManager, - vectorStore, - scanner, - fileWatcher, - ) + // A failed batch must prevent success even if another batch was indexed. + scanner.scanDirectory.mockImplementation( + async ( + _dir: string, + onBatchError: (error: Error) => void, + onBlocksIndexed: (count: number) => void, + onFileParsed: (count: number) => void, + ) => { + onFileParsed(3) + if (indexedCount > 0) { + onBlocksIndexed(indexedCount) + } + onBatchError(batchError) + return { stats: { processed: 1, skipped: 0 }, totalBlockCount: 3 } + }, + ) - await orchestrator.startIndexing() + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) - // Incremental scan doesn't gate on batch errors — Indexed state is still reached - const calls = stateManager.setSystemState.mock.calls.map((c: any[]) => c[0]) - expect(calls[calls.length - 1]).toBe("Indexed") - expect(calls).not.toContain("Error") - }) + await orchestrator.startIndexing() + + expect(orchestrator.state).toBe("Error") + expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() + expect(stateManager.setSystemState).not.toHaveBeenCalledWith("Indexed", expect.any(String)) + expect(vectorStore.clearCollection).not.toHaveBeenCalled() + expect(cacheManager.clearCacheFile).not.toHaveBeenCalled() + }, + ) }) describe("CodeIndexOrchestrator - stopIndexing", () => { @@ -261,6 +855,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { vectorStore = { initialize: vi.fn().mockResolvedValue(false), + hasCodePoints: vi.fn().mockResolvedValue(false), hasIndexedData: vi.fn().mockResolvedValue(false), markIndexingIncomplete: vi.fn().mockResolvedValue(undefined), markIndexingComplete: vi.fn().mockResolvedValue(undefined), @@ -280,6 +875,403 @@ describe("CodeIndexOrchestrator - stopIndexing", () => { } }) + it("preserves active indexing state when a repeated start finds missing configuration", async () => { + let releasePreparation!: () => void + const preparation = new Promise((resolve) => { + releasePreparation = resolve + }) + vectorStore.initialize.mockImplementationOnce(async () => { + await preparation + return false + }) + 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() + try { + configManager.isFeatureConfigured = false + stateManager.setSystemState.mockClear() + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Indexing") + expect(stateManager.setSystemState).not.toHaveBeenCalled() + expect(vectorStore.initialize).toHaveBeenCalledTimes(1) + } finally { + configManager.isFeatureConfigured = true + releasePreparation() + await indexing + } + }) + + it("retains a cancelled run until cache flushing finishes and allows a fresh run afterwards", async () => { + let finishPreparation!: () => void + const preparation = new Promise((resolve) => { + finishPreparation = resolve + }) + let notifyFlushing!: () => void + const flushing = new Promise((resolve) => { + notifyFlushing = resolve + }) + let finishFlushing!: () => void + const flushed = new Promise((resolve) => { + finishFlushing = resolve + }) + vectorStore.initialize.mockImplementationOnce(async () => { + await preparation + return false + }) + cacheManager.flush.mockImplementationOnce(async () => { + notifyFlushing() + await flushed + }) + 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() + orchestrator.stopIndexing() + orchestrator.stopIndexing() + finishPreparation() + await flushing + try { + // A visible state update must not release ownership while cleanup is pending. + stateManager.setSystemState("Standby", "External status update") + await orchestrator.startIndexing() + expect(vectorStore.initialize).toHaveBeenCalledTimes(1) + expect(vectorStore.clearCollection).not.toHaveBeenCalled() + expect(fileWatcher.initialize).not.toHaveBeenCalled() + } finally { + finishFlushing() + await indexing + } + + await orchestrator.startIndexing() + expect(vectorStore.initialize).toHaveBeenCalledTimes(2) + expect(orchestrator.state).toBe("Indexed") + }) + + it("releases a watcher created after cancellation before deleting index data", async () => { + let notifyStarting!: () => void + const starting = new Promise((resolve) => { + notifyStarting = resolve + }) + let releaseInitialization!: () => void + const initialization = new Promise((resolve) => { + releaseInitialization = resolve + }) + const events: string[] = [] + let watcherAlive = false + fileWatcher.initialize.mockImplementation(async () => { + notifyStarting() + await initialization + watcherAlive = true + events.push("watcher created") + }) + fileWatcher.dispose.mockImplementation(() => { + if (watcherAlive) events.push("watcher released") + watcherAlive = false + }) + const aliveDuringDeletion: boolean[] = [] + vectorStore.deleteCollection = vi.fn().mockImplementation(async () => { + aliveDuringDeletion.push(watcherAlive) + events.push("collection deleted") + }) + 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 starting + const clearing = orchestrator.clearIndexData() + try { + await Promise.resolve() + expect(vectorStore.deleteCollection).not.toHaveBeenCalled() + } finally { + releaseInitialization() + await Promise.all([indexing, clearing]) + } + + expect(aliveDuringDeletion).toEqual([false]) + expect(events).toEqual(["watcher created", "watcher released", "collection deleted"]) + expect(watcherAlive).toBe(false) + expect(fileWatcher.onBatchProgressUpdate).not.toHaveBeenCalled() + expect(fileWatcher.onDidFinishBatchProcessing).not.toHaveBeenCalled() + expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() + expect(orchestrator.state).toBe("Standby") + }) + + it("should abort and await an active scan before deleting index data", async () => { + const events: string[] = [] + let notifyScanStarted!: () => void + const scanStarted = new Promise((resolve) => { + notifyScanStarted = resolve + }) + let releaseScan!: () => void + const scanReleased = new Promise((resolve) => { + releaseScan = resolve + }) + let notifyClearAction!: () => void + const clearAction = new Promise((resolve) => { + notifyClearAction = resolve + }) + let scanSignal: AbortSignal | undefined + + vectorStore.deleteCollection = vi.fn().mockImplementation(async () => { + events.push("collection deleted") + notifyClearAction() + }) + scanner.scanDirectory.mockImplementation( + async ( + _dir: string, + _onError?: (error: Error) => void, + _onBlocksIndexed?: (count: number) => void, + _onFileParsed?: (count: number) => void, + signal?: AbortSignal, + ) => { + scanSignal = signal + signal?.addEventListener("abort", notifyClearAction, { once: true }) + notifyScanStarted() + // Even after cancellation, an in-flight scan needs time to settle. + await scanReleased + signal?.removeEventListener("abort", notifyClearAction) + events.push("scan finished") + return { stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 } + }, + ) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + + const indexing = orchestrator.startIndexing() + await scanStarted + const clearing = orchestrator.clearIndexData() + // Either cancellation (correct) or deletion (bug) lets the test proceed. + await clearAction + releaseScan() + await Promise.all([indexing, clearing]) + + expect(events).toEqual(["scan finished", "collection deleted"]) + expect(scanSignal?.aborted).toBe(true) + expect(fileWatcher.initialize).not.toHaveBeenCalled() + expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() + expect(orchestrator.state).toBe("Standby") + }) + + it("should reject indexing while collection deletion is pending and allow it after clearing", async () => { + let notifyDeletionStarted!: () => void + const deletionStarted = new Promise((resolve) => { + notifyDeletionStarted = resolve + }) + let finishDeletion!: () => void + const deletionFinished = new Promise((resolve) => { + finishDeletion = resolve + }) + vectorStore.deleteCollection = vi.fn().mockImplementation(async () => { + notifyDeletionStarted() + await deletionFinished + }) + scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + + const clearing = orchestrator.clearIndexData() + await deletionStarted + try { + await orchestrator.clearIndexData() + expect(vectorStore.deleteCollection).toHaveBeenCalledTimes(1) + await orchestrator.startIndexing() + expect(vectorStore.initialize).not.toHaveBeenCalled() + expect(cacheManager.clearCacheFile).not.toHaveBeenCalled() + } finally { + finishDeletion() + await clearing + } + expect(cacheManager.clearCacheFile).toHaveBeenCalledTimes(1) + await orchestrator.startIndexing() + expect(orchestrator.state).toBe("Indexed") + }) + + it("should preserve cache on deletion failure and allow clearing to be retried", async () => { + vectorStore.deleteCollection = vi + .fn() + .mockRejectedValueOnce(new Error("delete failed")) + .mockResolvedValue(undefined) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + + await orchestrator.clearIndexData() + expect(orchestrator.state).toBe("Error") + expect(cacheManager.clearCacheFile).not.toHaveBeenCalled() + + await orchestrator.clearIndexData() + expect(orchestrator.state).toBe("Standby") + expect(cacheManager.clearCacheFile).toHaveBeenCalledTimes(1) + }) + + it.each([false, true])( + "should remain stopped when watcher initialization finishes after cancellation (existing data: %s)", + async (hasExistingData) => { + vectorStore.hasIndexedData.mockResolvedValue(hasExistingData) + let notifyWatcherStarting!: () => void + const watcherStarting = new Promise((resolve) => { + notifyWatcherStarting = resolve + }) + let finishWatcherInitialization!: () => void + const watcherInitialization = new Promise((resolve) => { + finishWatcherInitialization = resolve + }) + scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 }) + fileWatcher.initialize.mockImplementation(async () => { + notifyWatcherStarting() + await watcherInitialization + }) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + + const indexing = orchestrator.startIndexing() + await watcherStarting + orchestrator.stopIndexing() + finishWatcherInitialization() + await indexing + + expect(TelemetryService.instance.captureEvent).not.toHaveBeenCalled() + expect(orchestrator.state).toBe("Standby") + expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled() + expect(stateManager.setSystemState).not.toHaveBeenCalledWith("Indexed", expect.any(String)) + }, + ) + + it.each([false, true])( + "should restore the incomplete marker when cancelled during completion persistence (existing data: %s)", + async (hasExistingData) => { + vectorStore.hasIndexedData.mockResolvedValue(hasExistingData) + let notifySaving!: () => void + const saving = new Promise((resolve) => { + notifySaving = resolve + }) + let finishSaving!: () => void + const saved = new Promise((resolve) => { + finishSaving = resolve + }) + let complete = false + vectorStore.markIndexingIncomplete.mockImplementation(async () => { + complete = false + }) + vectorStore.markIndexingComplete.mockImplementation(async () => { + notifySaving() + await saved + complete = true + }) + 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 saving + orchestrator.stopIndexing() + finishSaving() + await indexing + + expect(complete).toBe(false) + expect(orchestrator.state).toBe("Standby") + expect(stateManager.setSystemState).not.toHaveBeenCalledWith("Indexed", expect.any(String)) + }, + ) + + it.each([false, true])( + "should finish cancellation when cache flush fails (scanner throws: %s)", + async (throwsAbort) => { + let notifyStarted!: () => void + const started = new Promise((resolve) => { + notifyStarted = resolve + }) + let releaseScan!: () => void + const released = new Promise((resolve) => { + releaseScan = resolve + }) + cacheManager.flush.mockRejectedValue(new Error("flush failed")) + scanner.scanDirectory.mockImplementation(async () => { + notifyStarted() + await released + if (throwsAbort) throw new DOMException("Stopped", "AbortError") + return { stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 } + }) + const orchestrator = new CodeIndexOrchestrator( + configManager, + stateManager, + workspacePath, + cacheManager, + vectorStore, + scanner, + fileWatcher, + ) + const indexing = orchestrator.startIndexing() + await started + orchestrator.stopIndexing() + fileWatcher.dispose.mockClear() + releaseScan() + + await expect(indexing).resolves.toBeUndefined() + expect(orchestrator.state).toBe("Standby") + expect(fileWatcher.dispose).toHaveBeenCalled() + expect(cacheManager.flush).toHaveBeenCalledTimes(1) + expect(vectorStore.clearCollection).not.toHaveBeenCalled() + expect(cacheManager.clearCacheFile).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-recovery.ts b/src/services/code-index/code-index-recovery.ts new file mode 100644 index 0000000000..a56c52a69c --- /dev/null +++ b/src/services/code-index/code-index-recovery.ts @@ -0,0 +1,93 @@ +import type { CacheManager } from "./cache-manager" +import type { IVectorStore } from "./interfaces" +import type { CodeIndexStateManager } from "./state-manager" +import type { CodeIndexWatcherSession } from "./code-index-watcher-session" +import type { CodeIndexRun } from "./code-index-run" +import { TelemetryService } from "@roo-code/telemetry" +import { TelemetryEventName } from "@roo-code/types" +import { t } from "../../i18n" + +/** Restores a consistent index state after cancellation or failure. Does not release run ownership. */ +export class CodeIndexRecovery { + constructor( + private readonly cacheManager: Pick, + private readonly vectorStore: Pick, + private readonly stateManager: Pick, + private readonly watcherSession: Pick, + ) {} + + async handle(error: unknown, run: CodeIndexRun): Promise { + if (this._isCancellation(error, run.signal)) { + await this._finishCancellation() + return + } + + this._reportError("[CodeIndexOrchestrator] Error during indexing:", error, "startIndexing") + if (run.canCleanupFailedScan) { + await this._cleanupFailedFullScan() + } else { + console.log("[CodeIndexOrchestrator] Preserving existing index and cache for a future incremental scan.") + } + + const errorMessage = + typeof error === "object" && error !== null && "message" in error && error.message + ? error.message + : t("embeddings:orchestrator.unknownError") + this.stateManager.setSystemState( + "Error", + t("embeddings:orchestrator.failedDuringInitialScan", { errorMessage }), + ) + this.watcherSession.stop() + } + + private async _finishCancellation(): Promise { + console.log("[CodeIndexOrchestrator] Indexing aborted by user.") + try { + await this.cacheManager.flush() + } catch (flushError) { + console.error("[CodeIndexOrchestrator] Failed to flush cache after cancellation:", flushError) + } + this.watcherSession.stop() + this.stateManager.setSystemState("Standby", t("embeddings:orchestrator.indexingStopped")) + } + + private async _cleanupFailedFullScan(): Promise { + try { + await this.vectorStore.clearCollection() + } catch (cleanupError) { + this._reportError( + "[CodeIndexOrchestrator] Failed to clean up after error:", + cleanupError, + "startIndexing.cleanup", + ) + } + try { + await this.cacheManager.clearCacheFile() + } catch (cleanupError) { + console.error("[CodeIndexOrchestrator] Failed to clear cache after indexing error:", cleanupError) + } + console.log("[CodeIndexOrchestrator] Indexing failed after starting. Clearing cache to avoid inconsistency.") + } + + handleClearError(error: unknown): void { + const message = error instanceof Error ? error.message : String(error) + this._reportError("[CodeIndexOrchestrator] Failed to clear index data:", error, "clearIndexData") + this.stateManager.setSystemState("Error", `Failed to clear index data: ${message}`) + } + + private _isCancellation(error: unknown, signal: AbortSignal): boolean { + return ( + signal.aborted || + (typeof error === "object" && error !== null && "name" in error && error.name === "AbortError") + ) + } + + private _reportError(message: string, error: unknown, location: string): void { + console.error(message, error) + TelemetryService.instance.captureEvent(TelemetryEventName.CODE_INDEX_ERROR, { + error: error instanceof Error ? error.message : String(error), + stack: error instanceof Error ? error.stack : undefined, + location, + }) + } +} diff --git a/src/services/code-index/code-index-run.ts b/src/services/code-index/code-index-run.ts new file mode 100644 index 0000000000..f123122a9c --- /dev/null +++ b/src/services/code-index/code-index-run.ts @@ -0,0 +1,52 @@ +import { StateHolder, StateStream } from "../../utils/StateHolder" + +export type CodeIndexScanMode = "incremental" | "full" +export type CodeIndexRunState = "running" | "cancelling" | "finished" + +/** Owns cancellation and completion for one indexing attempt, including its cleanup. */ +export class CodeIndexRun { + private startedScanMode: CodeIndexScanMode | undefined + /** Unknown until the initialized collection has been queried successfully. */ + preexistingCodePoints: boolean | undefined + + constructor( + private readonly controller: AbortController, + private readonly codeIndexRunStateHolder: StateHolder, + ) {} + + get state(): StateStream { + return this.codeIndexRunStateHolder + } + + get signal(): AbortSignal { + return this.controller.signal + } + + get fullScanStarted(): boolean { + return this.startedScanMode === "full" + } + + get canCleanupFailedScan(): boolean { + return this.fullScanStarted && this.preexistingCodePoints === false + } + + /** Record this immediately before invoking the scanner, not during preparation. */ + markScanStarted(mode: CodeIndexScanMode): void { + this.startedScanMode = mode + } + + /** Cancellation does not notify completion listeners; cleanup still belongs to this run. */ + cancel(): void { + if (this.codeIndexRunStateHolder.value !== "running") return + this.codeIndexRunStateHolder.set("cancelling") + this.controller.abort() + } + + async waitUntilFinished(): Promise { + await this.codeIndexRunStateHolder.waitFor((state) => state === "finished") + } + + finish(): void { + this.codeIndexRunStateHolder.set("finished") + } +} diff --git a/src/services/code-index/code-index-scan-executor.ts b/src/services/code-index/code-index-scan-executor.ts new file mode 100644 index 0000000000..f6459178af --- /dev/null +++ b/src/services/code-index/code-index-scan-executor.ts @@ -0,0 +1,118 @@ +import type { IDirectoryScanner, IVectorStore } from "./interfaces" +import type { CodeIndexStateManager } from "./state-manager" +import { t } from "../../i18n" + +/** Executes workspace scans; lifecycle, cleanup and watcher ownership stay with the orchestrator. */ +export class CodeIndexScanExecutor { + constructor( + private readonly workspacePath: string, + private readonly scanner: IDirectoryScanner, + private readonly vectorStore: Pick, + private readonly stateManager: Pick, + ) {} + + public async runIncrementalScan(signal: AbortSignal): Promise { + // Collection exists with data - run incremental scan to catch any new/changed files + // This handles files added while workspace was closed or Qdrant was inactive + console.log( + "[CodeIndexOrchestrator] Collection already has indexed data. Running incremental scan for new/changed files...", + ) + this.stateManager.setSystemState("Indexing", "Checking for new or modified files...") + + const summary = await this._scanWorkspace(signal, "incremental") + const { indexed, found, batchErrors } = summary + + if (batchErrors.length > 0) { + throw new Error(`Incremental indexing failed: ${batchErrors[0].message}`) + } + + // If new files were found and indexed, log the results + if (found > 0) { + console.log( + `[CodeIndexOrchestrator] Incremental scan completed: ${indexed} blocks indexed from new/changed files`, + ) + } else { + console.log("[CodeIndexOrchestrator] No new or changed files found") + } + } + + public async runFullScan(signal: AbortSignal): Promise { + // No existing data or collection was just created - do a full scan + this.stateManager.setSystemState("Indexing", "Services ready. Starting workspace scan...") + + const summary = await this._scanWorkspace(signal, "full") + + this._validateFullScan(summary.indexed, summary.found, summary.batchErrors) + } + + private async _scanWorkspace( + signal: AbortSignal, + kind: "full" | "incremental", + ): Promise<{ indexed: number; found: number; batchErrors: Error[] }> { + // The scanner uses the cache to skip unchanged files in either scan mode. + await this.vectorStore.markIndexingIncomplete() + + let indexed = 0 + let found = 0 + const batchErrors: Error[] = [] + + const handleFileParsed = (fileBlockCount: number) => { + found += fileBlockCount + this.stateManager.reportBlockIndexingProgress(indexed, found) + } + + const handleBlocksIndexed = (indexedCount: number) => { + indexed += indexedCount + this.stateManager.reportBlockIndexingProgress(indexed, found) + } + + const result = await this.scanner.scanDirectory( + this.workspacePath, + (batchError: Error) => { + console.error( + `[CodeIndexOrchestrator] Error during ${kind === "full" ? "initial" : "incremental"} scan batch: ${batchError.message}`, + batchError, + ) + batchErrors.push(batchError) + }, + handleBlocksIndexed, + handleFileParsed, + signal, + ) + + signal.throwIfAborted() + + if (!result) { + throw new Error( + kind === "full" + ? "Scan failed, is scanner initialized?" + : "Incremental scan failed, is scanner initialized?", + ) + } + + return { indexed, found, batchErrors } + } + + private _validateFullScan(indexed: number, found: number, batchErrors: Error[]): void { + const firstError = batchErrors[0] + + if (indexed === 0 && found > 0) { + throw new Error( + firstError + ? `Indexing failed: ${firstError.message}` + : t("embeddings:orchestrator.indexingFailedNoBlocks"), + ) + } + + if (firstError && indexed === 0) { + throw new Error(`Indexing failed completely: ${firstError.message}`) + } + + // Preserve the full-scan policy: batch errors are fatal above 10% failed blocks. + if (firstError && found > 0 && (found - indexed) / found > 0.1) { + throw new Error( + `Indexing partially failed: Only ${indexed} of ${found} blocks were indexed. ${firstError.message}`, + ) + } + } +} 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..dfacdb3d0a --- /dev/null +++ b/src/services/code-index/code-index-watcher-session.ts @@ -0,0 +1,109 @@ +import type { Disposable } from "vscode" +import * as path from "path" +import type { CodeIndexStateManager } from "./state-manager" +import type { IFileWatcher, BatchProcessingSummary } from "./interfaces" + +/** Owns file watching, subscriptions and batch progress reporting. */ +export class CodeIndexWatcherSession { + private subscriptions: Disposable[] = [] + private phase: "stopped" | "starting" | "running" = "stopped" + private generation = 0 + + constructor( + private readonly fileWatcher: IFileWatcher, + private readonly stateManager: Pick< + CodeIndexStateManager, + "state" | "setSystemState" | "reportFileQueueProgress" + >, + ) {} + + get isRunning(): boolean { + return this.phase === "running" + } + + async start(signal: AbortSignal): Promise { + signal.throwIfAborted() + if (this.isRunning) return + if (this.phase === "starting") throw new Error("File watcher startup is already in progress.") + const generation = this.generation + this.phase = "starting" + try { + await this.fileWatcher.initialize() + signal.throwIfAborted() + if (generation !== this.generation) { + throw new DOMException("File watcher startup was stopped.", "AbortError") + } + this.subscriptions.push( + this.fileWatcher.onBatchProgressUpdate((progress) => this.handleBatchProgress(progress)), + ) + this.subscriptions.push( + this.fileWatcher.onDidFinishBatchProcessing((summary) => this.handleBatchFinished(summary)), + ) + this.phase = "running" + } catch (error) { + this.releaseResources() + this.phase = "stopped" + throw error + } + } + + private handleBatchProgress({ + processedInBatch, + totalInBatch, + currentFile, + }: { + processedInBatch: number + totalInBatch: number + currentFile?: string + }): void { + // Reporting terminal progress would reset the batch's final state to Indexing. + if (processedInBatch >= totalInBatch) return + if (this.stateManager.state !== "Indexing") { + this.stateManager.setSystemState("Indexing", "Processing file changes...") + } + this.stateManager.reportFileQueueProgress( + processedInBatch, + totalInBatch, + currentFile ? path.basename(currentFile) : undefined, + ) + } + + private handleBatchFinished(summary: BatchProcessingSummary): void { + if (summary.batchError) { + console.error("[CodeIndexWatcherSession] Batch processing failed:", summary.batchError) + this.stateManager.setSystemState("Error", summary.batchError.message) + return + } + const failedFile = summary.processedFiles.find( + (file) => file.status === "error" || file.status === "local_error", + ) + if (failedFile) { + this.stateManager.setSystemState( + "Error", + failedFile.error?.message ?? `Failed to index file: ${failedFile.path}`, + ) + return + } + this.stateManager.setSystemState("Indexed", "File changes processed. Index up-to-date.") + } + + stop(): void { + this.generation++ + // Keep startup exclusive until its await settles, even after a stop request. + if (this.phase !== "starting") this.phase = "stopped" + this.releaseResources() + } + + private releaseResources(): void { + const resources = [...this.subscriptions, this.fileWatcher] + this.subscriptions = [] + for (const resource of resources) { + try { + resource.dispose() + } catch (error) { + // Cleanup must attempt every resource and preserve the original startup error. + console.error("[CodeIndexWatcherSession] Failed to dispose watcher resource:", error) + } + } + } +} diff --git a/src/services/code-index/interfaces/vector-store.ts b/src/services/code-index/interfaces/vector-store.ts index 7946563fd5..5df109b5ae 100644 --- a/src/services/code-index/interfaces/vector-store.ts +++ b/src/services/code-index/interfaces/vector-store.ts @@ -10,7 +10,7 @@ export type PointStruct = { export interface IVectorStore { /** * Initializes the vector store - * @returns Promise resolving to boolean indicating if a new collection was created + * @returns Whether a new collection was created, including recreation after a dimension change */ initialize(): Promise @@ -64,11 +64,17 @@ export interface IVectorStore { collectionExists(): Promise /** - * Checks if the collection exists and has indexed points - * @returns Promise resolving to boolean indicating if the collection exists and has points + * Checks index readiness using completion metadata (or legacy point count without a marker). + * Returns false for incomplete indexes or read errors; not suitable for authorizing cleanup. */ hasIndexedData(): Promise + /** + * Checks for non-metadata points in an initialized collection, regardless of completion status. + * Rejects on query failure: unknown contents must not be treated as an empty collection. + */ + hasCodePoints(): Promise + /** * Marks the indexing process as complete by storing metadata * Should be called after a successful full workspace scan or incremental scan diff --git a/src/services/code-index/orchestrator.ts b/src/services/code-index/orchestrator.ts index 1efe647be9..c2db103b73 100644 --- a/src/services/code-index/orchestrator.ts +++ b/src/services/code-index/orchestrator.ts @@ -1,428 +1,217 @@ 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 { DirectoryScanner } from "./processors" import { CacheManager } from "./cache-manager" -import { TelemetryService } from "@roo-code/telemetry" -import { TelemetryEventName } from "@roo-code/types" +import { CodeIndexScanExecutor } from "./code-index-scan-executor" +import { CodeIndexRun, CodeIndexRunState, CodeIndexScanMode } from "./code-index-run" +import { StateHolder } from "../../utils/StateHolder" +import { CodeIndexWatcherSession } from "./code-index-watcher-session" +import { CodeIndexRecovery } from "./code-index-recovery" import { t } from "../../i18n" /** * Manages the code indexing workflow, coordinating between different services and managers. */ export class CodeIndexOrchestrator { - private _fileWatcherSubscriptions: vscode.Disposable[] = [] - private _isProcessing: boolean = false - private _abortController: AbortController | null = null + private readonly codeIndexWatcherSession: CodeIndexWatcherSession + private readonly codeIndexRecovery: CodeIndexRecovery + private _activeRun: CodeIndexRun | undefined + private _isClearing = false + private readonly codeIndexScanExecutor: CodeIndexScanExecutor constructor( private readonly configManager: CodeIndexConfigManager, private readonly stateManager: CodeIndexStateManager, - private readonly workspacePath: string, + workspacePath: string, private readonly cacheManager: CacheManager, private readonly vectorStore: IVectorStore, - private readonly scanner: DirectoryScanner, - private readonly fileWatcher: IFileWatcher, - ) {} + scanner: DirectoryScanner, + fileWatcher: IFileWatcher, + ) { + this.codeIndexScanExecutor = new CodeIndexScanExecutor(workspacePath, scanner, vectorStore, stateManager) + this.codeIndexWatcherSession = new CodeIndexWatcherSession(fileWatcher, stateManager) + this.codeIndexRecovery = new CodeIndexRecovery( + cacheManager, + vectorStore, + stateManager, + this.codeIndexWatcherSession, + ) + } /** - * Starts the file watcher if not already running. + * Gets the current state of the indexing system. */ - private async _startWatcher(): Promise { - if (!this.configManager.isFeatureConfigured) { - throw new Error("Cannot start watcher: Service not configured.") - } + public get state(): IndexingState { + return this.stateManager.state + } - this.stateManager.setSystemState("Indexing", "Initializing file watcher...") + /** Runs a scan and starts watching files, retaining ownership until cleanup finishes. */ + public async startIndexing(): Promise { + if (!this._canStartIndexing()) return - try { - await this.fileWatcher.initialize() + const run = new CodeIndexRun(new AbortController(), new StateHolder("running")) + this._activeRun = run - 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 - } - }), - ] + try { + this.stateManager.setSystemState("Indexing", "Initializing services...") + const scanMode = await this._prepareScan(run) + await this._runScan(run, scanMode) + await this._completeIndexing(run.signal) } catch (error) { - console.error("[CodeIndexOrchestrator] Failed to start file watcher:", error) - TelemetryService.instance.captureEvent(TelemetryEventName.CODE_INDEX_ERROR, { - error: error instanceof Error ? error.message : String(error), - stack: error instanceof Error ? error.stack : undefined, - location: "_startWatcher", - }) - throw error + await this.codeIndexRecovery.handle(error, run) + } finally { + this._activeRun = undefined + run.finish() } } /** - * Updates the status of a file in the state manager. + * Checks start preconditions and reports why a request was rejected. */ + private _canStartIndexing(): boolean { + if (this._isClearing) { + return false + } - /** - * Initiates the indexing process (initial scan and starts watcher). - */ - public async startIndexing(): Promise { - // Check if workspace is available first - if (!vscode.workspace.workspaceFolders || vscode.workspace.workspaceFolders.length === 0) { + // Rejected requests must not overwrite the active run's status. + if (this._activeRun) { + console.warn("[CodeIndexOrchestrator] Start rejected: An indexing run is still active.") + return false + } + + if (!vscode.workspace.workspaceFolders?.length) { this.stateManager.setSystemState("Error", t("embeddings:orchestrator.indexingRequiresWorkspace")) console.warn("[CodeIndexOrchestrator] Start rejected: No workspace folder open.") - return + return false } if (!this.configManager.isFeatureConfigured) { this.stateManager.setSystemState("Standby", "Missing configuration. Save your settings to start indexing.") console.warn("[CodeIndexOrchestrator] Start rejected: Missing configuration.") - return + return false } - if ( - this._isProcessing || - (this.stateManager.state !== "Standby" && - this.stateManager.state !== "Error" && - this.stateManager.state !== "Indexed") - ) { + if (!["Standby", "Error", "Indexed"].includes(this.stateManager.state)) { console.warn( `[CodeIndexOrchestrator] Start rejected: Already processing or in state ${this.stateManager.state}.`, ) - return + return false } - this._isProcessing = true - this._abortController = new AbortController() - const signal = this._abortController.signal - this.stateManager.setSystemState("Indexing", "Initializing services...") - - // Track whether we successfully connected to Qdrant and started indexing - // This helps us decide whether to preserve cache on error - let indexingStarted = false - - try { - const collectionCreated = await this.vectorStore.initialize() - - // Successfully connected to Qdrant - indexingStarted = true - - if (collectionCreated) { - await this.cacheManager.clearCacheFile() - } - - // Check if the collection already has indexed data - // If it does, we can skip the full scan and just start the watcher - const hasExistingData = await this.vectorStore.hasIndexedData() - - if (hasExistingData && !collectionCreated) { - // Collection exists with data - run incremental scan to catch any new/changed files - // This handles files added while workspace was closed or Qdrant was inactive - console.log( - "[CodeIndexOrchestrator] Collection already has indexed data. Running incremental scan for new/changed files...", - ) - this.stateManager.setSystemState("Indexing", "Checking for new or modified files...") - - // Mark as incomplete at the start of incremental scan - await this.vectorStore.markIndexingIncomplete() - - let cumulativeBlocksIndexed = 0 - let cumulativeBlocksFoundSoFar = 0 - const batchErrors: Error[] = [] - - const handleFileParsed = (fileBlockCount: number) => { - cumulativeBlocksFoundSoFar += fileBlockCount - this.stateManager.reportBlockIndexingProgress(cumulativeBlocksIndexed, cumulativeBlocksFoundSoFar) - } - - const handleBlocksIndexed = (indexedCount: number) => { - cumulativeBlocksIndexed += indexedCount - this.stateManager.reportBlockIndexingProgress(cumulativeBlocksIndexed, cumulativeBlocksFoundSoFar) - } - - // Run incremental scan - scanner will skip unchanged files using cache - const result = await this.scanner.scanDirectory( - this.workspacePath, - (batchError: Error) => { - console.error( - `[CodeIndexOrchestrator] Error during incremental scan batch: ${batchError.message}`, - batchError, - ) - batchErrors.push(batchError) - }, - handleBlocksIndexed, - handleFileParsed, - signal, - ) - - if (signal.aborted) { - await this.cacheManager.flush() - this.stopWatcher() - this.stateManager.setSystemState("Standby", t("embeddings:orchestrator.indexingStopped")) - return - } - - if (!result) { - throw new Error("Incremental scan failed, is scanner initialized?") - } - - // If new files were found and indexed, log the results - if (cumulativeBlocksFoundSoFar > 0) { - console.log( - `[CodeIndexOrchestrator] Incremental scan completed: ${cumulativeBlocksIndexed} blocks indexed from new/changed files`, - ) - } else { - console.log("[CodeIndexOrchestrator] No new or changed files found") - } - - await this._startWatcher() - - // Mark indexing as complete after successful incremental scan - await this.vectorStore.markIndexingComplete() - - this.stateManager.setSystemState("Indexed", t("embeddings:orchestrator.fileWatcherStarted")) - } else { - // No existing data or collection was just created - do a full scan - this.stateManager.setSystemState("Indexing", "Services ready. Starting workspace scan...") - - // Mark as incomplete at the start of full scan - await this.vectorStore.markIndexingIncomplete() - - let cumulativeBlocksIndexed = 0 - let cumulativeBlocksFoundSoFar = 0 - const batchErrors: Error[] = [] - - const handleFileParsed = (fileBlockCount: number) => { - cumulativeBlocksFoundSoFar += fileBlockCount - this.stateManager.reportBlockIndexingProgress(cumulativeBlocksIndexed, cumulativeBlocksFoundSoFar) - } - - const handleBlocksIndexed = (indexedCount: number) => { - cumulativeBlocksIndexed += indexedCount - this.stateManager.reportBlockIndexingProgress(cumulativeBlocksIndexed, cumulativeBlocksFoundSoFar) - } - - const result = await this.scanner.scanDirectory( - this.workspacePath, - (batchError: Error) => { - console.error( - `[CodeIndexOrchestrator] Error during initial scan batch: ${batchError.message}`, - batchError, - ) - batchErrors.push(batchError) - }, - handleBlocksIndexed, - handleFileParsed, - signal, - ) - - if (signal.aborted) { - await this.cacheManager.flush() - this.stopWatcher() - this.stateManager.setSystemState("Standby", t("embeddings:orchestrator.indexingStopped")) - return - } - - if (!result) { - throw new Error("Scan failed, is scanner initialized?") - } + return true + } - const { stats } = result + private async _prepareScan(run: CodeIndexRun): Promise { + const collectionCreated = await this.vectorStore.initialize() + // Read actual contents, not completion metadata, before granting cleanup ownership. + // Query even a recreated collection; a failed read must never authorize cleanup. + run.preexistingCodePoints = await this.vectorStore.hasCodePoints() + if (collectionCreated) { + await this.cacheManager.clearCacheFile() + } - // Check if any blocks were actually indexed successfully - // If no blocks were indexed but blocks were found, it means all batches failed - if (cumulativeBlocksIndexed === 0 && cumulativeBlocksFoundSoFar > 0) { - if (batchErrors.length > 0) { - // Use the first batch error as it's likely representative of the main issue - const firstError = batchErrors[0] - throw new Error(`Indexing failed: ${firstError.message}`) - } else { - throw new Error(t("embeddings:orchestrator.indexingFailedNoBlocks")) - } - } + // Existing data can be updated incrementally using the cache. + const hasExistingData = await this.vectorStore.hasIndexedData() + return hasExistingData && !collectionCreated ? "incremental" : "full" + } - // Check for partial failures - if a significant portion of blocks failed - const failureRate = (cumulativeBlocksFoundSoFar - cumulativeBlocksIndexed) / cumulativeBlocksFoundSoFar - if (batchErrors.length > 0 && failureRate > 0.1) { - // More than 10% of blocks failed to index - const firstError = batchErrors[0] - throw new Error( - `Indexing partially failed: Only ${cumulativeBlocksIndexed} of ${cumulativeBlocksFoundSoFar} blocks were indexed. ${firstError.message}`, - ) - } + private async _runScan(run: CodeIndexRun, mode: CodeIndexScanMode): Promise { + run.markScanStarted(mode) + if (mode === "incremental") { + await this.codeIndexScanExecutor.runIncrementalScan(run.signal) + } else { + await this.codeIndexScanExecutor.runFullScan(run.signal) + } + } - // CRITICAL: If there were ANY batch errors and NO blocks were successfully indexed, - // this is a complete failure regardless of the failure rate calculation - if (batchErrors.length > 0 && cumulativeBlocksIndexed === 0) { - const firstError = batchErrors[0] - throw new Error(`Indexing failed completely: ${firstError.message}`) - } + private async _completeIndexing(signal: AbortSignal): Promise { + await this._startWatcher(signal) + signal.throwIfAborted() + await this.vectorStore.markIndexingComplete() + if (signal.aborted) { + await this.vectorStore.markIndexingIncomplete() + } + signal.throwIfAborted() - // Final sanity check: If we found blocks but indexed none and somehow no errors were reported, - // this is still a failure - if (cumulativeBlocksFoundSoFar > 0 && cumulativeBlocksIndexed === 0) { - throw new Error(t("embeddings:orchestrator.indexingFailedCritical")) - } + this.stateManager.setSystemState("Indexed", t("embeddings:orchestrator.fileWatcherStarted")) + } - await this._startWatcher() + /** + * Starts the file watcher if not already running. + */ + private async _startWatcher(signal: AbortSignal): Promise { + signal.throwIfAborted() + if (!this.configManager.isFeatureConfigured) { + throw new Error("Cannot start watcher: Service not configured.") + } + if (this.codeIndexWatcherSession.isRunning) return - // Mark indexing as complete after successful full scan - await this.vectorStore.markIndexingComplete() + this.stateManager.setSystemState("Indexing", "Initializing file watcher...") + await this.codeIndexWatcherSession.start(signal) + } - this.stateManager.setSystemState("Indexed", t("embeddings:orchestrator.fileWatcherStarted")) - } - } catch (error: any) { - // Handle abort gracefully — not an error, just a user-initiated stop - if (error?.name === "AbortError" || signal.aborted) { - console.log("[CodeIndexOrchestrator] Indexing aborted by user.") - await this.cacheManager.flush() - this.stopWatcher() - this.stateManager.setSystemState("Standby", t("embeddings:orchestrator.indexingStopped")) - return - } + /** + * Clears all index data by stopping indexing, clearing the vector store, + * and resetting the cache file. + */ + public async clearIndexData(): Promise { + if (this._isClearing) { + return + } + this._isClearing = true - console.error("[CodeIndexOrchestrator] Error during indexing:", error) - TelemetryService.instance.captureEvent(TelemetryEventName.CODE_INDEX_ERROR, { - error: error instanceof Error ? error.message : String(error), - stack: error instanceof Error ? error.stack : undefined, - location: "startIndexing", - }) - if (indexingStarted) { - try { - await this.vectorStore.clearCollection() - } catch (cleanupError) { - console.error("[CodeIndexOrchestrator] Failed to clean up after error:", cleanupError) - TelemetryService.instance.captureEvent(TelemetryEventName.CODE_INDEX_ERROR, { - error: cleanupError instanceof Error ? cleanupError.message : String(cleanupError), - stack: cleanupError instanceof Error ? cleanupError.stack : undefined, - location: "startIndexing.cleanup", - }) - } - } + try { + await this._stopAndAwaitIndexing() + await this._deleteIndexData() + this.stateManager.setSystemState("Standby", "Index data cleared successfully.") + } catch (error) { + this.codeIndexRecovery.handleClearError(error) + } finally { + this._isClearing = false + } + } - // Only clear cache if indexing had started (Qdrant connection succeeded) - // If we never connected to Qdrant, preserve cache for incremental scan when it comes back - if (indexingStarted) { - // Indexing started but failed mid-way - clear cache to avoid cache-Qdrant mismatch - await this.cacheManager.clearCacheFile() - console.log( - "[CodeIndexOrchestrator] Indexing failed after starting. Clearing cache to avoid inconsistency.", - ) - } else { - // Never connected to Qdrant - preserve cache for future incremental scan - console.log( - "[CodeIndexOrchestrator] Failed to connect to Qdrant. Preserving cache for future incremental scan.", - ) - } + private async _stopAndAwaitIndexing(): Promise { + const run = this._activeRun + this._requestIndexingCancellation(run) + this.codeIndexWatcherSession.stop() + await run?.waitUntilFinished() + } - this.stateManager.setSystemState( - "Error", - t("embeddings:orchestrator.failedDuringInitialScan", { - errorMessage: error.message || t("embeddings:orchestrator.unknownError"), - }), - ) - this.stopWatcher() - } finally { - this._isProcessing = false - this._abortController = null + private async _deleteIndexData(): Promise { + if (this.configManager.isFeatureConfigured) { + await this.vectorStore.deleteCollection() + } else { + console.warn("[CodeIndexOrchestrator] Service not configured, skipping vector collection clear.") } + + await this.cacheManager.clearCacheFile() } /** * Stops any in-progress indexing by aborting the scan and stopping the file watcher. */ public stopIndexing(): void { - if (this._abortController) { - this.stateManager.setSystemState("Stopping", t("embeddings:orchestrator.indexingStoppedPartial")) - this._abortController.abort() - this._abortController = null - } + this._requestIndexingCancellation(this._activeRun) this.stopWatcher() } + private _requestIndexingCancellation(run: CodeIndexRun | undefined): void { + if (!run || run.signal.aborted) return + this.stateManager.setSystemState("Stopping", t("embeddings:orchestrator.indexingStoppedPartial")) + run.cancel() + } + /** * Stops the file watcher and cleans up resources. */ public stopWatcher(): void { - this.fileWatcher.dispose() - this._fileWatcherSubscriptions.forEach((sub) => sub.dispose()) - this._fileWatcherSubscriptions = [] + this.codeIndexWatcherSession.stop() - if (this.stateManager.state !== "Error" && this.stateManager.state !== "Stopping") { + if (!["Error", "Stopping"].includes(this.stateManager.state)) { this.stateManager.setSystemState("Standby", t("embeddings:orchestrator.fileWatcherStopped")) } - this._isProcessing = false - } - - /** - * Clears all index data by stopping the watcher, clearing the vector store, - * and resetting the cache file. - */ - public async clearIndexData(): Promise { - this._isProcessing = true - - try { - await this.stopWatcher() - - try { - if (this.configManager.isFeatureConfigured) { - await this.vectorStore.deleteCollection() - } else { - console.warn("[CodeIndexOrchestrator] Service not configured, skipping vector collection clear.") - } - } catch (error: any) { - console.error("[CodeIndexOrchestrator] Failed to clear vector collection:", error) - TelemetryService.instance.captureEvent(TelemetryEventName.CODE_INDEX_ERROR, { - error: error instanceof Error ? error.message : String(error), - stack: error instanceof Error ? error.stack : undefined, - location: "clearIndexData", - }) - this.stateManager.setSystemState("Error", `Failed to clear vector collection: ${error.message}`) - } - - await this.cacheManager.clearCacheFile() - - if (this.stateManager.state !== "Error") { - this.stateManager.setSystemState("Standby", "Index data cleared successfully.") - } - } finally { - this._isProcessing = false - } - } - - /** - * Gets the current state of the indexing system. - */ - public get state(): IndexingState { - return this.stateManager.state } } diff --git a/src/services/code-index/processors/__tests__/file-watcher.spec.ts b/src/services/code-index/processors/__tests__/file-watcher.spec.ts index fc61e687bd..fe21a893fe 100644 --- a/src/services/code-index/processors/__tests__/file-watcher.spec.ts +++ b/src/services/code-index/processors/__tests__/file-watcher.spec.ts @@ -39,21 +39,24 @@ vi.mock("../parser", () => ({ })) const createMockEventEmitter = () => { + let disposed = false const listeners = new Set<(event: any) => void>() return { event: vi.fn((listener: (event: any) => void) => { - listeners.add(listener) + if (!disposed) listeners.add(listener) return { dispose: () => listeners.delete(listener), } }), fire: vi.fn((event: any) => { + if (disposed) return for (const listener of listeners) { listener(event) } }), dispose: vi.fn(() => { + disposed = true listeners.clear() }), } @@ -88,6 +91,101 @@ vi.mock("vscode", () => ({ })) describe("FileWatcher", () => { + it("does not deliver an old batch into subscriptions created after restart", async () => { + await fileWatcher.initialize() + let release!: () => void + let notifyStarted!: () => void + const started = new Promise((resolve) => { + notifyStarted = resolve + }) + const blocked = new Promise((resolve) => { + release = resolve + }) + vi.spyOn(fileWatcher, "processFile").mockImplementationOnce(async (path) => { + notifyStarted() + await blocked + return { path, status: "skipped", reason: "test" } + }) + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/old.ts")) + const processing = flushBatch() + await started + fileWatcher.dispose() + await fileWatcher.initialize() + const progress = vi.fn() + const finished = vi.fn() + fileWatcher.onBatchProgressUpdate(progress) + fileWatcher.onDidFinishBatchProcessing(finished) + release() + await processing + expect(progress).not.toHaveBeenCalled() + expect(finished).not.toHaveBeenCalled() + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/new.ts")) + await flushBatch() + expect(progress).toHaveBeenCalled() + expect(finished).toHaveBeenCalledOnce() + }) + + it("restores all batch events after disposal and reinitialization", async () => { + vi.mocked(vscode.workspace.createFileSystemWatcher).mockClear() + await fileWatcher.initialize() + const oldListener = vi.fn() + fileWatcher.onDidFinishBatchProcessing(oldListener) + fileWatcher.dispose() + await fileWatcher.initialize() + const started = vi.fn() + const progress = vi.fn() + const finished = vi.fn() + fileWatcher.onDidStartBatchProcessing(started) + fileWatcher.onBatchProgressUpdate(progress) + fileWatcher.onDidFinishBatchProcessing(finished) + fileWatcher["_onDidStartBatchProcessing"].fire(["file.ts"]) + fileWatcher["_onBatchProgressUpdate"].fire({ processedInBatch: 1, totalInBatch: 1 }) + fileWatcher["_onDidFinishBatchProcessing"].fire({ processedFiles: [] }) + expect(started).toHaveBeenCalledOnce() + expect(progress).toHaveBeenCalledOnce() + expect(finished).toHaveBeenCalledOnce() + expect(oldListener).not.toHaveBeenCalled() + expect(vscode.workspace.createFileSystemWatcher).toHaveBeenCalledTimes(2) + }) + + it("does not create duplicate native watchers on repeated initialization", async () => { + vi.mocked(vscode.workspace.createFileSystemWatcher).mockClear() + await fileWatcher.initialize() + await fileWatcher.initialize() + expect(vscode.workspace.createFileSystemWatcher).toHaveBeenCalledOnce() + }) + it("reports progress and preserves cached hashes when batch deletion fails", async () => { + await fileWatcher.initialize() + const error = new Error("deletion failed") + mockVectorStore.deletePointsByMultipleFilePaths.mockRejectedValueOnce(error) + const progress = vi.fn() + const finished = vi.fn() + fileWatcher.onBatchProgressUpdate(progress) + fileWatcher.onDidFinishBatchProcessing(finished) + const paths = ["/mock/workspace/first.ts", "/mock/workspace/second.ts"] + + for (const path of paths) { + await mockOnDidDelete(vscode.Uri.file(path)) + } + await flushBatch() + + expect(mockVectorStore.deletePointsByMultipleFilePaths).toHaveBeenCalledExactlyOnceWith(paths) + expect(progress.mock.calls).toEqual([ + [{ processedInBatch: 0, totalInBatch: 2, currentFile: undefined }], + [{ processedInBatch: 1, totalInBatch: 2, currentFile: paths[0] }], + [{ processedInBatch: 2, totalInBatch: 2, currentFile: paths[1] }], + [{ processedInBatch: 2, totalInBatch: 2 }], + [{ processedInBatch: 0, totalInBatch: 0, currentFile: undefined }], + ]) + expect(finished).toHaveBeenCalledExactlyOnceWith({ + processedFiles: paths.map((path) => ({ path, status: "error", error })), + batchError: error, + }) + expect(mockCacheManager.deleteHash).not.toHaveBeenCalled() + expect(mockCacheManager.updateHash).not.toHaveBeenCalled() + expect(mockVectorStore.upsertPoints).not.toHaveBeenCalled() + }) + let fileWatcher: FileWatcher let mockWatcher: any let mockOnDidCreate: any diff --git a/src/services/code-index/processors/file-watcher.ts b/src/services/code-index/processors/file-watcher.ts index a6a3122c36..826464e860 100644 --- a/src/services/code-index/processors/file-watcher.ts +++ b/src/services/code-index/processors/file-watcher.ts @@ -41,28 +41,35 @@ export class FileWatcher implements IFileWatcher { private readonly FILE_PROCESSING_CONCURRENCY_LIMIT = 10 private readonly batchSegmentThreshold: number - private readonly _onDidStartBatchProcessing = new vscode.EventEmitter() - private readonly _onBatchProgressUpdate = new vscode.EventEmitter<{ + private eventsDisposed = false + private _onDidStartBatchProcessing = new vscode.EventEmitter() + private _onBatchProgressUpdate = new vscode.EventEmitter<{ processedInBatch: number totalInBatch: number currentFile?: string }>() - private readonly _onDidFinishBatchProcessing = new vscode.EventEmitter() + private _onDidFinishBatchProcessing = new vscode.EventEmitter() /** * Event emitted when a batch of files begins processing */ - public readonly onDidStartBatchProcessing = this._onDidStartBatchProcessing.event + public get onDidStartBatchProcessing() { + return this._onDidStartBatchProcessing.event + } /** * Event emitted to report progress during batch processing */ - public readonly onBatchProgressUpdate = this._onBatchProgressUpdate.event + public get onBatchProgressUpdate() { + return this._onBatchProgressUpdate.event + } /** * Event emitted when a batch of files has finished processing */ - public readonly onDidFinishBatchProcessing = this._onDidFinishBatchProcessing.event + public get onDidFinishBatchProcessing() { + return this._onDidFinishBatchProcessing.event + } /** * Creates a new file watcher @@ -106,6 +113,17 @@ export class FileWatcher implements IFileWatcher { * Initializes the file watcher */ async initialize(): Promise { + if (this.fileWatcher) return + if (this.eventsDisposed) { + this._onDidStartBatchProcessing = new vscode.EventEmitter() + this._onBatchProgressUpdate = new vscode.EventEmitter<{ + processedInBatch: number + totalInBatch: number + currentFile?: string + }>() + this._onDidFinishBatchProcessing = new vscode.EventEmitter() + this.eventsDisposed = false + } // Create file watcher const filePattern = new vscode.RelativePattern( this.workspacePath, @@ -124,12 +142,15 @@ export class FileWatcher implements IFileWatcher { */ dispose(): void { this.fileWatcher?.dispose() + this.fileWatcher = undefined if (this.batchProcessDebounceTimer) { clearTimeout(this.batchProcessDebounceTimer) + this.batchProcessDebounceTimer = undefined } this._onDidStartBatchProcessing.dispose() this._onBatchProgressUpdate.dispose() this._onDidFinishBatchProcessing.dispose() + this.eventsDisposed = true this.accumulatedEvents.clear() } @@ -181,10 +202,16 @@ export class FileWatcher implements IFileWatcher { const eventsToProcess = new Map(this.accumulatedEvents) this.accumulatedEvents.clear() + // Capture this session's emitters before notifying listeners or awaiting work. + // Disposed emitters drop late events instead of forwarding them to a restarted session. + const batchEvents = { + progress: this._onBatchProgressUpdate, + finished: this._onDidFinishBatchProcessing, + } const filePathsInBatch = Array.from(eventsToProcess.keys()) this._onDidStartBatchProcessing.fire(filePathsInBatch) - await this.processBatch(eventsToProcess) + await this.processBatch(eventsToProcess, batchEvents) } /** @@ -197,6 +224,7 @@ export class FileWatcher implements IFileWatcher { totalFilesInBatch: number, pathsToExplicitlyDelete: string[], filesToUpsertDetails: Array<{ path: string; uri: vscode.Uri; originalType: "create" | "change" }>, + progress: FileWatcher["_onBatchProgressUpdate"], ): Promise<{ overallBatchError?: Error; clearedPaths: Set; processedCount: number }> { let overallBatchError: Error | undefined const allPathsToClearFromDB = new Set(pathsToExplicitlyDelete) @@ -215,7 +243,7 @@ export class FileWatcher implements IFileWatcher { this.cacheManager.deleteHash(path) batchResults.push({ path, status: "success" }) processedCountInBatch++ - this._onBatchProgressUpdate.fire({ + progress.fire({ processedInBatch: processedCountInBatch, totalInBatch: totalFilesInBatch, currentFile: path, @@ -238,7 +266,7 @@ export class FileWatcher implements IFileWatcher { for (const path of pathsToExplicitlyDelete) { batchResults.push({ path, status: "error", error: error as Error }) processedCountInBatch++ - this._onBatchProgressUpdate.fire({ + progress.fire({ processedInBatch: processedCountInBatch, totalInBatch: totalFilesInBatch, currentFile: path, @@ -256,6 +284,7 @@ export class FileWatcher implements IFileWatcher { processedCountInBatch: number, totalFilesInBatch: number, pathsToExplicitlyDelete: string[], + progress: FileWatcher["_onBatchProgressUpdate"], ): Promise<{ pointsForBatchUpsert: PointStruct[] successfullyProcessedForUpsert: Array<{ path: string; newHash?: string }> @@ -269,7 +298,7 @@ export class FileWatcher implements IFileWatcher { const chunkToProcess = filesToProcessConcurrently.slice(i, i + this.FILE_PROCESSING_CONCURRENCY_LIMIT) const chunkProcessingPromises = chunkToProcess.map(async (fileDetail) => { - this._onBatchProgressUpdate.fire({ + progress.fire({ processedInBatch: processedCountInBatch, totalInBatch: totalFilesInBatch, currentFile: fileDetail.path, @@ -335,7 +364,7 @@ export class FileWatcher implements IFileWatcher { if (!pathsToExplicitlyDelete.includes(resultPath || "")) { processedCountInBatch++ } - this._onBatchProgressUpdate.fire({ + progress.fire({ processedInBatch: processedCountInBatch, totalInBatch: totalFilesInBatch, currentFile: resultPath, @@ -420,6 +449,10 @@ export class FileWatcher implements IFileWatcher { private async processBatch( eventsToProcess: Map, + batchEvents: { + progress: FileWatcher["_onBatchProgressUpdate"] + finished: FileWatcher["_onDidFinishBatchProcessing"] + }, ): Promise { const batchResults: FileProcessingResult[] = [] let processedCountInBatch = 0 @@ -427,7 +460,7 @@ export class FileWatcher implements IFileWatcher { let overallBatchError: Error | undefined // Initial progress update - this._onBatchProgressUpdate.fire({ + batchEvents.progress.fire({ processedInBatch: 0, totalInBatch: totalFilesInBatch, currentFile: undefined, @@ -456,6 +489,7 @@ export class FileWatcher implements IFileWatcher { totalFilesInBatch, pathsToExplicitlyDelete, filesToUpsertDetails, + batchEvents.progress, ) overallBatchError = deletionError processedCountInBatch = deletionCount @@ -471,6 +505,7 @@ export class FileWatcher implements IFileWatcher { processedCountInBatch, totalFilesInBatch, pathsToExplicitlyDelete, + batchEvents.progress, ) processedCountInBatch = upsertCount @@ -483,17 +518,17 @@ export class FileWatcher implements IFileWatcher { ) // Finalize - this._onDidFinishBatchProcessing.fire({ + batchEvents.finished.fire({ processedFiles: batchResults, batchError: overallBatchError, }) - this._onBatchProgressUpdate.fire({ + batchEvents.progress.fire({ processedInBatch: totalFilesInBatch, totalInBatch: totalFilesInBatch, }) if (this.accumulatedEvents.size === 0) { - this._onBatchProgressUpdate.fire({ + batchEvents.progress.fire({ processedInBatch: 0, totalInBatch: 0, currentFile: undefined, diff --git a/src/services/code-index/vector-store/__tests__/qdrant-client.spec.ts b/src/services/code-index/vector-store/__tests__/qdrant-client.spec.ts index 80a1fca835..9be419ddb7 100644 --- a/src/services/code-index/vector-store/__tests__/qdrant-client.spec.ts +++ b/src/services/code-index/vector-store/__tests__/qdrant-client.spec.ts @@ -39,6 +39,7 @@ const mockQdrantClientInstance = { createPayloadIndex: vitest.fn(), upsert: vitest.fn(), query: vitest.fn(), + scroll: vi.fn(), delete: vitest.fn(), } @@ -75,6 +76,39 @@ describe("QdrantVectorStore", () => { vectorStore = new QdrantVectorStore(mockWorkspacePath, mockQdrantUrl, mockVectorSize, mockApiKey) }) + it.each([ + { name: "empty", complete: undefined, hasCode: false }, + { name: "incomplete metadata only", complete: false, hasCode: false }, + { name: "complete metadata only", complete: true, hasCode: false }, + { name: "incomplete index", complete: false, hasCode: true }, + { name: "complete index", complete: true, hasCode: true }, + { name: "legacy index", complete: undefined, hasCode: true }, + ])("detects actual code points in $name", async ({ complete, hasCode }) => { + const points = [ + ...(complete === undefined + ? [] + : [{ id: "metadata", payload: { type: "metadata", indexing_complete: complete } }]), + ...(hasCode ? [{ id: "code", payload: { filePath: "existing.ts" } }] : []), + ] + mockQdrantClientInstance.scroll.mockImplementation(async (_collection, request) => { + expect(request).toEqual({ + filter: { must_not: [{ key: "type", match: { value: "metadata" } }] }, + limit: 1, + with_payload: false, + with_vector: false, + }) + return { points: points.filter((point) => !("type" in point.payload)).slice(0, 1) } + }) + await expect(vectorStore.hasCodePoints()).resolves.toBe(hasCode) + expect(mockQdrantClientInstance.scroll).toHaveBeenCalledWith(expectedCollectionName, expect.any(Object)) + }) + + it("propagates code-point query errors instead of reporting an empty collection", async () => { + const error = new Error("Qdrant unavailable") + mockQdrantClientInstance.scroll.mockRejectedValue(error) + await expect(vectorStore.hasCodePoints()).rejects.toBe(error) + }) + it("should correctly initialize QdrantClient and collectionName in constructor", () => { expect(QdrantClient).toHaveBeenCalledTimes(1) expect(QdrantClient).toHaveBeenCalledWith({ diff --git a/src/services/code-index/vector-store/qdrant-client.ts b/src/services/code-index/vector-store/qdrant-client.ts index 99e6b33deb..60632260d8 100644 --- a/src/services/code-index/vector-store/qdrant-client.ts +++ b/src/services/code-index/vector-store/qdrant-client.ts @@ -144,7 +144,7 @@ export class QdrantVectorStore implements IVectorStore { /** * Initializes the vector store - * @returns Promise resolving to boolean indicating if a new collection was created + * @returns Whether a new collection was created, including recreation after a dimension change */ async initialize(): Promise { let created = false @@ -581,9 +581,20 @@ export class QdrantVectorStore implements IVectorStore { } /** - * Checks if the collection exists and has indexed points - * @returns Promise resolving to boolean indicating if the collection exists and has points + * Checks actual data presence independently of the indexing completion marker. + * Errors propagate so recovery cannot mistake an unreadable collection for an empty one. */ + async hasCodePoints(): Promise { + const result = await this.client.scroll(this.collectionName, { + filter: { must_not: [{ key: "type", match: { value: "metadata" } }] }, + limit: 1, + with_payload: false, + with_vector: false, + }) + return result.points.length > 0 + } + + /** Checks index readiness, retaining legacy completion-marker semantics. */ async hasIndexedData(): Promise { try { const collectionInfo = await this.getCollectionInfo() diff --git a/src/utils/StateHolder.ts b/src/utils/StateHolder.ts new file mode 100644 index 0000000000..4c7719f4d8 --- /dev/null +++ b/src/utils/StateHolder.ts @@ -0,0 +1,64 @@ +export interface StateSubscription { + unsubscribe(): void +} + +/** A replaying, read-only view of state. Subscriptions are synchronous. */ +export interface StateStream { + readonly value: T + subscribe(listener: (value: T) => void): StateSubscription + waitFor(predicate: (value: T) => boolean): Promise +} + +/** Stores state without retaining promises. Use immutable values when updating it. */ +export class StateHolder implements StateStream { + private readonly listeners = new Set<(value: T) => void>() + + constructor(private current: T) {} + + get value(): T { + return this.current + } + + set(value: T): void { + if (Object.is(this.current, value)) return + this.current = value + for (const listener of [...this.listeners]) { + if (this.listeners.has(listener)) listener(value) + } + } + + subscribe(listener: (value: T) => void): StateSubscription { + // Each subscription owns its registration, even for the same callback. + const subscriber = (value: T) => listener(value) + this.listeners.add(subscriber) + try { + subscriber(this.current) + } catch (error) { + this.listeners.delete(subscriber) + throw error + } + return { + unsubscribe: () => { + this.listeners.delete(subscriber) + }, + } + } + + waitFor(predicate: (value: T) => boolean): Promise { + return new Promise((resolve, reject) => { + const listener = (value: T) => { + try { + if (!predicate(value)) return + this.listeners.delete(listener) + resolve(value) + } catch (error) { + this.listeners.delete(listener) + reject(error) + } + } + // Register before checking current state so synchronous replay cannot leak a listener. + this.listeners.add(listener) + listener(this.current) + }) + } +} diff --git a/src/utils/__tests__/StateHolder.spec.ts b/src/utils/__tests__/StateHolder.spec.ts new file mode 100644 index 0000000000..9cd2cf374d --- /dev/null +++ b/src/utils/__tests__/StateHolder.spec.ts @@ -0,0 +1,97 @@ +import { StateHolder } from "../StateHolder" + +describe("StateHolder", () => { + it("replays current state, publishes changes and suppresses identical values", () => { + const holder = new StateHolder(0) + const listener = vi.fn() + const subscription = holder.subscribe(listener) + holder.set(1) + holder.set(1) + expect(holder.value).toBe(1) + expect(listener.mock.calls).toEqual([[0], [1]]) + subscription.unsubscribe() + subscription.unsubscribe() + holder.set(2) + expect(listener).toHaveBeenCalledTimes(2) + }) + + it("registers the same callback independently", () => { + const holder = new StateHolder(0) + const listener = vi.fn() + const first = holder.subscribe(listener) + const second = holder.subscribe(listener) + first.unsubscribe() + listener.mockClear() + holder.set(1) + expect(listener).toHaveBeenCalledExactlyOnceWith(1) + second.unsubscribe() + }) + + it("skips a subscriber removed by an earlier listener during notification", () => { + const holder = new StateHolder(0) + const first = holder.subscribe((value) => { + if (value === 1) second.unsubscribe() + }) + const removedListener = vi.fn() + const second = holder.subscribe(removedListener) + const remainingListener = vi.fn() + const third = holder.subscribe(remainingListener) + removedListener.mockClear() + remainingListener.mockClear() + + holder.set(1) + holder.set(2) + + expect(removedListener).not.toHaveBeenCalled() + expect(remainingListener.mock.calls).toEqual([[1], [2]]) + expect(holder.value).toBe(2) + first.unsubscribe() + third.unsubscribe() + }) + + it("resolves immediately for matching current state without retaining a listener", async () => { + const holder = new StateHolder("finished") + await expect(holder.waitFor((value) => value === "finished")).resolves.toBe("finished") + expect(holder["listeners"].size).toBe(0) + }) + + it("waits for matching state and removes each completed waiter", async () => { + const holder = new StateHolder(0) + const notified = vi.fn() + const first = holder.waitFor((value) => value === 2).then(notified) + const second = holder.waitFor((value) => value === 3) + holder.set(1) + await Promise.resolve() + expect(notified).not.toHaveBeenCalled() + holder.set(2) + await first + expect(notified).toHaveBeenCalledExactlyOnceWith(2) + expect(holder["listeners"].size).toBe(1) + holder.set(3) + await expect(second).resolves.toBe(3) + expect(holder["listeners"].size).toBe(0) + }) + + it.each([0, 1])("rejects and unsubscribes when a predicate throws at state %s", async (failureState) => { + const holder = new StateHolder(0) + const error = new Error("predicate failed") + const waiting = holder.waitFor((value) => { + if (value === failureState) throw error + return false + }) + const assertion = expect(waiting).rejects.toBe(error) + holder.set(1) + await assertion + expect(holder["listeners"].size).toBe(0) + }) + + it("removes a subscription if initial replay throws", () => { + const holder = new StateHolder(0) + expect(() => + holder.subscribe(() => { + throw new Error("replay failed") + }), + ).toThrow("replay failed") + expect(holder["listeners"].size).toBe(0) + }) +})