From 1b05145d69a568800aa8855ccfdc1d7889c3457b Mon Sep 17 00:00:00 2001 From: SalvageA Date: Mon, 28 Sep 2026 15:56:02 +0900 Subject: [PATCH 1/2] fix(concurrency): make lifecycle and durable state transitions ownership-safe Salvages the unique work of draft PRs #766 and #767 into one change. - fileLock: add withFileLockSync (Atomics wait) for sync RMW callers, with the ownership-safe stale reclaim and ENOTEMPTY-tolerant unlink; lock waits use the real timer so a test that fakes setTimeout cannot freeze lock hand-off. - pairPipeline: replace the per-instance abortSignal/stuckDetector with an AsyncLocalStorage RunControl, so concurrent run() calls on one instance cannot observe or overwrite each other's cancellation state. - processRegistry: taskId+spawnedAt identity ownership on activity, close and health-check, plus a SIGKILL escalation timer that is unref'd, cancelled on child close, and cleared by a force kill. - store: tryClaimTaskAdmission claims a task in one locked read-modify-write; stale-lock reclaim tolerates ENOTEMPTY; main's updateTaskLinearState lock invariant is preserved. - runnerState: cross-process locks on the rejection, decomposition, pace, project-selection and pipeline-history RMW paths; main's withRunnerStateLock and reserveDailyCreations are kept, and the daily reset runs outside the decomposition lock because withFileLockSync is not reentrant. - decisionEngine: admission claims on the auto-execute paths and an in_progress filter; file-locked engine state and backlog appends. - oauthStore, dev, logRotation, codex, reembed, taskParser, gitInfo, locale: the locked or scoped variant of the same invariant. Dropped: telemetry.ts (main has withTelemetryLock), workflow.ts (owned by another salvage), prProcessor.ts (its state lock cannot fit the repo's 1500-line gate, which main's copy already sits on), task_state_model.py, agt3420-probe.txt, and package.json/lock version churn. --- src/adapters/processRegistry.test.ts | 9 +- src/adapters/processRegistry.ts | 128 ++++++---- .../pairPipeline.cancelIsolation.test.ts | 94 +++++++ src/agents/pairPipeline.ts | 40 +-- src/agents/pairPipelineTypes.ts | 4 + src/auth/oauthStore.ts | 161 ++++++++---- .../runnerState.concurrency.test.ts | 92 +++++++ .../runnerState.rejection.fixture.ts | 15 ++ src/automation/runnerState.ts | 239 ++++++++++-------- src/knowledge/gitInfo.test.ts | 84 +++++- src/knowledge/gitInfo.ts | 24 +- src/locale/index.ts | 53 +++- src/locale/locale.scope.test.ts | 38 +++ src/memory/codex.ts | 8 + src/memory/memoryCore.ts | 26 +- src/memory/memoryWriteRetry.test.ts | 48 ++-- src/memory/reembed.ts | 126 ++++----- .../decisionEngine.admission.test.ts | 43 ++++ .../decisionEngine.coverage.test.ts | 69 ++++- src/orchestration/decisionEngine.ts | 130 +++++++--- src/orchestration/taskParser.coverage.test.ts | 23 +- src/orchestration/taskParser.ts | 61 ++++- src/support/dev.ts | 43 +++- src/support/fileLock.test.ts | 44 +++- src/support/fileLock.ts | 150 ++++++++++- src/support/logRotation.test.ts | 44 +++- src/support/logRotation.ts | 25 +- src/taskState/store.test.ts | 30 +++ src/taskState/store.ts | 49 +++- 29 files changed, 1528 insertions(+), 372 deletions(-) create mode 100644 src/agents/pairPipeline.cancelIsolation.test.ts create mode 100644 src/automation/runnerState.concurrency.test.ts create mode 100644 src/automation/runnerState.rejection.fixture.ts create mode 100644 src/locale/locale.scope.test.ts create mode 100644 src/orchestration/decisionEngine.admission.test.ts diff --git a/src/adapters/processRegistry.test.ts b/src/adapters/processRegistry.test.ts index 8f12cda5..2aa9d82d 100644 --- a/src/adapters/processRegistry.test.ts +++ b/src/adapters/processRegistry.test.ts @@ -21,7 +21,14 @@ function fakeChild(pid: number): EventEmitter & { pid: number } { function register(pid: number): EventEmitter & { pid: number } { const proc = fakeChild(pid); registerProcess( - { pid, adapter: 'codex', role: 'worker', command: 'codex', spawnedAt: Date.now() } as never, + { + pid, + taskId: 't', + stage: 'worker', + projectPath: '/tmp', + spawnedAt: Date.now(), + lastActivityAt: Date.now(), + }, proc as never, ); return proc; diff --git a/src/adapters/processRegistry.ts b/src/adapters/processRegistry.ts index e74259dd..a3aa8704 100644 --- a/src/adapters/processRegistry.ts +++ b/src/adapters/processRegistry.ts @@ -50,13 +50,13 @@ export function registerProcess(info: ProcessInfo, proc: ChildProcess): void { }, }); - // Track activity from stdout/stderr - const updateActivity = () => { + // Activity tracking + const updateActivity = (): void => { const entry = registry.get(info.pid); - if (!entry) return; - entry.lastActivityAt = Date.now(); + // Ownership check: only update if this exact identity still owns the PID + if (!entry || entry.taskId !== info.taskId || entry.spawnedAt !== info.spawnedAt) return; - // Throttled broadcast + entry.lastActivityAt = Date.now(); const lastBroadcast = activityThrottle.get(info.pid) ?? 0; if (Date.now() - lastBroadcast >= ACTIVITY_THROTTLE_MS) { activityThrottle.set(info.pid, Date.now()); @@ -67,13 +67,26 @@ export function registerProcess(info: ProcessInfo, proc: ChildProcess): void { proc.stdout?.on('data', updateActivity); proc.stderr?.on('data', updateActivity); - // Cleanup on close + // Cleanup on close — verify ownership before acting so a reused PID does not + // cause this handler to clean up a newer process that happens to share the same + // numeric PID (PID reuse is real on busy systems). proc.on('close', (code, signal) => { const entry = registry.get(info.pid); - const durationMs = entry ? Date.now() - entry.spawnedAt : 0; + // Ownership check: the entry must match this exact process identity (taskId + + // spawnedAt), not just the PID. A reused PID would have a different spawnedAt + // or taskId. + if (!entry || entry.taskId !== info.taskId || entry.spawnedAt !== info.spawnedAt) { + return; + } + const durationMs = Date.now() - entry.spawnedAt; registry.delete(info.pid); processHandles.delete(info.pid); activityThrottle.delete(info.pid); + const pending = escalationTimers.get(info.pid); + if (pending) { + clearTimeout(pending); + escalationTimers.delete(info.pid); + } broadcastEvent({ type: 'process:exit', @@ -105,61 +118,90 @@ export function getAllProcesses(): ProcessInfo[] { return Array.from(registry.values()); } +const KILL_ESCALATE_MS = 5_000; + +/** Pending SIGKILL escalations per pid, so a force kill can cancel one. */ +const escalationTimers = new Map(); + /** - * Kill a tracked process. Sends SIGTERM first, then SIGKILL after 5s. + * Kill a tracked process by PID. + * Signals through the retained ChildProcess handle — never by raw PID — so a + * recycled numeric PID cannot be escalated into after the original child exits. + * Soft kill schedules SIGKILL escalation after {@link KILL_ESCALATE_MS} only + * while this exact registry identity still owns the handle; the caller does not + * wait for that timer (so a recycled PID in those seconds is never signalled). */ export async function killProcess(pid: number, force = false): Promise { - const entry = registry.get(pid); - if (!entry) return false; + const info = registry.get(pid); + const proc = processHandles.get(pid); + if (!info || !proc) return false; + + const stillOurs = (): boolean => { + const current = registry.get(pid); + return ( + !!current + && current.taskId === info.taskId + && current.spawnedAt === info.spawnedAt + && processHandles.get(pid) === proc + ); + }; + + if (!stillOurs()) return false; try { - const proc = processHandles.get(pid); - if (proc) { - if (force) { - terminateCliProcessTree(proc); - } else { - signalCliProcessTree(proc, 'SIGTERM'); - // Escalate the same process group. The npm wrapper may exit before its - // native CLI/MCP descendants, so checking only proc.exitCode is unsafe. - const escalation = setTimeout(() => terminateCliProcessTree(proc), 5000); - escalation.unref(); + if (force) { + if (!stillOurs()) return false; + // A soft kill's pending escalation must not outlive this call: it would + // signal again 5s later, and stillOurs() alone cannot see that. + if (escalationTimers.get(pid)) { + clearTimeout(escalationTimers.get(pid)); + escalationTimers.delete(pid); } + terminateCliProcessTree(proc); return true; } - // No handle means we cannot prove this PID is still the child we spawned. - // Signalling it anyway is the ownership hazard this registry exists to - // remove: PIDs are recycled, and a delayed escalation in particular would - // land on whatever unrelated process inherited the number. Registration and - // handle are written and cleared together, so this is a defensive branch — - // report it rather than guessing. - console.warn(`[ProcessRegistry] No process handle for pid ${pid}; refusing to signal by PID`); - registry.delete(pid); - activityThrottle.delete(pid); - return false; + + signalCliProcessTree(proc, 'SIGTERM'); + const timer = setTimeout(() => { + escalationTimers.delete(pid); + if (stillOurs()) terminateCliProcessTree(proc); + }, KILL_ESCALATE_MS); + // Unref: a pending escalation must not keep the event loop alive. + timer.unref(); + escalationTimers.set(pid, timer); + proc.once('close', () => { + clearTimeout(timer); + escalationTimers.delete(pid); + }); + return true; } catch { - // Process already gone - registry.delete(pid); - processHandles.delete(pid); - activityThrottle.delete(pid); return false; } } /** - * Start periodic health checker that verifies processes are still alive. - * Removes stale entries from the registry. + * Start periodic health checker that removes stale entries + * where the process handle has already exited. + * + * Identity-safe: re-reads the registry entry before acting so a reused PID + * does not cause this handler to clean up a newer process. */ export function startHealthChecker(intervalMs = 30000): void { if (healthCheckTimer) return; + healthCheckTimer = setInterval(() => { + const now = Date.now(); for (const [pid, info] of registry) { - try { - process.kill(pid, 0); // No-op signal — just checks if alive - } catch { - // Process is dead but wasn't cleaned up - const durationMs = Date.now() - info.spawnedAt; + const proc = processHandles.get(pid); + if (!proc) { + // Ownership check: re-read the entry to confirm it hasn't been replaced + // by a reused PID with a different identity (taskId + spawnedAt). + const current = registry.get(pid); + if (!current || current.taskId !== info.taskId || current.spawnedAt !== info.spawnedAt) { + continue; + } + const durationMs = now - info.spawnedAt; registry.delete(pid); - processHandles.delete(pid); activityThrottle.delete(pid); broadcastEvent({ type: 'process:exit', @@ -188,4 +230,4 @@ export function stopHealthChecker(): void { clearInterval(healthCheckTimer); healthCheckTimer = null; } -} +} \ No newline at end of file diff --git a/src/agents/pairPipeline.cancelIsolation.test.ts b/src/agents/pairPipeline.cancelIsolation.test.ts new file mode 100644 index 00000000..c070642c --- /dev/null +++ b/src/agents/pairPipeline.cancelIsolation.test.ts @@ -0,0 +1,94 @@ +// Regression: the per-run abort/stuck controls live in AsyncLocalStorage, so a +// run() call can no longer observe or overwrite another run's cancellation +// state. Before the fix both were instance fields, and the second run's already +// aborted signal made the FIRST run cancel itself at its next stage boundary. +import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'; +import type { TaskItem } from '../orchestration/decisionEngine.js'; + +const runWorker = vi.fn(); +const broadcastEvent = vi.fn(); +const getDefaultModel = vi.fn(); +const hasRepoSnapshot = vi.fn(); +const scanAndCache = vi.fn(); +const analyzeIssue = vi.fn(); +const recallRepoKnowledge = vi.fn(); + +vi.mock('./worker.js', async () => { + const actual = await vi.importActual('./worker.js'); + return { ...actual, runWorker }; +}); +vi.mock('./tester.js', async () => { + const actual = await vi.importActual('./tester.js'); + return { ...actual, runTester: vi.fn() }; +}); +vi.mock('../knowledge/index.js', () => ({ hasRepoSnapshot, scanAndCache, analyzeIssue, recallRepoKnowledge })); +vi.mock('../core/eventHub.js', () => ({ broadcastEvent })); +vi.mock('../adapters/index.js', async () => { + const actual = await vi.importActual('../adapters/index.js'); + return { ...actual, getAdapter: () => ({ getDefaultModel }) }; +}); + +describe('PairPipeline concurrent run isolation', () => { + beforeEach(() => { + vi.spyOn(console, 'log').mockImplementation(() => undefined); + vi.spyOn(console, 'warn').mockImplementation(() => undefined); + vi.spyOn(console, 'error').mockImplementation(() => undefined); + hasRepoSnapshot.mockReturnValue(true); + scanAndCache.mockResolvedValue(undefined); + analyzeIssue.mockResolvedValue(undefined); + recallRepoKnowledge.mockResolvedValue([]); + getDefaultModel.mockResolvedValue('model'); + runWorker.mockResolvedValue({ + success: true, + summary: 'done', + filesChanged: ['src/a.ts'], + commands: [], + output: '', + confidencePercent: 95, + }); + }); + afterEach(() => { vi.restoreAllMocks(); vi.clearAllMocks(); }); + + function task(id: string): TaskItem { + return { id, source: 'linear', title: id, description: 'd', priority: 1, createdAt: Date.now(), estimatedMinutes: 30 }; + } + + it('does not cancel a run just because a concurrent run on the same instance is aborted', async () => { + const { PairPipeline } = await import('./pairPipeline.js'); + const pipeline = new PairPipeline({ + stages: ['worker'], + maxIterations: 2, + roles: { worker: { enabled: true, timeoutMs: 0 } }, + }); + + let releaseFirstCall: () => void = () => {}; + const firstCallGate = new Promise((resolve) => { releaseFirstCall = resolve; }); + let workerCalls = 0; + runWorker.mockImplementation(async () => { + workerCalls++; + if (workerCalls === 1) await firstCallGate; + return { + success: true, + summary: 'done', + filesChanged: ['src/a.ts'], + commands: [], + output: '', + confidencePercent: 95, + }; + }); + + const first = pipeline.run(task('TASK-A'), process.cwd()); + + // A second run on the SAME instance that is already cancelled. It must not + // install its aborted signal anywhere the first run can read it. + const controller = new AbortController(); + controller.abort(); + const second = await pipeline.run(task('TASK-B'), process.cwd(), { signal: controller.signal }); + expect(second.finalStatus).toBe('cancelled'); + + releaseFirstCall(); + const firstResult = await first; + expect(firstResult.finalStatus).not.toBe('cancelled'); + expect(firstResult.success).toBe(true); + }); +}); diff --git a/src/agents/pairPipeline.ts b/src/agents/pairPipeline.ts index 3ea7ca07..d1b0e98b 100644 --- a/src/agents/pairPipeline.ts +++ b/src/agents/pairPipeline.ts @@ -3,6 +3,7 @@ // Worker → Reviewer → Tester → Documenter pipeline // ============================================ import { EventEmitter } from 'node:events'; +import { AsyncLocalStorage } from 'node:async_hooks'; import { taskAttributionKey, taskEventKey, type TaskItem } from '../orchestration/decisionEngine.js'; import { rejectedWorkerPaths } from '../support/rejectedWorkerPaths.js'; import { scratchNotesSection, workerScratchpadRunId } from './workerScratchpad.js'; @@ -107,20 +108,26 @@ import { reviewWorkerBlocker } from './workerBlockerReview.js'; * rather than an inline assignment so the loop's later reads are not narrowed * to `undefined`. */ +type RunControl = { signal?: AbortSignal; stuck: StuckDetector }; + function dropStaleTestVerdict(context: { testerResult?: TesterResult }): void { context.testerResult = undefined; } export class PairPipeline extends EventEmitter { private config: PipelineConfig; + /** Fallback for callers outside run(); each run() installs a fresh control via AsyncLocalStorage. */ private stuckDetector: StuckDetector; - /** Set per run() — aborts the pipeline + in-flight adapter call on cancel/disable. */ - private abortSignal?: AbortSignal; + /** Per-run abort + stuck controls — never share across concurrent run() calls. */ + private static runControl = new AsyncLocalStorage(); /** Cache of adapter default models (heavy: OAuth + live catalog) keyed by adapter name. (INT-2393) */ private defaultModelCache = new Map>(); /** Throw if this run has been cancelled. Called at iteration/stage boundaries. */ private throwIfAborted(): void { - if (this.abortSignal?.aborted) throw new PipelineCancelledError(); + if (PairPipeline.runControl.getStore()?.signal?.aborted) throw new PipelineCancelledError(); + } + private getStuck(): StuckDetector { + return PairPipeline.runControl.getStore()?.stuck ?? this.stuckDetector; } constructor(config: PipelineConfig) { @@ -153,13 +160,13 @@ export class PairPipeline extends EventEmitter { * On failure at any stage, returns to Worker (up to maxIterations) */ async run(task: TaskItem, projectPath: string, opts?: { signal?: AbortSignal }): Promise { + // A fresh detector per run: two run() calls on one instance must not share + // the per-run identity that runControl exists to keep separate. + const stuckDetector = createStuckDetector({ sameErrorRepeat: 3, revisionLoop: 4 }); + return PairPipeline.runControl.run({ signal: opts?.signal, stuck: stuckDetector }, async () => { const startTime = Date.now(); const stages: StageResult[] = []; const maxIterations = this.config.maxIterations ?? 3; - this.abortSignal = opts?.signal; - - // Reset stuck detector (new pipeline run) - this.stuckDetector.reset(); // Ensure repo graph snapshot exists (first-time scan if needed) if (!hasRepoSnapshot(projectPath)) { @@ -193,6 +200,8 @@ export class PairPipeline extends EventEmitter { currentIteration: 0, taskPrefix, reflection: createReflectionState(), + abortSignal: opts?.signal, + stuckDetector, }; try { if (this.config.verify?.enabled) try { @@ -247,7 +256,7 @@ export class PairPipeline extends EventEmitter { } catch (error) { // Cancellation (project disable / manual stop) is not a failure — surface it // as 'cancelled' so the scheduler doesn't count it failed or trigger a retry. - const cancelled = error instanceof PipelineCancelledError || !!this.abortSignal?.aborted; + const cancelled = error instanceof PipelineCancelledError || !!PairPipeline.runControl.getStore()?.signal?.aborted; // A 429/usage-limit propagates up here from any stage (worker/reviewer/…). // Surface it as its own finalStatus so the runner pauses until quota resets // instead of counting a failure and spamming Linear comments. (INT-1906) @@ -294,6 +303,7 @@ export class PairPipeline extends EventEmitter { }, }; } + }); } /** * Worker에 주입할 코드 컨텍스트 수집 @@ -308,7 +318,7 @@ export class PairPipeline extends EventEmitter { private async runPostSuccessStage(stage: PipelineStage, context: PipelineContext, stages: StageResult[]): Promise { try { stages.push(await this.runStage(stage, context)); } catch (err) { - if (err instanceof PipelineCancelledError || this.abortSignal?.aborted) throw err; + if (err instanceof PipelineCancelledError || PairPipeline.runControl.getStore()?.signal?.aborted) throw err; safeConsole.warn(`[${context.taskPrefix}] ${stage} skipped (non-blocking failure): ${err instanceof Error ? err.message : String(err)}`); } } @@ -489,7 +499,7 @@ export class PairPipeline extends EventEmitter { onLog, processContext: { taskId: taskAttributionKey(context.task), stage: 'worker' }, workerContext, - signal: this.abortSignal, + signal: PairPipeline.runControl.getStore()?.signal, instructionCapsule: this.config.instructionCapsule, mcpTools: this.config.roleMcpTools?.worker, adapterRouting: this.config.adapterRouting, @@ -557,7 +567,7 @@ export class PairPipeline extends EventEmitter { // scaffolded task was getting LESS review — exactly the wrong incentive. The // completion-criteria hard gate is the real check now. const reviewerOptions = await buildReviewerStageOptions({ - config: this.config, context, prefix, overrides, abortSignal: this.abortSignal, + config: this.config, context, prefix, overrides, abortSignal: PairPipeline.runControl.getStore()?.signal, }); safeConsole.log(`[${prefix}] Running full review...`); @@ -845,7 +855,7 @@ export class PairPipeline extends EventEmitter { await captureBeforeIteration(context); // Stuck detection check (before iteration starts) - const stuckCheck = this.stuckDetector.check(); + const stuckCheck = this.getStuck().check(); if (stuckCheck.isStuck) { context.stuckReason = stuckCheck.reason; safeConsole.error(`[${context.taskPrefix}] STUCK DETECTED: ${stuckCheck.reason}`); @@ -902,7 +912,7 @@ export class PairPipeline extends EventEmitter { stages.push(workerResult); // Record Worker result in stuck detector - this.stuckDetector.addEntry({ + this.getStuck().addEntry({ stage: 'worker', success: workerResult.success, output: (workerResult.result as WorkerResult).summary, @@ -954,7 +964,7 @@ export class PairPipeline extends EventEmitter { if (detail?.startsWith('worker-scope:')) context.repeatedScopeRejection = detail; // A no-edit stop with a stated reason is a claim about the task: put it // to the reviewer once instead of retrying blind (AGT-4535). - const blocker = await reviewWorkerBlocker(this.config, context, failedWorker, this.abortSignal); + const blocker = await reviewWorkerBlocker(this.config, context, failedWorker, PairPipeline.runControl.getStore()?.signal); if (blocker?.confirmed) { agentPair.updateSessionStatus(context.session.id, 'waiting_on_operator'); this.emit('halt', { confidence: failedWorker.confidencePercent ?? 0, haltReason: blocker.reason, sessionId: context.session.id, iteration: context.currentIteration, context }); @@ -1269,7 +1279,7 @@ export class PairPipeline extends EventEmitter { const decision = (reviewerResult.result as ReviewResult).decision; // Record Reviewer result in stuck detector - this.stuckDetector.addEntry({ + this.getStuck().addEntry({ stage: 'reviewer', success: reviewerResult.success, decision: decision, diff --git a/src/agents/pairPipelineTypes.ts b/src/agents/pairPipelineTypes.ts index 247df69f..71178396 100644 --- a/src/agents/pairPipelineTypes.ts +++ b/src/agents/pairPipelineTypes.ts @@ -209,6 +209,10 @@ export interface PipelineContext { verifiedBlocker?: string; /** The reviewer's verdict on a blocker claim, kept for cost accounting. */ blockerReview?: ReviewResult; + /** Abort signal for this run — set per run() so concurrent runs do not share it. */ + abortSignal?: AbortSignal; + /** Per-run stuck detector — never share across concurrent run() calls. */ + stuckDetector?: import('../support/stuckDetector.js').StuckDetector; } export type PipelineEventType = 'stage:start' | 'stage:complete' | 'stage:fail' | 'iteration:start' | 'iteration:complete' | 'iteration:fail' | 'pipeline:complete' | 'pipeline:fail' | 'fanout:gate' | 'halt'; diff --git a/src/auth/oauthStore.ts b/src/auth/oauthStore.ts index eda793be..ee867b7e 100644 --- a/src/auth/oauthStore.ts +++ b/src/auth/oauthStore.ts @@ -7,6 +7,7 @@ import { readFileSync, existsSync, renameSync } from 'node:fs'; import { join } from 'node:path'; import { homedir } from 'node:os'; import { atomicWriteFileSync } from '../support/atomicFile.js'; +import { withFileLock, withFileLockSync } from '../support/fileLock.js'; import { parseTokenResponse } from './tokenResponse.js'; // Types @@ -61,6 +62,7 @@ function isAuthProfile(value: unknown): value is AuthProfile { // Constants const STORE_PATH = join(homedir(), '.openswarm', 'auth-profiles.json'); +const STORE_LOCK = `${STORE_PATH}.lock`; const REFRESH_BUFFER_MS = 5 * 60 * 1000; // 5분 전에 갱신 const OPENAI_TOKEN_ENDPOINT = 'https://auth.openai.com/oauth/token'; const LINEAR_TOKEN_ENDPOINT = 'https://api.linear.app/oauth/token'; @@ -136,14 +138,20 @@ export class AuthProfileStore { * token fails the next refresh with invalid_grant. Only the keys this * instance actually touched are applied on top of the current file. * - * This narrows the race rather than removing it: two writers rotating the - * *same* key concurrently still read-then-write, so the later one can land on - * a snapshot taken before the earlier write. Closing that needs a lock around - * read-modify-write, which is a larger change than this fix. Different keys — - * the common CLI-beside-daemon case, and the one that used to lose unrelated - * providers' credentials — are now safe. + * Same-key concurrent writers are serialized by {@link STORE_LOCK} around + * this merge+write (and around refresh in {@link ensureValidToken}). */ save(): void { + withFileLockSync(STORE_LOCK, () => { + this.saveUnlocked(); + }); + } + + /** + * Merge touched keys onto disk and write. Caller must hold {@link STORE_LOCK} + * (or accept a race). Used by {@link save} and by refresh under an async lock. + */ + saveUnlocked(): void { const onDisk = existsSync(STORE_PATH) ? this.readProfilesQuietly() : {}; for (const key of this.touched) { const profile = this.data.profiles[key]; @@ -155,6 +163,34 @@ export class AuthProfileStore { atomicWriteFileSync(STORE_PATH, `${JSON.stringify(this.data, null, 2)}\n`, 0o600); } + /** + * Re-read `key` from disk into memory. Used under lock before deciding to refresh + * so a concurrent refresh that already completed is visible. + */ + reloadProfileFromDisk(key: string): AuthProfile | null { + const onDisk = this.readProfilesQuietly()[key]; + if (onDisk) { + this.data.profiles[key] = onDisk; + return onDisk; + } + return this.data.profiles[key] ?? null; + } + + /** + * Like {@link setProfile}, but does not acquire {@link STORE_LOCK}. + * Caller must already hold the lock (e.g. refresh inside {@link withFileLock}). + */ + setProfileUnlocked(key: string, profile: AuthProfile): void { + if (!isAuthProfile(profile)) { + throw new Error( + `Refusing to store an invalid auth profile for "${key}" — it would make the store unloadable.`, + ); + } + this.data.profiles[key] = profile; + this.touched.add(key); + this.saveUnlocked(); + } + /** Current on-disk profiles, or an empty map if the file is unreadable. */ private readProfilesQuietly(): Record { try { @@ -249,6 +285,11 @@ export class TokenRefreshError extends Error { /** * 유효한 access token 반환. 만료 임박 시 자동 refresh. + * + * Concurrent refreshes of the same key serialize on {@link STORE_LOCK}: the + * second waiter re-reads disk under the lock and returns the already-rotated + * access token instead of refreshing again (which would invalidate the first + * writer's refresh_token). */ export async function ensureValidToken(store: AuthProfileStore, profileKey: string): Promise { const profile = store.getProfile(profileKey); @@ -261,59 +302,73 @@ export async function ensureValidToken(store: AuthProfileStore, profileKey: stri return profile.access; } - const now = Date.now(); - if (now < profile.expires - REFRESH_BUFFER_MS) { + if (Date.now() < profile.expires - REFRESH_BUFFER_MS) { return profile.access; } - // Token 갱신 - console.log(`[Auth] Refreshing token for ${profileKey}...`); - - const body = new URLSearchParams({ - grant_type: 'refresh_token', - refresh_token: profile.refresh, - client_id: profile.clientId, - }); + return withFileLock(STORE_LOCK, async () => { + // Re-check under lock — another process may have refreshed already. + const fresh = store.reloadProfileFromDisk(profileKey); + if (!fresh) { + throw new Error(`Auth profile "${profileKey}" not found. Run: openswarm auth login --provider gpt`); + } + if (fresh.type === 'apiKey') { + return fresh.access; + } + if (Date.now() < fresh.expires - REFRESH_BUFFER_MS) { + return fresh.access; + } - const endpoint = TOKEN_ENDPOINTS[profile.provider]; - if (!endpoint) { - throw new Error(`Unknown OAuth provider "${profile.provider}" for auth profile "${profileKey}". Re-run auth login for this provider.`); - } + console.log(`[Auth] Refreshing token for ${profileKey}...`); - const res = await fetch(endpoint, { - method: 'POST', - headers: { 'Content-Type': 'application/x-www-form-urlencoded' }, - body: body.toString(), - }); - - if (!res.ok) { - const errText = await res.text().catch(() => ''); - const reauth = profile.provider === 'linear' ? 'linear' : 'gpt'; - throw new TokenRefreshError( - `Token refresh failed (${res.status}): ${errText.slice(0, 200)}. Run: openswarm auth login --provider ${reauth}`, - res.status, - ); - } + const body = new URLSearchParams({ + grant_type: 'refresh_token', + refresh_token: fresh.refresh, + client_id: fresh.clientId, + }); - // Validated, not cast. A 200 carrying an error body — or a proxy's HTML — - // would otherwise put `undefined` into access and `NaN` into expires, and - // that profile gets written to disk like any other, where it fails the - // whole-file schema check on the next load and takes every other provider's - // credentials down with it. refresh_token stays optional here: providers may - // legitimately keep the existing one on a refresh. - const tokens = parseTokenResponse(await res.json(), { - provider: profile.provider, - requireRefreshToken: false, - }); - - profile.access = tokens.accessToken; - if (tokens.refreshToken) { - profile.refresh = tokens.refreshToken; - } - profile.expires = Date.now() + tokens.expiresIn * 1000; + const endpoint = TOKEN_ENDPOINTS[fresh.provider]; + if (!endpoint) { + throw new Error(`Unknown OAuth provider "${fresh.provider}" for auth profile "${profileKey}". Re-run auth login for this provider.`); + } - store.setProfile(profileKey, profile); - console.log(`[Auth] Token refreshed successfully.`); + const res = await fetch(endpoint, { + method: 'POST', + headers: { 'Content-Type': 'application/x-www-form-urlencoded' }, + body: body.toString(), + }); + + if (!res.ok) { + const errText = await res.text().catch(() => ''); + const reauth = fresh.provider === 'linear' ? 'linear' : 'gpt'; + throw new TokenRefreshError( + `Token refresh failed (${res.status}): ${errText.slice(0, 200)}. Run: openswarm auth login --provider ${reauth}`, + res.status, + ); + } - return profile.access; + // Validated, not cast. A 200 carrying an error body — or a proxy's HTML — + // would otherwise put `undefined` into access and `NaN` into expires, and + // that profile gets written to disk like any other, where it fails the + // whole-file schema check on the next load and takes every other provider's + // credentials down with it. refresh_token stays optional here: providers may + // legitimately keep the existing one on a refresh. + const tokens = parseTokenResponse(await res.json(), { + provider: fresh.provider, + requireRefreshToken: false, + }); + + const updated: AuthProfile = { + ...fresh, + access: tokens.accessToken, + refresh: tokens.refreshToken ?? fresh.refresh, + expires: Date.now() + tokens.expiresIn * 1000, + }; + + // save() would re-enter STORE_LOCK and deadlock; persist under the held lock. + store.setProfileUnlocked(profileKey, updated); + console.log(`[Auth] Token refreshed successfully.`); + + return updated.access; + }, { timeoutMs: 30_000 }); } diff --git a/src/automation/runnerState.concurrency.test.ts b/src/automation/runnerState.concurrency.test.ts new file mode 100644 index 00000000..822d841e --- /dev/null +++ b/src/automation/runnerState.concurrency.test.ts @@ -0,0 +1,92 @@ +// Concurrent runner-state RMW: without a cross-process lock, two writers reload +// the same snapshot and the later write drops the earlier increment. (AGT-3420) + +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { mkdirSync, mkdtempSync, readFileSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { spawn } from 'node:child_process'; +import { fileURLToPath } from 'node:url'; + +let home: string; +let rejectionFile: string; + +async function loadRunnerState() { + vi.resetModules(); + vi.stubEnv('HOME', home); + vi.stubEnv('USERPROFILE', home); + vi.stubEnv('OPENSWARM_RUNNER_REJECTION_STATE_FILE', rejectionFile); + vi.stubEnv('OPENSWARM_RUNNER_PIPELINE_HISTORY_FILE', join(home, '.claude', 'openswarm-pipeline-history.json')); + vi.stubEnv('OPENSWARM_RUNNER_TASK_STATE_FILE', join(home, '.claude', 'openswarm-task-state.json')); + vi.stubEnv('OPENSWARM_RUNNER_DECOMPOSITION_STATE_FILE', join(home, '.claude', 'openswarm-decomposition-state.json')); + return await import('./runnerState.js'); +} + +beforeEach(() => { + home = mkdtempSync(join(tmpdir(), 'openswarm-runner-conc-')); + mkdirSync(join(home, '.claude'), { recursive: true }); + rejectionFile = join(home, '.claude', 'openswarm-rejection-state.json'); +}); + +afterEach(() => { + vi.unstubAllEnvs(); + rmSync(home, { recursive: true, force: true }); +}); + +describe('concurrent runner-state updates', () => { + it('keeps every rejection increment when child processes race', async () => { + const { incrementRejection } = await loadRunnerState(); + expect(incrementRejection('ISSUE-1', 'seed')).toBe(1); + + const fixture = fileURLToPath(new URL('./runnerState.rejection.fixture.ts', import.meta.url)); + const run = (reason: string) => new Promise((resolve, reject) => { + const child = spawn( + process.execPath, + ['--import', 'tsx', fixture, rejectionFile, 'ISSUE-1', reason], + { + stdio: ['ignore', 'pipe', 'pipe'], + env: { + ...process.env, + HOME: home, + USERPROFILE: home, + OPENSWARM_RUNNER_REJECTION_STATE_FILE: rejectionFile, + }, + }, + ); + let stdout = ''; + let stderr = ''; + child.stdout.on('data', (chunk) => { stdout += String(chunk); }); + child.stderr.on('data', (chunk) => { stderr += String(chunk); }); + child.on('error', reject); + child.on('exit', (code) => { + if (code !== 0) reject(new Error(stderr || `child exited ${code}`)); + else resolve(Number(stdout.trim())); + }); + }); + + await Promise.all(Array.from({ length: 8 }, (_, i) => run(`reason-${i}`))); + + const raw = JSON.parse(readFileSync(rejectionFile, 'utf8')) as { + rejections: Record; + }; + expect(raw.rejections['ISSUE-1'].count).toBe(9); + expect(raw.rejections['ISSUE-1'].reasons.length).toBeLessThanOrEqual(5); + }, 30_000); + + it('keeps every pipeline history entry across sequential locked appends', async () => { + const { appendPipelineHistory, getPipelineHistory } = await loadRunnerState(); + for (let i = 0; i < 6; i++) { + appendPipelineHistory({ + sessionId: `s-${i}`, + taskTitle: `t-${i}`, + success: true, + finalStatus: 'done', + iterations: 1, + totalDuration: 1, + stages: [], + completedAt: new Date(2026, 0, 1, 0, 0, i).toISOString(), + }); + } + expect(getPipelineHistory(20)).toHaveLength(6); + }); +}); diff --git a/src/automation/runnerState.rejection.fixture.ts b/src/automation/runnerState.rejection.fixture.ts new file mode 100644 index 00000000..53f24030 --- /dev/null +++ b/src/automation/runnerState.rejection.fixture.ts @@ -0,0 +1,15 @@ +const rejectionFile = process.argv[2]; +const issueId = process.argv[3]; +const reason = process.argv[4]; +if (!rejectionFile || !issueId || !reason) { + console.error('usage: fixture '); + process.exit(2); +} + +process.env.OPENSWARM_RUNNER_REJECTION_STATE_FILE = rejectionFile; +process.env.HOME = process.env.HOME || '/tmp'; +process.env.USERPROFILE = process.env.USERPROFILE || process.env.HOME; + +const { incrementRejection } = await import('./runnerState.js'); +const count = incrementRejection(issueId, reason); +process.stdout.write(String(count)); diff --git a/src/automation/runnerState.ts b/src/automation/runnerState.ts index aba656ca..f1748756 100644 --- a/src/automation/runnerState.ts +++ b/src/automation/runnerState.ts @@ -21,6 +21,7 @@ import { taskEventKey, type TaskItem } from '../orchestration/decisionEngine.js' import type { PipelineResult } from '../agents/pairPipelineTypes.js'; import type { VerifyEvidence } from '../verify/runner.js'; import { atomicWriteFileSync } from '../support/atomicFile.js'; +import { withFileLockSync } from '../support/fileLock.js'; import { isProofCapableSpace, processAppearsAlive, @@ -241,7 +242,9 @@ export function loadProjectSelection(file: string = PROJECT_SELECTION_FILE): Pro export function saveProjectSelection(sel: ProjectSelection, file: string = PROJECT_SELECTION_FILE): void { try { ensureParentDir(file); - atomicWriteFileSync(file, JSON.stringify(sel, null, 2)); + withFileLockSync(`${file}.lock`, () => { + atomicWriteFileSync(file, JSON.stringify(sel, null, 2)); + }); } catch (err) { console.warn('[ProjectSelection] Failed to save:', err); } @@ -257,13 +260,18 @@ function pruneOldEntries(entries: ProjectPaceEntry[]): ProjectPaceEntry[] { // below — daily-pace.json remains useful as a cost/throughput telemetry trail. export function recordProjectCompletion(projectName: string, costUsd?: number): void { - const state = ensurePaceLoaded(); - if (!state.projects[projectName]) state.projects[projectName] = []; - state.projects[projectName] = pruneOldEntries(state.projects[projectName]); - state.projects[projectName].push({ completedAt: new Date().toISOString(), costUsd }); - state.updatedAt = new Date().toISOString(); - savePace(); - console.log(`[Pace] ${projectName}: ${state.projects[projectName].length} tasks in 5h window`); + withFileLockSync(`${DAILY_PACE_FILE}.lock`, () => { + // Reload under the lock: two processes appending completions would otherwise + // each write a snapshot taken before the other's append, dropping it. + paceState = null; + const state = ensurePaceLoaded(); + if (!state.projects[projectName]) state.projects[projectName] = []; + state.projects[projectName] = pruneOldEntries(state.projects[projectName]); + state.projects[projectName].push({ completedAt: new Date().toISOString(), costUsd }); + state.updatedAt = new Date().toISOString(); + savePace(); + console.log(`[Pace] ${projectName}: ${state.projects[projectName].length} tasks in 5h window`); + }); } export function getDailyPaceInfo(): DailyPaceState { @@ -572,48 +580,57 @@ export function getRejectionCount(issueId: string): number { } export function incrementRejection(issueId: string, reason: string): number { - const state = ensureRejectionStateLoaded(); - const entry = state.rejections[issueId] || { - issueId, - count: 0, - lastRejection: new Date().toISOString(), - reasons: [], - }; + return withFileLockSync(`${REJECTION_STATE_FILE}.lock`, () => { + // Cross-process read-modify-write: reload disk state under the lock, or a + // concurrent increment's count and reasons are dropped by this write. + rejectionState = null; + const state = ensureRejectionStateLoaded(); + const entry = state.rejections[issueId] || { + issueId, + count: 0, + lastRejection: new Date().toISOString(), + reasons: [], + }; - entry.count++; - entry.lastRejection = new Date().toISOString(); - entry.reasons.push(reason); + entry.count++; + entry.lastRejection = new Date().toISOString(); + entry.reasons.push(reason); - // Keep only last 5 reasons - if (entry.reasons.length > 5) { - entry.reasons = entry.reasons.slice(-5); - } + // Keep only last 5 reasons + if (entry.reasons.length > 5) { + entry.reasons = entry.reasons.slice(-5); + } - state.rejections[issueId] = entry; - state.updatedAt = new Date().toISOString(); + state.rejections[issueId] = entry; + state.updatedAt = new Date().toISOString(); - // Persist to disk - try { - ensureParentDir(REJECTION_STATE_FILE); - atomicWriteFileSync(REJECTION_STATE_FILE, JSON.stringify(state, null, 2)); - } catch (err) { - console.warn('[RejectionState] Failed to save:', err); - } + try { + ensureParentDir(REJECTION_STATE_FILE); + atomicWriteFileSync(REJECTION_STATE_FILE, JSON.stringify(state, null, 2)); + } catch (err) { + console.warn('[RejectionState] Failed to save:', err); + } - return entry.count; + return entry.count; + }, { timeoutMs: 10_000 }); } export function clearRejection(issueId: string): void { - const state = ensureRejectionStateLoaded(); - delete state.rejections[issueId]; - state.updatedAt = new Date().toISOString(); + withFileLockSync(`${REJECTION_STATE_FILE}.lock`, () => { + // Reload under the lock so the clear applies to the current disk state, + // not a snapshot that another process has already incremented. + rejectionState = null; + const state = ensureRejectionStateLoaded(); + delete state.rejections[issueId]; + state.updatedAt = new Date().toISOString(); - try { - ensureParentDir(REJECTION_STATE_FILE); - atomicWriteFileSync(REJECTION_STATE_FILE, JSON.stringify(state, null, 2)); - } catch (err) { - console.warn('[RejectionState] Failed to save:', err); - } + try { + ensureParentDir(REJECTION_STATE_FILE); + atomicWriteFileSync(REJECTION_STATE_FILE, JSON.stringify(state, null, 2)); + } catch (err) { + console.warn('[RejectionState] Failed to save:', err); + } + }, { timeoutMs: 10_000 }); } export function isRejectionLimitReached(issueId: string): boolean { @@ -694,9 +711,16 @@ export function getChildrenCount(issueId: string): number { * this function ensures the counter resets even when using the in-memory cache. */ function resetDailyCounterIfNeeded(): void { - const state = ensureDecompositionStateLoaded(); const today = new Date().toLocaleDateString('en-CA'); - if (state.dailyCreationDate !== today) { + // Fast path: this runs on every canCreateMoreIssues call, and no lock is + // needed when the cached date is already current. Only an actual reset takes + // the cross-process lock (re-checked under it) so a concurrent decomposition's + // count is not clobbered by a stale reset write. + if (ensureDecompositionStateLoaded().dailyCreationDate === today) return; + withFileLockSync(`${DECOMPOSITION_STATE_FILE}.lock`, () => { + decompositionState = null; + const state = ensureDecompositionStateLoaded(); + if (state.dailyCreationDate === today) return; console.log(`[DecompositionState] Daily counter reset: ${state.dailyCreationCount} → 0 (date: ${state.dailyCreationDate} → ${today})`); state.dailyCreationCount = 0; state.dailyCreationDate = today; @@ -707,7 +731,7 @@ function resetDailyCounterIfNeeded(): void { } catch (err) { console.warn('[DecompositionState] Failed to persist daily reset:', err); } - } + }, { timeoutMs: 10_000 }); } export function getDailyCreationCount(): number { @@ -772,59 +796,67 @@ export function registerDecomposition( parentId: string | undefined, childrenIds: string[] ): void { - const state = ensureDecompositionStateLoaded(); - const now = new Date().toISOString(); - const parentDepth = parentId ? (state.decompositions[parentId]?.depth ?? 0) : -1; - const issueDepth = parentDepth + 1; - const uniqueChildren = [...new Set(childrenIds)]; - const existingChildren = new Set( - Object.values(state.decompositions) - .filter((entry) => entry.parentId === issueId) - .map((entry) => entry.issueId), - ); - - // Validate the full batch before mutating the in-memory projection. A child - // identity collision must leave no half-created parent entry behind. - for (const childId of uniqueChildren) { - const existing = state.decompositions[childId]; - if (existing && existing.parentId !== issueId) { - throw new Error(`Decomposition child ${childId} is already owned by ${existing.parentId ?? 'no parent'}`); + // Date rollover first, and outside the lock: resetDailyCounterIfNeeded takes + // the same lock and withFileLockSync is not reentrant. + resetDailyCounterIfNeeded(); + withFileLockSync(`${DECOMPOSITION_STATE_FILE}.lock`, () => { + // Reload under the lock: concurrent decompositions in sibling processes would + // otherwise each write a snapshot missing the other's child links and daily + // budget spend. + decompositionState = null; + const state = ensureDecompositionStateLoaded(); + const now = new Date().toISOString(); + const parentDepth = parentId ? (state.decompositions[parentId]?.depth ?? 0) : -1; + const issueDepth = parentDepth + 1; + const uniqueChildren = [...new Set(childrenIds)]; + const existingChildren = new Set( + Object.values(state.decompositions) + .filter((entry) => entry.parentId === issueId) + .map((entry) => entry.issueId), + ); + + // Validate the full batch before mutating the in-memory projection. A child + // identity collision must leave no half-created parent entry behind. + for (const childId of uniqueChildren) { + const existing = state.decompositions[childId]; + if (existing && existing.parentId !== issueId) { + throw new Error(`Decomposition child ${childId} is already owned by ${existing.parentId ?? 'no parent'}`); + } } - } - - const existingIssue = state.decompositions[issueId]; - state.decompositions[issueId] = { - issueId, - parentId, - depth: issueDepth, - childrenCount: new Set([...existingChildren, ...uniqueChildren]).size, - createdAt: existingIssue?.createdAt ?? now, - }; - let newlyRegistered = 0; - for (const childId of uniqueChildren) { - const existing = state.decompositions[childId]; - if (!existingChildren.has(childId)) newlyRegistered++; - state.decompositions[childId] = { - issueId: childId, - parentId: issueId, - depth: issueDepth + 1, - childrenCount: existing?.childrenCount ?? 0, - createdAt: existing?.createdAt ?? now, + const existingIssue = state.decompositions[issueId]; + state.decompositions[issueId] = { + issueId, + parentId, + depth: issueDepth, + childrenCount: new Set([...existingChildren, ...uniqueChildren]).size, + createdAt: existingIssue?.createdAt ?? now, }; - } - // Retried deterministic children do not consume the daily budget twice. - state.dailyCreationCount += newlyRegistered; - state.updatedAt = new Date().toISOString(); + let newlyRegistered = 0; + for (const childId of uniqueChildren) { + const existing = state.decompositions[childId]; + if (!existingChildren.has(childId)) newlyRegistered++; + state.decompositions[childId] = { + issueId: childId, + parentId: issueId, + depth: issueDepth + 1, + childrenCount: existing?.childrenCount ?? 0, + createdAt: existing?.createdAt ?? now, + }; + } - // Persist to disk - try { - ensureParentDir(DECOMPOSITION_STATE_FILE); - atomicWriteFileSync(DECOMPOSITION_STATE_FILE, JSON.stringify(state, null, 2)); - } catch (err) { - console.warn('[DecompositionState] Failed to save:', err); - } + // Retried deterministic children do not consume the daily budget twice. + state.dailyCreationCount += newlyRegistered; + state.updatedAt = new Date().toISOString(); + + try { + ensureParentDir(DECOMPOSITION_STATE_FILE); + atomicWriteFileSync(DECOMPOSITION_STATE_FILE, JSON.stringify(state, null, 2)); + } catch (err) { + console.warn('[DecompositionState] Failed to save:', err); + } + }, { timeoutMs: 10_000 }); } // Pipeline History (persistent, time-ordered) @@ -848,14 +880,23 @@ function ensureHistoryLoaded(): PipelineHistoryEntry[] { } export function appendPipelineHistory(entry: PipelineHistoryEntry): void { - const history = ensureHistoryLoaded(); - history.unshift(entry); // newest first - if (history.length > MAX_PIPELINE_HISTORY) { - history.length = MAX_PIPELINE_HISTORY; - } try { - ensureParentDir(PIPELINE_HISTORY_FILE); - atomicWriteFileSync(PIPELINE_HISTORY_FILE, JSON.stringify(history, null, 2)); + withFileLockSync(`${PIPELINE_HISTORY_FILE}.lock`, () => { + // Read-modify-write against disk, not the process-local cache: a second + // process appending its own entry would otherwise overwrite this one. + let history: PipelineHistoryEntry[] = []; + try { + if (existsSync(PIPELINE_HISTORY_FILE)) { + const parsed = JSON.parse(readFileSync(PIPELINE_HISTORY_FILE, 'utf8')) as PipelineHistoryEntry[]; + if (Array.isArray(parsed)) history = parsed; + } + } catch { history = []; } + history.unshift(entry); // newest first + if (history.length > MAX_PIPELINE_HISTORY) history.length = MAX_PIPELINE_HISTORY; + ensureParentDir(PIPELINE_HISTORY_FILE); + atomicWriteFileSync(PIPELINE_HISTORY_FILE, JSON.stringify(history, null, 2)); + pipelineHistory = history; + }, { timeoutMs: 10_000 }); } catch (err) { console.warn('[PipelineHistory] Failed to save:', err); } diff --git a/src/knowledge/gitInfo.test.ts b/src/knowledge/gitInfo.test.ts index 87f5db92..2f8ccefb 100644 --- a/src/knowledge/gitInfo.test.ts +++ b/src/knowledge/gitInfo.test.ts @@ -1,5 +1,20 @@ -import { describe, expect, it } from 'vitest'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { EventEmitter } from 'node:events'; import { parseNulDelimitedChurnOutput } from './gitInfo.js'; +import { KnowledgeGraph } from './graph.js'; +import type { GraphNode } from './types.js'; + +const saveGraphMock = vi.hoisted(() => vi.fn(async () => {})); +vi.mock('./store.js', () => ({ + saveGraph: saveGraphMock, +})); + +const spawnMock = vi.hoisted(() => vi.fn()); +vi.mock('node:child_process', () => ({ + spawn: spawnMock, +})); + +import { enrichWithGitInfo } from './gitInfo.js'; /** * Real `git log -z --format=%x1e%ct --name-only` output: commits are @@ -50,3 +65,70 @@ describe('parseNulDelimitedChurnOutput', () => { expect(parseNulDelimitedChurnOutput('\0\0').size).toBe(0); }); }); + +function moduleNode(id: string, path: string): GraphNode { + return { + id, + type: 'module', + name: id, + path, + metrics: { loc: 10, exportCount: 1, importCount: 1, language: 'typescript' }, + }; +} + +function mockGitLogOutput(filePath: string, timestampSec: number): string { + // Matches `git log -z --format=%x1e%ct --name-only`, which + // parseNulDelimitedChurnOutput consumes: each commit is RS+ts, its first + // filename carries the format-terminating newline, and there is no empty + // token between commits. + return `${RS}${timestampSec}\0\n${filePath}\0`; +} + +beforeEach(() => { + vi.clearAllMocks(); + vi.spyOn(console, 'log').mockImplementation(() => undefined); + vi.spyOn(console, 'warn').mockImplementation(() => undefined); +}); + +describe('enrichWithGitInfo', () => { + it('does not mutate original node objects and persists via saveGraph', async () => { + const graph = new KnowledgeGraph('test-project', '/repo'); + const node = moduleNode('mod-a', 'src/a.ts'); + graph.addNode(node); + const nodeRefBefore = graph.getNode('mod-a'); + expect(nodeRefBefore?.gitInfo).toBeUndefined(); + + const nowSec = Math.floor(Date.now() / 1000); + spawnMock.mockImplementation((_cmd: string, _args: string[], _opts: unknown) => { + const proc = new EventEmitter() as EventEmitter & { + stdout: EventEmitter; + stderr: EventEmitter; + kill: () => void; + }; + proc.stdout = new EventEmitter(); + proc.stderr = new EventEmitter(); + proc.kill = vi.fn(); + queueMicrotask(() => { + proc.stdout.emit('data', mockGitLogOutput('src/a.ts', nowSec)); + proc.emit('close', 0); + }); + return proc; + }); + + await enrichWithGitInfo(graph, '/repo'); + + expect(node.gitInfo).toBeUndefined(); + expect(nodeRefBefore?.gitInfo).toBeUndefined(); + + const enriched = graph.getNode('mod-a'); + expect(enriched).toBeDefined(); + expect(enriched).not.toBe(node); + expect(enriched?.metrics).not.toBe(node.metrics); + expect(enriched?.gitInfo).toEqual({ + lastCommitDate: nowSec * 1000, + commitCount30d: 1, + churnScore: 1, + }); + expect(saveGraphMock).toHaveBeenCalledWith(graph); + }); +}); diff --git a/src/knowledge/gitInfo.ts b/src/knowledge/gitInfo.ts index 004314cc..cb049114 100644 --- a/src/knowledge/gitInfo.ts +++ b/src/knowledge/gitInfo.ts @@ -6,6 +6,7 @@ import { spawn } from 'node:child_process'; import type { KnowledgeGraph } from './graph.js'; import type { GitInfo } from './types.js'; +import { saveGraph } from './store.js'; // Git Command Runner (same pattern as gitTracker.ts) @@ -139,23 +140,30 @@ export async function enrichWithGitInfo( for (const mod of modules) { const churn = churns.get(mod.path); - if (churn) { - const gitInfo: GitInfo = { + const gitInfo: GitInfo = churn + ? { lastCommitDate: churn.lastCommitDate, commitCount30d: churn.commitCount, churnScore: Math.round((churn.commitCount / maxCommits) * 1000) / 1000, - }; - mod.gitInfo = gitInfo; - } else { - // File not in git history (no changes in 30 days) - mod.gitInfo = { + } + : { lastCommitDate: 0, commitCount30d: 0, churnScore: 0, }; - } + // addNode REPLACES the entry — writing mod.gitInfo directly only mutated the + // copy callers already held, so churn data was recomputed on every load and + // never survived a process restart. (AGT-3420) metrics is copied so the graph + // never shares a mutable object with the node a caller still holds. + graph.addNode({ + ...mod, + metrics: mod.metrics ? { ...mod.metrics } : undefined, + gitInfo, + }); } + await saveGraph(graph); + console.log(`[GitInfo] Enriched ${modules.length} modules with git data (${churns.size} files had changes in ${sinceDays}d)`); } diff --git a/src/locale/index.ts b/src/locale/index.ts index 11e24c5c..e5efc8d0 100644 --- a/src/locale/index.ts +++ b/src/locale/index.ts @@ -3,6 +3,7 @@ // t() helper, initLocale(), getPrompts(), getDateLocale() // ============================================ +import { AsyncLocalStorage } from 'node:async_hooks'; import type { LocaleMessages, PromptTemplates, SupportedLocale } from './types.js'; import { en } from './en.js'; import { ko } from './ko.js'; @@ -13,9 +14,12 @@ export type { LocaleMessages, PromptTemplates, SupportedLocale } from './types.j // ── State ───────────────────────────────── -let currentLocale: SupportedLocale = 'en'; -let currentMessages: LocaleMessages = en; -let currentPrompts: PromptTemplates = enPrompts; +// Process-global default locale. This is only the fallback for code that runs +// outside any `withLocale` scope (e.g. top-level CLI setup). Concurrent +// executions must not mutate it — they should use `withLocale` to scope their +// locale choice to the current async execution instead, so one runner's locale +// never leaks into another's. (AGT-3420) +let defaultLocale: SupportedLocale = 'en'; const catalogs: Record = { en, ko }; const promptCatalogs: Record = { @@ -23,6 +27,12 @@ const promptCatalogs: Record = { ko: koPrompts, }; +// Execution-scoped locale. `withLocale` sets this for the duration of an async +// execution; `t`/`getPrompts`/`getDateLocale`/`getLocale` read it first and only +// fall back to `defaultLocale` when no scope is active. This removes the +// mutable process-global from concurrent execution paths. +const localeScope = new AsyncLocalStorage(); + type LocaleLeafKey = { [K in Extract]: T[K] extends string @@ -38,24 +48,41 @@ type LocaleLookupKey = K extends LocaleKey ? K : string extend // ── Public API ──────────────────────────── /** - * Initialize the locale module. Call once at startup. + * Initialize the default locale module. Call once at startup. + * + * This sets the process-global fallback locale. For concurrent executions that + * need a specific locale, prefer `withLocale` so the choice is scoped to that + * execution and does not leak into sibling runners. */ export function initLocale(locale: SupportedLocale = 'en'): void { if (!catalogs[locale]) { console.warn(`[Locale] Unknown locale "${locale}", falling back to "en"`); locale = 'en'; } - currentLocale = locale; - currentMessages = catalogs[locale]; - currentPrompts = promptCatalogs[locale]; + defaultLocale = locale; console.log(`[Locale] Initialized: ${locale}`); } +/** + * Run `fn` with `locale` scoped to the current async execution. + * + * Any `t`/`getPrompts`/`getDateLocale`/`getLocale` call made (synchronously or + * through awaited work) inside `fn` resolves to `locale`, and the previous + * scope is restored when `fn` returns — so concurrent executions each see their + * own locale and never mutate a shared process-global. + */ +export async function withLocale( + locale: SupportedLocale, + fn: () => Promise | T, +): Promise { + return localeScope.run(locale, async () => fn()); +} + /** * Get the current locale identifier. */ export function getLocale(): SupportedLocale { - return currentLocale; + return localeScope.getStore() ?? defaultLocale; } /** @@ -67,9 +94,11 @@ export function getLocale(): SupportedLocale { * t('discord.errors.sessionNotFound', { name: 'main' }) */ export function t(key: LocaleLookupKey, params?: Record): string { - const value = resolvePath(currentMessages, key); + const locale = getLocale(); + const messages = catalogs[locale]; + const value = resolvePath(messages, key); if (value === undefined) { - console.warn(`[Locale] Missing key: "${key}" for locale "${currentLocale}"`); + console.warn(`[Locale] Missing key: "${key}" for locale "${locale}"`); return key; } if (typeof value !== 'string') { @@ -84,14 +113,14 @@ export function t(key: LocaleLookupKey, params?: Reco * Return the current locale's prompt templates. */ export function getPrompts(): PromptTemplates { - return currentPrompts; + return promptCatalogs[getLocale()]; } /** * Return the BCP 47 locale tag for Date.toLocaleString() etc. */ export function getDateLocale(): string { - return currentLocale === 'ko' ? 'ko-KR' : 'en-US'; + return getLocale() === 'ko' ? 'ko-KR' : 'en-US'; } // ── Internals ───────────────────────────── diff --git a/src/locale/locale.scope.test.ts b/src/locale/locale.scope.test.ts new file mode 100644 index 00000000..e22e3c4d --- /dev/null +++ b/src/locale/locale.scope.test.ts @@ -0,0 +1,38 @@ +import { afterEach, describe, expect, it } from 'vitest'; +import { getLocale, initLocale, t, withLocale } from './index.js'; + +afterEach(() => { + initLocale('en'); +}); + +describe('execution-scoped locale (AGT-3420)', () => { + it('isolates concurrent withLocale scopes so siblings do not leak', async () => { + initLocale('en'); + + const seen: string[] = []; + await Promise.all([ + withLocale('ko', async () => { + await new Promise((r) => setTimeout(r, 20)); + seen.push(`ko:${getLocale()}:${t('common.timeAgo.justNow')}`); + }), + withLocale('en', async () => { + await new Promise((r) => setTimeout(r, 5)); + seen.push(`en:${getLocale()}:${t('common.timeAgo.justNow')}`); + }), + ]); + + expect(seen).toContain('en:en:just now'); + expect(seen).toContain('ko:ko:방금 전'); + // Process default remains whatever initLocale last set — scopes must not + // permanently flip it for other concurrent work. + expect(getLocale()).toBe('en'); + }); + + it('restores the outer locale after withLocale returns', async () => { + initLocale('en'); + await withLocale('ko', async () => { + expect(getLocale()).toBe('ko'); + }); + expect(getLocale()).toBe('en'); + }); +}); diff --git a/src/memory/codex.ts b/src/memory/codex.ts index 27026b77..8f1a9a36 100644 --- a/src/memory/codex.ts +++ b/src/memory/codex.ts @@ -16,6 +16,7 @@ import { resolve, basename, join } from 'path'; import { getDateLocale } from '../locale/index.js'; import { homedir } from 'os'; import { createHash } from 'crypto'; +import { withFileLock } from '../support/fileLock.js'; import { atomicWriteFile } from '../support/atomicFile.js'; // Codex storage path @@ -308,9 +309,15 @@ export async function saveSession( /** * Update index.md + * + * Cross-process read-modify-write: a second runner saving a session while this + * one holds the parsed index would insert its entry into a snapshot that this + * write then overwrites, losing that session from the index entirely. The lock + * is taken before the read, not before the write. */ async function updateIndex(session: CodexSession, summaryPath: string): Promise { const indexPath = join(CODEX_DIR, 'index.md'); + await withFileLock(indexPath + '.lock', async () => { let content = await fs.readFile(indexPath, 'utf-8'); const relativePath = summaryPath.replace(CODEX_DIR + '/', ''); @@ -354,6 +361,7 @@ async function updateIndex(session: CodexSession, summaryPath: string): Promise< await atomicWriteFile(indexPath, content, 0o644); console.log('[Codex] Updated index.md'); + }); } /** diff --git a/src/memory/memoryCore.ts b/src/memory/memoryCore.ts index 143f4581..f31a4da9 100644 --- a/src/memory/memoryCore.ts +++ b/src/memory/memoryCore.ts @@ -5,9 +5,10 @@ */ import { connect, Table, Connection } from '@lancedb/lancedb'; import { pipeline, env as transformersEnv, type FeatureExtractionPipeline } from '@huggingface/transformers'; -import { resolve } from 'path'; +import { join, resolve } from 'path'; import { homedir } from 'os'; import { c, status } from '../support/colors.js'; +import { withFileLock } from '../support/fileLock.js'; import { safeConsole as console } from '../support/safeLog.js'; import { randomUUID } from 'node:crypto'; import { @@ -417,6 +418,24 @@ export async function getMemoryIdsByDerivedFrom(derivedFrom: string, limit = 100 return rows.map((row: any) => String(row.id)); } +/** + * Cross-process lock serializing memory table mutations (writes and, in + * compaction, the build-then-swap that replaces the table). + * + * {@link withMemoryWriteLock} only orders writers inside one process. The + * competing writers that motivated the retry below are separate `openswarm` + * processes sharing one on-disk Lance table, which an in-process chain cannot + * see; this lock is what makes compaction and a writer exclude each other. + */ +export function crossProcessMemoryMutationLockPath(): string { + return process.env.OPENSWARM_MEMORY_MUTATION_LOCK + ?? join(homedir(), '.openswarm', 'memory-mutation.lock'); +} + +export async function withCrossProcessMemoryMutationLock(operation: () => Promise): Promise { + return withFileLock(crossProcessMemoryMutationLockPath(), operation, { timeoutMs: 120_000 }); +} + /** * Retry a Lance write (add/update/delete) on optimistic-concurrency conflict. * @@ -429,13 +448,16 @@ export async function getMemoryIdsByDerivedFrom(derivedFrom: string, limit = 100 * jitter to desynchronize the competing writers. Appends and predicated * update/delete are safe to re-run: Lance re-commits against the latest version on * each attempt, so a retry is not a double-apply. + * + * Each attempt also runs under {@link withCrossProcessMemoryMutationLock}, so a compaction + * cannot swap the table out from under a write mid-retry. */ export async function withMemoryWriteRetry(op: () => Promise, label = 'write'): Promise { return withMemoryWriteLock(async () => { const MAX_ATTEMPTS = 8; for (let attempt = 1; ; attempt++) { try { - return await op(); + return await withCrossProcessMemoryMutationLock(() => op()); } catch (err) { const msg = err instanceof Error ? err.message : String(err); // Keep this matcher tight to genuine optimistic-concurrency conflicts. "Too diff --git a/src/memory/memoryWriteRetry.test.ts b/src/memory/memoryWriteRetry.test.ts index 35e154df..a4c69e1b 100644 --- a/src/memory/memoryWriteRetry.test.ts +++ b/src/memory/memoryWriteRetry.test.ts @@ -1,13 +1,28 @@ -import { describe, it, expect, vi, afterEach } from 'vitest'; -import { withMemoryWriteRetry } from './memoryCore.js'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { describe, it, expect, vi, afterEach, afterAll } from 'vitest'; + +const lockDir = mkdtempSync(join(tmpdir(), 'openswarm-mem-lock-')); +process.env.OPENSWARM_MEMORY_MUTATION_LOCK = join(lockDir, 'mutation.lock'); + +const { withMemoryWriteRetry } = await import('./memoryCore.js'); // withMemoryWriteRetry wraps Lance writes so `openswarm review --max` (up to 16 // concurrent reviewer processes sharing one on-disk table) survives Lance's // optimistic-concurrency conflicts instead of surfacing "Too many concurrent -// writers". Fake timers skip the real backoff sleeps. +// writers". +// +// Real timers, deliberately: each attempt now also takes a cross-process file +// lock, whose acquisition is real fs I/O. Fake timers drain the backoff queue +// before that I/O settles, so the next backoff timer is armed after the drain +// has already finished and the test hangs. The sleeps are ~25ms + jitter. describe('withMemoryWriteRetry (INT-2817 store-path concurrency)', () => { + afterAll(() => { + rmSync(lockDir, { recursive: true, force: true }); + }); + afterEach(() => { - vi.useRealTimers(); vi.restoreAllMocks(); }); @@ -18,16 +33,13 @@ describe('withMemoryWriteRetry (INT-2817 store-path concurrency)', () => { }); it('retries on a concurrent-writer conflict and then succeeds', async () => { - vi.useFakeTimers(); const op = vi.fn() .mockRejectedValueOnce(new Error('lance error: Too many concurrent writers.')) .mockRejectedValueOnce(new Error('Commit conflict: version conflict detected')) .mockResolvedValue('stored'); - const p = withMemoryWriteRetry(op, 'test'); - await vi.runAllTimersAsync(); - await expect(p).resolves.toBe('stored'); + await expect(withMemoryWriteRetry(op, 'test')).resolves.toBe('stored'); expect(op).toHaveBeenCalledTimes(3); - }); + }, 30_000); it('rethrows a non-retryable error immediately (no retry)', async () => { const op = vi.fn().mockRejectedValue(new Error('schema mismatch: column not found')); @@ -36,12 +48,18 @@ describe('withMemoryWriteRetry (INT-2817 store-path concurrency)', () => { }); it('gives up after the attempt cap when the conflict never clears', async () => { - vi.useFakeTimers(); const op = vi.fn().mockRejectedValue(new Error('Too many concurrent writers.')); - const p = withMemoryWriteRetry(op, 'test'); - const assertion = expect(p).rejects.toThrow('concurrent writers'); - await vi.runAllTimersAsync(); - await assertion; + await expect(withMemoryWriteRetry(op, 'test')).rejects.toThrow('concurrent writers'); expect(op).toHaveBeenCalledTimes(8); // MAX_ATTEMPTS - }); + }, 30_000); + + it('releases the mutation lock when an attempt fails, so the next writer can proceed', async () => { + // A lock leaked on the error path would wedge every later memory write + // until the lock timed out (120s); the next retry has to be able to take it. + const op = vi.fn() + .mockRejectedValueOnce(new Error('Commit conflict: version conflict detected')) + .mockResolvedValue('stored'); + await expect(withMemoryWriteRetry(op, 'test')).resolves.toBe('stored'); + expect(op).toHaveBeenCalledTimes(2); + }, 30_000); }); diff --git a/src/memory/reembed.ts b/src/memory/reembed.ts index 9ecf1c30..c37ddf49 100644 --- a/src/memory/reembed.ts +++ b/src/memory/reembed.ts @@ -9,6 +9,8 @@ // table in one pass, following compaction's build-then-swap shape so a failure // leaves the original table intact. +import { join } from 'node:path'; +import { withFileLock } from '../support/fileLock.js'; import { c, status } from '../support/colors.js'; import { EMBEDDING_DIM, @@ -45,77 +47,83 @@ export interface ReembedOptions { } export async function reembedMemoryTable(options: ReembedOptions = {}): Promise { - await initDatabase(); - const db = getDb(); - const table = getTable(); - if (!db || !table) throw new Error('Memory database is not initialized'); + const memoryDir = options.memoryDir ?? MEMORY_DIR; + const lockPath = join(memoryDir, '.reembed.lock'); - const spec = resolveEmbeddingConfig(); - const signature = embeddingSignature(spec); - const progressEvery = options.progressEvery ?? 50; + return withFileLock(lockPath, async () => { + await initDatabase(); + const db = getDb(); + const table = getTable(); + if (!db || !table) throw new Error('Memory database is not initialized'); - const rows = (await table.query().limit(1_000_000).toArray()) as unknown as CognitiveMemoryRecord[]; - const total = rows.length; - console.log(`${status.info('[Reembed]')} ${c.dim('rebuilding')} ${c.cyan(String(total))} ${c.dim('vectors with')} ${c.yellow(spec.id)}`); + const spec = resolveEmbeddingConfig(); + const signature = embeddingSignature(spec); + const progressEvery = options.progressEvery ?? 50; - // normalizeRecords first so the rewritten table lands on the lean v3 schema, - // exactly like compaction does; vectors are replaced immediately after. - const normalized = normalizeRecords(rows); - let reembedded = 0; - let empty = 0; + const rows = (await table.query().limit(1_000_000).toArray()) as unknown as CognitiveMemoryRecord[]; + rows.sort((a, b) => String(a.id).localeCompare(String(b.id))); + const total = rows.length; + console.log(`${status.info('[Reembed]')} ${c.dim('rebuilding')} ${c.cyan(String(total))} ${c.dim('vectors with')} ${c.yellow(spec.id)}`); - for (let i = 0; i < normalized.length; i++) { - const record = normalized[i]; - const text = embeddingTextFor(String(record.title ?? ''), String(record.content ?? '')); - if (!text) { - record.vector = Array.from({ length: EMBEDDING_DIM }, () => 0); - empty++; - } else { - record.vector = await embedPassage(text); - reembedded++; - } - if ((i + 1) % progressEvery === 0) { - options.onProgress?.(i + 1, total); - console.log(`${c.dim(`[Reembed] ${i + 1}/${total}`)}`); - } - } - options.onProgress?.(total, total); + // normalizeRecords first so the rewritten table lands on the lean v3 schema, + // exactly like compaction does; vectors are replaced immediately after. + const normalized = normalizeRecords(rows); + let reembedded = 0; + let empty = 0; - const targetTableName = table.name; - const tempTableName = `${targetTableName}_reembed_${Date.now()}`; + for (let i = 0; i < normalized.length; i++) { + const record = normalized[i]; + const text = embeddingTextFor(String(record.title ?? ''), String(record.content ?? '')); + if (!text) { + record.vector = Array.from({ length: EMBEDDING_DIM }, () => 0); + empty++; + } else { + record.vector = await embedPassage(text); + reembedded++; + } + if ((i + 1) % progressEvery === 0) { + options.onProgress?.(i + 1, total); + console.log(`${c.dim(`[Reembed] ${i + 1}/${total}`)}`); + } + } + options.onProgress?.(total, total); - // Build a validated replacement before touching the live table. - if (normalized.length > 0) { - await db.createTable(tempTableName, normalized); - } else { - await db.createEmptyTable(tempTableName, await table.schema()); - } + const targetTableName = table.name; + const tempTableName = `${targetTableName}_reembed_${Date.now()}`; - let replaced = false; - try { + // Build a validated replacement before touching the live table. if (normalized.length > 0) { - await db.createTable(targetTableName, normalized, { mode: 'overwrite' }); + await db.createTable(tempTableName, normalized); } else { - await db.createEmptyTable(targetTableName, await table.schema(), { mode: 'overwrite' }); + await db.createEmptyTable(tempTableName, await table.schema()); } - setTable(await db.openTable(targetTableName)); - replaced = true; - } finally { - if (replaced) { - try { - await db.dropTable(tempTableName); - } catch (cleanupError) { - console.warn(`[Reembed] Failed to drop temporary table ${tempTableName}:`, cleanupError); + + let replaced = false; + try { + if (normalized.length > 0) { + await db.createTable(targetTableName, normalized, { mode: 'overwrite' }); + } else { + await db.createEmptyTable(targetTableName, await table.schema(), { mode: 'overwrite' }); + } + setTable(await db.openTable(targetTableName)); + replaced = true; + } finally { + if (replaced) { + try { + await db.dropTable(tempTableName); + } catch (cleanupError) { + console.warn(`[Reembed] Failed to drop temporary table ${tempTableName}:`, cleanupError); + } + } else { + console.warn(`[Reembed] Replacement failed; retained recoverable table ${tempTableName}`); } - } else { - console.warn(`[Reembed] Replacement failed; retained recoverable table ${tempTableName}`); } - } - // Only claim the new signature once the swap actually succeeded — otherwise the - // store would advertise vectors it does not have. - writeStoredSignature(options.memoryDir ?? MEMORY_DIR, signature); + // Only claim the new signature once the swap actually succeeded — otherwise the + // store would advertise vectors it does not have. + writeStoredSignature(memoryDir, signature); - console.log(`${status.ok('[Reembed] done')} ${c.dim('records:')} ${c.cyan(String(total))} ${c.dim('signature:')} ${c.yellow(signature)}`); - return { total, reembedded, empty, signature }; + console.log(`${status.ok('[Reembed] done')} ${c.dim('records:')} ${c.cyan(String(total))} ${c.dim('signature:')} ${c.yellow(signature)}`); + return { total, reembedded, empty, signature }; + }); } diff --git a/src/orchestration/decisionEngine.admission.test.ts b/src/orchestration/decisionEngine.admission.test.ts new file mode 100644 index 00000000..aabe0f26 --- /dev/null +++ b/src/orchestration/decisionEngine.admission.test.ts @@ -0,0 +1,43 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { join } from 'node:path'; +import { tmpdir } from 'node:os'; +import { + tryClaimTaskAdmission, + getTaskState, + resetTaskStateStoreForTests, +} from '../taskState/store.js'; + +describe('tryClaimTaskAdmission', () => { + let stateDir: string; + let stateFile: string; + + beforeEach(() => { + stateDir = mkdtempSync(join(tmpdir(), 'openswarm-admission-')); + stateFile = join(stateDir, 'state.json'); + process.env.OPENSWARM_TASK_STATE_FILE = stateFile; + resetTaskStateStoreForTests(); + }); + + afterEach(() => { + delete process.env.OPENSWARM_TASK_STATE_FILE; + rmSync(stateDir, { recursive: true, force: true }); + }); + + it('first claim succeeds, second returns null', () => { + const first = tryClaimTaskAdmission('AGT-3420', { + issueIdentifier: 'AGT-3420', + title: 'Atomic admission', + }); + expect(first).not.toBeNull(); + expect(first?.execution.status).toBe('in_progress'); + + const second = tryClaimTaskAdmission('AGT-3420', { + issueIdentifier: 'AGT-3420', + title: 'Atomic admission', + }); + expect(second).toBeNull(); + + expect(getTaskState('AGT-3420')?.execution.status).toBe('in_progress'); + }); +}); diff --git a/src/orchestration/decisionEngine.coverage.test.ts b/src/orchestration/decisionEngine.coverage.test.ts index e114db50..2ed7f11c 100644 --- a/src/orchestration/decisionEngine.coverage.test.ts +++ b/src/orchestration/decisionEngine.coverage.test.ts @@ -17,10 +17,27 @@ const fsMock = vi.hoisted(() => ({ })); vi.mock('fs/promises', () => fsMock); +const fileLockMock = vi.hoisted(() => ({ + withFileLock: vi.fn((_path: string, operation: () => unknown) => operation()), +})); +vi.mock('../support/fileLock.js', () => fileLockMock); + +const atomicFileMock = vi.hoisted(() => ({ + atomicWriteFile: vi.fn(), +})); +vi.mock('../support/atomicFile.js', () => atomicFileMock); + const timeWindowMock = vi.hoisted(() => ({ checkWorkAllowed: vi.fn() })); vi.mock('../support/timeWindow.js', () => timeWindowMock); -const taskStateMock = vi.hoisted(() => ({ getTaskReadiness: vi.fn() })); +const taskStateMock = vi.hoisted(() => ({ + getTaskReadiness: vi.fn(), + getTaskState: vi.fn(() => undefined), + tryClaimTaskAdmission: vi.fn((_id: string, _patch?: unknown) => ({ + issueId: _id, + execution: { status: 'in_progress' }, + })), +})); vi.mock('../taskState/store.js', () => taskStateMock); const memoryMock = vi.hoisted(() => ({ saveCognitiveMemory: vi.fn() })); @@ -114,6 +131,13 @@ beforeEach(() => { fsMock.readFile.mockRejectedValue(Object.assign(new Error('ENOENT'), { code: 'ENOENT' })); fsMock.writeFile.mockResolvedValue(undefined); fsMock.mkdir.mockResolvedValue(undefined); + fileLockMock.withFileLock.mockImplementation((_path: string, operation: () => unknown) => operation()); + atomicFileMock.atomicWriteFile.mockResolvedValue(undefined); + taskStateMock.getTaskState.mockReturnValue(undefined); + taskStateMock.tryClaimTaskAdmission.mockImplementation((_id: string) => ({ + issueId: _id, + execution: { status: 'in_progress' }, + })); workflowMock.loadWorkflow.mockResolvedValue(workflow()); workflowMock.listWorkflows.mockResolvedValue([]); workflowMock.createCIPipelineTemplate.mockReturnValue(workflow({ id: 'ci-fallback' })); @@ -198,6 +222,17 @@ describe('DecisionEngine.heartbeat', () => { expect(console.log).toHaveBeenCalledWith(expect.stringContaining('waiting on dependency (blocked by: blocker-1)')); }); + it('filters out a task that is already executing locally', async () => { + taskStateMock.getTaskState.mockReturnValueOnce({ + issueId: 'issue-1', + execution: { status: 'in_progress' }, + }); + const engine = new DecisionEngine(); + const result = await engine.heartbeat([task()]); + expect(result).toEqual({ action: 'skip', reason: 'No executable tasks in backlog' }); + expect(console.log).toHaveBeenCalledWith(expect.stringContaining('already executing')); + }); + it('rejects a task whose source is outside backlog scope', async () => { const engine = new DecisionEngine(); const result = await engine.heartbeat([task({ source: 'github_pr' })]); @@ -230,6 +265,20 @@ describe('DecisionEngine.heartbeat', () => { workflowMock.loadWorkflow.mockResolvedValueOnce(wf); const result = await engine.heartbeat([t]); expect(result).toEqual({ action: 'execute', task: t, workflow: wf, reason: `Auto-executing: ${t.title}` }); + expect(taskStateMock.tryClaimTaskAdmission).toHaveBeenCalledWith('issue-1', expect.objectContaining({ + issueIdentifier: 'INT-1', + title: t.title, + })); + }); + + it('skips when autoExecute claim fails (already claimed)', async () => { + taskStateMock.tryClaimTaskAdmission.mockReturnValueOnce(null); + const engine = new DecisionEngine({ autoExecute: true }); + const t = task({ workflowId: 'wf-1' }); + workflowMock.loadWorkflow.mockResolvedValueOnce(workflow()); + const result = await engine.heartbeat([t]); + expect(result.action).toBe('skip'); + expect(result.reason).toContain('already claimed'); }); it('defers (requires approval) when autoExecute is false', async () => { @@ -383,10 +432,11 @@ describe('DecisionEngine.addToBacklog', () => { fsMock.readFile.mockResolvedValueOnce(JSON.stringify([discoveredTask({ title: 'existing' })])); await engine.addToBacklog(discoveredTask({ title: 'new finding' })); - expect(fsMock.writeFile).toHaveBeenCalledTimes(1); - const written = JSON.parse(fsMock.writeFile.mock.calls[0][1] as string); + expect(atomicFileMock.atomicWriteFile).toHaveBeenCalledTimes(1); + const written = JSON.parse(atomicFileMock.atomicWriteFile.mock.calls[0][1] as string); expect(written).toHaveLength(2); expect(written[1].title).toBe('new finding'); + expect(fileLockMock.withFileLock).toHaveBeenCalled(); expect(memoryMock.saveCognitiveMemory).toHaveBeenCalledWith( 'belief', expect.stringContaining('new finding'), @@ -398,7 +448,7 @@ describe('DecisionEngine.addToBacklog', () => { const engine = new DecisionEngine(); fsMock.readFile.mockRejectedValueOnce(new Error('ENOENT')); await engine.addToBacklog(discoveredTask()); - const written = JSON.parse(fsMock.writeFile.mock.calls[0][1] as string); + const written = JSON.parse(atomicFileMock.atomicWriteFile.mock.calls[0][1] as string); expect(written).toHaveLength(1); }); @@ -465,6 +515,17 @@ describe('DecisionEngine.heartbeatMultiple', () => { expect(result.tasks.map((s) => s.task.id)).toEqual(['a', 'b']); expect(result.reason).toBe('Auto-executing 2 tasks'); expect(result.skippedCount).toBe(0); + expect(taskStateMock.tryClaimTaskAdmission).toHaveBeenCalledTimes(2); + }); + + it('skips when all selected tasks fail admission claim', async () => { + taskStateMock.tryClaimTaskAdmission.mockReturnValue(null); + const engine = new DecisionEngine({ autoExecute: true }); + const a = task({ id: 'a', issueId: 'a', workflowId: 'wf-1' }); + const result = await engine.heartbeatMultiple([a], 3); + expect(result.action).toBe('skip'); + expect(result.reason).toContain('already claimed'); + expect(result.tasks).toEqual([]); }); it('defers multiple selected tasks (requires approval) when autoExecute is false', async () => { diff --git a/src/orchestration/decisionEngine.ts b/src/orchestration/decisionEngine.ts index 48fd7782..607cd157 100644 --- a/src/orchestration/decisionEngine.ts +++ b/src/orchestration/decisionEngine.ts @@ -6,6 +6,8 @@ import { isAbsolute, relative, resolve } from 'path'; import { homedir } from 'os'; import * as fs from 'fs/promises'; +import { withFileLock } from '../support/fileLock.js'; +import { atomicWriteFile } from '../support/atomicFile.js'; import { WorkflowConfig, ExecutorResult, @@ -18,7 +20,7 @@ import { checkWorkAllowed } from '../support/timeWindow.js'; import { saveCognitiveMemory } from '../memory/index.js'; import { analyzeIssue } from '../knowledge/index.js'; import type { ImpactAnalysis } from '../knowledge/index.js'; -import { getTaskReadiness } from '../taskState/store.js'; +import { getTaskReadiness, getTaskState, tryClaimTaskAdmission } from '../taskState/store.js'; import { applyDurablePriorityCouncilRanking, resolvePriorityCouncilRepositoryScopes, @@ -436,20 +438,22 @@ interface EngineState { } async function loadState(): Promise { - try { - const content = await fs.readFile(ENGINE_STATE_FILE, 'utf-8'); - const saved = JSON.parse(content) as Partial; - return { - lastTaskId: saved.lastTaskId, - totalTasksCompleted: saved.totalTasksCompleted ?? 0, - totalTasksFailed: saved.totalTasksFailed ?? 0, - }; - } catch { - return { - totalTasksCompleted: 0, - totalTasksFailed: 0, - }; - } + return withFileLock(ENGINE_STATE_FILE + '.lock', async () => { + try { + const content = await fs.readFile(ENGINE_STATE_FILE, 'utf-8'); + const saved = JSON.parse(content) as Partial; + return { + lastTaskId: saved.lastTaskId, + totalTasksCompleted: saved.totalTasksCompleted ?? 0, + totalTasksFailed: saved.totalTasksFailed ?? 0, + }; + } catch { + return { + totalTasksCompleted: 0, + totalTasksFailed: 0, + }; + } + }); } // Decision Engine @@ -536,13 +540,32 @@ export class DecisionEngine { // 8. Return decision console.log(`[DecisionEngine] Returning decision: autoExecute=${this.config.autoExecute}`); + if (this.config.autoExecute) { + const issueId = selectedTask.issueId || selectedTask.id; + const claimed = tryClaimTaskAdmission(issueId, { + issueIdentifier: selectedTask.issueIdentifier, + title: selectedTask.title, + projectId: selectedTask.linearProject?.id, + projectName: selectedTask.linearProject?.name, + }); + if (!claimed) { + return { + action: 'skip', + reason: `Task ${selectedTask.issueIdentifier || issueId} already claimed by another instance`, + }; + } + return { + action: 'execute', + task: selectedTask, + workflow, + reason: `Auto-executing: ${selectedTask.title}`, + }; + } return { - action: this.config.autoExecute ? 'execute' : 'defer', + action: 'defer', task: selectedTask, workflow, - reason: this.config.autoExecute - ? `Auto-executing: ${selectedTask.title}` - : `Ready to execute (requires approval): ${selectedTask.title}`, + reason: `Ready to execute (requires approval): ${selectedTask.title}`, }; } @@ -620,12 +643,41 @@ export class DecisionEngine { } console.log(`[DecisionEngine] Selected ${selectedTasks.length} tasks for parallel execution`); + if (this.config.autoExecute) { + const claimedTasks: Array<{ task: TaskItem; workflow: WorkflowConfig }> = []; + for (const item of selectedTasks) { + const issueId = item.task.issueId || item.task.id; + const claimed = tryClaimTaskAdmission(issueId, { + issueIdentifier: item.task.issueIdentifier, + title: item.task.title, + projectId: item.task.linearProject?.id, + projectName: item.task.linearProject?.name, + }); + if (claimed) { + claimedTasks.push(item); + } else { + console.log(`[DecisionEngine] Skipping ${item.task.issueIdentifier}: already claimed`); + } + } + if (claimedTasks.length === 0) { + return { + action: 'skip', + tasks: [], + reason: 'All selected tasks already claimed by another instance', + skippedCount: sorted.length, + }; + } + return { + action: 'execute', + tasks: claimedTasks, + reason: `Auto-executing ${claimedTasks.length} tasks`, + skippedCount, + }; + } return { - action: this.config.autoExecute ? 'execute' : 'defer', + action: 'defer', tasks: selectedTasks, - reason: this.config.autoExecute - ? `Auto-executing ${selectedTasks.length} tasks` - : `Ready to execute ${selectedTasks.length} tasks (requires approval)`, + reason: `Ready to execute ${selectedTasks.length} tasks (requires approval)`, skippedCount, }; } @@ -679,6 +731,13 @@ export class DecisionEngine { return false; } + const issueId = task.issueId || task.id; + const localState = getTaskState(issueId); + if (localState?.execution.status === 'in_progress') { + console.log(`[DecisionEngine] Filtered out ${task.issueIdentifier}: already executing`); + return false; + } + const readiness = getTaskReadiness(task); if (!readiness.ready) { console.log(`[DecisionEngine] Filtered out ${task.issueIdentifier}: ${readiness.reason || 'not ready'} (blocked by: ${readiness.blockedBy.join(', ') || 'none'})`); @@ -856,20 +915,21 @@ export class DecisionEngine { async addToBacklog(discovered: DiscoveredTask): Promise { console.log(`[DecisionEngine] Adding to backlog: ${discovered.title}`); - // Save to local file (sync to Linear later) - let discoveredTasks: DiscoveredTask[] = []; - try { - const content = await fs.readFile(DISCOVERED_TASKS_FILE, 'utf-8'); - discoveredTasks = JSON.parse(content); - } catch { - discoveredTasks = []; - } + await withFileLock(DISCOVERED_TASKS_FILE + '.lock', async () => { + let discoveredTasks: DiscoveredTask[] = []; + try { + const content = await fs.readFile(DISCOVERED_TASKS_FILE, 'utf-8'); + discoveredTasks = JSON.parse(content); + } catch { + discoveredTasks = []; + } - discoveredTasks.push({ - ...discovered, - }); + discoveredTasks.push({ + ...discovered, + }); - await fs.writeFile(DISCOVERED_TASKS_FILE, JSON.stringify(discoveredTasks, null, 2)); + await atomicWriteFile(DISCOVERED_TASKS_FILE, JSON.stringify(discoveredTasks, null, 2)); + }); // Also record in memory try { diff --git a/src/orchestration/taskParser.coverage.test.ts b/src/orchestration/taskParser.coverage.test.ts index 64c99525..5d04c4cc 100644 --- a/src/orchestration/taskParser.coverage.test.ts +++ b/src/orchestration/taskParser.coverage.test.ts @@ -10,6 +10,16 @@ vi.mock('fs/promises', () => ({ readFile: vi.fn(), })); +const fileLockMock = vi.hoisted(() => ({ + withFileLock: vi.fn((_path: string, operation: () => unknown) => operation()), +})); +vi.mock('../support/fileLock.js', () => fileLockMock); + +const atomicFileMock = vi.hoisted(() => ({ + atomicWriteFile: vi.fn(), +})); +vi.mock('../support/atomicFile.js', () => atomicFileMock); + import * as fs from 'fs/promises'; import { formatParsedTaskSummary, @@ -185,6 +195,8 @@ describe('parsedTaskFilePath validation (via saveParsedTask/loadParsedTask)', () vi.mocked(fs.mkdir).mockReset().mockResolvedValue(undefined as never); vi.mocked(fs.writeFile).mockReset().mockResolvedValue(undefined as never); vi.mocked(fs.readFile).mockReset(); + fileLockMock.withFileLock.mockImplementation((_path: string, operation: () => unknown) => operation()); + atomicFileMock.atomicWriteFile.mockReset().mockResolvedValue(undefined); }); it.each([ @@ -209,7 +221,8 @@ describe('parsedTaskFilePath validation (via saveParsedTask/loadParsedTask)', () await saveParsedTask(parsed); expect(fs.mkdir).toHaveBeenCalledWith(PARSED_TASKS_DIR, { recursive: true }); - expect(fs.writeFile).toHaveBeenCalledWith( + expect(fileLockMock.withFileLock).toHaveBeenCalled(); + expect(atomicFileMock.atomicWriteFile).toHaveBeenCalledWith( resolve(PARSED_TASKS_DIR, 'INT-209.json'), JSON.stringify(parsed, null, 2), ); @@ -235,6 +248,14 @@ describe('parsedTaskFilePath validation (via saveParsedTask/loadParsedTask)', () expect(result).toBeNull(); }); + + it('resolves to null when stored JSON fails schema validation', async () => { + vi.mocked(fs.readFile).mockResolvedValue(JSON.stringify({ invalid: true }) as never); + + const result = await loadParsedTask('INT-212'); + + expect(result).toBeNull(); + }); }); describe('formatParsedTaskSummary', () => { diff --git a/src/orchestration/taskParser.ts b/src/orchestration/taskParser.ts index ecb2bb4f..d673a6bf 100644 --- a/src/orchestration/taskParser.ts +++ b/src/orchestration/taskParser.ts @@ -6,6 +6,9 @@ import { basename, isAbsolute, relative, resolve } from 'path'; import { homedir } from 'os'; import * as fs from 'fs/promises'; +import { z } from 'zod'; +import { withFileLock } from '../support/fileLock.js'; +import { atomicWriteFile } from '../support/atomicFile.js'; import { WorkflowConfig, WorkflowStep } from './workflow.js'; // Types @@ -622,6 +625,53 @@ function subtasksToWorkflow( const PARSED_TASKS_DIR = resolve(homedir(), '.openswarm/parsed-tasks'); +const SubtaskSchema = z.object({ + id: z.string(), + order: z.number(), + title: z.string(), + description: z.string(), + prompt: z.string(), + dependsOn: z.array(z.string()), + type: z.enum(['analysis', 'implementation', 'test', 'review', 'documentation']), + optional: z.boolean(), +}); + +const WorkflowConfigSchema = z.object({ + id: z.string(), + name: z.string(), + description: z.string().optional(), + projectPath: z.string(), + steps: z.array(z.object({ + id: z.string(), + name: z.string(), + prompt: z.string(), + dependsOn: z.array(z.string()).optional(), + onFailure: z.enum(['rollback', 'retry', 'skip', 'abort', 'notify']).optional(), + }).passthrough()), + onFailure: z.enum(['rollback', 'retry', 'skip', 'abort', 'notify']).optional(), + trigger: z.object({}).passthrough().optional(), + linearIssue: z.string().optional(), + tags: z.array(z.string()).optional(), +}).passthrough(); + +const ParsedTaskSchema = z.object({ + original: z.object({ + id: z.string(), + title: z.string(), + description: z.string(), + }), + analysis: z.object({ + type: z.enum(['bug_fix', 'feature', 'refactor', 'docs', 'test', 'ci_cd', 'investigation', 'unknown']), + complexity: z.enum(['simple', 'medium', 'complex']), + estimatedSteps: z.number(), + requiresHumanReview: z.boolean(), + risks: z.array(z.string()), + }), + subtasks: z.array(SubtaskSchema), + workflow: WorkflowConfigSchema, + parsedAt: z.number(), +}); + function parsedTaskFilePath(issueId: string): string { if ( !issueId || @@ -648,9 +698,12 @@ function parsedTaskFilePath(issueId: string): string { * Save parsed result */ export async function saveParsedTask(parsed: ParsedTask): Promise { + const validated = ParsedTaskSchema.parse(parsed); await fs.mkdir(PARSED_TASKS_DIR, { recursive: true }); - const filePath = parsedTaskFilePath(parsed.original.id); - await fs.writeFile(filePath, JSON.stringify(parsed, null, 2)); + const filePath = parsedTaskFilePath(validated.original.id); + await withFileLock(filePath + '.lock', async () => { + await atomicWriteFile(filePath, JSON.stringify(validated, null, 2)); + }); } /** @@ -660,7 +713,9 @@ export async function loadParsedTask(issueId: string): Promise { + if (finalized) return; + finalized = true; + activeTasks.delete(taskId); + try { + onComplete?.(resultText, code); + } catch (callbackError) { + console.error(`[Dev] onComplete failed for ${taskId}:`, callbackError); + } + }; + // Collect stdout claudeProcess.stdout?.on('data', (data: Buffer) => { const chunk = data.toString(); devTask.output += chunk; - onProgress?.(chunk); + try { + onProgress?.(chunk); + } catch (callbackError) { + console.error(`[Dev] onProgress failed for ${taskId}:`, callbackError); + } }); // Collect stderr @@ -226,20 +246,18 @@ export async function runDevTask( // Generate report file const duration = Math.floor((Date.now() - devTask.startedAt) / 1000); generateReport(devTask, code, duration); - } catch (err) { - console.error(`[Dev] close-handler reporting failed for ${taskId}:`, err); + } catch (reportError) { + // Reporting must not decide whether the task is finalized. + console.error(`[Dev] close-handler reporting failed for ${taskId}:`, reportError); } finally { - // Cleanup must run even when reporting throws - onComplete?.(resultText, code); - activeTasks.delete(taskId); + finalize(resultText, code); } }); // Handle errors claudeProcess.on('error', (err) => { devTask.output += `\nError: ${err.message}`; - onComplete?.(devTask.output, -1); - activeTasks.delete(taskId); + finalize(devTask.output, -1); }); return { taskId, path }; @@ -265,8 +283,13 @@ export function cancelTask(taskId: string): boolean { const task = activeTasks.get(taskId); if (!task) return false; - task.process.kill('SIGTERM'); - activeTasks.delete(taskId); + try { + task.process.kill('SIGTERM'); + } catch { + // Already gone; the close handler still finalizes and frees the slot. + } + // The entry stays until 'close' — getActiveTasks must not advertise the repo + // as free while the SIGTERM'd child is still running. return true; } diff --git a/src/support/fileLock.test.ts b/src/support/fileLock.test.ts index e64e2266..f3da4c69 100644 --- a/src/support/fileLock.test.ts +++ b/src/support/fileLock.test.ts @@ -1,10 +1,10 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; -import { chmodSync, mkdtempSync, readFileSync, statSync, utimesSync, writeFileSync, existsSync } from 'node:fs'; +import { chmodSync, existsSync, mkdirSync, mkdtempSync, readFileSync, statSync, utimesSync, writeFileSync } from 'node:fs'; import { rm } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; -import { withFileLock } from './fileLock.js'; +import { withFileLock, withFileLockSync } from './fileLock.js'; /** A pid that cannot be running: above the platform maximum. */ const DEAD_PID = 0x7fffffff; @@ -165,3 +165,43 @@ describe('withFileLock stale takeover', () => { .rejects.toThrow(/Timed out waiting for file lock/); }); }); + +describe('withFileLockSync', () => { + let dir: string; + let lockPath: string; + + beforeEach(() => { + dir = mkdtempSync(join(tmpdir(), 'openswarm-filelock-sync-')); + lockPath = join(dir, 'nested', 'resource.lock'); + }); + + afterEach(async () => { + await rm(dir, { recursive: true, force: true }); + }); + + it('runs the operation under a held lock and releases afterward', () => { + const result = withFileLockSync(lockPath, () => { + expect(existsSync(lockPath)).toBe(true); + expect(JSON.parse(readFileSync(lockPath, 'utf8'))).toMatchObject({ pid: process.pid }); + return 42; + }); + expect(result).toBe(42); + expect(existsSync(lockPath)).toBe(false); + }); + + it('releases the lock when the sync operation throws', () => { + expect(() => withFileLockSync(lockPath, () => { + throw new Error('sync boom'); + })).toThrow('sync boom'); + expect(existsSync(lockPath)).toBe(false); + }); + + it('takes over a lock left behind by a dead process', () => { + // mkdir the parent first: the lock path is nested, and the dead owner's lock + // is written directly rather than acquired through withFileLockSync. + mkdirSync(join(dir, 'nested'), { recursive: true }); + writeFileSync(lockPath, JSON.stringify({ pid: DEAD_PID, token: 'dead' }), { mode: 0o600 }); + expect(withFileLockSync(lockPath, () => 'taken', { timeoutMs: 500 })).toBe('taken'); + expect(existsSync(lockPath)).toBe(false); + }); +}); diff --git a/src/support/fileLock.ts b/src/support/fileLock.ts index 71b93290..431c144e 100644 --- a/src/support/fileLock.ts +++ b/src/support/fileLock.ts @@ -1,9 +1,29 @@ import { randomUUID } from 'node:crypto'; +import { + closeSync, + fsyncSync, + mkdirSync, + openSync, + readFileSync, + statSync, + unlinkSync, + writeFileSync, +} from 'node:fs'; import { mkdir, open, readFile, stat, unlink } from 'node:fs/promises'; import { dirname } from 'node:path'; type LockOwner = { pid: number; token: string }; +const lockWaitBuffer = new Int32Array(new SharedArrayBuffer(4)); + +/** + * The real setTimeout, captured at module load. Lock waits are real-time + * scheduling, not application timers: a test that fakes setTimeout to skip a + * retry backoff must not also freeze lock hand-off, which would hang the + * acquire loop until the test's own timeout. + */ +const lockWaitTimer = globalThis.setTimeout; + function alive(pid: number): boolean { try { process.kill(pid, 0); @@ -13,9 +33,9 @@ function alive(pid: number): boolean { } } -async function owner(path: string): Promise { +function parseOwner(raw: string): LockOwner | null { try { - const value = JSON.parse(await readFile(path, 'utf8')) as Partial; + const value = JSON.parse(raw) as Partial; return Number.isInteger(value.pid) && (value.pid ?? 0) > 0 && typeof value.token === 'string' ? { pid: value.pid!, token: value.token } : null; @@ -24,6 +44,74 @@ async function owner(path: string): Promise { } } +async function owner(path: string): Promise { + try { + return parseOwner(await readFile(path, 'utf8')); + } catch { + return null; + } +} + +function ownerSync(path: string): LockOwner | null { + try { + return parseOwner(readFileSync(path, 'utf8')); + } catch { + return null; + } +} + +/** + * Synchronous counterpart of the async reclaim below, sharing its ownership + * judgement: only the exact lock observed as stale may be removed, so a live + * holder that took the file between judgement and unlink is left alone. + * + * Returns true when the caller should retry the acquire loop. + */ +function reclaimStaleLockSync(path: string, current: LockOwner | null, malformedStaleMs: number): boolean { + let judgedMtimeMs: number | undefined; + let malformedAndStale = false; + if (current === null) { + try { + judgedMtimeMs = statSync(path).mtimeMs; + malformedAndStale = Date.now() - judgedMtimeMs > malformedStaleMs; + } catch (statError) { + // The holder released the lock between our failed open and this stat. + // That is the normal hand-off, not an error: retry the open. + if ((statError as NodeJS.ErrnoException).code !== 'ENOENT') throw statError; + return true; + } + } else { + try { + judgedMtimeMs = statSync(path).mtimeMs; + } catch (statError) { + if ((statError as NodeJS.ErrnoException).code !== 'ENOENT') throw statError; + return true; + } + } + // Reclaim only the lock we judged — same reasoning as the async path. + if (((current !== null && !alive(current.pid)) || malformedAndStale) && judgedMtimeMs !== undefined) { + const judgedToken = current?.token; + try { + const currentOwner = ownerSync(path); + const currentMtimeMs = statSync(path).mtimeMs; + if (currentMtimeMs === judgedMtimeMs && currentOwner?.token === judgedToken) { + try { + unlinkSync(path); + } catch (unlinkError) { + const code = (unlinkError as NodeJS.ErrnoException).code; + // ENOENT: another reclaim won. ENOTEMPTY: directory-style locks with + // concurrent claim markers — retry the acquire loop after a wait. + if (code !== 'ENOENT' && code !== 'ENOTEMPTY') throw unlinkError; + } + } + } catch (statError) { + if ((statError as NodeJS.ErrnoException).code !== 'ENOENT') throw statError; + } + return true; + } + return false; +} + export async function withFileLock( path: string, operation: () => Promise, @@ -77,7 +165,10 @@ export async function withFileLock( && currentOwner?.token === judgedToken; if (sameLock) { await unlink(path).catch((unlinkError) => { - if ((unlinkError as NodeJS.ErrnoException).code !== 'ENOENT') throw unlinkError; + const code = (unlinkError as NodeJS.ErrnoException).code; + // ENOENT: another reclaim raced us. ENOTEMPTY: directory-style + // locks carrying concurrent claim markers — retry the loop. + if (code !== 'ENOENT' && code !== 'ENOTEMPTY') throw unlinkError; }); } } catch (statError) { @@ -86,13 +177,14 @@ export async function withFileLock( continue; } if (Date.now() >= deadline) throw new Error(`Timed out waiting for file lock: ${path}`); - await new Promise((resolve) => setTimeout(resolve, 10)); + await new Promise((resolve) => lockWaitTimer(resolve, 10)); } } try { return await operation(); } finally { + // Ownership-safe release: only the token that created this lock may unlink. if ((await owner(path))?.token === token) { await unlink(path).catch((error) => { if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; @@ -100,3 +192,53 @@ export async function withFileLock( } } } + +/** + * Synchronous counterpart for callers that must stay sync (durable runner state, + * auth store writes). Same ownership rules as {@link withFileLock}; the wait is + * an Atomics wait on a shared buffer rather than a busy loop. + */ +export function withFileLockSync( + path: string, + operation: () => T, + options: { timeoutMs?: number; malformedStaleMs?: number } = {}, +): T { + const timeoutMs = options.timeoutMs ?? 5_000; + const malformedStaleMs = options.malformedStaleMs ?? 30_000; + const deadline = Date.now() + timeoutMs; + const token = randomUUID(); + mkdirSync(dirname(path), { recursive: true }); + + for (;;) { + let fd: number | undefined; + try { + fd = openSync(path, 'wx', 0o600); + writeFileSync(fd, JSON.stringify({ pid: process.pid, token }), 'utf8'); + fsyncSync(fd); + closeSync(fd); + fd = undefined; + break; + } catch (error) { + if (fd !== undefined) closeSync(fd); + if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error; + if (reclaimStaleLockSync(path, ownerSync(path), malformedStaleMs)) { + Atomics.wait(lockWaitBuffer, 0, 0, 10); + continue; + } + if (Date.now() >= deadline) throw new Error(`Timed out waiting for file lock: ${path}`); + Atomics.wait(lockWaitBuffer, 0, 0, 10); + } + } + + try { + return operation(); + } finally { + if (ownerSync(path)?.token === token) { + try { + unlinkSync(path); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; + } + } + } +} diff --git a/src/support/logRotation.test.ts b/src/support/logRotation.test.ts index eba9c9e8..87af835c 100644 --- a/src/support/logRotation.test.ts +++ b/src/support/logRotation.test.ts @@ -1,9 +1,31 @@ import { existsSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; -import { afterEach, describe, expect, it } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import Database from 'better-sqlite3'; -import { rotateServiceLogs } from './logRotation.js'; + +// Named `fsyncSync` imports in logRotation.ts are closed over at load time; spyOn +// on the namespace is unreliable under Vitest ESM. Intercept via vi.mock instead. +const fsyncControl = vi.hoisted(() => ({ + failOnCall: null as number | null, + calls: 0, +})); + +vi.mock('node:fs', async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + fsyncSync: (fd: number) => { + fsyncControl.calls += 1; + if (fsyncControl.failOnCall !== null && fsyncControl.calls === fsyncControl.failOnCall) { + throw Object.assign(new Error('fsync failed'), { code: 'EIO' }); + } + return actual.fsyncSync(fd); + }, + }; +}); + +const { rotateServiceLogs } = await import('./logRotation.js'); const roots: string[] = []; function logDir(): string { @@ -12,6 +34,11 @@ function logDir(): string { return root; } +beforeEach(() => { + fsyncControl.failOnCall = null; + fsyncControl.calls = 0; +}); + afterEach(() => { for (const root of roots.splice(0)) rmSync(root, { recursive: true, force: true }); }); @@ -57,4 +84,17 @@ describe('rotateServiceLogs', () => { expect(rotateServiceLogs({ logDir: dir, maxBytes: 4 })) .toEqual({ rotated: ['stdout.log'], skippedLocked: false }); }); + + it('restores the active log from the staged archive when post-truncate fsync fails', () => { + const dir = logDir(); + const logPath = join(dir, 'stdout.log'); + const precious = 'precious-log-data-that-must-survive'; + writeFileSync(logPath, precious); + + // First fsync is the archive; second is the active fd after truncate. + fsyncControl.failOnCall = 2; + + expect(() => rotateServiceLogs({ logDir: dir, maxBytes: 4 })).toThrow(/fsync failed|EIO/); + expect(readFileSync(logPath, 'utf8')).toBe(precious); + }); }); diff --git a/src/support/logRotation.ts b/src/support/logRotation.ts index 6c07d121..e134afc2 100644 --- a/src/support/logRotation.ts +++ b/src/support/logRotation.ts @@ -12,6 +12,7 @@ import { readFileSync, unlinkSync, writeFileSync, + writeSync, } from 'node:fs'; import { randomUUID } from 'node:crypto'; import { homedir } from 'node:os'; @@ -88,8 +89,28 @@ function rotateOne(logPath: string, maxBytes: number, generations: number): bool // launchd opens stdout/stderr before Node starts. Truncating the already // opened descriptor keeps that inode (and inherited descriptors) valid. - ftruncateSync(activeFd, 0); - fsyncSync(activeFd); + // If truncate/fsync fails after the archive rename, restore the active + // log from the staged `.1` so the rotation does not leave an empty file + // while the only copy of the data sits in an incomplete commit. + try { + ftruncateSync(activeFd, 0); + fsyncSync(activeFd); + } catch (commitError) { + try { + // activeFd's offset may still sit at the pre-truncate EOF (readFileSync + // above advanced it); write at absolute position 0 or restore is sparse. + const archived = readFileSync(`${logPath}.1`); + ftruncateSync(activeFd, 0); + writeSync(activeFd, archived, 0, archived.byteLength, 0); + fsyncSync(activeFd); + } catch (restoreError) { + throw new AggregateError( + [commitError, restoreError], + `Log rotation commit failed and restore from ${logPath}.1 also failed`, + ); + } + throw commitError; + } return true; } catch (error) { safeUnlink(temporary); diff --git a/src/taskState/store.test.ts b/src/taskState/store.test.ts index bd6ef177..5ec2696d 100644 --- a/src/taskState/store.test.ts +++ b/src/taskState/store.test.ts @@ -184,6 +184,36 @@ describe('task state store', () => { expect(getTaskState('PROCESS-B')?.title).toBe('PROCESS-B'); }); + it('lets simultaneous stale-lock reclaimers both proceed', async () => { + const lockPath = `${stateFile}.lock`; + // Expired abandoned lock: mtime past LOCK_ABANDON_MS so expiredLock fires + // regardless of whether pid 999999 is judgeable in this namespace. + writeFileSync(lockPath, JSON.stringify({ + pid: 999_999, + token: 'stale', + ns: 'other-space', + instance: 'dead', + })); + const longAgo = new Date(Date.now() - 900_000); + utimesSync(lockPath, longAgo, longAgo); + + const fixture = fileURLToPath(new URL('./storeClaimProcess.fixture.ts', import.meta.url)); + const run = (issueId: string) => new Promise((resolve, reject) => { + const child = spawn(process.execPath, ['--import', 'tsx', fixture, stateFile, issueId, '0'], { + stdio: 'pipe', + }); + let stderr = ''; + child.stderr.on('data', (c) => { stderr += String(c); }); + child.on('error', reject); + child.on('exit', (code) => code === 0 ? resolve() : reject(new Error(stderr || `exit ${code}`))); + }); + + await Promise.all([run('RECLAIM-A'), run('RECLAIM-B')]); + resetTaskStateStoreForTests(); + expect(getTaskState('RECLAIM-A')?.title).toBe('RECLAIM-A'); + expect(getTaskState('RECLAIM-B')?.title).toBe('RECLAIM-B'); + }); + it('keeps concurrent execution upsert and Linear reconciliation consistent under the store lock', async () => { // Seed an in_progress row, then race a local execution bump against a Linear // Done reconciliation. Both paths take withStoreLock; the store must remain diff --git a/src/taskState/store.ts b/src/taskState/store.ts index c0cabb15..fc8b6688 100644 --- a/src/taskState/store.ts +++ b/src/taskState/store.ts @@ -304,7 +304,17 @@ function withStoreLock(operation: () => T): T { const currentOwner = readStoreLockOwner(lockPath); const sameLock = currentMtimeMs === judgedMtimeMs && currentOwner?.token === owner?.token; - if (sameLock) unlinkSync(lockPath); + if (sameLock) { + try { + unlinkSync(lockPath); + } catch (unlinkError) { + const code = (unlinkError as NodeJS.ErrnoException).code; + // Another reclaim raced (ENOENT) or claim markers make a dir lock + // non-empty (ENOTEMPTY) — retry the acquire loop instead of + // failing the whole store write. + if (code !== 'ENOENT' && code !== 'ENOTEMPTY') throw unlinkError; + } + } continue; } } catch (statError) { @@ -561,6 +571,43 @@ export function markTaskInProgress( }); } +/** + * Claim a task for execution in one locked read-modify-write, or return null + * when another actor already holds it. + * + * `markTaskInProgress` alone is not an admission check: two runners (or two + * daemon instances sharing one state file) both read `backlog`, both write + * `in_progress`, and the same issue is executed twice. Deciding and writing + * under one store lock is what makes the claim exclusive. + */ +export function tryClaimTaskAdmission( + issueId: string, + patch: Parameters[1] = {}, +): OpenSwarmTaskState | null { + return withStoreLock(() => { + const store = ensureStoreLoaded(); + if (store.tasks[issueId]?.execution.status === 'in_progress') return null; + const result = upsertTaskStateUnlocked(issueId, { + issueIdentifier: patch.issueIdentifier, + title: patch.title, + projectId: patch.projectId, + projectName: patch.projectName, + linearState: patch.linearState ?? 'In Progress', + execution: { + status: 'in_progress', + blockedReason: undefined, + retryCount: 0, + lastSessionId: patch.sessionId, + }, + worktree: { + branchName: patch.branchName, + worktreePath: patch.worktreePath, + }, + }); + return result; + }); +} + export function markTaskBacklog( issueId: string, patch: { From ca2805adb809380d8e2ad26cda3e33b2f9474550 Mon Sep 17 00:00:00 2001 From: SalvageA Date: Mon, 28 Sep 2026 16:10:12 +0900 Subject: [PATCH 2/2] fix(concurrency): make lifecycle and durable state transitions ownership-safe MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Salvages the worthwhile work from draft PRs #766 and #767 (both 'failed 4 times … very likely incomplete', 174-247 commits behind) into one change, rebased onto current main. Kept: unified withFileLockSync (Atomics.wait + ownership-safe reclaim + ENOTEMPTY tolerance); AsyncLocalStorage per-run RunControl in pairPipeline; processRegistry taskId/spawnedAt ownership + cancellable unref'd timer; taskState tryClaimTaskAdmission + ENOTEMPTY reclaim; decisionEngine admission claims + in_progress filter; taskParser zod schemas; cross-process locks in memoryCore/codex/reembed/runnerState; gitInfo enrichment now persists via graph.addNode+saveGraph; locale AsyncLocalStorage scope; oauthStore refresh serialized on a store lock with reload-under-lock. Also updates the duck-typed AuthProfileStore double in codexResponses.test.ts (reloadProfileFromDisk/setProfileUnlocked) so it implements the interface its declared collaborator now requires — the adapter's tests are the only caller that is not a real AuthProfileStore instance. --- src/adapters/codexResponses.test.ts | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/src/adapters/codexResponses.test.ts b/src/adapters/codexResponses.test.ts index 38bfcdac..b988d720 100644 --- a/src/adapters/codexResponses.test.ts +++ b/src/adapters/codexResponses.test.ts @@ -664,6 +664,12 @@ describe('401 refresh scope', () => { // refreshAndRetry expires the token through the store rather than writing // a snapshot back, so the fake has to offer the same door. (INT-2961) expireProfile: (_k: string) => { profile = { ...profile, expires: 0 }; return true; }, + // The refresh path now runs under the store lock: it re-reads the profile + // from disk under that lock and persists without re-entering it + // (AuthProfileStore.saveUnlocked/setProfileUnlocked/reloadProfileFromDisk). + // The fake is in-memory, so both simply read/write the one profile. + reloadProfileFromDisk: (_k: string) => profile, + setProfileUnlocked: (_k: string, p: typeof profile) => { profile = p; }, }; }