From 16e5bcdc1311fb822e8a545cbdb9f7307bdd6123 Mon Sep 17 00:00:00 2001 From: gubin-dev Date: Mon, 28 Sep 2026 16:51:02 +0300 Subject: [PATCH 1/3] fix(code-index): restore watcher events after restart --- src/eslint-suppressions.json | 2 +- .../processors/__tests__/file-watcher.spec.ts | 125 +++++++++++++++++- .../code-index/processors/file-watcher.ts | 65 ++++++--- 3 files changed, 171 insertions(+), 21 deletions(-) diff --git a/src/eslint-suppressions.json b/src/eslint-suppressions.json index 24db0bf433..0fd2904fe1 100644 --- a/src/eslint-suppressions.json +++ b/src/eslint-suppressions.json @@ -1396,7 +1396,7 @@ }, "services/code-index/processors/__tests__/file-watcher.spec.ts": { "@typescript-eslint/no-explicit-any": { - "count": 25 + "count": 22 } }, "services/code-index/processors/__tests__/parser.spec.ts": { 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..8af87983e9 100644 --- a/src/services/code-index/processors/__tests__/file-watcher.spec.ts +++ b/src/services/code-index/processors/__tests__/file-watcher.spec.ts @@ -38,22 +38,25 @@ vi.mock("../parser", () => ({ }, })) -const createMockEventEmitter = () => { - const listeners = new Set<(event: any) => void>() +const createMockEventEmitter = () => { + let disposed = false + const listeners = new Set<(event: T) => void>() return { - event: vi.fn((listener: (event: any) => void) => { - listeners.add(listener) + event: vi.fn((listener: (event: T) => void) => { + if (!disposed) listeners.add(listener) return { dispose: () => listeners.delete(listener), } }), - fire: vi.fn((event: any) => { + fire: vi.fn((event: T) => { + if (disposed) return for (const listener of listeners) { listener(event) } }), dispose: vi.fn(() => { + disposed = true listeners.clear() }), } @@ -88,6 +91,118 @@ 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) + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/restarted.ts")) + await flushBatch() + expect(started).toHaveBeenCalledExactlyOnceWith(["/mock/workspace/restarted.ts"]) + expect(progress).toHaveBeenCalled() + 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("discards queued events and the pending timer when disposed before a batch starts", async () => { + await fileWatcher.initialize() + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/old.ts")) + fileWatcher.dispose() + fileWatcher.dispose() + expect(mockWatcher.dispose).toHaveBeenCalledOnce() + expect(vi.getTimerCount()).toBe(0) + await fileWatcher.initialize() + const started = vi.fn() + fileWatcher.onDidStartBatchProcessing(started) + await flushBatch() + expect(started).not.toHaveBeenCalled() + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/new.ts")) + await flushBatch() + expect(started).toHaveBeenCalledExactlyOnceWith(["/mock/workspace/new.ts"]) + }) + + 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, From 770f0eb624b37590e281e73b8e5abd42644064ca Mon Sep 17 00:00:00 2001 From: gubin-dev Date: Mon, 28 Sep 2026 17:16:07 +0300 Subject: [PATCH 2/3] fix(code-index): drain active watcher batches before restart --- .../code-index/interfaces/file-processor.ts | 7 + .../processors/__tests__/file-watcher.spec.ts | 137 +++++++++++++++++- .../code-index/processors/file-watcher.ts | 51 ++++++- 3 files changed, 186 insertions(+), 9 deletions(-) diff --git a/src/services/code-index/interfaces/file-processor.ts b/src/services/code-index/interfaces/file-processor.ts index 8ecdc518c8..ec9bb41c36 100644 --- a/src/services/code-index/interfaces/file-processor.ts +++ b/src/services/code-index/interfaces/file-processor.ts @@ -56,6 +56,13 @@ export interface IFileWatcher extends vscode.Disposable { */ initialize(): Promise + /** + * Waits for all started watcher batches to settle, including vector writes and hash updates. + * Call dispose() first to stop intake and discard queued events. Does not drain scans + * or the cache manager's independent debounced disk saves. + */ + waitForIdle(): Promise + /** * Event emitted when a batch of files begins processing. * The event payload is an array of file paths included in the batch. 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 8af87983e9..773a5fa866 100644 --- a/src/services/code-index/processors/__tests__/file-watcher.spec.ts +++ b/src/services/code-index/processors/__tests__/file-watcher.spec.ts @@ -3,6 +3,8 @@ import * as vscode from "vscode" import { FileWatcher } from "../file-watcher" +import type { PointStruct } from "../../interfaces" +import { createHash } from "crypto" import { clearAllMocks } from "../../../../test-utils/reset" @@ -91,6 +93,136 @@ vi.mock("vscode", () => ({ })) describe("FileWatcher", () => { + const deferred = () => { + let resolve!: () => void + const promise = new Promise((done) => { + resolve = done + }) + return { promise, resolve } + } + + it("drains delayed persistence before concurrent restarts can write newer points and hashes", async () => { + await fileWatcher.initialize() + vi.mocked(vscode.workspace.createFileSystemWatcher).mockClear() + const blocked = deferred() + const entered = deferred() + const persisted = new Map() + const hashes = new Map() + mockCacheManager.updateHash.mockImplementation((path: string, hash: string) => hashes.set(path, hash)) + mockVectorStore.upsertPoints.mockImplementation(async (points: PointStruct[]) => { + entered.resolve() + await blocked.promise + for (const point of points) persisted.set(point.id, point) + }) + const path = "/mock/workspace/same.ts" + await mockOnDidCreate(vscode.Uri.file(path)) + await flushBatch() + await entered.promise + fileWatcher.dispose() + const restarted = vi.fn() + const restarts = [fileWatcher.initialize().then(restarted), fileWatcher.initialize()] + await Promise.resolve() + expect(restarted).not.toHaveBeenCalled() + expect(vscode.workspace.createFileSystemWatcher).not.toHaveBeenCalled() + expect(hashes.size).toBe(0) + blocked.resolve() + await Promise.all(restarts) + expect(vscode.workspace.createFileSystemWatcher).toHaveBeenCalledOnce() + vi.spyOn(fileWatcher, "processFile").mockResolvedValueOnce({ + path, + status: "processed_for_batching", + newHash: "new-hash", + pointsToUpsert: [...persisted.values()].map((point) => ({ ...point, vector: [9, 9, 9] })), + }) + await mockOnDidChange(vscode.Uri.file(path)) + await flushBatch() + await fileWatcher.waitForIdle() + expect([...persisted.values()].map((point) => point.vector)).toEqual([[9, 9, 9]]) + expect(hashes.get(path)).toBe("new-hash") + expect(mockCacheManager.updateHash.mock.calls.map((call: string[]) => call[1])).toEqual([ + createHash("sha256").update("test content").digest("hex"), + "new-hash", + ]) + }) + + it("waits for every overlapping batch, not only the latest, and releases all idle callers", async () => { + await fileWatcher.initialize() + const first = deferred() + const second = deferred() + mockVectorStore.upsertPoints + .mockImplementationOnce(() => first.promise) + .mockImplementationOnce(() => second.promise) + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/first.ts")) + await flushBatch() + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/second.ts")) + await flushBatch() + expect(mockVectorStore.upsertPoints).toHaveBeenCalledTimes(2) + fileWatcher.dispose() + const idle = vi.fn() + const waiters = [fileWatcher.waitForIdle().then(idle), fileWatcher.waitForIdle().then(idle)] + second.resolve() + await flushBatch() + expect(idle).not.toHaveBeenCalled() + first.resolve() + await Promise.all(waiters) + expect(idle).toHaveBeenCalledTimes(2) + expect(fileWatcher["activeBatches"].size).toBe(0) + await fileWatcher.waitForIdle() + }) + + it("does not revive a pending restart after another disposal or accept stale native callbacks", async () => { + await fileWatcher.initialize() + const staleCreate = mockOnDidCreate + const blocked = deferred() + mockVectorStore.upsertPoints.mockImplementationOnce(() => blocked.promise) + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/old.ts")) + await flushBatch() + fileWatcher.dispose() + vi.mocked(vscode.workspace.createFileSystemWatcher).mockClear() + const restart = fileWatcher.initialize() + fileWatcher.dispose() + blocked.resolve() + await restart + expect(vscode.workspace.createFileSystemWatcher).not.toHaveBeenCalled() + await fileWatcher.initialize() + await staleCreate(vscode.Uri.file("/mock/workspace/stale.ts")) + await flushBatch() + expect(mockVectorStore.upsertPoints).toHaveBeenCalledOnce() + expect(vscode.workspace.createFileSystemWatcher).toHaveBeenCalledOnce() + }) + + it("registers work before a batch-start listener restarts the watcher", async () => { + await fileWatcher.initialize() + const blocked = deferred() + mockVectorStore.upsertPoints.mockImplementationOnce(() => blocked.promise) + let restart: Promise | undefined + fileWatcher.onDidStartBatchProcessing(() => { + fileWatcher.dispose() + restart = fileWatcher.initialize() + }) + vi.mocked(vscode.workspace.createFileSystemWatcher).mockClear() + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/old.ts")) + await flushBatch() + expect(restart).toBeDefined() + expect(vscode.workspace.createFileSystemWatcher).not.toHaveBeenCalled() + blocked.resolve() + await restart + expect(vscode.workspace.createFileSystemWatcher).toHaveBeenCalledOnce() + }) + + it("handles timer batch rejection and removes rejected work from the drain set", async () => { + await fileWatcher.initialize() + const error = new Error("unexpected batch failure") + const log = vi.spyOn(console, "error").mockImplementation(() => {}) + fileWatcher["processBatch"] = vi.fn().mockRejectedValueOnce(error) + await mockOnDidCreate(vscode.Uri.file("/mock/workspace/error.ts")) + await flushBatch() + await fileWatcher.waitForIdle() + expect(log).toHaveBeenCalledWith("[FileWatcher] Unhandled batch processing error:", error) + expect(fileWatcher["activeBatches"].size).toBe(0) + log.mockRestore() + }) + it("does not deliver an old batch into subscriptions created after restart", async () => { await fileWatcher.initialize() let release!: () => void @@ -110,12 +242,13 @@ describe("FileWatcher", () => { const processing = flushBatch() await started fileWatcher.dispose() - await fileWatcher.initialize() + const restarting = fileWatcher.initialize() + release() + await restarting 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() diff --git a/src/services/code-index/processors/file-watcher.ts b/src/services/code-index/processors/file-watcher.ts index 826464e860..107cddf76c 100644 --- a/src/services/code-index/processors/file-watcher.ts +++ b/src/services/code-index/processors/file-watcher.ts @@ -42,6 +42,8 @@ export class FileWatcher implements IFileWatcher { private readonly batchSegmentThreshold: number private eventsDisposed = false + private sessionGeneration = 0 + private readonly activeBatches = new Set>() private _onDidStartBatchProcessing = new vscode.EventEmitter() private _onBatchProgressUpdate = new vscode.EventEmitter<{ processedInBatch: number @@ -114,6 +116,10 @@ export class FileWatcher implements IFileWatcher { */ async initialize(): Promise { if (this.fileWatcher) return + const generation = this.sessionGeneration + await this.waitForIdle() + // A stop invalidates pending restarts; concurrent initializers share the same watcher. + if (generation !== this.sessionGeneration || this.fileWatcher) return if (this.eventsDisposed) { this._onDidStartBatchProcessing = new vscode.EventEmitter() this._onBatchProgressUpdate = new vscode.EventEmitter<{ @@ -132,15 +138,32 @@ export class FileWatcher implements IFileWatcher { this.fileWatcher = vscode.workspace.createFileSystemWatcher(filePattern) // Register event handlers - this.fileWatcher.onDidCreate(this.handleFileCreated.bind(this)) - this.fileWatcher.onDidChange(this.handleFileChanged.bind(this)) - this.fileWatcher.onDidDelete(this.handleFileDeleted.bind(this)) + this.fileWatcher.onDidCreate((uri) => { + if (generation === this.sessionGeneration) void this.handleFileCreated(uri) + }) + this.fileWatcher.onDidChange((uri) => { + if (generation === this.sessionGeneration) void this.handleFileChanged(uri) + }) + this.fileWatcher.onDidDelete((uri) => { + if (generation === this.sessionGeneration) void this.handleFileDeleted(uri) + }) + } + + /** Waits for all started batches, including persistence and hash updates, to settle. + * Dispose first to discard queued events and prevent more work from arriving. + * This does not drain independent cache-save timers or directory scans. + */ + async waitForIdle(): Promise { + while (this.activeBatches.size > 0) { + await Promise.allSettled([...this.activeBatches]) + } } /** * Disposes the file watcher */ dispose(): void { + this.sessionGeneration++ this.fileWatcher?.dispose() this.fileWatcher = undefined if (this.batchProcessDebounceTimer) { @@ -188,7 +211,12 @@ export class FileWatcher implements IFileWatcher { if (this.batchProcessDebounceTimer) { clearTimeout(this.batchProcessDebounceTimer) } - this.batchProcessDebounceTimer = setTimeout(() => this.triggerBatchProcessing(), this.BATCH_DEBOUNCE_DELAY_MS) + this.batchProcessDebounceTimer = setTimeout(() => { + this.batchProcessDebounceTimer = undefined + void this.triggerBatchProcessing().catch((error: unknown) => { + console.error("[FileWatcher] Unhandled batch processing error:", error) + }) + }, this.BATCH_DEBOUNCE_DELAY_MS) } /** @@ -205,13 +233,22 @@ export class FileWatcher implements IFileWatcher { // 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 = { + started: this._onDidStartBatchProcessing, progress: this._onBatchProgressUpdate, finished: this._onDidFinishBatchProcessing, } const filePathsInBatch = Array.from(eventsToProcess.keys()) - this._onDidStartBatchProcessing.fire(filePathsInBatch) - - await this.processBatch(eventsToProcess, batchEvents) + // Register before invoking listeners: a listener may synchronously stop/restart us. + const batch = Promise.resolve().then(async () => { + batchEvents.started.fire(filePathsInBatch) + await this.processBatch(eventsToProcess, batchEvents) + }) + this.activeBatches.add(batch) + try { + await batch + } finally { + this.activeBatches.delete(batch) + } } /** From bae7be83bc163d6fcea6899682185ca0d60ecccc Mon Sep 17 00:00:00 2001 From: gubin-dev Date: Mon, 28 Sep 2026 17:22:03 +0300 Subject: [PATCH 3/3] revert(code-index): keep watcher restart at PR 1821 scope --- .../code-index/interfaces/file-processor.ts | 7 - .../processors/__tests__/file-watcher.spec.ts | 137 +----------------- .../code-index/processors/file-watcher.ts | 51 +------ 3 files changed, 9 insertions(+), 186 deletions(-) diff --git a/src/services/code-index/interfaces/file-processor.ts b/src/services/code-index/interfaces/file-processor.ts index ec9bb41c36..8ecdc518c8 100644 --- a/src/services/code-index/interfaces/file-processor.ts +++ b/src/services/code-index/interfaces/file-processor.ts @@ -56,13 +56,6 @@ export interface IFileWatcher extends vscode.Disposable { */ initialize(): Promise - /** - * Waits for all started watcher batches to settle, including vector writes and hash updates. - * Call dispose() first to stop intake and discard queued events. Does not drain scans - * or the cache manager's independent debounced disk saves. - */ - waitForIdle(): Promise - /** * Event emitted when a batch of files begins processing. * The event payload is an array of file paths included in the batch. 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 773a5fa866..8af87983e9 100644 --- a/src/services/code-index/processors/__tests__/file-watcher.spec.ts +++ b/src/services/code-index/processors/__tests__/file-watcher.spec.ts @@ -3,8 +3,6 @@ import * as vscode from "vscode" import { FileWatcher } from "../file-watcher" -import type { PointStruct } from "../../interfaces" -import { createHash } from "crypto" import { clearAllMocks } from "../../../../test-utils/reset" @@ -93,136 +91,6 @@ vi.mock("vscode", () => ({ })) describe("FileWatcher", () => { - const deferred = () => { - let resolve!: () => void - const promise = new Promise((done) => { - resolve = done - }) - return { promise, resolve } - } - - it("drains delayed persistence before concurrent restarts can write newer points and hashes", async () => { - await fileWatcher.initialize() - vi.mocked(vscode.workspace.createFileSystemWatcher).mockClear() - const blocked = deferred() - const entered = deferred() - const persisted = new Map() - const hashes = new Map() - mockCacheManager.updateHash.mockImplementation((path: string, hash: string) => hashes.set(path, hash)) - mockVectorStore.upsertPoints.mockImplementation(async (points: PointStruct[]) => { - entered.resolve() - await blocked.promise - for (const point of points) persisted.set(point.id, point) - }) - const path = "/mock/workspace/same.ts" - await mockOnDidCreate(vscode.Uri.file(path)) - await flushBatch() - await entered.promise - fileWatcher.dispose() - const restarted = vi.fn() - const restarts = [fileWatcher.initialize().then(restarted), fileWatcher.initialize()] - await Promise.resolve() - expect(restarted).not.toHaveBeenCalled() - expect(vscode.workspace.createFileSystemWatcher).not.toHaveBeenCalled() - expect(hashes.size).toBe(0) - blocked.resolve() - await Promise.all(restarts) - expect(vscode.workspace.createFileSystemWatcher).toHaveBeenCalledOnce() - vi.spyOn(fileWatcher, "processFile").mockResolvedValueOnce({ - path, - status: "processed_for_batching", - newHash: "new-hash", - pointsToUpsert: [...persisted.values()].map((point) => ({ ...point, vector: [9, 9, 9] })), - }) - await mockOnDidChange(vscode.Uri.file(path)) - await flushBatch() - await fileWatcher.waitForIdle() - expect([...persisted.values()].map((point) => point.vector)).toEqual([[9, 9, 9]]) - expect(hashes.get(path)).toBe("new-hash") - expect(mockCacheManager.updateHash.mock.calls.map((call: string[]) => call[1])).toEqual([ - createHash("sha256").update("test content").digest("hex"), - "new-hash", - ]) - }) - - it("waits for every overlapping batch, not only the latest, and releases all idle callers", async () => { - await fileWatcher.initialize() - const first = deferred() - const second = deferred() - mockVectorStore.upsertPoints - .mockImplementationOnce(() => first.promise) - .mockImplementationOnce(() => second.promise) - await mockOnDidCreate(vscode.Uri.file("/mock/workspace/first.ts")) - await flushBatch() - await mockOnDidCreate(vscode.Uri.file("/mock/workspace/second.ts")) - await flushBatch() - expect(mockVectorStore.upsertPoints).toHaveBeenCalledTimes(2) - fileWatcher.dispose() - const idle = vi.fn() - const waiters = [fileWatcher.waitForIdle().then(idle), fileWatcher.waitForIdle().then(idle)] - second.resolve() - await flushBatch() - expect(idle).not.toHaveBeenCalled() - first.resolve() - await Promise.all(waiters) - expect(idle).toHaveBeenCalledTimes(2) - expect(fileWatcher["activeBatches"].size).toBe(0) - await fileWatcher.waitForIdle() - }) - - it("does not revive a pending restart after another disposal or accept stale native callbacks", async () => { - await fileWatcher.initialize() - const staleCreate = mockOnDidCreate - const blocked = deferred() - mockVectorStore.upsertPoints.mockImplementationOnce(() => blocked.promise) - await mockOnDidCreate(vscode.Uri.file("/mock/workspace/old.ts")) - await flushBatch() - fileWatcher.dispose() - vi.mocked(vscode.workspace.createFileSystemWatcher).mockClear() - const restart = fileWatcher.initialize() - fileWatcher.dispose() - blocked.resolve() - await restart - expect(vscode.workspace.createFileSystemWatcher).not.toHaveBeenCalled() - await fileWatcher.initialize() - await staleCreate(vscode.Uri.file("/mock/workspace/stale.ts")) - await flushBatch() - expect(mockVectorStore.upsertPoints).toHaveBeenCalledOnce() - expect(vscode.workspace.createFileSystemWatcher).toHaveBeenCalledOnce() - }) - - it("registers work before a batch-start listener restarts the watcher", async () => { - await fileWatcher.initialize() - const blocked = deferred() - mockVectorStore.upsertPoints.mockImplementationOnce(() => blocked.promise) - let restart: Promise | undefined - fileWatcher.onDidStartBatchProcessing(() => { - fileWatcher.dispose() - restart = fileWatcher.initialize() - }) - vi.mocked(vscode.workspace.createFileSystemWatcher).mockClear() - await mockOnDidCreate(vscode.Uri.file("/mock/workspace/old.ts")) - await flushBatch() - expect(restart).toBeDefined() - expect(vscode.workspace.createFileSystemWatcher).not.toHaveBeenCalled() - blocked.resolve() - await restart - expect(vscode.workspace.createFileSystemWatcher).toHaveBeenCalledOnce() - }) - - it("handles timer batch rejection and removes rejected work from the drain set", async () => { - await fileWatcher.initialize() - const error = new Error("unexpected batch failure") - const log = vi.spyOn(console, "error").mockImplementation(() => {}) - fileWatcher["processBatch"] = vi.fn().mockRejectedValueOnce(error) - await mockOnDidCreate(vscode.Uri.file("/mock/workspace/error.ts")) - await flushBatch() - await fileWatcher.waitForIdle() - expect(log).toHaveBeenCalledWith("[FileWatcher] Unhandled batch processing error:", error) - expect(fileWatcher["activeBatches"].size).toBe(0) - log.mockRestore() - }) - it("does not deliver an old batch into subscriptions created after restart", async () => { await fileWatcher.initialize() let release!: () => void @@ -242,13 +110,12 @@ describe("FileWatcher", () => { const processing = flushBatch() await started fileWatcher.dispose() - const restarting = fileWatcher.initialize() - release() - await restarting + 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() diff --git a/src/services/code-index/processors/file-watcher.ts b/src/services/code-index/processors/file-watcher.ts index 107cddf76c..826464e860 100644 --- a/src/services/code-index/processors/file-watcher.ts +++ b/src/services/code-index/processors/file-watcher.ts @@ -42,8 +42,6 @@ export class FileWatcher implements IFileWatcher { private readonly batchSegmentThreshold: number private eventsDisposed = false - private sessionGeneration = 0 - private readonly activeBatches = new Set>() private _onDidStartBatchProcessing = new vscode.EventEmitter() private _onBatchProgressUpdate = new vscode.EventEmitter<{ processedInBatch: number @@ -116,10 +114,6 @@ export class FileWatcher implements IFileWatcher { */ async initialize(): Promise { if (this.fileWatcher) return - const generation = this.sessionGeneration - await this.waitForIdle() - // A stop invalidates pending restarts; concurrent initializers share the same watcher. - if (generation !== this.sessionGeneration || this.fileWatcher) return if (this.eventsDisposed) { this._onDidStartBatchProcessing = new vscode.EventEmitter() this._onBatchProgressUpdate = new vscode.EventEmitter<{ @@ -138,32 +132,15 @@ export class FileWatcher implements IFileWatcher { this.fileWatcher = vscode.workspace.createFileSystemWatcher(filePattern) // Register event handlers - this.fileWatcher.onDidCreate((uri) => { - if (generation === this.sessionGeneration) void this.handleFileCreated(uri) - }) - this.fileWatcher.onDidChange((uri) => { - if (generation === this.sessionGeneration) void this.handleFileChanged(uri) - }) - this.fileWatcher.onDidDelete((uri) => { - if (generation === this.sessionGeneration) void this.handleFileDeleted(uri) - }) - } - - /** Waits for all started batches, including persistence and hash updates, to settle. - * Dispose first to discard queued events and prevent more work from arriving. - * This does not drain independent cache-save timers or directory scans. - */ - async waitForIdle(): Promise { - while (this.activeBatches.size > 0) { - await Promise.allSettled([...this.activeBatches]) - } + this.fileWatcher.onDidCreate(this.handleFileCreated.bind(this)) + this.fileWatcher.onDidChange(this.handleFileChanged.bind(this)) + this.fileWatcher.onDidDelete(this.handleFileDeleted.bind(this)) } /** * Disposes the file watcher */ dispose(): void { - this.sessionGeneration++ this.fileWatcher?.dispose() this.fileWatcher = undefined if (this.batchProcessDebounceTimer) { @@ -211,12 +188,7 @@ export class FileWatcher implements IFileWatcher { if (this.batchProcessDebounceTimer) { clearTimeout(this.batchProcessDebounceTimer) } - this.batchProcessDebounceTimer = setTimeout(() => { - this.batchProcessDebounceTimer = undefined - void this.triggerBatchProcessing().catch((error: unknown) => { - console.error("[FileWatcher] Unhandled batch processing error:", error) - }) - }, this.BATCH_DEBOUNCE_DELAY_MS) + this.batchProcessDebounceTimer = setTimeout(() => this.triggerBatchProcessing(), this.BATCH_DEBOUNCE_DELAY_MS) } /** @@ -233,22 +205,13 @@ export class FileWatcher implements IFileWatcher { // 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 = { - started: this._onDidStartBatchProcessing, progress: this._onBatchProgressUpdate, finished: this._onDidFinishBatchProcessing, } const filePathsInBatch = Array.from(eventsToProcess.keys()) - // Register before invoking listeners: a listener may synchronously stop/restart us. - const batch = Promise.resolve().then(async () => { - batchEvents.started.fire(filePathsInBatch) - await this.processBatch(eventsToProcess, batchEvents) - }) - this.activeBatches.add(batch) - try { - await batch - } finally { - this.activeBatches.delete(batch) - } + this._onDidStartBatchProcessing.fire(filePathsInBatch) + + await this.processBatch(eventsToProcess, batchEvents) } /**