From 2bdf1d75a4c625ab36949c3e9cf9d1c90b5c7bfe Mon Sep 17 00:00:00 2001 From: qinsehm1128 <1048014860@qq.com> Date: Sun, 12 Apr 2026 14:17:33 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=A2=9E=E5=8A=A0=20MCP=20=E7=B4=A2?= =?UTF-8?q?=E5=BC=95=E5=86=99=E5=B7=A5=E5=85=B7=E4=B8=8E=20embedding=20?= =?UTF-8?q?=E9=85=8D=E7=BD=AE=E8=A7=A3=E6=9E=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- gitnexus/src/core/embeddings/config.ts | 142 ++++++ gitnexus/src/core/embeddings/embedder.ts | 5 +- gitnexus/src/core/embeddings/http-client.ts | 66 +-- gitnexus/src/core/embeddings/index.ts | 1 + .../core/index-jobs/analyze-job-service.ts | 275 ++++++++++++ .../src/core/index-jobs/embed-job-service.ts | 138 ++++++ gitnexus/src/core/index-jobs/index.ts | 5 + gitnexus/src/core/index-jobs/job-manager.ts | 190 ++++++++ .../src/core/index-jobs/job-query-service.ts | 51 +++ .../src/core/index-jobs/repo-lock-manager.ts | 26 ++ gitnexus/src/core/lbug/schema.ts | 9 +- gitnexus/src/mcp/core/embedder.ts | 3 +- gitnexus/src/mcp/local/local-backend.ts | 152 +++++++ gitnexus/src/mcp/tools.ts | 85 ++++ gitnexus/src/server/analyze-job.ts | 204 +-------- gitnexus/src/server/api.ts | 416 +++++------------- gitnexus/src/storage/repo-manager.ts | 25 ++ gitnexus/test/unit/calltool-dispatch.test.ts | 88 ++++ gitnexus/test/unit/embedding-config.test.ts | 189 ++++++++ gitnexus/test/unit/http-embedder.test.ts | 73 ++- gitnexus/test/unit/tools.test.ts | 31 +- 21 files changed, 1608 insertions(+), 566 deletions(-) create mode 100644 gitnexus/src/core/embeddings/config.ts create mode 100644 gitnexus/src/core/index-jobs/analyze-job-service.ts create mode 100644 gitnexus/src/core/index-jobs/embed-job-service.ts create mode 100644 gitnexus/src/core/index-jobs/index.ts create mode 100644 gitnexus/src/core/index-jobs/job-manager.ts create mode 100644 gitnexus/src/core/index-jobs/job-query-service.ts create mode 100644 gitnexus/src/core/index-jobs/repo-lock-manager.ts create mode 100644 gitnexus/test/unit/embedding-config.test.ts diff --git a/gitnexus/src/core/embeddings/config.ts b/gitnexus/src/core/embeddings/config.ts new file mode 100644 index 000000000..012c982aa --- /dev/null +++ b/gitnexus/src/core/embeddings/config.ts @@ -0,0 +1,142 @@ +import { + loadCLIConfig, + loadCLIConfigSync, + type CLIConfig, + type CLIEmbeddingConfig, + type CLIEmbeddingProvider, +} from '../../storage/repo-manager.js'; +import { DEFAULT_EMBEDDING_CONFIG } from './types.js'; + +export type EmbeddingConfigSource = 'overrides' | 'config' | 'env' | 'default'; +export type ResolvedEmbeddingMode = 'local' | 'http'; + +export interface EmbeddingConfigOverrides extends Partial {} + +export interface ResolvedEmbeddingConfig { + mode: ResolvedEmbeddingMode; + provider: CLIEmbeddingProvider; + model: string; + dimensions: number; + baseUrl?: string; + apiKey: string; + explicitDimensionsSource?: Exclude; +} + +interface ResolvedValue { + value: T; + source: EmbeddingConfigSource; +} + +const trimToUndefined = (value: string | undefined): string | undefined => { + if (value === undefined) return undefined; + const trimmed = value.trim(); + return trimmed.length > 0 ? trimmed : undefined; +}; + +const resolveValue = ( + overridesValue: T | undefined, + configValue: T | undefined, + envValue: T | undefined, + defaultValue: T, +): ResolvedValue => { + if (overridesValue !== undefined) return { value: overridesValue, source: 'overrides' }; + if (configValue !== undefined) return { value: configValue, source: 'config' }; + if (envValue !== undefined) return { value: envValue, source: 'env' }; + return { value: defaultValue, source: 'default' }; +}; + +const parsePositiveInteger = ( + rawValue: number | string | undefined, + source: Exclude, +): number | undefined => { + if (rawValue === undefined) return undefined; + + const parsed = + typeof rawValue === 'number' + ? Number.isInteger(rawValue) + ? rawValue + : Number.NaN + : parseInt(rawValue, 10); + + if (Number.isNaN(parsed) || parsed <= 0) { + const fieldLabel = + source === 'config' + ? 'embedding.dimensions' + : source === 'overrides' + ? 'embedding dimensions override' + : 'GITNEXUS_EMBEDDING_DIMS'; + throw new Error(`${fieldLabel} must be a positive integer, got "${rawValue}"`); + } + + return parsed; +}; + +const getSavedEmbeddingConfig = (savedConfig: CLIConfig): CLIEmbeddingConfig => + savedConfig.embedding ?? {}; + +const resolveEmbeddingConfigFromSaved = ( + savedConfig: CLIConfig, + overrides: EmbeddingConfigOverrides = {}, +): ResolvedEmbeddingConfig => { + const savedEmbeddingConfig = getSavedEmbeddingConfig(savedConfig); + + const provider = resolveValue( + overrides.provider, + savedEmbeddingConfig.provider, + undefined, + undefined, + ); + const baseUrl = resolveValue( + trimToUndefined(overrides.baseUrl), + trimToUndefined(savedEmbeddingConfig.baseUrl), + trimToUndefined(process.env.GITNEXUS_EMBEDDING_URL), + undefined, + ); + const httpModel = resolveValue( + trimToUndefined(overrides.model), + trimToUndefined(savedEmbeddingConfig.model), + trimToUndefined(process.env.GITNEXUS_EMBEDDING_MODEL), + undefined, + ); + + const dimensions = resolveValue( + parsePositiveInteger(overrides.dimensions, 'overrides'), + parsePositiveInteger(savedEmbeddingConfig.dimensions, 'config'), + parsePositiveInteger(process.env.GITNEXUS_EMBEDDING_DIMS, 'env'), + DEFAULT_EMBEDDING_CONFIG.dimensions, + ); + + const forceLocal = provider.value === 'local'; + const mode: ResolvedEmbeddingMode = + !forceLocal && baseUrl.value && httpModel.value ? 'http' : 'local'; + + const apiKey = resolveValue( + trimToUndefined(overrides.apiKey), + trimToUndefined(savedEmbeddingConfig.apiKey), + trimToUndefined(process.env.GITNEXUS_EMBEDDING_API_KEY), + undefined, + ); + + return { + mode, + provider: mode === 'http' ? (provider.value ?? 'custom') : 'local', + model: mode === 'http' ? httpModel.value! : DEFAULT_EMBEDDING_CONFIG.modelId, + dimensions: mode === 'http' ? dimensions.value : DEFAULT_EMBEDDING_CONFIG.dimensions, + baseUrl: mode === 'http' ? baseUrl.value!.replace(/\/+$/, '') : undefined, + apiKey: mode === 'http' ? (apiKey.value ?? 'unused') : '', + explicitDimensionsSource: + mode === 'http' && dimensions.source !== 'default' ? dimensions.source : undefined, + }; +}; + +export const resolveEmbeddingConfigSync = ( + overrides: EmbeddingConfigOverrides = {}, +): ResolvedEmbeddingConfig => { + return resolveEmbeddingConfigFromSaved(loadCLIConfigSync(), overrides); +}; + +export const resolveEmbeddingConfig = async ( + overrides: EmbeddingConfigOverrides = {}, +): Promise => { + return resolveEmbeddingConfigFromSaved(await loadCLIConfig(), overrides); +}; diff --git a/gitnexus/src/core/embeddings/embedder.ts b/gitnexus/src/core/embeddings/embedder.ts index e2d26c9e7..ee736e5fd 100644 --- a/gitnexus/src/core/embeddings/embedder.ts +++ b/gitnexus/src/core/embeddings/embedder.ts @@ -20,6 +20,7 @@ import { execFileSync } from 'child_process'; import { join, dirname } from 'path'; import { createRequire } from 'module'; import { DEFAULT_EMBEDDING_CONFIG, type EmbeddingConfig, type ModelProgress } from './types.js'; +import { resolveEmbeddingConfigSync } from './config.js'; import { isHttpMode, getHttpDimensions, httpEmbed } from './http-client.js'; /** @@ -250,13 +251,13 @@ export const isEmbedderReady = (): boolean => { /** * Get the effective embedding dimensions. - * In HTTP mode, uses GITNEXUS_EMBEDDING_DIMS if set, otherwise the default. + * Resolved from overrides/config/env with local defaults as fallback. */ export const getEmbeddingDimensions = (): number => { if (isHttpMode()) { return getHttpDimensions() ?? DEFAULT_EMBEDDING_CONFIG.dimensions; } - return DEFAULT_EMBEDDING_CONFIG.dimensions; + return resolveEmbeddingConfigSync().dimensions; }; /** diff --git a/gitnexus/src/core/embeddings/http-client.ts b/gitnexus/src/core/embeddings/http-client.ts index 85ad79111..f2e2e033f 100644 --- a/gitnexus/src/core/embeddings/http-client.ts +++ b/gitnexus/src/core/embeddings/http-client.ts @@ -5,45 +5,20 @@ * Imported by both the core embedder (batch) and MCP embedder (query). */ +import { resolveEmbeddingConfigSync, type ResolvedEmbeddingConfig } from './config.js'; + const HTTP_TIMEOUT_MS = 30_000; const HTTP_MAX_RETRIES = 2; const HTTP_RETRY_BACKOFF_MS = 1_000; const HTTP_BATCH_SIZE = 64; -const DEFAULT_DIMS = 384; - -interface HttpConfig { - baseUrl: string; - model: string; - apiKey: string; - dimensions?: number; -} /** - * Build config from the current process.env snapshot. - * Returns null when GITNEXUS_EMBEDDING_URL + GITNEXUS_EMBEDDING_MODEL are unset. - * Not cached — env vars are read fresh so late configuration takes effect. + * Resolve HTTP embedding config from overrides/config/env on every call. + * Returns null when HTTP mode is not active. */ -const readConfig = (): HttpConfig | null => { - const baseUrl = process.env.GITNEXUS_EMBEDDING_URL; - const model = process.env.GITNEXUS_EMBEDDING_MODEL; - if (!baseUrl || !model) return null; - - const rawDims = process.env.GITNEXUS_EMBEDDING_DIMS; - let dimensions: number | undefined; - if (rawDims !== undefined) { - const parsed = parseInt(rawDims, 10); - if (Number.isNaN(parsed) || parsed <= 0) { - throw new Error(`GITNEXUS_EMBEDDING_DIMS must be a positive integer, got "${rawDims}"`); - } - dimensions = parsed; - } - - return { - baseUrl: baseUrl.replace(/\/+$/, ''), - model, - apiKey: process.env.GITNEXUS_EMBEDDING_API_KEY ?? 'unused', - dimensions, - }; +const readConfig = (): ResolvedEmbeddingConfig | null => { + const config = resolveEmbeddingConfigSync(); + return config.mode === 'http' ? config : null; }; /** @@ -57,6 +32,19 @@ export const isHttpMode = (): boolean => readConfig() !== null; */ export const getHttpDimensions = (): number | undefined => readConfig()?.dimensions; +const getDimensionHint = (config: ResolvedEmbeddingConfig, actualDimensions: number): string => { + switch (config.explicitDimensionsSource) { + case 'overrides': + return 'Update the embedding dimensions override to match your model output.'; + case 'config': + return 'Update embedding.dimensions in ~/.gitnexus/config.json to match your model output.'; + case 'env': + return 'Update GITNEXUS_EMBEDDING_DIMS to match your model output.'; + default: + return `Set GITNEXUS_EMBEDDING_DIMS=${actualDimensions} to match your model output.`; + } +}; + /** * Return a safe representation of a URL for error messages. * Strips query string (may contain tokens) and userinfo. @@ -168,14 +156,11 @@ export const httpEmbed = async (texts: string[]): Promise => { const vec = new Float32Array(item.embedding); // Fail fast on dimension mismatch rather than inserting bad vectors // into the FLOAT[N] column which would cause a cryptic Kuzu error. - const expected = config.dimensions ?? DEFAULT_DIMS; + const expected = config.dimensions; if (vec.length !== expected) { - const hint = config.dimensions - ? 'Update GITNEXUS_EMBEDDING_DIMS to match your model output.' - : `Set GITNEXUS_EMBEDDING_DIMS=${vec.length} to match your model output.`; throw new Error( `Embedding dimension mismatch: endpoint returned ${vec.length}d vector, ` + - `but expected ${expected}d. ${hint}`, + `but expected ${expected}d. ${getDimensionHint(config, vec.length)}`, ); } @@ -206,14 +191,11 @@ export const httpEmbedQuery = async (text: string): Promise => { const embedding = items[0].embedding; // Same dimension checks as httpEmbed — catch mismatches before they // reach the Kuzu FLOAT[N] cast in search queries. - const expected = config.dimensions ?? DEFAULT_DIMS; + const expected = config.dimensions; if (embedding.length !== expected) { - const hint = config.dimensions - ? 'Update GITNEXUS_EMBEDDING_DIMS to match your model output.' - : `Set GITNEXUS_EMBEDDING_DIMS=${embedding.length} to match your model output.`; throw new Error( `Embedding dimension mismatch: endpoint returned ${embedding.length}d vector, ` + - `but expected ${expected}d. ${hint}`, + `but expected ${expected}d. ${getDimensionHint(config, embedding.length)}`, ); } return embedding; diff --git a/gitnexus/src/core/embeddings/index.ts b/gitnexus/src/core/embeddings/index.ts index d4de4e0a8..abeaf0c70 100644 --- a/gitnexus/src/core/embeddings/index.ts +++ b/gitnexus/src/core/embeddings/index.ts @@ -5,6 +5,7 @@ */ export * from './types.js'; +export * from './config.js'; export * from './http-client.js'; export * from './embedder.js'; export * from './text-generator.js'; diff --git a/gitnexus/src/core/index-jobs/analyze-job-service.ts b/gitnexus/src/core/index-jobs/analyze-job-service.ts new file mode 100644 index 000000000..6bd3553da --- /dev/null +++ b/gitnexus/src/core/index-jobs/analyze-job-service.ts @@ -0,0 +1,275 @@ +import path from 'path'; +import { fork } from 'child_process'; +import { createRequire } from 'node:module'; +import { fileURLToPath, pathToFileURL } from 'url'; +import { getStoragePath } from '../../storage/repo-manager.js'; +import { cloneOrPull, extractRepoName, getCloneDir } from '../../server/git-clone.js'; +import { type AnalyzeOptions } from '../run-analyze.js'; +import { JobManager, type AnalyzeJob } from './job-manager.js'; +import { RepoLockManager } from './repo-lock-manager.js'; + +const _require = createRequire(import.meta.url); +const MAX_WORKER_RETRIES = 2; + +interface StartMessage { + type: 'start'; + repoPath: string; + options: AnalyzeOptions; +} + +interface ProgressMessage { + type: 'progress'; + phase: string; + percent: number; + message: string; +} + +interface CompleteMessage { + type: 'complete'; + result: { repoName: string }; +} + +interface ErrorMessage { + type: 'error'; + message: string; +} + +type WorkerMessage = ProgressMessage | CompleteMessage | ErrorMessage; + +export interface StartAnalyzeJobParams { + repoUrl?: string; + repoPath?: string; + force?: boolean; + embeddings?: boolean; +} + +export interface AnalyzeJobServiceOptions { + backendInit: () => Promise; + repoLocks: RepoLockManager; + jobManager?: JobManager; + logger?: Pick; +} + +export class AnalyzeJobService { + readonly jobManager: JobManager; + + private readonly backendInit: () => Promise; + private readonly repoLocks: RepoLockManager; + private readonly logger: Pick; + private readonly workerPath: string; + private readonly tsxHookArgs: string[]; + + constructor(options: AnalyzeJobServiceOptions) { + this.backendInit = options.backendInit; + this.repoLocks = options.repoLocks; + this.jobManager = options.jobManager ?? new JobManager(); + this.logger = options.logger ?? console; + + const callerPath = fileURLToPath(import.meta.url); + const isDev = callerPath.endsWith('.ts'); + const workerFile = isDev ? 'analyze-worker.ts' : 'analyze-worker.js'; + this.workerPath = path.join(path.dirname(callerPath), '..', '..', 'server', workerFile); + this.tsxHookArgs = isDev ? ['--import', pathToFileURL(_require.resolve('tsx/esm')).href] : []; + } + + startJob(params: StartAnalyzeJobParams): AnalyzeJob { + const job = this.jobManager.createJob({ + repoUrl: params.repoUrl, + repoPath: params.repoPath, + }); + + if (job.status !== 'queued') { + return job; + } + + this.jobManager.updateJob(job.id, { status: 'cloning' }); + void this.runJob(job.id, params); + + return job; + } + + getJob(jobId: string): AnalyzeJob | undefined { + return this.jobManager.getJob(jobId); + } + + listJobs(): AnalyzeJob[] { + return this.jobManager.listJobs(); + } + + cancelJob(jobId: string, reason?: string): boolean { + return this.jobManager.cancelJob(jobId, reason); + } + + dispose(): void { + this.jobManager.dispose(); + } + + private async runJob(jobId: string, params: StartAnalyzeJobParams): Promise { + let targetPath = params.repoPath; + let repoLockKey: string | null = null; + + const releaseRepoLock = () => { + if (repoLockKey) { + this.repoLocks.release(repoLockKey); + repoLockKey = null; + } + }; + + try { + if (params.repoUrl && !params.repoPath) { + const repoName = extractRepoName(params.repoUrl); + targetPath = getCloneDir(repoName); + + this.jobManager.updateJob(jobId, { + status: 'cloning', + repoName, + progress: { phase: 'cloning', percent: 0, message: `Cloning ${params.repoUrl}...` }, + }); + + await cloneOrPull(params.repoUrl, targetPath, (progress) => { + this.jobManager.updateJob(jobId, { + progress: { phase: progress.phase, percent: 5, message: progress.message }, + }); + }); + } + + if (!targetPath) { + throw new Error('No target path resolved'); + } + + repoLockKey = getStoragePath(targetPath); + const lockErr = this.repoLocks.acquire(repoLockKey); + if (lockErr) { + this.jobManager.updateJob(jobId, { status: 'failed', error: lockErr }); + return; + } + + this.jobManager.updateJob(jobId, { repoPath: targetPath, status: 'analyzing' }); + this.forkWorker(jobId, targetPath, repoLockKey, { + force: !!params.force, + embeddings: !!params.embeddings, + }); + } catch (err: any) { + releaseRepoLock(); + this.jobManager.updateJob(jobId, { + status: 'failed', + error: err.message || 'Analysis failed', + }); + } + } + + private forkWorker( + jobId: string, + targetPath: string, + repoLockKey: string, + options: AnalyzeOptions, + ) { + const currentJob = this.jobManager.getJob(jobId); + if (!currentJob || this.isTerminal(currentJob.status)) { + return; + } + + const child = fork(this.workerPath, [], { + execArgv: [...this.tsxHookArgs, '--max-old-space-size=8192'], + stdio: ['ignore', 'pipe', 'pipe', 'ipc'], + }); + + let stderrChunks = ''; + child.stderr?.on('data', (chunk: Buffer) => { + stderrChunks += chunk.toString(); + if (stderrChunks.length > 4096) { + stderrChunks = stderrChunks.slice(-4096); + } + }); + + child.on('message', (msg: WorkerMessage) => { + if (msg.type === 'progress') { + this.jobManager.updateJob(jobId, { + status: 'analyzing', + progress: { phase: msg.phase, percent: msg.percent, message: msg.message }, + }); + return; + } + + this.repoLocks.release(repoLockKey); + + if (msg.type === 'complete') { + this.backendInit() + .then(() => { + this.jobManager.updateJob(jobId, { + status: 'complete', + repoName: msg.result.repoName, + }); + }) + .catch((err) => { + this.logger.error('backend.init() failed after analyze:', err); + this.jobManager.updateJob(jobId, { + status: 'failed', + error: 'Server failed to reload after analysis. Try again.', + }); + }); + return; + } + + this.jobManager.updateJob(jobId, { + status: 'failed', + error: msg.message, + }); + }); + + child.on('error', (err) => { + this.repoLocks.release(repoLockKey); + this.jobManager.updateJob(jobId, { + status: 'failed', + error: `Worker process error: ${err.message}`, + }); + }); + + child.on('exit', (code) => { + const job = this.jobManager.getJob(jobId); + if (!job || this.isTerminal(job.status)) { + return; + } + + if (job.retryCount < MAX_WORKER_RETRIES) { + job.retryCount++; + const delay = 1000 * Math.pow(2, job.retryCount - 1); + const lastErr = stderrChunks.trim().split('\n').pop() || ''; + this.logger.warn( + `Analyze worker crashed (code ${code}), retry ${job.retryCount}/${MAX_WORKER_RETRIES} in ${delay}ms` + + (lastErr ? `: ${lastErr}` : ''), + ); + this.jobManager.updateJob(jobId, { + status: 'analyzing', + progress: { + phase: 'retrying', + percent: job.progress.percent, + message: `Worker crashed, retrying (${job.retryCount}/${MAX_WORKER_RETRIES})...`, + }, + }); + stderrChunks = ''; + setTimeout(() => this.forkWorker(jobId, targetPath, repoLockKey, options), delay); + return; + } + + this.repoLocks.release(repoLockKey); + this.jobManager.updateJob(jobId, { + status: 'failed', + error: `Worker crashed ${MAX_WORKER_RETRIES + 1} times (code ${code})${stderrChunks ? ': ' + stderrChunks.trim().split('\n').pop() : ''}`, + }); + }); + + this.jobManager.registerChild(jobId, child); + + const startMessage: StartMessage = { + type: 'start', + repoPath: targetPath, + options, + }; + child.send(startMessage); + } + + private isTerminal(status: AnalyzeJob['status']): boolean { + return status === 'complete' || status === 'failed'; + } +} diff --git a/gitnexus/src/core/index-jobs/embed-job-service.ts b/gitnexus/src/core/index-jobs/embed-job-service.ts new file mode 100644 index 000000000..08dc7636b --- /dev/null +++ b/gitnexus/src/core/index-jobs/embed-job-service.ts @@ -0,0 +1,138 @@ +import path from 'path'; +import { executeQuery, executeWithReusedStatement, withLbugDb } from '../lbug/lbug-adapter.js'; +import { + loadCLIConfig, + saveCLIConfig, + type CLIEmbeddingConfig, +} from '../../storage/repo-manager.js'; +import { JobManager, type AnalyzeJob } from './job-manager.js'; +import { RepoLockManager } from './repo-lock-manager.js'; + +export interface StartEmbedJobParams { + repoName: string; + storagePath: string; + embedding?: CLIEmbeddingConfig; + saveConfig?: boolean; +} + +export interface EmbedJobServiceOptions { + repoLocks: RepoLockManager; + jobManager?: JobManager; +} + +export class EmbedJobService { + readonly jobManager: JobManager; + + private readonly repoLocks: RepoLockManager; + + constructor(options: EmbedJobServiceOptions) { + this.repoLocks = options.repoLocks; + this.jobManager = options.jobManager ?? new JobManager(); + } + + async startJob(params: StartEmbedJobParams): Promise { + if (params.saveConfig && params.embedding) { + const existing = await loadCLIConfig(); + await saveCLIConfig({ + ...existing, + embedding: { + ...(existing.embedding ?? {}), + ...params.embedding, + }, + }); + } + + const job = this.jobManager.createJob({ repoPath: params.storagePath }); + if (job.status !== 'queued') { + return job; + } + + const repoLockPath = params.storagePath; + const lockErr = this.repoLocks.acquire(repoLockPath); + if (lockErr) { + this.jobManager.updateJob(job.id, { status: 'failed', error: lockErr }); + return job; + } + + this.jobManager.updateJob(job.id, { + repoName: params.repoName, + status: 'analyzing', + progress: { phase: 'analyzing', percent: 0, message: 'Starting embedding generation...' }, + }); + + const EMBED_TIMEOUT_MS = 30 * 60 * 1000; + const embedTimeout = setTimeout(() => { + const current = this.jobManager.getJob(job.id); + if (current && current.status !== 'complete' && current.status !== 'failed') { + this.repoLocks.release(repoLockPath); + this.jobManager.updateJob(job.id, { + status: 'failed', + error: 'Embedding timed out (30 minute limit)', + }); + } + }, EMBED_TIMEOUT_MS); + + void (async () => { + try { + const lbugPath = path.join(params.storagePath, 'lbug'); + await withLbugDb(lbugPath, async () => { + const { runEmbeddingPipeline } = await import('../embeddings/embedding-pipeline.js'); + await runEmbeddingPipeline(executeQuery, executeWithReusedStatement, (progress) => { + this.jobManager.updateJob(job.id, { + progress: { + phase: + progress.phase === 'ready' + ? 'complete' + : progress.phase === 'error' + ? 'failed' + : progress.phase, + percent: progress.percent, + message: + progress.phase === 'loading-model' + ? 'Loading embedding model...' + : progress.phase === 'embedding' + ? `Embedding nodes (${progress.percent}%)...` + : progress.phase === 'indexing' + ? 'Creating vector index...' + : progress.phase === 'ready' + ? 'Embeddings complete' + : `${progress.phase} (${progress.percent}%)`, + }, + }); + }); + }); + + clearTimeout(embedTimeout); + this.repoLocks.release(repoLockPath); + const current = this.jobManager.getJob(job.id); + if (!current || current.status !== 'failed') { + this.jobManager.updateJob(job.id, { status: 'complete' }); + } + } catch (err: any) { + clearTimeout(embedTimeout); + this.repoLocks.release(repoLockPath); + const current = this.jobManager.getJob(job.id); + if (!current || current.status !== 'failed') { + this.jobManager.updateJob(job.id, { + status: 'failed', + error: err.message || 'Embedding generation failed', + }); + } + } + })(); + + return job; + } + + getJob(jobId: string): AnalyzeJob | undefined { + return this.jobManager.getJob(jobId); + } + + cancelJob(jobId: string, reason?: string): boolean { + return this.jobManager.cancelJob(jobId, reason); + } + + dispose(): void { + this.jobManager.dispose(); + } +} diff --git a/gitnexus/src/core/index-jobs/index.ts b/gitnexus/src/core/index-jobs/index.ts new file mode 100644 index 000000000..785939209 --- /dev/null +++ b/gitnexus/src/core/index-jobs/index.ts @@ -0,0 +1,5 @@ +export * from './job-manager.js'; +export * from './repo-lock-manager.js'; +export * from './job-query-service.js'; +export * from './analyze-job-service.js'; +export * from './embed-job-service.js'; diff --git a/gitnexus/src/core/index-jobs/job-manager.ts b/gitnexus/src/core/index-jobs/job-manager.ts new file mode 100644 index 000000000..4fa18019b --- /dev/null +++ b/gitnexus/src/core/index-jobs/job-manager.ts @@ -0,0 +1,190 @@ +/** + * Index Job Manager + * + * Tracks server-side jobs with: + * - In-memory Map storage + * - Single-slot concurrency (one active job at a time) + * - Same-repo deduplication (returns existing job) + * - Progress event emission for SSE relay + * - 1-hour TTL cleanup for completed/failed jobs + */ + +import { randomUUID } from 'crypto'; +import { EventEmitter } from 'events'; +import type { ChildProcess } from 'child_process'; + +export interface AnalyzeJobProgress { + phase: string; + percent: number; + message: string; +} + +export interface AnalyzeJob { + id: string; + status: 'queued' | 'cloning' | 'analyzing' | 'loading' | 'complete' | 'failed'; + repoUrl?: string; + repoPath?: string; + repoName?: string; + progress: AnalyzeJobProgress; + error?: string; + startedAt: number; + completedAt?: number; + /** Number of times the worker has been retried after a crash. */ + retryCount: number; +} + +const JOB_TTL_MS = 60 * 60 * 1000; // 1 hour +const CLEANUP_INTERVAL_MS = 5 * 60 * 1000; // 5 minutes +const JOB_TIMEOUT_MS = 30 * 60 * 1000; // 30 minutes + +export class JobManager { + private jobs = new Map(); + private children = new Map(); + private timeouts = new Map>(); + private emitter = new EventEmitter(); + private cleanupTimer: ReturnType; + + constructor() { + this.cleanupTimer = setInterval(() => this.cleanup(), CLEANUP_INTERVAL_MS); + } + + /** Create a new job, or return existing active job for the same repo. */ + createJob(params: { repoUrl?: string; repoPath?: string }): AnalyzeJob { + for (const job of this.jobs.values()) { + if (!this.isTerminal(job.status)) { + const isSameRepo = + (params.repoUrl && job.repoUrl === params.repoUrl) || + (params.repoPath && job.repoPath === params.repoPath); + if (isSameRepo) { + return job; + } + } + } + + for (const job of this.jobs.values()) { + if (!this.isTerminal(job.status)) { + throw new Error(`Analysis already in progress (job ${job.id})`); + } + } + + const job: AnalyzeJob = { + id: randomUUID(), + status: 'queued', + repoUrl: params.repoUrl, + repoPath: params.repoPath, + progress: { phase: 'queued', percent: 0, message: 'Waiting to start...' }, + startedAt: Date.now(), + retryCount: 0, + }; + + this.jobs.set(job.id, job); + return job; + } + + getJob(id: string): AnalyzeJob | undefined { + return this.jobs.get(id); + } + + listJobs(): AnalyzeJob[] { + return Array.from(this.jobs.values()); + } + + updateJob( + id: string, + update: Partial< + Pick + >, + ) { + const job = this.jobs.get(id); + if (!job) return; + + Object.assign(job, update); + + if (this.isTerminal(job.status)) { + job.completedAt = job.completedAt ?? Date.now(); + } + + if (update.status === 'complete' || update.status === 'failed') { + this.emitter.emit(`progress:${id}`, { + phase: update.status, + percent: update.status === 'complete' ? 100 : job.progress.percent, + message: update.status === 'complete' ? 'Complete' : update.error || 'Failed', + }); + } else if (update.progress) { + this.emitter.emit(`progress:${id}`, update.progress); + } + } + + /** Register a child process for a job — enables cancellation and timeout. */ + registerChild(jobId: string, child: ChildProcess) { + this.children.set(jobId, child); + + const timer = setTimeout(() => { + const job = this.jobs.get(jobId); + if (job && !this.isTerminal(job.status)) { + this.cancelJob(jobId, 'Analysis timed out (30 minute limit)'); + } + }, JOB_TIMEOUT_MS); + this.timeouts.set(jobId, timer); + + child.on('exit', () => { + this.children.delete(jobId); + const currentTimer = this.timeouts.get(jobId); + if (currentTimer) { + clearTimeout(currentTimer); + this.timeouts.delete(jobId); + } + }); + } + + cancelJob(jobId: string, reason?: string): boolean { + const job = this.jobs.get(jobId); + if (!job || this.isTerminal(job.status)) return false; + + const child = this.children.get(jobId); + if (child) { + child.kill('SIGTERM'); + } + + this.updateJob(jobId, { + status: 'failed', + error: reason || 'Analysis cancelled', + }); + + return true; + } + + onProgress(jobId: string, listener: (progress: AnalyzeJobProgress) => void): () => void { + const event = `progress:${jobId}`; + this.emitter.on(event, listener); + return () => this.emitter.off(event, listener); + } + + dispose() { + for (const child of this.children.values()) { + child.kill('SIGTERM'); + } + this.children.clear(); + + for (const timer of this.timeouts.values()) { + clearTimeout(timer); + } + this.timeouts.clear(); + + clearInterval(this.cleanupTimer); + this.emitter.removeAllListeners(); + } + + private isTerminal(status: AnalyzeJob['status']): boolean { + return status === 'complete' || status === 'failed'; + } + + private cleanup() { + const now = Date.now(); + for (const [id, job] of this.jobs) { + if (this.isTerminal(job.status) && job.completedAt && now - job.completedAt > JOB_TTL_MS) { + this.jobs.delete(id); + } + } + } +} diff --git a/gitnexus/src/core/index-jobs/job-query-service.ts b/gitnexus/src/core/index-jobs/job-query-service.ts new file mode 100644 index 000000000..605e22c4a --- /dev/null +++ b/gitnexus/src/core/index-jobs/job-query-service.ts @@ -0,0 +1,51 @@ +import { JobManager, type AnalyzeJob } from './job-manager.js'; + +export type IndexJobKind = 'analyze' | 'embed'; + +export interface IndexJobSnapshot extends AnalyzeJob { + kind: IndexJobKind; +} + +/** + * Read-only facade across multiple in-memory job managers. + * + * Keeps HTTP and future MCP write tools on a single query surface without + * forcing them to understand where each job is stored. + */ +export class IndexJobQueryService { + private readonly managers = new Map(); + + register(kind: IndexJobKind, manager: JobManager): this { + this.managers.set(kind, manager); + return this; + } + + getJob(jobId: string): IndexJobSnapshot | undefined { + for (const [kind, manager] of this.managers) { + const job = manager.getJob(jobId); + if (job) { + return { kind, ...job }; + } + } + + return undefined; + } + + listJobs(kind?: IndexJobKind): IndexJobSnapshot[] { + if (kind) { + return this.snapshotJobs(kind, this.managers.get(kind)); + } + + return Array.from(this.managers.entries()) + .flatMap(([currentKind, manager]) => this.snapshotJobs(currentKind, manager)) + .sort((a, b) => b.startedAt - a.startedAt); + } + + private snapshotJobs(kind: IndexJobKind, manager?: JobManager): IndexJobSnapshot[] { + if (!manager) return []; + return manager + .listJobs() + .map((job: AnalyzeJob) => ({ kind, ...job })) + .sort((a, b) => b.startedAt - a.startedAt); + } +} diff --git a/gitnexus/src/core/index-jobs/repo-lock-manager.ts b/gitnexus/src/core/index-jobs/repo-lock-manager.ts new file mode 100644 index 000000000..fdc897c15 --- /dev/null +++ b/gitnexus/src/core/index-jobs/repo-lock-manager.ts @@ -0,0 +1,26 @@ +/** + * Shared repository lock for background index jobs. + * + * The lock key should be stable across analyze/embed callers. In practice we + * use the storage path so both tasks serialize access to the same LadybugDB. + */ +export class RepoLockManager { + private activeRepoPaths = new Set(); + + acquire(repoPath: string): string | null { + if (this.activeRepoPaths.has(repoPath)) { + return 'Another job is already active for this repository'; + } + + this.activeRepoPaths.add(repoPath); + return null; + } + + release(repoPath: string): void { + this.activeRepoPaths.delete(repoPath); + } + + isLocked(repoPath: string): boolean { + return this.activeRepoPaths.has(repoPath); + } +} diff --git a/gitnexus/src/core/lbug/schema.ts b/gitnexus/src/core/lbug/schema.ts index 257938a01..c95c194cd 100644 --- a/gitnexus/src/core/lbug/schema.ts +++ b/gitnexus/src/core/lbug/schema.ts @@ -11,6 +11,7 @@ // Import from shared package (single source of truth) — used in DDL templates below import { NODE_TABLES, REL_TABLE_NAME, REL_TYPES, EMBEDDING_TABLE_NAME } from 'gitnexus-shared'; +import { resolveEmbeddingConfigSync } from '../embeddings/config.js'; // Re-export so downstream consumers keep the same import path export { NODE_TABLES, REL_TABLE_NAME, REL_TYPES, EMBEDDING_TABLE_NAME }; export type { NodeTableName, RelType } from 'gitnexus-shared'; @@ -428,13 +429,7 @@ CREATE REL TABLE ${REL_TABLE_NAME} ( // ============================================================================ /** Embedding vector dimensions. Default 384 (snowflake-arctic-embed-xs). */ -const _rawDims = parseInt(process.env.GITNEXUS_EMBEDDING_DIMS ?? '384', 10); -if (Number.isNaN(_rawDims) || _rawDims <= 0) { - throw new Error( - `GITNEXUS_EMBEDDING_DIMS must be a positive integer, got "${process.env.GITNEXUS_EMBEDDING_DIMS}"`, - ); -} -export const EMBEDDING_DIMS = _rawDims; +export const EMBEDDING_DIMS = resolveEmbeddingConfigSync().dimensions; export const EMBEDDING_SCHEMA = ` CREATE NODE TABLE ${EMBEDDING_TABLE_NAME} ( diff --git a/gitnexus/src/mcp/core/embedder.ts b/gitnexus/src/mcp/core/embedder.ts index e53a46542..97f1bf893 100644 --- a/gitnexus/src/mcp/core/embedder.ts +++ b/gitnexus/src/mcp/core/embedder.ts @@ -6,6 +6,7 @@ */ import { pipeline, env, type FeatureExtractionPipeline } from '@huggingface/transformers'; +import { resolveEmbeddingConfigSync } from '../../core/embeddings/config.js'; import { isHttpMode, getHttpDimensions, @@ -116,7 +117,7 @@ export const embedQuery = async (query: string): Promise => { * Get embedding dimensions */ export const getEmbeddingDims = (): number => { - return getHttpDimensions() ?? 384; + return isHttpMode() ? (getHttpDimensions() ?? 384) : resolveEmbeddingConfigSync().dimensions; }; /** diff --git a/gitnexus/src/mcp/local/local-backend.ts b/gitnexus/src/mcp/local/local-backend.ts index 041bd27eb..8a62f4c38 100644 --- a/gitnexus/src/mcp/local/local-backend.ts +++ b/gitnexus/src/mcp/local/local-backend.ts @@ -25,9 +25,17 @@ import { parseDiffHunks, type FileDiff } from '../../storage/git.js'; import { listRegisteredRepos, cleanupOldKuzuFiles, + type CLIEmbeddingConfig, type RegistryEntry, } from '../../storage/repo-manager.js'; import { GroupService, type GroupToolPort } from '../../core/group/service.js'; +import { + AnalyzeJobService, + EmbedJobService, + IndexJobQueryService, + RepoLockManager, + type IndexJobSnapshot, +} from '../../core/index-jobs/index.js'; // AI context generation is CLI-only (gitnexus analyze) // import { generateAIContextFiles } from '../../cli/ai-context.js'; @@ -182,6 +190,25 @@ export class LocalBackend { private reinitPromises: Map> = new Map(); private lastStalenessCheck: Map = new Map(); private groupToolSvc: GroupService | null = null; + private readonly repoLocks: RepoLockManager; + private readonly analyzeJobs: AnalyzeJobService; + private readonly embedJobs: EmbedJobService; + private readonly jobQueries: IndexJobQueryService; + + constructor() { + this.repoLocks = new RepoLockManager(); + this.analyzeJobs = new AnalyzeJobService({ + backendInit: async () => { + await this.refreshRepos(); + }, + repoLocks: this.repoLocks, + logger: console, + }); + this.embedJobs = new EmbedJobService({ repoLocks: this.repoLocks }); + this.jobQueries = new IndexJobQueryService() + .register('analyze', this.analyzeJobs.jobManager) + .register('embed', this.embedJobs.jobManager); + } /** * Cross-repo group tools (CLI). Shares logic with MCP `group_*` handlers. @@ -201,6 +228,8 @@ export class LocalBackend { /** Close all pooled LadybugDB connections (CLI one-shot; optional for long-lived MCP). */ async dispose(): Promise { + this.analyzeJobs.dispose(); + this.embedJobs.dispose(); await closeLbug(); } @@ -458,6 +487,10 @@ export class LocalBackend { return this.listRepos(); } + if (this.isWriteTool(method)) { + return this.handleWriteTool(method, params || {}); + } + if (method.startsWith('group_')) { return this.handleGroupTool(method, params || {}); } @@ -500,6 +533,125 @@ export class LocalBackend { } } + private isWriteTool(method: string): boolean { + return ( + method === 'refresh_repos' || + method === 'get_index_job' || + method === 'analyze_repo' || + method === 'rebuild_embeddings' + ); + } + + private serializeIndexJob(job: IndexJobSnapshot) { + return { + kind: job.kind, + id: job.id, + status: job.status, + repoUrl: job.repoUrl, + repoPath: job.repoPath, + repoName: job.repoName, + progress: job.progress, + error: job.error, + startedAt: job.startedAt, + completedAt: job.completedAt, + }; + } + + private async handleWriteTool(method: string, params: any): Promise { + switch (method) { + case 'refresh_repos': { + await this.refreshRepos(); + return { + refreshed: true, + repoCount: this.repos.size, + repos: [...this.repos.values()].map((repo) => ({ + name: repo.name, + path: repo.repoPath, + indexedAt: repo.indexedAt, + })), + }; + } + case 'get_index_job': { + const jobId = typeof params?.jobId === 'string' ? params.jobId.trim() : ''; + if (!jobId) { + throw new Error('"jobId" is required'); + } + const job = this.jobQueries.getJob(jobId); + if (!job) { + throw new Error(`Job "${jobId}" not found`); + } + return this.serializeIndexJob(job); + } + case 'analyze_repo': { + const repoUrl = typeof params?.url === 'string' ? params.url.trim() : undefined; + const repoPath = typeof params?.path === 'string' ? params.path.trim() : undefined; + if (!repoUrl && !repoPath) { + throw new Error('Provide "url" (git URL) or "path" (local path)'); + } + if (repoPath) { + if (!path.isAbsolute(repoPath)) { + throw new Error('"path" must be an absolute path'); + } + if (path.normalize(repoPath) !== path.resolve(repoPath)) { + throw new Error('"path" must not contain traversal sequences'); + } + } + + const job = this.analyzeJobs.startJob({ + repoUrl, + repoPath, + force: !!params?.force, + embeddings: !!params?.embeddings, + }); + return { + kind: 'analyze', + jobId: job.id, + status: job.status, + repoUrl, + repoPath, + }; + } + case 'rebuild_embeddings': { + const repo = await this.resolveRepo(params?.repo); + const embedding: CLIEmbeddingConfig = {}; + if (typeof params?.provider === 'string' && params.provider.trim()) { + embedding.provider = params.provider.trim() as CLIEmbeddingConfig['provider']; + } + if (typeof params?.baseUrl === 'string' && params.baseUrl.trim()) { + embedding.baseUrl = params.baseUrl.trim(); + } + if (typeof params?.model === 'string' && params.model.trim()) { + embedding.model = params.model.trim(); + } + if (typeof params?.apiKey === 'string' && params.apiKey.trim()) { + embedding.apiKey = params.apiKey.trim(); + } + if (typeof params?.dimensions === 'number') { + embedding.dimensions = params.dimensions; + } + + const hasOverrides = Object.keys(embedding).length > 0; + const saveConfig = hasOverrides ? params?.saveConfig !== false : false; + + const job = await this.embedJobs.startJob({ + repoName: repo.name, + storagePath: repo.storagePath, + embedding: hasOverrides ? embedding : undefined, + saveConfig, + }); + return { + kind: 'embed', + jobId: job.id, + status: job.status, + repo: repo.name, + savedConfig: saveConfig, + }; + } + default: + throw new Error(`Unknown tool: ${method}`); + } + } + // ─── Tool Implementations ──────────────────────────────────────── /** diff --git a/gitnexus/src/mcp/tools.ts b/gitnexus/src/mcp/tools.ts index 2c884d10b..541408c2d 100644 --- a/gitnexus/src/mcp/tools.ts +++ b/gitnexus/src/mcp/tools.ts @@ -377,6 +377,91 @@ Returns: single route object when one match, or { routes: [...], total: N } for required: [], }, }, + { + name: 'refresh_repos', + description: `Force GitNexus to re-read the global registry and refresh its in-memory repository cache. + +WHEN TO USE: After indexing a new repo, removing a repo, or when you want to confirm the MCP server sees the latest registry state without restarting. + +Returns: refreshed flag, repo count, and the current repo list snapshot.`, + inputSchema: { + type: 'object', + properties: {}, + required: [], + }, + }, + { + name: 'get_index_job', + description: `Fetch the current status of a background index job created by analyze_repo or rebuild_embeddings. + +WHEN TO USE: Poll long-running analyze/embed jobs started from MCP. + +Returns: job kind, status, repo metadata, progress, timestamps, and any error message.`, + inputSchema: { + type: 'object', + properties: { + jobId: { + type: 'string', + description: 'Background job identifier returned by a write tool', + }, + }, + required: ['jobId'], + }, + }, + { + name: 'analyze_repo', + description: `Start a background repository analysis job from MCP. + +WHEN TO USE: Initialize or rebuild a project index without leaving the MCP workflow. + +Provide either a local absolute path or a git URL. Returns a job ID for polling via get_index_job.`, + inputSchema: { + type: 'object', + properties: { + path: { type: 'string', description: 'Absolute local path to analyze' }, + url: { type: 'string', description: 'Git URL to clone/analyze' }, + force: { type: 'boolean', description: 'Force a full rebuild', default: false }, + embeddings: { + type: 'boolean', + description: 'Generate embeddings as part of analysis', + default: false, + }, + }, + required: [], + }, + }, + { + name: 'rebuild_embeddings', + description: `Start a background embedding rebuild job for an indexed repository. + +WHEN TO USE: Enable or regenerate semantic search vectors after analysis, or switch to a configured remote embedding provider without using shell environment variables. + +Optional embedding fields can be persisted to ~/.gitnexus/config.json before the rebuild. Returns a job ID for polling via get_index_job.`, + inputSchema: { + type: 'object', + properties: { + repo: { + type: 'string', + description: 'Repository name or path. Omit if only one repo is indexed.', + }, + provider: { + type: 'string', + description: 'Embedding provider (local, openai, openrouter, azure, custom)', + enum: ['local', 'openai', 'openrouter', 'azure', 'custom'], + }, + baseUrl: { type: 'string', description: 'OpenAI-compatible base URL for /v1 embeddings' }, + model: { type: 'string', description: 'Embedding model name' }, + apiKey: { type: 'string', description: 'Embedding API key / bearer token' }, + dimensions: { type: 'number', description: 'Expected embedding dimensions' }, + saveConfig: { + type: 'boolean', + description: + 'Persist provided embedding settings into ~/.gitnexus/config.json before rebuilding (default: true when overrides are provided)', + }, + }, + required: [], + }, + }, { name: 'group_list', description: `List all configured repository groups, or return details for one group (repos, manifest links). diff --git a/gitnexus/src/server/analyze-job.ts b/gitnexus/src/server/analyze-job.ts index f7d97e97d..482dec6c4 100644 --- a/gitnexus/src/server/analyze-job.ts +++ b/gitnexus/src/server/analyze-job.ts @@ -1,201 +1,11 @@ /** - * Analyze Job Manager + * Backward-compatible server export for index job primitives. * - * Tracks server-side analysis jobs with: - * - In-memory Map storage - * - Single-slot concurrency (one active job at a time) - * - Same-repo deduplication (returns existing job) - * - Progress event emission for SSE relay - * - 1-hour TTL cleanup for completed/failed jobs + * New code should import from ../core/index-jobs. */ -import { randomUUID } from 'crypto'; -import { EventEmitter } from 'events'; -import type { ChildProcess } from 'child_process'; - -export interface AnalyzeJobProgress { - phase: string; - percent: number; - message: string; -} - -export interface AnalyzeJob { - id: string; - status: 'queued' | 'cloning' | 'analyzing' | 'loading' | 'complete' | 'failed'; - repoUrl?: string; - repoPath?: string; - repoName?: string; - progress: AnalyzeJobProgress; - error?: string; - startedAt: number; - completedAt?: number; - /** Number of times the worker has been retried after a crash. */ - retryCount: number; -} - -const JOB_TTL_MS = 60 * 60 * 1000; // 1 hour -const CLEANUP_INTERVAL_MS = 5 * 60 * 1000; // 5 minutes -const JOB_TIMEOUT_MS = 30 * 60 * 1000; // 30 minutes - -export class JobManager { - private jobs = new Map(); - private children = new Map(); - private timeouts = new Map>(); - private emitter = new EventEmitter(); - private cleanupTimer: ReturnType; - - constructor() { - this.cleanupTimer = setInterval(() => this.cleanup(), CLEANUP_INTERVAL_MS); - } - - /** Create a new job, or return existing active job for the same repo. */ - createJob(params: { repoUrl?: string; repoPath?: string }): AnalyzeJob { - // Dedup: return existing active job for the same repo (by URL or path) - for (const job of this.jobs.values()) { - if (!this.isTerminal(job.status)) { - const isSameRepo = - (params.repoUrl && job.repoUrl === params.repoUrl) || - (params.repoPath && job.repoPath === params.repoPath); - if (isSameRepo) { - return job; - } - } - } - - // Single-slot: reject if another job is active (different repo) - for (const job of this.jobs.values()) { - if (!this.isTerminal(job.status)) { - throw new Error(`Analysis already in progress (job ${job.id})`); - } - } - - const job: AnalyzeJob = { - id: randomUUID(), - status: 'queued', - repoUrl: params.repoUrl, - repoPath: params.repoPath, - progress: { phase: 'queued', percent: 0, message: 'Waiting to start...' }, - startedAt: Date.now(), - retryCount: 0, - }; - - this.jobs.set(job.id, job); - return job; - } - - getJob(id: string): AnalyzeJob | undefined { - return this.jobs.get(id); - } - - /** Return a snapshot of all tracked jobs for inspection. */ - listJobs(): AnalyzeJob[] { - return Array.from(this.jobs.values()); - } - - updateJob( - id: string, - update: Partial< - Pick - >, - ) { - const job = this.jobs.get(id); - if (!job) return; - - Object.assign(job, update); - - if (this.isTerminal(job.status)) { - job.completedAt = job.completedAt ?? Date.now(); - } - - // Emit exactly one event per updateJob call to prevent SSE double-write - if (update.status === 'complete' || update.status === 'failed') { - // Terminal event takes precedence — don't also emit the progress event - this.emitter.emit(`progress:${id}`, { - phase: update.status, - percent: update.status === 'complete' ? 100 : job.progress.percent, - message: update.status === 'complete' ? 'Complete' : update.error || 'Failed', - }); - } else if (update.progress) { - this.emitter.emit(`progress:${id}`, update.progress); - } - } - - /** Register a child process for a job — enables cancellation and timeout. */ - registerChild(jobId: string, child: ChildProcess) { - this.children.set(jobId, child); - - // 30-minute timeout - const timer = setTimeout(() => { - const job = this.jobs.get(jobId); - if (job && !this.isTerminal(job.status)) { - this.cancelJob(jobId, 'Analysis timed out (30 minute limit)'); - } - }, JOB_TIMEOUT_MS); - this.timeouts.set(jobId, timer); - - // Clean up tracking when child exits - child.on('exit', () => { - this.children.delete(jobId); - const t = this.timeouts.get(jobId); - if (t) { - clearTimeout(t); - this.timeouts.delete(jobId); - } - }); - } - - /** Cancel a running job — sends SIGTERM to child process. */ - cancelJob(jobId: string, reason?: string): boolean { - const job = this.jobs.get(jobId); - if (!job || this.isTerminal(job.status)) return false; - - const child = this.children.get(jobId); - if (child) { - child.kill('SIGTERM'); - } - - this.updateJob(jobId, { - status: 'failed', - error: reason || 'Analysis cancelled', - }); - - return true; - } - - /** Subscribe to progress events for a job. Returns unsubscribe function. */ - onProgress(jobId: string, listener: (progress: AnalyzeJobProgress) => void): () => void { - const event = `progress:${jobId}`; - this.emitter.on(event, listener); - return () => this.emitter.off(event, listener); - } - - dispose() { - // Kill all active child processes - for (const child of this.children.values()) { - child.kill('SIGTERM'); - } - this.children.clear(); - - // Clear all timeouts - for (const timer of this.timeouts.values()) { - clearTimeout(timer); - } - this.timeouts.clear(); - - clearInterval(this.cleanupTimer); - this.emitter.removeAllListeners(); - } - - private isTerminal(status: AnalyzeJob['status']): boolean { - return status === 'complete' || status === 'failed'; - } - - private cleanup() { - const now = Date.now(); - for (const [id, job] of this.jobs) { - if (this.isTerminal(job.status) && job.completedAt && now - job.completedAt > JOB_TTL_MS) { - this.jobs.delete(id); - } - } - } -} +export { + JobManager, + type AnalyzeJob, + type AnalyzeJobProgress, +} from '../core/index-jobs/job-manager.js'; diff --git a/gitnexus/src/server/api.ts b/gitnexus/src/server/api.ts index 3d4cf9a6a..e98785f08 100644 --- a/gitnexus/src/server/api.ts +++ b/gitnexus/src/server/api.ts @@ -14,10 +14,19 @@ import path from 'path'; import fs from 'fs/promises'; import { createRequire } from 'node:module'; import { loadMeta, listRegisteredRepos, getStoragePath } from '../storage/repo-manager.js'; +import { + AnalyzeJobService, + EmbedJobService, + IndexJobQueryService, + RepoLockManager, + type AnalyzeJob, + type IndexJobKind, + type IndexJobSnapshot, + type JobManager, +} from '../core/index-jobs/index.js'; import { executeQuery, executePrepared, - executeWithReusedStatement, streamQuery, closeLbug, withLbugDb, @@ -30,10 +39,7 @@ import { hybridSearch } from '../core/search/hybrid-search.js'; // at server startup — crashes on unsupported Node ABI versions (#89) import { LocalBackend } from '../mcp/local/local-backend.js'; import { mountMCPEndpoints } from './mcp-http.js'; -import { fork } from 'child_process'; -import { fileURLToPath, pathToFileURL } from 'url'; -import { JobManager } from './analyze-job.js'; -import { extractRepoName, getCloneDir, cloneOrPull } from './git-clone.js'; +import { getCloneDir } from './git-clone.js'; const _require = createRequire(import.meta.url); const pkg = _require('../../package.json'); @@ -424,6 +430,26 @@ const requestedRepo = (req: express.Request): string | undefined => { return undefined; }; +const serializeJob = (job: AnalyzeJob) => ({ + id: job.id, + status: job.status, + repoUrl: job.repoUrl, + repoPath: job.repoPath, + repoName: job.repoName, + progress: job.progress, + error: job.error, + startedAt: job.startedAt, + completedAt: job.completedAt, +}); + +const serializeIndexJob = (job: IndexJobSnapshot) => ({ + kind: job.kind, + ...serializeJob(job), +}); + +const isIndexJobKind = (value: string): value is IndexJobKind => + value === 'analyze' || value === 'embed'; + export const createServer = async (port: number, host: string = '127.0.0.1') => { const app = express(); app.disable('x-powered-by'); @@ -462,23 +488,18 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => const backend = new LocalBackend(); await backend.init(); const cleanupMcp = mountMCPEndpoints(app, backend); - const jobManager = new JobManager(); - - // Shared repo lock — prevents concurrent analyze + embed on the same repo path, - // which would corrupt LadybugDB (analyze calls closeLbug + initLbug while embed has queries in flight). - const activeRepoPaths = new Set(); - - const acquireRepoLock = (repoPath: string): string | null => { - if (activeRepoPaths.has(repoPath)) { - return `Another job is already active for this repository`; - } - activeRepoPaths.add(repoPath); - return null; - }; - - const releaseRepoLock = (repoPath: string): void => { - activeRepoPaths.delete(repoPath); - }; + const repoLocks = new RepoLockManager(); + const analyzeJobs = new AnalyzeJobService({ + backendInit: async () => { + await backend.init(); + }, + repoLocks, + logger: console, + }); + const embedJobs = new EmbedJobService({ repoLocks }); + const jobQueries = new IndexJobQueryService() + .register('analyze', analyzeJobs.jobManager) + .register('embed', embedJobs.jobManager); /** * Maximum time the hold-queue will wait for an active analysis job to complete. @@ -518,7 +539,7 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => clientGone = true; }); - for (const job of jobManager.listJobs()) { + for (const job of jobQueries.listJobs('analyze')) { const isMatch = job.repoName?.toLowerCase() === lower || (job.repoUrl && path.basename(job.repoUrl).replace('.git', '').toLowerCase() === lower) || @@ -532,7 +553,7 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => } for (let wait = 0; wait < HOLD_QUEUE_TIMEOUT_SECS; wait++) { if (clientGone) return null; // client disconnected — stop polling - const currentJob = jobManager.getJob(job.id); + const currentJob = analyzeJobs.getJob(job.id); if (!currentJob || currentJob.status === 'failed') break; if (currentJob.status === 'complete') { await backend.init(); @@ -660,7 +681,7 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => // Acquire repo lock — prevents deleting while analyze/embed is in flight const lockKey = getStoragePath(entry.path); - const lockErr = acquireRepoLock(lockKey); + const lockErr = repoLocks.acquire(lockKey); if (lockErr) { res.status(409).json({ error: lockErr }); return; @@ -696,7 +717,7 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => res.json({ deleted: entry.name }); } finally { - releaseRepoLock(lockKey); + repoLocks.release(lockKey); } } catch (err: any) { res.status(500).json({ error: err.message || 'Failed to delete repo' }); @@ -1174,185 +1195,12 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => } } - const job = jobManager.createJob({ repoUrl, repoPath: repoLocalPath }); - - // If job was already running (dedup), just return its id - if (job.status !== 'queued') { - res.status(202).json({ jobId: job.id, status: job.status }); - return; - } - - // Mark as active synchronously to prevent race with concurrent requests - jobManager.updateJob(job.id, { status: 'cloning' }); - - // Start async work — don't await - (async () => { - let targetPath = repoLocalPath; - try { - // Clone if URL provided - if (repoUrl && !repoLocalPath) { - const repoName = extractRepoName(repoUrl); - targetPath = getCloneDir(repoName); - - jobManager.updateJob(job.id, { - status: 'cloning', - repoName, - progress: { phase: 'cloning', percent: 0, message: `Cloning ${repoUrl}...` }, - }); - - await cloneOrPull(repoUrl, targetPath, (progress) => { - jobManager.updateJob(job.id, { - progress: { phase: progress.phase, percent: 5, message: progress.message }, - }); - }); - } - - if (!targetPath) { - throw new Error('No target path resolved'); - } - - // Acquire shared repo lock (keyed on storagePath to match embed handler) - const analyzeLockKey = getStoragePath(targetPath); - const lockErr = acquireRepoLock(analyzeLockKey); - if (lockErr) { - jobManager.updateJob(job.id, { status: 'failed', error: lockErr }); - return; - } - - jobManager.updateJob(job.id, { repoPath: targetPath, status: 'analyzing' }); - - // ── Worker fork with auto-retry ────────────────────────────── - // - // Forks a child process with 8GB heap. If the worker crashes - // (OOM, native addon segfault, etc.), it retries up to - // MAX_WORKER_RETRIES times with exponential backoff before - // marking the job as permanently failed. - // - // In dev mode (tsx), registers the tsx ESM hook via a file:// - // URL so the child can compile TypeScript on-the-fly. - - const MAX_WORKER_RETRIES = 2; - const callerPath = fileURLToPath(import.meta.url); - const isDev = callerPath.endsWith('.ts'); - const workerFile = isDev ? 'analyze-worker.ts' : 'analyze-worker.js'; - const workerPath = path.join(path.dirname(callerPath), workerFile); - const tsxHookArgs: string[] = isDev - ? ['--import', pathToFileURL(_require.resolve('tsx/esm')).href] - : []; - - const forkWorker = () => { - const currentJob = jobManager.getJob(job.id); - if (!currentJob || currentJob.status === 'complete' || currentJob.status === 'failed') - return; - - const child = fork(workerPath, [], { - execArgv: [...tsxHookArgs, '--max-old-space-size=8192'], - stdio: ['ignore', 'pipe', 'pipe', 'ipc'], - }); - - // Capture stderr for crash diagnostics - let stderrChunks = ''; - child.stderr?.on('data', (chunk: Buffer) => { - stderrChunks += chunk.toString(); - if (stderrChunks.length > 4096) stderrChunks = stderrChunks.slice(-4096); - }); - - child.on('message', (msg: any) => { - if (msg.type === 'progress') { - jobManager.updateJob(job.id, { - status: 'analyzing', - progress: { phase: msg.phase, percent: msg.percent, message: msg.message }, - }); - } else if (msg.type === 'complete') { - releaseRepoLock(analyzeLockKey); - // Reinitialize backend BEFORE marking complete — ensures the new - // repo is queryable when the client receives the SSE complete event. - backend - .init() - .then(() => { - jobManager.updateJob(job.id, { - status: 'complete', - repoName: msg.result.repoName, - }); - }) - .catch((err) => { - console.error('backend.init() failed after analyze:', err); - jobManager.updateJob(job.id, { - status: 'failed', - error: 'Server failed to reload after analysis. Try again.', - }); - }); - } else if (msg.type === 'error') { - releaseRepoLock(analyzeLockKey); - jobManager.updateJob(job.id, { - status: 'failed', - error: msg.message, - }); - } - }); - - child.on('error', (err) => { - releaseRepoLock(analyzeLockKey); - jobManager.updateJob(job.id, { - status: 'failed', - error: `Worker process error: ${err.message}`, - }); - }); - - child.on('exit', (code) => { - const j = jobManager.getJob(job.id); - if (!j || j.status === 'complete' || j.status === 'failed') return; - - // Worker crashed — attempt retry if under the limit - if (j.retryCount < MAX_WORKER_RETRIES) { - j.retryCount++; - const delay = 1000 * Math.pow(2, j.retryCount - 1); // 1s, 2s - const lastErr = stderrChunks.trim().split('\n').pop() || ''; - console.warn( - `Analyze worker crashed (code ${code}), retry ${j.retryCount}/${MAX_WORKER_RETRIES} in ${delay}ms` + - (lastErr ? `: ${lastErr}` : ''), - ); - jobManager.updateJob(job.id, { - status: 'analyzing', - progress: { - phase: 'retrying', - percent: j.progress.percent, - message: `Worker crashed, retrying (${j.retryCount}/${MAX_WORKER_RETRIES})...`, - }, - }); - stderrChunks = ''; - setTimeout(forkWorker, delay); - } else { - // Exhausted retries — permanent failure - releaseRepoLock(analyzeLockKey); - jobManager.updateJob(job.id, { - status: 'failed', - error: `Worker crashed ${MAX_WORKER_RETRIES + 1} times (code ${code})${stderrChunks ? ': ' + stderrChunks.trim().split('\n').pop() : ''}`, - }); - } - }); - - // Register child for cancellation + timeout tracking - jobManager.registerChild(job.id, child); - - // Send start command to child - child.send({ - type: 'start', - repoPath: targetPath, - options: { force: !!force, embeddings: !!embeddings }, - }); - }; - - forkWorker(); - } catch (err: any) { - if (targetPath) releaseRepoLock(getStoragePath(targetPath)); - jobManager.updateJob(job.id, { - status: 'failed', - error: err.message || 'Analysis failed', - }); - } - })(); - + const job = analyzeJobs.startJob({ + repoUrl, + repoPath: repoLocalPath, + force: !!force, + embeddings: !!embeddings, + }); res.status(202).json({ jobId: job.id, status: job.status }); } catch (err: any) { if (err.message?.includes('already in progress')) { @@ -1365,30 +1213,20 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => // GET /api/analyze/:jobId — poll job status app.get('/api/analyze/:jobId', (req, res) => { - const job = jobManager.getJob(req.params.jobId); + const job = analyzeJobs.getJob(req.params.jobId); if (!job) { res.status(404).json({ error: 'Job not found' }); return; } - res.json({ - id: job.id, - status: job.status, - repoUrl: job.repoUrl, - repoPath: job.repoPath, - repoName: job.repoName, - progress: job.progress, - error: job.error, - startedAt: job.startedAt, - completedAt: job.completedAt, - }); + res.json(serializeJob(job)); }); // GET /api/analyze/:jobId/progress — SSE stream (shared helper) - mountSSEProgress(app, '/api/analyze/:jobId/progress', jobManager); + mountSSEProgress(app, '/api/analyze/:jobId/progress', analyzeJobs.jobManager); // DELETE /api/analyze/:jobId — cancel a running analysis job app.delete('/api/analyze/:jobId', (req, res) => { - const job = jobManager.getJob(req.params.jobId); + const job = analyzeJobs.getJob(req.params.jobId); if (!job) { res.status(404).json({ error: 'Job not found' }); return; @@ -1397,13 +1235,47 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => res.status(400).json({ error: `Job already ${job.status}` }); return; } - jobManager.cancelJob(req.params.jobId, 'Cancelled by user'); + analyzeJobs.cancelJob(req.params.jobId, 'Cancelled by user'); res.json({ id: job.id, status: 'failed', error: 'Cancelled by user' }); }); - // ── Embedding endpoints ──────────────────────────────────────────── + // GET /api/jobs — lightweight job discovery across analyze/embed managers + app.get('/api/jobs', (req, res) => { + const kindParam = typeof req.query.kind === 'string' ? req.query.kind : undefined; + if (kindParam && !isIndexJobKind(kindParam)) { + res.status(400).json({ error: 'Invalid "kind" query parameter' }); + return; + } + const kind = kindParam as IndexJobKind | undefined; - const embedJobManager = new JobManager(); + const repo = typeof req.query.repo === 'string' ? req.query.repo.trim().toLowerCase() : ''; + const jobs = jobQueries + .listJobs(kind) + .filter((job) => { + if (!repo) return true; + return ( + job.repoName?.toLowerCase().includes(repo) || + job.repoPath?.toLowerCase().includes(repo) || + job.repoUrl?.toLowerCase().includes(repo) + ); + }) + .map(serializeIndexJob); + + res.json({ jobs }); + }); + + // GET /api/jobs/:jobId — fetch a job without knowing its manager + app.get('/api/jobs/:jobId', (req, res) => { + const job = jobQueries.getJob(req.params.jobId); + if (!job) { + res.status(404).json({ error: 'Job not found' }); + return; + } + + res.json(serializeIndexJob(job)); + }); + + // ── Embedding endpoints ──────────────────────────────────────────── // POST /api/embed — trigger server-side embedding generation app.post('/api/embed', async (req, res) => { @@ -1414,83 +1286,11 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => return; } - // Check shared repo lock — prevent concurrent analyze + embed on same repo - const repoLockPath = entry.storagePath; - const lockErr = acquireRepoLock(repoLockPath); - if (lockErr) { - res.status(409).json({ error: lockErr }); - return; - } - - const job = embedJobManager.createJob({ repoPath: entry.storagePath }); - embedJobManager.updateJob(job.id, { + const job = await embedJobs.startJob({ repoName: entry.name, - status: 'analyzing' as any, - progress: { phase: 'analyzing', percent: 0, message: 'Starting embedding generation...' }, + storagePath: entry.storagePath, }); - - // 30-minute timeout for embedding jobs (same as analyze jobs) - const EMBED_TIMEOUT_MS = 30 * 60 * 1000; - const embedTimeout = setTimeout(() => { - const current = embedJobManager.getJob(job.id); - if (current && current.status !== 'complete' && current.status !== 'failed') { - releaseRepoLock(repoLockPath); - embedJobManager.updateJob(job.id, { - status: 'failed', - error: 'Embedding timed out (30 minute limit)', - }); - } - }, EMBED_TIMEOUT_MS); - - // Run embedding pipeline asynchronously - (async () => { - try { - const lbugPath = path.join(entry.storagePath, 'lbug'); - await withLbugDb(lbugPath, async () => { - const { runEmbeddingPipeline } = - await import('../core/embeddings/embedding-pipeline.js'); - await runEmbeddingPipeline(executeQuery, executeWithReusedStatement, (p) => { - embedJobManager.updateJob(job.id, { - progress: { - phase: - p.phase === 'ready' ? 'complete' : p.phase === 'error' ? 'failed' : p.phase, - percent: p.percent, - message: - p.phase === 'loading-model' - ? 'Loading embedding model...' - : p.phase === 'embedding' - ? `Embedding nodes (${p.percent}%)...` - : p.phase === 'indexing' - ? 'Creating vector index...' - : p.phase === 'ready' - ? 'Embeddings complete' - : `${p.phase} (${p.percent}%)`, - }, - }); - }); - }); - - clearTimeout(embedTimeout); - releaseRepoLock(repoLockPath); - // Don't overwrite 'failed' if the job was cancelled while the pipeline was running - const current = embedJobManager.getJob(job.id); - if (!current || current.status !== 'failed') { - embedJobManager.updateJob(job.id, { status: 'complete' }); - } - } catch (err: any) { - clearTimeout(embedTimeout); - releaseRepoLock(repoLockPath); - const current = embedJobManager.getJob(job.id); - if (!current || current.status !== 'failed') { - embedJobManager.updateJob(job.id, { - status: 'failed', - error: err.message || 'Embedding generation failed', - }); - } - } - })(); - - res.status(202).json({ jobId: job.id, status: 'analyzing' }); + res.status(202).json({ jobId: job.id, status: job.status }); } catch (err: any) { if (err.message?.includes('already in progress')) { res.status(409).json({ error: err.message }); @@ -1502,28 +1302,20 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => // GET /api/embed/:jobId — poll embedding job status app.get('/api/embed/:jobId', (req, res) => { - const job = embedJobManager.getJob(req.params.jobId); + const job = embedJobs.getJob(req.params.jobId); if (!job) { res.status(404).json({ error: 'Job not found' }); return; } - res.json({ - id: job.id, - status: job.status, - repoName: job.repoName, - progress: job.progress, - error: job.error, - startedAt: job.startedAt, - completedAt: job.completedAt, - }); + res.json(serializeJob(job)); }); // GET /api/embed/:jobId/progress — SSE stream (shared helper) - mountSSEProgress(app, '/api/embed/:jobId/progress', embedJobManager); + mountSSEProgress(app, '/api/embed/:jobId/progress', embedJobs.jobManager); // DELETE /api/embed/:jobId — cancel embedding job app.delete('/api/embed/:jobId', (req, res) => { - const job = embedJobManager.getJob(req.params.jobId); + const job = embedJobs.getJob(req.params.jobId); if (!job) { res.status(404).json({ error: 'Job not found' }); return; @@ -1532,7 +1324,7 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => res.status(400).json({ error: `Job already ${job.status}` }); return; } - embedJobManager.cancelJob(req.params.jobId, 'Cancelled by user'); + embedJobs.cancelJob(req.params.jobId, 'Cancelled by user'); res.json({ id: job.id, status: 'failed', error: 'Cancelled by user' }); }); @@ -1556,8 +1348,8 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => const shutdown = async () => { console.log('\nShutting down...'); server.close(); - jobManager.dispose(); - embedJobManager.dispose(); + analyzeJobs.dispose(); + embedJobs.dispose(); await cleanupMcp(); await closeLbug(); await backend.disconnect(); diff --git a/gitnexus/src/storage/repo-manager.ts b/gitnexus/src/storage/repo-manager.ts index b233d44fb..fae9786a9 100644 --- a/gitnexus/src/storage/repo-manager.ts +++ b/gitnexus/src/storage/repo-manager.ts @@ -6,6 +6,7 @@ * so the MCP server can discover indexed repos from any cwd. */ +import { readFileSync } from 'fs'; import fs from 'fs/promises'; import path from 'path'; import os from 'os'; @@ -320,6 +321,16 @@ export const listRegisteredRepos = async (opts?: { // ─── Global CLI Config (~/.gitnexus/config.json) ───────────────────────── +export type CLIEmbeddingProvider = 'local' | 'openai' | 'openrouter' | 'azure' | 'custom'; + +export interface CLIEmbeddingConfig { + provider?: CLIEmbeddingProvider; + baseUrl?: string; + model?: string; + apiKey?: string; + dimensions?: number; +} + export interface CLIConfig { apiKey?: string; model?: string; @@ -330,6 +341,7 @@ export interface CLIConfig { apiVersion?: string; /** Set true when the deployment is a reasoning model (o1, o3, o4-mini). Auto-detected for OpenAI; must be set for Azure deployments. */ isReasoningModel?: boolean; + embedding?: CLIEmbeddingConfig; } /** @@ -351,6 +363,19 @@ export const loadCLIConfig = async (): Promise => { } }; +/** + * Load CLI config synchronously from ~/.gitnexus/config.json. + * Used by startup-time configuration paths that cannot await. + */ +export const loadCLIConfigSync = (): CLIConfig => { + try { + const raw = readFileSync(getGlobalConfigPath(), 'utf-8'); + return JSON.parse(raw) as CLIConfig; + } catch { + return {}; + } +}; + /** * Save CLI config to ~/.gitnexus/config.json */ diff --git a/gitnexus/test/unit/calltool-dispatch.test.ts b/gitnexus/test/unit/calltool-dispatch.test.ts index ebccb42a2..aff149360 100644 --- a/gitnexus/test/unit/calltool-dispatch.test.ts +++ b/gitnexus/test/unit/calltool-dispatch.test.ts @@ -37,6 +37,9 @@ vi.mock('../../src/mcp/core/lbug-adapter.js', async (importOriginal) => { vi.mock('../../src/storage/repo-manager.js', () => ({ listRegisteredRepos: vi.fn().mockResolvedValue([]), cleanupOldKuzuFiles: vi.fn().mockResolvedValue({ found: false, needsReindex: false }), + loadCLIConfig: vi.fn().mockResolvedValue({}), + loadCLIConfigSync: vi.fn().mockReturnValue({}), + saveCLIConfig: vi.fn().mockResolvedValue(undefined), })); // Also mock the search modules to avoid loading onnxruntime @@ -159,6 +162,91 @@ describe('LocalBackend.callTool', () => { expect(result[0].name).toBe('test-project'); }); + it('dispatches refresh_repos without repo resolution', async () => { + const result = await backend.callTool('refresh_repos', {}); + expect(result.refreshed).toBe(true); + expect(result.repoCount).toBe(1); + expect(result.repos[0].name).toBe('test-project'); + }); + + it('dispatches get_index_job', async () => { + (backend as any).jobQueries.getJob = vi.fn().mockReturnValue({ + kind: 'analyze', + id: 'job-1', + status: 'analyzing', + repoName: 'test-project', + progress: { phase: 'analyzing', percent: 50, message: 'Halfway' }, + startedAt: 1, + retryCount: 0, + }); + + const result = await backend.callTool('get_index_job', { jobId: 'job-1' }); + expect(result.kind).toBe('analyze'); + expect(result.id).toBe('job-1'); + expect(result.progress.percent).toBe(50); + }); + + it('dispatches analyze_repo with an absolute path', async () => { + (backend as any).analyzeJobs.startJob = vi.fn().mockReturnValue({ + id: 'job-2', + status: 'queued', + }); + + const result = await backend.callTool('analyze_repo', { + path: 'C:\\tmp\\new-repo', + force: true, + embeddings: true, + }); + + expect((backend as any).analyzeJobs.startJob).toHaveBeenCalledWith({ + repoUrl: undefined, + repoPath: 'C:\\tmp\\new-repo', + force: true, + embeddings: true, + }); + expect(result.jobId).toBe('job-2'); + expect(result.kind).toBe('analyze'); + }); + + it('rejects analyze_repo when path is relative', async () => { + await expect(backend.callTool('analyze_repo', { path: './relative' })).rejects.toThrow( + '"path" must be an absolute path', + ); + }); + + it('dispatches rebuild_embeddings with provider overrides', async () => { + (backend as any).embedJobs.startJob = vi.fn().mockResolvedValue({ + id: 'job-3', + status: 'analyzing', + }); + + const result = await backend.callTool('rebuild_embeddings', { + repo: 'test-project', + provider: 'openai', + baseUrl: 'https://api.openai.com/v1', + model: 'text-embedding-3-large', + apiKey: 'secret', + dimensions: 3072, + saveConfig: true, + }); + + expect((backend as any).embedJobs.startJob).toHaveBeenCalledWith({ + repoName: 'test-project', + storagePath: '/tmp/.gitnexus/test-project', + embedding: { + provider: 'openai', + baseUrl: 'https://api.openai.com/v1', + model: 'text-embedding-3-large', + apiKey: 'secret', + dimensions: 3072, + }, + saveConfig: true, + }); + expect(result.kind).toBe('embed'); + expect(result.jobId).toBe('job-3'); + expect(result.savedConfig).toBe(true); + }); + it('throws for unknown tool name', async () => { await expect(backend.callTool('nonexistent_tool', {})).rejects.toThrow( 'Unknown tool: nonexistent_tool', diff --git a/gitnexus/test/unit/embedding-config.test.ts b/gitnexus/test/unit/embedding-config.test.ts new file mode 100644 index 000000000..4a2dde055 --- /dev/null +++ b/gitnexus/test/unit/embedding-config.test.ts @@ -0,0 +1,189 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +const loadCLIConfig = vi.fn(); +const loadCLIConfigSync = vi.fn(); + +vi.mock('../../src/storage/repo-manager.js', () => ({ + loadCLIConfig, + loadCLIConfigSync, +})); + +const ENV_KEYS = [ + 'GITNEXUS_EMBEDDING_URL', + 'GITNEXUS_EMBEDDING_MODEL', + 'GITNEXUS_EMBEDDING_API_KEY', + 'GITNEXUS_EMBEDDING_DIMS', +] as const; + +describe('embedding config resolver', () => { + const savedEnv = Object.fromEntries(ENV_KEYS.map((key) => [key, process.env[key]])); + + beforeEach(() => { + vi.resetModules(); + loadCLIConfig.mockResolvedValue({}); + loadCLIConfigSync.mockReturnValue({}); + for (const key of ENV_KEYS) delete process.env[key]; + }); + + afterEach(() => { + vi.clearAllMocks(); + for (const key of ENV_KEYS) { + if (savedEnv[key] === undefined) { + delete process.env[key]; + } else { + process.env[key] = savedEnv[key]; + } + } + }); + + it('prefers explicit overrides over saved config and env', async () => { + loadCLIConfigSync.mockReturnValue({ + embedding: { + provider: 'custom', + baseUrl: 'https://config.example/v1', + model: 'config-model', + apiKey: 'config-key', + dimensions: 1024, + }, + }); + process.env.GITNEXUS_EMBEDDING_URL = 'https://env.example/v1'; + process.env.GITNEXUS_EMBEDDING_MODEL = 'env-model'; + process.env.GITNEXUS_EMBEDDING_DIMS = '2048'; + + const { resolveEmbeddingConfigSync } = await import('../../src/core/embeddings/config.js'); + const resolved = resolveEmbeddingConfigSync({ + provider: 'openai', + baseUrl: 'https://override.example/v1', + model: 'override-model', + apiKey: 'override-key', + dimensions: 3072, + }); + + expect(resolved).toMatchObject({ + mode: 'http', + provider: 'openai', + baseUrl: 'https://override.example/v1', + model: 'override-model', + apiKey: 'override-key', + dimensions: 3072, + explicitDimensionsSource: 'overrides', + }); + }); + + it('uses saved config before env fallback', async () => { + loadCLIConfigSync.mockReturnValue({ + embedding: { + provider: 'custom', + baseUrl: 'https://config.example/v1', + model: 'config-model', + apiKey: 'config-key', + dimensions: 1536, + }, + }); + process.env.GITNEXUS_EMBEDDING_URL = 'https://env.example/v1'; + process.env.GITNEXUS_EMBEDDING_MODEL = 'env-model'; + process.env.GITNEXUS_EMBEDDING_DIMS = '2048'; + + const { resolveEmbeddingConfigSync } = await import('../../src/core/embeddings/config.js'); + const resolved = resolveEmbeddingConfigSync(); + + expect(resolved).toMatchObject({ + mode: 'http', + provider: 'custom', + baseUrl: 'https://config.example/v1', + model: 'config-model', + apiKey: 'config-key', + dimensions: 1536, + explicitDimensionsSource: 'config', + }); + }); + + it('falls back to environment configuration when saved config is empty', async () => { + process.env.GITNEXUS_EMBEDDING_URL = 'https://env.example/v1'; + process.env.GITNEXUS_EMBEDDING_MODEL = 'env-model'; + process.env.GITNEXUS_EMBEDDING_API_KEY = 'env-key'; + process.env.GITNEXUS_EMBEDDING_DIMS = '768'; + + const { resolveEmbeddingConfigSync } = await import('../../src/core/embeddings/config.js'); + const resolved = resolveEmbeddingConfigSync(); + + expect(resolved).toMatchObject({ + mode: 'http', + provider: 'custom', + baseUrl: 'https://env.example/v1', + model: 'env-model', + apiKey: 'env-key', + dimensions: 768, + explicitDimensionsSource: 'env', + }); + }); + + it('falls back to local defaults when no HTTP backend is configured', async () => { + const { resolveEmbeddingConfigSync } = await import('../../src/core/embeddings/config.js'); + const resolved = resolveEmbeddingConfigSync(); + + expect(resolved).toMatchObject({ + mode: 'local', + provider: 'local', + model: 'Snowflake/snowflake-arctic-embed-xs', + dimensions: 384, + apiKey: '', + }); + expect(resolved.baseUrl).toBeUndefined(); + expect(resolved.explicitDimensionsSource).toBeUndefined(); + }); + + it('ignores custom dimensions when local mode is forced', async () => { + loadCLIConfigSync.mockReturnValue({ + embedding: { + provider: 'local', + dimensions: 3072, + }, + }); + process.env.GITNEXUS_EMBEDDING_DIMS = '2048'; + + const { resolveEmbeddingConfigSync } = await import('../../src/core/embeddings/config.js'); + const resolved = resolveEmbeddingConfigSync(); + + expect(resolved.mode).toBe('local'); + expect(resolved.dimensions).toBe(384); + expect(resolved.explicitDimensionsSource).toBeUndefined(); + }); + + it('supports async resolution from saved config', async () => { + loadCLIConfig.mockResolvedValue({ + embedding: { + provider: 'openrouter', + baseUrl: 'https://openrouter.ai/api/v1', + model: 'text-embedding-model', + apiKey: 'saved-key', + dimensions: 1024, + }, + }); + + const { resolveEmbeddingConfig } = await import('../../src/core/embeddings/config.js'); + const resolved = await resolveEmbeddingConfig(); + + expect(resolved).toMatchObject({ + mode: 'http', + provider: 'openrouter', + baseUrl: 'https://openrouter.ai/api/v1', + model: 'text-embedding-model', + apiKey: 'saved-key', + dimensions: 1024, + explicitDimensionsSource: 'config', + }); + }); + + it('throws when dimensions are invalid', async () => { + process.env.GITNEXUS_EMBEDDING_URL = 'https://env.example/v1'; + process.env.GITNEXUS_EMBEDDING_MODEL = 'env-model'; + process.env.GITNEXUS_EMBEDDING_DIMS = '0'; + + const { resolveEmbeddingConfigSync } = await import('../../src/core/embeddings/config.js'); + + expect(() => resolveEmbeddingConfigSync()).toThrow( + 'GITNEXUS_EMBEDDING_DIMS must be a positive integer', + ); + }); +}); diff --git a/gitnexus/test/unit/http-embedder.test.ts b/gitnexus/test/unit/http-embedder.test.ts index 0dea77758..39a8381d9 100644 --- a/gitnexus/test/unit/http-embedder.test.ts +++ b/gitnexus/test/unit/http-embedder.test.ts @@ -1,7 +1,11 @@ +import fs from 'fs/promises'; +import os from 'os'; +import path from 'path'; import { describe, it, expect, vi, afterEach } from 'vitest'; import { getEmbeddingDims, isEmbedderReady } from '../../src/mcp/core/embedder.js'; const ENV_KEYS = [ + 'GITNEXUS_HOME', 'GITNEXUS_EMBEDDING_URL', 'GITNEXUS_EMBEDDING_MODEL', 'GITNEXUS_EMBEDDING_API_KEY', @@ -15,9 +19,13 @@ describe('HTTP embedding backend', () => { // Save original env state before any test mutates it const savedEnv = Object.fromEntries(ENV_KEYS.map((k) => [k, process.env[k]])); - afterEach(() => { + afterEach(async () => { + const tmpHome = process.env.GITNEXUS_HOME; vi.unstubAllGlobals(); vi.resetModules(); + if (tmpHome && tmpHome !== savedEnv.GITNEXUS_HOME) { + await fs.rm(tmpHome, { recursive: true, force: true }); + } // Restore env vars to pre-test state so a mid-test throw can't leak for (const key of ENV_KEYS) { if (savedEnv[key] === undefined) { @@ -44,6 +52,26 @@ describe('HTTP embedding backend', () => { expect(mod.isEmbedderReady()).toBe(true); }); + it('reads HTTP mode from saved embedding config', async () => { + const tmpHome = await fs.mkdtemp(path.join(os.tmpdir(), 'gitnexus-http-config-')); + process.env.GITNEXUS_HOME = tmpHome; + await fs.writeFile( + path.join(tmpHome, 'config.json'), + JSON.stringify({ + embedding: { + baseUrl: 'http://localhost:8080/v1', + model: 'config-model', + dimensions: 1024, + }, + }), + 'utf-8', + ); + + const mod = await import('../../src/mcp/core/embedder.js'); + expect(mod.isEmbedderReady()).toBe(true); + expect(mod.getEmbeddingDims()).toBe(1024); + }); + it('reads custom dimensions from environment', async () => { process.env.GITNEXUS_EMBEDDING_URL = 'http://localhost:8080/v1'; process.env.GITNEXUS_EMBEDDING_MODEL = 'test-model'; @@ -96,6 +124,45 @@ describe('HTTP embedding backend', () => { expect(result.length).toBe(384); }); + it('sends requests using saved embedding config', async () => { + const tmpHome = await fs.mkdtemp(path.join(os.tmpdir(), 'gitnexus-http-config-')); + process.env.GITNEXUS_HOME = tmpHome; + await fs.writeFile( + path.join(tmpHome, 'config.json'), + JSON.stringify({ + embedding: { + baseUrl: 'http://config-test:8080/v1/', + model: 'config-model', + apiKey: 'config-key', + dimensions: 384, + }, + }), + 'utf-8', + ); + + const mockEmbedding = Array.from({ length: 384 }, (_, i) => i * 0.001); + vi.stubGlobal( + 'fetch', + vi.fn().mockResolvedValue({ + ok: true, + json: async () => ({ data: [{ embedding: mockEmbedding }] }), + }), + ); + + const { embedText } = await import('../../src/core/embeddings/embedder.js'); + const result = await embedText('config text'); + + expect(fetch).toHaveBeenCalledOnce(); + const [url, init] = (fetch as any).mock.calls[0]; + const body = JSON.parse(init.body); + expect(url).toBe('http://config-test:8080/v1/embeddings'); + expect(body.model).toBe('config-model'); + expect(body.input).toEqual(['config text']); + expect(init.headers.Authorization).toBe('Bearer config-key'); + expect(result).toBeInstanceOf(Float32Array); + expect(result.length).toBe(384); + }); + it('retries on server error', async () => { process.env.GITNEXUS_EMBEDDING_URL = 'http://test:8080/v1'; process.env.GITNEXUS_EMBEDDING_MODEL = 'test-model'; @@ -265,7 +332,9 @@ describe('HTTP embedding backend', () => { expect(EMBEDDING_DIMS).toBe(384); }); - it('reads dimensions from environment variable', async () => { + it('reads dimensions from HTTP embedding environment configuration', async () => { + process.env.GITNEXUS_EMBEDDING_URL = 'http://test:8080/v1'; + process.env.GITNEXUS_EMBEDDING_MODEL = 'test-model'; process.env.GITNEXUS_EMBEDDING_DIMS = '1024'; const { EMBEDDING_DIMS } = await import('../../src/core/lbug/schema.js'); expect(EMBEDDING_DIMS).toBe(1024); diff --git a/gitnexus/test/unit/tools.test.ts b/gitnexus/test/unit/tools.test.ts index 4274716a7..349b1a2cd 100644 --- a/gitnexus/test/unit/tools.test.ts +++ b/gitnexus/test/unit/tools.test.ts @@ -2,7 +2,7 @@ * Unit Tests: MCP Tool Definitions * * Tests: GITNEXUS_TOOLS from tools.ts - * - All 16 tools are defined (per-repo + group_*) + * - All 20 tools are defined (read tools + write tools + group_*) * - Each tool has valid name, description, inputSchema * - Required fields are correct * - Optional repo parameter is present on tools that need it @@ -19,8 +19,8 @@ const GROUP_TOOLS = new Set([ ]); describe('GITNEXUS_TOOLS', () => { - it('exports all tools (7 base + 3 route/tool/shape + 1 api_impact + 5 group)', () => { - expect(GITNEXUS_TOOLS).toHaveLength(16); + it('exports all tools (11 core + 4 write + 5 group)', () => { + expect(GITNEXUS_TOOLS).toHaveLength(20); }); it('contains all expected tool names', () => { @@ -35,6 +35,10 @@ describe('GITNEXUS_TOOLS', () => { 'rename', 'impact', 'api_impact', + 'refresh_repos', + 'get_index_job', + 'analyze_repo', + 'rebuild_embeddings', ]), ); }); @@ -95,12 +99,33 @@ describe('GITNEXUS_TOOLS', () => { for (const tool of GITNEXUS_TOOLS) { if (tool.name === 'list_repos') continue; if (GROUP_TOOLS.has(tool.name)) continue; + if (['refresh_repos', 'get_index_job', 'analyze_repo'].includes(tool.name)) continue; expect(tool.inputSchema.properties.repo).toBeDefined(); expect(tool.inputSchema.properties.repo.type).toBe('string'); expect(tool.inputSchema.required).not.toContain('repo'); } }); + it('write tools expose the expected required fields', () => { + const getJobTool = GITNEXUS_TOOLS.find((t) => t.name === 'get_index_job')!; + expect(getJobTool.inputSchema.required).toEqual(['jobId']); + + const analyzeRepoTool = GITNEXUS_TOOLS.find((t) => t.name === 'analyze_repo')!; + expect(analyzeRepoTool.inputSchema.properties.path).toBeDefined(); + expect(analyzeRepoTool.inputSchema.properties.url).toBeDefined(); + expect(analyzeRepoTool.inputSchema.properties.force).toBeDefined(); + expect(analyzeRepoTool.inputSchema.properties.embeddings).toBeDefined(); + + const rebuildTool = GITNEXUS_TOOLS.find((t) => t.name === 'rebuild_embeddings')!; + expect(rebuildTool.inputSchema.properties.repo).toBeDefined(); + expect(rebuildTool.inputSchema.properties.provider).toBeDefined(); + expect(rebuildTool.inputSchema.properties.baseUrl).toBeDefined(); + expect(rebuildTool.inputSchema.properties.model).toBeDefined(); + expect(rebuildTool.inputSchema.properties.apiKey).toBeDefined(); + expect(rebuildTool.inputSchema.properties.dimensions).toBeDefined(); + expect(rebuildTool.inputSchema.properties.saveConfig).toBeDefined(); + }); + it('group_contracts has optional repo filter', () => { const groupContracts = GITNEXUS_TOOLS.find((t) => t.name === 'group_contracts')!; expect(groupContracts.inputSchema.properties).toHaveProperty('repo');