From 7c4e5700ffec19cff26667890d829e7d76b5ec96 Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Tue, 9 Jun 2026 09:54:50 +0000 Subject: [PATCH] fix(review): resolve tri-review findings on the upload flow MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A multi-agent review of the upload implementation surfaced a P0 plus several P2/P3s; all are addressed here. - P0: the upload handler took the single analysis slot (createJob) before validating/promoting, so any failure in that window left a queued job that was never failed — wedging ALL analysis until restart (trivially triggered by a single-segment manifest). Now: validate the folder before taking the slot, release it via failJob on any pre-launch error, and reject single-segment / multi-top manifests during ingest (also fixes a silent file-drop). - CI: rate-limit.test's source-regex broke when Prettier wrapped the /api/analyze registration; made it wrapping-tolerant. - Resource: the startup sweep now also removes stale promoted upload dirs with no .gitnexus index (orphans from analyses that failed before registering). - Frontend: guard against post-unmount SSE opening, reset upload state on cancel/mode-change, guard concurrent uploads, fall back to the folder name, add aria-busy, and fix the {{count}} plural ("1 files"). - Maintainability: extract launchAnalysisWorker into analyze-launch.ts (DI + typed WorkerMessage IPC), move requireLocalhostOrigin to middleware.ts, share REPO_NAME_PATTERN, tighten UploadJobRef, name the collision-retry constant. Co-Authored-By: Claude Opus 4.8 (1M context) --- gitnexus-web/src/components/RepoAnalyzer.tsx | 29 ++- gitnexus-web/src/locales/en/onboarding.json | 2 +- .../src/locales/zh-CN/onboarding.json | 2 +- gitnexus/src/server/analyze-launch.ts | 170 ++++++++++++++++++ gitnexus/src/server/analyze-upload.ts | 89 +++++---- gitnexus/src/server/analyze-worker.ts | 9 +- gitnexus/src/server/api.ts | 156 ++-------------- gitnexus/src/server/git-clone.ts | 2 +- gitnexus/src/server/middleware.ts | 29 +++ gitnexus/src/server/upload-ingest.ts | 17 +- gitnexus/src/server/upload-paths.ts | 8 +- gitnexus/src/server/upload-sweep.ts | 19 +- gitnexus/test/unit/api-analyze-upload.test.ts | 86 ++++++++- gitnexus/test/unit/rate-limit.test.ts | 4 +- gitnexus/test/unit/upload-sweep.test.ts | 19 ++ 15 files changed, 415 insertions(+), 226 deletions(-) create mode 100644 gitnexus/src/server/analyze-launch.ts create mode 100644 gitnexus/src/server/middleware.ts diff --git a/gitnexus-web/src/components/RepoAnalyzer.tsx b/gitnexus-web/src/components/RepoAnalyzer.tsx index 42aa97cf4..0e4a55d90 100644 --- a/gitnexus-web/src/components/RepoAnalyzer.tsx +++ b/gitnexus-web/src/components/RepoAnalyzer.tsx @@ -188,9 +188,11 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp const sseControllerRef = useRef(null); const completeTimerRef = useRef | null>(null); const folderInputRef = useRef(null); + const isMountedRef = useRef(true); useEffect(() => { return () => { + isMountedRef.current = false; sseControllerRef.current?.abort(); if (completeTimerRef.current) clearTimeout(completeTimerRef.current); }; @@ -202,13 +204,13 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp setGitlabUrl(''); setLocalPath(''); setValidationError(null); + setUploadSummary(null); + setUploading(false); }; - // Use the browser's native directory picker (webkitdirectory doesn't give paths, - // so we use a text input + a "Browse" button that opens a standard file input - // to let users pick files from the folder — the path is typed manually since - // browsers don't expose absolute paths for security reasons). - // For local paths, the user types or pastes the absolute path. + // Local-folder mode uploads the selected folder's files (the browser never + // exposes an absolute path, so the old typed-path/browse approach couldn't + // work — see handleFolderUpload). A typed server path is also still accepted. const canSubmit = mode === 'github' @@ -259,6 +261,9 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp // Drive an already-created analysis job through the SSE progress stream to // completion. Shared by the path/URL analyze flow and the folder-upload flow. const trackJob = (jobId: string, fallbackNameSource: string | null) => { + // The component may have unmounted while an upload POST was in flight; don't + // open an SSE stream nothing will ever abort. + if (!isMountedRef.current) return; jobIdRef.current = jobId; setPhase('analyzing'); const controller = streamAnalyzeProgress( @@ -290,6 +295,7 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp // Upload a browser-selected folder (webkitdirectory) and start analysis. The // upload endpoint returns a jobId, which then joins the normal SSE flow. const handleFolderUpload = async (fileList: FileList) => { + if (uploading || isLoading) return; // guard against a concurrent upload const { files, manifest, droppedCount } = filterRepoFiles(fileList); if (files.length === 0) { setValidationError(t('onboarding:repoAnalyzer.upload.empty')); @@ -299,11 +305,16 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp setUploadSummary({ count: files.length, dropped: droppedCount }); setUploading(true); setPhase('starting'); + // The selected folder's name (manifest entries are `/`) is a + // sensible fallback if the server's complete event omits repoName. + const folderName = manifest[0]?.split('/')[0] ?? null; try { const { jobId } = await uploadFolder(files, manifest); + if (!isMountedRef.current) return; setUploading(false); - trackJob(jobId, null); + trackJob(jobId, folderName); } catch (err) { + if (!isMountedRef.current) return; setUploading(false); setValidationError(err instanceof Error ? err.message : t('errors:startAnalysisFailed')); setPhase('error'); @@ -321,6 +332,8 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp } setPhase('input'); setProgress({ phase: 'queued', percent: 0, message: t('common:analyzePhases.queued') }); + setUploading(false); + setUploadSummary(null); }; const isLoading = phase === 'starting'; @@ -508,7 +521,7 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp {t('onboarding:repoAnalyzer.upload.button')} {uploading && ( -
+
@@ -520,7 +533,7 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp {uploadSummary && !uploading && phase !== 'error' && (

{t('onboarding:repoAnalyzer.upload.selected', { - count: uploadSummary.count, + fileCount: uploadSummary.count, dropped: uploadSummary.dropped, })}

diff --git a/gitnexus-web/src/locales/en/onboarding.json b/gitnexus-web/src/locales/en/onboarding.json index fd746a290..8be58efce 100644 --- a/gitnexus-web/src/locales/en/onboarding.json +++ b/gitnexus-web/src/locales/en/onboarding.json @@ -65,7 +65,7 @@ "upload": { "button": "Upload a folder", "uploading": "Uploading…", - "selected": "{{count}} files ready ({{dropped}} skipped: .git, node_modules, build output)", + "selected": "{{fileCount}} files ready ({{dropped}} skipped: .git, node_modules, build output)", "empty": "No analyzable files found in that folder." } } diff --git a/gitnexus-web/src/locales/zh-CN/onboarding.json b/gitnexus-web/src/locales/zh-CN/onboarding.json index 4ac2c44d8..f20e60b6d 100644 --- a/gitnexus-web/src/locales/zh-CN/onboarding.json +++ b/gitnexus-web/src/locales/zh-CN/onboarding.json @@ -65,7 +65,7 @@ "upload": { "button": "上传文件夹", "uploading": "上传中…", - "selected": "已准备 {{count}} 个文件(已跳过 {{dropped}} 个:.git、node_modules、构建产物)", + "selected": "已准备 {{fileCount}} 个文件(已跳过 {{dropped}} 个:.git、node_modules、构建产物)", "empty": "该文件夹中未找到可分析的文件。" } } diff --git a/gitnexus/src/server/analyze-launch.ts b/gitnexus/src/server/analyze-launch.ts new file mode 100644 index 000000000..184c5a813 --- /dev/null +++ b/gitnexus/src/server/analyze-launch.ts @@ -0,0 +1,170 @@ +/** + * Shared analyze-worker launcher. + * + * Forks the analyze worker for an already-resolved repo directory and owns the + * lock + auto-retry + IPC machinery. Used by both the JSON `/api/analyze` route + * and the multipart `/api/analyze/upload` route. Dependency-injected (like + * createAnalyzeUploadHandler) so the seam is testable and api.ts stays smaller. + * + * NOTE: this module must live alongside analyze-worker.{ts,js} — the worker + * path is resolved relative to `import.meta.url`. + */ + +import path from 'path'; +import { fork } from 'child_process'; +import { fileURLToPath, pathToFileURL } from 'url'; +import { createRequire } from 'node:module'; +import { getStoragePath } from '../storage/repo-manager.js'; +import { logger } from '../core/logger.js'; +import type { JobManager } from './analyze-job.js'; +import type { WorkerMessage } from './analyze-worker.js'; + +const _require = createRequire(import.meta.url); + +export interface LaunchDeps { + jobManager: JobManager; + backend: { init: () => Promise }; + acquireRepoLock: (key: string) => string | null; + releaseRepoLock: (key: string) => void; +} + +export interface LaunchOptions { + force?: boolean; + embeddings?: boolean; + dropEmbeddings?: boolean; + registryName?: string; +} + +const MAX_WORKER_RETRIES = 2; + +export function createLaunchAnalysisWorker(deps: LaunchDeps) { + const { jobManager, backend, acquireRepoLock, releaseRepoLock } = deps; + + return function launchAnalysisWorker( + job: { id: string }, + targetPath: string, + opts: LaunchOptions, + ): void { + // 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 ────────────────────────────── + 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: WorkerMessage) => { + 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) => { + logger.error({ err }, 'backend.init() failed after analyze:'); + 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() || ''; + logger.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: !!opts.force, + embeddings: !!opts.embeddings, + dropEmbeddings: !!opts.dropEmbeddings, + ...(opts.registryName ? { registryName: opts.registryName } : {}), + }, + }); + }; + + forkWorker(); + }; +} diff --git a/gitnexus/src/server/analyze-upload.ts b/gitnexus/src/server/analyze-upload.ts index ee58c58e8..adb8b861b 100644 --- a/gitnexus/src/server/analyze-upload.ts +++ b/gitnexus/src/server/analyze-upload.ts @@ -15,17 +15,25 @@ import type { IncomingMessage } from 'http'; import { ingestUpload } from './upload-ingest.js'; import { UPLOAD_ROOT, getUploadDir, deriveUploadName } from './upload-paths.js'; import { BadRequestError } from './validation.js'; +import type { AnalyzeJob } from './analyze-job.js'; -export interface UploadJobRef { - id: string; - status: string; -} +/** Minimal job shape the handler needs (a subset of the real AnalyzeJob). */ +export type UploadJobRef = Pick; + +/** Cap on collision-suffix attempts when allocating an upload dir name. */ +const MAX_NAME_COLLISION_TRIES = 100; export interface AnalyzeUploadDeps { /** Create (or throw on busy) an analysis job for the given upload dir. */ createJob: (params: { repoPath: string }) => UploadJobRef; /** Launch the analyze worker against an already-resolved repo directory. */ launch: (job: UploadJobRef, targetPath: string, opts: { registryName: string }) => void; + /** + * Mark a created job failed. The job occupies the single analysis slot from + * createJob onward, so ANY error before launch must release it — otherwise a + * leaked non-terminal job wedges all future analyses until restart. + */ + failJob: (jobId: string, error: string) => void; /** Injectable for tests (defaults to the real ingestUpload). */ ingest?: typeof ingestUpload; } @@ -35,7 +43,7 @@ export interface AnalyzeUploadDeps { * collision with an existing upload. Bounded to avoid an unbounded scan. */ async function pickAvailableName(base: string): Promise { - for (let i = 0; i < 100; i++) { + for (let i = 0; i < MAX_NAME_COLLISION_TRIES; i++) { const name = i === 0 ? base : `${base}-${i + 1}`; let dir: string; try { @@ -50,7 +58,10 @@ async function pickAvailableName(base: string): Promise { return name; // ENOENT → available } } - throw new BadRequestError('Could not allocate an upload directory', 409); + throw new BadRequestError( + `Could not allocate an upload directory after ${MAX_NAME_COLLISION_TRIES} attempts`, + 409, + ); } export function createAnalyzeUploadHandler(deps: AnalyzeUploadDeps) { @@ -59,20 +70,38 @@ export function createAnalyzeUploadHandler(deps: AnalyzeUploadDeps) { return async function handleAnalyzeUploadRequest(req: Request, res: Response): Promise { let stageRoot: string | undefined; let promotedDir: string | undefined; + let createdJobId: string | undefined; let launched = false; try { - const result = await ingest(req as unknown as IncomingMessage); + const result = await ingest(req as IncomingMessage); stageRoot = result.stageRoot; const baseName = deriveUploadName(result.topLevelName); if (!baseName) { throw new BadRequestError('Uploaded folder has no usable name'); } + + // webkitRelativePath prefixes every entry with the picked folder, so the + // real repo root is stageRoot/. Validate it is a directory + // BEFORE taking the single analysis slot — a malformed (non-folder) + // upload must not be able to occupy the slot. + const innerRoot = path.join(result.stageRoot, result.topLevelName); + let innerIsDir = false; + try { + innerIsDir = (await fsp.stat(innerRoot)).isDirectory(); + } catch { + innerIsDir = false; + } + if (!innerIsDir) { + throw new BadRequestError('Upload must be a folder'); + } + const finalName = await pickAvailableName(baseName); const finalDir = getUploadDir(finalName); - // createJob BEFORE promote: a busy server throws → 409 and nothing has - // been moved into place yet, so the staging dir is the only thing to clean. + // createJob occupies the single analysis slot (throws 'already in + // progress' → 409). From here on, ANY error before launch MUST release + // the slot via failJob in the catch, or the server wedges all analyses. let job: UploadJobRef; try { job = deps.createJob({ repoPath: finalDir }); @@ -83,18 +112,7 @@ export function createAnalyzeUploadHandler(deps: AnalyzeUploadDeps) { } throw err; } - - // webkitRelativePath prefixes every entry with the picked folder, so the - // real repo root is stageRoot/. Promote that inner dir. - const innerRoot = path.join(result.stageRoot, result.topLevelName); - try { - if (!(await fsp.stat(innerRoot)).isDirectory()) { - throw new BadRequestError('Upload must be a folder'); - } - } catch (err) { - if (err instanceof BadRequestError) throw err; - throw new BadRequestError('Upload must be a folder'); - } + createdJobId = job.id; // Promote staging → persistent upload dir. Both live under UPLOAD_ROOT's // filesystem, so this rename stays atomic (no EXDEV). @@ -116,6 +134,11 @@ export function createAnalyzeUploadHandler(deps: AnalyzeUploadDeps) { res.status(202).json({ jobId: job.id, status: job.status }); } catch (err) { + // Release the single analysis slot if a job was created but never + // launched — otherwise the leaked queued job blocks all future analyses. + if (createdJobId && !launched) { + deps.failJob(createdJobId, err instanceof Error ? err.message : 'Upload failed'); + } if (stageRoot) { await fsp.rm(stageRoot, { recursive: true, force: true }).catch(() => {}); } @@ -130,27 +153,3 @@ export function createAnalyzeUploadHandler(deps: AnalyzeUploadDeps) { } }; } - -/** - * Per-route guard that restricts an endpoint to localhost browser origins. - * Non-browser requests (no Origin header, e.g. curl) pass through. This closes - * cross-origin reach (e.g. the allow-listed public deploy + Private Network - * Access) to write routes without affecting read routes. - */ -export function requireLocalhostOrigin(req: Request, res: Response, next: () => void): void { - const origin = req.headers.origin; - if (origin === undefined) { - next(); - return; - } - try { - const hostname = new URL(origin).hostname; - if (hostname === 'localhost' || hostname === '127.0.0.1' || hostname === '::1') { - next(); - return; - } - } catch { - /* malformed origin → reject */ - } - res.status(403).json({ error: 'This endpoint is restricted to localhost origins' }); -} diff --git a/gitnexus/src/server/analyze-worker.ts b/gitnexus/src/server/analyze-worker.ts index a87b34179..c53893c2d 100644 --- a/gitnexus/src/server/analyze-worker.ts +++ b/gitnexus/src/server/analyze-worker.ts @@ -20,24 +20,25 @@ interface StartMessage { options: AnalyzeOptions; } -interface ProgressMessage { +export interface ProgressMessage { type: 'progress'; phase: string; percent: number; message: string; } -interface CompleteMessage { +export interface CompleteMessage { type: 'complete'; result: AnalyzeResult; } -interface ErrorMessage { +export interface ErrorMessage { type: 'error'; message: string; } -type WorkerMessage = ProgressMessage | CompleteMessage | ErrorMessage; +/** Child → parent IPC messages. Shared with the parent-side launcher. */ +export type WorkerMessage = ProgressMessage | CompleteMessage | ErrorMessage; function send(msg: WorkerMessage) { process.send?.(msg); diff --git a/gitnexus/src/server/api.ts b/gitnexus/src/server/api.ts index 16f3b2c6f..67ae63009 100644 --- a/gitnexus/src/server/api.ts +++ b/gitnexus/src/server/api.ts @@ -30,12 +30,13 @@ import { searchFTSFromLbug } from '../core/search/bm25-index.js'; import { hybridSearch } from '../core/search/hybrid-search.js'; 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 { fileURLToPath } from 'url'; import { JobManager } from './analyze-job.js'; import { assertString, escapeRegExp, BadRequestError, createRouteLimiter } from './validation.js'; import { extractRepoName, getCloneDir, cloneOrPull } from './git-clone.js'; -import { createAnalyzeUploadHandler, requireLocalhostOrigin } from './analyze-upload.js'; +import { createAnalyzeUploadHandler } from './analyze-upload.js'; +import { requireLocalhostOrigin } from './middleware.js'; +import { createLaunchAnalysisWorker } from './analyze-launch.js'; import { UPLOAD_ROOT } from './upload-paths.js'; import { sweepStaleUploads } from './upload-sweep.js'; import { logger, flushLoggerSync } from '../core/logger.js'; @@ -764,147 +765,13 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => }; // Launch the analyze worker for an already-resolved repo directory. Shared by - // the JSON /api/analyze route and the multipart /api/analyze/upload route so - // the lock + fork + auto-retry + IPC machinery lives in one place. - const launchAnalysisWorker = ( - job: { id: string }, - targetPath: string, - opts: { - force?: boolean; - embeddings?: boolean; - dropEmbeddings?: boolean; - registryName?: string; - }, - ): void => { - // 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 ────────────────────────────── - 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) => { - logger.error({ err }, 'backend.init() failed after analyze:'); - 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() || ''; - logger.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: !!opts.force, - embeddings: !!opts.embeddings, - dropEmbeddings: !!opts.dropEmbeddings, - ...(opts.registryName ? { registryName: opts.registryName } : {}), - }, - }); - }; - - forkWorker(); - }; + // the JSON /api/analyze route and the multipart /api/analyze/upload route. + const launchAnalysisWorker = createLaunchAnalysisWorker({ + jobManager, + backend, + acquireRepoLock, + releaseRepoLock, + }); /** * Maximum time the hold-queue will wait for an active analysis job to complete. @@ -1691,6 +1558,7 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => createAnalyzeUploadHandler({ createJob: (params) => jobManager.createJob(params), launch: (job, targetPath, opts) => launchAnalysisWorker(job, targetPath, opts), + failJob: (jobId, error) => jobManager.updateJob(jobId, { status: 'failed', error }), }), ); diff --git a/gitnexus/src/server/git-clone.ts b/gitnexus/src/server/git-clone.ts index d92e9c28f..eae13b94e 100644 --- a/gitnexus/src/server/git-clone.ts +++ b/gitnexus/src/server/git-clone.ts @@ -20,7 +20,7 @@ const CLONE_ROOT = path.resolve(path.join(os.homedir(), '.gitnexus', 'repos')); // Rejecting anything else (including `..`, `/`, `\`, shell metacharacters) // guarantees getCloneDir(repoName) cannot escape CLONE_ROOT regardless of // how the caller derived repoName. -const REPO_NAME_PATTERN = /^[a-zA-Z0-9._-]+$/; +export const REPO_NAME_PATTERN = /^[a-zA-Z0-9._-]+$/; /** * Extract the repository name from a git URL (HTTPS or SSH). diff --git a/gitnexus/src/server/middleware.ts b/gitnexus/src/server/middleware.ts new file mode 100644 index 000000000..73d3edab7 --- /dev/null +++ b/gitnexus/src/server/middleware.ts @@ -0,0 +1,29 @@ +/** + * Shared Express route guards (alongside createRouteLimiter in validation.ts). + */ + +import type { Request, Response } from 'express'; + +/** + * Restrict a route to localhost browser origins. Non-browser requests (no + * Origin header, e.g. curl / the CLI) pass through. This closes cross-origin + * reach (the allow-listed public deploy + Private Network Access) to write + * routes without affecting read routes. + */ +export function requireLocalhostOrigin(req: Request, res: Response, next: () => void): void { + const origin = req.headers.origin; + if (origin === undefined) { + next(); + return; + } + try { + const hostname = new URL(origin).hostname; + if (hostname === 'localhost' || hostname === '127.0.0.1' || hostname === '::1') { + next(); + return; + } + } catch { + /* malformed origin → reject */ + } + res.status(403).json({ error: 'This endpoint is restricted to localhost origins' }); +} diff --git a/gitnexus/src/server/upload-ingest.ts b/gitnexus/src/server/upload-ingest.ts index b88f31c91..b726966fb 100644 --- a/gitnexus/src/server/upload-ingest.ts +++ b/gitnexus/src/server/upload-ingest.ts @@ -229,11 +229,20 @@ export async function ingestUpload( let dest: string; try { dest = resolveContainedDest(stageRoot, rel); + // A folder upload is exactly one top-level directory: every entry must + // have ≥2 segments and share the same first segment. This rejects a + // bare file at the root (which would make the promote target a file) + // and a multi-top manifest (which would silently drop all but the + // first folder). Validated here, before any job is created. + const segs = String(rel) + .split('/') + .filter((s) => s.length > 0); + const firstSeg = (segs[0] ?? '').normalize('NFC'); if (!topLevelName) { - // The on-disk top folder uses the NFC-normalized first segment - // (matching how each path segment is written); the promote step - // renames stageRoot/ into place. - topLevelName = (String(rel).split('/').filter(Boolean)[0] ?? '').normalize('NFC'); + topLevelName = firstSeg; + } + if (segs.length < 2 || firstSeg !== topLevelName) { + throw new BadRequestError('Upload must be a single folder of files'); } mkdirContained(stageRoot, dest, dirState); } catch (err) { diff --git a/gitnexus/src/server/upload-paths.ts b/gitnexus/src/server/upload-paths.ts index cae4b6d8e..ef957c85b 100644 --- a/gitnexus/src/server/upload-paths.ts +++ b/gitnexus/src/server/upload-paths.ts @@ -14,6 +14,7 @@ import path from 'path'; import os from 'os'; import { sanitizeRepoName } from '../storage/git.js'; +import { REPO_NAME_PATTERN } from './git-clone.js'; /** Root directory for all uploaded repositories. Targets must resolve inside this. */ export const UPLOAD_ROOT = path.resolve(path.join(os.homedir(), '.gitnexus', 'uploads')); @@ -21,11 +22,6 @@ export const UPLOAD_ROOT = path.resolve(path.join(os.homedir(), '.gitnexus', 'up /** Prefix for per-upload staging directories created under UPLOAD_ROOT. */ export const STAGING_PREFIX = '.staging-'; -// Filesystem-safe repo name: alphanumerics plus `. _ -`. Mirrors -// git-clone.ts REPO_NAME_PATTERN so getUploadDir(name) cannot escape -// UPLOAD_ROOT regardless of how the caller derived the name. -const UPLOAD_NAME_PATTERN = /^[a-zA-Z0-9._-]+$/; - /** * Get the upload target directory for a repo name. * @@ -41,7 +37,7 @@ export function getUploadDir(repoName: string): string { repoName === '..' || repoName === 'unknown' || repoName.startsWith('.') || - !UPLOAD_NAME_PATTERN.test(repoName) + !REPO_NAME_PATTERN.test(repoName) ) { throw new Error('Invalid repository name'); } diff --git a/gitnexus/src/server/upload-sweep.ts b/gitnexus/src/server/upload-sweep.ts index 3f2747dc7..46dcd4aed 100644 --- a/gitnexus/src/server/upload-sweep.ts +++ b/gitnexus/src/server/upload-sweep.ts @@ -35,13 +35,28 @@ export async function sweepStaleUploads(opts: SweepOptions = {}): Promise<{ remo } for (const entry of entries) { - if (!entry.isDirectory() || !entry.name.startsWith(STAGING_PREFIX)) continue; + if (!entry.isDirectory()) continue; const full = path.join(root, entry.name); try { const st = await fsp.stat(full); - if (now - st.mtimeMs > maxAgeMs) { + if (now - st.mtimeMs <= maxAgeMs) continue; // recent — keep + + if (entry.name.startsWith(STAGING_PREFIX)) { + // Transient staging dir orphaned by a crash — always removable. await fsp.rm(full, { recursive: true, force: true }).catch(() => {}); removed.push(full); + } else { + // Promoted upload dir. A successfully-analyzed (registered) repo always + // has a `.gitnexus` index inside it; a stale promoted dir WITHOUT one is + // an orphan from an analysis that failed before registering — remove it. + const hasIndex = await fsp + .access(path.join(full, '.gitnexus')) + .then(() => true) + .catch(() => false); + if (!hasIndex) { + await fsp.rm(full, { recursive: true, force: true }).catch(() => {}); + removed.push(full); + } } } catch { /* stat race — skip */ diff --git a/gitnexus/test/unit/api-analyze-upload.test.ts b/gitnexus/test/unit/api-analyze-upload.test.ts index 35786b4dd..af3bbfd23 100644 --- a/gitnexus/test/unit/api-analyze-upload.test.ts +++ b/gitnexus/test/unit/api-analyze-upload.test.ts @@ -3,10 +3,8 @@ import path from 'node:path'; import fs from 'node:fs/promises'; import { Readable } from 'node:stream'; import type { IncomingMessage } from 'node:http'; -import { - createAnalyzeUploadHandler, - requireLocalhostOrigin, -} from '../../src/server/analyze-upload.js'; +import { createAnalyzeUploadHandler } from '../../src/server/analyze-upload.js'; +import { requireLocalhostOrigin } from '../../src/server/middleware.js'; const BOUNDARY = '----gitnexusuploadtest'; @@ -81,7 +79,8 @@ describe('createAnalyzeUploadHandler', () => { const top = uniqueTop(); const createJob = vi.fn(() => ({ id: 'job-1', status: 'queued' })); const launch = vi.fn((_j, dir: string) => promoted.push(dir)); - const handler = createAnalyzeUploadHandler({ createJob, launch }); + const failJob = vi.fn(); + const handler = createAnalyzeUploadHandler({ createJob, launch, failJob }); const res = mockRes(); await handler( @@ -112,7 +111,8 @@ describe('createAnalyzeUploadHandler', () => { throw new Error('Analysis already in progress for another repository'); }); const launch = vi.fn((_j, dir: string) => promoted.push(dir)); - const handler = createAnalyzeUploadHandler({ createJob, launch }); + const failJob = vi.fn(); + const handler = createAnalyzeUploadHandler({ createJob, launch, failJob }); const res = mockRes(); await handler( @@ -133,7 +133,8 @@ describe('createAnalyzeUploadHandler', () => { it('rejects a traversal path in the manifest (400) without launching', async () => { const createJob = vi.fn(() => ({ id: 'j', status: 'queued' })); const launch = vi.fn(); - const handler = createAnalyzeUploadHandler({ createJob, launch }); + const failJob = vi.fn(); + const handler = createAnalyzeUploadHandler({ createJob, launch, failJob }); const res = mockRes(); await handler( @@ -152,7 +153,8 @@ describe('createAnalyzeUploadHandler', () => { it('rejects an un-nameable top folder (Windows-reserved → 400)', async () => { const createJob = vi.fn(() => ({ id: 'j', status: 'queued' })); const launch = vi.fn(); - const handler = createAnalyzeUploadHandler({ createJob, launch }); + const failJob = vi.fn(); + const handler = createAnalyzeUploadHandler({ createJob, launch, failJob }); const res = mockRes(); await handler( @@ -171,7 +173,8 @@ describe('createAnalyzeUploadHandler', () => { const top = uniqueTop(); const createJob = vi.fn(() => ({ id: 'job-x', status: 'queued' })); const launch = vi.fn((_j, dir: string) => promoted.push(dir)); - const handler = createAnalyzeUploadHandler({ createJob, launch }); + const failJob = vi.fn(); + const handler = createAnalyzeUploadHandler({ createJob, launch, failJob }); const res = mockRes(); await handler( @@ -188,6 +191,71 @@ describe('createAnalyzeUploadHandler', () => { await expect(fs.access(path.join(dir, '.gitnexus'))).rejects.toBeTruthy(); expect(await fs.readFile(path.join(dir, 'a.js'), 'utf8')).toBe('real'); }); + + it('rejects a single-segment manifest before creating a job (no slot taken)', async () => { + const createJob = vi.fn(() => ({ id: 'j', status: 'queued' })); + const launch = vi.fn(); + const failJob = vi.fn(); + const handler = createAnalyzeUploadHandler({ createJob, launch, failJob }); + + const res = mockRes(); + await handler( + mockReq([ + { name: 'manifest', value: JSON.stringify(['loosefile.js']) }, + { name: 'files', filename: 'blob', data: Buffer.from('x') }, + ]) as never, + res as never, + ); + + expect(res.statusCode).toBe(400); + expect(createJob).not.toHaveBeenCalled(); // slot never taken → no wedge + expect(launch).not.toHaveBeenCalled(); + }); + + it('rejects a multi-top-folder manifest (would silently drop folders)', async () => { + const createJob = vi.fn(() => ({ id: 'j', status: 'queued' })); + const launch = vi.fn(); + const failJob = vi.fn(); + const handler = createAnalyzeUploadHandler({ createJob, launch, failJob }); + + const res = mockRes(); + await handler( + mockReq([ + { name: 'manifest', value: JSON.stringify(['aaa/x.js', 'bbb/y.js']) }, + { name: 'files', filename: 'blob', data: Buffer.from('1') }, + { name: 'files', filename: 'blob', data: Buffer.from('2') }, + ]) as never, + res as never, + ); + + expect(res.statusCode).toBe(400); + expect(createJob).not.toHaveBeenCalled(); + }); + + it('releases the single slot (failJob) when a step fails after createJob', async () => { + const top = uniqueTop(); + const createJob = vi.fn(() => ({ id: 'job-fail', status: 'queued' })); + // launch throws AFTER createJob + promote — the slot must be released. + const launch = vi.fn((_j, dir: string) => { + promoted.push(dir); + throw new Error('worker fork blew up'); + }); + const failJob = vi.fn(); + const handler = createAnalyzeUploadHandler({ createJob, launch, failJob }); + + const res = mockRes(); + await handler( + mockReq([ + { name: 'manifest', value: JSON.stringify([`${top}/a.js`]) }, + { name: 'files', filename: 'blob', data: Buffer.from('x') }, + ]) as never, + res as never, + ); + + expect(createJob).toHaveBeenCalledOnce(); + expect(failJob).toHaveBeenCalledWith('job-fail', expect.any(String)); + expect(res.statusCode).toBe(500); + }); }); describe('requireLocalhostOrigin', () => { diff --git a/gitnexus/test/unit/rate-limit.test.ts b/gitnexus/test/unit/rate-limit.test.ts index 6a6d83a25..285a8f8c3 100644 --- a/gitnexus/test/unit/rate-limit.test.ts +++ b/gitnexus/test/unit/rate-limit.test.ts @@ -241,7 +241,9 @@ describe('production routes — rate-limit middleware wiring', () => { }); it('POST /api/analyze is wired with createRouteLimiter', () => { - expect(apiSource).toMatch(/app\.post\('\/api\/analyze',\s*createRouteLimiter\(/); + // Tolerate Prettier wrapping the registration across lines (it does once + // the route carries extra middleware like requireLocalhostOrigin). + expect(apiSource).toMatch(/app\.post\(\s*'\/api\/analyze',\s*createRouteLimiter\(/); }); it('POST /api/embed is wired with createRouteLimiter', () => { diff --git a/gitnexus/test/unit/upload-sweep.test.ts b/gitnexus/test/unit/upload-sweep.test.ts index 118f75ee0..a2667bbcb 100644 --- a/gitnexus/test/unit/upload-sweep.test.ts +++ b/gitnexus/test/unit/upload-sweep.test.ts @@ -36,6 +36,25 @@ describe('sweepStaleUploads', () => { await expect(fs.access(path.join(root, 'myrepo'))).resolves.toBeUndefined(); }); + it('removes a stale promoted dir without a .gitnexus index, keeps one with it', async () => { + const now = 2_000_000_000_000; + const old = new Date(now - 10 * 60 * 60 * 1000); + + // Orphan: a failed analysis that never wrote an index. + await fs.mkdir(path.join(root, 'orphan')); + await fs.utimes(path.join(root, 'orphan'), old, old); + + // Registered: stale but carries the .gitnexus index → must be kept. + await fs.mkdir(path.join(root, 'registered', '.gitnexus'), { recursive: true }); + await fs.utimes(path.join(root, 'registered'), old, old); + + const { removed } = await sweepStaleUploads({ root, now, maxAgeMs: 6 * 60 * 60 * 1000 }); + + expect(removed.some((r) => r.endsWith('orphan'))).toBe(true); + await expect(fs.access(path.join(root, 'orphan'))).rejects.toBeTruthy(); + await expect(fs.access(path.join(root, 'registered'))).resolves.toBeUndefined(); + }); + it('tolerates a missing root', async () => { const { removed } = await sweepStaleUploads({ root: path.join(root, 'does-not-exist') }); expect(removed).toEqual([]);