diff --git a/.restore-from-main.mjs b/.restore-from-main.mjs new file mode 100644 index 00000000..2467d7b5 --- /dev/null +++ b/.restore-from-main.mjs @@ -0,0 +1,6 @@ +import { copyFileSync, readFileSync, writeFileSync } from 'node:fs'; +const WT = '/work/OpenSwarm/worktree/743ccc9e-7b29-4268-ad09-c64586dad683'; +const MAIN = '/work/OpenSwarm/src'; +const rel = 'automation/prProcessor.ts'; +writeFileSync(`${WT}/src/${rel}`, readFileSync(`${MAIN}/${rel}`)); +console.log('copied', rel); diff --git a/_patch_pr_processor.py b/_patch_pr_processor.py new file mode 100644 index 00000000..f0614aa6 --- /dev/null +++ b/_patch_pr_processor.py @@ -0,0 +1,177 @@ +#!/usr/bin/env python3 +"""Copy and patch prProcessor.ts: withStoreLock -> withFileLock lease.""" +from __future__ import annotations + +from pathlib import Path + +SRC = Path("/work/OpenSwarm/src/automation/prProcessor.ts") +DST = Path( + "/work/OpenSwarm/worktree/743ccc9e-7b29-4268-ad09-c64586dad683" + "/src/automation/prProcessor.ts" +) + +text = SRC.read_text(encoding="utf-8") + +# 1. Add withFileLock import after atomicWriteFileSync import +old_import = "import { atomicWriteFileSync } from '../support/atomicFile.js';\n" +new_import = ( + "import { atomicWriteFileSync } from '../support/atomicFile.js';\n" + "import { withFileLock } from '../support/fileLock.js';\n" +) +if "withFileLock" not in text.split("from '../support/fileLock.js'")[0] if "fileLock" in text else True: + if "import { withFileLock } from '../support/fileLock.js';" not in text: + if old_import not in text: + raise SystemExit("atomicWriteFileSync import not found") + text = text.replace(old_import, new_import, 1) + +# 2. Add PR_STATE_LOCK_PATH after PR_STATE_PATH +old_path = ( + "const PR_STATE_PATH = resolve(homedir(), '.openswarm', 'pr-state.json');\n" +) +new_path = ( + "const PR_STATE_PATH = resolve(homedir(), '.openswarm', 'pr-state.json');\n" + "const PR_STATE_LOCK_PATH = `${PR_STATE_PATH}.lock`;\n" +) +if "PR_STATE_LOCK_PATH" not in text: + if old_path not in text: + raise SystemExit("PR_STATE_PATH not found") + text = text.replace(old_path, new_path, 1) + +# 3. Wrap fixOne processPR call +old_fix = """ await this.processPR(pr, projectPath, state, key); + const entry = state.prs[key]; + return { + success: entry?.status === 'completed', + error: entry?.lastError, + iterations: entry?.iterations ?? 0, + }; + } + + /** + * One-shot review-feedback pass for a single PR (CLI `openswarm pr review`). +""" + +new_fix = """ // Cross-process lease shared with the cron path (state + checkout serialization) + await withFileLock(PR_STATE_LOCK_PATH, async () => { + await this.processPR(pr, projectPath, state, key); + }, { timeoutMs: 30 * 60_000 }); + const entry = state.prs[key]; + return { + success: entry?.status === 'completed', + error: entry?.lastError, + iterations: entry?.iterations ?? 0, + }; + } + + /** + * One-shot review-feedback pass for a single PR (CLI `openswarm pr review`). +""" + +if old_fix not in text: + raise SystemExit("fixOne processPR block not found") +text = text.replace(old_fix, new_fix, 1) + +# 4. Wrap reviewOne load/process/save +old_review = """ async reviewOne( + pr: PRInfo, + 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, + }; + await this.processReviewFeedback(pr, projectPath, state, key, 0); + await this.saveState(state); + const entry = state.prs[key]; + return { + success: entry?.status === 'completed', + error: entry?.lastError, + iterations: entry?.iterations ?? 0, + }; + } +""" + +new_review = """ async reviewOne( + pr: PRInfo, + projectPath: string, + ): Promise<{ success: boolean; error?: string; iterations: number }> { + const key = `${pr.repo}#${pr.number}`; + return await withFileLock(PR_STATE_LOCK_PATH, async () => { + const state = await this.loadState(); + state.prs[key] = { + ...state.prs[key], + repo: pr.repo, + prNumber: pr.number, + status: 'processing', + iterations: 0, + }; + await this.processReviewFeedback(pr, projectPath, state, key, 0); + await this.saveState(state); + const entry = state.prs[key]; + return { + success: entry?.status === 'completed', + error: entry?.lastError, + iterations: entry?.iterations ?? 0, + }; + }, { timeoutMs: 30 * 60_000 }); + } +""" + +if old_review not in text: + raise SystemExit("reviewOne block not found") +text = text.replace(old_review, new_review, 1) + +# 5. Wrap processPRs try body in withFileLock +old_try = """ try { + const state = await this.loadState(); + + for (const repo of this.config.repos) { +""" + +new_try = """ try { + await withFileLock(PR_STATE_LOCK_PATH, async () => { + const state = await this.loadState(); + + for (const repo of this.config.repos) { +""" + +if old_try not in text: + raise SystemExit("processPRs try start not found") +text = text.replace(old_try, new_try, 1) + +# Close the withFileLock before catch — indent the saveState and close callback +old_end = """ await this.saveState(state); + } catch (err) { + console.error('[PRProcessor] Error:', err); + } finally { +""" + +new_end = """ await this.saveState(state); + }, { timeoutMs: 30 * 60_000 }); + } catch (err) { + console.error('[PRProcessor] Error:', err); + } finally { +""" + +if old_end not in text: + raise SystemExit("processPRs try end not found") +text = text.replace(old_end, new_end, 1) + +if "withStoreLock" in text: + raise SystemExit("withStoreLock still present after patch") + +DST.write_text(text, encoding="utf-8") +lines = text.count("\n") + (0 if text.endswith("\n") else 1) +print(f"Wrote {DST}") +print(f"lines: {lines}") +print("head:") +print("\n".join(text.splitlines()[:3])) +print("--- verify ---") +for i, line in enumerate(text.splitlines(), 1): + if any(s in line for s in ("withFileLock", "withStoreLock", "PR_STATE_LOCK")): + print(f"{i}:{line}") diff --git a/_tester_probe.txt b/_tester_probe.txt new file mode 100644 index 00000000..9daeafb9 --- /dev/null +++ b/_tester_probe.txt @@ -0,0 +1 @@ +test diff --git a/package-lock.json b/package-lock.json index ba2b8fa8..5713a7b7 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1300,9 +1300,6 @@ "cpu": [ "arm" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1319,9 +1316,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1338,9 +1332,6 @@ "cpu": [ "ppc64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1357,9 +1348,6 @@ "cpu": [ "riscv64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1376,9 +1364,6 @@ "cpu": [ "s390x" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1395,9 +1380,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1414,9 +1396,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1433,9 +1412,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1452,9 +1428,6 @@ "cpu": [ "arm" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1477,9 +1450,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1502,9 +1472,6 @@ "cpu": [ "ppc64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1527,9 +1494,6 @@ "cpu": [ "riscv64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1552,9 +1516,6 @@ "cpu": [ "s390x" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1577,9 +1538,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1602,9 +1560,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1627,9 +1582,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "Apache-2.0", "optional": true, "os": [ diff --git a/src/automation/prProcessor.ts b/src/automation/prProcessor.ts index cf12348d..38424c3a 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_PATH = `${PR_STATE_PATH}.lock`; // PR Processor @@ -338,7 +340,10 @@ export class PRProcessor { integrationBaselines: {}, updatedAt: new Date().toISOString(), }; - await this.processPR(pr, projectPath, state, key); + // Cross-process lease shared with cron processPRs (serialize state + checkout) + await withFileLock(PR_STATE_LOCK_PATH, async () => { + await this.processPR(pr, projectPath, state, key); + }, { timeoutMs: 30 * 60_000 }); const entry = state.prs[key]; return { success: entry?.status === 'completed', @@ -367,22 +372,24 @@ 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, - }; - await this.processReviewFeedback(pr, projectPath, state, key, 0); - await this.saveState(state); - const entry = state.prs[key]; - return { - success: entry?.status === 'completed', - error: entry?.lastError, - iterations: entry?.iterations ?? 0, - }; + return await withFileLock(PR_STATE_LOCK_PATH, async () => { + const state = await this.loadState(); + state.prs[key] = { + ...state.prs[key], + repo: pr.repo, + prNumber: pr.number, + status: 'processing', + iterations: 0, + }; + await this.processReviewFeedback(pr, projectPath, state, key, 0); + await this.saveState(state); + const entry = state.prs[key]; + return { + success: entry?.status === 'completed', + error: entry?.lastError, + iterations: entry?.iterations ?? 0, + }; + }, { timeoutMs: 30 * 60_000 }); } /** @@ -489,6 +496,10 @@ export class PRProcessor { const review = await runReviewCommand({ path: scratchWorktree, base: mergeBase, + // The scratch checkout is the reviewed repository, not OpenSwarm, so + // config discovery there falls back to the unavailable `codex` CLI. + // Preserve the daemon's explicitly configured PR reviewer adapter. + adapter: this.config.roles?.reviewer?.adapter, // The checked-out content is another PR's diff — untrusted the same // way review-gate.yml's CI run is (INT-3189). Denying mutating tools, // including bash, keeps a malicious PR from using the reviewer's @@ -630,6 +641,7 @@ export class PRProcessor { broadcastEvent({ type: 'pr_processor_start', data: { repos: this.config.repos } }); try { + await withFileLock(PR_STATE_LOCK_PATH, async () => { const state = await this.loadState(); for (const repo of this.config.repos) { @@ -754,6 +766,7 @@ export class PRProcessor { } await this.saveState(state); + }, { timeoutMs: 30 * 60_000 }); } catch (err) { console.error('[PRProcessor] Error:', err); } finally { diff --git a/src/automation/runnerState.ts b/src/automation/runnerState.ts index efa113d6..b6b25366 100644 --- a/src/automation/runnerState.ts +++ b/src/automation/runnerState.ts @@ -219,6 +219,14 @@ export function pickFailureDetail(candidates: Array): string /** Prefer the stage that actually failed over earlier successful feedback. */ export function pickPipelineFailureDetail(result: PipelineResult): string | undefined { + const workerFailure = result.workerResult?.success === false + ? pickFailureDetail([ + result.workerResult.error, + result.workerResult.haltReason, + result.workerResult.noChangesReason, + result.workerResult.summary, + ]) + : undefined; const testerFailure = result.testerResult?.success === false ? pickFailureDetail([ result.testerResult.error, @@ -227,11 +235,25 @@ export function pickPipelineFailureDetail(result: PipelineResult): string | unde ]) : undefined; + // Guards, security audit, verification, worktree setup and publication + // report through `stages[]` rather than a typed sub-result. Without this + // fallback the ledger recorded 57% of one day's failures with no message + // at all (vela, 2026-09-01), and the reason was unrecoverable once the + // container's log was gone. + const failedStage = [...result.stages].reverse().find((stage) => !stage.success); + const stageError = failedStage && 'error' in failedStage.result && typeof failedStage.result.error === 'string' + ? `${failedStage.stage}: ${failedStage.result.error}` + : undefined; + return pickFailureDetail([ + // Publication failed after every stage passed: nothing below describes it. + result.failureDetail, testerFailure, result.lastReviewFeedback, result.reviewResult?.feedback, - result.workerResult?.error, + workerFailure, + stageError, + result.stuckReason, ]); } @@ -599,6 +621,8 @@ export function registerDecomposition( // Validate the full batch before mutating the in-memory projection. A child // identity collision must leave no half-created parent entry behind. + // Daily creation capacity is reserved atomically via reserveDailyCreations() + // in the caller (runnerExecution) before child issues are created. for (const childId of uniqueChildren) { const existing = state.decompositions[childId]; if (existing && existing.parentId !== issueId) { diff --git a/src/issues/linearBridge.ts b/src/issues/linearBridge.ts index 62924e9f..9cf91655 100644 --- a/src/issues/linearBridge.ts +++ b/src/issues/linearBridge.ts @@ -12,6 +12,7 @@ import type { Issue, IssueStatus, IssuePriority } from './schema.js'; let linearClient: any = null; let linearTeamId: string = ''; let linearInitPromise: Promise | null = null; +let linearInitAttempts = 0; /** * Linear 브릿지 초기화 @@ -23,10 +24,14 @@ export function initLinearBridge(apiKey: string, teamId: string): Promise linearClient = null; linearInitPromise = import('@linear/sdk').then(({ LinearClient }) => { linearClient = new LinearClient({ apiKey }); + linearInitAttempts = 0; // Reset on success console.log('[LinearBridge] 초기화 완료 — team:', teamId); }).catch((err) => { linearClient = null; - console.warn('[LinearBridge] Linear SDK 로드 실패:', err); + linearInitAttempts++; + console.warn(`[LinearBridge] Linear SDK 로드 실패 (시도 ${linearInitAttempts}):`, err); + // Allow retry on next call by clearing the rejected promise + linearInitPromise = null; }); return linearInitPromise; } diff --git a/src/issues/memoryBridge.ts b/src/issues/memoryBridge.ts index f142f47b..b8a2c0b9 100644 --- a/src/issues/memoryBridge.ts +++ b/src/issues/memoryBridge.ts @@ -129,12 +129,9 @@ export async function saveBlockingConstraint( derivedFrom: `issue:${issue.id}`, }); + // linkMemory already emits a single memory_linked event — do not add a second one. if (memoryId) { store.linkMemory(issue.id, memoryId); - store.addEvent(issue.id, 'memory_linked', { - memoryId, - content: `블로킹 제약 조건 기억 저장: ${reason}`, - }); } return memoryId; diff --git a/src/orchestration/workflow.test.ts b/src/orchestration/workflow.test.ts index 94e85e9b..df818e46 100644 --- a/src/orchestration/workflow.test.ts +++ b/src/orchestration/workflow.test.ts @@ -4,6 +4,7 @@ import { loadWorkflow, saveExecution, saveWorkflow, + validateExecution, type WorkflowConfig, type WorkflowExecution, } from './workflow.js'; @@ -36,3 +37,45 @@ describe('workflow storage IDs', () => { await expect(loadExecution('../outside')).resolves.toBeNull(); }); }); + +describe('validateExecution', () => { + it('rejects completed steps missing completedAt', () => { + expect(() => validateExecution({ + workflowId: 'wf', + executionId: 'ex', + status: 'running', + startedAt: 0, + stepResults: { + step: { stepId: 'step', status: 'completed', startedAt: 0 }, + }, + })).toThrow(/completed but has no completedAt/); + }); + + it('rejects failed steps missing error', () => { + expect(() => validateExecution({ + workflowId: 'wf', + executionId: 'ex', + status: 'failed', + startedAt: 0, + stepResults: { + step: { stepId: 'step', status: 'failed', startedAt: 0, completedAt: 1 }, + }, + })).toThrow(/failed but has no error/); + }); + + it('rejects DAG-illegal advancement past pending dependencies', () => { + expect(() => validateExecution({ + workflowId: 'wf', + executionId: 'ex', + status: 'running', + startedAt: 0, + stepResults: { + a: { stepId: 'a', status: 'pending', startedAt: 0 }, + b: { stepId: 'b', status: 'completed', startedAt: 0, completedAt: 1 }, + }, + }, [ + { id: 'a', name: 'A', prompt: 'a' }, + { id: 'b', name: 'B', prompt: 'b', dependsOn: ['a'] }, + ])).toThrow(/cannot be completed when dependency a is pending/); + }); +}); diff --git a/src/orchestration/workflow.ts b/src/orchestration/workflow.ts index 021bdaca..4d666cb7 100644 --- a/src/orchestration/workflow.ts +++ b/src/orchestration/workflow.ts @@ -271,10 +271,64 @@ function storageFilePath(rootDir: string, id: string, extension: string): string return filePath; } +/** + * Reject incomplete step results and DAG-illegal lifecycle states before persist. + */ +export function validateExecution( + execution: WorkflowExecution, + workflowSteps?: WorkflowStep[], +): void { + const results = execution.stepResults; + + for (const [id, result] of Object.entries(results)) { + if (result.stepId !== id) { + throw new Error(`Step result key "${id}" does not match stepId "${result.stepId}"`); + } + if (result.status === 'completed' && result.completedAt == null) { + throw new Error(`Step ${id} is completed but has no completedAt`); + } + if (result.status === 'failed' && (result.error == null || result.error === '')) { + throw new Error(`Step ${id} is failed but has no error`); + } + if ( + (result.status === 'failed' || result.status === 'skipped') && + result.completedAt == null + ) { + throw new Error(`Step ${id} is ${result.status} but has no completedAt`); + } + } + + if (!workflowSteps || workflowSteps.length === 0) return; + + const stepById = new Map(workflowSteps.map((step) => [step.id, step])); + for (const [id, result] of Object.entries(results)) { + if (result.status === 'pending') continue; + const step = stepById.get(id); + if (!step?.dependsOn) continue; + for (const dep of step.dependsOn) { + const depResult = results[dep]; + if (!depResult || depResult.status === 'pending' || depResult.status === 'running') { + throw new Error( + `Step ${id} cannot be ${result.status} when dependency ${dep} is ${depResult?.status ?? 'missing'}`, + ); + } + if (depResult.status === 'failed' && result.status !== 'failed' && result.status !== 'skipped') { + throw new Error( + `Step ${id} cannot be ${result.status} when dependency ${dep} is failed`, + ); + } + } + } +} + /** * Save workflow */ export async function saveWorkflow(workflow: WorkflowConfig): Promise { + const validation = validateWorkflow(workflow); + if (!validation.valid) { + throw new Error(`Invalid workflow: ${validation.errors.join(', ')}`); + } const filePath = storageFilePath(WORKFLOW_DIR, workflow.id, '.yaml'); await fs.mkdir(WORKFLOW_DIR, { recursive: true }); await fs.writeFile(filePath, yaml.stringify(workflow), 'utf-8'); @@ -326,6 +380,8 @@ export async function listWorkflows(): Promise { * Save execution state */ export async function saveExecution(execution: WorkflowExecution): Promise { + const workflow = await loadWorkflow(execution.workflowId); + validateExecution(execution, workflow?.steps); 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'); diff --git a/src/support/dev.ts b/src/support/dev.ts index 2c73a62e..27eb2424 100644 --- a/src/support/dev.ts +++ b/src/support/dev.ts @@ -205,29 +205,34 @@ export async function runDevTask( // Handle completion claudeProcess.on('close', (code) => { - // Extract cost from stream-json output - const costInfo = extractCostFromStreamJson(devTask.output); - if (costInfo) { - console.log(`[Dev] ${repo} cost: ${formatCost(costInfo)}`); - } - - // Extract result text from stream-json for output let resultText = devTask.output; try { - const lines = devTask.output.split('\n').filter(Boolean); - const resultLine = lines.find((l) => l.includes('"type":"result"')); - if (resultLine) { - const parsed = JSON.parse(resultLine); - if (parsed.result) resultText = parsed.result; + // Extract cost from stream-json output + const costInfo = extractCostFromStreamJson(devTask.output); + if (costInfo) { + console.log(`[Dev] ${repo} cost: ${formatCost(costInfo)}`); } - } catch { /* use original */ } - // Generate report file - const duration = Math.floor((Date.now() - devTask.startedAt) / 1000); - generateReport(devTask, code, duration); + // Extract result text from stream-json for output + try { + const lines = devTask.output.split('\n').filter(Boolean); + const resultLine = lines.find((l) => l.includes('"type":"result"')); + if (resultLine) { + const parsed = JSON.parse(resultLine); + if (parsed.result) resultText = parsed.result; + } + } catch { /* use original */ } - onComplete?.(resultText, code); - activeTasks.delete(taskId); + // Generate report file + const duration = Math.floor((Date.now() - devTask.startedAt) / 1000); + generateReport(devTask, code, duration); + } catch (err) { + console.error(`[Dev] close-handler reporting failed for ${taskId}:`, err); + } finally { + // Cleanup must run even when reporting throws + onComplete?.(resultText, code); + activeTasks.delete(taskId); + } }); // Handle errors diff --git a/src/taskState/store.test.ts b/src/taskState/store.test.ts index 01ed9377..54116718 100644 --- a/src/taskState/store.test.ts +++ b/src/taskState/store.test.ts @@ -18,6 +18,7 @@ import { hydrateTaskStateFromComments, markTaskBacklog, planLinearStateReconciliation, + reconcileDependencyBlockers, resetTaskStateStoreForTests, buildLockPayload, type OpenSwarmTaskState, @@ -183,6 +184,22 @@ describe('task state store', () => { expect(getTaskState('PROCESS-B')?.title).toBe('PROCESS-B'); }); + it('reloads after a same-size state-file replacement within one mtime tick', async () => { + // Titles must be equal length so the on-disk JSON stays the same byte size. + const fixture = fileURLToPath(new URL('./storeClaimProcess.fixture.ts', import.meta.url)); + await new Promise((resolve, reject) => { + const child = spawn( + process.execPath, + ['--import', 'tsx', fixture, stateFile, 'SWAP-SRC', '0', '--same-size-replace'], + { stdio: 'pipe' }, + ); + let stderr = ''; + child.stderr.on('data', (chunk) => { stderr += String(chunk); }); + child.on('error', reject); + child.on('exit', (code) => code === 0 ? resolve() : reject(new Error(stderr || `child exited ${code}`))); + }); + }); + it('keeps tasks blocked until dependencies are done, then releases them', () => { upsertTaskState('ISSUE-1', { execution: { status: 'in_progress', retryCount: 0 }, @@ -247,6 +264,387 @@ describe('task state store', () => { expect(ready.blockedBy).toEqual([]); }); + it('reconcileDependencyBlockers releases a task whose blocker finished outside this daemon', async () => { + // AGT-4241-class bug (measured live on vela 2026-09-09): a blocker completed + // via a path that never called releaseDependentTasks (PR merged, Linear moved + // to Done by something other than this daemon's own pipeline). Its local + // taskState is stuck at a stale non-terminal status. The blocker also never + // appears in a future slim fetch again (Done issues are excluded from it), so + // task.blockedBy comes back empty and getTaskReadiness falls back to the + // stale dependencyIssueIds forever — this is the only path that can unstick it. + // upsertTaskState always stamps updatedAt to the real current time, so + // staleness is simulated by advancing the reconciler's `now` instead of + // trying to backdate the stored timestamp. + upsertTaskState('AGT-BLOCKER', { + execution: { status: 'in_progress', retryCount: 0 }, + linearState: 'In Review', + }); + // Live vela data (AGT-4209/AGT-4208, 2026-09-09): the stuck dependents' own + // execution.status was 'todo', not 'blocked' — getTaskReadiness gates on + // dependencyIssueIds at read time regardless of the stored status, and only + // releaseDependentTasks itself ever writes 'blocked'. The eligibility filter + // must therefore cover 'todo' too, not just 'blocked'. + upsertTaskState('AGT-DEPENDENT', { + dependencyIssueIds: ['AGT-BLOCKER'], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + const future = Date.now() + 2 * 60 * 60_000; + + const lookupIssueState = async (id: string) => id === 'AGT-BLOCKER' + ? { ok: true as const, issue: { state: 'Done', stateType: 'completed' } } + : { ok: false as const, error: 'unexpected id' }; + + const result = await reconcileDependencyBlockers({ source: { lookupIssueState }, now: future }); + + expect(result.eligible).toBe(1); + expect(result.lookedUp).toBe(1); + expect(result.resolved).toBe(1); + expect(result.released).toBe(1); + expect(getTaskState('AGT-BLOCKER')?.linearState).toBe('Done'); + expect(getTaskState('AGT-DEPENDENT')?.execution.status).toBe('todo'); + + const dependentTask = { + id: 'AGT-DEPENDENT', source: 'linear' as const, title: 'dependent', priority: 2, + createdAt: Date.now(), issueId: 'AGT-DEPENDENT', + }; + expect(getTaskReadiness(dependentTask).ready).toBe(true); + }); + + it('reconcileDependencyBlockers does not touch a task that already moved past todo/ready/blocked', async () => { + // Regression for a real finding from independent review: dependencyIssueIds + // is never cleared once a task moves on. A task that's already 'done' (or + // in_progress/decomposed/...) can still carry a stale, unresolved dependency + // entry from long before it finished. Without the execution.status filter, + // resolving that leftover dependency would call releaseDependentTasks and + // incorrectly reset the already-finished task back to 'todo'. + upsertTaskState('AGT-OLD-BLOCKER', { + execution: { status: 'in_progress', retryCount: 0 }, + linearState: 'In Review', + }); + upsertTaskState('AGT-ALREADY-DONE', { + dependencyIssueIds: ['AGT-OLD-BLOCKER'], + execution: { status: 'done', retryCount: 0 }, + linearState: 'Done', + }); + const future = Date.now() + 2 * 60 * 60_000; + + const lookupIssueState = async () => ({ ok: true as const, issue: { state: 'Done', stateType: 'completed' } }); + + const result = await reconcileDependencyBlockers({ source: { lookupIssueState }, now: future }); + + expect(result.eligible).toBe(0); + expect(result.lookedUp).toBe(0); + expect(result.released).toBe(0); + expect(getTaskState('AGT-ALREADY-DONE')?.execution.status).toBe('done'); + expect(getTaskState('AGT-ALREADY-DONE')?.linearState).toBe('Done'); + }); + + it('reconcileDependencyBlockers skips a dependency still present in the fresh fetch', async () => { + upsertTaskState('AGT-STILL-OPEN', { + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + upsertTaskState('AGT-DEPENDENT-2', { + dependencyIssueIds: ['AGT-STILL-OPEN'], + execution: { status: 'blocked', retryCount: 0 }, + linearState: 'Todo', + }); + const future = Date.now() + 2 * 60 * 60_000; + + const lookupIssueState = async () => { + throw new Error('should not be called — dependency is in knownTaskIds'); + }; + + const result = await reconcileDependencyBlockers({ + source: { lookupIssueState }, + knownTaskIds: new Set(['AGT-STILL-OPEN']), + now: future, + }); + + expect(result.eligible).toBe(0); + expect(result.lookedUp).toBe(0); + expect(result.released).toBe(0); + }); + + it('reconcileDependencyBlockers leaves a task blocked when the blocker is not yet stale', async () => { + upsertTaskState('AGT-FRESH-BLOCKER', { + execution: { status: 'in_progress', retryCount: 0 }, + linearState: 'In Progress', + updatedAt: new Date().toISOString(), + }); + upsertTaskState('AGT-FRESH-DEPENDENT', { + dependencyIssueIds: ['AGT-FRESH-BLOCKER'], + execution: { status: 'blocked', retryCount: 0 }, + linearState: 'Todo', + updatedAt: new Date().toISOString(), + }); + + const lookupIssueState = async () => { + throw new Error('should not be called — dependent was updated too recently to be stale'); + }; + + const result = await reconcileDependencyBlockers({ source: { lookupIssueState } }); + + expect(result.eligible).toBe(0); + expect(result.lookedUp).toBe(0); + }); + + it('reconcileDependencyBlockers still looks up blockers of a current-fetch dependent that was just upserted', async () => { + upsertTaskState('AGT-4114', { + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Backlog', + }); + upsertTaskState('AGT-4121', { + dependencyIssueIds: ['AGT-4114'], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + + const lookupIssueState = async (id: string) => id === 'AGT-4114' + ? { ok: true as const, issue: { state: 'Done', stateType: 'completed' } } + : { ok: false as const, error: 'unexpected id' }; + + const result = await reconcileDependencyBlockers({ + source: { lookupIssueState }, + knownTaskIds: new Set(['AGT-4121']), + }); + + expect(result.lookedUp).toBe(1); + expect(result.resolved).toBe(1); + expect(getTaskReadiness({ + id: 'AGT-4121', source: 'linear' as const, title: 'waiting', priority: 2, + createdAt: Date.now(), issueId: 'AGT-4121', + }).ready).toBe(true); + }); + + it('reconcileDependencyBlockers fails closed when the lookup errors', async () => { + upsertTaskState('AGT-ERR-BLOCKER', { + execution: { status: 'in_progress', retryCount: 0 }, + linearState: 'In Review', + }); + upsertTaskState('AGT-ERR-DEPENDENT', { + dependencyIssueIds: ['AGT-ERR-BLOCKER'], + execution: { status: 'blocked', retryCount: 0 }, + linearState: 'Todo', + }); + const future = Date.now() + 2 * 60 * 60_000; + + const lookupIssueState = async () => ({ ok: false as const, error: 'rate limited' }); + const result = await reconcileDependencyBlockers({ source: { lookupIssueState }, now: future }); + + expect(result.lookedUp).toBe(1); + expect(result.resolved).toBe(0); + expect(result.released).toBe(0); + expect(getTaskState('AGT-ERR-BLOCKER')?.linearState).toBe('In Review'); + expect(getTaskState('AGT-ERR-BLOCKER')?.dependencyLookupFailed).toBe(true); + + const retried: string[] = []; + const samePass = await reconcileDependencyBlockers({ + source: { lookupIssueState: async (id) => { retried.push(id); return { ok: false as const, error: 'rate limited' }; } }, + now: future, + }); + expect(retried).toEqual([]); + expect(samePass.lookedUp).toBe(0); + + const afterErrorWindow = await reconcileDependencyBlockers({ + source: { lookupIssueState: async (id) => { retried.push(id); return { ok: false as const, error: 'rate limited' }; } }, + now: future + 15 * 60_000, + }); + expect(retried).toEqual(['AGT-ERR-BLOCKER']); + expect(afterErrorWindow.lookedUp).toBe(1); + }); + + it('treats a Canceled blocker as terminal so dependents are not waiting on dead work', () => { + // vela 2026-09-09: 20 locally-Canceled STO-* ids were the first + // reconcileDependencyBlockers candidates. isResolved ignored Canceled, so + // they never left the set and maxLookups never reached AGT-4207 (Done on + // Linear, 157th). DecisionEngine gates on getTaskReadiness, which must + // agree — otherwise even a perfect reconciler cannot unblock the queue. + upsertTaskState('STO-CANCELED', { + execution: { status: 'backlog', retryCount: 0 }, + linearState: 'Canceled', + }); + upsertTaskState('AGT-WAITING', { + dependencyIssueIds: ['STO-CANCELED'], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + + const waiting = { + id: 'AGT-WAITING', source: 'linear' as const, title: 'waiting', priority: 2, + createdAt: Date.now(), issueId: 'AGT-WAITING', + }; + expect(getTaskReadiness(waiting).ready).toBe(true); + expect(getTaskReadiness(waiting).blockedBy).toEqual([]); + }); + + it('reconcileDependencyBlockers skips locally-Canceled blockers and looks up a later Done one under maxLookups', async () => { + const future = Date.now() + 2 * 60 * 60_000; + for (let i = 0; i < 20; i++) { + const id = `STO-CANCELED-${String(i).padStart(2, '0')}`; + upsertTaskState(id, { + execution: { status: 'blocked', retryCount: 0 }, + linearState: 'Canceled', + }); + upsertTaskState(`DEP-ON-CANCELED-${i}`, { + dependencyIssueIds: [id], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + } + upsertTaskState('AGT-4207', { + execution: { status: 'in_progress', retryCount: 0 }, + linearState: 'In Review', + }); + upsertTaskState('AGT-4209', { + dependencyIssueIds: ['AGT-4207'], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + + const lookedUp: string[] = []; + const lookupIssueState = async (id: string) => { + lookedUp.push(id); + if (id === 'AGT-4207') return { ok: true as const, issue: { state: 'Done', stateType: 'completed' } }; + throw new Error(`should not look up ${id}`); + }; + + const result = await reconcileDependencyBlockers({ + source: { lookupIssueState }, + now: future, + maxLookups: 20, + }); + + expect(lookedUp).toEqual(['AGT-4207']); + expect(result.eligible).toBe(1); + expect(result.lookedUp).toBe(1); + expect(result.resolved).toBe(1); + expect(result.released).toBe(1); + expect(getTaskReadiness({ + id: 'AGT-4209', source: 'linear' as const, title: 'dependent', priority: 2, + createdAt: Date.now(), issueId: 'AGT-4209', + }).ready).toBe(true); + }); + + it('reconcileDependencyBlockers rotates past recently-checked open blockers on the next pass', async () => { + const future = Date.now() + 2 * 60 * 60_000; + for (let i = 0; i < 20; i++) { + const id = `OPEN-${String(i).padStart(2, '0')}`; + upsertTaskState(id, { + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + upsertTaskState(`DEP-ON-OPEN-${i}`, { + dependencyIssueIds: [id], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + } + upsertTaskState('DONE-BLOCKER-ZZZ', { + execution: { status: 'in_progress', retryCount: 0 }, + linearState: 'In Review', + }); + upsertTaskState('DEP-ON-DONE', { + dependencyIssueIds: ['DONE-BLOCKER-ZZZ'], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + + const lookedUp: string[] = []; + const lookupIssueState = async (id: string) => { + lookedUp.push(id); + if (id === 'DONE-BLOCKER-ZZZ') return { ok: true as const, issue: { state: 'Done', stateType: 'completed' } }; + return { ok: true as const, issue: { state: 'Todo', stateType: 'unstarted' } }; + }; + + const first = await reconcileDependencyBlockers({ + source: { lookupIssueState }, + now: future, + maxLookups: 20, + }); + expect(first.lookedUp).toBe(20); + expect(first.resolved).toBe(0); + expect(lookedUp).not.toContain('DONE-BLOCKER-ZZZ'); + + lookedUp.length = 0; + const second = await reconcileDependencyBlockers({ + source: { lookupIssueState }, + now: future, + maxLookups: 20, + }); + expect(lookedUp).toEqual(['DONE-BLOCKER-ZZZ']); + expect(second.lookedUp).toBe(1); + expect(second.resolved).toBe(1); + expect(second.released).toBe(1); + expect(getTaskReadiness({ + id: 'DEP-ON-DONE', source: 'linear' as const, title: 'dependent', priority: 2, + createdAt: Date.now(), issueId: 'DEP-ON-DONE', + }).ready).toBe(true); + }); + + it('treats a Duplicate blocker as terminal so dependents are not waiting on it', () => { + upsertTaskState('AGT-4115', { + execution: { status: 'in_progress', retryCount: 0 }, + linearState: 'Duplicate', + }); + upsertTaskState('AGT-4121', { + dependencyIssueIds: ['AGT-4115'], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + expect(getTaskReadiness({ + id: 'AGT-4121', source: 'linear' as const, title: 'waiting', priority: 2, + createdAt: Date.now(), issueId: 'AGT-4121', + }).ready).toBe(true); + }); + + it('reconcileDependencyBlockers looks up heartbeat-priority deps before older out-of-scope ones', async () => { + const future = Date.now() + 2 * 60 * 60_000; + for (let i = 0; i < 20; i++) { + const id = `OLD-${String(i).padStart(2, '0')}`; + upsertTaskState(id, { + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + upsertTaskState(`DEP-ON-OLD-${i}`, { + dependencyIssueIds: [id], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + } + upsertTaskState('AGT-4114', { + execution: { status: 'backlog', retryCount: 0 }, + linearState: 'Backlog', + }); + upsertTaskState('AGT-4121-LIVE', { + dependencyIssueIds: ['AGT-4114'], + execution: { status: 'todo', retryCount: 0 }, + linearState: 'Todo', + }); + + const lookedUp: string[] = []; + const lookupIssueState = async (id: string) => { + lookedUp.push(id); + if (id === 'AGT-4114') return { ok: true as const, issue: { state: 'Done', stateType: 'completed' } }; + return { ok: true as const, issue: { state: 'Todo', stateType: 'unstarted' } }; + }; + + const result = await reconcileDependencyBlockers({ + source: { lookupIssueState }, + now: future, + maxLookups: 20, + priorityDepIds: new Set(['AGT-4114']), + }); + + expect(lookedUp[0]).toBe('AGT-4114'); + expect(result.resolved).toBe(1); + expect(getTaskReadiness({ + id: 'AGT-4121-LIVE', source: 'linear' as const, title: 'live', priority: 2, + createdAt: Date.now(), issueId: 'AGT-4121-LIVE', + }).ready).toBe(true); + }); + it('reconciles stale in_progress against Linear state (R5)', () => { // Operator parks an actively-running issue → local in_progress is stale. markTaskInProgress('KT-400', { linearState: 'In Progress' }); @@ -334,6 +732,30 @@ describe('task state store', () => { expect(parent?.linearState).toBe('Done'); }); + it('does not complete a parent when a child is Canceled — that is not Done', () => { + upsertTaskState('PARENT-C', { + childIssueIds: ['CHILD-C1', 'CHILD-C2'], + execution: { status: 'decomposed', retryCount: 0 }, + linearState: 'In Progress', + updatedAt: new Date().toISOString(), + }); + upsertTaskState('CHILD-C1', { + parentIssueId: 'PARENT-C', + execution: { status: 'done', retryCount: 0 }, + linearState: 'Done', + updatedAt: new Date().toISOString(), + }); + upsertTaskState('CHILD-C2', { + parentIssueId: 'PARENT-C', + execution: { status: 'backlog', retryCount: 0 }, + linearState: 'Canceled', + updatedAt: new Date().toISOString(), + }); + + expect(completeParentIfChildrenDone('CHILD-C1')).toBeNull(); + expect(getTaskState('PARENT-C')?.execution.status).toBe('decomposed'); + }); + it('hydrates canonical state from the latest Linear sync comment', () => { const older = buildTaskStateSyncComment( upsertTaskState('ISSUE-9', { diff --git a/src/taskState/store.ts b/src/taskState/store.ts index a3ca1146..8d632eda 100644 --- a/src/taskState/store.ts +++ b/src/taskState/store.ts @@ -72,6 +72,10 @@ export const OpenSwarmTaskStateSchema = z.object({ execution: ExecutionStateSchema.default({ status: 'backlog', retryCount: 0 }), worktree: WorktreeStateSchema.default({}), updatedAt: z.string(), + /** Last time reconcileDependencyBlockers looked this id up. Distinct from + * updatedAt so a worker touching the blocker does not look like a lookup. */ + dependencyCheckedAt: z.string().optional(), + dependencyLookupFailed: z.boolean().optional(), }); function createTaskMap( @@ -149,6 +153,14 @@ function lockPidIsJudgeable(owner: StoreLockOwner): boolean { return sameProcessNamespace(owner.ns); } +/** Cache identity for the on-disk store. Includes inode so a same-size + * cross-process replacement (atomic rename) within one mtime tick still + * invalidates — mtime+size alone can miss that case. */ +function storeFileStamp(path: string): string { + const stat = statSync(path); + return `${stat.mtimeMs}:${stat.size}:${stat.ino}`; +} + function readStoreLockOwner(lockPath: string): StoreLockOwner | null { try { const value = JSON.parse(readFileSync(lockPath, 'utf8')) as Partial; @@ -171,9 +183,7 @@ function getStorePath(): string { function ensureStoreLoaded(): TaskStateStore { const path = getStorePath(); - const currentStamp = existsSync(path) - ? (() => { const stat = statSync(path); return `${stat.mtimeMs}:${stat.size}`; })() - : 'missing'; + const currentStamp = existsSync(path) ? storeFileStamp(path) : 'missing'; if (cache && cacheStamp === currentStamp) return cache; if (existsSync(path)) { @@ -348,7 +358,7 @@ function persistStore(): void { renameSync(temporaryPath, path); chmodSync(path, 0o600); const persistedStat = statSync(path); - cacheStamp = `${persistedStat.mtimeMs}:${persistedStat.size}`; + cacheStamp = `${persistedStat.mtimeMs}:${persistedStat.size}:${persistedStat.ino}`; // Persist the directory entry where the platform supports directory fsync. let directoryFd: number | undefined; @@ -648,6 +658,21 @@ function isResolved(state: OpenSwarmTaskState | undefined): boolean { return state.execution.status === 'done' || state.linearState === 'Done'; } +function isDependencyTerminalLinearState(linearState: string | undefined): boolean { + const name = linearState?.trim().toLowerCase(); + return name === 'canceled' || name === 'cancelled' || name === 'duplicate'; +} + +/** A blocker that can no longer move work forward. Done is the historical + * isResolved contract (parent-completion still uses that). Canceled/Duplicate + * are terminal for dependencies — vela 2026-09-09: Canceled STO-* ids occupied + * maxLookups, and AGT-4115 (Duplicate on Linear, In Progress locally) kept + * AGT-4121 blocked after AGT-4253. Matches trackerTerminalReconciler. */ +function isDependencyTerminal(state: OpenSwarmTaskState | undefined): boolean { + if (isResolved(state)) return true; + return isDependencyTerminalLinearState(state?.linearState); +} + export function getTaskReadiness(task: TaskItem): { ready: boolean; blockedBy: string[]; @@ -667,7 +692,8 @@ export function getTaskReadiness(task: TaskItem): { const reactivated = linearState === 'Todo' || linearState === 'In Progress' || - linearState === 'In Review'; + linearState === 'In Review' || + linearState === 'Backlog'; if (!reactivated) { return { ready: false, @@ -683,7 +709,7 @@ export function getTaskReadiness(task: TaskItem): { return { ready: true, blockedBy: [] }; } - const unresolved = dependencyIssueIds.filter((depId) => !isResolved(getTaskState(depId))); + const unresolved = dependencyIssueIds.filter((depId) => !isDependencyTerminal(getTaskState(depId))); if (unresolved.length > 0) { return { ready: false, @@ -702,7 +728,7 @@ export function releaseDependentTasks(completedIssueId: string): OpenSwarmTaskSt for (const state of Object.values(store.tasks)) { if (!state.dependencyIssueIds.includes(completedIssueId)) continue; - const unresolved = state.dependencyIssueIds.filter((depId) => !isResolved(store.tasks[depId])); + const unresolved = state.dependencyIssueIds.filter((depId) => !isDependencyTerminal(store.tasks[depId])); if (unresolved.length > 0) { upsertTaskState(state.issueId, { execution: { @@ -728,6 +754,168 @@ export function releaseDependentTasks(completedIssueId: string): OpenSwarmTaskSt return released; } +/** Structural subset of ITaskSource — avoids importing automation/taskSource.js, which + * already imports enrichTaskFromState from this module. */ +interface DependencyLookupSource { + lookupIssueState(issueIdOrIdentifier: string): Promise< + | { ok: true; issue: { state: string; stateType?: string } | null } + | { ok: false; error: string } + >; +} + +export interface DependencyBlockerReconcileOptions { + source: DependencyLookupSource | null; + /** Issue ids present in the current heartbeat's fresh fetch. A dependency id found + * here is still Todo/In Progress/In Review/Backlog — genuinely open, skip the lookup. */ + knownTaskIds?: ReadonlySet; + now?: number; + /** Only reconsider a dependent task that has sat blocked at least this long — + * a decomposition just created seconds ago is not yet stale. */ + staleAfterMs?: number; + maxLookups?: number; + /** Skip a blocker looked up this recently — same role as + * reconcileTrackerTerminalRuns' recheckAfterMs. Default 6h. */ + recheckAfterMs?: number; + /** Failed lookups retry sooner than a confirmed-open skip. Default 15m. */ + errorRecheckAfterMs?: number; + /** Dependency ids blocking tasks in this heartbeat's Linear fetch. + * Looked up before the rest of the store so a live queue cannot starve + * behind July KT-* Backlog from projects this daemon does not run. */ + priorityDepIds?: ReadonlySet; +} + +export interface DependencyBlockerReconcileResult { + /** Distinct unresolved dependency ids that were candidates this pass. */ + eligible: number; + lookedUp: number; + /** Dependencies confirmed terminal (Done/Cancelled) by a live lookup. */ + resolved: number; + /** Dependent tasks released as a result. */ + released: number; +} + +/** + * Reconcile dependency ids that a Done Linear issue leaves behind. + * + * `releaseDependentTasks` only fires when the daemon's own pipeline observes a + * completion (runnerExecution.ts). A blocker finished through any other path — + * merged by a human, completed in a different session — never triggers it, and + * getTaskReadiness's fallback to the locally cached `dependencyIssueIds` (used + * whenever the fresh Linear fetch's `blockedBy` comes back empty, which it always + * does once the blocker is Done and drops out of the slim fetch) then blocks the + * dependent forever. Terminal issues are invisible to the regular slim fetch + * (Todo/In Progress/In Review/Backlog only), so this is the explicit per-issue + * read that can see them — same shape as reconcileTrackerTerminalRuns, applied to + * taskState instead of the durable ledger. + * + * Canceled/Cancelled blockers are terminal for dependents (isDependencyTerminal) + * even though isResolved stays Done-only for completeParentIfChildrenDone. + * Live vela 2026-09-09: 20 locally-Canceled STO-* ids occupied maxLookups + * every heartbeat, so AGT-4207 (157th, Linear already Done) was never reached. + * + * Lookups are capped. Candidates due for a check are sorted heartbeat-priority + * first (deps of tasks in this fetch), then never-checked, then oldest-checked. + * A lookup — success or fail-closed — stamps dependencyCheckedAt so the same + * 20 cannot consume the cap on the next heartbeat. + * + * Only tasks still in a not-yet-executed phase (todo/ready/blocked) are + * considered. dependencyIssueIds is never cleared once a task moves on (done, + * in_progress, decomposed, ...) — without this filter, an unrelated stale + * dependency entry left on an already-finished task would get resolved by this + * sweep and incorrectly reset that task's status back to 'todo' via + * releaseDependentTasks. A dependent isn't necessarily marked 'blocked' up + * front — getTaskReadiness gates on dependencyIssueIds at read time regardless + * of the stored execution.status, and only releaseDependentTasks itself writes + * 'blocked' when it finds a dependency still outstanding — so 'todo'/'ready' + * both need to stay in scope, not just 'blocked'. + */ +const DEPENDENCY_RECONCILE_ELIGIBLE_STATUSES = new Set(['todo', 'ready', 'blocked']); + +function blockerCheckedAtMs(state: OpenSwarmTaskState | undefined): number { + if (!state?.dependencyCheckedAt) return 0; + const ms = Date.parse(state.dependencyCheckedAt); + return Number.isFinite(ms) ? ms : 0; +} + +function blockerUpdatedAtMs(state: OpenSwarmTaskState | undefined): number { + if (!state?.updatedAt) return 0; + const ms = Date.parse(state.updatedAt); + return Number.isFinite(ms) ? ms : 0; +} + +export async function reconcileDependencyBlockers( + options: DependencyBlockerReconcileOptions, +): Promise { + const result: DependencyBlockerReconcileResult = { eligible: 0, lookedUp: 0, resolved: 0, released: 0 }; + const { source } = options; + if (!source) return result; + + const now = options.now ?? Date.now(); + const staleAfterMs = options.staleAfterMs ?? 60 * 60_000; + const recheckAfterMs = options.recheckAfterMs ?? 6 * 60 * 60_000; + const errorRecheckAfterMs = options.errorRecheckAfterMs ?? 15 * 60_000; + const maxLookups = Math.max(1, Math.floor(options.maxLookups ?? 20)); + const knownTaskIds = options.knownTaskIds ?? new Set(); + const priorityDepIds = options.priorityDepIds ?? new Set(); + + const byId = new Map(); + for (const state of listTaskStates()) { + if (!DEPENDENCY_RECONCILE_ELIGIBLE_STATUSES.has(state.execution.status)) continue; + if (state.dependencyIssueIds.length === 0) continue; + const updatedAtMs = Date.parse(state.updatedAt); + // Current-fetch dependents are upserted this heartbeat, so updatedAt is + // always fresh. Skipping them is why AGT-4121 stayed blocked after AGT-4254 + // shipped: its six Linear-Done blockers were never eligible for lookup. + const fromCurrentFetch = knownTaskIds.has(state.issueId); + if (!fromCurrentFetch && Number.isFinite(updatedAtMs) && now - updatedAtMs < staleAfterMs) continue; + for (const depId of state.dependencyIssueIds) { + if (isDependencyTerminal(getTaskState(depId))) continue; + if (knownTaskIds.has(depId)) continue; // still open per this fetch — not stale + if (byId.has(depId)) continue; + const blocker = getTaskState(depId); + const lastChecked = blockerCheckedAtMs(blocker); + const skipFor = blocker?.dependencyLookupFailed ? errorRecheckAfterMs : recheckAfterMs; + if (lastChecked > 0 && now - lastChecked < skipFor) continue; + byId.set(depId, { + depId, + lastChecked, + lastSeen: blockerUpdatedAtMs(blocker), + priority: priorityDepIds.has(depId) ? 1 : 0, + }); + } + } + const candidates = [...byId.values()].sort( + (a, b) => b.priority - a.priority || a.lastChecked - b.lastChecked || a.lastSeen - b.lastSeen || a.depId.localeCompare(b.depId), + ); + result.eligible = candidates.length; + + const checkedAt = new Date(now).toISOString(); + for (const { depId } of candidates) { + if (result.lookedUp >= maxLookups) break; + let lookup: Awaited>; + try { + lookup = await source.lookupIssueState(depId); + } catch (error) { + lookup = { ok: false, error: error instanceof Error ? error.message : String(error) }; + } + result.lookedUp++; + if (lookup.ok && lookup.issue) { + updateTaskLinearState(depId, lookup.issue.state); + upsertTaskState(depId, { dependencyCheckedAt: checkedAt, dependencyLookupFailed: false }); + } else { + // Stamp on fail-closed so a persistent lookup error cannot monopolize + // the cap. Retries after errorRecheckAfterMs (15m), not recheckAfterMs (6h). + upsertTaskState(depId, { dependencyCheckedAt: checkedAt, dependencyLookupFailed: true }); + } + if (isDependencyTerminal(getTaskState(depId))) { + result.resolved++; + result.released += releaseDependentTasks(depId).length; + } + } + + return result; +} + export function completeParentIfChildrenDone(childIssueId: string): OpenSwarmTaskState | null { const childState = getTaskState(childIssueId); if (!childState?.parentIssueId) return null; diff --git a/src/taskState/storeClaimProcess.fixture.ts b/src/taskState/storeClaimProcess.fixture.ts index a9ebc716..4140adaf 100644 --- a/src/taskState/storeClaimProcess.fixture.ts +++ b/src/taskState/storeClaimProcess.fixture.ts @@ -1,8 +1,44 @@ -import { resetTaskStateStoreForTests, upsertTaskState } from './store.js'; +import { resetTaskStateStoreForTests, upsertTaskState, getTaskState } from './store.js'; +import { promises as fs } from 'node:fs'; -const [stateFile, issueId, delayText = '0'] = process.argv.slice(2); +const [stateFile, issueId, delayText = '0', mode] = process.argv.slice(2); if (!stateFile || !issueId) throw new Error('state file and issue id are required'); process.env.OPENSWARM_TASK_STATE_FILE = stateFile; resetTaskStateStoreForTests(); await new Promise((resolve) => setTimeout(resolve, Number(delayText))); upsertTaskState(issueId, { title: issueId }); + +// Same-size replacement within one mtime tick (new inode via unlink+write). +// Proves ensureStoreLoaded invalidates when mtime+size alone would match. +if (mode === '--same-size-replace') { + const warmed = getTaskState(issueId); + if (warmed?.title !== issueId) throw new Error('expected warm cache before replacement'); + + // Keep the JSON byte length identical (e.g. SWAP-SRC → SWAP-DST). + const swappedTitle = issueId.replace(/SRC$/, 'DST'); + if (swappedTitle === issueId || swappedTitle.length !== issueId.length) { + throw new Error(`issueId must end with SRC for same-size replace (got ${issueId})`); + } + + const original = await fs.readFile(stateFile, 'utf8'); + const marker = `"title": ${JSON.stringify(issueId)}`; + const replacement = `"title": ${JSON.stringify(swappedTitle)}`; + if (marker.length !== replacement.length) { + throw new Error(`titles must be same length for same-size replace (${marker.length} vs ${replacement.length})`); + } + const replaced = original.replace(marker, replacement); + if (replaced.length !== original.length) { + throw new Error(`replacement must keep the same byte length (${original.length} → ${replaced.length})`); + } + + const stat = await fs.stat(stateFile); + await fs.unlink(stateFile); + await fs.writeFile(stateFile, replaced, 'utf8'); + await fs.utimes(stateFile, stat.atime, stat.mtime); + + // Do not reset the cache — invalidation must come from the stamp (incl. ino). + const after = getTaskState(issueId); + if (after?.title !== swappedTitle) { + throw new Error(`cache missed replacement: got ${JSON.stringify(after?.title)}`); + } +}