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
223 changes: 223 additions & 0 deletions src/services/code-index/__tests__/code-index-watcher-session.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,223 @@
import type { Event } from "vscode"
import type { BatchProcessingSummary, IFileWatcher } from "../interfaces"
import { CodeIndexStateManager } from "../state-manager"
import { CodeIndexWatcherSession } from "../code-index-watcher-session"

vi.mock("vscode", async () => {
const { makeEventEmitter } = await import("../../../test-utils/vscode")
return {
EventEmitter: vi.fn().mockImplementation(function () {
return makeEventEmitter()
}),
}
})

function eventSource<T>() {
const listeners = new Set<(value: T) => void>()
const dispose = vi.fn(() => listeners.clear())
const event: Event<T> = (listener) => {
listeners.add(listener)
return { dispose }
}
return { event: vi.fn(event), fire: (value: T) => listeners.forEach((listener) => listener(value)), dispose }
}

function setup() {
const start = eventSource<string[]>()
const progress = eventSource<{ processedInBatch: number; totalInBatch: number; currentFile?: string }>()
const finish = eventSource<BatchProcessingSummary>()
const watcher = {
initialize: vi.fn<() => Promise<void>>().mockResolvedValue(undefined),
dispose: vi.fn(),
processFile: vi.fn<IFileWatcher["processFile"]>(),
onDidStartBatchProcessing: start.event,
onBatchProgressUpdate: progress.event,
onDidFinishBatchProcessing: finish.event,
} satisfies IFileWatcher
const state = new CodeIndexStateManager()
const factory = { create: vi.fn(() => watcher) }
const session = new CodeIndexWatcherSession(factory, state)
return { start, progress, finish, watcher, state, session, factory }
}

