mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-10 03:27:59 +00:00
feat: 增加 MCP 索引写工具与 embedding 配置解析
This commit is contained in:
parent
5be0537ce4
commit
2bdf1d75a4
21 changed files with 1608 additions and 566 deletions
142
gitnexus/src/core/embeddings/config.ts
Normal file
142
gitnexus/src/core/embeddings/config.ts
Normal file
|
|
@ -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<CLIEmbeddingConfig> {}
|
||||
|
||||
export interface ResolvedEmbeddingConfig {
|
||||
mode: ResolvedEmbeddingMode;
|
||||
provider: CLIEmbeddingProvider;
|
||||
model: string;
|
||||
dimensions: number;
|
||||
baseUrl?: string;
|
||||
apiKey: string;
|
||||
explicitDimensionsSource?: Exclude<EmbeddingConfigSource, 'default'>;
|
||||
}
|
||||
|
||||
interface ResolvedValue<T> {
|
||||
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 = <T>(
|
||||
overridesValue: T | undefined,
|
||||
configValue: T | undefined,
|
||||
envValue: T | undefined,
|
||||
defaultValue: T,
|
||||
): ResolvedValue<T> => {
|
||||
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<EmbeddingConfigSource, 'default'>,
|
||||
): 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<CLIEmbeddingProvider | undefined>(
|
||||
overrides.provider,
|
||||
savedEmbeddingConfig.provider,
|
||||
undefined,
|
||||
undefined,
|
||||
);
|
||||
const baseUrl = resolveValue<string | undefined>(
|
||||
trimToUndefined(overrides.baseUrl),
|
||||
trimToUndefined(savedEmbeddingConfig.baseUrl),
|
||||
trimToUndefined(process.env.GITNEXUS_EMBEDDING_URL),
|
||||
undefined,
|
||||
);
|
||||
const httpModel = resolveValue<string | undefined>(
|
||||
trimToUndefined(overrides.model),
|
||||
trimToUndefined(savedEmbeddingConfig.model),
|
||||
trimToUndefined(process.env.GITNEXUS_EMBEDDING_MODEL),
|
||||
undefined,
|
||||
);
|
||||
|
||||
const dimensions = resolveValue<number>(
|
||||
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<string | undefined>(
|
||||
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<ResolvedEmbeddingConfig> => {
|
||||
return resolveEmbeddingConfigFromSaved(await loadCLIConfig(), overrides);
|
||||
};
|
||||
|
|
@ -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;
|
||||
};
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -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<Float32Array[]> => {
|
|||
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<number[]> => {
|
|||
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;
|
||||
|
|
|
|||
|
|
@ -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';
|
||||
|
|
|
|||
275
gitnexus/src/core/index-jobs/analyze-job-service.ts
Normal file
275
gitnexus/src/core/index-jobs/analyze-job-service.ts
Normal file
|
|
@ -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<void>;
|
||||
repoLocks: RepoLockManager;
|
||||
jobManager?: JobManager;
|
||||
logger?: Pick<Console, 'warn' | 'error'>;
|
||||
}
|
||||
|
||||
export class AnalyzeJobService {
|
||||
readonly jobManager: JobManager;
|
||||
|
||||
private readonly backendInit: () => Promise<void>;
|
||||
private readonly repoLocks: RepoLockManager;
|
||||
private readonly logger: Pick<Console, 'warn' | 'error'>;
|
||||
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<void> {
|
||||
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';
|
||||
}
|
||||
}
|
||||
138
gitnexus/src/core/index-jobs/embed-job-service.ts
Normal file
138
gitnexus/src/core/index-jobs/embed-job-service.ts
Normal file
|
|
@ -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<AnalyzeJob> {
|
||||
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();
|
||||
}
|
||||
}
|
||||
5
gitnexus/src/core/index-jobs/index.ts
Normal file
5
gitnexus/src/core/index-jobs/index.ts
Normal file
|
|
@ -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';
|
||||
190
gitnexus/src/core/index-jobs/job-manager.ts
Normal file
190
gitnexus/src/core/index-jobs/job-manager.ts
Normal file
|
|
@ -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<string, AnalyzeJob>();
|
||||
private children = new Map<string, ChildProcess>();
|
||||
private timeouts = new Map<string, ReturnType<typeof setTimeout>>();
|
||||
private emitter = new EventEmitter();
|
||||
private cleanupTimer: ReturnType<typeof setInterval>;
|
||||
|
||||
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<AnalyzeJob, 'status' | 'progress' | 'error' | 'repoPath' | 'repoName' | 'completedAt'>
|
||||
>,
|
||||
) {
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
51
gitnexus/src/core/index-jobs/job-query-service.ts
Normal file
51
gitnexus/src/core/index-jobs/job-query-service.ts
Normal file
|
|
@ -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<IndexJobKind, JobManager>();
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
26
gitnexus/src/core/index-jobs/repo-lock-manager.ts
Normal file
26
gitnexus/src/core/index-jobs/repo-lock-manager.ts
Normal file
|
|
@ -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<string>();
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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} (
|
||||
|
|
|
|||
|
|
@ -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<number[]> => {
|
|||
* Get embedding dimensions
|
||||
*/
|
||||
export const getEmbeddingDims = (): number => {
|
||||
return getHttpDimensions() ?? 384;
|
||||
return isHttpMode() ? (getHttpDimensions() ?? 384) : resolveEmbeddingConfigSync().dimensions;
|
||||
};
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -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<string, Promise<void>> = new Map();
|
||||
private lastStalenessCheck: Map<string, number> = 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<void> {
|
||||
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<any> {
|
||||
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 ────────────────────────────────────────
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -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).
|
||||
|
|
|
|||
|
|
@ -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<string, AnalyzeJob>();
|
||||
private children = new Map<string, ChildProcess>();
|
||||
private timeouts = new Map<string, ReturnType<typeof setTimeout>>();
|
||||
private emitter = new EventEmitter();
|
||||
private cleanupTimer: ReturnType<typeof setInterval>;
|
||||
|
||||
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<AnalyzeJob, 'status' | 'progress' | 'error' | 'repoPath' | 'repoName' | 'completedAt'>
|
||||
>,
|
||||
) {
|
||||
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';
|
||||
|
|
|
|||
|
|
@ -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<string>();
|
||||
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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<CLIConfig> => {
|
|||
}
|
||||
};
|
||||
|
||||
/**
|
||||
* 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
|
||||
*/
|
||||
|
|
|
|||
|
|
@ -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',
|
||||
|
|
|
|||
189
gitnexus/test/unit/embedding-config.test.ts
Normal file
189
gitnexus/test/unit/embedding-config.test.ts
Normal file
|
|
@ -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',
|
||||
);
|
||||
});
|
||||
});
|
||||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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');
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue