fix(review): resolve tri-review findings on the upload flow

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) <noreply@anthropic.com>
This commit is contained in:
Gergo Magyar 2026-06-09 09:54:50 +00:00
parent cc8b2cfb45
commit 7c4e5700ff
15 changed files with 415 additions and 226 deletions

View file

@ -188,9 +188,11 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp
const sseControllerRef = useRef<AbortController | null>(null);
const completeTimerRef = useRef<ReturnType<typeof setTimeout> | null>(null);
const folderInputRef = useRef<HTMLInputElement>(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 `<folder>/<rest>`) 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')}
</button>
{uploading && (
<div role="status" data-testid="upload-progress" className="space-y-1">
<div role="status" aria-busy="true" data-testid="upload-progress" className="space-y-1">
<div className="h-1.5 w-full overflow-hidden rounded-full bg-elevated">
<div className="h-full w-1/3 animate-pulse rounded-full bg-accent" />
</div>
@ -520,7 +533,7 @@ export const RepoAnalyzer = ({ variant, onComplete, onCancel }: RepoAnalyzerProp
{uploadSummary && !uploading && phase !== 'error' && (
<p className="text-xs text-text-muted" data-testid="upload-summary">
{t('onboarding:repoAnalyzer.upload.selected', {
count: uploadSummary.count,
fileCount: uploadSummary.count,
dropped: uploadSummary.dropped,
})}
</p>

View file

@ -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."
}
}

View file

@ -65,7 +65,7 @@
"upload": {
"button": "上传文件夹",
"uploading": "上传中…",
"selected": "已准备 {{count}} 个文件(已跳过 {{dropped}} 个:.git、node_modules、构建产物)",
"selected": "已准备 {{fileCount}} 个文件(已跳过 {{dropped}} 个:.git、node_modules、构建产物)",
"empty": "该文件夹中未找到可分析的文件。"
}
}

View file

@ -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<unknown> };
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();
};
}

View file

@ -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<AnalyzeJob, 'id' | 'status'>;
/** 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<string> {
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<string> {
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<void> {
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/<topLevelName>. 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/<topLevelName>. 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' });
}

View file

@ -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);

View file

@ -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 }),
}),
);

View file

@ -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).

View file

@ -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' });
}

View file

@ -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/<topLevelName> 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) {

View file

@ -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');
}

View file

@ -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 */

View file

@ -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', () => {

View file

@ -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', () => {

View file

@ -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([]);