Merge branch 'linux_stress_test'
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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.
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
|
||||
@@ -308,7 +308,7 @@ export async function evalServerCommand(options?: EvalServerOptions): Promise<vo
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
const repos = backend.listRepos();
|
||||
const repos = await backend.listRepos();
|
||||
console.error(`GitNexus eval-server: ${repos.length} repo(s) loaded: ${repos.map(r => r.name).join(', ')}`);
|
||||
|
||||
let idleTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
|
||||
@@ -45,7 +45,7 @@ export const mcpCommand = async () => {
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
const repoNames = backend.listRepos().map(r => r.name);
|
||||
const repoNames = (await backend.listRepos()).map(r => r.name);
|
||||
console.error(`GitNexus: MCP server starting with ${repoNames.length} repo(s): ${repoNames.join(', ')}`);
|
||||
|
||||
// Start MCP server (serves all repos)
|
||||
|
||||
@@ -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';
|
||||
|
||||
|
||||
@@ -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;
|
||||
};
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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 });
|
||||
}
|
||||
});
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -76,10 +76,23 @@ export class LocalBackend {
|
||||
* Returns true if at least one repo is available.
|
||||
*/
|
||||
async init(): Promise<boolean> {
|
||||
await this.refreshRepos();
|
||||
return this.repos.size > 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* Re-read the global registry and update the in-memory repo map.
|
||||
* New repos are added, existing repos are updated, removed repos are pruned.
|
||||
* KuzuDB connections for removed repos are NOT closed (they idle-timeout naturally).
|
||||
*/
|
||||
private async refreshRepos(): Promise<void> {
|
||||
const entries = await listRegisteredRepos({ validate: true });
|
||||
const freshIds = new Set<string>();
|
||||
|
||||
for (const entry of entries) {
|
||||
const id = this.repoId(entry.name, entry.path);
|
||||
freshIds.add(id);
|
||||
|
||||
const storagePath = entry.storagePath;
|
||||
const kuzuPath = path.join(storagePath, 'kuzu');
|
||||
|
||||
@@ -109,7 +122,14 @@ export class LocalBackend {
|
||||
});
|
||||
}
|
||||
|
||||
return this.repos.size > 0;
|
||||
// Prune repos that no longer exist in the registry
|
||||
for (const id of this.repos.keys()) {
|
||||
if (!freshIds.has(id)) {
|
||||
this.repos.delete(id);
|
||||
this.contextCache.delete(id);
|
||||
this.initializedRepos.delete(id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -136,11 +156,38 @@ export class LocalBackend {
|
||||
* - If repoParam is given, match by name or path
|
||||
* - If only 1 repo, use it
|
||||
* - If 0 or multiple without param, throw with helpful message
|
||||
*
|
||||
* On a miss, re-reads the registry once in case a new repo was indexed
|
||||
* while the MCP server was running.
|
||||
*/
|
||||
resolveRepo(repoParam?: string): RepoHandle {
|
||||
async resolveRepo(repoParam?: string): Promise<RepoHandle> {
|
||||
const result = this.resolveRepoFromCache(repoParam);
|
||||
if (result) return result;
|
||||
|
||||
// Miss — refresh registry and try once more
|
||||
await this.refreshRepos();
|
||||
const retried = this.resolveRepoFromCache(repoParam);
|
||||
if (retried) return retried;
|
||||
|
||||
// Still no match — throw with helpful message
|
||||
if (this.repos.size === 0) {
|
||||
throw new Error('No indexed repositories. Run: gitnexus analyze');
|
||||
}
|
||||
if (repoParam) {
|
||||
const names = [...this.repos.values()].map(h => h.name);
|
||||
throw new Error(`Repository "${repoParam}" not found. Available: ${names.join(', ')}`);
|
||||
}
|
||||
const names = [...this.repos.values()].map(h => h.name);
|
||||
throw new Error(
|
||||
`Multiple repositories indexed. Specify which one with the "repo" parameter. Available: ${names.join(', ')}`
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Try to resolve a repo from the in-memory cache. Returns null on miss.
|
||||
*/
|
||||
private resolveRepoFromCache(repoParam?: string): RepoHandle | null {
|
||||
if (this.repos.size === 0) return null;
|
||||
|
||||
if (repoParam) {
|
||||
const paramLower = repoParam.toLowerCase();
|
||||
@@ -159,19 +206,14 @@ export class LocalBackend {
|
||||
for (const handle of this.repos.values()) {
|
||||
if (handle.name.toLowerCase().includes(paramLower)) return handle;
|
||||
}
|
||||
|
||||
const names = [...this.repos.values()].map(h => h.name);
|
||||
throw new Error(`Repository "${repoParam}" not found. Available: ${names.join(', ')}`);
|
||||
return null;
|
||||
}
|
||||
|
||||
if (this.repos.size === 1) {
|
||||
return this.repos.values().next().value!;
|
||||
}
|
||||
|
||||
const names = [...this.repos.values()].map(h => h.name);
|
||||
throw new Error(
|
||||
`Multiple repositories indexed. Specify which one with the "repo" parameter. Available: ${names.join(', ')}`
|
||||
);
|
||||
return null; // Multiple repos, no param — ambiguous
|
||||
}
|
||||
|
||||
// ─── Lazy KuzuDB Init ────────────────────────────────────────────
|
||||
@@ -210,8 +252,11 @@ export class LocalBackend {
|
||||
|
||||
/**
|
||||
* List all registered repos with their metadata.
|
||||
* Re-reads the global registry so newly indexed repos are discovered
|
||||
* without restarting the MCP server.
|
||||
*/
|
||||
listRepos(): Array<{ name: string; path: string; indexedAt: string; lastCommit: string; stats?: any }> {
|
||||
async listRepos(): Promise<Array<{ name: string; path: string; indexedAt: string; lastCommit: string; stats?: any }>> {
|
||||
await this.refreshRepos();
|
||||
return [...this.repos.values()].map(h => ({
|
||||
name: h.name,
|
||||
path: h.repoPath,
|
||||
@@ -228,8 +273,8 @@ export class LocalBackend {
|
||||
return this.listRepos();
|
||||
}
|
||||
|
||||
// Resolve repo from optional param
|
||||
const repo = this.resolveRepo(params?.repo);
|
||||
// Resolve repo from optional param (re-reads registry on miss)
|
||||
const repo = await this.resolveRepo(params?.repo);
|
||||
|
||||
switch (method) {
|
||||
case 'query':
|
||||
@@ -1289,7 +1334,7 @@ export class LocalBackend {
|
||||
* Used by getClustersResource — avoids legacy overview() dispatch.
|
||||
*/
|
||||
async queryClusters(repoName?: string, limit = 100): Promise<{ clusters: any[] }> {
|
||||
const repo = this.resolveRepo(repoName);
|
||||
const repo = await this.resolveRepo(repoName);
|
||||
await this.ensureInitialized(repo.id);
|
||||
|
||||
try {
|
||||
@@ -1318,7 +1363,7 @@ export class LocalBackend {
|
||||
* Used by getProcessesResource — avoids legacy overview() dispatch.
|
||||
*/
|
||||
async queryProcesses(repoName?: string, limit = 50): Promise<{ processes: any[] }> {
|
||||
const repo = this.resolveRepo(repoName);
|
||||
const repo = await this.resolveRepo(repoName);
|
||||
await this.ensureInitialized(repo.id);
|
||||
|
||||
try {
|
||||
@@ -1347,7 +1392,7 @@ export class LocalBackend {
|
||||
* Used by getClusterDetailResource.
|
||||
*/
|
||||
async queryClusterDetail(name: string, repoName?: string): Promise<any> {
|
||||
const repo = this.resolveRepo(repoName);
|
||||
const repo = await this.resolveRepo(repoName);
|
||||
await this.ensureInitialized(repo.id);
|
||||
|
||||
const escaped = name.replace(/'/g, "''");
|
||||
@@ -1398,7 +1443,7 @@ export class LocalBackend {
|
||||
* Used by getProcessDetailResource.
|
||||
*/
|
||||
async queryProcessDetail(name: string, repoName?: string): Promise<any> {
|
||||
const repo = this.resolveRepo(repoName);
|
||||
const repo = await this.resolveRepo(repoName);
|
||||
await this.ensureInitialized(repo.id);
|
||||
|
||||
const escaped = name.replace(/'/g, "''");
|
||||
|
||||
@@ -153,8 +153,8 @@ export async function readResource(uri: string, backend: LocalBackend): Promise<
|
||||
/**
|
||||
* Repos resource — list all indexed repositories
|
||||
*/
|
||||
function getReposResource(backend: LocalBackend): string {
|
||||
const repos = backend.listRepos();
|
||||
async function getReposResource(backend: LocalBackend): Promise<string> {
|
||||
const repos = await backend.listRepos();
|
||||
|
||||
if (repos.length === 0) {
|
||||
return 'repos: []\n# No repositories indexed. Run: gitnexus analyze';
|
||||
@@ -187,7 +187,7 @@ function getReposResource(backend: LocalBackend): string {
|
||||
*/
|
||||
async function getContextResource(backend: LocalBackend, repoName?: string): Promise<string> {
|
||||
// Resolve repo
|
||||
const repo = backend.resolveRepo(repoName);
|
||||
const repo = await backend.resolveRepo(repoName);
|
||||
const repoId = repo.name.toLowerCase();
|
||||
const context = backend.getContext(repoId) || backend.getContext();
|
||||
|
||||
@@ -432,8 +432,8 @@ async function getProcessDetailResource(name: string, backend: LocalBackend, rep
|
||||
* Useful for `gitnexus setup` onboarding or dynamic content injection.
|
||||
*/
|
||||
async function getSetupResource(backend: LocalBackend): Promise<string> {
|
||||
const repos = backend.listRepos();
|
||||
|
||||
const repos = await backend.listRepos();
|
||||
|
||||
if (repos.length === 0) {
|
||||
return '# GitNexus\n\nNo repositories indexed. Run: `npx gitnexus analyze` in a repository.';
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user