From bf72d774a3e0a03952a525e7eaab52c9a6b36c6d Mon Sep 17 00:00:00 2001 From: unohee Date: Mon, 28 Sep 2026 15:11:46 +0900 Subject: [PATCH 1/3] fix(state-integrity): salvage durable state fences and the review quality harness MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Salvages the still-unique work from four abandoned draft PRs (#757, #771, #773, #777) onto current main. Each PR's merge base is 172-247 commits behind, so only hunks absent from main were taken; work main already implemented independently was kept as-is. #757 (automation-state concurrency/recovery) - storeFileStamp(): the task-state cache stamp now includes the inode, so a same-size cross-process replacement (atomic rename) inside one mtime tick invalidates the cache instead of serving a stale snapshot forever. - validateExecution(): rejects incomplete step results and DAG-illegal lifecycle states before a workflow execution is persisted. - dev.ts: the close handler's reporting is wrapped in try/catch/finally so onComplete and activeTasks.delete always run. - memoryBridge.ts: one memory_linked event, not two (linkMemory already emits it). - linearBridge.ts: a failed SDK load no longer pins the rejected init promise, so the next call can retry. #771 (review quality harness) - New src/verify/qualityHarness.ts: full-tree, fail-closed static scan over every tracked source file plus isolated verify commands, wired into `review --max`, `--harness-only`, and the markdown audit report. - The harness test fixture put the empty catch body on its own line, which the per-line bs detector can never match — the detection assertion was vacuous. The fixture is single-line now and asserts a real finding. #773 (transactional state updates) - changeStatus()/updateIssue() read the effective status inside the write transaction so the event log's oldValue cannot be stamped stale. - saveExecution() gains a definitionStamp fence that refuses a snapshot whose workflow definition was replaced underneath it. - dailyReporter: bounded per-project retry; the watermark advances only when every project succeeded after that retry. - projectUpdater: buildBoundedProjectDescription() reserves room for the summary before truncating, so the 255-char limit cannot chop it off. - projectHandler: a failed quarantine surfaces as its own error instead of being reported as a successful preserve. #777 (issue/delivery mutation scope) - linearBridge: pendingLinearMappings recovery + idempotencyKey, so a failed local mapping persist cannot orphan-recreate a Linear issue. - sqliteStore.createIssue() honors the documented idempotent-id contract. - backlogGrooming: validIssueIds is mandatory at parse time and apply time, so a hallucinated id can never reach a mutation path. - workerValidationEvidence: EXECUTABLE_SOURCE_FILE_RE keeps shell/sql/etc. source under data dirs in scope for validation evidence. - projectUpdater: fetchProjectOverviewIssues surfaces missing/repeated cursors instead of silently truncating the page walk. --- src/__tests__/issueStore.test.ts | 13 + src/agents/workerValidationEvidence.test.ts | 17 + src/agents/workerValidationEvidence.ts | 9 +- .../backlogGrooming.coverage.test.ts | 30 +- src/automation/backlogGrooming.test.ts | 20 +- src/automation/backlogGrooming.ts | 20 +- src/automation/dailyReporter.retry.test.ts | 65 +++ src/automation/dailyReporter.ts | 68 ++- .../dailyReporter.watermark.test.ts | 6 +- src/cli.ts | 8 + src/cli/projectHandler.coverage.test.ts | 10 + src/cli/projectHandler.ts | 17 +- src/cli/reviewAudit.test.ts | 52 +++ src/cli/reviewAudit.ts | 47 ++ src/cli/reviewMaxCommand.tsx | 69 +++ src/cli/reviewMaxHarness.smoke.test.ts | 96 ++++ src/issues/linearBridge.recovery.test.ts | 115 +++++ src/issues/linearBridge.ts | 122 +++++- src/issues/memoryBridge.ts | 5 +- src/issues/sqliteStore.test.ts | 10 + src/issues/sqliteStore.ts | 37 +- src/linear/index.ts | 2 +- src/linear/projectUpdater.boundedDesc.test.ts | 24 + src/linear/projectUpdater.pagination.test.ts | 107 +++++ src/linear/projectUpdater.ts | 68 ++- src/orchestration/workflow.coverage.test.ts | 40 ++ src/orchestration/workflow.test.ts | 43 ++ src/orchestration/workflow.ts | 102 ++++- src/support/dev.ts | 41 +- src/taskState/store.ts | 14 +- src/taskState/storeFileStamp.test.ts | 79 ++++ src/verify/qualityHarness.test.ts | 134 ++++++ src/verify/qualityHarness.ts | 410 ++++++++++++++++++ 33 files changed, 1791 insertions(+), 109 deletions(-) create mode 100644 src/automation/dailyReporter.retry.test.ts create mode 100644 src/cli/reviewMaxHarness.smoke.test.ts create mode 100644 src/issues/linearBridge.recovery.test.ts create mode 100644 src/linear/projectUpdater.boundedDesc.test.ts create mode 100644 src/linear/projectUpdater.pagination.test.ts create mode 100644 src/taskState/storeFileStamp.test.ts create mode 100644 src/verify/qualityHarness.test.ts create mode 100644 src/verify/qualityHarness.ts diff --git a/src/__tests__/issueStore.test.ts b/src/__tests__/issueStore.test.ts index ee390ddd..771ad039 100644 --- a/src/__tests__/issueStore.test.ts +++ b/src/__tests__/issueStore.test.ts @@ -129,6 +129,19 @@ describe('SqliteIssueStore', () => { expect(done?.closedAt).toBeDefined(); }); + + it('event oldValue reflects the status read inside the write transaction', () => { + const issue = store.createIssue({ projectId: 'p1', title: 'task' }); + store.changeStatus(issue.id, 'todo'); + store.changeStatus(issue.id, 'in_progress'); + + // getEvents returns newest-first (created_at DESC, rowid DESC). + const events = store.getEvents(issue.id).filter((e) => e.type === 'status_changed'); + expect(events.map((e) => [e.oldValue, e.newValue])).toEqual([ + ['todo', 'in_progress'], + ['backlog', 'todo'], + ]); + }); }); describe('listIssues', () => { diff --git a/src/agents/workerValidationEvidence.test.ts b/src/agents/workerValidationEvidence.test.ts index a5dc3db6..f8b8d95d 100644 --- a/src/agents/workerValidationEvidence.test.ts +++ b/src/agents/workerValidationEvidence.test.ts @@ -67,6 +67,23 @@ describe('missingWorkerValidationIssues', () => { })).length).toBeGreaterThan(0); }); + it('requires validation for executable formats under locale/i18n dirs', () => { + // sh/swift/sql (and other VALIDATION_RELEVANT executables) must not inherit + // the data-only exemption that applies to json locale strings. + expect(missingWorkerValidationIssues(worker({ + filesChanged: ['src/locales/format.sh'], + commands: [], + })).length).toBeGreaterThan(0); + expect(missingWorkerValidationIssues(worker({ + filesChanged: ['src/i18n/Localizable.swift'], + commands: [], + })).length).toBeGreaterThan(0); + expect(missingWorkerValidationIssues(worker({ + filesChanged: ['src/locales/seed.sql'], + commands: [], + })).length).toBeGreaterThan(0); + }); + it('treats a source module named readme.ts as code, not docs', () => { // README.md is docs; readme.ts is a real module and must hit the gate. expect(missingWorkerValidationIssues(worker({ diff --git a/src/agents/workerValidationEvidence.ts b/src/agents/workerValidationEvidence.ts index 8b5a9232..a78d04b2 100644 --- a/src/agents/workerValidationEvidence.ts +++ b/src/agents/workerValidationEvidence.ts @@ -11,6 +11,10 @@ const DOC_ONLY_FILE_RE = /(^|\/)(README|CHANGELOG|LICENSE|NOTICE)(\.(md|mdx|txt| // nothing to build or test on their own; exempt them so a data-only edit does // not get bounced for "no validation command". const DATA_ONLY_DIR_RE = /(^|\/)(locales?|i18n|fixtures?|__fixtures__|__snapshots__|snapshots?|__mocks__|mocks?|testdata|test-data)\//i; +// Executable / source formats under data dirs still require validation evidence. +// Broader than TESTER_CODE_FILE_RE: includes sh/swift/sql/etc. that VALIDATION_RELEVANT +// already tracks, but excludes pure data formats (json/yaml/toml). +const EXECUTABLE_SOURCE_FILE_RE = /\.(ts|tsx|mts|cts|js|jsx|mjs|cjs|py|rs|go|java|rb|c|cc|cpp|h|hpp|swift|kt|kts|scala|cs|php|sh|bash|zsh|sql)$/i; const VALIDATION_COMMAND_RE = /\b(npm\s+(?:test|run\s+(?:test|build|lint|typecheck|check|ci|verify|validate|smoke))|pnpm\s+(?:test|run\s+(?:test|build|lint|typecheck|check|ci|verify|validate|smoke))|yarn\s+(?:test|run\s+(?:test|build|lint|typecheck|check|ci|verify|validate|smoke))|bun\s+(?:test|run\s+(?:test|build|lint|typecheck|check|ci|verify|validate|smoke))|vitest|jest|mocha|pytest|ruff|mypy|pyright|tsc|eslint|oxlint|cargo\s+(?:check|test|clippy|build)|go\s+(?:test|vet|build)|swift\s+test|gradle\s+(?:test|build|check)|mvn\s+(?:test|verify)|make\b|cmake\b|py_compile|compileall|clippy|fmt\s+--check)\b/i; // Anchored at each segment start: a leading inspection verb means that segment // ran no validation (e.g. `rg "npm test"` searches for the string, it does not @@ -23,8 +27,9 @@ export function isValidationRelevantFile(file: string): boolean { if (/(^|\/)docs?\//i.test(file)) return false; if (VALIDATION_RELEVANT_BASENAME_RE.test(file)) return true; // Data/asset trees (locale, fixtures, snapshots, mocks) are exempt ONLY for - // non-code assets. A real source module under such a dir still needs a check. - if (DATA_ONLY_DIR_RE.test(file) && !TESTER_CODE_FILE_RE.test(file)) return false; + // non-executable assets. Source/executable modules under such a dir (including + // sh/swift/sql and locale i18n modules) still need a validation check. + if (DATA_ONLY_DIR_RE.test(file) && !EXECUTABLE_SOURCE_FILE_RE.test(file)) return false; return VALIDATION_RELEVANT_FILE_RE.test(file) && !DOC_ONLY_FILE_RE.test(file); } diff --git a/src/automation/backlogGrooming.coverage.test.ts b/src/automation/backlogGrooming.coverage.test.ts index 1245b04b..b0169b3e 100644 --- a/src/automation/backlogGrooming.coverage.test.ts +++ b/src/automation/backlogGrooming.coverage.test.ts @@ -91,7 +91,7 @@ describe('parseBacklogGroomingOutput edge branches', () => { it('drops non-object decision entries', () => { const result = parseBacklogGroomingOutput(`\`\`\`json {"decisions": [null, "not-an-object", 42, {"issueId":"id-1","status":"active","reason":"ok"}]} -\`\`\``); +\`\`\``, new Set(['id-1'])); expect(result.success).toBe(true); expect(result.decisions.map(d => d.issueId)).toEqual(['id-1']); }); @@ -99,13 +99,13 @@ describe('parseBacklogGroomingOutput edge branches', () => { it('drops decision entries missing required fields', () => { const result = parseBacklogGroomingOutput(`\`\`\`json {"decisions": [{"identifier":"INT-9"}, {"issueId":"id-1","status":"active"}, {"issueId":"id-2","reason":"no status"}]} -\`\`\``); +\`\`\``, new Set(['id-1', 'id-2'])); expect(result.success).toBe(true); expect(result.decisions).toEqual([]); }); it('returns a failure result when the output cannot be parsed as JSON', () => { - const result = parseBacklogGroomingOutput('not json at all, no brace here'); + const result = parseBacklogGroomingOutput('not json at all, no brace here', new Set()); expect(result.success).toBe(false); expect(result.decisions).toEqual([]); expect(result.error).toBeTruthy(); @@ -131,7 +131,10 @@ describe('runBacklogGroomingPlanner', () => { stdout: '```json\n{"decisions":[{"issueId":"id-1","status":"active","reason":"fine"}]}\n```', stderr: '', }); - const result = await runBacklogGroomingPlanner({ tasks: [baseTask()], projectPath: '/repo' }); + const result = await runBacklogGroomingPlanner({ + tasks: [baseTask({ issueId: 'id-1' })], + projectPath: '/repo', + }); expect(result.success).toBe(true); expect(result.decisions.map(d => d.issueId)).toEqual(['id-1']); }); @@ -142,11 +145,28 @@ describe('runBacklogGroomingPlanner', () => { stdout: '```json\n{"decisions":[{"issueId":"id-1","status":"active","reason":"fine"}]}\n```', stderr: 'warning: partial output', }); - const result = await runBacklogGroomingPlanner({ tasks: [baseTask()], projectPath: '/repo' }); + const result = await runBacklogGroomingPlanner({ + tasks: [baseTask({ issueId: 'id-1' })], + projectPath: '/repo', + }); expect(result.success).toBe(true); expect(result.decisions).toHaveLength(1); }); + it('drops out-of-scope decision ids from planner output', async () => { + spawnCli.mockResolvedValueOnce({ + exitCode: 0, + stdout: '```json\n{"decisions":[{"issueId":"other","status":"stale","reason":"nope","closeState":"Done"}]}\n```', + stderr: '', + }); + const result = await runBacklogGroomingPlanner({ + tasks: [baseTask({ issueId: 'id-1' })], + projectPath: '/repo', + }); + expect(result.success).toBe(true); + expect(result.decisions).toEqual([]); + }); + it('reports stderr as the error when exit code is non-zero and stdout is empty', async () => { spawnCli.mockResolvedValueOnce({ exitCode: 1, stdout: ' ', stderr: 'adapter blew up' }); const result = await runBacklogGroomingPlanner({ tasks: [baseTask()], projectPath: '/repo' }); diff --git a/src/automation/backlogGrooming.test.ts b/src/automation/backlogGrooming.test.ts index dce33d8c..95b1a6da 100644 --- a/src/automation/backlogGrooming.test.ts +++ b/src/automation/backlogGrooming.test.ts @@ -28,10 +28,11 @@ describe('backlogGrooming (INT-1609)', () => { "decisions": [ {"issueId":"id-1","identifier":"INT-1","status":"stale","reason":"implemented","evidence":["src/a.ts:10"],"closeState":"Done"}, {"issueId":"id-2","status":"bogus","reason":"bad"}, - {"issueId":"id-3","status":"needs_update","reason":"drifted","updatedDescription":"new body"} + {"issueId":"id-3","status":"needs_update","reason":"drifted","updatedDescription":"new body"}, + {"issueId":"hallucinated","status":"stale","reason":"out of scope","closeState":"Done"} ] } -\`\`\``); +\`\`\``, new Set(['id-1', 'id-2', 'id-3'])); expect(result.success).toBe(true); expect(result.decisions.map(d => d.issueId)).toEqual(['id-1', 'id-3']); expect(result.decisions[0].closeState).toBe('Done'); @@ -63,12 +64,25 @@ describe('backlogGrooming (INT-1609)', () => { { issueId: 'id-2', status: 'needs_update', reason: 'drifted', evidence: ['src/a.ts:1'], updatedDescription: 'new body' }, { issueId: 'id-3', status: 'stale', reason: 'implemented', evidence: ['src/b.ts:2'], closeState: 'Done' }, ], - }, 'apply'); + }, 'apply', new Set(['id-1', 'id-2', 'id-3'])); expect(applied).toEqual({ commented: 2, failedComments: 0, updatedDescriptions: 1, moved: 1, movedIssueIds: ['id-3'], skippedUnknown: 0 }); expect(src.updateDescription).toHaveBeenCalledWith('id-2', 'new body'); expect(src.updateState).toHaveBeenCalledWith('id-3', 'Done'); }); + it('refuses apply mutations when no scope Set is supplied', async () => { + const src = source(); + const applied = await applyBacklogGrooming(src, { + success: true, + decisions: [ + { issueId: 'id-1', status: 'stale', reason: 'implemented', evidence: ['src/a.ts:1'], closeState: 'Done' }, + ], + }, 'apply'); + expect(applied.moved).toBe(0); + expect(applied.skippedUnknown).toBe(1); + expect(src.updateState).not.toHaveBeenCalled(); + }); + it('does not count a stale issue as moved when state transition fails', async () => { const src = source(); src.updateState.mockResolvedValueOnce(false); diff --git a/src/automation/backlogGrooming.ts b/src/automation/backlogGrooming.ts index 07e83186..89a3a7f4 100644 --- a/src/automation/backlogGrooming.ts +++ b/src/automation/backlogGrooming.ts @@ -131,7 +131,10 @@ Rules: - Keep updatedDescription concise and implementation-ready.`; } -export function parseBacklogGroomingOutput(output: string): BacklogGroomingResult { +export function parseBacklogGroomingOutput( + output: string, + validIssueIds: Set, +): BacklogGroomingResult { try { const fence = output.match(/```json\s*([\s\S]*?)```/i); const jsonText = fence?.[1] ?? output.slice(output.indexOf('{')); @@ -142,8 +145,11 @@ export function parseBacklogGroomingOutput(output: string): BacklogGroomingResul const d = item as Partial; if (!d.issueId || !d.status || !d.reason) return []; if (!['active', 'needs_update', 'stale'].includes(d.status)) return []; + const issueId = String(d.issueId); + // Drop hallucinated IDs before any downstream mutation path can see them. + if (!validIssueIds.has(issueId)) return []; return [{ - issueId: String(d.issueId), + issueId, identifier: d.identifier ? String(d.identifier) : undefined, status: d.status, reason: String(d.reason), @@ -177,7 +183,10 @@ export async function runBacklogGroomingPlanner(options: RunBacklogGroomingOptio if (raw.exitCode !== 0 && !raw.stdout.trim()) { return { success: false, decisions: [], error: raw.stderr.slice(0, 500) || `Planner adapter exited with code ${raw.exitCode}` }; } - return parseBacklogGroomingOutput(raw.stdout); + const validIssueIds = new Set( + options.tasks.map(task => task.issueId || task.id).filter(Boolean), + ); + return parseBacklogGroomingOutput(raw.stdout, validIssueIds); } catch (error) { return { success: false, decisions: [], error: error instanceof Error ? error.message : String(error) }; } @@ -198,7 +207,7 @@ export async function applyBacklogGrooming( source: ITaskSource, result: BacklogGroomingResult, mode: BacklogGroomingMode = 'comment', - validIssueIds?: Set, + validIssueIds: Set = new Set(), ): Promise { const applied: ApplyBacklogGroomingResult = { commented: 0, @@ -209,8 +218,9 @@ export async function applyBacklogGrooming( skippedUnknown: 0, }; if (!result.success) return applied; + // Scope is mandatory: an empty/missing Set must not mutate arbitrary IDs. for (const decision of result.decisions) { - if (validIssueIds && !validIssueIds.has(decision.issueId)) { + if (!validIssueIds.has(decision.issueId)) { applied.skippedUnknown++; continue; } diff --git a/src/automation/dailyReporter.retry.test.ts b/src/automation/dailyReporter.retry.test.ts new file mode 100644 index 00000000..aa9b6fe6 --- /dev/null +++ b/src/automation/dailyReporter.retry.test.ts @@ -0,0 +1,65 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { rmSync } from 'node:fs'; + +// Isolate the watermark file: generateDailyReports skips a window whose +// watermark already exists, and the real file lives in the developer's home. +const homeState = vi.hoisted(() => ({ + home: `/tmp/openswarm-daily-retry-${process.pid}`, +})); + +vi.mock('node:os', async (importOriginal) => ({ + ...(await importOriginal()), + homedir: () => homeState.home, +})); + +const postStatusUpdateMock = vi.fn(); +vi.mock('../linear/index.js', () => ({ + postStatusUpdate: (...args: unknown[]) => postStatusUpdateMock(...args), +})); + +import { + generateDailyReports, + setLinearClient, + setTeamId, +} from './dailyReporter.js'; + +describe('generateDailyReports retry targeting', () => { + beforeEach(() => { + rmSync(homeState.home, { recursive: true, force: true }); + postStatusUpdateMock.mockReset(); + setTeamId('team-1'); + }); + + it('retries only projects whose first publication update failed', async () => { + const projects = [ + { id: 'p-ok', name: 'Alpha', state: 'started' }, + { id: 'p-fail', name: 'Beta', state: 'started' }, + { id: 'p-ok2', name: 'Gamma', state: 'started' }, + ]; + + setLinearClient({ + team: async () => ({ + projects: async () => ({ + nodes: projects, + pageInfo: { hasNextPage: false, endCursor: null }, + }), + }), + } as never); + + postStatusUpdateMock.mockImplementation(async (id: string) => { + if (id === 'p-fail') { + // Fail once, then succeed on retry. + if (postStatusUpdateMock.mock.calls.filter((c) => c[0] === 'p-fail').length === 1) { + throw new Error('transient Linear error'); + } + } + }); + + await generateDailyReports(); + + const callsByProject = postStatusUpdateMock.mock.calls.map((c) => c[0] as string); + expect(callsByProject.filter((id) => id === 'p-ok')).toHaveLength(1); + expect(callsByProject.filter((id) => id === 'p-ok2')).toHaveLength(1); + expect(callsByProject.filter((id) => id === 'p-fail')).toHaveLength(2); + }); +}); diff --git a/src/automation/dailyReporter.ts b/src/automation/dailyReporter.ts index 2cc41c30..7eb133b2 100644 --- a/src/automation/dailyReporter.ts +++ b/src/automation/dailyReporter.ts @@ -120,17 +120,13 @@ export function stopDailyReporter(): void { } /** - * Generate daily status reports for all active projects - * Watermark is persisted ONLY after all reports succeed. + * Generate daily status reports for all active projects. + * Watermark is persisted ONLY after all reports succeed; failed projects get + * one bounded retry so a transient Linear error is not reported as a failure. */ export async function generateDailyReports(): Promise { - if (!linearClient) { - console.warn('[DailyReporter] LinearClient not set, skipping reports'); - return; - } - - if (!teamId) { - console.warn('[DailyReporter] Team ID not set, skipping reports'); + if (!linearClient || !teamId) { + console.warn('[DailyReporter] Linear client or team ID not configured'); return; } @@ -143,7 +139,6 @@ export async function generateDailyReports(): Promise { console.log('[DailyReporter] Generating daily reports...'); try { - // Fetch all active projects from Linear const team = await linearClient.team(teamId); if (!team) { console.warn('[DailyReporter] Team not found'); @@ -165,36 +160,63 @@ export async function generateDailyReports(): Promise { console.log(`[DailyReporter] Found ${activeProjects.length} active projects`); - // Generate status update for each project - let successCount = 0; - let failCount = 0; + // Track per-project publication outcome so retries target only failed projects + const projectResults: { id: string; name: string; ok: boolean }[] = []; for (const project of activeProjects) { try { const projectPath = projectPathMapping.get(project.id); await postStatusUpdate(project.id, project.name, projectPath); - successCount++; + projectResults.push({ id: project.id, name: project.name, ok: true }); } catch (err) { console.error(`[DailyReporter] Failed to post update for "${project.name}":`, err); - failCount++; + projectResults.push({ id: project.id, name: project.name, ok: false }); } } + const successCount = projectResults.filter(r => r.ok).length; + const failCount = projectResults.filter(r => !r.ok).length; + const failedProjects = projectResults.filter(r => !r.ok).map(r => r.name); + console.log(`[DailyReporter] Reports completed: ${successCount} success, ${failCount} failed`); - // Only persist watermark when ALL reports succeeded. - // If any failed, the watermark stays at the previous value so the next - // run retries the same window instead of skipping it. - if (failCount === 0) { + // Retry only failed projects (up to 1 retry each) + if (failCount > 0) { + console.log(`[DailyReporter] Retrying ${failCount} failed project(s): ${failedProjects.join(', ')}`); + for (const result of projectResults) { + if (!result.ok) { + try { + const projectPath = projectPathMapping.get(result.id); + await postStatusUpdate(result.id, result.name, projectPath); + result.ok = true; + console.log(`[DailyReporter] Retry succeeded for "${result.name}"`); + } catch (err) { + console.error(`[DailyReporter] Retry also failed for "${result.name}":`, err); + } + } + } + } + + // Outcome counts must reflect post-retry state so Discord/summary stay accurate. + const finalSuccessCount = projectResults.filter(r => r.ok).length; + const finalFailCount = projectResults.filter(r => !r.ok).length; + if (finalSuccessCount !== successCount || finalFailCount !== failCount) { + console.log(`[DailyReporter] After retry: ${finalSuccessCount} success, ${finalFailCount} failed`); + } + + // Only persist watermark when EVERY project succeeded (after the retry pass). + // A remaining failure keeps the previous watermark so the next run retries + // the same window instead of skipping it. + if (finalFailCount === 0) { writeWatermark(today); console.log(`[DailyReporter] Watermark persisted: ${today}`); } else { - console.warn(`[DailyReporter] ${failCount} report(s) failed — watermark NOT updated`); + console.warn(`[DailyReporter] ${finalFailCount} report(s) failed — watermark NOT updated`); } // Send summary to Discord - if (discordReporter && successCount > 0) { - await sendDiscordSummary(activeProjects.length, successCount, failCount); + if (discordReporter && finalSuccessCount > 0) { + await sendDiscordSummary(activeProjects.length, finalSuccessCount, finalFailCount); } } catch (error) { console.error('[DailyReporter] Failed to generate reports:', error); @@ -230,4 +252,4 @@ async function sendDiscordSummary( } catch (err) { console.error('[DailyReporter] Failed to send Discord summary:', err); } -} +} \ No newline at end of file diff --git a/src/automation/dailyReporter.watermark.test.ts b/src/automation/dailyReporter.watermark.test.ts index 47b036a4..c6129132 100644 --- a/src/automation/dailyReporter.watermark.test.ts +++ b/src/automation/dailyReporter.watermark.test.ts @@ -66,6 +66,10 @@ describe('daily reporter watermark', () => { it('retries after a partial failure because no watermark is written', async () => { testState.postStatusUpdate + // Both attempts in the first run fail: the bounded retry is one extra + // attempt inside that run, so a persistent failure still leaves the + // watermark unadvanced and the window is retried on the next run. + .mockRejectedValueOnce(new Error('temporary Linear failure')) .mockRejectedValueOnce(new Error('temporary Linear failure')) .mockResolvedValueOnce(undefined); configureReporter(); @@ -74,7 +78,7 @@ describe('daily reporter watermark', () => { expect(existsSync(watermarkFile)).toBe(false); await generateDailyReports(); - expect(testState.postStatusUpdate).toHaveBeenCalledTimes(2); + expect(testState.postStatusUpdate).toHaveBeenCalledTimes(3); expect(existsSync(watermarkFile)).toBe(true); }); }); diff --git a/src/cli.ts b/src/cli.ts index 09d1ce3e..39a8e845 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -359,12 +359,14 @@ program .option('--in-place', 'For --max --fix: edit the current working tree instead of an isolated worktree (no branch, no PR)') .option('--fix-rounds ', 'For --max --fix: optional round cap (default: until clean, with a two-hour safety budget)', parsePositiveIntegerOption) .option('--no-security-audit', 'For --max --fix: disable the default CodeQL audit gate') + .option('--harness-only', 'For --max: skip LLM reviewers; run only the deterministic quality harness (static scan + verify commands)') .option('--no-learn', 'For --max: do not record the audit findings into the repo knowledge memory') .action(async (opts: { path?: string; base?: string; issues?: string | boolean; issuesPerArea?: string | boolean; file?: string | boolean; adapter?: string; model?: string; debug?: boolean; json?: boolean; sarif?: string; readOnly?: boolean; maxTurns?: number; timeout?: number; max?: boolean; concurrency?: number; maxFilesPerArea?: number; yes?: boolean; dryRun?: boolean; out?: string; linear?: boolean; fallback?: string | boolean; fix?: boolean; inPlace?: boolean; fixRounds?: number; learn?: boolean; securityAudit?: boolean; + harnessOnly?: boolean; }) => { try { // --json/--sarif are declared on the shared `review` command but only the @@ -384,6 +386,11 @@ program process.exitCode = 2; return; } + if (opts.harnessOnly && !opts.max) { + console.error('--harness-only requires --max.'); + process.exitCode = 2; + return; + } if (opts.max) { const { runReviewMaxCommand, reviewMaxResultFailed } = await import('./cli/reviewMaxCommand.js'); const result = await runReviewMaxCommand({ @@ -409,6 +416,7 @@ program fixRounds: opts.fixRounds, learn: opts.learn, securityAudit: opts.securityAudit, + harnessOnly: opts.harnessOnly, }); // Exit contract (INT-3100): 2 = the gate did not run at all (no area // reviewed — quota/infra), 1 = it ran and failed. CI reads only this. diff --git a/src/cli/projectHandler.coverage.test.ts b/src/cli/projectHandler.coverage.test.ts index 1297256b..a06609b4 100644 --- a/src/cli/projectHandler.coverage.test.ts +++ b/src/cli/projectHandler.coverage.test.ts @@ -222,6 +222,16 @@ describe('loadRepos malformed-JSON recovery (via handleProjectList)', () => { expect(() => handleProjectList()).toThrow(/preserved as/); expect(renameSyncMock).toHaveBeenCalledOnce(); }); + + it('surfaces a quarantine failure when the corrupt file cannot be moved aside', () => { + readFileSyncMock.mockReturnValue('{ not valid json ,, }'); + existsSyncMock.mockImplementation((p: string) => typeof p === 'string' && p.endsWith('openswarm-repos.json')); + renameSyncMock.mockImplementation(() => { + throw new Error('EACCES: permission denied'); + }); + expect(() => handleProjectList()).toThrow(/quarantine failure/); + expect(errors.join('\n')).toMatch(/quarantine failed/); + }); }); describe('loadRepos defaults missing fields (via handleProjectList)', () => { diff --git a/src/cli/projectHandler.ts b/src/cli/projectHandler.ts index 1234bad3..4c37e248 100644 --- a/src/cli/projectHandler.ts +++ b/src/cli/projectHandler.ts @@ -48,8 +48,21 @@ export function loadRepos(file: string = REPOS_FILE): ReposConfig { }; } catch (error) { const recoveryPath = `${file}.corrupt-${Date.now()}`; - try { renameSync(file, recoveryPath); } catch { /* preserve original error below */ } - throw new Error(`Repository registry is malformed at ${file}; preserved as ${recoveryPath}: ${error instanceof Error ? error.message : String(error)}`); + let quarantined = false; + try { + renameSync(file, recoveryPath); + quarantined = true; + } catch { + // Quarantine itself failed — leave the corrupt file in place and surface that. + } + const detail = error instanceof Error ? error.message : String(error); + if (!quarantined) { + console.error(`Repository registry is malformed at ${file}; quarantine failed (left in place): ${detail}`); + throw new Error(`Repository registry quarantine failure at ${file}: ${detail}`, { cause: error }); + } + console.error(`Repository registry is malformed at ${file}; preserved as ${recoveryPath}: ${detail}`); + // Corrupt-but-quarantined is still a load failure for callers; do not silently recover. + throw new Error(`Repository registry is corrupt; preserved as ${recoveryPath}: ${detail}`, { cause: error }); } } diff --git a/src/cli/reviewAudit.test.ts b/src/cli/reviewAudit.test.ts index 3ec54ca9..467edda4 100644 --- a/src/cli/reviewAudit.test.ts +++ b/src/cli/reviewAudit.test.ts @@ -10,6 +10,8 @@ import { runMaxReview, mergeFallback, mergeSecurityAuditFindings, + mergeQualityHarnessResult, + QUALITY_HARNESS_AREA, type AuditArea, type AuditAreaResult, type AuditProgress, @@ -217,6 +219,56 @@ describe('mergeSecurityAuditFindings', () => { }); }); +describe('mergeQualityHarnessResult', () => { + it('always injects a harness area so clean scans still gate-ran', () => { + const base: AuditRun = { + results: [], + summary: aggregateAuditResults([]), + }; + const merged = mergeQualityHarnessResult(base, { + status: 'passed', + filesListed: 3, + filesScanned: 3, + findings: [], + commands: [{ name: 'typecheck', kind: 'typecheck', status: 'pass', detail: 'ok' }], + }); + expect(merged.summary.decision).toBe('approve'); + expect(merged.summary.completed).toBe(1); + expect(merged.results[0]?.area.label).toBe(QUALITY_HARNESS_AREA); + const md = formatAuditReport(merged.summary, 'repo', 'ts'); + expect(md).toContain(QUALITY_HARNESS_AREA); + expect(md).toContain('Verdict: APPROVE'); + }); + + it('rejects when static or command findings are errors', () => { + const base: AuditRun = { + results: [{ + area: { label: 'src', dir: 'src', files: ['src/a.ts'] }, + review: { decision: 'approve', feedback: '', issues: [], recommendedActions: [] }, + }], + summary: aggregateAuditResults([{ + area: { label: 'src', dir: 'src', files: ['src/a.ts'] }, + review: { decision: 'approve', feedback: '', issues: [], recommendedActions: [] }, + }]), + }; + const merged = mergeQualityHarnessResult(base, { + status: 'failed', + filesListed: 1, + filesScanned: 1, + findings: [{ + ruleId: 'openswarm/quality-truncated', + level: 'error', + message: 'too large', + filePath: 'src/a.ts', + }], + commands: [], + }); + expect(merged.summary.decision).toBe('reject'); + expect(merged.summary.issues.some((i) => i.includes('openswarm/quality-truncated'))).toBe(true); + expect(formatAuditReport(merged.summary, 'repo', 'ts')).toContain('Quality openswarm/quality-truncated'); + }); +}); + describe('formatAuditReport (INT-2022)', () => { it('renders markdown with verdict, failures, typed follow-ups, and issues', () => { const summary: AuditSummary = { diff --git a/src/cli/reviewAudit.ts b/src/cli/reviewAudit.ts index 000ac043..0a3a4c37 100644 --- a/src/cli/reviewAudit.ts +++ b/src/cli/reviewAudit.ts @@ -19,10 +19,14 @@ import { isInfraError } from '../adapters/errorClassification.js'; import { c, status } from '../support/colors.js'; import { sanitizeTerminalText } from '../tui/sanitize.js'; import type { SecurityFinding } from '../verify/securityAudit.js'; +import type { QualityHarnessResult } from '../verify/qualityHarness.js'; /** Synthetic area label carrying deterministic CodeQL findings into a review run. */ export const SECURITY_AUDIT_AREA = '.openswarm/codeql-security'; +/** Synthetic area label for the deterministic quality harness (static + verify cmds). */ +export const QUALITY_HARNESS_AREA = '.openswarm/quality-harness'; + /** Add deterministic CodeQL findings as a fixable, synthetic review area. */ export function mergeSecurityAuditFindings(run: AuditRun, findings: readonly SecurityFinding[]): AuditRun { const results = run.results.filter((result) => result.area.label !== SECURITY_AUDIT_AREA); @@ -48,6 +52,49 @@ export function mergeSecurityAuditFindings(run: AuditRun, findings: readonly Sec return { ...run, results, summary: aggregateAuditResults(results) }; } +/** + * Fold the deterministic quality harness into the audit run. Always injects an + * area so `--harness-only` still produces a gate-ran verdict and the markdown + * report records coverage even when the scan is clean. + */ +export function mergeQualityHarnessResult(run: AuditRun, harness: QualityHarnessResult): AuditRun { + const results = run.results.filter((result) => result.area.label !== QUALITY_HARNESS_AREA); + const files = [...new Set( + harness.findings.map((finding) => finding.filePath).filter((file): file is string => Boolean(file)), + )].sort(); + const issues = harness.findings.map((finding) => { + const location = finding.filePath + ? `${finding.filePath}${finding.line ? `:${finding.line}` : ''}` + : 'repository'; + return `Quality ${finding.ruleId} (${location}): ${finding.message}`; + }); + const commandLine = harness.commands.length === 0 + ? 'no verify commands discovered' + : harness.commands.map((c) => `${c.name}:${c.status}`).join(', '); + const coverage = `static ${harness.filesScanned}/${harness.filesListed}; commands: ${commandLine}`; + const decision: ReviewResult['decision'] = harness.findings.some((f) => f.level === 'error') + ? 'reject' + : harness.findings.length > 0 + ? 'revise' + : 'approve'; + const feedback = decision === 'approve' + ? `Deterministic quality harness passed (${coverage}).` + : `Deterministic quality harness findings (${coverage}):\n${issues.join('\n')}`; + results.push({ + area: { label: QUALITY_HARNESS_AREA, dir: '.', files }, + review: { + decision, + feedback, + issues, + recommendedActions: harness.findings.map((finding) => ({ + type: finding.ruleId.includes('command') ? 'test' : 'quality', + title: `Address ${finding.ruleId}${finding.filePath ? ` at ${finding.filePath}${finding.line ? `:${finding.line}` : ''}` : ''}`, + })), + }, + }); + return { ...run, results, summary: aggregateAuditResults(results) }; +} + // Source extensions and test patterns mirror src/knowledge/scanner.ts. Kept // local (not imported) because those are unexported module consts; the audit // only needs the stable subset and drift here is low-risk. diff --git a/src/cli/reviewMaxCommand.tsx b/src/cli/reviewMaxCommand.tsx index 984d8ea0..87db3559 100644 --- a/src/cli/reviewMaxCommand.tsx +++ b/src/cli/reviewMaxCommand.tsx @@ -23,6 +23,7 @@ import { oneLineError, mergeFallback, mergeSecurityAuditFindings, + mergeQualityHarnessResult, type AuditArea, type AuditRun, type AuditSummary, @@ -56,6 +57,7 @@ import { loadTrustedVerifyPlan, runDeterministicTester } from '../agents/determi import { buildFixRepositoryContext } from './fixPlanning.js'; import { collectFixRuntimePreflightIssues } from './fixPreflight.js'; import { DEFAULT_SECURITY_AUDIT_CONFIG, listTrackedSecurityFiles, runSecurityAudit, type SecurityFinding } from '../verify/securityAudit.js'; +import { runQualityHarness } from '../verify/qualityHarness.js'; /** * Best-effort verify config: `review --max` must still run in a repo with no — @@ -125,6 +127,11 @@ export interface ReviewMaxOptions { learn?: boolean; /** Disable the default-on CodeQL audit gate. */ securityAudit?: boolean; + /** + * Skip LLM area fan-out and run only the deterministic quality harness + * (static scan + isolated verify commands). (M0 / PLATFORM_ROADMAP) + */ + harnessOnly?: boolean; } export interface ReviewMaxCommandResult { @@ -420,6 +427,39 @@ export async function runReviewMaxCommand(rawOpts: ReviewMaxOptions = {}): Promi const concurrency = positiveIntegerOption(opts.concurrency, 4, '--concurrency'); const maxFilesPerArea = positiveIntegerOption(opts.maxFilesPerArea, 12, '--max-files-per-area'); + // Deterministic-only path: no LLM cost, no area fan-out. Still writes the + // audit report and participates in the same exit-code contract. (M0) + if (opts.harnessOnly) { + if (opts.fix) { + throw new Error('--harness-only cannot be combined with --fix'); + } + const verifyConfig = loadVerifyConfigBestEffort(); + console.log(status.running('Quality harness') + c.dim(' — static scan + isolated verify commands (no LLM)')); + const harness = await runQualityHarness(cwd, { verify: verifyConfig }); + console.log(c.dim( + ` Quality harness: ${harness.status}, scanned ${harness.filesScanned}/${harness.filesListed}, ` + + `${harness.findings.length} finding(s), ${harness.commands.length} command(s).`, + )); + let run: AuditRun = { results: [], summary: aggregateAuditResults([]) }; + run = mergeQualityHarnessResult(run, harness); + console.log(formatAuditSummary(run.summary)); + + const ts = new Date().toISOString().replace(/[:.]/g, '-').slice(0, 19); + const report = formatAuditReport(run.summary, basename(cwd) || cwd, ts); + const outPath = opts.out ?? join(cwd, '.openswarm', 'audit', `audit-${ts}.md`); + try { + await mkdir(dirname(outPath), { recursive: true }); + await writeFile(outPath, report, 'utf8'); + console.log(`\nReport saved: ${outPath}`); + } catch (e) { + console.warn(`Could not save report: ${e instanceof Error ? e.message : String(e)}`); + } + return { + decision: run.summary.decision, + gateRan: true, + }; + } + let files: string[]; try { files = listSourceFiles(cwd); @@ -778,6 +818,35 @@ export async function runReviewMaxCommand(rawOpts: ReviewMaxOptions = {}): Promi } } + // Deterministic quality harness (static full-tree + isolated verify commands). + // Runs for every --max so the final verdict and markdown report always carry + // CodeQL-style coverage evidence, not only LLM area notes. (M0 / AGT-3619) + try { + console.log(`\n${status.running('Quality harness')} ${c.dim('static scan + isolated verify commands')}`); + const harness = await runQualityHarness(workCwd, { verify: verifyConfig }); + console.log(c.dim( + ` Quality harness: ${harness.status}, scanned ${harness.filesScanned}/${harness.filesListed}, ` + + `${harness.findings.length} finding(s), ${harness.commands.length} command(s).`, + )); + run = mergeQualityHarnessResult(run, harness); + if (harness.findings.length > 0) { + console.log(formatAuditSummary(run.summary)); + } + } catch (error) { + run = mergeQualityHarnessResult(run, { + status: 'failed', + filesListed: 0, + filesScanned: 0, + findings: [{ + ruleId: 'openswarm/quality-runtime', + level: 'error', + message: `Quality harness aborted: ${error instanceof Error ? error.message : String(error)}`, + }], + commands: [], + }); + console.warn(status.warn(`Quality harness aborted — recorded as an explicit failure.`)); + } + // (3.6) Persist a markdown report so the result isn't lost to the scrollback. // Built here (after --fix) so it reflects the verified post-fix verdicts — // and so the Linear master issue below embeds the same final state. (INT-2022 / INT-2443) diff --git a/src/cli/reviewMaxHarness.smoke.test.ts b/src/cli/reviewMaxHarness.smoke.test.ts new file mode 100644 index 00000000..ec814cb6 --- /dev/null +++ b/src/cli/reviewMaxHarness.smoke.test.ts @@ -0,0 +1,96 @@ +// CLI smoke for `openswarm review --max --harness-only` (AGT-3619 / M0). +import { execFileSync } from 'node:child_process'; +import { mkdir, mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { dirname, join } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { afterEach, describe, expect, it } from 'vitest'; + +const roots: string[] = []; +const repoRoot = join(dirname(fileURLToPath(import.meta.url)), '../..'); +const cliEntry = join(repoRoot, 'src/cli.ts'); + +async function gitRepo(files: Record): Promise { + const root = await mkdtemp(join(tmpdir(), 'openswarm-harness-cli-')); + roots.push(root); + execFileSync('git', ['init'], { cwd: root, stdio: 'ignore' }); + execFileSync('git', ['config', 'user.email', 'harness@example.test'], { cwd: root, stdio: 'ignore' }); + execFileSync('git', ['config', 'user.name', 'Harness'], { cwd: root, stdio: 'ignore' }); + for (const [name, content] of Object.entries(files)) { + const path = join(root, name); + await mkdir(dirname(path), { recursive: true }); + await writeFile(path, content); + } + execFileSync('git', ['add', '-A'], { cwd: root, stdio: 'ignore' }); + execFileSync('git', ['commit', '-m', 'init'], { cwd: root, stdio: 'ignore' }); + return root; +} + +afterEach(async () => { + await Promise.all(roots.splice(0).map((root) => rm(root, { recursive: true, force: true }))); +}); + +describe('review --max --harness-only CLI smoke', () => { + it('runs the quality harness without LLM reviewers and writes a markdown report', async () => { + const root = await gitRepo({ + 'src/ok.ts': 'export const value = 1;\n', + 'package.json': JSON.stringify({ name: 'fixture', private: true }), + }); + const out = join(root, 'audit-report.md'); + const stdout = execFileSync( + process.execPath, + ['--import', 'tsx', cliEntry, 'review', '--max', '--harness-only', '--yes', '--no-linear', '--path', root, '--out', out], + { + cwd: repoRoot, + encoding: 'utf8', + env: { + ...process.env, + NO_COLOR: '1', + OPENSWARM_DISABLE_TELEMETRY: '1', + }, + stdio: ['ignore', 'pipe', 'pipe'], + timeout: 120_000, + }, + ); + expect(stdout).toMatch(/Quality harness/i); + expect(stdout).toMatch(/Verdict:\s*APPROVE/i); + const report = await readFile(out, 'utf8'); + expect(report).toContain('.openswarm/quality-harness'); + expect(report).toContain('Verdict: APPROVE'); + }, 120_000); + + it('fails closed when a tracked source file cannot be fully scanned', async () => { + const root = await gitRepo({ + 'src/huge.ts': 'x'.repeat(512 * 1024 + 32), + }); + const out = join(root, 'audit-report.md'); + let code = 0; + let combined = ''; + try { + combined = execFileSync( + process.execPath, + ['--import', 'tsx', cliEntry, 'review', '--max', '--harness-only', '--yes', '--no-linear', '--path', root, '--out', out], + { + cwd: repoRoot, + encoding: 'utf8', + env: { + ...process.env, + NO_COLOR: '1', + OPENSWARM_DISABLE_TELEMETRY: '1', + }, + stdio: ['ignore', 'pipe', 'pipe'], + timeout: 120_000, + }, + ); + } catch (error) { + const failure = error as { status?: number; stdout?: string; stderr?: string }; + code = failure.status ?? 1; + combined = `${failure.stdout ?? ''}\n${failure.stderr ?? ''}`; + } + expect(code).toBe(1); + expect(combined).toMatch(/quality-truncated|Verdict:\s*REJECT/i); + const report = await readFile(out, 'utf8'); + expect(report).toContain('openswarm/quality-truncated'); + expect(report).toContain('Verdict: REJECT'); + }, 120_000); +}); diff --git a/src/issues/linearBridge.recovery.test.ts b/src/issues/linearBridge.recovery.test.ts new file mode 100644 index 00000000..7c92d071 --- /dev/null +++ b/src/issues/linearBridge.recovery.test.ts @@ -0,0 +1,115 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { + __clearPendingLinearMappingsForTests, + __setLinearBridgeClientForTests, + pushToLinear, +} from './linearBridge.js'; +import { SqliteIssueStore } from './sqliteStore.js'; + +let dir: string | undefined; + +function dbPath(): string { + dir ??= mkdtempSync(join(tmpdir(), 'openswarm-linear-bridge-')); + return join(dir, 'issues.db'); +} + +function installFakeLinear(createIssue = vi.fn()) { + const fakeClient = { + createIssue, + team: vi.fn(async () => ({ + states: async () => ({ + nodes: [ + { id: 'state-todo', name: 'Todo' }, + { id: 'state-backlog', name: 'Backlog' }, + ], + }), + })), + }; + createIssue.mockResolvedValue({ + issue: Promise.resolve({ + id: 'lin-uuid-1', + identifier: 'AGT-1', + url: 'https://linear.app/agt-1', + }), + }); + __setLinearBridgeClientForTests(fakeClient, 'team-test'); + return { fakeClient, createIssue }; +} + +beforeEach(() => { + __clearPendingLinearMappingsForTests(); +}); + +afterEach(() => { + __clearPendingLinearMappingsForTests(); + __setLinearBridgeClientForTests(null); + if (dir) rmSync(dir, { recursive: true, force: true }); + dir = undefined; +}); + +describe('pushToLinear mapping recovery', () => { + it('returns the Linear id when local mapping persist fails and does not recreate on retry', async () => { + const { createIssue } = installFakeLinear(); + const store = new SqliteIssueStore(dbPath()); + const issue = store.createIssue({ projectId: 'p', title: 'recover-me', status: 'todo' }); + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const error = vi.spyOn(console, 'error').mockImplementation(() => {}); + + const updateIssue = vi.spyOn(store, 'updateIssue').mockImplementation(() => { + throw new Error('persist boom'); + }); + + const first = await pushToLinear(store, issue.id); + expect(first).toBe('lin-uuid-1'); + expect(createIssue).toHaveBeenCalledTimes(1); + // Mapping never landed locally. + expect(store.getIssue(issue.id)?.linearId).toBeUndefined(); + + const second = await pushToLinear(store, issue.id); + expect(second).toBe('lin-uuid-1'); + // Pending-map recovery must not call Linear create again. + expect(createIssue).toHaveBeenCalledTimes(1); + + updateIssue.mockRestore(); + warn.mockRestore(); + error.mockRestore(); + store.close(); + }); + + it('retries pending mapping persist and recovers without a second Linear create', async () => { + const { createIssue } = installFakeLinear(); + const store = new SqliteIssueStore(dbPath()); + const issue = store.createIssue({ projectId: 'p', title: 'retry-map', status: 'todo' }); + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const error = vi.spyOn(console, 'error').mockImplementation(() => {}); + + let failuresLeft = 4; // 3 persist attempts + 1 updateIssue-only recovery path + const realUpdate = store.updateIssue.bind(store); + const updateIssue = vi.spyOn(store, 'updateIssue').mockImplementation((id, patch) => { + if (failuresLeft > 0) { + failuresLeft -= 1; + throw new Error(`persist fail ${failuresLeft}`); + } + return realUpdate(id, patch); + }); + + const first = await pushToLinear(store, issue.id); + expect(first).toBe('lin-uuid-1'); + expect(store.getIssue(issue.id)?.linearId).toBeUndefined(); + expect(createIssue).toHaveBeenCalledTimes(1); + + updateIssue.mockRestore(); + const recovered = await pushToLinear(store, issue.id); + expect(recovered).toBe('lin-uuid-1'); + expect(createIssue).toHaveBeenCalledTimes(1); + expect(store.getIssue(issue.id)?.linearId).toBe('lin-uuid-1'); + expect(store.getIssue(issue.id)?.linearIdentifier).toBe('AGT-1'); + + warn.mockRestore(); + error.mockRestore(); + store.close(); + }); +}); diff --git a/src/issues/linearBridge.ts b/src/issues/linearBridge.ts index a9285bde..4bbee884 100644 --- a/src/issues/linearBridge.ts +++ b/src/issues/linearBridge.ts @@ -32,6 +32,7 @@ import { let linearClient: any = null; let linearTeamId: string = ''; let linearInitPromise: Promise | null = null; +let linearInitAttempts = 0; /** Per-local-issue in-process queue — serializes createOutboundIssue callers. */ const outboundQueues = new Map>(); @@ -148,10 +149,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; } @@ -216,8 +221,83 @@ export async function syncFromLinear( } /** - * 로컬 → Linear: 로컬 이슈를 Linear에 생성 (durable claim + per-issue serialize) + * 로컬 → Linear: 로컬 이슈를 Linear에 생성 (durable claim + per-issue serialize). + * Linear create와 로컬 mapping persist를 분리해, mapping 실패 시 linearId로 + * 재연결/재시도하고 동일 프로세스 재호출에서 중복 create를 막는다. */ +const pendingLinearMappings = new Map(); + +const MAPPING_PERSIST_ATTEMPTS = 3; + +/** @internal Test-only: install a fake client without loading the SDK. */ +export function __setLinearBridgeClientForTests(client: unknown, teamId = 'team-test'): void { + linearClient = client; + linearTeamId = teamId; + linearInitPromise = Promise.resolve(); +} + +/** @internal Test-only: drop in-process pending mapping recovery state. */ +export function __clearPendingLinearMappingsForTests(): void { + pendingLinearMappings.clear(); +} + +function persistLinearMapping( + store: SqliteIssueStore, + issueId: string, + mapping: { linearId: string; linearIdentifier: string; linearUrl: string }, +): void { + store.updateIssue(issueId, { + linearId: mapping.linearId, + linearIdentifier: mapping.linearIdentifier, + linearUrl: mapping.linearUrl, + }); + store.addEvent(issueId, 'linked', { + content: `Linear에 생성: ${mapping.linearIdentifier}`, + newValue: mapping.linearIdentifier, + idempotencyKey: `linear-linked:${mapping.linearId}`, + }); +} + +function persistLinearMappingWithRetry( + store: SqliteIssueStore, + issueId: string, + mapping: { linearId: string; linearIdentifier: string; linearUrl: string }, +): boolean { + let lastErr: unknown; + for (let attempt = 1; attempt <= MAPPING_PERSIST_ATTEMPTS; attempt++) { + try { + persistLinearMapping(store, issueId, mapping); + pendingLinearMappings.delete(issueId); + return true; + } catch (err) { + lastErr = err; + console.warn( + `[LinearBridge] 로컬 mapping persist 실패 (${attempt}/${MAPPING_PERSIST_ATTEMPTS}):`, + err, + ); + } + } + // Best-effort reconnect: updateIssue alone may succeed even if addEvent failed. + try { + store.updateIssue(issueId, { + linearId: mapping.linearId, + linearIdentifier: mapping.linearIdentifier, + linearUrl: mapping.linearUrl, + }); + pendingLinearMappings.delete(issueId); + console.warn('[LinearBridge] mapping recovered via updateIssue-only path'); + return true; + } catch (err) { + lastErr = err; + } + console.error('[LinearBridge] 로컬 mapping persist 복구 실패:', lastErr); + return false; +} + export async function pushToLinear( store: SqliteIssueStore, issueId: string, @@ -245,6 +325,18 @@ export async function createOutboundIssue( outboundQueues.set(issueId, prev.then(() => held, () => held)); await prev; + // In-process recovery: a prior create succeeded but local mapping failed. + const pending = pendingLinearMappings.get(issueId); + if (pending) { + if (persistLinearMappingWithRetry(store, issueId, pending)) { + console.log(`[LinearBridge] 이슈 ${issueId} → Linear ${pending.linearIdentifier} (recovered)`); + return pending.linearId; + } + // Still unrecovered — return known linearId to avoid a duplicate create. + return pending.linearId; + } + + let mapping: { linearId: string; linearIdentifier: string; linearUrl: string }; try { return await withOutboundClaim(issueId, async () => { await waitForLinearBridgeInit(); @@ -271,23 +363,29 @@ export async function createOutboundIssue( const linearIssue = await created.issue; if (!linearIssue) return null; - store.updateIssue(issueId, { + mapping = { linearId: linearIssue.id, linearIdentifier: linearIssue.identifier, linearUrl: linearIssue.url, - }); - - store.addEvent(issueId, 'linked', { - content: `Linear에 생성: ${linearIssue.identifier}`, - newValue: linearIssue.identifier, - }); - - console.log(`[LinearBridge] 이슈 ${issueId} → Linear ${linearIssue.identifier}`); - return linearIssue.id; + }; } catch (err) { console.error('[LinearBridge] Linear 생성 실패:', err); return null; } + + // Remember the external id before local persist so a crash or a failed + // write cannot orphan the Linear issue behind a duplicate recreate. + pendingLinearMappings.set(issueId, mapping); + if (persistLinearMappingWithRetry(store, issueId, mapping)) { + console.log(`[LinearBridge] 이슈 ${issueId} → Linear ${mapping.linearIdentifier}`); + } else { + // External issue exists; return its id so callers do not treat this as + // "not created" and retry into a duplicate. + console.error( + `[LinearBridge] Linear ${mapping.linearIdentifier} 생성됨 but local mapping incomplete for ${issueId}`, + ); + } + return mapping.linearId; }); } finally { unlock(); diff --git a/src/issues/memoryBridge.ts b/src/issues/memoryBridge.ts index d9cc59e9..a5e38759 100644 --- a/src/issues/memoryBridge.ts +++ b/src/issues/memoryBridge.ts @@ -144,12 +144,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/issues/sqliteStore.test.ts b/src/issues/sqliteStore.test.ts index f0912132..76557ff3 100644 --- a/src/issues/sqliteStore.test.ts +++ b/src/issues/sqliteStore.test.ts @@ -25,6 +25,16 @@ describe('SqliteIssueStore durable semantics', () => { store.close(); }); + it('returns the existing row when createIssue is called again with the same id', () => { + const store = new SqliteIssueStore(path()); + const first = store.createIssue({ id: 'stable-1', projectId: 'p', title: 'first' }); + const second = store.createIssue({ id: 'stable-1', projectId: 'p', title: 'ignored duplicate' }); + expect(second.id).toBe(first.id); + expect(second.title).toBe('first'); + expect(store.listIssues().total).toBe(1); + store.close(); + }); + it('emits memory_linked only for a newly inserted link', () => { const store = new SqliteIssueStore(path()); const issue = store.createIssue({ projectId: 'p', title: 'link' }); diff --git a/src/issues/sqliteStore.ts b/src/issues/sqliteStore.ts index 0ac8c00f..11ab9644 100644 --- a/src/issues/sqliteStore.ts +++ b/src/issues/sqliteStore.ts @@ -326,6 +326,13 @@ export class SqliteIssueStore implements IIssueStore { // ============ 이슈 CRUD ============ createIssue(input: CreateIssueInput): Issue { + // Honor the documented idempotent-ID contract: a caller-supplied stable id + // returns the existing row instead of colliding on UNIQUE(id). + if (input.id) { + const existing = this.getIssue(input.id); + if (existing) return existing; + } + const id = input.id ?? nanoid(12); const now = new Date().toISOString(); @@ -438,6 +445,12 @@ export class SqliteIssueStore implements IIssueStore { return transaction() as string; } catch (error) { const message = error instanceof Error ? error.message : String(error); + // Concurrent create with the same caller-supplied id: return the winner's + // row instead of surfacing a UNIQUE(id) violation to the caller. + if (input.id) { + const existing = this.getIssue(input.id); + if (existing) return existing.id; + } // Cross-process inbound sync can both pass the pre-insert SELECT and then // collide on the unique Linear indexes — reclaim the winner's row. if (/UNIQUE/i.test(message) && (input.linearId || input.linearIdentifier)) { @@ -562,7 +575,12 @@ export class SqliteIssueStore implements IIssueStore { } if (patch.status !== undefined) { - this.applyStatusChange(id, existing.status, patch.status, 'system'); + // Re-read inside the write txn so event oldValue matches effective DB state. + const current = this.db.prepare('SELECT status FROM issues WHERE id = ?').get(id) as + | { status: IssueStatus } + | undefined; + if (!current) return; + this.applyStatusChange(id, current.status, patch.status, 'system'); } }); @@ -654,11 +672,18 @@ export class SqliteIssueStore implements IIssueStore { // ============ 상태 전이 ============ changeStatus(id: string, status: IssueStatus, actor?: string): Issue | null { - const existing = this.getIssue(id); - if (!existing) return null; - - this.applyStatusChange(id, existing.status, status, actor ?? 'system'); - return this.getIssue(id); + const run = this.db.transaction(() => { + // Read effective status inside the write transaction so concurrent + // transitions cannot stamp a stale oldValue onto the event log. + const row = this.db.prepare('SELECT status FROM issues WHERE id = ?').get(id) as + | { status: IssueStatus } + | undefined; + if (!row) return null; + + this.applyStatusChange(id, row.status, status, actor ?? 'system'); + return this.getIssue(id); + }); + return run(); } private applyStatusChange(id: string, oldStatus: IssueStatus, status: IssueStatus, actor: string): void { diff --git a/src/linear/index.ts b/src/linear/index.ts index 433d999e..61360740 100644 --- a/src/linear/index.ts +++ b/src/linear/index.ts @@ -1,2 +1,2 @@ export * from './linear.js'; -export { updateProjectAfterTask, postStatusUpdate, setLinearClient } from './projectUpdater.js'; +export { updateProjectAfterTask, postStatusUpdate, setLinearClient, fetchProjectOverviewIssues } from './projectUpdater.js'; diff --git a/src/linear/projectUpdater.boundedDesc.test.ts b/src/linear/projectUpdater.boundedDesc.test.ts new file mode 100644 index 00000000..ca193887 --- /dev/null +++ b/src/linear/projectUpdater.boundedDesc.test.ts @@ -0,0 +1,24 @@ +import { describe, expect, it } from 'vitest'; +import { buildBoundedProjectDescription } from './projectUpdater.js'; + +describe('buildBoundedProjectDescription', () => { + it('keeps the compact automation summary when the base description is long', () => { + const base = 'A'.repeat(400); + const desc = buildBoundedProjectDescription(base, { done: 3, inProgress: 2, todo: 7 }); + + expect(desc.length).toBeLessThanOrEqual(255); + expect(desc.endsWith('[Done:3 InProgress:2 Todo:7]')).toBe(true); + expect(desc).toContain('...'); + }); + + it('fits short base text and summary without truncation', () => { + const desc = buildBoundedProjectDescription('Ship it', { done: 1, inProgress: 0, todo: 0 }); + expect(desc).toBe('Ship it\n\n[Done:1 InProgress:0 Todo:0]'); + expect(desc.length).toBeLessThanOrEqual(255); + }); + + it('returns only the summary when base text is empty', () => { + const desc = buildBoundedProjectDescription('', { done: 0, inProgress: 1, todo: 2 }); + expect(desc).toBe('[Done:0 InProgress:1 Todo:2]'); + }); +}); diff --git a/src/linear/projectUpdater.pagination.test.ts b/src/linear/projectUpdater.pagination.test.ts new file mode 100644 index 00000000..1956c9b1 --- /dev/null +++ b/src/linear/projectUpdater.pagination.test.ts @@ -0,0 +1,107 @@ +import { describe, expect, it } from 'vitest'; +import { LinearClient } from '@linear/sdk'; +import { fetchProjectOverviewIssues } from './projectUpdater.js'; + +describe('fetchProjectOverviewIssues pagination', () => { + it('collects every page until hasNextPage is false', async () => { + let page = 0; + const linear = { + client: { + rawRequest: async () => { + const current = page++; + return { + data: { + project: { + issues: { + nodes: [{ priority: current + 1, state: { name: `S${current}` } }], + pageInfo: { + hasNextPage: current === 0, + endCursor: current === 0 ? 'cursor-1' : null, + }, + }, + }, + }, + }; + }, + }, + } as unknown as LinearClient; + + const nodes = await fetchProjectOverviewIssues(linear, 'proj-1'); + expect(nodes.map((n) => n.state?.name)).toEqual(['S0', 'S1']); + }); + + it('rejects a missing endCursor while more pages are claimed', async () => { + const linear = { + client: { + rawRequest: async () => ({ + data: { + project: { + issues: { + nodes: [{ priority: 1, state: { name: 'Todo' } }], + pageInfo: { hasNextPage: true, endCursor: null }, + }, + }, + }, + }), + }, + } as unknown as LinearClient; + + await expect(fetchProjectOverviewIssues(linear, 'proj-1')).rejects.toThrow( + /missing or repeated cursor/, + ); + }); + + it('rejects a repeated endCursor that cannot progress', async () => { + const linear = { + client: { + rawRequest: async () => ({ + data: { + project: { + issues: { + nodes: [{ priority: 2, state: { name: 'Todo' } }], + pageInfo: { hasNextPage: true, endCursor: 'same-cursor' }, + }, + }, + }, + }), + }, + } as unknown as LinearClient; + + // First page sets after=same-cursor; second page returns the same cursor again. + await expect(fetchProjectOverviewIssues(linear, 'proj-1')).rejects.toThrow( + /missing or repeated cursor/, + ); + }); + + it('reports explicit truncation instead of silently returning a partial set', async () => { + let page = 0; + const linear = { + client: { + rawRequest: async () => ({ + data: { + project: { + issues: { + nodes: [{ priority: 1, state: { name: 'Todo' } }], + pageInfo: { hasNextPage: true, endCursor: `cursor-${page++}` }, + }, + }, + }, + }), + }, + } as unknown as LinearClient; + + await expect(fetchProjectOverviewIssues(linear, 'proj-1')).rejects.toThrow(/safety cap/); + }); + + it('rejects a null issues connection', async () => { + const linear = { + client: { + rawRequest: async () => ({ data: { project: { issues: null } } }), + }, + } as unknown as LinearClient; + + await expect(fetchProjectOverviewIssues(linear, 'proj-1')).rejects.toThrow( + /no issues connection/, + ); + }); +}); diff --git a/src/linear/projectUpdater.ts b/src/linear/projectUpdater.ts index 10fd4ca5..8c605cd9 100644 --- a/src/linear/projectUpdater.ts +++ b/src/linear/projectUpdater.ts @@ -383,6 +383,32 @@ export async function postStatusUpdate( const AUTOMATION_SECTION_MARKER = '## Automation Status'; const PROJECT_OVERVIEW_PAGE_SIZE = 100; const PROJECT_OVERVIEW_MAX_PAGES = 10; +const PROJECT_DESCRIPTION_LIMIT = 255; + +/** + * Build a Linear project description that always retains the compact automation + * summary within the 255-character hard limit (truncates the base text first). + */ +export function buildBoundedProjectDescription( + baseDesc: string, + counts: { done: number; inProgress: number; todo: number }, +): string { + const bracketedSummary = `[Done:${counts.done} InProgress:${counts.inProgress} Todo:${counts.todo}]`; + const sep = baseDesc ? '\n\n' : ''; + const baseBudget = PROJECT_DESCRIPTION_LIMIT - bracketedSummary.length - sep.length; + let clippedBase = ''; + if (baseBudget > 0 && baseDesc) { + clippedBase = baseDesc.length > baseBudget + ? `${baseDesc.slice(0, Math.max(0, baseBudget - 3))}...` + : baseDesc; + } + const finalDesc = clippedBase + ? `${clippedBase}${sep}${bracketedSummary}` + : bracketedSummary.slice(0, PROJECT_DESCRIPTION_LIMIT); + return finalDesc.length > PROJECT_DESCRIPTION_LIMIT + ? finalDesc.slice(0, PROJECT_DESCRIPTION_LIMIT) + : finalDesc; +} interface ProjectOverviewIssueNode { priority: number; @@ -402,7 +428,7 @@ const PROJECT_OVERVIEW_ISSUES_QUERY = ` } }`; -async function fetchProjectOverviewIssues( +export async function fetchProjectOverviewIssues( linear: LinearClient, projectId: string, ): Promise { @@ -411,7 +437,8 @@ async function fetchProjectOverviewIssues( }).client; const issueNodes: ProjectOverviewIssueNode[] = []; let after: string | undefined; - let complete = false; + let hasNextPage = false; + const safetyCap = PROJECT_OVERVIEW_MAX_PAGES * PROJECT_OVERVIEW_PAGE_SIZE; for (let page = 0; page < PROJECT_OVERVIEW_MAX_PAGES; page++) { const res = await withRateLimit('linear', () => @@ -429,19 +456,23 @@ async function fetchProjectOverviewIssues( }), ); const issues = res.data.project?.issues; - if (!issues) break; + if (!issues) { + throw new Error('Project overview pagination returned no issues connection'); + } issueNodes.push(...issues.nodes); - if (!issues.pageInfo.hasNextPage) { - complete = true; - break; + hasNextPage = issues.pageInfo.hasNextPage === true; + if (!hasNextPage) break; + + const endCursor = issues.pageInfo.endCursor ?? undefined; + if (!endCursor || endCursor === after) { + throw new Error('Project overview pagination returned a missing or repeated cursor'); } - after = issues.pageInfo.endCursor ?? undefined; - if (!after) break; + after = endCursor; } - if (!complete && issueNodes.length >= PROJECT_OVERVIEW_MAX_PAGES * PROJECT_OVERVIEW_PAGE_SIZE) { - throw new Error(`Project overview exceeds the ${PROJECT_OVERVIEW_MAX_PAGES * PROJECT_OVERVIEW_PAGE_SIZE}-issue safety cap`); + if (hasNextPage) { + throw new Error(`Project overview exceeds the ${safetyCap}-issue safety cap`); } return issueNodes; @@ -484,19 +515,16 @@ async function refreshProjectOverview(projectId: string, projectPath?: string): // Strip any previously-appended compact summary so it isn't doubled on each call. const baseDesc = stripped.replace(/\s*\[Done:\d+ InProgress:\d+ Todo:\d+\]$/, '').trimEnd(); - // Build a compact summary line for description (fits within 255 chars) + // Build a compact summary line for description (fits within 255 chars). + // Reserve capacity for the summary first so truncation never chops it off. const doneCount = stateCounts.get('Done') ?? 0; const inProgressCount = stateCounts.get('In Progress') ?? 0; const todoCount = stateCounts.get('Todo') ?? 0; - const compactSummary = `Done:${doneCount} InProgress:${inProgressCount} Todo:${todoCount}`; - const descWithSummary = baseDesc - ? `${baseDesc}\n\n[${compactSummary}]` - : compactSummary; - - // Truncate to 255 chars (Linear hard limit) - const finalDesc = descWithSummary.length > 255 - ? descWithSummary.slice(0, 252) + '...' - : descWithSummary; + const finalDesc = buildBoundedProjectDescription(baseDesc, { + done: doneCount, + inProgress: inProgressCount, + todo: todoCount, + }); await linear.updateProject(projectId, { description: finalDesc }); console.log(`[ProjectUpdater] Project overview updated for "${project.name}"`); diff --git a/src/orchestration/workflow.coverage.test.ts b/src/orchestration/workflow.coverage.test.ts index f2a78ec1..7e534bfb 100644 --- a/src/orchestration/workflow.coverage.test.ts +++ b/src/orchestration/workflow.coverage.test.ts @@ -218,9 +218,49 @@ describe('workflow storage round trips', () => { await saveExecution(execution); const loaded = await loadExecution(executionId); + expect(execution.definitionStamp).toBe('missing'); expect(loaded).toEqual(execution); }); + it('refuses to persist an execution after the workflow definition is replaced', async () => { + const workflowId = uniqueId('cov-wf-fence'); + const executionId = uniqueId('cov-exec-fence'); + cleanupWorkflowIds.push(workflowId); + cleanupExecutionIds.push(executionId); + + await saveWorkflow({ + id: workflowId, + name: 'Fence me', + projectPath: '/tmp/project', + steps: [{ id: 'step', name: 'Step', prompt: 'run' }], + }); + + const execution: WorkflowExecution = { + workflowId, + executionId, + status: 'running', + startedAt: Date.now(), + stepResults: {}, + }; + await saveExecution(execution); + expect(execution.definitionStamp).toMatch(/^\d+(\.\d+)?:\d+$/); + + // Replace the definition so mtime/size change under the live execution. + await saveWorkflow({ + id: workflowId, + name: 'Fence me — replaced', + projectPath: '/tmp/project', + steps: [ + { id: 'step', name: 'Step', prompt: 'run' }, + { id: 'extra', name: 'Extra', prompt: 'also run' }, + ], + }); + + await expect( + saveExecution({ ...execution, status: 'completed', completedAt: Date.now() }), + ).rejects.toThrow(/Workflow definition changed/); + }); + it('returns null when loading a well-formed but nonexistent execution ID', async () => { const loaded = await loadExecution(uniqueId('cov-execution-missing')); expect(loaded).toBeNull(); diff --git a/src/orchestration/workflow.test.ts b/src/orchestration/workflow.test.ts index f0730f32..f7b5bbf9 100644 --- a/src/orchestration/workflow.test.ts +++ b/src/orchestration/workflow.test.ts @@ -7,6 +7,7 @@ import { loadWorkflow, saveExecution, saveWorkflow, + validateExecution, WorkflowConfigSchema, WorkflowExecutionSchema, type WorkflowConfig, @@ -90,3 +91,45 @@ describe('workflow schema validation', () => { await expect(loadExecution(id)).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 9a71a781..36f7bf4b 100644 --- a/src/orchestration/workflow.ts +++ b/src/orchestration/workflow.ts @@ -113,6 +113,12 @@ export interface WorkflowExecution { completedAt?: number; stepResults: Record; checkpoint?: string; // git commit hash for rollback + /** + * Fence against concurrent workflow-definition replacement. + * Format matches filesystem identity: `${mtimeMs}:${size}` or `missing`. + * Captured on first persist; later saves refuse if the definition file changed. + */ + definitionStamp?: string; } /** @@ -174,6 +180,8 @@ export const WorkflowExecutionSchema = z.object({ completedAt: z.number().optional(), stepResults: z.record(z.string(), StepResultSchema), checkpoint: z.string().optional(), + /** Fence against concurrent workflow-definition replacement (see saveExecution). */ + definitionStamp: z.string().optional(), }); // DAG Utilities @@ -322,6 +330,56 @@ 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 */ @@ -341,6 +399,22 @@ export async function saveWorkflow(workflow: WorkflowConfig): Promise { console.log(`[Workflow] Saved: ${parsed.data.name} (${parsed.data.id})`); } +/** + * Loader-compatible stamp for a workflow definition file (`mtimeMs:size` or + * `missing`). Same convention as storeFileStamp/atomicWriteFile consumers. + */ +export async function workflowDefinitionStamp(workflowId: string): Promise { + try { + const filePath = storageFilePath(WORKFLOW_DIR, workflowId, '.yaml'); + const st = await fs.stat(filePath); + return `${st.mtimeMs}:${st.size}`; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return 'missing'; + // Invalid IDs throw from storageFilePath; propagate those. + throw error; + } +} + /** * Load workflow */ @@ -389,8 +463,21 @@ export async function listWorkflows(): Promise { /** * Save execution state — validated against the workflow definition when present. + * Refuses to persist when the linked workflow definition was replaced under this + * execution (definitionStamp fence), so a stale snapshot cannot be reported as a + * success for a definition it never ran. */ export async function saveExecution(execution: WorkflowExecution): Promise { + const currentStamp = await workflowDefinitionStamp(execution.workflowId); + if (execution.definitionStamp !== undefined && execution.definitionStamp !== currentStamp) { + throw new Error( + `Workflow definition changed under execution ${execution.executionId} ` + + `(expected ${execution.definitionStamp}, found ${currentStamp})`, + ); + } + // Stamp the caller's object so subsequent in-memory saves keep the fence. + execution.definitionStamp = execution.definitionStamp ?? currentStamp; + const parsed = WorkflowExecutionSchema.safeParse(execution); if (!parsed.success) { const details = parsed.error.issues.map((issue) => `${issue.path.join('.') || 'root'}: ${issue.message}`); @@ -399,7 +486,12 @@ export async function saveExecution(execution: WorkflowExecution): Promise // Both checks, because they see different things: the schema is structural and // cannot know whether a step id exists in the workflow DEFINITION, which is // what this one reads from disk to compare against. (AGT-3457 + AGT-4288) - await assertExecutionPersistable(parsed.data); + const workflow = await loadWorkflow(parsed.data.workflowId); + await assertExecutionPersistable(parsed.data, workflow); + // Lifecycle/DAG validation: a step cannot be persisted as completed without a + // completion time, failed without an error, or advanced past a dependency that + // is still pending/running. (AGT-3489) + validateExecution(parsed.data, workflow?.steps); const filePath = storageFilePath(EXECUTION_DIR, parsed.data.executionId, '.json'); await fs.mkdir(EXECUTION_DIR, { recursive: true }); await fs.writeFile(filePath, JSON.stringify(parsed.data, null, 2), 'utf-8'); @@ -408,8 +500,12 @@ export async function saveExecution(execution: WorkflowExecution): Promise /** * Reject execution snapshots that are structurally invalid or incompatible with * their workflow definition (unknown step ids, failed definition validation). + * `definition` may be passed by a caller that already loaded it. */ -export async function assertExecutionPersistable(execution: WorkflowExecution): Promise { +export async function assertExecutionPersistable( + execution: WorkflowExecution, + definition?: WorkflowConfig | null, +): Promise { const allowedStatuses = new Set(['running', 'completed', 'failed', 'aborted']); if (!execution.workflowId) throw new Error('Execution workflowId is required'); if (!execution.executionId) throw new Error('Execution executionId is required'); @@ -420,7 +516,7 @@ export async function assertExecutionPersistable(execution: WorkflowExecution): throw new Error('Execution stepResults must be an object'); } - const workflow = await loadWorkflow(execution.workflowId); + const workflow = definition ?? await loadWorkflow(execution.workflowId); if (!workflow) { // Definition not on disk yet (common in unit tests that only exercise // execution storage IDs). Structural checks above still apply. 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.ts b/src/taskState/store.ts index 395a99c8..c0cabb15 100644 --- a/src/taskState/store.ts +++ b/src/taskState/store.ts @@ -169,15 +169,21 @@ function readStoreLockOwner(lockPath: string): StoreLockOwner | null { } } +/** 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 getStorePath(): string { return process.env.OPENSWARM_TASK_STATE_FILE || join(homedir(), '.openswarm', 'task-state.json'); } 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)) { @@ -352,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; diff --git a/src/taskState/storeFileStamp.test.ts b/src/taskState/storeFileStamp.test.ts new file mode 100644 index 00000000..bcf87cd4 --- /dev/null +++ b/src/taskState/storeFileStamp.test.ts @@ -0,0 +1,79 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { + mkdtempSync, + readFileSync, + renameSync, + rmSync, + statSync, + utimesSync, + writeFileSync, +} from 'node:fs'; +import { join } from 'node:path'; +import { tmpdir } from 'node:os'; +import { getTaskState, resetTaskStateStoreForTests, upsertTaskState } from './store.js'; + +/** + * The store cache is keyed on a file stamp. Including only `mtimeMs:size` misses + * a same-size cross-process replacement (atomic rename) that lands inside the + * same mtime tick — the cache then keeps serving the pre-replacement snapshot + * forever. These cases pin the inode component of the stamp. + */ +describe('task state store file stamp', () => { + let stateDir: string; + let stateFile: string; + + beforeEach(() => { + stateDir = mkdtempSync(join(tmpdir(), 'openswarm-store-stamp-')); + stateFile = join(stateDir, 'state.json'); + process.env.OPENSWARM_TASK_STATE_FILE = stateFile; + resetTaskStateStoreForTests(); + }); + + afterEach(() => { + resetTaskStateStoreForTests(); + rmSync(stateDir, { recursive: true, force: true }); + }); + + it('reloads after a same-size replacement that reuses the original mtime', () => { + upsertTaskState('STAMP-1', { title: 'SWAP-SRC' }); + const original = readFileSync(stateFile, 'utf8'); + + // A whole-second mtime has no sub-millisecond component, so both files can be + // given byte-identical `mtimeMs` and `size` — leaving the inode as the only + // stamp component that can distinguish them. + const sharedMtime = new Date(Math.floor(Date.now() / 1000) * 1000); + utimesSync(stateFile, sharedMtime, sharedMtime); + const before = statSync(stateFile); + + // Warm the cache on the original snapshot. + expect(getTaskState('STAMP-1')?.title).toBe('SWAP-SRC'); + + const replaced = original.replace('"title": "SWAP-SRC"', '"title": "SWAP-DST"'); + expect(replaced).not.toBe(original); + expect(Buffer.byteLength(replaced)).toBe(Buffer.byteLength(original)); + + // Cross-process replacement by atomic rename: a different inode with the + // same byte length and the same mtime. + const incoming = `${stateFile}.incoming`; + writeFileSync(incoming, replaced, 'utf8'); + utimesSync(incoming, sharedMtime, sharedMtime); + const incomingStat = statSync(incoming); + expect(incomingStat.size).toBe(before.size); + expect(incomingStat.mtimeMs).toBe(before.mtimeMs); + expect(incomingStat.ino).not.toBe(before.ino); + renameSync(incoming, stateFile); + + // The cache is NOT reset here: invalidation must come from the stamp alone. + expect(getTaskState('STAMP-1')?.title).toBe('SWAP-DST'); + }); + + it('keeps serving the cached snapshot when nothing about the file changed', () => { + upsertTaskState('STAMP-2', { title: 'stable' }); + const first = getTaskState('STAMP-2'); + expect(first?.title).toBe('stable'); + + // Same inode, same mtime, same size: the stamp matches and the cached object + // is returned rather than re-parsed. + expect(getTaskState('STAMP-2')).toBe(first); + }); +}); diff --git a/src/verify/qualityHarness.test.ts b/src/verify/qualityHarness.test.ts new file mode 100644 index 00000000..8cc7c523 --- /dev/null +++ b/src/verify/qualityHarness.test.ts @@ -0,0 +1,134 @@ +import { mkdir, mkdtemp, symlink, writeFile, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { dirname, join } from 'node:path'; +import { afterEach, describe, expect, it } from 'vitest'; +import { + runQualityHarness, + scanStaticQuality, + selectQualitySourceFiles, + type QualityCommandResult, +} from './qualityHarness.js'; + +const roots: string[] = []; + +async function fixture(files: Record): Promise { + const root = await mkdtemp(join(tmpdir(), 'openswarm-quality-harness-')); + roots.push(root); + for (const [name, content] of Object.entries(files)) { + const path = join(root, name); + await mkdir(dirname(path), { recursive: true }); + await writeFile(path, content); + } + return root; +} + +afterEach(async () => { + await Promise.all(roots.splice(0).map((root) => rm(root, { recursive: true, force: true }))); +}); + +describe('selectQualitySourceFiles', () => { + it('keeps tracked source and drops junk / non-source paths', () => { + expect(selectQualitySourceFiles([ + 'src/a.ts', + 'src/a.ts', + 'README.md', + 'node_modules/x.js', + 'dist/out.js', + 'pkg/main.py', + ])).toEqual(['pkg/main.py', 'src/a.ts']); + }); +}); + +describe('scanStaticQuality', () => { + it('scans every listed file and surfaces critical BS as error findings', async () => { + const root = await fixture({ + 'src/clean.ts': 'export const ok = 1;\n', + // The bsDetector rule matches per line, so the empty catch body must sit + // on the same line as `catch` for the detection under test to fire. + 'src/bad.ts': 'try { doWork(); } catch {}\n', + }); + const { findings, filesScanned } = await scanStaticQuality(root, ['src/clean.ts', 'src/bad.ts']); + expect(filesScanned).toBe(2); + expect(findings.some((f) => f.ruleId === 'openswarm/quality-bs/exception_hiding' && f.filePath === 'src/bad.ts')).toBe(true); + }); + + it('fails closed on oversize (truncation) rather than skipping the file', async () => { + const root = await fixture({ + 'src/huge.ts': Buffer.alloc(512 * 1024 + 8, 0x61), + }); + const { findings, filesScanned } = await scanStaticQuality(root, ['src/huge.ts']); + expect(filesScanned).toBe(0); + expect(findings).toEqual([expect.objectContaining({ + ruleId: 'openswarm/quality-truncated', + level: 'error', + filePath: 'src/huge.ts', + })]); + }); + + it('fails closed on scope escape and unreadable paths', async () => { + const root = await fixture({ 'src/ok.ts': 'export {};\n' }); + const escaped = await scanStaticQuality(root, ['../outside.ts']); + expect(escaped.findings.some((f) => f.ruleId === 'openswarm/quality-scope')).toBe(true); + + const missing = await scanStaticQuality(root, ['src/missing.ts']); + expect(missing.findings.some((f) => f.ruleId === 'openswarm/quality-read' && f.filePath === 'src/missing.ts')).toBe(true); + }); + + it('refuses symlinked source as a non-regular read failure', async () => { + const root = await fixture({ 'src/real.ts': 'export {};\n' }); + await symlink(join(root, 'src/real.ts'), join(root, 'src/link.ts')); + const { findings } = await scanStaticQuality(root, ['src/link.ts']); + expect(findings.some((f) => f.ruleId === 'openswarm/quality-read' && f.filePath === 'src/link.ts')).toBe(true); + }); +}); + +describe('runQualityHarness', () => { + it('combines static findings with isolated command results', async () => { + const root = await fixture({ + 'src/a.ts': 'export const x = 1;\n', + 'package.json': JSON.stringify({ scripts: { typecheck: 'tsc --noEmit' } }), + }); + const executeCommands = async (): Promise => ([ + { name: 'typecheck', kind: 'typecheck', status: 'fail', detail: 'error TS2304' }, + { name: 'test', kind: 'test', status: 'pass', detail: 'ok' }, + ]); + const result = await runQualityHarness(root, { + sourceFiles: ['src/a.ts'], + staticOnly: false, + verify: { enabled: true, blockOnNewFailures: true, maxCommands: 4 }, + executeCommands, + }); + expect(result.filesScanned).toBe(1); + expect(result.commands).toHaveLength(2); + expect(result.findings.some((f) => f.ruleId === 'openswarm/quality-command/typecheck')).toBe(true); + expect(result.status).toBe('failed'); + }); + + it('passes when static scan is clean and commands pass', async () => { + const root = await fixture({ 'src/a.ts': 'export const x = 1;\n' }); + const result = await runQualityHarness(root, { + sourceFiles: ['src/a.ts'], + staticOnly: true, + }); + expect(result).toMatchObject({ + status: 'passed', + filesListed: 1, + filesScanned: 1, + findings: [], + commands: [], + }); + }); + + it('records verify-plan failures as explicit harness errors', async () => { + const root = await fixture({ + 'src/a.ts': 'export {};\n', + '.openswarm/verify.yaml': 'version: 1\ncommands: []\n', + }); + const result = await runQualityHarness(root, { + sourceFiles: ['src/a.ts'], + verify: { enabled: true, blockOnNewFailures: true, maxCommands: 4 }, + }); + expect(result.status).toBe('failed'); + expect(result.findings.some((f) => f.ruleId === 'openswarm/quality-commands')).toBe(true); + }); +}); diff --git a/src/verify/qualityHarness.ts b/src/verify/qualityHarness.ts new file mode 100644 index 00000000..dc05bd2d --- /dev/null +++ b/src/verify/qualityHarness.ts @@ -0,0 +1,410 @@ +// ============================================ +// OpenSwarm - deterministic CodeQL-style quality harness +// ============================================ +// +// Full-tree, fail-closed inspection for `openswarm review --max`: +// 1. Enumerate every tracked source file (no silent skips). +// 2. Read each file with a hard byte ceiling; read / scope / truncation +// failures become explicit error findings — never "passed with gaps". +// 3. Run discover / `.openswarm/verify.yaml` quality commands inside the +// existing isolated verify sandbox. +// +// This is the M0 engine behind hygiene-style inspection (PLATFORM_ROADMAP). + +import { execFile } from 'node:child_process'; +import { constants } from 'node:fs'; +import { access, open } from 'node:fs/promises'; +import { delimiter, extname, isAbsolute, join, relative, resolve, sep } from 'node:path'; +import { promisify } from 'node:util'; + +import { resolveBaseRef } from '../support/worktreeManager.js'; +import type { VerifyConfig } from '../core/types.js'; +import { scanFileContent } from '../registry/bsDetector.js'; +import { loadTrustedVerifyPlan } from '../agents/deterministicTester.js'; +import type { VerifyCommand } from './manifest.js'; +import { runVerify } from './runner.js'; + +const execFileAsync = promisify(execFile); +const GIT_TIMEOUT_MS = 30_000; +const MAX_SOURCE_BYTES = 512 * 1024; +const OUTPUT_TAIL = 2_000; + +/** Source extensions inspected by the harness (mirrors reviewAudit coverage). */ +const SOURCE_EXTENSIONS = new Set([ + '.ts', '.tsx', '.js', '.jsx', '.mjs', '.cjs', + '.py', '.pyw', + '.rs', '.go', + '.java', '.kt', '.kts', '.scala', '.groovy', + '.c', '.cc', '.cpp', '.cxx', '.h', '.hpp', '.hxx', '.cs', + '.rb', '.php', '.swift', '.m', '.mm', + '.ex', '.exs', '.clj', '.cljs', '.ml', '.mli', '.hs', '.dart', '.lua', '.jl', '.zig', '.nim', +]); + +const SKIP_DIR_SEGMENTS = new Set([ + 'node_modules', 'dist', 'build', 'trash', '.openswarm', 'htmlcov', 'coverage', 'vendor', + 'target', '__pycache__', 'bin', 'obj', +]); + +export type QualityFindingLevel = 'error' | 'warning' | 'note'; + +export interface QualityFinding { + ruleId: string; + level: QualityFindingLevel; + message: string; + filePath?: string; + line?: number; +} + +export interface QualityCommandResult { + name: string; + kind: VerifyCommand['kind']; + status: 'pass' | 'fail' | 'infra' | 'skipped'; + detail: string; +} + +export interface QualityHarnessResult { + status: 'passed' | 'findings' | 'failed'; + filesListed: number; + filesScanned: number; + findings: QualityFinding[]; + commands: QualityCommandResult[]; + detail?: string; +} + +export type QualityCommandExecutor = ( + projectPath: string, + commands: VerifyCommand[], + packageJsonByDirectory: Record, +) => Promise; + +export interface QualityHarnessOptions { + verify?: VerifyConfig; + /** Skip isolated quality commands (static scan only). */ + staticOnly?: boolean; + /** Override the tracked-source listing (tests). */ + sourceFiles?: readonly string[]; + /** Inject isolated command execution (tests). */ + executeCommands?: QualityCommandExecutor; +} + +function shortened(value: string, limit = OUTPUT_TAIL): string { + const flat = value.replace(/\s+/g, ' ').trim(); + if (flat.length <= limit) return flat; + return `${flat.slice(0, limit - 1)}…`; +} + +function inside(root: string, candidate: string): boolean { + const rel = relative(root, candidate); + return rel === '' || (rel !== '..' && !rel.startsWith(`..${sep}`) && !isAbsolute(rel)); +} + +function languageForExtension(ext: string): string | null { + const map: Record = { + '.ts': 'typescript', '.tsx': 'typescript', '.js': 'javascript', '.jsx': 'javascript', + '.mjs': 'javascript', '.cjs': 'javascript', + '.py': 'python', '.pyw': 'python', + '.go': 'go', '.rs': 'rust', '.java': 'java', + '.c': 'c', '.h': 'c', '.cpp': 'cpp', '.cxx': 'cpp', '.cc': 'cpp', + '.hpp': 'cpp', '.hxx': 'cpp', '.cs': 'csharp', + }; + return map[ext] ?? null; +} + +export function selectQualitySourceFiles(paths: readonly string[]): string[] { + return [...new Set(paths.filter((file) => { + if (!file || file.includes('\0')) return false; + const ext = extname(file).toLowerCase(); + if (!SOURCE_EXTENSIONS.has(ext)) return false; + if (file.split(/[/\\]/).some((seg) => SKIP_DIR_SEGMENTS.has(seg))) return false; + return true; + }))].sort(); +} + +async function findGitExecutable(): Promise { + const binary = process.platform === 'win32' ? 'git.exe' : 'git'; + const candidates: string[] = []; + for (const directory of (process.env.PATH ?? '').split(delimiter)) { + if (!isAbsolute(directory)) continue; + candidates.push(join(directory, binary)); + } + const seen = new Set(); + for (const candidate of candidates) { + if (seen.has(candidate)) continue; + seen.add(candidate); + try { + await access(candidate, constants.X_OK); + return candidate; + } catch { + // keep looking + } + } + return undefined; +} + +/** + * Every tracked source path (git index). Missing coverage is an explicit failure — + * the harness must not silently audit a subset. + */ +export async function listTrackedQualitySourceFiles(projectPath: string): Promise { + const git = await findGitExecutable(); + if (!git) throw new Error('git is not available on an absolute PATH entry.'); + try { + const { stdout } = await execFileAsync(git, ['ls-files', '-z'], { + cwd: projectPath, + timeout: GIT_TIMEOUT_MS, + maxBuffer: 64 * 1024 * 1024, + windowsHide: true, + }); + return selectQualitySourceFiles(stdout.split('\u0000')); + } catch (error) { + throw new Error('Could not enumerate tracked source for the quality harness.', { cause: error }); + } +} + +async function readBoundedSource(absolutePath: string): Promise { + const handle = await open(absolutePath, constants.O_RDONLY | constants.O_NOFOLLOW); + try { + const info = await handle.stat(); + if (!info.isFile() || info.isSymbolicLink()) { + throw new Error('source must be a regular file'); + } + if (info.size > MAX_SOURCE_BYTES) { + const err = new Error(`source exceeds ${MAX_SOURCE_BYTES} bytes`); + (err as NodeJS.ErrnoException & { code?: string }).code = 'QUALITY_TRUNCATED'; + throw err; + } + const buffer = Buffer.alloc(MAX_SOURCE_BYTES + 1); + let offset = 0; + while (offset < buffer.length) { + const { bytesRead } = await handle.read(buffer, offset, buffer.length - offset, null); + if (bytesRead === 0) break; + offset += bytesRead; + } + if (offset > MAX_SOURCE_BYTES) { + const err = new Error(`source exceeds ${MAX_SOURCE_BYTES} bytes`); + (err as NodeJS.ErrnoException & { code?: string }).code = 'QUALITY_TRUNCATED'; + throw err; + } + const bytes = buffer.subarray(0, offset); + if (bytes.includes(0)) { + const err = new Error('source contains NUL bytes'); + (err as NodeJS.ErrnoException & { code?: string }).code = 'QUALITY_BINARY'; + throw err; + } + return bytes.toString('utf8'); + } finally { + await handle.close(); + } +} + +/** + * Static pass over every listed path. A path that cannot be fully read yields an + * error finding so the harness cannot report a clean scan with holes. + */ +export async function scanStaticQuality( + projectPath: string, + sourceFiles: readonly string[], +): Promise<{ findings: QualityFinding[]; filesScanned: number }> { + const root = resolve(projectPath); + const findings: QualityFinding[] = []; + let filesScanned = 0; + + for (const file of sourceFiles) { + if (!file || file.includes('\0') || isAbsolute(file) || /(^|\/)\.\.(\/|$)/.test(file)) { + findings.push({ + ruleId: 'openswarm/quality-scope', + level: 'error', + message: `Source path escapes or is invalid for the quality harness: ${file || ''}`, + filePath: file || undefined, + }); + continue; + } + const absolute = resolve(root, file); + if (!inside(root, absolute)) { + findings.push({ + ruleId: 'openswarm/quality-scope', + level: 'error', + message: `Source path escapes repository root: ${file}`, + filePath: file, + }); + continue; + } + + try { + const content = await readBoundedSource(absolute); + filesScanned += 1; + const language = languageForExtension(extname(file).toLowerCase()); + if (!language) continue; + for (const issue of scanFileContent(content, file, language)) { + if (issue.severity === 'minor') continue; + findings.push({ + ruleId: `openswarm/quality-bs/${issue.category}`, + level: issue.severity === 'critical' ? 'error' : 'warning', + message: issue.message, + filePath: file, + line: issue.line > 0 ? issue.line : undefined, + }); + } + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + const message = error instanceof Error ? error.message : String(error); + if (code === 'QUALITY_TRUNCATED') { + findings.push({ + ruleId: 'openswarm/quality-truncated', + level: 'error', + message: `Read truncated — file exceeds the ${MAX_SOURCE_BYTES}-byte quality harness ceiling.`, + filePath: file, + }); + } else if (code === 'QUALITY_BINARY') { + findings.push({ + ruleId: 'openswarm/quality-binary', + level: 'error', + message: 'Source read failed — file contains NUL bytes and cannot be scanned as text.', + filePath: file, + }); + } else { + findings.push({ + ruleId: 'openswarm/quality-read', + level: 'error', + message: `Source read failed — scan coverage is incomplete: ${shortened(message, 240)}`, + filePath: file, + }); + } + } + } + + if (sourceFiles.length > 0 && filesScanned === 0 && findings.every((f) => f.ruleId !== 'openswarm/quality-read')) { + // Enumeration produced paths but none were readable without an explicit finding — + // treat as coverage failure so we never claim a vacuous pass. + const hasCoverageFinding = findings.some((f) => + f.ruleId === 'openswarm/quality-truncated' + || f.ruleId === 'openswarm/quality-binary' + || f.ruleId === 'openswarm/quality-scope' + || f.ruleId === 'openswarm/quality-read'); + if (!hasCoverageFinding) { + findings.push({ + ruleId: 'openswarm/quality-coverage', + level: 'error', + message: `Listed ${sourceFiles.length} source file(s) but scanned none.`, + }); + } + } + + return { findings, filesScanned }; +} + +export async function defaultExecuteQualityCommands( + projectPath: string, + commands: VerifyCommand[], + packageJsonByDirectory: Record, +): Promise { + if (commands.length === 0) return []; + const base = await resolveBaseRef(projectPath).catch((error) => { + throw new Error(`quality-harness: failed to resolve base ref: ${error instanceof Error ? error.message : String(error)}`); + }); + const evidence = await runVerify({ + projectPath, + commands, + baseRef: base.ref, + trustedPackageJsonByDirectory: packageJsonByDirectory, + }); + return evidence.map((item) => { + let status: QualityCommandResult['status']; + if (item.securityFailure) status = 'fail'; + else if (item.headStatus === 'pass') status = 'pass'; + else if (item.headStatus === 'infra') status = 'infra'; + else status = 'fail'; + return { + name: item.command.name, + kind: item.command.kind, + status, + detail: shortened(item.rawOutputTail), + }; + }); +} + +function commandFindings(commands: readonly QualityCommandResult[]): QualityFinding[] { + const findings: QualityFinding[] = []; + for (const command of commands) { + if (command.status === 'pass' || command.status === 'skipped') continue; + findings.push({ + ruleId: `openswarm/quality-command/${command.kind}`, + level: 'error', + message: command.status === 'infra' + ? `Quality command "${command.name}" hit an infrastructure failure: ${command.detail || 'no detail'}` + : `Quality command "${command.name}" failed in isolation: ${command.detail || 'non-zero exit'}`, + }); + } + return findings; +} + +function harnessStatus(findings: readonly QualityFinding[]): QualityHarnessResult['status'] { + if (findings.some((f) => f.level === 'error')) return 'failed'; + if (findings.length > 0) return 'findings'; + return 'passed'; +} + +/** + * Deterministic quality harness: complete static coverage of tracked source, + * then isolated typecheck/lint/test/build commands from verify discovery. + */ +export async function runQualityHarness( + projectPath: string, + options: QualityHarnessOptions = {}, +): Promise { + let filesListed = 0; + let sourceFiles: string[]; + try { + sourceFiles = options.sourceFiles + ? selectQualitySourceFiles(options.sourceFiles) + : await listTrackedQualitySourceFiles(projectPath); + filesListed = sourceFiles.length; + } catch (error) { + const cause = shortened(error instanceof Error ? error.message : String(error)); + return { + status: 'failed', + filesListed: 0, + filesScanned: 0, + findings: [{ + ruleId: 'openswarm/quality-enumerate', + level: 'error', + message: `Could not list tracked source for the quality harness: ${cause}`, + }], + commands: [], + detail: cause, + }; + } + + const staticScan = await scanStaticQuality(projectPath, sourceFiles); + const findings = [...staticScan.findings]; + const commands: QualityCommandResult[] = []; + + const verify = options.verify ?? { enabled: true, blockOnNewFailures: true, maxCommands: 4 }; + if (!options.staticOnly && verify.enabled) { + try { + const plan = await loadTrustedVerifyPlan(projectPath, verify); + const execute = options.executeCommands ?? defaultExecuteQualityCommands; + const results = await execute(projectPath, plan.commands, plan.packageJsonByDirectory); + commands.push(...results); + findings.push(...commandFindings(results)); + } catch (error) { + const cause = shortened(error instanceof Error ? error.message : String(error)); + findings.push({ + ruleId: 'openswarm/quality-commands', + level: 'error', + message: `Quality command planning/execution failed: ${cause}`, + }); + } + } + + return { + status: harnessStatus(findings), + filesListed, + filesScanned: staticScan.filesScanned, + findings, + commands, + ...(filesListed !== staticScan.filesScanned + ? { detail: `Scanned ${staticScan.filesScanned}/${filesListed} tracked source file(s).` } + : {}), + }; +} From 0ba0e3a1491f3cc47cfad97aafa832798669761f Mon Sep 17 00:00:00 2001 From: unohee Date: Mon, 28 Sep 2026 15:26:11 +0900 Subject: [PATCH 2/3] fix(issues): reject an idempotent create whose stored row is a different artifact MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CI caught a regression from the salvaged idempotent-createIssue guard: src/automation/taskSource.test.ts:174 "rejects an idempotent child collision when a retried plan changed" (6799 passed, 1 failed). The guard returned the existing row for any caller-supplied id, which erased SqliteTaskSource's deterministic duplicate-sibling guard (AGT-2908): its createSubIssue wraps createIssue in try/catch and relies on the UNIQUE(id) collision to distinguish "the same decomposition retried" (identical title+ description → reuse the child) from "the retried plan changed" (→ report `existing artifact does not match the requested plan`). Returning the stored row unconditionally made every changed-plan retry look like a successful reuse, so the caller got the OLD child's title back instead of `{ error }`. Reconciled rather than reverted: a caller-supplied id found in the store is only treated as an idempotent retry when the identity-bearing fields agree (title, description, parentId, projectId — via sameIssueContent). Otherwise createIssue throws an explicit collision error naming the id, so: - an identical retry returns the existing row (keeps #777's intent), and - a materially different artifact under the same id is rejected and NOT overwritten (keeps AGT-2908's invariant and its test green). The same comparison is applied to the concurrent-create catch branch, which previously also returned the winner's row unconditionally. sqliteStore.test.ts's salvaged case asserted the over-permissive behavior ("second.title === 'first'" for a changed title); it now asserts the real contract instead of being deleted, plus a new case pinning the rejection. --- src/issues/sqliteStore.test.ts | 26 ++++++++++++++++---- src/issues/sqliteStore.ts | 43 +++++++++++++++++++++++++++++----- 2 files changed, 58 insertions(+), 11 deletions(-) diff --git a/src/issues/sqliteStore.test.ts b/src/issues/sqliteStore.test.ts index 76557ff3..9e900a3b 100644 --- a/src/issues/sqliteStore.test.ts +++ b/src/issues/sqliteStore.test.ts @@ -25,12 +25,28 @@ describe('SqliteIssueStore durable semantics', () => { store.close(); }); - it('returns the existing row when createIssue is called again with the same id', () => { + it('returns the existing row when an identical create is retried with the same id', () => { const store = new SqliteIssueStore(path()); - const first = store.createIssue({ id: 'stable-1', projectId: 'p', title: 'first' }); - const second = store.createIssue({ id: 'stable-1', projectId: 'p', title: 'ignored duplicate' }); - expect(second.id).toBe(first.id); - expect(second.title).toBe('first'); + const first = store.createIssue({ id: 'stable-1', projectId: 'p', title: 'first', description: 'plan' }); + const retried = store.createIssue({ id: 'stable-1', projectId: 'p', title: 'first', description: 'plan' }); + expect(retried.id).toBe(first.id); + expect(retried.title).toBe('first'); + expect(store.listIssues().total).toBe(1); + store.close(); + }); + + it('rejects a caller-supplied id whose stored row describes a different artifact', () => { + const store = new SqliteIssueStore(path()); + store.createIssue({ id: 'stable-2', projectId: 'p', title: 'first plan', description: 'plan' }); + // A re-planned create under the same id must not be silently accepted as the + // stored row (and must not overwrite it) — the caller needs to see the clash. + expect(() => store.createIssue({ + id: 'stable-2', + projectId: 'p', + title: 'changed plan', + description: 'plan', + })).toThrow(/already exists with different content/); + expect(store.getIssue('stable-2')?.title).toBe('first plan'); expect(store.listIssues().total).toBe(1); store.close(); }); diff --git a/src/issues/sqliteStore.ts b/src/issues/sqliteStore.ts index 11ab9644..5ac21bc0 100644 --- a/src/issues/sqliteStore.ts +++ b/src/issues/sqliteStore.ts @@ -115,6 +115,22 @@ function restrictDatabasePermissions(path: string): void { } } +/** + * Whether a caller-supplied `id` already in the store describes the same + * artifact the caller is asking for. + * + * Compares the identity-bearing fields only: title, description, parent and + * project. Optional metadata (priority, estimate, assignee) is deliberately not + * compared — a retry that omits or re-derives it is still the same artifact, + * while a changed title/description/parent is a different one. + */ +function sameIssueContent(existing: Issue, input: CreateIssueInput): boolean { + return existing.title === input.title + && (existing.description ?? '') === (input.description ?? '') + && (existing.parentId ?? undefined) === (input.parentId ?? undefined) + && existing.projectId === input.projectId; +} + export class SqliteIssueStore implements IIssueStore { private db: Database.Database; @@ -326,14 +342,23 @@ export class SqliteIssueStore implements IIssueStore { // ============ 이슈 CRUD ============ createIssue(input: CreateIssueInput): Issue { - // Honor the documented idempotent-ID contract: a caller-supplied stable id - // returns the existing row instead of colliding on UNIQUE(id). + const id = input.id ?? nanoid(12); + // Honor the documented idempotent-ID contract, but do not let it mask a real + // collision: a caller-supplied id that already exists is only "the same + // create retried" when the identity-bearing fields agree. A materially + // different row under the same id (a re-planned decomposition, AGT-2908) must + // be rejected — returning the stored row would silently accept the new plan, + // and overwriting it would destroy the artifact the first plan produced. if (input.id) { const existing = this.getIssue(input.id); - if (existing) return existing; + if (existing) { + if (sameIssueContent(existing, input)) return existing; + throw new Error( + `Issue ${input.id} already exists with different content: existing artifact does not match the requested create`, + ); + } } - const id = input.id ?? nanoid(12); const now = new Date().toISOString(); const insertIssue = this.db.prepare(` @@ -446,10 +471,16 @@ export class SqliteIssueStore implements IIssueStore { } catch (error) { const message = error instanceof Error ? error.message : String(error); // Concurrent create with the same caller-supplied id: return the winner's - // row instead of surfacing a UNIQUE(id) violation to the caller. + // row when it is the same artifact, otherwise report the collision rather + // than surfacing a raw UNIQUE(id) violation. if (input.id) { const existing = this.getIssue(input.id); - if (existing) return existing.id; + if (existing) { + if (sameIssueContent(existing, input)) return existing.id; + throw new Error( + `Issue ${input.id} already exists with different content: existing artifact does not match the requested create`, + ); + } } // Cross-process inbound sync can both pass the pre-insert SELECT and then // collide on the unique Linear indexes — reclaim the winner's row. From a8d42c23e669d22ecad8b9ae86b45e69006bbd17 Mon Sep 17 00:00:00 2001 From: unohee Date: Mon, 28 Sep 2026 15:29:42 +0900 Subject: [PATCH 3/3] docs(changelog): record the salvaged state-integrity and quality-harness changes --- CHANGELOG.md | 21 +++++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 567f6b09..d3919ee5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,27 @@ ## [Unreleased] +### Added + +- **`openswarm review --max --harness-only` runs a deterministic quality harness with no LLM cost.** The new harness (`src/verify/qualityHarness.ts`) enumerates every **tracked** source file via the git index and scans each one, then runs the discovered typecheck/lint/test/build commands inside the existing isolated verify sandbox. It is fail-closed by construction: a file over the 512 KiB ceiling, one containing NUL bytes, an unreadable path, a symlinked source, or a path that escapes the repository root each become an explicit error finding rather than a silent skip, and a listing that scanned nothing is itself an error — so the gate cannot report "passed" over a subset it never read. Findings are folded into the audit run as a synthetic `.openswarm/quality-harness` area, which means the markdown report and the exit-code contract carry the evidence on every `--max` run, not only when an LLM area happened to notice something. + +### Fixed + +- **A workflow execution can no longer be persisted in a state its DAG forbids.** `saveExecution` now rejects a step marked `completed` without a `completedAt`, `failed` without an `error`, or advanced past a dependency that is still pending, running or failed — states that previously reached disk and were then read back as if the pipeline had progressed. It also gains a `definitionStamp` fence: a definition replaced underneath a live execution makes the next save refuse rather than record a snapshot that never ran against it. +- **The local issue store no longer serves a stale snapshot after a same-size replacement.** The cache stamp was `mtimeMs:size`, which is unchanged when another process replaces the file by atomic rename within the same mtime tick and writes the same number of bytes — the process then kept serving the old contents indefinitely. The stamp now includes the inode. +- **A status transition's event log records the status the write actually saw.** `changeStatus` and `updateIssue` read the current status inside the write transaction instead of from a pre-transaction read, so a concurrent transition can no longer stamp a stale `oldValue` into the audit trail. +- **The daily reporter retries a failed project once, and only advances its watermark when every project succeeded (after that retry).** A transient Linear error previously left a project unreported until the next day, while a permanently failing project could hold the window open; the watermark now reflects post-retry reality. +- **A Linear project description keeps its automation summary.** The compact `[Done:n InProgress:n Todo:n]` line was appended before truncation, so a long base description pushed it past Linear's 255-character limit and it was silently cut. The builder now reserves room for the summary first. +- **Repository-registry corruption that could not be quarantined says so.** `loadRepos` reported a failed `renameSync` as if the corrupt file had been preserved, naming a recovery path that does not exist; the quarantine failure is now surfaced as its own error. +- **A failed local Linear mapping can no longer orphan-recreate the remote issue.** If `updateIssue`/`addEvent` throws after Linear's `createIssue` succeeded, the new `linearId` is remembered in-process and re-persisted on the next call (with an `idempotencyKey` on the `linked` event), instead of the retry creating a second Linear issue. +- **Backlog grooming refuses to act on an issue id it was not given.** `validIssueIds` is now mandatory at both parse time and apply time, so a hallucinated or out-of-scope id from the planner is dropped before any mutation path sees it, and an omitted scope means "nothing is in scope" rather than "anything goes". +- **Shell/SQL source under data directories still requires validation evidence.** `isValidationRelevantFile` exempts locale/fixtures/snapshots trees as non-code, but its idea of "code" excluded `sh`, `bash`, `sql`, `swift` and others that a worker can meaningfully break; those now count. +- **Linear project-overview pagination surfaces a broken page walk.** A response with no issues connection, or a missing/repeated `endCursor`, now throws instead of silently returning a truncated issue list that the overview reports as complete. +- **An idempotent `createIssue` rejects a materially different artifact under the same id.** A caller-supplied `id` that already exists is treated as a retry only when title, description, parent and project agree; otherwise the store throws an explicit collision error instead of returning the stored row (which masked a re-planned decomposition) or overwriting it (which would destroy the first plan's artifact). +- **The `dev.ts` close handler always releases its task.** Reporting ran before `onComplete` and `activeTasks.delete`, so a throw while formatting cost/output left the task registered forever; the reporting is now contained and the cleanup runs in `finally`. +- **`memoryBridge` emits one `memory_linked` event per link**, not two (`linkMemory` already emits one). +- **A failed Linear SDK load is retryable.** The rejected init promise stayed cached, so every later `initLinearBridge` call awaited the same rejection and the bridge never recovered within the process. + ## 0.24.3 — 2026-09-28 ### Added