Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion src/eslint-suppressions.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down
125 changes: 120 additions & 5 deletions src/services/code-index/processors/__tests__/file-watcher.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,22 +38,25 @@ vi.mock("../parser", () => ({
},
}))

const createMockEventEmitter = () => {
const listeners = new Set<(event: any) => void>()
const createMockEventEmitter = <T>() => {
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()
}),
}
Expand Down Expand Up @@ -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<void>((resolve) => {
notifyStarted = resolve
})
const blocked = new Promise<void>((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
Expand Down
65 changes: 50 additions & 15 deletions src/services/code-index/processors/file-watcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string[]>()
private readonly _onBatchProgressUpdate = new vscode.EventEmitter<{
private eventsDisposed = false
private _onDidStartBatchProcessing = new vscode.EventEmitter<string[]>()
private _onBatchProgressUpdate = new vscode.EventEmitter<{
processedInBatch: number
totalInBatch: number
currentFile?: string
}>()
private readonly _onDidFinishBatchProcessing = new vscode.EventEmitter<BatchProcessingSummary>()
private _onDidFinishBatchProcessing = new vscode.EventEmitter<BatchProcessingSummary>()

/**
* 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
Expand Down Expand Up @@ -106,6 +113,17 @@ export class FileWatcher implements IFileWatcher {
* Initializes the file watcher
*/
async initialize(): Promise<void> {
if (this.fileWatcher) return
if (this.eventsDisposed) {
this._onDidStartBatchProcessing = new vscode.EventEmitter<string[]>()
this._onBatchProgressUpdate = new vscode.EventEmitter<{
processedInBatch: number
totalInBatch: number
currentFile?: string
}>()
this._onDidFinishBatchProcessing = new vscode.EventEmitter<BatchProcessingSummary>()
this.eventsDisposed = false
}
// Create file watcher
const filePattern = new vscode.RelativePattern(
this.workspacePath,
Expand All @@ -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()
}

Expand Down Expand Up @@ -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)
}

/**
Expand All @@ -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<string>; processedCount: number }> {
let overallBatchError: Error | undefined
const allPathsToClearFromDB = new Set<string>(pathsToExplicitlyDelete)
Expand All @@ -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,
Expand All @@ -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,
Expand All @@ -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 }>
Expand All @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -420,14 +449,18 @@ export class FileWatcher implements IFileWatcher {

private async processBatch(
eventsToProcess: Map<string, { uri: vscode.Uri; type: "create" | "change" | "delete" }>,
batchEvents: {
progress: FileWatcher["_onBatchProgressUpdate"]
finished: FileWatcher["_onDidFinishBatchProcessing"]
},
): Promise<void> {
const batchResults: FileProcessingResult[] = []
let processedCountInBatch = 0
const totalFilesInBatch = eventsToProcess.size
let overallBatchError: Error | undefined

// Initial progress update
this._onBatchProgressUpdate.fire({
batchEvents.progress.fire({
processedInBatch: 0,
totalInBatch: totalFilesInBatch,
currentFile: undefined,
Expand Down Expand Up @@ -456,6 +489,7 @@ export class FileWatcher implements IFileWatcher {
totalFilesInBatch,
pathsToExplicitlyDelete,
filesToUpsertDetails,
batchEvents.progress,
)
overallBatchError = deletionError
processedCountInBatch = deletionCount
Expand All @@ -471,6 +505,7 @@ export class FileWatcher implements IFileWatcher {
processedCountInBatch,
totalFilesInBatch,
pathsToExplicitlyDelete,
batchEvents.progress,
)
processedCountInBatch = upsertCount

Expand All @@ -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,
Expand Down
Loading