-
Notifications
You must be signed in to change notification settings - Fork 300
refactor(code-index): isolate watcher sessions and preserve batch outcomes #1844
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
WebMad
wants to merge
6
commits into
Zoo-Code-Org:main
Choose a base branch
from
WebMad:refactor/code-index-watcher-session
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from all commits
Commits
Show all changes
6 commits
Select commit
Hold shift + click to select a range
0c1d18c
refactor(code-index): isolate watcher sessions and preserve batch out…
WebMad 93a37c1
fix(code-index): recreate watcher sessions on restart
WebMad d2a34b0
Revert "fix(code-index): recreate watcher sessions on restart"
WebMad d07446c
fix(code-index): recreate watchers through injected factory on restart
WebMad 5d02a51
test(code-index): cover empty watcher batch preserving error state
WebMad f33afe3
refactor(code-index): clarify watcher session naming and types
WebMad File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
223 changes: 223 additions & 0 deletions
223
src/services/code-index/__tests__/code-index-watcher-session.spec.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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") | ||
| }) | ||
| }) | ||
| }) |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.