describe("CodeIndexWatcherSession", () => {
it("allows startup after stopping an idle owner without allocating resources", async () => {
const { session, watcher, factory } = setup()
session.stop()
expect(factory.create).not.toHaveBeenCalled()
expect(watcher.dispose).not.toHaveBeenCalled()
await session.start()
expect(watcher.initialize).toHaveBeenCalledTimes(1)
})

it("reuses pending and active sessions without duplicating subscriptions", async () => {
const { session, watcher, start } = setup()
const pending = session.start()
expect(session.start()).toBe(pending)
await pending
await session.start()
expect(watcher.initialize).toHaveBeenCalledTimes(1)
expect(start.event).toHaveBeenCalledTimes(1)
})

it.each(["initialize", "progress", "finish"])("cleans partial startup when %s fails", async (stage) => {
const { session, watcher, start, progress, finish } = setup()
const error = new Error("startup failed")
if (stage === "initialize") watcher.initialize.mockRejectedValue(error)
else
(stage === "progress" ? progress : finish).event.mockImplementation(() => {
throw error
})
await expect(session.start()).rejects.toBe(error)
expect(watcher.dispose).toHaveBeenCalledTimes(1)
if (stage !== "initialize") expect(start.dispose).toHaveBeenCalledTimes(1)
if (stage === "finish") expect(progress.dispose).toHaveBeenCalledTimes(1)
session.stop()
expect(watcher.dispose).toHaveBeenCalledTimes(1)
})

it("does not revive initialization stopped while pending", async () => {
const { session, watcher, start } = setup()
let resolve!: () => void
watcher.initialize.mockReturnValue(
new Promise<void>((done) => {
resolve = done
}),
)
const pending = session.start()
session.stop()
resolve()
await expect(pending).rejects.toMatchObject({ name: "AbortError" })
expect(start.event).not.toHaveBeenCalled()
expect(watcher.dispose).toHaveBeenCalledTimes(2)
expect(watcher.initialize).toHaveBeenCalledTimes(1)
})

it("unsubscribes once and ignores retained callbacks after stop", async () => {
const { session, watcher, start, progress, finish, state } = setup()
await session.start()
const callback = finish.event.mock.calls[0][0]
session.stop()
session.stop()
callback({ processedFiles: [], batchError: new Error("late") })
expect(state.state).toBe("Standby")
for (const source of [start, progress, finish]) expect(source.dispose).toHaveBeenCalledTimes(1)
expect(watcher.dispose).toHaveBeenCalledTimes(1)
})

it.each(["stop", "failure"])("creates a working replacement after %s", async (reason) => {
const first = setup()
const next = setup()
first.factory.create.mockReturnValueOnce(first.watcher).mockReturnValue(next.watcher)
if (reason === "failure") {
first.watcher.initialize.mockRejectedValueOnce(new Error("startup failed"))
await expect(first.session.start()).rejects.toThrow("startup failed")
} else {
await first.session.start()
first.session.stop()
}
await first.session.start()
expect(first.factory.create).toHaveBeenCalledTimes(2)
expect(first.watcher.dispose).toHaveBeenCalledTimes(1)
next.progress.fire({ processedInBatch: 0, totalInBatch: 1 })
expect(first.state.state).toBe("Indexing")
next.finish.fire({ processedFiles: [] })
expect(first.state.state).toBe("Indexed")
})

it.each(["resolve", "reject"])("preserves replacement when old startup later %ss", async (outcome) => {
const first = setup()
const next = setup()
let resolve!: () => void
let reject!: (error: Error) => void
first.watcher.initialize.mockReturnValue(
new Promise<void>((done, fail) => {
resolve = done
reject = fail
}),
)
first.factory.create.mockReturnValueOnce(first.watcher).mockReturnValue(next.watcher)
const pending = first.session.start()
const rejected = expect(pending).rejects.toBeInstanceOf(Error)
first.session.stop()
await first.session.start()
if (outcome === "resolve") resolve()
else reject(new Error("late failure"))
await rejected
await first.session.start()
expect(first.factory.create).toHaveBeenCalledTimes(2)
expect(next.watcher.initialize).toHaveBeenCalledTimes(1)
expect(next.watcher.dispose).not.toHaveBeenCalled()
next.progress.fire({ processedInBatch: 0, totalInBatch: 1 })
expect(first.state.state).toBe("Indexing")
next.finish.fire({ processedFiles: [] })
expect(first.state.state).toBe("Indexed")
})

describe("real state manager integration", () => {
it("preserves the current outcome when an empty batch starts", async () => {
const { session, start, finish, state } = setup()
await session.start()
finish.fire({ processedFiles: [], batchError: new Error("database unavailable") })
const outcome = state.getCurrentStatus()
expect(outcome.systemStatus).toBe("Error")
start.fire([])
expect(state.getCurrentStatus()).toEqual(outcome)
})

it.each(["success", "skipped", "error", "local_error"] as const)(
"preserves the %s outcome through terminal and empty progress",
async (status) => {
const { session, start, progress, finish, state } = setup()
await session.start()
start.fire(["/workspace/file.ts"])
progress.fire({ processedInBatch: 0, totalInBatch: 1, currentFile: "/workspace/file.ts" })
expect(state.getCurrentStatus()).toMatchObject({ systemStatus: "Indexing", currentItemUnit: "files" })
expect(state.getCurrentStatus().message).toContain("Current: file.ts")
progress.fire({ processedInBatch: 1, totalInBatch: 1 })
expect(state.state).toBe("Indexing")
finish.fire({ processedFiles: [{ path: "/workspace/file.ts", status }] })
const outcome = state.getCurrentStatus()
expect(outcome.systemStatus).toBe(status === "error" || status === "local_error" ? "Error" : "Indexed")
progress.fire({ processedInBatch: 1, totalInBatch: 1 })
progress.fire({ processedInBatch: 0, totalInBatch: 0 })
expect(state.getCurrentStatus()).toEqual(outcome)
},
)

it("reports batch errors, then allows a subsequent successful batch to recover", async () => {
const { session, start, finish, state } = setup()
await session.start()
finish.fire({ processedFiles: [], batchError: new Error("database unavailable") })
expect(state.state).toBe("Error")
expect(state.getCurrentStatus().message).toContain("database unavailable")
start.fire(["next.ts"])
expect(state.state).toBe("Indexing")
finish.fire({ processedFiles: [{ path: "next.ts", status: "success" }] })
expect(state.state).toBe("Indexed")
})

it("reports file errors even in a mixed successful batch", async () => {
const { session, finish, state } = setup()
await session.start()
finish.fire({
processedFiles: [
{ path: "good.ts", status: "success" },
{ path: "bad.ts", status: "local_error", error: new Error("parse failed") },
],
})
expect(state.state).toBe("Error")
expect(state.getCurrentStatus().message).toContain("parse failed")
})

it("does not override Stopping with any watcher event", async () => {
const { session, start, progress, finish, state } = setup()
await session.start()
state.setSystemState("Stopping", "Stopped")
start.fire(["file.ts"])
progress.fire({ processedInBatch: 0, totalInBatch: 1 })
finish.fire({ processedFiles: [] })
expect(state.state).toBe("Stopping")
})
})
})
90 changes: 81 additions & 9 deletions src/services/code-index/__tests__/orchestrator.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => {
cacheManager,
vectorStore,
scanner,
fileWatcher,
{ create: () => fileWatcher },
)
await orchestrator.startIndexing()
expect(events).toEqual(["incomplete", "scan", "watcher", "complete"])
Expand All @@ -155,7 +155,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => {
cacheManager,
vectorStore,
scanner,
fileWatcher,
{ create: () => fileWatcher },
)
scanner.scanDirectory.mockImplementation(async () => {
orchestrator.stopIndexing()
Expand Down Expand Up @@ -184,7 +184,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => {
cacheManager,
vectorStore,
scanner,
fileWatcher,
{ create: () => fileWatcher },
)

// Act
Expand Down Expand Up @@ -213,7 +213,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => {
cacheManager,
vectorStore,
scanner,
fileWatcher,
{ create: () => fileWatcher },
)

// Act
Expand Down Expand Up @@ -249,7 +249,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => {
cacheManager,
vectorStore,
scanner,
fileWatcher,
{ create: () => fileWatcher },
)

await orchestrator.startIndexing()
Expand Down Expand Up @@ -279,7 +279,7 @@ describe("CodeIndexOrchestrator - error path cleanup gating", () => {
cacheManager,
vectorStore,
scanner,
fileWatcher,
{ create: () => fileWatcher },
)

await orchestrator.startIndexing()
Expand Down Expand Up @@ -346,6 +346,78 @@ describe("CodeIndexOrchestrator - stopIndexing", () => {
}
})

it("does not publish Indexed when stopped during watcher initialization", async () => {
let finishInitialization!: () => void
let enteredInitialization!: () => void
const entered = new Promise<void>((resolve) => {
enteredInitialization = resolve
})
fileWatcher.initialize.mockImplementation(() => {
enteredInitialization()
return new Promise<void>((resolve) => {
finishInitialization = resolve
})
})
scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 })
const orchestrator = new CodeIndexOrchestrator(
configManager,
stateManager,
workspacePath,
cacheManager,
vectorStore,
scanner,
{ create: () => fileWatcher },
)
const indexing = orchestrator.startIndexing()
await entered
orchestrator.stopIndexing()
finishInitialization()
await indexing
expect(stateManager.state).toBe("Standby")
expect(stateManager.setSystemState).not.toHaveBeenCalledWith("Indexed", expect.anything())
expect(fileWatcher.onDidStartBatchProcessing).not.toHaveBeenCalled()
expect(vectorStore.markIndexingComplete).not.toHaveBeenCalled()
})

