fix: worker pool robustness — timeout, exit handler, 512KB guards, error reporting

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
abhigyanpatwari 2026-02-20 20:24:48 +05:30
parent 79501abb6c
commit 33aa774f03
8 changed files with 99 additions and 24 deletions

View file

@ -65,7 +65,13 @@
"WebFetch(domain:smartlogic.io)",
"Bash(ls:*)",
"Bash(wc:*)",
"Bash(grep:*)"
"Bash(grep:*)",
"Bash(powershell -Command:*)",
"Bash(cmd /c \"dir /s C:\\\\Users\\\\ADMIN\\\\.cache\\\\huggingface 2>nul | findstr /i \"\"File\\(s\\)\"\"\")",
"Bash(du:*)",
"mcp__desktop-commander__list_directory",
"Bash(python3 -c \":*)",
"mcp__gitnexus__detect_changes"
]
},
"enableAllProjectMcpServers": true,

View file

@ -3,7 +3,7 @@
<!-- gitnexus:start -->
# GitNexus MCP
This project is indexed by GitNexus as **GitnexusV2** (1295 symbols, 3262 relationships, 99 execution flows).
This project is indexed by GitNexus as **GitnexusV2** (1312 symbols, 3315 relationships, 101 execution flows).
GitNexus provides a knowledge graph over this codebase — call chains, blast radius, execution flows, and semantic search.

View file

@ -1,7 +1,7 @@
<!-- gitnexus:start -->
# GitNexus MCP
This project is indexed by GitNexus as **GitnexusV2** (1295 symbols, 3262 relationships, 99 execution flows).
This project is indexed by GitNexus as **GitnexusV2** (1312 symbols, 3315 relationships, 101 execution flows).
GitNexus provides a knowledge graph over this codebase — call chains, blast radius, execution flows, and semantic search.

View file

@ -1,12 +1,19 @@
/**
* Embedder Module
*
*
* Singleton factory for transformers.js embedding pipeline.
* Handles model loading, caching, and both single and batch embedding operations.
*
*
* Uses snowflake-arctic-embed-xs by default (22M params, 384 dims, ~90MB)
*/
// Suppress ONNX Runtime native warnings (e.g. VerifyEachNodeIsAssignedToAnEp)
// Must be set BEFORE onnxruntime-node is imported by transformers.js
// Level 3 = Error only (skips Warning/Info)
if (!process.env.ORT_LOG_LEVEL) {
process.env.ORT_LOG_LEVEL = '3';
}
import { pipeline, env, type FeatureExtractionPipeline } from '@huggingface/transformers';
import { DEFAULT_EMBEDDING_CONFIG, type EmbeddingConfig, type ModelProgress } from './types.js';

View file

