/** * Worker pool for parallel query processing. * Dispatches queries to workers via shortest-queue for concurrent execution. * Falls back to sequential processing when workers are unavailable. */ import { HNSWIndex } from './HNSWIndex.js'; export declare class WorkerPool { private workers; private nextQueryId; private initialized; private initializing; private fallbackIndex; private activeIndex; private defaultTimeout; /** Dimension of vectors, cached from init for batch dispatch */ private dimension; private metric; private quantizationEnabled; private indexLeaseAcquired; private numWorkers; /** * Create a worker pool for parallel HNSW search. * * @param numWorkers Number of workers (default: available hardware concurrency or 4) * @param timeout Default per-query timeout in ms (default: 30000) */ constructor(numWorkers?: number, timeout?: number); /** * Create and initialize a worker pool in one step. * * @param index The HNSW index to distribute to workers * @param numWorkers Number of workers (default: auto-detect) * @param timeout Default per-query timeout in ms (default: 30000) * @returns Initialized WorkerPool */ static create(index: HNSWIndex, numWorkers?: number, timeout?: number): Promise; private assertInitialized; private validateQueryVector; private validateSearchParams; private validateCandidateMultiplier; private validateBatchQueries; private rejectHandle; private terminateHandles; private handleWorkerFailure; /** * Initialize workers with shared index data. * When the index uses SharedArrayBuffer, workers receive zero-copy * views into the same vector memory. Otherwise, data is copied. * * @param index The HNSW index to distribute to workers */ init(index: HNSWIndex): Promise; /** * Get the worker with the fewest pending queries (shortest-queue dispatch). */ private getShortestQueueWorker; /** * Search for k nearest neighbors using worker pool. * Dispatches to the worker with fewest pending queries. * * @param query Query vector * @param k Number of results * @param efSearch Search effort parameter * @param timeout Per-query timeout in ms (default: pool default) * @returns Array of {id, distance} results */ search(query: Float32Array, k: number, efSearch?: number, timeout?: number): Promise>; /** * Quantized search using worker pool. * Workers perform int8 candidate scan + float32 rescore. * * @param query Query vector * @param k Number of results * @param candidateMultiplier Multiplier for rescore candidates * @param efSearch Search effort parameter * @returns Array of {id, distance} results */ searchQuantized(query: Float32Array, k: number, candidateMultiplier?: number, efSearch?: number): Promise>; /** * Batch search: dispatch multiple queries across workers using batch messages. * Sends one postMessage per worker instead of one per query. * * @param queries Array of query vectors * @param k Number of results per query * @param efSearch Search effort parameter * @returns Array of results, one per query */ searchBatch(queries: Float32Array[], k: number, efSearch?: number): Promise>>; /** * Batch quantized search across workers. */ searchBatchQuantized(queries: Float32Array[], k: number, candidateMultiplier?: number, efSearch?: number): Promise>>; /** * Internal batch dispatch: partition queries across workers and send one * postMessage per worker with all queries concatenated into a single Float32Array. */ private _dispatchBatch; /** * Send incremental graph updates to all workers. * When shared graph SABs are used, this only updates metadata. * Workers see neighbor data changes via SharedArrayBuffer directly. */ broadcastGraphUpdate(newNodes: Array<{ id: number; neighbors: number[][]; }>, entryPointId?: number, maxLevel?: number, nodeCount?: number, graphGeneration?: number, publishedGraphGeneration?: number): void; /** * Broadcast only metadata updates (entry point, max level, node count). * Used when shared graph SABs handle neighbor data directly. */ broadcastMetadataUpdate(entryPointId: number, maxLevel: number, nodeCount: number, graphGeneration?: number): void; /** * Terminate all workers and clean up. */ destroy(reason?: Error): void; private assertIndexLease; /** * Get the number of active workers. */ get workerCount(): number; } //# sourceMappingURL=WorkerPool.d.ts.map