fix: cancel abandoned highlight() semantic work and harden WASM engine recovery

highlight() used Promise.race with a 10s timeout, but the losing
semantic phase promise continued running 25 WASM micro-batches,
saturating the event loop and degrading all subsequent operations
(find() going from ~200ms to ~10,000ms).

Add AbortController to highlight() so the semantic phase stops
immediately on timeout or error. Pass abort signal through
embedBatch() → EmbeddingManager → micro-batch loop.

Also add defensive hardening:
- CandleEmbeddingEngine: try/catch around WASM calls resets engine
  state on failure so next call triggers re-initialization
- WASMEmbeddingEngine: initialize() now checks underlying Candle
  engine state, not just its own flag, completing the recovery chain
This commit is contained in:
David Snelling 2026-01-27 18:26:37 -08:00
parent 279fccebfe
commit f8dd93c93c
4 changed files with 54 additions and 23 deletions

View file

@ -4647,7 +4647,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* // embeddings.length === 3
* // embeddings[0].length === 384
*/
async embedBatch(texts: string[]): Promise<number[][]> {
async embedBatch(texts: string[], options?: { signal?: AbortSignal }): Promise<number[][]> {
await this.ensureInitialized()
if (texts.length === 0) {
@ -4655,7 +4655,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
}
// Use native WASM batch API: single forward pass instead of N individual calls
return await embeddingManager.embedBatch(texts)
return await embeddingManager.embedBatch(texts, options)
}
/**
@ -4784,20 +4784,25 @@ export class Brainy<T = any> implements BrainyInterface<T> {
}
// === PHASE 2: Find semantic matches with timeout fallback ===
// AbortController ensures the background semantic work (WASM batch embedding)
// is cancelled on timeout or error, preventing event loop saturation and
// WASM engine crashes from abandoned promises.
const SEMANTIC_TIMEOUT_MS = 10_000
const abortController = new AbortController()
try {
const semanticResult = await Promise.race([
this.highlightSemanticPhase(query, chunks, threshold, highlightMap),
this.highlightSemanticPhase(query, chunks, threshold, highlightMap, abortController.signal),
new Promise<'timeout'>((resolve) => setTimeout(() => resolve('timeout'), SEMANTIC_TIMEOUT_MS))
])
if (semanticResult === 'timeout') {
// Return Phase 1 text-only results on timeout
abortController.abort()
const textHighlights = Array.from(highlightMap.values())
return textHighlights.sort((a, b) => b.score - a.score)
}
} catch {
// On any error in semantic phase, return Phase 1 results
abortController.abort()
const textHighlights = Array.from(highlightMap.values())
return textHighlights.sort((a, b) => b.score - a.score)
}
@ -4815,14 +4820,19 @@ export class Brainy<T = any> implements BrainyInterface<T> {
query: string,
chunks: Array<{ text: string, position: [number, number], contentCategory?: import('./types/brainy.types.js').ContentCategory }>,
threshold: number,
highlightMap: Map<string, import('./types/brainy.types.js').Highlight>
highlightMap: Map<string, import('./types/brainy.types.js').Highlight>,
signal?: AbortSignal
): Promise<void> {
if (signal?.aborted) return
// Get query embedding
const queryVector = await this.embed(query)
if (signal?.aborted) return
// Batch embed all chunks using native WASM batch API
const chunkTexts = chunks.map(c => c.text)
const chunkVectors = await this.embedBatch(chunkTexts)
const chunkVectors = await this.embedBatch(chunkTexts, { signal })
if (signal?.aborted) return
// Calculate semantic similarities
for (let i = 0; i < chunks.length; i++) {

View file

@ -219,7 +219,7 @@ export class EmbeddingManager {
* @param texts Array of strings to embed
* @returns Array of embedding vectors (384 dimensions each)
*/
async embedBatch(texts: string[]): Promise<number[][]> {
async embedBatch(texts: string[], options?: { signal?: AbortSignal }): Promise<number[][]> {
if (texts.length === 0) return []
const isTestMode =
@ -248,6 +248,10 @@ export class EmbeddingManager {
// so other requests (HTTP, timers, I/O) can proceed
const allResults: number[][] = []
for (let i = 0; i < texts.length; i += MICRO_BATCH_SIZE) {
if (options?.signal?.aborted) {
return allResults
}
const batch = texts.slice(i, i + MICRO_BATCH_SIZE)
const batchResults = await this.engine.embedBatch(batch)
allResults.push(...batchResults)

View file

@ -205,19 +205,27 @@ export class CandleEmbeddingEngine {
throw new Error('Engine not properly initialized')
}
const startTime = Date.now()
try {
const startTime = Date.now()
const embedding = this.engine.embed(text)
const embeddingArray = Array.from(embedding)
const embedding = this.engine.embed(text)
const embeddingArray = Array.from(embedding)
const processingTimeMs = Date.now() - startTime
this.embedCount++
this.totalProcessingTimeMs += processingTimeMs
const processingTimeMs = Date.now() - startTime
this.embedCount++
this.totalProcessingTimeMs += processingTimeMs
return {
embedding: embeddingArray,
tokenCount: 0, // Candle handles tokenization internally
processingTimeMs,
return {
embedding: embeddingArray,
tokenCount: 0, // Candle handles tokenization internally
processingTimeMs,
}
} catch (error) {
console.error('WASM embed failed, marking engine for re-initialization:', error)
this.initialized = false
this.engine = null
this.wasmModule = null
throw error
}
}
@ -237,10 +245,18 @@ export class CandleEmbeddingEngine {
return []
}
const embeddings = this.engine.embed_batch(texts)
this.embedCount += texts.length
try {
const embeddings = this.engine.embed_batch(texts)
this.embedCount += texts.length
return embeddings.map((e) => Array.from(e))
return embeddings.map((e) => Array.from(e))
} catch (error) {
console.error('WASM embedBatch failed, marking engine for re-initialization:', error)
this.initialized = false
this.engine = null
this.wasmModule = null
throw error
}
}
/**

View file

@ -52,8 +52,9 @@ export class WASMEmbeddingEngine {
* Initialize the engine
*/
async initialize(): Promise<void> {
// Already initialized
if (this.initialized) {
// Only skip if BOTH this layer and the underlying Candle engine are initialized.
// If Candle crashed and reset itself, we must re-initialize even though our flag is true.
if (this.initialized && this.candleEngine.isInitialized()) {
return
}