@ -10,6 +10,9 @@ export interface FileEntry {
const READ_CONCURRENCY = 32;
/** Skip files larger than 512KB — they're usually generated/vendored and crash tree-sitter */
const MAX_FILE_SIZE = 512 * 1024;
export const walkRepository = async (
repoPath: string,
onProgress?: (current: number, total: number, filePath: string) => void
@ -23,19 +26,26 @@ export const walkRepository = async (
const filtered = files.filter(file => !shouldIgnorePath(file));
const entries: FileEntry[] = [];
let processed = 0;
let skippedLarge = 0;
for (let start = 0; start < filtered.length; start += READ_CONCURRENCY) {
const batch = filtered.slice(start, start + READ_CONCURRENCY);
const results = await Promise.allSettled(
batch.map(relativePath =>
fs.readFile(path.join(repoPath, relativePath), 'utf-8')
.then(content => ({ path: relativePath.replace(/\\/g, '/'), content }))
)
batch.map(async relativePath => {
const fullPath = path.join(repoPath, relativePath);
const stat = await fs.stat(fullPath);
if (stat.size > MAX_FILE_SIZE) {
skippedLarge++;
return null;
}
const content = await fs.readFile(fullPath, 'utf-8');
return { path: relativePath.replace(/\\/g, '/'), content };
})
);
for (const result of results) {
processed++;
if (result.status === 'fulfilled') {
if (result.status === 'fulfilled' && result.value !== null) {
entries.push(result.value);
onProgress?.(processed, filtered.length, result.value.path);
} else {
@ -44,5 +54,9 @@ export const walkRepository = async (
}
}
if (skippedLarge > 0) {
console.warn(` Skipped ${skippedLarge} files larger than ${MAX_FILE_SIZE / 1024}KB`);
}
return entries;
};

View file

@ -209,6 +209,9 @@ const processParsingSequential = async (
if (!language) continue;
// Skip very large files — they can crash tree-sitter or cause OOM
if (file.content.length > 512 * 1024) continue;
await loadLanguage(language, file.path);
let tree;
@ -335,7 +338,7 @@ export const processParsing = async (
try {
return await processParsingWithWorkers(graph, files, symbolTable, astCache, workerPool, onFileProgress);
} catch (err) {
console.warn('Worker pool parsing failed, falling back to sequential:', err);
console.warn('Worker pool parsing failed, falling back to sequential:', err instanceof Error ? err.message : err);
}
}

View file

@ -402,6 +402,9 @@ const processFileGroup = (
}
for (const file of files) {
// Skip very large files — they can crash tree-sitter or cause OOM
if (file.content.length > 512 * 1024) continue;
let tree;
try {
tree = parser.parse(file.content, undefined, { bufferSize: 1024 * 256 });
@ -528,8 +531,13 @@ const processFileGroup = (
// ============================================================================
parentPort!.on('message', (files: ParseWorkerInput[]) => {
const result = processBatch(files, (filesProcessed) => {
parentPort!.postMessage({ type: 'progress', filesProcessed });
});
parentPort!.postMessage({ type: 'result', data: result });
try {
const result = processBatch(files, (filesProcessed) => {
parentPort!.postMessage({ type: 'progress', filesProcessed });
});
parentPort!.postMessage({ type: 'result', data: result });
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
parentPort!.postMessage({ type: 'error', error: message });
}
});

View file

@ -50,29 +50,66 @@ export const createWorkerPool = (workerUrl: URL, poolSize?: number): WorkerPool
const promises = chunks.map((chunk, i) => {
const worker = workers[i];
return new Promise<TResult>((resolve, reject) => {
let settled = false;
const cleanup = () => {
clearTimeout(timer);
worker.removeListener('message', handler);
worker.removeListener('error', errorHandler);
worker.removeListener('exit', exitHandler);
};
const timer = setTimeout(() => {
if (!settled) {
settled = true;
cleanup();
reject(new Error(`Worker ${i} timed out after 5 minutes (chunk: ${chunk.length} items). Worker may have crashed or is processing too much data.`));
}
}, 5 * 60 * 1000);
const handler = (msg: any) => {
if (settled) return;
if (msg && msg.type === 'progress') {
// Intermediate progress from worker
workerProgress[i] = msg.filesProcessed;
if (onProgress) {
const total = workerProgress.reduce((a, b) => a + b, 0);
onProgress(total);
}
} else if (msg && msg.type === 'error') {
// Error reported by worker via postMessage
settled = true;
cleanup();
reject(new Error(`Worker ${i} error: ${msg.error}`));
} else if (msg && msg.type === 'result') {
// Final result
worker.removeListener('message', handler);
settled = true;
cleanup();
resolve(msg.data);
} else {
// Legacy: treat any non-typed message as result (backward compat)
worker.removeListener('message', handler);
// Legacy: treat any non-typed message as result
settled = true;
cleanup();
resolve(msg);
}
};
const errorHandler = (err: any) => {
if (!settled) {
settled = true;
cleanup();
reject(err);
}
};
const exitHandler = (code: number) => {
if (!settled) {
settled = true;
cleanup();
reject(new Error(`Worker ${i} exited unexpectedly with code ${code}. This usually indicates an out-of-memory crash or native addon failure.`));
}
};
worker.on('message', handler);
worker.once('error', (err) => {
worker.removeListener('message', handler);
reject(err);
});
worker.once('error', errorHandler);
worker.once('exit', exitHandler);
worker.postMessage(chunk);
});
});