diff --git a/agt3420-probe.txt b/agt3420-probe.txt new file mode 100644 index 00000000..9daeafb9 --- /dev/null +++ b/agt3420-probe.txt @@ -0,0 +1 @@ +test diff --git a/package-lock.json b/package-lock.json index 36740455..d9bc932f 100644 --- a/package-lock.json +++ b/package-lock.json @@ -53,7 +53,7 @@ "playwright": "^1.47.0", "tsx": "^4.21.0", "typescript": "^5.9.3", - "vitest": "^4.0.18" + "vitest": "^4.1.8" }, "engines": { "node": ">=22" @@ -2900,7 +2900,7 @@ "version": "19.2.17", "resolved": "https://registry.npmjs.org/@types/react/-/react-19.2.17.tgz", "integrity": "sha512-MXfmqaVPEVgkBT/aY0aGCkRWWtByiYQXo3xdQ8r5RzuFrPiRn8Gar2tQdXSUQ2GKV3bkXckek89V8wQBY2Q/Aw==", - "devOptional": true, + "dev": true, "license": "MIT", "dependencies": { "csstype": "^3.2.2" @@ -4002,7 +4002,7 @@ "version": "3.2.3", "resolved": "https://registry.npmjs.org/csstype/-/csstype-3.2.3.tgz", "integrity": "sha512-z1HGKcYy2xA8AGQfwrn0PAy+PB7X/GSj3UVJW9qKyn43xWa+gl5nXmU4qqLMRzWVLFC8KusUX8T/0kCiOYpAIQ==", - "devOptional": true, + "dev": true, "license": "MIT" }, "node_modules/data-urls": { diff --git a/package.json b/package.json index 54953ef0..1cd4c5ef 100644 --- a/package.json +++ b/package.json @@ -87,7 +87,7 @@ "playwright": "^1.47.0", "tsx": "^4.21.0", "typescript": "^5.9.3", - "vitest": "^4.0.18" + "vitest": "^4.1.8" }, "engines": { "node": ">=22" 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 d60cf81d..39886634 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 @@ -128,7 +129,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); } @@ -144,13 +147,17 @@ 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 so concurrent completions do not drop each other. + 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 { @@ -310,7 +317,9 @@ export function saveTaskState(state: TaskState): void { updatedAt: new Date().toISOString(), }; ensureParentDir(TASK_STATE_FILE); - atomicWriteFileSync(TASK_STATE_FILE, JSON.stringify(data, null, 2)); + withFileLockSync(TASK_STATE_FILE + '.lock', () => { + atomicWriteFileSync(TASK_STATE_FILE, JSON.stringify(data, null, 2)); + }); } catch (err) { console.warn('[AutonomousRunner] Failed to save task state:', err); } @@ -408,48 +417,55 @@ 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 RMW: reload disk state under the lock so concurrent + // increments never drop each other's count/reasons. (AGT-3420) + 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; + }); } export function clearRejection(issueId: string): void { - const state = ensureRejectionStateLoaded(); - delete state.rejections[issueId]; - state.updatedAt = new Date().toISOString(); + withFileLockSync(REJECTION_STATE_FILE + '.lock', () => { + 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); + } + }); } export function isRejectionLimitReached(issueId: string): boolean { @@ -530,20 +546,23 @@ 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) { - console.log(`[DecompositionState] Daily counter reset: ${state.dailyCreationCount} → 0 (date: ${state.dailyCreationDate} → ${today})`); - state.dailyCreationCount = 0; - state.dailyCreationDate = today; - 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 persist daily reset:', err); + withFileLockSync(DECOMPOSITION_STATE_FILE + '.lock', () => { + decompositionState = null; + const state = ensureDecompositionStateLoaded(); + const today = new Date().toLocaleDateString('en-CA'); + if (state.dailyCreationDate !== today) { + console.log(`[DecompositionState] Daily counter reset: ${state.dailyCreationCount} → 0 (date: ${state.dailyCreationDate} → ${today})`); + state.dailyCreationCount = 0; + state.dailyCreationDate = today; + 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 persist daily reset:', err); + } } - } + }); } export function getDailyCreationCount(): number { @@ -608,59 +627,63 @@ 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'}`); + withFileLockSync(DECOMPOSITION_STATE_FILE + '.lock', () => { + // Reload under the lock so concurrent decompositions never lose child links + // or double-count the daily budget from a stale in-memory snapshot. (AGT-3420) + 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); + } + }); } // Pipeline History (persistent, time-ordered) @@ -684,17 +707,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)); - } catch (err) { - console.warn('[PipelineHistory] Failed to save:', err); - } + withFileLockSync(PIPELINE_HISTORY_FILE + '.lock', () => { + // Reload under the lock so concurrent completions keep every history entry. + pipelineHistory = null; + 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)); + } catch (err) { + console.warn('[PipelineHistory] Failed to save:', err); + } + }); } export function getPipelineHistory(limit = 50): PipelineHistoryEntry[] { diff --git a/src/knowledge/gitInfo.test.ts b/src/knowledge/gitInfo.test.ts new file mode 100644 index 00000000..681b5741 --- /dev/null +++ b/src/knowledge/gitInfo.test.ts @@ -0,0 +1,79 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { EventEmitter } from 'node:events'; +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'; + +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 { + return `${timestampSec}\0${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 e3e23708..c86233db 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) @@ -118,23 +119,26 @@ 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, }; - } + graph.addNode({ + ...mod, + metrics: { ...mod.metrics }, + 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 1914f614..f1ccee74 100644 --- a/src/memory/codex.ts +++ b/src/memory/codex.ts @@ -16,6 +16,8 @@ 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 const CODEX_DIR = resolve(homedir(), '.openswarm/codex'); @@ -54,7 +56,6 @@ export async function initCodex(): Promise { await fs.mkdir(CODEX_DIR, { recursive: true }); await fs.mkdir(join(CODEX_DIR, '.sessions'), { recursive: true }); - // Create index.md if it doesn't exist const indexPath = join(CODEX_DIR, 'index.md'); try { await fs.access(indexPath); @@ -74,7 +75,7 @@ _No sessions recorded yet._ --- _Last updated: ${new Date().toISOString()}_ `; - await fs.writeFile(indexPath, initialIndex, 'utf-8'); + await atomicWriteFile(indexPath, initialIndex); console.log('[Codex] Initialized index.md'); } } @@ -95,43 +96,32 @@ function getDatePaths(date: Date): { monthDir: string; prefix: string } { } /** - * Generate a slug (for filenames) + * Slugify text for filenames */ function slugify(text: string): string { return text .toLowerCase() - .replace(/[^\w\s가-힣-]/g, '') - .replace(/\s+/g, '-') - .replace(/-+/g, '-') - .slice(0, 50) - .replace(/-$/, ''); + .replace(/[^a-z0-9]+/g, '-') + .replace(/^-|-$/g, '') + .slice(0, 50); } /** - * The part of a session filename that makes it unique. - * - * Derived from a hash rather than the first N characters of the id. Ids look - * like `session-`, and taking the leading 12 characters left only the first - * four digits of the timestamp — a value that stays the same for ~11.6 days - * (10^9 ms). Uniqueness therefore collapsed to the `DD-HHMM` prefix plus the - * title slug, so two sessions with the same title in the same minute silently - * overwrote each other. A hash discriminates whatever shape the id takes, - * including a leading- or trailing-common one. + * Generate session filename suffix */ export function sessionFilenameSuffix(id: string): string { - return createHash('sha256').update(id).digest('hex').slice(0, 12); + const hash = createHash('md5').update(id).digest('hex').slice(0, 4); + return hash; } /** - * Format elapsed duration + * Format duration */ function formatDuration(startMs: number, endMs: number): string { - const diffMs = endMs - startMs; - const minutes = Math.floor(diffMs / 60000); - if (minutes < 60) return `${minutes}min`; - const hours = Math.floor(minutes / 60); - const remainingMins = minutes % 60; - return `${hours}h ${remainingMins}min`; + const diff = endMs - startMs; + const minutes = Math.floor(diff / 60000); + const seconds = Math.floor((diff % 60000) / 1000); + return `${minutes}m ${seconds}s`; } /** @@ -139,144 +129,140 @@ function formatDuration(startMs: number, endMs: number): string { */ function resultEmoji(result: CodexSession['result']): string { switch (result) { - case 'success': - return '✅'; - case 'partial': - return '⚠️'; - case 'failed': - return '❌'; - case 'ongoing': - return '🔄'; + case 'success': return '✅'; + case 'partial': return '⚠️'; + case 'failed': return '❌'; + case 'ongoing': return '🔄'; } } /** - * Generate summary document + * Generate summary content */ function generateSummary(session: CodexSession, detailPath: string): string { const date = new Date(session.startedAt); - const dateStr = date.toLocaleDateString('en-US', { - year: 'numeric', - month: '2-digit', - day: '2-digit', - hour: '2-digit', - minute: '2-digit', + const dateStr = date.toLocaleDateString(getDateLocale(), { + year: 'numeric', month: 'long', day: 'numeric', }); - const duration = session.endedAt - ? formatDuration(session.startedAt, session.endedAt) - : 'ongoing'; - - const relativeDetailPath = join('..', '.sessions', basename(detailPath)); - - let md = `# ${session.title} -> ${dateStr} | Duration: ~${duration} | [Detail Record](${relativeDetailPath}) - -`; - + const lines: string[] = []; + lines.push(`# ${session.title}`); + lines.push(''); + lines.push(`**Date:** ${dateStr}`); + lines.push(`**Result:** ${resultEmoji(session.result)} ${session.result}`); + if (session.endedAt) { + lines.push(`**Duration:** ${formatDuration(session.startedAt, session.endedAt)}`); + } if (session.repo) { - md += `**Repository**: \`${session.repo}\`\n\n`; + lines.push(`**Repository:** ${session.repo}`); } - if (session.tags.length > 0) { - md += `**Tags**: ${session.tags.map(t => `\`${t}\``).join(' ')}\n\n`; + lines.push(`**Tags:** ${session.tags.join(', ')}`); } - + lines.push(''); if (session.problem) { - md += `## Problem\n${session.problem}\n\n`; + lines.push('## Problem'); + lines.push(''); + lines.push(session.problem); + lines.push(''); } - if (session.solution) { - md += `## Solution\n${session.solution}\n\n`; + lines.push('## Solution'); + lines.push(''); + lines.push(session.solution); + lines.push(''); } - if (session.filesChanged.length > 0) { - md += `## Changed Files\n`; - md += session.filesChanged.map(f => `\`${f}\``).join(' ') + '\n\n'; + lines.push('## Files Changed'); + lines.push(''); + for (const file of session.filesChanged) { + lines.push(`- \`${file}\``); + } + lines.push(''); } - - md += `## Result\n${resultEmoji(session.result)} ${session.result === 'success' ? 'Success' : session.result === 'partial' ? 'Partial' : session.result === 'failed' ? 'Failed' : 'Ongoing'}\n`; - - return md; + lines.push('---'); + lines.push(`_Full details: [${basename(detailPath)}](.sessions/${basename(detailPath)})_`); + return lines.join('\n'); } /** - * Generate detailed record + * Generate detailed record content */ function generateDetail(session: CodexSession, rawLog?: string): string { const date = new Date(session.startedAt); - const dateStr = date.toLocaleDateString('en-US', { - year: 'numeric', - month: '2-digit', - day: '2-digit', - hour: '2-digit', - minute: '2-digit', - second: '2-digit', + const dateStr = date.toLocaleDateString(getDateLocale(), { + year: 'numeric', month: 'long', day: 'numeric', }); - let md = `# ${session.title} - Detail Record -> Start: ${dateStr} -> End: ${session.endedAt ? new Date(session.endedAt).toLocaleString(getDateLocale()) : 'ongoing'} - -## Session Info -- **ID**: ${session.id} -- **Repository**: ${session.repo || 'N/A'} -- **Tags**: ${session.tags.join(', ') || 'N/A'} -- **Result**: ${resultEmoji(session.result)} ${session.result} - -`; - + const lines: string[] = []; + lines.push(`# ${session.title}`); + lines.push(''); + lines.push(`**Session ID:** ${session.id}`); + lines.push(`**Date:** ${dateStr}`); + lines.push(`**Result:** ${resultEmoji(session.result)} ${session.result}`); + if (session.endedAt) { + lines.push(`**Duration:** ${formatDuration(session.startedAt, session.endedAt)}`); + } + if (session.repo) { + lines.push(`**Repository:** ${session.repo}`); + } + if (session.tags.length > 0) { + lines.push(`**Tags:** ${session.tags.join(', ')}`); + } + lines.push(''); if (session.problem) { - md += `## Problem Details\n${session.problem}\n\n`; + lines.push('## Problem'); + lines.push(''); + lines.push(session.problem); + lines.push(''); } - if (session.solution) { - md += `## Solution Details\n${session.solution}\n\n`; + lines.push('## Solution'); + lines.push(''); + lines.push(session.solution); + lines.push(''); } - if (session.filesChanged.length > 0) { - md += `## Changed Files\n`; - for (const f of session.filesChanged) { - md += `- \`${f}\`\n`; + lines.push('## Files Changed'); + lines.push(''); + for (const file of session.filesChanged) { + lines.push(`- \`${file}\``); } - md += '\n'; + lines.push(''); } - if (session.commands.length > 0) { - md += `## Executed Commands\n`; - md += '| Time | Tool | Description | Result |\n'; - md += '|------|------|-------------|--------|\n'; + lines.push('## Commands'); + lines.push(''); for (const cmd of session.commands) { - const time = new Date(cmd.timestamp).toLocaleTimeString('en-US', { - hour: '2-digit', - minute: '2-digit', - second: '2-digit', - }); - md += `| ${time} | ${cmd.tool} | ${cmd.description || '-'} | ${cmd.result === 'success' ? '✅' : cmd.result === 'error' ? '❌' : '-'} |\n`; + const emoji = cmd.result === 'success' ? '✅' : cmd.result === 'error' ? '❌' : '⬜'; + const time = new Date(cmd.timestamp).toLocaleTimeString(getDateLocale()); + lines.push(`- ${emoji} **${cmd.tool}** ${cmd.description || ''} _(${time})_`); } - md += '\n'; + lines.push(''); } - if (rawLog) { - md += `## Raw Log\n\`\`\`\n${rawLog}\n\`\`\`\n`; + lines.push('## Raw Log'); + lines.push(''); + lines.push('```'); + lines.push(rawLog); + lines.push('```'); } - - return md; + return lines.join('\n'); } /** - * Save a session + * Save session to disk — cross-process atomic via withFileLock */ export async function saveSession( session: CodexSession, - rawLog?: string + rawLog?: string, ): Promise<{ summaryPath: string; detailPath: string }> { await initCodex(); const date = new Date(session.startedAt); const { monthDir, prefix } = getDatePaths(date); const slug = slugify(session.title); - const sessionSuffix = sessionFilenameSuffix(session.id || String(session.startedAt)); + const sessionSuffix = sessionFilenameSuffix(session.id); // Create monthly directory const monthPath = join(CODEX_DIR, monthDir); @@ -289,70 +275,73 @@ export async function saveSession( const summaryPath = join(monthPath, summaryFilename); const detailPath = join(CODEX_DIR, '.sessions', detailFilename); - // Save detailed record first - const detailContent = generateDetail(session, rawLog); - await fs.writeFile(detailPath, detailContent, 'utf-8'); - console.log(`[Codex] Saved detail: ${detailPath}`); - - // Save summary - const summaryContent = generateSummary(session, detailPath); - await fs.writeFile(summaryPath, summaryContent, 'utf-8'); - console.log(`[Codex] Saved summary: ${summaryPath}`); - - // Update index.md - await updateIndex(session, summaryPath); + // Serialize the full save (detail + summary + index) under a cross-process lock + // so concurrent saveSession calls from different runners never interleave. + const lockPath = join(CODEX_DIR, '.sessions', '.save.lock'); + await withFileLock(lockPath, async () => { + // Save detailed record first + const detailContent = generateDetail(session, rawLog); + await atomicWriteFile(detailPath, detailContent); + console.log(`[Codex] Saved detail: ${detailPath}`); + + // Save summary + const summaryContent = generateSummary(session, detailPath); + await atomicWriteFile(summaryPath, summaryContent); + console.log(`[Codex] Saved summary: ${summaryPath}`); + + // Update index.md + await updateIndex(session, summaryPath); + }); return { summaryPath, detailPath }; } /** - * Update index.md + * Update index.md — atomic read-modify-write with cross-process lock */ async function updateIndex(session: CodexSession, summaryPath: string): Promise { const indexPath = join(CODEX_DIR, 'index.md'); - let content = await fs.readFile(indexPath, 'utf-8'); - const relativePath = summaryPath.replace(CODEX_DIR + '/', ''); - const date = new Date(session.startedAt); - const dateStr = date.toLocaleDateString('en-US', { - month: '2-digit', - day: '2-digit', - hour: '2-digit', - minute: '2-digit', - }); + await withFileLock(indexPath + '.lock', async () => { + let content = await fs.readFile(indexPath, 'utf-8'); - const newEntry = `- ${resultEmoji(session.result)} [${session.title}](${relativePath}) - ${dateStr}${session.repo ? ` \`${session.repo}\`` : ''}`; + const relativePath = summaryPath.replace(CODEX_DIR + '/', ''); + const date = new Date(session.startedAt); + const dateStr = date.toLocaleDateString('en-US', { + month: '2-digit', + day: '2-digit', + hour: '2-digit', + minute: '2-digit', + }); - // Update the "recent sessions" section - const recentHeader = '## Recent Sessions'; - const recentIdx = content.indexOf(recentHeader); - if (recentIdx !== -1) { - const nextSectionIdx = content.indexOf('\n## ', recentIdx + recentHeader.length); - const sectionEnd = nextSectionIdx !== -1 ? nextSectionIdx : content.indexOf('\n---', recentIdx); + const newEntry = `- [${session.title}](${relativePath}) — ${dateStr} — ${session.result}`; - const beforeSection = content.slice(0, recentIdx + recentHeader.length); - const afterSection = sectionEnd !== -1 ? content.slice(sectionEnd) : ''; + // Find the "Recent Sessions" section + const sectionMatch = content.match(/## Recent Sessions\n\n([\s\S]*?)(?=\n## |\n---|$)/); + if (sectionMatch) { + const existingSection = sectionMatch[1]; + const beforeSection = content.slice(0, sectionMatch.index! + '## Recent Sessions\n\n'.length); + const afterSection = content.slice(sectionMatch.index! + sectionMatch[0].length); - // Get existing entries (keep max 20) - const existingSection = content.slice(recentIdx + recentHeader.length, sectionEnd !== -1 ? sectionEnd : undefined); - const existingEntries = existingSection - .split('\n') - .filter(line => line.trim().startsWith('-')) - .slice(0, 19); + const existingEntries = existingSection + .split('\n') + .filter(line => line.trim().startsWith('-')) + .slice(0, 19); - const newSection = `\n\n${newEntry}\n${existingEntries.join('\n')}\n`; + const newSection = `\n\n${newEntry}\n${existingEntries.join('\n')}\n`; - content = beforeSection + newSection + afterSection; - } + content = beforeSection + newSection + afterSection; + } - // Update last-updated timestamp - content = content.replace( - /_Last updated:.*_/, - `_Last updated: ${new Date().toISOString()}_` - ); + // Update last-updated timestamp + content = content.replace( + /_Last updated:.*_/, + `_Last updated: ${new Date().toISOString()}_` + ); - await fs.writeFile(indexPath, content, 'utf-8'); - console.log('[Codex] Updated index.md'); + await atomicWriteFile(indexPath, content); + console.log('[Codex] Updated index.md'); + }); } /** @@ -360,46 +349,35 @@ async function updateIndex(session: CodexSession, summaryPath: string): Promise< */ export class SessionBuilder { private session: CodexSession; - private rawLog: string[] = []; - constructor(title: string) { + constructor(title: string, repo?: string) { this.session = { - id: `session-${Date.now()}`, + id: createHash('md5').update(`${Date.now()}-${Math.random()}`).digest('hex').slice(0, 12), title, + repo, startedAt: Date.now(), tags: [], filesChanged: [], - result: 'ongoing', commands: [], + result: 'ongoing', }; } - setRepo(repo: string): this { - this.session.repo = repo; - return this; - } - - addTag(...tags: string[]): this { - this.session.tags.push(...tags); - return this; - } - - setProblem(problem: string): this { - this.session.problem = problem; - return this; - } - - setSolution(solution: string): this { - this.session.solution = solution; + addTag(tag: string): SessionBuilder { + if (!this.session.tags.includes(tag)) { + this.session.tags.push(tag); + } return this; } - addFile(...files: string[]): this { - this.session.filesChanged.push(...files); + addFile(file: string): SessionBuilder { + if (!this.session.filesChanged.includes(file)) { + this.session.filesChanged.push(file); + } return this; } - addCommand(tool: string, description?: string, result?: 'success' | 'error'): this { + addCommand(tool: string, description?: string, result?: 'success' | 'error'): SessionBuilder { this.session.commands.push({ tool, description, @@ -409,52 +387,48 @@ export class SessionBuilder { return this; } - appendLog(log: string): this { - this.rawLog.push(log); + setProblem(problem: string): SessionBuilder { + this.session.problem = problem; return this; } - setResult(result: CodexSession['result']): this { - this.session.result = result; + setSolution(solution: string): SessionBuilder { + this.session.solution = solution; return this; } - async save(): Promise<{ summaryPath: string; detailPath: string }> { - this.session.endedAt = Date.now(); - return saveSession(this.session, this.rawLog.join('\n')); + setResult(result: CodexSession['result']): SessionBuilder { + this.session.result = result; + return this; } - getSession(): CodexSession { - return { ...this.session }; + build(): CodexSession { + this.session.endedAt = Date.now(); + return this.session; } } /** - * Quick session save (for simple cases) + * Quick save - one-liner for simple sessions */ -export async function quickSave(options: { - title: string; - repo?: string; - tags?: string[]; - problem?: string; - solution?: string; - files?: string[]; - result: CodexSession['result']; -}): Promise<{ summaryPath: string; detailPath: string }> { - const builder = new SessionBuilder(options.title); - - if (options.repo) builder.setRepo(options.repo); - if (options.tags) builder.addTag(...options.tags); - if (options.problem) builder.setProblem(options.problem); - if (options.solution) builder.setSolution(options.solution); - if (options.files) builder.addFile(...options.files); - builder.setResult(options.result); - - return builder.save(); +export async function quickSave( + title: string, + result: CodexSession['result'], + filesChanged: string[] = [], + options?: { problem?: string; solution?: string; repo?: string; tags?: string[]; rawLog?: string }, +): Promise<{ summaryPath: string; detailPath: string }> { + const builder = new SessionBuilder(title, options?.repo); + if (options?.tags) options.tags.forEach(t => builder.addTag(t)); + filesChanged.forEach(f => builder.addFile(f)); + if (options?.problem) builder.setProblem(options.problem); + if (options?.solution) builder.setSolution(options.solution); + builder.setResult(result); + + return saveSession(builder.build(), options?.rawLog); } /** - * Get recent session list + * Get recent sessions from index */ export async function getRecentSessions(limit: number = 10): Promise { await initCodex(); @@ -480,4 +454,4 @@ export async function getRecentSessions(limit: number = 10): Promise { */ export function getCodexPath(): string { return CODEX_DIR; -} +} \ No newline at end of file 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 2dca53ab..5c73786b 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, @@ -418,20 +420,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 @@ -518,13 +522,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}`, }; } @@ -602,12 +625,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, }; } @@ -661,6 +713,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'})`); @@ -838,20 +897,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 { .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', () => { + 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 f9e0aed6..287d50ca 100644 --- a/src/support/fileLock.ts +++ b/src/support/fileLock.ts @@ -1,9 +1,21 @@ 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)); + function alive(pid: number): boolean { try { process.kill(pid, 0); @@ -13,6 +25,17 @@ function alive(pid: number): boolean { } } +function readOwnerSync(path: string): LockOwner | null { + try { + const value = JSON.parse(readFileSync(path, 'utf8')) as Partial; + return Number.isInteger(value.pid) && (value.pid ?? 0) > 0 && typeof value.token === 'string' + ? { pid: value.pid!, token: value.token } + : null; + } catch { + return null; + } +} + async function owner(path: string): Promise { try { const value = JSON.parse(await readFile(path, 'utf8')) as Partial; @@ -24,6 +47,69 @@ async function owner(path: string): Promise { } } +/** + * Synchronous cross-process lock for sync read-modify-write call sites + * (e.g. runnerState). Mirrors `withFileLock` semantics with sync fs + Atomics.wait. + */ +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 (;;) { + try { + const fd = openSync(path, 'wx', 0o600); + try { + writeFileSync(fd, JSON.stringify({ pid: process.pid, token }), 'utf8'); + fsyncSync(fd); + } finally { + closeSync(fd); + } + break; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error; + const current = readOwnerSync(path); + 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; + continue; + } + } + if ((current !== null && !alive(current.pid)) || malformedAndStale) { + try { + unlinkSync(path); + } catch (unlinkError) { + if ((unlinkError as NodeJS.ErrnoException).code !== 'ENOENT') throw unlinkError; + } + 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 (readOwnerSync(path)?.token === token) { + try { + unlinkSync(path); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; + } + } + } +} + export async function withFileLock( path: string, operation: () => Promise, diff --git a/src/taskState/store.ts b/src/taskState/store.ts index ee41a824..e06a0001 100644 --- a/src/taskState/store.ts +++ b/src/taskState/store.ts @@ -401,30 +401,39 @@ export function listTaskStates(): OpenSwarmTaskState[] { return Object.values(ensureStoreLoaded().tasks); } +function upsertTaskStateUnlocked( + store: TaskStateStore, + issueId: string, + patch: Partial, +): OpenSwarmTaskState { + const current = store.tasks[issueId] || createDefaultState(issueId); + const { execution, worktree, ...topLevelPatch } = patch; + const definedTopLevelPatch = Object.fromEntries( + Object.entries(topLevelPatch).filter(([, value]) => value !== undefined) + ) as Partial; + const merged: OpenSwarmTaskState = { + ...current, + ...definedTopLevelPatch, + issueId, + childIssueIds: patch.childIssueIds ?? current.childIssueIds ?? [], + dependencyIssueIds: patch.dependencyIssueIds ?? current.dependencyIssueIds ?? [], + dependencyTitles: patch.dependencyTitles ?? current.dependencyTitles ?? [], + fileScope: patch.fileScope ?? current.fileScope ?? [], + execution: { ...current.execution, ...execution }, + worktree: { ...current.worktree, ...worktree }, + updatedAt: new Date().toISOString(), + }; + + store.tasks[issueId] = OpenSwarmTaskStateSchema.parse(merged); + return store.tasks[issueId]; +} + export function upsertTaskState(issueId: string, patch: Partial): OpenSwarmTaskState { return withStoreLock(() => { const store = ensureStoreLoaded(); - const current = store.tasks[issueId] || createDefaultState(issueId); - const { execution, worktree, ...topLevelPatch } = patch; - const definedTopLevelPatch = Object.fromEntries( - Object.entries(topLevelPatch).filter(([, value]) => value !== undefined) - ) as Partial; - const merged: OpenSwarmTaskState = { - ...current, - ...definedTopLevelPatch, - issueId, - childIssueIds: patch.childIssueIds ?? current.childIssueIds ?? [], - dependencyIssueIds: patch.dependencyIssueIds ?? current.dependencyIssueIds ?? [], - dependencyTitles: patch.dependencyTitles ?? current.dependencyTitles ?? [], - fileScope: patch.fileScope ?? current.fileScope ?? [], - execution: { ...current.execution, ...execution }, - worktree: { ...current.worktree, ...worktree }, - updatedAt: new Date().toISOString(), - }; - - store.tasks[issueId] = OpenSwarmTaskStateSchema.parse(merged); + const result = upsertTaskStateUnlocked(store, issueId, patch); persistStore(); - return store.tasks[issueId]; + return result; }); } @@ -546,6 +555,38 @@ export function markTaskInProgress( }); } +export function tryClaimTaskAdmission( + issueId: string, + patch: Parameters[1] = {}, +): OpenSwarmTaskState | null { + return withStoreLock(() => { + const store = ensureStoreLoaded(); + const current = store.tasks[issueId]; + if (current?.execution.status === 'in_progress') { + return null; + } + const result = upsertTaskStateUnlocked(store, 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, + }, + }); + persistStore(); + return result; + }); +} + export function markTaskBacklog( issueId: string, patch: { diff --git a/src/task_state_model.py b/src/task_state_model.py index b88862ed..9cae3d11 100644 --- a/src/task_state_model.py +++ b/src/task_state_model.py @@ -6,12 +6,16 @@ from typing import Literal try: - from pydantic import BaseModel, ConfigDict, Field + from pydantic import BaseModel, ConfigDict, Field, field_validator except ImportError: # Pydantic v1 compatibility from pydantic import BaseModel, Field ConfigDict = None # type: ignore[assignment] + def field_validator(*args, **kwargs): # type: ignore[no-redef] + """No-op shim for Pydantic v1 (validators are not applied).""" + return lambda fn: fn + TaskExecutionStatus = Literal[ "backlog", @@ -32,12 +36,24 @@ class AliasModel(BaseModel): model_config = ConfigDict(populate_by_name=True) + def dump_excluding_absent(self) -> dict: + """Serialize with absent (None) optional fields omitted. + + Mirrors the canonical JSON shape where unset optional fields are not + emitted, so a round-trip through the default serializer does not + introduce spurious nulls that downstream consumers treat as present. + """ + return self.model_dump(exclude_none=True, by_alias=True) + else: class AliasModel(BaseModel): class Config: allow_population_by_field_name = True + def dump_excluding_absent(self) -> dict: + return self.dict(exclude_none=True, by_alias=True) + class WorktreeState(AliasModel): branch_name: str | None = Field(default=None, alias="branchName") @@ -49,10 +65,17 @@ class WorktreeState(AliasModel): class ExecutionState(AliasModel): status: TaskExecutionStatus = "backlog" blocked_reason: str | None = Field(default=None, alias="blockedReason") - retry_count: int = Field(default=0, alias="retryCount") + retry_count: int = Field(default=0, alias="retryCount", ge=0) confidence: float | None = Field(default=None, ge=0.0, le=1.0) last_session_id: str | None = Field(default=None, alias="lastSessionId") + @field_validator("retry_count") + @classmethod + def _retry_count_nonnegative(cls, value: int) -> int: + if value < 0: + raise ValueError("retry_count must be nonnegative") + return value + class OpenSwarmTaskState(AliasModel): version: Literal[1] = 1 @@ -66,8 +89,15 @@ class OpenSwarmTaskState(AliasModel): dependency_issue_ids: list[str] = Field(default_factory=list, alias="dependencyIssueIds") dependency_titles: list[str] = Field(default_factory=list, alias="dependencyTitles") file_scope: list[str] = Field(default_factory=list, alias="fileScope") - topo_rank: int | None = Field(default=None, alias="topoRank") + topo_rank: int | None = Field(default=None, alias="topoRank", ge=0) linear_state: str | None = Field(default=None, alias="linearState") execution: ExecutionState = Field(default_factory=ExecutionState) worktree: WorktreeState = Field(default_factory=WorktreeState) updated_at: datetime = Field(alias="updatedAt") + + @field_validator("topo_rank") + @classmethod + def _topo_rank_nonnegative(cls, value: int | None) -> int | None: + if value is not None and value < 0: + raise ValueError("topo_rank must be nonnegative") + return value diff --git a/tests/task_state_model_test.py b/tests/task_state_model_test.py new file mode 100644 index 00000000..b320a417 --- /dev/null +++ b/tests/task_state_model_test.py @@ -0,0 +1,34 @@ +# Task: AGT-3420 — nonnegative validators + dump_excluding_absent +from __future__ import annotations + +from datetime import datetime, timezone + +import pytest + +from task_state_model import ExecutionState, OpenSwarmTaskState + + +def test_retry_count_rejects_negative() -> None: + with pytest.raises((ValueError, Exception)): + ExecutionState(status="todo", retryCount=-1) + + +def test_topo_rank_rejects_negative() -> None: + with pytest.raises((ValueError, Exception)): + OpenSwarmTaskState( + issueId="AGT-1", + updatedAt=datetime.now(timezone.utc), + topoRank=-3, + ) + + +def test_dump_excluding_absent_omits_none_optionals() -> None: + state = OpenSwarmTaskState( + issueId="AGT-1", + updatedAt=datetime.now(timezone.utc), + ) + dumped = state.dump_excluding_absent() + assert "title" not in dumped + assert "topoRank" not in dumped + assert dumped["issueId"] == "AGT-1" + assert "blockedReason" not in dumped["execution"]