Comment thread
WebMad marked this conversation as resolved.
it.each(["stop", "clear", "error"])("restarts the same orchestrator after %s", async (reason) => {
scanner.scanDirectory.mockResolvedValue({ stats: { processed: 0, skipped: 0 }, totalBlockCount: 0 })
vectorStore.deleteCollection = vi.fn().mockResolvedValue(undefined)
const nextWatcher = {
initialize: vi.fn().mockResolvedValue(undefined),
processFile: vi.fn(),
onDidStartBatchProcessing: vi.fn().mockReturnValue({ dispose: vi.fn() }),
onBatchProgressUpdate: vi.fn().mockReturnValue({ dispose: vi.fn() }),
onDidFinishBatchProcessing: vi.fn().mockReturnValue({ dispose: vi.fn() }),
dispose: vi.fn(),
}
const factory = { create: vi.fn().mockReturnValueOnce(fileWatcher).mockReturnValue(nextWatcher) }
const orchestrator = new CodeIndexOrchestrator(
configManager,
stateManager,
workspacePath,
cacheManager,
vectorStore,
scanner,
factory,
)
if (reason === "error") vectorStore.initialize.mockRejectedValueOnce(new Error("Qdrant unavailable"))
await orchestrator.startIndexing()
if (reason === "error") {
expect(stateManager.state).toBe("Error")
expect(factory.create).not.toHaveBeenCalled()
} else {
expect(stateManager.state).toBe("Indexed")
if (reason === "stop") orchestrator.stopIndexing()
else await orchestrator.clearIndexData()
expect(fileWatcher.dispose).toHaveBeenCalledTimes(1)
}
await orchestrator.startIndexing()
expect(stateManager.state).toBe("Indexed")
expect(factory.create).toHaveBeenCalledTimes(reason === "error" ? 1 : 2)
expect(fileWatcher.initialize).toHaveBeenCalledTimes(1)
expect(nextWatcher.initialize).toHaveBeenCalledTimes(reason === "error" ? 0 : 1)
})

it("should abort indexing when stopIndexing() is called", async () => {
// Make scanner hang until aborted
scanner.scanDirectory.mockImplementation(
Expand All @@ -369,7 +441,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => {
cacheManager,
vectorStore,
scanner,
fileWatcher,
{ create: () => fileWatcher },
)

// Start indexing (async, don't await)
Expand Down Expand Up @@ -412,7 +484,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => {
cacheManager,
vectorStore,
scanner,
fileWatcher,
{ create: () => fileWatcher },
)

const indexingPromise = orchestrator.startIndexing()
Expand Down Expand Up @@ -450,7 +522,7 @@ describe("CodeIndexOrchestrator - stopIndexing", () => {
cacheManager,
vectorStore,
scanner,
fileWatcher,
{ create: () => fileWatcher },
)

const indexingPromise = orchestrator.startIndexing()
Expand Down
Loading
Loading