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..3f429c99 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,10 +67,18 @@ 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); @@ -105,61 +113,76 @@ export function getAllProcesses(): ProcessInfo[] { return Array.from(registry.values()); } +const KILL_ESCALATE_MS = 5_000; + /** - * 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; + 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(() => { + if (stillOurs()) terminateCliProcessTree(proc); + }, KILL_ESCALATE_MS); + proc.once('close', () => { + clearTimeout(timer); + }); + 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 +211,4 @@ export function stopHealthChecker(): void { clearInterval(healthCheckTimer); healthCheckTimer = null; } -} +} \ No newline at end of file diff --git a/src/agents/pairPipeline.ts b/src/agents/pairPipeline.ts index 983500aa..6864b66e 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 { taskEventKey, type TaskItem } from '../orchestration/decisionEngine.js'; import { enforcedFileScope } from '../orchestration/writeScope.js'; import type { WorkerResult, ReviewResult } from './agentPair.js'; @@ -77,6 +78,7 @@ export { buildTaskPrefix } from './pipelineTaskPrefix.js'; export { stageTimeoutMs } from './stageTimeouts.js'; import { stageTimeoutMs } from './stageTimeouts.js'; +type RunControl = { signal?: AbortSignal; stuck: StuckDetector }; /** * Resolve the coordination identity for one stage of a task. @@ -89,14 +91,18 @@ import { stageTimeoutMs } from './stageTimeouts.js'; */ export class PairPipeline extends EventEmitter { private config: PipelineConfig; + /** Fallback for callers outside run(); each run() installs a fresh detector 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) { @@ -129,13 +135,11 @@ 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 { + 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)) { @@ -168,6 +172,8 @@ export class PairPipeline extends EventEmitter { currentIteration: 0, taskPrefix, reflection: createReflectionState(), + abortSignal: opts?.signal, + stuckDetector, }; try { if (this.config.verify?.enabled) try { @@ -215,7 +221,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) @@ -260,6 +266,7 @@ export class PairPipeline extends EventEmitter { }, }; } + }); } /** * Worker에 주입할 코드 컨텍스트 수집 @@ -274,7 +281,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)}`); } } @@ -415,7 +422,7 @@ export class PairPipeline extends EventEmitter { onLog, processContext: { taskId: taskEventKey(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, @@ -512,7 +519,7 @@ export class PairPipeline extends EventEmitter { type: 'log', data: { taskId: taskEventKey(context.task), stage: 'reviewer', line: `[${prefix}] ${line}` }, }), - signal: this.abortSignal, + signal: PairPipeline.runControl.getStore()?.signal, instructionCapsule: this.config.instructionCapsule, mcpTools: this.config.roleMcpTools?.reviewer, coordinationContext: coordinationContextFor(context, 'reviewer'), @@ -777,7 +784,7 @@ export class PairPipeline extends EventEmitter { context.currentIteration++; // 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}`); @@ -833,7 +840,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, @@ -1143,7 +1150,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 93169333..aefef4f1 100644 --- a/src/agents/pairPipelineTypes.ts +++ b/src/agents/pairPipelineTypes.ts @@ -179,6 +179,9 @@ export interface PipelineContext { newSecurityFindings?: import('../verify/securityAudit.js').SecurityFinding[]; /** Pair-level stagnation detector reason, preserved so the scheduler does not rerun the same loop. */ stuckReason?: string; + 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/prProcessor.ts b/src/automation/prProcessor.ts index dcd8c260..f8e08126 100644 --- a/src/automation/prProcessor.ts +++ b/src/automation/prProcessor.ts @@ -13,6 +13,7 @@ import { randomUUID } from 'node:crypto'; import { promisify } from 'node:util'; import { z } from 'zod'; import { atomicWriteFileSync } from '../support/atomicFile.js'; +import { withFileLock } from '../support/fileLock.js'; import { safeConsole as console } from '../support/safeLog.js'; const execFileAsync = promisify(execFile); @@ -249,6 +250,7 @@ const PRStateSchema = z.object({ // Constants const PR_STATE_PATH = resolve(homedir(), '.openswarm', 'pr-state.json'); +const PR_STATE_LOCK = `${PR_STATE_PATH}.lock`; // PR Processor @@ -367,16 +369,29 @@ export class PRProcessor { projectPath: string, ): Promise<{ success: boolean; error?: string; iterations: number }> { const key = `${pr.repo}#${pr.number}`; - const state = await this.loadState(); - state.prs[key] = { - ...state.prs[key], - repo: pr.repo, - prNumber: pr.number, - status: 'processing', - iterations: 0, - }; + // Lease only the durable RMW bookends — not the long pipeline — so cron + // saveState can still proceed while review feedback is being addressed. + const state = await this.withStateLock(async () => { + const current = await this.loadState(); + current.prs[key] = { + ...current.prs[key], + repo: pr.repo, + prNumber: pr.number, + status: 'processing', + iterations: 0, + }; + await this.saveStateUnlocked(current); + return current; + }); + await this.processReviewFeedback(pr, projectPath, state, key, 0); - await this.saveState(state); + + await this.withStateLock(async () => { + const disk = await this.loadState(); + disk.prs[key] = state.prs[key]; + await this.saveStateUnlocked(disk); + }); + const entry = state.prs[key]; return { success: entry?.status === 'completed', @@ -1477,6 +1492,10 @@ export class PRProcessor { // State Persistence // ============================================ + private async withStateLock(operation: () => Promise): Promise { + return withFileLock(PR_STATE_LOCK, operation, { timeoutMs: 60_000 }); + } + private async loadState(): Promise { try { const data = await readFile(PR_STATE_PATH, 'utf-8'); @@ -1489,8 +1508,15 @@ export class PRProcessor { } } - private async saveState(state: PRState): Promise { + /** Atomic write only — caller must already hold {@link PR_STATE_LOCK}. */ + private async saveStateUnlocked(state: PRState): Promise { state.updatedAt = new Date().toISOString(); atomicWriteFileSync(PR_STATE_PATH, `${JSON.stringify(state, null, 2)}\n`); } + + private async saveState(state: PRState): Promise { + await withFileLock(PR_STATE_LOCK, async () => { + await this.saveStateUnlocked(state); + }, { timeoutMs: 15_000 }); + } } diff --git a/src/automation/runnerState.ts b/src/automation/runnerState.ts index d60cf81d..44280b09 100644 --- a/src/automation/runnerState.ts +++ b/src/automation/runnerState.ts @@ -9,6 +9,7 @@ import { join, dirname, isAbsolute, relative, sep } from 'node:path'; import { taskEventKey, type TaskItem } from '../orchestration/decisionEngine.js'; import type { PipelineResult } from '../agents/pairPipelineTypes.js'; import { atomicWriteFileSync } from '../support/atomicFile.js'; +import { withFileLockSync } from '../support/fileLock.js'; /** * Write-temp-then-rename instead of an in-place write, so a crash mid-write (or @@ -684,14 +685,21 @@ 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`, () => { + let history: PipelineHistoryEntry[] = []; + try { + if (existsSync(PIPELINE_HISTORY_FILE)) { + history = JSON.parse(readFileSync(PIPELINE_HISTORY_FILE, 'utf8')) as PipelineHistoryEntry[]; + if (!Array.isArray(history)) history = []; + } + } catch { history = []; } + history.unshift(entry); + 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/memory/compaction.ts b/src/memory/compaction.ts index 2a88f5f5..73a60f90 100644 --- a/src/memory/compaction.ts +++ b/src/memory/compaction.ts @@ -2,7 +2,16 @@ // OpenSwarm - Memory Compaction // ============================================ -import { getDb, getTable, initDatabase, EMBEDDING_DIM, PERMANENT_EXPIRY, normalizeRecords, setTable } from './memoryCore.js'; +import { + getDb, + getTable, + initDatabase, + EMBEDDING_DIM, + PERMANENT_EXPIRY, + normalizeRecords, + setTable, + withMemoryMutationLock, +} from './memoryCore.js'; import type { CognitiveMemoryRecord } from './memoryCore.js'; import { isTransientReviewRejectionMemory } from './memoryFilters.js'; @@ -112,105 +121,106 @@ export async function compactMemoryTable(): Promise<{ console.log('[Compaction] Starting memory table compaction...'); try { - await initDatabase(); - const table = getTable(); - const db = getDb(); - - if (!table || !db) { - console.error('[Compaction] Database not initialized'); - return { before: 0, after: 0, removed: 0, deduplicated: 0 }; - } - - // 1. Read all records - const queryLimit = 100_000; - const allRecords = await table - .search(Array.from({ length: EMBEDDING_DIM }, () => 0)) - .limit(queryLimit) - .toArray(); + return await withMemoryMutationLock(async () => { + await initDatabase(); + const table = getTable(); + const db = getDb(); + + if (!table || !db) { + console.error('[Compaction] Database not initialized'); + return { before: 0, after: 0, removed: 0, deduplicated: 0 }; + } - if (allRecords.length >= queryLimit) { - throw new Error(`Memory compaction refused: query reached the ${queryLimit}-row safety limit`); - } + // 1. Read all records + const queryLimit = 100_000; + const allRecords = await table + .search(Array.from({ length: EMBEDDING_DIM }, () => 0)) + .limit(queryLimit) + .toArray(); - const beforeCount = allRecords.length; - console.log(`[Compaction] Found ${beforeCount} records`); + if (allRecords.length >= queryLimit) { + throw new Error(`Memory compaction refused: query reached the ${queryLimit}-row safety limit`); + } - if (beforeCount === 0) { - console.log('[Compaction] No records to compact'); - return { before: 0, after: 0, removed: 0, deduplicated: 0 }; - } + const beforeCount = allRecords.length; + console.log(`[Compaction] Found ${beforeCount} records`); - // 2. Filter valid records - const now = Date.now(); - const validRecords = allRecords.filter((r: any) => { - if (r.id === 'init') return true; + if (beforeCount === 0) { + console.log('[Compaction] No records to compact'); + return { before: 0, after: 0, removed: 0, deduplicated: 0 }; + } - // Remove transient infrastructure failures that were previously stored as - // high-importance reviewer constraints. - if (isTransientReviewRejectionMemory(r)) return false; + // 2. Filter valid records + const now = Date.now(); + const validRecords = allRecords.filter((r: any) => { + if (r.id === 'init') return true; - // Remove if expired - if (r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) return false; + // Remove transient infrastructure failures that were previously stored as + // high-importance reviewer constraints. + if (isTransientReviewRejectionMemory(r)) return false; - // Remove if unimportant - if (r.importance < MIN_IMPORTANCE) return false; + // Remove if expired + if (r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) return false; - return true; - }); + // Remove if unimportant + if (r.importance < MIN_IMPORTANCE) return false; - const afterFilter = validRecords.length; - console.log(`[Compaction] After filtering: ${afterFilter} records (removed ${beforeCount - afterFilter})`); + return true; + }); - // 3. Deduplicate - const deduplicated = removeDuplicates(validRecords as CognitiveMemoryRecord[]); - const afterDedup = deduplicated.length; - console.log(`[Compaction] After deduplication: ${afterDedup} records (merged ${afterFilter - afterDedup})`); + const afterFilter = validRecords.length; + console.log(`[Compaction] After filtering: ${afterFilter} records (removed ${beforeCount - afterFilter})`); - // 4. Validate replacement before touching the live table - const normalized = normalizeRecords(deduplicated); - const targetTableName = table.name; - const tempTableName = `${targetTableName}_compact_${Date.now()}`; + // 3. Deduplicate + const deduplicated = removeDuplicates(validRecords as CognitiveMemoryRecord[]); + const afterDedup = deduplicated.length; + console.log(`[Compaction] After deduplication: ${afterDedup} records (merged ${afterFilter - afterDedup})`); - console.log(`[Compaction] Creating validated replacement for ${targetTableName}...`); - if (normalized.length > 0) { - await db.createTable(tempTableName, normalized); - } else { - await db.createEmptyTable(tempTableName, await table.schema()); - } + // 4. Validate replacement before touching the live table + const normalized = normalizeRecords(deduplicated); + const targetTableName = table.name; + const tempTableName = `${targetTableName}_compact_${Date.now()}`; - let replaced = false; - try { - console.log(`[Compaction] Replacing ${targetTableName} with compacted data...`); + console.log(`[Compaction] Creating validated replacement for ${targetTableName}...`); 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()); } - const newTable = await db.openTable(targetTableName); - setTable(newTable); - replaced = true; - } finally { - if (replaced) { - try { - await db.dropTable(tempTableName); - } catch (cleanupError) { - console.warn(`[Compaction] Failed to drop temporary table ${tempTableName}:`, cleanupError); + + let replaced = false; + try { + console.log(`[Compaction] Replacing ${targetTableName} with compacted data...`); + if (normalized.length > 0) { + await db.createTable(targetTableName, normalized, { mode: 'overwrite' }); + } else { + await db.createEmptyTable(targetTableName, await table.schema(), { mode: 'overwrite' }); + } + const newTable = await db.openTable(targetTableName); + setTable(newTable); + replaced = true; + } finally { + if (replaced) { + try { + await db.dropTable(tempTableName); + } catch (cleanupError) { + console.warn(`[Compaction] Failed to drop temporary table ${tempTableName}:`, cleanupError); + } + } else { + console.warn(`[Compaction] Replacement failed; retained recoverable table ${tempTableName}`); } - } else { - console.warn(`[Compaction] Replacement failed; retained recoverable table ${tempTableName}`); } - } - const stats = { - before: beforeCount, - after: afterDedup, - removed: beforeCount - afterDedup, - deduplicated: afterFilter - afterDedup, - }; - - console.log('[Compaction] Complete:', stats); - return stats; + const stats = { + before: beforeCount, + after: afterDedup, + removed: beforeCount - afterDedup, + deduplicated: afterFilter - afterDedup, + }; + console.log('[Compaction] Complete:', stats); + return stats; + }); } catch (error) { console.error('[Compaction] Failed:', error); throw error; diff --git a/src/memory/memoryCore.ts b/src/memory/memoryCore.ts index 01cc715e..392ea28e 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 { @@ -404,6 +405,16 @@ export async function getMemoryIdsByDerivedFrom(derivedFrom: string, limit = 100 return rows.map((row: any) => String(row.id)); } +/** Cross-process lock serializing memory table mutations (writes + compaction). */ +export function memoryMutationLockPath(): string { + return process.env.OPENSWARM_MEMORY_MUTATION_LOCK + ?? join(homedir(), '.openswarm', 'memory-mutation.lock'); +} + +export async function withMemoryMutationLock(operation: () => Promise): Promise { + return withFileLock(memoryMutationLockPath(), operation, { timeoutMs: 120_000 }); +} + /** * Retry a Lance write (add/update/delete) on optimistic-concurrency conflict. * @@ -416,12 +427,15 @@ 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 runs under {@link withMemoryMutationLock} so compaction cannot + * interleave with a write mid-retry. */ export async function withMemoryWriteRetry(op: () => Promise, label = 'write'): Promise { const MAX_ATTEMPTS = 8; for (let attempt = 1; ; attempt++) { try { - return await op(); + return await withMemoryMutationLock(() => 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..8728b70b 100644 --- a/src/memory/memoryWriteRetry.test.ts +++ b/src/memory/memoryWriteRetry.test.ts @@ -1,11 +1,22 @@ -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. describe('withMemoryWriteRetry (INT-2817 store-path concurrency)', () => { + afterAll(() => { + rmSync(lockDir, { recursive: true, force: true }); + }); + afterEach(() => { vi.useRealTimers(); vi.restoreAllMocks(); diff --git a/src/orchestration/workflow.ts b/src/orchestration/workflow.ts index 021bdaca..32e9a5ac 100644 --- a/src/orchestration/workflow.ts +++ b/src/orchestration/workflow.ts @@ -7,6 +7,8 @@ import { basename, isAbsolute, relative, resolve } from 'path'; import { homedir } from 'os'; import * as fs from 'fs/promises'; import * as yaml from 'yaml'; +import { atomicWriteFile } from '../support/atomicFile.js'; +import { withFileLock } from '../support/fileLock.js'; // Types & Interfaces @@ -277,7 +279,9 @@ function storageFilePath(rootDir: string, id: string, extension: string): string export async function saveWorkflow(workflow: WorkflowConfig): Promise { const filePath = storageFilePath(WORKFLOW_DIR, workflow.id, '.yaml'); await fs.mkdir(WORKFLOW_DIR, { recursive: true }); - await fs.writeFile(filePath, yaml.stringify(workflow), 'utf-8'); + await withFileLock(`${filePath}.lock`, async () => { + await atomicWriteFile(filePath, yaml.stringify(workflow)); + }, { timeoutMs: 10_000 }); console.log(`[Workflow] Saved: ${workflow.name} (${workflow.id})`); } @@ -328,7 +332,9 @@ export async function listWorkflows(): Promise { export async function saveExecution(execution: WorkflowExecution): Promise { const filePath = storageFilePath(EXECUTION_DIR, execution.executionId, '.json'); await fs.mkdir(EXECUTION_DIR, { recursive: true }); - await fs.writeFile(filePath, JSON.stringify(execution, null, 2), 'utf-8'); + await withFileLock(`${filePath}.lock`, async () => { + await atomicWriteFile(filePath, JSON.stringify(execution, null, 2)); + }, { timeoutMs: 10_000 }); } /** diff --git a/src/support/dev.ts b/src/support/dev.ts index 2c73a62e..185aa65a 100644 --- a/src/support/dev.ts +++ b/src/support/dev.ts @@ -190,11 +190,30 @@ export async function runDevTask( activeTasks.set(taskId, devTask); + // Idempotent finalization: close and error can both fire; cancelTask may + // remove the map entry before the child exits. Callback failures must not + // leave the task stuck in activeTasks or prevent cleanup. + let finalized = false; + const finalize = (resultText: string, code: number | null): void => { + 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 @@ -224,17 +243,19 @@ export async function runDevTask( // Generate report file const duration = Math.floor((Date.now() - devTask.startedAt) / 1000); - generateReport(devTask, code, duration); + try { + generateReport(devTask, code, duration); + } catch (reportError) { + console.error(`[Dev] generateReport failed for ${taskId}:`, reportError); + } - 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 }; @@ -260,8 +281,14 @@ export function cancelTask(taskId: string): boolean { const task = activeTasks.get(taskId); if (!task) return false; - task.process.kill('SIGTERM'); + // Free the repo slot immediately; finalize on close still runs once (idempotent) + // so onComplete is isolated from double close/error delivery. activeTasks.delete(taskId); + try { + task.process.kill('SIGTERM'); + } catch { + // Process may already be gone; close handler still finalizes. + } return true; } diff --git a/src/support/fileLock.ts b/src/support/fileLock.ts index f9e0aed6..0697cbc6 100644 --- a/src/support/fileLock.ts +++ b/src/support/fileLock.ts @@ -1,4 +1,14 @@ 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'; @@ -13,9 +23,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 +34,83 @@ 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; + } +} + +function sleepSync(ms: number): void { + const end = Date.now() + ms; + while (Date.now() < end) { + // Busy-wait is acceptable for short cross-process lock hand-offs (≤10ms). + } +} + +async function reclaimStaleLock( + path: string, + current: LockOwner | null, + malformedStaleMs: number, +): Promise { + let malformedAndStale = false; + if (current === null) { + try { + malformedAndStale = Date.now() - (await stat(path)).mtimeMs > 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; + } + } + if ((current !== null && !alive(current.pid)) || malformedAndStale) { + await unlink(path).catch((unlinkError) => { + const code = (unlinkError as NodeJS.ErrnoException).code; + // ENOENT: another reclaim won. ENOTEMPTY: directory-style locks with + // concurrent claim markers — retry the open loop after a brief wait. + if (code !== 'ENOENT' && code !== 'ENOTEMPTY') throw unlinkError; + }); + return true; + } + return false; +} + +function reclaimStaleLockSync( + path: string, + current: LockOwner | null, + malformedStaleMs: number, +): boolean { + let malformedAndStale = false; + if (current === null) { + try { + malformedAndStale = Date.now() - statSync(path).mtimeMs > malformedStaleMs; + } catch (statError) { + if ((statError as NodeJS.ErrnoException).code !== 'ENOENT') throw statError; + return true; + } + } + if ((current !== null && !alive(current.pid)) || malformedAndStale) { + try { + unlinkSync(path); + } catch (unlinkError) { + const code = (unlinkError as NodeJS.ErrnoException).code; + if (code !== 'ENOENT' && code !== 'ENOTEMPTY') throw unlinkError; + } + return true; + } + return false; +} + export async function withFileLock( path: string, operation: () => Promise, @@ -45,21 +132,8 @@ export async function withFileLock( } catch (error) { if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error; const current = await owner(path); - let malformedAndStale = false; - if (current === null) { - try { - malformedAndStale = Date.now() - (await stat(path)).mtimeMs > 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; - continue; - } - } - if ((current !== null && !alive(current.pid)) || malformedAndStale) { - await unlink(path).catch((unlinkError) => { - if ((unlinkError as NodeJS.ErrnoException).code !== 'ENOENT') throw unlinkError; - }); + if (await reclaimStaleLock(path, current, malformedStaleMs)) { + await new Promise((resolve) => setTimeout(resolve, 10)); continue; } if (Date.now() >= deadline) throw new Error(`Timed out waiting for file lock: ${path}`); @@ -70,6 +144,7 @@ export async function withFileLock( 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; @@ -77,3 +152,53 @@ export async function withFileLock( } } } + +/** + * Synchronous counterpart for callers that must stay sync (telemetry, oauth + * save, pipeline history). Same ownership rules as {@link withFileLock}. + */ +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; + const current = ownerSync(path); + if (reclaimStaleLockSync(path, current, malformedStaleMs)) { + sleepSync(10); + continue; + } + if (Date.now() >= deadline) throw new Error(`Timed out waiting for file lock: ${path}`); + sleepSync(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 ff57c19d..8f6a5d1f 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 tasks blocked until dependencies are done, then releases them', () => { upsertTaskState('ISSUE-1', { execution: { status: 'in_progress', retryCount: 0 }, diff --git a/src/taskState/store.ts b/src/taskState/store.ts index ee41a824..462e9e80 100644 --- a/src/taskState/store.ts +++ b/src/taskState/store.ts @@ -298,7 +298,16 @@ 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. + if (code !== 'ENOENT' && code !== 'ENOTEMPTY') throw unlinkError; + } + } continue; } } catch (statError) { diff --git a/src/telemetry/telemetry.ts b/src/telemetry/telemetry.ts index cdb41193..b10f58ce 100644 --- a/src/telemetry/telemetry.ts +++ b/src/telemetry/telemetry.ts @@ -19,6 +19,7 @@ import { join } from 'node:path'; import { readFileSync } from 'node:fs'; import { nanoid } from 'nanoid'; import { atomicWriteFileSync } from '../support/atomicFile.js'; +import { withFileLockSync } from '../support/fileLock.js'; const STATE_DIR = join(homedir(), '.config', 'openswarm'); const TELEMETRY_FILE = join(STATE_DIR, 'telemetry.json'); @@ -96,7 +97,9 @@ export function mergeState( function writeState(state: TelemetryState): void { try { - atomicWriteFileSync(TELEMETRY_FILE, JSON.stringify(mergeState(readState(), state), null, 2)); + withFileLockSync(`${TELEMETRY_FILE}.lock`, () => { + atomicWriteFileSync(TELEMETRY_FILE, JSON.stringify(mergeState(readState(), state), null, 2)); + }); } catch { // A read-only home or race is non-fatal: telemetry just stays best-effort. }