mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-09 03:17:54 +00:00
fix(watch): harden auto-sync lifecycle
This commit is contained in:
parent
9a2134770f
commit
5e92b1518f
17 changed files with 235 additions and 121 deletions
|
|
@ -135,7 +135,7 @@ export const en = {
|
|||
'help.command.watch.description':
|
||||
'Control scheduled repository clone/pull and analysis from GITNEXUS_HOME/watch_config.yml',
|
||||
'help.watch.details':
|
||||
'\nActions: init, start (default), restart, stop, status, reset\nConfiguration: GITNEXUS_HOME/watch_config.yml\nRuntime files: GITNEXUS_HOME/watch/watch.pid, watch.mutex, watch.owner.json, watch.status.json, auto-sync-state.json\nRecovery: mutexes with verified dead owners are reclaimed automatically; invalid or legacy mutexes fail closed and require manual removal after confirming no watch process is running.\nWrites: GITNEXUS_HOME/watch/project_commit_info.txt\nRemote URLs: only git@github.com:owner/repo.git, git@gitlab.com:group/repo.git, and git@gitee.com:owner/repo.git are allowed.\nRuns once immediately, then repeats on sync_interval_minutes.',
|
||||
'\nActions: init, start (default), restart, stop, status, reset\nConfiguration: GITNEXUS_HOME/watch_config.yml\nRuntime files: GITNEXUS_HOME/watch/watch.pid, watch.mutex, watch.owner.json, watch.status.json, auto-sync-state.json\nRecovery: mutexes with verified dead owners are reclaimed automatically; invalid or legacy mutexes fail closed and require manual removal after confirming no watch process is running.\nWrites: GITNEXUS_HOME/watch/project_commit_info.txt\nRemote URLs: only SSH URLs on github.com, gitlab.com, and gitee.com are allowed.\nRuns once immediately, then repeats on sync_interval_minutes.',
|
||||
'help.command.analyze.description': 'Index a repository (full analysis)',
|
||||
'help.command.index.description':
|
||||
'Register an existing .gitnexus/ folder into the global registry (no re-analysis needed)',
|
||||
|
|
|
|||
|
|
@ -136,7 +136,7 @@ export const zhCN = {
|
|||
'help.command.watch.description':
|
||||
'控制基于 GITNEXUS_HOME/watch_config.yml 的定时 clone/pull 和分析',
|
||||
'help.watch.details':
|
||||
'\n操作:init、start(默认)、restart、stop、status、reset\n配置:GITNEXUS_HOME/watch_config.yml\n运行时文件:GITNEXUS_HOME/watch/watch.pid、watch.mutex、watch.owner.json、watch.status.json、auto-sync-state.json\n恢复:已验证 owner 退出的 mutex 会自动回收;无效或旧版 mutex 会安全拒绝,确认没有 watch 进程运行后再手动删除。\n写入:GITNEXUS_HOME/watch/project_commit_info.txt\n远程地址:仅允许 git@github.com:owner/repo.git、git@gitlab.com:group/repo.git 和 git@gitee.com:owner/repo.git。\n启动后立即运行一次,之后按 sync_interval_minutes 重复。',
|
||||
'\n操作:init、start(默认)、restart、stop、status、reset\n配置:GITNEXUS_HOME/watch_config.yml\n运行时文件:GITNEXUS_HOME/watch/watch.pid、watch.mutex、watch.owner.json、watch.status.json、auto-sync-state.json\n恢复:已验证 owner 退出的 mutex 会自动回收;无效或旧版 mutex 会安全拒绝,确认没有 watch 进程运行后再手动删除。\n写入:GITNEXUS_HOME/watch/project_commit_info.txt\n远程地址:仅允许 github.com、gitlab.com 和 gitee.com 上的 SSH 地址。\n启动后立即运行一次,之后按 sync_interval_minutes 重复。',
|
||||
'help.command.analyze.description': '索引仓库(完整分析)',
|
||||
'help.command.index.description': '将现有 .gitnexus/ 文件夹注册到全局注册表(无需重新分析)',
|
||||
'help.command.serve.description': '启动供 Web UI 连接的本地 HTTP 服务器',
|
||||
|
|
|
|||
|
|
@ -67,6 +67,7 @@ export function createAutoSyncAnalysisRunner(
|
|||
|
||||
let terminalOutcome: WorkerMessage | undefined;
|
||||
let terminationGrace: ReturnType<typeof setTimeout> | undefined;
|
||||
let terminationReason: 'timeout' | 'cancelled' | undefined;
|
||||
let settled = false;
|
||||
const cleanup = () => {
|
||||
deps.clearTimeoutFn(timeout);
|
||||
|
|
@ -80,28 +81,37 @@ export function createAutoSyncAnalysisRunner(
|
|||
if (error) reject(error);
|
||||
else resolve(result!);
|
||||
};
|
||||
const timeout = deps.setTimeoutFn(() => {
|
||||
const requestTermination = (reason: 'timeout' | 'cancelled') => {
|
||||
if (settled || terminationReason) return;
|
||||
terminationReason = reason;
|
||||
deps.clearTimeoutFn(timeout);
|
||||
child.kill('SIGTERM');
|
||||
terminationGrace = deps.setTimeoutFn(() => {
|
||||
child.kill('SIGKILL');
|
||||
settle(new Error(`Analysis timed out after ${timeoutMs}ms.`));
|
||||
settle(
|
||||
new Error(
|
||||
reason === 'timeout' ? `Analysis timed out after ${timeoutMs}ms.` : 'Analysis cancelled.',
|
||||
),
|
||||
);
|
||||
}, TERMINATION_GRACE_MS);
|
||||
}, timeoutMs);
|
||||
const onAbort = () => {
|
||||
child.kill('SIGKILL');
|
||||
settle(new Error('Analysis cancelled.'));
|
||||
};
|
||||
const timeout = deps.setTimeoutFn(() => requestTermination('timeout'), timeoutMs);
|
||||
const onAbort = () => requestTermination('cancelled');
|
||||
signal?.addEventListener('abort', onAbort, { once: true });
|
||||
|
||||
child.on('message', (message: WorkerMessage) => {
|
||||
if (message.type !== 'progress') terminalOutcome ??= message;
|
||||
// Once timeout/cancellation requested shutdown, its reason owns the
|
||||
// result. A terminal IPC can already be queued behind SIGTERM.
|
||||
if (message.type === 'progress' || terminalOutcome || terminationReason) return;
|
||||
terminalOutcome = message;
|
||||
deps.clearTimeoutFn(timeout);
|
||||
});
|
||||
child.on('error', (error) => {
|
||||
settle(new Error(`Auto-sync analyze worker error: ${error.message}`));
|
||||
});
|
||||
child.on('exit', (code, childSignal) => {
|
||||
if (settled) return;
|
||||
if (terminationGrace) {
|
||||
if (terminationReason === 'timeout') {
|
||||
settle(
|
||||
new Error(
|
||||
`Analysis timed out after ${timeoutMs}ms and worker exited (${childSignal ?? code ?? 'unknown'}).`,
|
||||
|
|
@ -109,6 +119,10 @@ export function createAutoSyncAnalysisRunner(
|
|||
);
|
||||
return;
|
||||
}
|
||||
if (terminationReason === 'cancelled') {
|
||||
settle(new Error('Analysis cancelled.'));
|
||||
return;
|
||||
}
|
||||
if (terminalOutcome?.type === 'complete') {
|
||||
settle(undefined, { stats: terminalOutcome.result.stats });
|
||||
return;
|
||||
|
|
|
|||
|
|
@ -249,7 +249,7 @@ export function validateAutoSyncRemoteUrl(remoteUrl: string): void {
|
|||
const match = /^git@([^:\s/]+):([^\s]+)$/.exec(remoteUrl.trim());
|
||||
if (!match) {
|
||||
throw new Error(
|
||||
'must use git@github.com:owner/repo.git, git@gitlab.com:group/repo.git, or git@gitee.com:owner/repo.git',
|
||||
'must use an SSH URL on github.com, gitlab.com, or gitee.com',
|
||||
);
|
||||
}
|
||||
const host = match[1].toLowerCase();
|
||||
|
|
@ -270,6 +270,12 @@ export function validateAutoSyncBranchName(branch: string): void {
|
|||
if (branch.startsWith('-')) throw new Error('must not start with "-"');
|
||||
if (branch.includes('..')) throw new Error('must not contain ".."');
|
||||
if (branch.includes('`')) throw new Error('must not contain backticks');
|
||||
if (branch.endsWith('/') || branch.endsWith('.'))
|
||||
throw new Error('must not end with "/" or "."');
|
||||
if (branch.includes('//')) throw new Error('must not contain consecutive slashes');
|
||||
if (branch.includes('@{')) throw new Error('must not contain "@{"');
|
||||
if (branch.split('/').some((component) => component.startsWith('.') || component.endsWith('.lock')))
|
||||
throw new Error('must not contain hidden or .lock path components');
|
||||
}
|
||||
|
||||
export function parseDurationMs(value: unknown): number {
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
import fs from 'node:fs/promises';
|
||||
import { randomUUID } from 'node:crypto';
|
||||
import os from 'node:os';
|
||||
import path from 'node:path';
|
||||
import { getGlobalDir } from '../../storage/repo-manager.js';
|
||||
|
|
@ -115,7 +116,7 @@ export async function quarantineAutoSyncPartial(
|
|||
await fs.mkdir(quarantineRoot, { recursive: true, mode: 0o700 });
|
||||
const base = path.basename(targetDir);
|
||||
const stamp = new Date().toISOString().replace(/[:.]/g, '-');
|
||||
const destination = path.join(quarantineRoot, `auto-sync-${stamp}-${process.pid}-${base}`);
|
||||
const destination = path.join(quarantineRoot, `auto-sync-${stamp}-${process.pid}-${randomUUID()}-${base}`);
|
||||
try {
|
||||
await fs.rename(targetDir, destination);
|
||||
} catch (err: unknown) {
|
||||
|
|
|
|||
|
|
@ -301,8 +301,9 @@ export async function runAutoSyncOnce(
|
|||
|
||||
if (repoResult.project.groupName) {
|
||||
let groupMembershipOk = false;
|
||||
let membershipAdded = false;
|
||||
try {
|
||||
await deps.addRepoToGroup(
|
||||
membershipAdded = await deps.addRepoToGroup(
|
||||
repoResult.project,
|
||||
getAutoSyncRepoIdentity(repoResult.remoteUrl),
|
||||
getAutoSyncRepoIdentity(repoResult.remoteUrl),
|
||||
|
|
@ -314,7 +315,10 @@ export async function runAutoSyncOnce(
|
|||
`[auto-sync] Group update failed for ${repoResult.project.groupName}: ${(err as Error).message}`,
|
||||
);
|
||||
}
|
||||
if (groupMembershipOk && analyzeStatus === 'success') {
|
||||
if (
|
||||
groupMembershipOk &&
|
||||
(analyzeStatus === 'success' || (membershipAdded && analyzeStatus === 'skipped'))
|
||||
) {
|
||||
groupsToSync.add(repoResult.project.groupName);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -27,7 +27,6 @@ import { IndexLockTimeoutError } from '../storage/index-lock.js';
|
|||
export interface WorkerAnalysisDeps {
|
||||
runFullAnalysis: typeof import('../core/run-analyze.js').runFullAnalysis;
|
||||
assertAnalysisFinalized: typeof import('../storage/repo-manager.js').assertAnalysisFinalized;
|
||||
acquireAnalysisLock: (repoPath: string) => Promise<() => Promise<void>>;
|
||||
send: (msg: WorkerMessage) => void;
|
||||
/**
|
||||
* Claim the single terminal-outcome slot. Returns `true` for the first caller
|
||||
|
|
@ -52,38 +51,33 @@ export async function runWorkerAnalysis(
|
|||
): Promise<void> {
|
||||
let terminal: WorkerMessage;
|
||||
try {
|
||||
const releaseAnalysisLock = await deps.acquireAnalysisLock(repoPath);
|
||||
try {
|
||||
const bootstrapArgs: [] | [AnalyzerRunnerIdentity] = runnerIdentityAtBootstrap
|
||||
? [runnerIdentityAtBootstrap]
|
||||
: [];
|
||||
const result = await deps.runFullAnalysis(
|
||||
repoPath,
|
||||
// This worker force-exits right after reporting, so skip the native close
|
||||
// (it can double-free in LadybugDB's ClientContext destructor after --pdg
|
||||
// writes); flushWAL still persists the index, process.exit reclaims handles.
|
||||
{ ...options, skipNativeCloseOnExit: true },
|
||||
{
|
||||
onProgress: (phase, percent, message) =>
|
||||
deps.send({ type: 'progress', phase, percent, message }),
|
||||
onLog: (message) => deps.send({ type: 'progress', phase: 'log', percent: -1, message }),
|
||||
},
|
||||
...bootstrapArgs,
|
||||
);
|
||||
// P2 (#2264): a half-finalized repo — meta.json written but the global
|
||||
// registry entry missing (e.g. a prior collision-aborted run, or a wiped
|
||||
// registry) — must NOT be reported as a successful analysis. Mirror the CLI's
|
||||
// assertAnalysisFinalized guard so the worker surfaces it as an error instead
|
||||
// of a false `complete` that leaves the repo invisible to list_repos.
|
||||
await deps.assertAnalysisFinalized(repoPath);
|
||||
const bootstrapArgs: [] | [AnalyzerRunnerIdentity] = runnerIdentityAtBootstrap
|
||||
? [runnerIdentityAtBootstrap]
|
||||
: [];
|
||||
const result = await deps.runFullAnalysis(
|
||||
repoPath,
|
||||
// This worker force-exits right after reporting, so skip the native close
|
||||
// (it can double-free in LadybugDB's ClientContext destructor after --pdg
|
||||
// writes); flushWAL still persists the index, process.exit reclaims handles.
|
||||
{ ...options, skipNativeCloseOnExit: true },
|
||||
{
|
||||
onProgress: (phase, percent, message) =>
|
||||
deps.send({ type: 'progress', phase, percent, message }),
|
||||
onLog: (message) => deps.send({ type: 'progress', phase: 'log', percent: -1, message }),
|
||||
},
|
||||
...bootstrapArgs,
|
||||
);
|
||||
// P2 (#2264): a half-finalized repo — meta.json written but the global
|
||||
// registry entry missing (e.g. a prior collision-aborted run, or a wiped
|
||||
// registry) — must NOT be reported as a successful analysis. Mirror the CLI's
|
||||
// assertAnalysisFinalized guard so the worker surfaces it as an error instead
|
||||
// of a false `complete` that leaves the repo invisible to list_repos.
|
||||
await deps.assertAnalysisFinalized(repoPath);
|
||||
|
||||
// Send a JSON-safe projection, NOT the raw result: the IPC channel is
|
||||
// default-JSON serialization and `result.pipelineResult` carries the live
|
||||
// KnowledgeGraph. See analyze-worker-ipc.ts.
|
||||
terminal = { type: 'complete', result: projectAnalyzeResultForIpc(result) };
|
||||
} finally {
|
||||
await releaseAnalysisLock();
|
||||
}
|
||||
// Send a JSON-safe projection, NOT the raw result: the IPC channel is
|
||||
// default-JSON serialization and `result.pipelineResult` carries the live
|
||||
// KnowledgeGraph. See analyze-worker-ipc.ts.
|
||||
terminal = { type: 'complete', result: projectAnalyzeResultForIpc(result) };
|
||||
} catch (err: unknown) {
|
||||
// Report the failure to the parent over IPC (the parent surfaces the message).
|
||||
const message = err instanceof Error ? err.message : 'Analysis failed';
|
||||
|
|
|
|||
|
|
@ -11,8 +11,6 @@
|
|||
* Child -> Parent: { type: 'error', message: string }
|
||||
*/
|
||||
|
||||
import path from 'path';
|
||||
import { createHash } from 'crypto';
|
||||
import type { StartMessage, WorkerMessage } from './analyze-worker-protocol.js';
|
||||
import { runWorkerAnalysis, createTerminalClaim } from './analyze-worker-core.js';
|
||||
type BoundedCheckpointBeforeExit =
|
||||
|
|
@ -105,13 +103,12 @@ process.on('message', async (msg: StartMessage) => {
|
|||
const prepared = await identityModule.captureAnalyzerIdentityBeforeLoad(
|
||||
import.meta.url,
|
||||
async () => {
|
||||
const [analysisModule, repoManager, shutdownHelpers, fileLock] = await Promise.all([
|
||||
const [analysisModule, repoManager, shutdownHelpers] = await Promise.all([
|
||||
import('../core/run-analyze.js'),
|
||||
import('../storage/repo-manager.js'),
|
||||
import('../core/lbug/shutdown-helpers.js'),
|
||||
import('../storage/file-lock.js'),
|
||||
]);
|
||||
return { analysisModule, repoManager, shutdownHelpers, fileLock };
|
||||
return { analysisModule, repoManager, shutdownHelpers };
|
||||
},
|
||||
);
|
||||
boundedCheckpointBeforeExit = prepared.loaded.shutdownHelpers.boundedCheckpointBeforeExit;
|
||||
|
|
@ -125,18 +122,6 @@ process.on('message', async (msg: StartMessage) => {
|
|||
{
|
||||
runFullAnalysis: prepared.loaded.analysisModule.runFullAnalysis,
|
||||
assertAnalysisFinalized: prepared.loaded.repoManager.assertAnalysisFinalized,
|
||||
acquireAnalysisLock: (repoPath) => {
|
||||
const repoKey = createHash('sha256')
|
||||
.update(prepared.loaded.repoManager.canonicalizePath(repoPath))
|
||||
.digest('hex');
|
||||
return prepared.loaded.fileLock.acquireFileLock(
|
||||
path.join(
|
||||
prepared.loaded.repoManager.getGlobalDir(),
|
||||
'locks',
|
||||
`analyze-${repoKey}.lock`,
|
||||
),
|
||||
);
|
||||
},
|
||||
send,
|
||||
claimTerminal,
|
||||
},
|
||||
|
|
|
|||
|
|
@ -523,6 +523,8 @@ export async function cloneOrPull(
|
|||
await assertDirectoryOwnerAndPermissions(cloneRoot);
|
||||
}
|
||||
await assertNoSymlinkPath(cloneRoot, safeTarget, Boolean(options?.allowedCloneRoot));
|
||||
await fs.mkdir(path.dirname(safeTarget), { recursive: true });
|
||||
await assertNoSymlinkPath(cloneRoot, safeTarget, Boolean(options?.allowedCloneRoot));
|
||||
await assertPreRealpathContainment(cloneRoot, safeTarget);
|
||||
|
||||
const exists = await fs.access(path.join(safeTarget, '.git')).then(
|
||||
|
|
@ -600,9 +602,6 @@ export async function cloneOrPull(
|
|||
if (targetExists && (await fs.readdir(safeTarget)).length > 0) {
|
||||
throw new Error(`Clone target already exists but is not a git repository: ${safeTarget}`);
|
||||
}
|
||||
await fs.mkdir(path.dirname(safeTarget), { recursive: true });
|
||||
await assertNoSymlinkPath(cloneRoot, safeTarget, Boolean(options?.allowedCloneRoot));
|
||||
await assertPreRealpathContainment(cloneRoot, safeTarget);
|
||||
onProgress?.({ phase: 'cloning', message: `Cloning ${url}...` });
|
||||
try {
|
||||
const runGitImpl = options?.runGitForTest ?? runGit;
|
||||
|
|
|
|||
|
|
@ -61,6 +61,7 @@ export async function acquireFileLock(
|
|||
if (
|
||||
await reclaimStaleLock(
|
||||
resolvedPath,
|
||||
owner,
|
||||
options.isProcessAlive ?? isProcessAlive,
|
||||
options.readProcessStartTime ?? readProcessStartTime,
|
||||
)
|
||||
|
|
@ -81,14 +82,21 @@ export async function acquireFileLock(
|
|||
|
||||
async function reclaimStaleLock(
|
||||
lockPath: string,
|
||||
guardOwner: FileLockOwner,
|
||||
ownerIsAlive: (pid: number) => boolean,
|
||||
getProcessStartTime: (pid: number) => string | undefined,
|
||||
): Promise<boolean> {
|
||||
const reclaimGuardPath = `${lockPath}.reclaim`;
|
||||
let releaseReclaimGuard: () => Promise<void>;
|
||||
try {
|
||||
await fs.mkdir(reclaimGuardPath);
|
||||
releaseReclaimGuard = await acquireFileLock(reclaimGuardPath, {
|
||||
pid: guardOwner.pid,
|
||||
processStartTime: guardOwner.processStartTime,
|
||||
isProcessAlive: ownerIsAlive,
|
||||
readProcessStartTime: getProcessStartTime,
|
||||
});
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code === 'EEXIST') return false;
|
||||
if (error instanceof FileLockBusyError) return false;
|
||||
throw error;
|
||||
}
|
||||
|
||||
|
|
@ -103,7 +111,7 @@ async function reclaimStaleLock(
|
|||
await fs.rm(lockPath, { force: true });
|
||||
return true;
|
||||
} finally {
|
||||
await fs.rmdir(reclaimGuardPath);
|
||||
await releaseReclaimGuard();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -18,7 +18,6 @@ import fs from 'fs/promises';
|
|||
import { realpathSync } from 'fs';
|
||||
import path from 'path';
|
||||
import os from 'os';
|
||||
import { randomBytes } from 'crypto';
|
||||
import { getInferredRepoName, resolveRepoIdentityRoot, stripUrlCredentials } from './git.js';
|
||||
import { stripWindowsLongPathPrefix } from '../lib/utils.js';
|
||||
import { writeFileAtomic } from './fs-atomic.js';
|
||||
|
|
@ -642,17 +641,10 @@ export const readRegistry = async (): Promise<RegistryEntry[]> => {
|
|||
const writeRegistry = async (entries: RegistryEntry[]): Promise<void> => {
|
||||
const dir = getGlobalDir();
|
||||
await fs.mkdir(dir, { recursive: true });
|
||||
// Atomic tmp+rename (mirrors saveMeta): a crash mid-write can never leave a
|
||||
// truncated/half-written registry.json that the next load would treat as
|
||||
// empty and silently drop every registered repo (#2106 R9).
|
||||
const target = getGlobalRegistryPath();
|
||||
const tmp = `${target}.${process.pid}.${randomBytes(8).toString('hex')}.tmp`;
|
||||
try {
|
||||
await fs.writeFile(tmp, JSON.stringify(sanitizeEntries(entries), null, 2), 'utf-8');
|
||||
await fs.rename(tmp, target);
|
||||
} finally {
|
||||
await fs.unlink(tmp).catch(() => {});
|
||||
}
|
||||
await writeFileAtomic(
|
||||
getGlobalRegistryPath(),
|
||||
JSON.stringify(sanitizeEntries(entries), null, 2),
|
||||
);
|
||||
};
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -32,7 +32,6 @@ const baseResult: AnalyzeResult = {
|
|||
|
||||
const okRun: WorkerAnalysisDeps['runFullAnalysis'] = vi.fn(async () => baseResult);
|
||||
const okFinalize: WorkerAnalysisDeps['assertAnalysisFinalized'] = vi.fn(async () => undefined);
|
||||
const okLock: WorkerAnalysisDeps['acquireAnalysisLock'] = vi.fn(async () => async () => undefined);
|
||||
const alwaysClaim: WorkerAnalysisDeps['claimTerminal'] = () => true;
|
||||
|
||||
describe('runWorkerAnalysis — finalize guard (#2264 P2)', () => {
|
||||
|
|
@ -50,7 +49,6 @@ describe('runWorkerAnalysis — finalize guard (#2264 P2)', () => {
|
|||
{
|
||||
runFullAnalysis: okRun,
|
||||
assertAnalysisFinalized,
|
||||
acquireAnalysisLock: okLock,
|
||||
send,
|
||||
claimTerminal: alwaysClaim,
|
||||
},
|
||||
|
|
@ -72,7 +70,6 @@ describe('runWorkerAnalysis — finalize guard (#2264 P2)', () => {
|
|||
{
|
||||
runFullAnalysis: okRun,
|
||||
assertAnalysisFinalized: okFinalize,
|
||||
acquireAnalysisLock: okLock,
|
||||
send,
|
||||
claimTerminal: alwaysClaim,
|
||||
},
|
||||
|
|
@ -93,7 +90,6 @@ describe('runWorkerAnalysis — finalize guard (#2264 P2)', () => {
|
|||
{
|
||||
runFullAnalysis: run,
|
||||
assertAnalysisFinalized: okFinalize,
|
||||
acquireAnalysisLock: okLock,
|
||||
send,
|
||||
claimTerminal: alwaysClaim,
|
||||
},
|
||||
|
|
@ -103,31 +99,6 @@ describe('runWorkerAnalysis — finalize guard (#2264 P2)', () => {
|
|||
expect(run.mock.calls[0]?.[3]).toBe(receipt);
|
||||
});
|
||||
|
||||
it('does not enter analysis when another worker holds the repo lock', async () => {
|
||||
const send = vi.fn<(msg: WorkerMessage) => void>();
|
||||
const run = vi.fn<WorkerAnalysisDeps['runFullAnalysis']>(async () => baseResult);
|
||||
|
||||
await runWorkerAnalysis(
|
||||
'/repo',
|
||||
{},
|
||||
{
|
||||
runFullAnalysis: run,
|
||||
assertAnalysisFinalized: okFinalize,
|
||||
acquireAnalysisLock: vi.fn(async () => {
|
||||
throw new Error('Lock is already held for /repo');
|
||||
}),
|
||||
send,
|
||||
claimTerminal: alwaysClaim,
|
||||
},
|
||||
);
|
||||
|
||||
expect(run).not.toHaveBeenCalled();
|
||||
expect(send).toHaveBeenCalledWith({
|
||||
type: 'error',
|
||||
message: 'Lock is already held for /repo',
|
||||
});
|
||||
});
|
||||
|
||||
it('reports error when finalization passes but the analysis itself throws', async () => {
|
||||
const send = vi.fn<(msg: WorkerMessage) => void>();
|
||||
const failingRun: WorkerAnalysisDeps['runFullAnalysis'] = vi.fn(async () => {
|
||||
|
|
@ -143,7 +114,6 @@ describe('runWorkerAnalysis — finalize guard (#2264 P2)', () => {
|
|||
{
|
||||
runFullAnalysis: failingRun,
|
||||
assertAnalysisFinalized: finalize,
|
||||
acquireAnalysisLock: okLock,
|
||||
send,
|
||||
claimTerminal: alwaysClaim,
|
||||
},
|
||||
|
|
@ -196,7 +166,6 @@ describe('runWorkerAnalysis — terminal-claim coordination (#2264 P3)', () => {
|
|||
{
|
||||
runFullAnalysis: okRun,
|
||||
assertAnalysisFinalized: okFinalize,
|
||||
acquireAnalysisLock: okLock,
|
||||
send,
|
||||
claimTerminal: alreadyClaimed,
|
||||
},
|
||||
|
|
|
|||
|
|
@ -80,7 +80,67 @@ describe('auto-sync analysis worker', () => {
|
|||
await expect(result).rejects.toThrow('Analysis timed out after 50ms');
|
||||
});
|
||||
|
||||
it('kills an active worker immediately when watch is stopped', async () => {
|
||||
it('keeps the timeout outcome when complete arrives after termination begins', async () => {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
send: vi.fn(),
|
||||
kill: vi.fn(),
|
||||
});
|
||||
const timers: Array<() => void> = [];
|
||||
const run = createAutoSyncAnalysisRunner({
|
||||
forkWorker: vi.fn(() => child as any),
|
||||
setTimeoutFn: vi.fn((callback: () => void) => {
|
||||
timers.push(callback);
|
||||
return timers.length as any;
|
||||
}) as any,
|
||||
clearTimeoutFn: vi.fn() as any,
|
||||
});
|
||||
|
||||
const result = run('/tmp/repo', { branch: 'main' }, 50);
|
||||
timers[0]();
|
||||
child.emit('message', { type: 'complete', result: { stats: { files: 3 } } });
|
||||
child.emit('exit', 0, null);
|
||||
|
||||
await expect(result).rejects.toThrow('Analysis timed out after 50ms');
|
||||
});
|
||||
|
||||
it('clears the analysis deadline after complete before the worker exits', async () => {
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
send: vi.fn(),
|
||||
kill: vi.fn(),
|
||||
});
|
||||
const run = createAutoSyncAnalysisRunner({ forkWorker: vi.fn(() => child as any) });
|
||||
|
||||
const result = run('/tmp/repo', { branch: 'main' }, 50);
|
||||
child.emit('message', { type: 'complete', result: { stats: { files: 3 } } });
|
||||
await vi.advanceTimersByTimeAsync(50);
|
||||
|
||||
expect(child.kill).not.toHaveBeenCalled();
|
||||
child.emit('exit', 0, null);
|
||||
await expect(result).resolves.toEqual({ stats: { files: 3 } });
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
||||
it('keeps the cancellation outcome when complete arrives after abort', async () => {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
send: vi.fn(),
|
||||
kill: vi.fn(),
|
||||
});
|
||||
const run = createAutoSyncAnalysisRunner({ forkWorker: vi.fn(() => child as any) });
|
||||
const controller = new AbortController();
|
||||
|
||||
const result = run('/tmp/repo', { branch: 'main' }, 50, controller.signal);
|
||||
controller.abort();
|
||||
child.emit('message', { type: 'complete', result: { stats: { files: 3 } } });
|
||||
child.emit('exit', 0, null);
|
||||
|
||||
await expect(result).rejects.toThrow('Analysis cancelled');
|
||||
});
|
||||
|
||||
it('asks an active worker to stop gracefully when watch is stopped', async () => {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
send: vi.fn(),
|
||||
kill: vi.fn(),
|
||||
|
|
@ -93,8 +153,9 @@ describe('auto-sync analysis worker', () => {
|
|||
const result = run('/tmp/repo', { branch: 'main' }, 50, controller.signal);
|
||||
controller.abort();
|
||||
|
||||
expect(child.kill).toHaveBeenCalledWith('SIGKILL');
|
||||
child.emit('exit', null, 'SIGKILL');
|
||||
expect(child.kill).toHaveBeenCalledWith('SIGTERM');
|
||||
expect(child.kill).not.toHaveBeenCalledWith('SIGKILL');
|
||||
child.emit('exit', null, 'SIGTERM');
|
||||
await expect(result).rejects.toThrow('Analysis cancelled');
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -295,7 +295,7 @@ describe('auto-sync runner', () => {
|
|||
})),
|
||||
saveState: vi.fn(async () => {}),
|
||||
writeCommitInfo: vi.fn(async () => {}),
|
||||
addRepoToGroup: vi.fn(async () => false),
|
||||
addRepoToGroup: vi.fn(async () => true),
|
||||
syncGroupByName: vi.fn(async () => {}),
|
||||
getAvailableMemoryGB: vi.fn(() => 8),
|
||||
});
|
||||
|
|
@ -308,7 +308,7 @@ describe('auto-sync runner', () => {
|
|||
expect(result.analyzed).toBe(0);
|
||||
expect(result.skippedAnalysis).toBe(1);
|
||||
expect(deps.runAnalysis).not.toHaveBeenCalled();
|
||||
expect(deps.syncGroupByName).not.toHaveBeenCalled();
|
||||
expect(deps.syncGroupByName).toHaveBeenCalledWith('back_end');
|
||||
});
|
||||
|
||||
it('uses remote identity under local_path as the clone target', async () => {
|
||||
|
|
@ -1115,13 +1115,14 @@ describe('auto-sync starter', () => {
|
|||
return timer;
|
||||
}) as unknown as typeof setInterval;
|
||||
const stderr = { write: vi.fn() };
|
||||
let releaseRun: (() => void) | undefined;
|
||||
const releaseRuns: Array<() => void> = [];
|
||||
const runOnce = vi.fn(
|
||||
() =>
|
||||
new Promise<any>((resolve) => {
|
||||
releaseRun = () => resolve({ synced: 0, analyzed: 0, skippedAnalysis: 0, failed: 0 });
|
||||
releaseRuns.push(() => resolve({ synced: 0, analyzed: 0, skippedAnalysis: 0, failed: 0 }));
|
||||
}),
|
||||
);
|
||||
let handle: Awaited<ReturnType<typeof startAutoSyncWatch>> | undefined;
|
||||
|
||||
try {
|
||||
process.env.GITNEXUS_HOME = tempDir;
|
||||
|
|
@ -1137,7 +1138,7 @@ describe('auto-sync starter', () => {
|
|||
].join('\n'),
|
||||
);
|
||||
|
||||
await startAutoSyncWatch({ setIntervalFn, runOnce, stderr });
|
||||
handle = await startAutoSyncWatch({ setIntervalFn, runOnce, stderr });
|
||||
scheduled?.();
|
||||
|
||||
expect(runOnce).toHaveBeenCalledTimes(1);
|
||||
|
|
@ -1145,12 +1146,17 @@ describe('auto-sync starter', () => {
|
|||
'[auto-sync] Previous run is still active; skipping overlapping run.\n',
|
||||
);
|
||||
|
||||
releaseRun?.();
|
||||
releaseRuns.shift()?.();
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
scheduled?.();
|
||||
|
||||
expect(runOnce).toHaveBeenCalledTimes(2);
|
||||
releaseRuns.shift()?.();
|
||||
await handle?.stop();
|
||||
handle = undefined;
|
||||
} finally {
|
||||
releaseRuns.splice(0).forEach((release) => release());
|
||||
await handle?.stop();
|
||||
if (previousHome === undefined) delete process.env.GITNEXUS_HOME;
|
||||
else process.env.GITNEXUS_HOME = previousHome;
|
||||
await fs.rm(tempDir, { recursive: true, force: true });
|
||||
|
|
|
|||
|
|
@ -351,6 +351,24 @@ describe('auto-sync', () => {
|
|||
await expect(fs.access(target)).rejects.toThrow();
|
||||
});
|
||||
|
||||
it('gives concurrent partial clone quarantines unique destinations', async () => {
|
||||
const quarantineRoot = path.join(gitnexusHome, 'watch', 'quarantine');
|
||||
const first = path.join(tempDir, 'one', 'partial-repo');
|
||||
const second = path.join(tempDir, 'two', 'partial-repo');
|
||||
await Promise.all([fs.mkdir(first, { recursive: true }), fs.mkdir(second, { recursive: true })]);
|
||||
|
||||
const [firstDestination, secondDestination] = await Promise.all([
|
||||
quarantineAutoSyncPartial(first, quarantineRoot),
|
||||
quarantineAutoSyncPartial(second, quarantineRoot),
|
||||
]);
|
||||
|
||||
expect(firstDestination).not.toBe(secondDestination);
|
||||
await expect(fs.access(firstDestination)).resolves.toBeUndefined();
|
||||
await expect(fs.access(secondDestination)).resolves.toBeUndefined();
|
||||
await expect(fs.access(first)).rejects.toThrow();
|
||||
await expect(fs.access(second)).rejects.toThrow();
|
||||
});
|
||||
|
||||
it('rejects group-writable configured clone roots', async () => {
|
||||
if (process.platform === 'win32') return;
|
||||
const root = path.join(tempDir, 'group-writable-repos');
|
||||
|
|
@ -388,10 +406,17 @@ describe('auto-sync', () => {
|
|||
|
||||
it('rejects unsafe auto-sync branch names', () => {
|
||||
expect(() => validateAutoSyncBranchName('feature/good-branch')).not.toThrow();
|
||||
expect(() => validateAutoSyncBranchName('foo./bar')).not.toThrow();
|
||||
expect(() => validateAutoSyncBranchName('-upload-pack=evil')).toThrow('must not start');
|
||||
expect(() => validateAutoSyncBranchName('feature bad')).toThrow('whitespace');
|
||||
expect(() => validateAutoSyncBranchName('feature..bad')).toThrow('must not contain ".."');
|
||||
expect(() => validateAutoSyncBranchName('bad:ref')).toThrow('not allowed');
|
||||
expect(() => validateAutoSyncBranchName('feature.')).toThrow('must not end');
|
||||
expect(() => validateAutoSyncBranchName('feature/')).toThrow('must not end');
|
||||
expect(() => validateAutoSyncBranchName('feature//branch')).toThrow('consecutive');
|
||||
expect(() => validateAutoSyncBranchName('feature@{x')).toThrow('must not contain "@{"');
|
||||
expect(() => validateAutoSyncBranchName('.hidden')).toThrow('hidden');
|
||||
expect(() => validateAutoSyncBranchName('foo/bar.lock')).toThrow('hidden or .lock');
|
||||
});
|
||||
|
||||
it('extracts safe repository names from remote URLs', () => {
|
||||
|
|
@ -411,6 +436,7 @@ describe('auto-sync', () => {
|
|||
});
|
||||
|
||||
it('allows only github, gitlab, and gitee SSH SCP remote URLs', () => {
|
||||
expect(() => validateAutoSyncRemoteUrl('git@github.com:owner/repo')).not.toThrow();
|
||||
expect(() => validateAutoSyncRemoteUrl('git@github.com:im-fan/multica.git')).not.toThrow();
|
||||
expect(() => validateAutoSyncRemoteUrl('git@gitlab.com:group/subgroup/repo.git')).not.toThrow();
|
||||
expect(() =>
|
||||
|
|
|
|||
|
|
@ -65,6 +65,15 @@ describe('file lock', () => {
|
|||
const lockPath = await tempLockPath();
|
||||
await acquireFileLock(lockPath, { pid: 111, processStartTime: 'old-start' });
|
||||
|
||||
await expect(
|
||||
acquireFileLock(lockPath, {
|
||||
pid: 222,
|
||||
processStartTime: 'next-start',
|
||||
isProcessAlive: () => true,
|
||||
readProcessStartTime: () => 'old-start',
|
||||
}),
|
||||
).rejects.toBeInstanceOf(FileLockBusyError);
|
||||
|
||||
const nextRelease = await acquireFileLock(lockPath, {
|
||||
pid: 222,
|
||||
processStartTime: 'next-start',
|
||||
|
|
@ -83,6 +92,23 @@ describe('file lock', () => {
|
|||
await expect(fs.access(lockPath)).resolves.toBeUndefined();
|
||||
});
|
||||
|
||||
it('recovers when a stale reclaim guard was left by a crashed contender', async () => {
|
||||
const lockPath = await tempLockPath();
|
||||
await acquireFileLock(lockPath, { pid: 999, processStartTime: 'abandoned' });
|
||||
await acquireFileLock(`${lockPath}.reclaim`, {
|
||||
pid: 998,
|
||||
processStartTime: 'abandoned-reclaimer',
|
||||
});
|
||||
|
||||
const release = await acquireFileLock(lockPath, {
|
||||
pid: 1000,
|
||||
processStartTime: 'next',
|
||||
isProcessAlive: () => false,
|
||||
});
|
||||
|
||||
await release();
|
||||
});
|
||||
|
||||
it('waits for the current holder when retries are configured', async () => {
|
||||
const lockPath = await tempLockPath();
|
||||
const release = await acquireFileLock(lockPath);
|
||||
|
|
|
|||
|
|
@ -637,6 +637,29 @@ describe('git-clone', () => {
|
|||
}
|
||||
});
|
||||
|
||||
it('creates missing nested parents before checking controlled clone containment', async () => {
|
||||
const root = await mkControlledRoot('gitnexus-controlled-root-');
|
||||
const target = path.join(root, 'github.com', 'owner', 'repo');
|
||||
const runGitForTest = vi.fn(async () => {
|
||||
await fs.mkdir(path.join(target, '.git'), { recursive: true });
|
||||
return '';
|
||||
});
|
||||
try {
|
||||
await expect(
|
||||
cloneOrPull('git@github.com:owner/repo', target, undefined, {
|
||||
allowedCloneRoot: root,
|
||||
expectedRepoName: 'repo',
|
||||
allowAutoSyncSsh: true,
|
||||
runGitForTest,
|
||||
}),
|
||||
).resolves.toBe(target);
|
||||
|
||||
expect(runGitForTest).toHaveBeenCalledOnce();
|
||||
} finally {
|
||||
await fs.rm(root, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
|
||||
it('allows auto-sync SSH SCP clone URLs with a per-repo timeout', async () => {
|
||||
const root = await mkControlledRoot('gitnexus-controlled-root-');
|
||||
const target = path.join(root, 'repo');
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue