diff --git a/src/extension.ts b/src/extension.ts index bd43bcbf8a..afb356ba67 100644 --- a/src/extension.ts +++ b/src/extension.ts @@ -215,6 +215,7 @@ export async function activate(context: vscode.ExtensionContext) { export async function deactivate() { outputChannel.appendLine(`${Package.name} extension deactivated`) await McpServerManager.cleanup(extensionContext) + CodeIndexManager.disposeAll() TelemetryService.instance.shutdown() TerminalRegistry.cleanup() } diff --git a/src/services/code-index/interfaces/manager.ts b/src/services/code-index/interfaces/manager.ts index 70e3fd9765..dfc00ca111 100644 --- a/src/services/code-index/interfaces/manager.ts +++ b/src/services/code-index/interfaces/manager.ts @@ -42,7 +42,7 @@ export interface ICodeIndexManager { /** * Stops the file watcher */ - stopWatcher(): void + stopWatcher(): Promise /** * Clears the index data @@ -66,7 +66,7 @@ export interface ICodeIndexManager { /** * Disposes of resources used by the manager */ - dispose(): void + dispose(): Promise } export type IndexingState = "Standby" | "Indexing" | "Indexed" | "Error" diff --git a/src/services/code-index/manager.ts b/src/services/code-index/manager.ts index 02c55f72e2..5ce803e069 100644 --- a/src/services/code-index/manager.ts +++ b/src/services/code-index/manager.ts @@ -1,14 +1,12 @@ import * as vscode from "vscode" +import { Worker } from "worker_threads" import { getWorkspacePath } from "../../utils/path" import { ContextProxy } from "../../core/config/ContextProxy" import { VectorStoreSearchResult } from "./interfaces" import { IndexingState } from "./interfaces/manager" import { CodeIndexConfigManager } from "./config-manager" import { CodeIndexStateManager } from "./state-manager" -import { CodeIndexServiceFactory } from "./service-factory" -import { CodeIndexSearchService } from "./search-service" -import { CodeIndexOrchestrator } from "./orchestrator" -import { CacheManager } from "./cache-manager" +import { WorkerCommand, WorkerResponse, WorkerMessage, WorkerInitConfig } from "./worker-messenger" import fs from "fs/promises" import ignore from "ignore" import path from "path" @@ -20,10 +18,10 @@ export class CodeIndexManager { // Specialized class instances private _configManager: CodeIndexConfigManager | undefined private readonly _stateManager: CodeIndexStateManager - private _serviceFactory: CodeIndexServiceFactory | undefined - private _orchestrator: CodeIndexOrchestrator | undefined - private _searchService: CodeIndexSearchService | undefined - private _cacheManager: CacheManager | undefined + private _worker: Worker | undefined + private _workerReady: boolean = false + private _messageId: number = 0 + private _pendingMessages: Map void; reject: (error: any) => void }> = new Map() public static getInstance(context: vscode.ExtensionContext): CodeIndexManager | undefined { const workspacePath = getWorkspacePath() // Assumes single workspace for now @@ -62,7 +60,7 @@ export class CodeIndexManager { } private assertInitialized() { - if (!this._configManager || !this._orchestrator || !this._searchService || !this._cacheManager) { + if (!this._configManager || !this._worker || !this._workerReady) { throw new Error("CodeIndexManager not initialized. Call initialize() first.") } } @@ -71,8 +69,7 @@ export class CodeIndexManager { if (!this.isFeatureEnabled) { return "Standby" } - this.assertInitialized() - return this._orchestrator!.state + return this._stateManager.state } public get isFeatureEnabled(): boolean { @@ -107,34 +104,25 @@ export class CodeIndexManager { // 2. Check if feature is enabled if (!this.isFeatureEnabled) { - if (this._orchestrator) { - this._orchestrator.stopWatcher() + if (this._worker) { + await this.stopWorker() } return { requiresRestart } } - // 3. CacheManager Initialization - if (!this._cacheManager) { - this._cacheManager = new CacheManager(this.context, this.workspacePath) - await this._cacheManager.initialize() + // 3. Determine if Worker Needs Recreation + const needsWorkerRecreation = !this._worker || requiresRestart + + if (needsWorkerRecreation) { + await this._recreateWorker() } - // 4. Determine if Core Services Need Recreation - const needsServiceRecreation = !this._serviceFactory || requiresRestart - - if (needsServiceRecreation) { - await this._recreateServices() - } - - // 5. Handle Indexing Start/Restart - // The enhanced vectorStore.initialize() in startIndexing() now handles dimension changes automatically - // by detecting incompatible collections and recreating them, so we rely on that for dimension changes + // 4. Handle Indexing Start/Restart const shouldStartOrRestartIndexing = - requiresRestart || - (needsServiceRecreation && (!this._orchestrator || this._orchestrator.state !== "Indexing")) + requiresRestart || (needsWorkerRecreation && this._stateManager.state !== "Indexing") if (shouldStartOrRestartIndexing) { - this._orchestrator?.startIndexing() // This method is async, but we don't await it here + await this.startIndexing() } return { requiresRestart } @@ -149,28 +137,26 @@ export class CodeIndexManager { return } this.assertInitialized() - await this._orchestrator!.startIndexing() + await this.sendWorkerCommand({ type: "start" }) } /** * Stops the file watcher and potentially cleans up resources. */ - public stopWatcher(): void { + public async stopWatcher(): Promise { if (!this.isFeatureEnabled) { return } - if (this._orchestrator) { - this._orchestrator.stopWatcher() + if (this._worker) { + await this.sendWorkerCommand({ type: "stop" }) } } /** * Cleans up the manager instance. */ - public dispose(): void { - if (this._orchestrator) { - this.stopWatcher() - } + public async dispose(): Promise { + await this.stopWorker() this._stateManager.dispose() } @@ -183,8 +169,7 @@ export class CodeIndexManager { return } this.assertInitialized() - await this._orchestrator!.clearIndexData() - await this._cacheManager!.clearCacheFile() + await this.sendWorkerCommand({ type: "clear" }) } // --- Private Helpers --- @@ -198,62 +183,141 @@ export class CodeIndexManager { return [] } this.assertInitialized() - return this._searchService!.searchIndex(query, directoryPrefix) + const response = await this.sendWorkerCommand({ type: "search", query, directoryPrefix }) + if (response.type === "searchResult") { + return response.results + } + throw new Error("Unexpected response from worker") } /** - * Private helper method to recreate services with current configuration. + * Private helper method to recreate worker with current configuration. * Used by both initialize() and handleSettingsChange(). */ - private async _recreateServices(): Promise { - // Stop watcher if it exists - if (this._orchestrator) { - this.stopWatcher() + private async _recreateWorker(): Promise { + // Stop existing worker if it exists + await this.stopWorker() + + // Create worker configuration + const config = this._configManager!.getConfig() + const workerConfig: WorkerInitConfig = { + workspacePath: this.workspacePath, + contextPath: this.context.globalStorageUri.fsPath, + isFeatureEnabled: this._configManager!.isFeatureEnabled, + isFeatureConfigured: this._configManager!.isFeatureConfigured, + qdrantUrl: config.qdrantUrl, + embedderProvider: config.embedderProvider, + embedderBaseUrl: config.openAiCompatibleOptions?.baseUrl, + embedderModelId: config.modelId, + embedderApiKey: this._getEmbedderApiKey(config), + searchMinScore: config.searchMinScore, } - // (Re)Initialize service factory - this._serviceFactory = new CodeIndexServiceFactory( - this._configManager!, - this.workspacePath, - this._cacheManager!, - ) + // Create new worker + const workerPath = path.join(__dirname, "../../workers/indexing-worker.js") + this._worker = new Worker(workerPath) + this._workerReady = false - const ignoreInstance = ignore() - const ignorePath = path.join(getWorkspacePath(), ".gitignore") - try { - const content = await fs.readFile(ignorePath, "utf8") - ignoreInstance.add(content) - ignoreInstance.add(".gitignore") - } catch (error) { - // Should never happen: reading file failed even though it exists - console.error("Unexpected error loading .gitignore:", error) + // Set up message handling + this._worker.on("message", this.handleWorkerMessage.bind(this)) + this._worker.on("error", this.handleWorkerError.bind(this)) + this._worker.on("exit", this.handleWorkerExit.bind(this)) + + // Initialize the worker + const response = await this.sendWorkerCommand({ type: "initialize", config: workerConfig }) + if (response.type === "initialized" && response.success) { + this._workerReady = true + } else { + throw new Error("Failed to initialize worker") + } + } + + private _getEmbedderApiKey(config: any): string | undefined { + if (config.embedderProvider === "openai") { + return config.openAiOptions?.openAiNativeApiKey + } else if (config.embedderProvider === "openai-compatible") { + return config.openAiCompatibleOptions?.apiKey + } + return undefined + } + + private async stopWorker(): Promise { + if (this._worker) { + // Clear pending messages + for (const [, pending] of this._pendingMessages) { + pending.reject(new Error("Worker stopped")) + } + this._pendingMessages.clear() + + // Terminate the worker + await this._worker.terminate() + this._worker = undefined + this._workerReady = false + } + } + + private async sendWorkerCommand(command: WorkerCommand): Promise { + if (!this._worker) { + throw new Error("Worker not initialized") } - // (Re)Create shared service instances - const { embedder, vectorStore, scanner, fileWatcher } = this._serviceFactory.createServices( - this.context, - this._cacheManager!, - ignoreInstance, - ) + const id = `msg_${this._messageId++}` + const message: WorkerMessage = { id, payload: command } - // (Re)Initialize orchestrator - this._orchestrator = new CodeIndexOrchestrator( - this._configManager!, - this._stateManager, - this.workspacePath, - this._cacheManager!, - vectorStore, - scanner, - fileWatcher, - ) + return new Promise((resolve, reject) => { + this._pendingMessages.set(id, { resolve, reject }) + this._worker!.postMessage(message) + }) + } - // (Re)Initialize search service - this._searchService = new CodeIndexSearchService( - this._configManager!, - this._stateManager, - embedder, - vectorStore, - ) + private handleWorkerMessage(message: WorkerMessage) { + const { id, payload } = message + + // Handle responses to specific commands + const pending = this._pendingMessages.get(id) + if (pending) { + this._pendingMessages.delete(id) + if (payload.type === "error") { + pending.reject(new Error(payload.error)) + } else { + pending.resolve(payload) + } + return + } + + // Handle unsolicited messages (progress updates, status changes) + switch (payload.type) { + case "progress": + // File processing progress - no need to update state manager + // as it's already updated in the worker + break + case "blockProgress": + // Block indexing progress - no need to update state manager + // as it's already updated in the worker + break + case "status": + // Update local state manager to match worker state + this._stateManager.setSystemState(payload.state, payload.message) + break + case "error": + console.error("[CodeIndexManager] Worker error:", payload.error) + this._stateManager.setSystemState("Error", payload.error) + break + } + } + + private handleWorkerError(error: Error) { + console.error("[CodeIndexManager] Worker error:", error) + this._stateManager.setSystemState("Error", `Worker error: ${error.message}`) + } + + private handleWorkerExit(code: number) { + if (code !== 0) { + console.error(`[CodeIndexManager] Worker exited with code ${code}`) + this._stateManager.setSystemState("Error", `Worker exited unexpectedly with code ${code}`) + } + this._worker = undefined + this._workerReady = false } /** @@ -271,11 +335,11 @@ export class CodeIndexManager { // If configuration changes require a restart and the manager is initialized, restart the service if (requiresRestart && isFeatureEnabled && isFeatureConfigured && this.isInitialized) { - // Recreate services with new configuration - await this._recreateServices() + // Recreate worker with new configuration + await this._recreateWorker() - // Start indexing with new services - this.startIndexing() + // Start indexing with new worker + await this.startIndexing() } } } diff --git a/src/services/code-index/worker-messenger.ts b/src/services/code-index/worker-messenger.ts new file mode 100644 index 0000000000..ab5ea6b3ae --- /dev/null +++ b/src/services/code-index/worker-messenger.ts @@ -0,0 +1,42 @@ +import { VectorStoreSearchResult } from "./interfaces" +import { IndexingState } from "./interfaces/manager" + +// Command messages sent from main thread to worker +export type WorkerCommand = + | { type: "start" } + | { type: "stop" } + | { type: "clear" } + | { type: "search"; query: string; directoryPrefix?: string } + | { type: "initialize"; config: WorkerInitConfig } + +// Response messages sent from worker to main thread +export type WorkerResponse = + | { type: "progress"; processedInBatch: number; totalInBatch: number; currentFile?: string } + | { type: "status"; state: IndexingState; message: string } + | { type: "searchResult"; results: VectorStoreSearchResult[] } + | { type: "error"; error: string } + | { type: "initialized"; success: boolean } + | { type: "stopped"; success: boolean } + | { type: "cleared"; success: boolean } + | { type: "blockProgress"; blocksIndexed: number; totalBlocks: number } + +// Configuration passed to worker on initialization +export interface WorkerInitConfig { + workspacePath: string + contextPath: string + isFeatureEnabled: boolean + isFeatureConfigured: boolean + // Add other necessary config from CodeIndexConfigManager + qdrantUrl?: string + embedderProvider?: string + embedderBaseUrl?: string + embedderModelId?: string + embedderApiKey?: string + searchMinScore?: number +} + +// Message wrapper for type safety +export interface WorkerMessage { + id: string + payload: T +} diff --git a/src/workers/indexing-worker.ts b/src/workers/indexing-worker.ts new file mode 100644 index 0000000000..fbf60f3a39 --- /dev/null +++ b/src/workers/indexing-worker.ts @@ -0,0 +1,219 @@ +import { parentPort } from "worker_threads" +import { WorkerCommand, WorkerResponse, WorkerMessage, WorkerInitConfig } from "../services/code-index/worker-messenger" +import { CodeIndexConfigManager } from "../services/code-index/config-manager" +import { CodeIndexStateManager, IndexingState } from "../services/code-index/state-manager" +import { CodeIndexServiceFactory } from "../services/code-index/service-factory" +import { CodeIndexOrchestrator } from "../services/code-index/orchestrator" +import { CodeIndexSearchService } from "../services/code-index/search-service" +import { CacheManager } from "../services/code-index/cache-manager" +import { DirectoryScanner } from "../services/code-index/processors" +import { IEmbedder, IVectorStore, IFileWatcher, VectorStoreSearchResult } from "../services/code-index/interfaces" +import ignore from "ignore" +import * as fs from "fs/promises" +import * as path from "path" + +class IndexingWorker { + private config: WorkerInitConfig | null = null + private orchestrator: CodeIndexOrchestrator | null = null + private searchService: CodeIndexSearchService | null = null + private stateManager: CodeIndexStateManager | null = null + private cacheManager: CacheManager | null = null + private embedder: IEmbedder | null = null + private vectorStore: IVectorStore | null = null + private scanner: DirectoryScanner | null = null + private fileWatcher: IFileWatcher | null = null + + constructor() { + if (!parentPort) { + throw new Error("This file must be run as a worker thread") + } + + parentPort.on("message", this.handleMessage.bind(this)) + } + + private async handleMessage(message: WorkerMessage) { + const { id, payload } = message + + try { + switch (payload.type) { + case "initialize": + await this.initialize(payload.config) + this.sendResponse(id, { type: "initialized", success: true }) + break + + case "start": + await this.startIndexing() + break + + case "stop": + await this.stopIndexing() + this.sendResponse(id, { type: "stopped", success: true }) + break + + case "clear": + await this.clearIndex() + break + + case "search": { + const results = await this.search(payload.query, payload.directoryPrefix) + this.sendResponse(id, { type: "searchResult", results }) + break + } + + default: + this.sendResponse(id, { + type: "error", + error: `Unknown command type: ${(payload as any).type}`, + }) + } + } catch (error) { + this.sendResponse(id, { + type: "error", + error: error instanceof Error ? error.message : String(error), + }) + } + } + + private async initialize(config: WorkerInitConfig) { + this.config = config + + // Initialize state manager + this.stateManager = new CodeIndexStateManager() + + // Set up state change listeners to forward to main thread + this.stateManager.onProgressUpdate((status) => { + // Send progress updates based on the current unit + if (status.currentItemUnit === "files") { + this.sendResponse("progress", { + type: "progress", + processedInBatch: status.processedItems, + totalInBatch: status.totalItems, + currentFile: status.message.includes("Current:") ? status.message.split("Current: ")[1] : undefined, + }) + } else if (status.currentItemUnit === "blocks") { + this.sendResponse("blockProgress", { + type: "blockProgress", + blocksIndexed: status.processedItems, + totalBlocks: status.totalItems, + }) + } + + // Always send status updates + this.sendResponse("status", { + type: "status", + state: status.systemStatus, + message: status.message, + }) + }) + + // Initialize cache manager + this.cacheManager = new CacheManager( + { globalStorageUri: { fsPath: config.contextPath } } as any, + config.workspacePath, + ) + await this.cacheManager.initialize() + + // Create a mock config manager with the provided config + const configManager = this.createMockConfigManager(config) + + // Initialize service factory + const serviceFactory = new CodeIndexServiceFactory(configManager, config.workspacePath, this.cacheManager) + + // Load .gitignore + const ignoreInstance = ignore() + const ignorePath = path.join(config.workspacePath, ".gitignore") + try { + const content = await fs.readFile(ignorePath, "utf8") + ignoreInstance.add(content) + ignoreInstance.add(".gitignore") + } catch (error) { + console.error("Failed to load .gitignore:", error) + } + + // Create services + const services = serviceFactory.createServices( + { globalStorageUri: { fsPath: config.contextPath } } as any, + this.cacheManager, + ignoreInstance, + ) + + this.embedder = services.embedder + this.vectorStore = services.vectorStore + this.scanner = services.scanner + this.fileWatcher = services.fileWatcher + + // Initialize orchestrator + this.orchestrator = new CodeIndexOrchestrator( + configManager, + this.stateManager, + config.workspacePath, + this.cacheManager, + this.vectorStore, + this.scanner, + this.fileWatcher, + ) + + // Initialize search service + this.searchService = new CodeIndexSearchService( + configManager, + this.stateManager, + this.embedder, + this.vectorStore, + ) + } + + private createMockConfigManager(config: WorkerInitConfig): CodeIndexConfigManager { + // Create a mock config manager that returns the worker config + return { + isFeatureEnabled: config.isFeatureEnabled, + isFeatureConfigured: config.isFeatureConfigured, + currentQdrantUrl: config.qdrantUrl || "http://localhost:6333", + currentEmbedderProvider: config.embedderProvider || "openai", + currentEmbedderBaseUrl: config.embedderBaseUrl, + currentEmbedderModelId: config.embedderModelId, + currentEmbedderApiKey: config.embedderApiKey, + currentSearchMinScore: config.searchMinScore || 0.7, + loadConfiguration: async () => ({ requiresRestart: false }), + } as any + } + + private async startIndexing() { + if (!this.orchestrator) { + throw new Error("Worker not initialized") + } + await this.orchestrator.startIndexing() + } + + private async stopIndexing() { + if (!this.orchestrator) { + throw new Error("Worker not initialized") + } + await this.orchestrator.stopWatcher() + } + + private async clearIndex() { + if (!this.orchestrator || !this.cacheManager) { + throw new Error("Worker not initialized") + } + await this.orchestrator.clearIndexData() + await this.cacheManager.clearCacheFile() + this.sendResponse("clear", { type: "cleared", success: true }) + } + + private async search(query: string, directoryPrefix?: string): Promise { + if (!this.searchService) { + throw new Error("Worker not initialized") + } + return await this.searchService.searchIndex(query, directoryPrefix) + } + + private sendResponse(id: string, response: WorkerResponse) { + if (parentPort) { + const message: WorkerMessage = { id, payload: response } + parentPort.postMessage(message) + } + } +} + +// Start the worker +new IndexingWorker()