From 01e3e04f0a4e19eca30f4f87290ece16e88b8211 Mon Sep 17 00:00:00 2001 From: Miya Date: Tue, 28 Jul 2026 20:52:54 +0200 Subject: [PATCH 1/3] fix(core): make workflow runs durable across resumes Persist runWorkflow state to the JSONL database, expose builder resume options, and retry stale running steps after interruption. --- .../src/__tests__/resume-fallback.test.ts | 27 ++++++++ .../src/__tests__/run-persistence.test.ts | 68 +++++++++++++++++++ packages/core/src/builder.ts | 8 ++- packages/core/src/run.ts | 5 ++ packages/core/src/runner.ts | 5 +- 5 files changed, 108 insertions(+), 5 deletions(-) create mode 100644 packages/core/src/__tests__/run-persistence.test.ts diff --git a/packages/core/src/__tests__/resume-fallback.test.ts b/packages/core/src/__tests__/resume-fallback.test.ts index 657e1e6..95fc039 100644 --- a/packages/core/src/__tests__/resume-fallback.test.ts +++ b/packages/core/src/__tests__/resume-fallback.test.ts @@ -262,6 +262,33 @@ describe('resume fallback to step-output cache', () => { expect(startedSteps).toContain('step-c'); }); + it('should reset stale running steps to pending and re-execute them', async () => { + const runId = 'resume-stale-running-run'; + const config = makeResumeConfig(); + + await db.insertRun(makeRunRow(runId, config, 'running')); + await db.insertStep(makeStepRow(runId, 'step-a', 'Do step A', [], 'running')); + await db.insertStep(makeStepRow(runId, 'step-b', 'Do step B', ['step-a'], 'pending')); + await db.insertStep(makeStepRow(runId, 'step-c', 'Do step C', ['step-b'], 'pending')); + + const events: Array<{ type: string; stepName?: string }> = []; + runner.on((event) => { + if ('stepName' in event) { + events.push({ type: event.type, stepName: event.stepName }); + } + }); + + const run = await runner.resume(runId); + expect(run.status, run.error).toBe('completed'); + + expect(db.updateStep).toHaveBeenCalledWith( + `${runId}-step-a`, + expect.objectContaining({ status: 'pending', error: undefined, completionReason: undefined }) + ); + const startedSteps = events.filter((event) => event.type === 'step:started').map((event) => event.stepName); + expect(startedSteps).toContain('step-a'); + }); + it('should handle empty step-output directory gracefully', async () => { const runId = 'resume-empty-cache'; const config = makeResumeConfig(); diff --git a/packages/core/src/__tests__/run-persistence.test.ts b/packages/core/src/__tests__/run-persistence.test.ts new file mode 100644 index 0000000..8f71ebf --- /dev/null +++ b/packages/core/src/__tests__/run-persistence.test.ts @@ -0,0 +1,68 @@ +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { mkdtempSync, rmSync, writeFileSync } from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { workflow } from '../builder.js'; +import { JsonFileWorkflowDb } from '../file-db.js'; +import { InMemoryWorkflowDb } from '../memory-db.js'; +import { runWorkflow } from '../run.js'; +import { WorkflowRunner } from '../runner.js'; +import type { WorkflowRunRow } from '../types.js'; + +describe('workflow run persistence', () => { + const tmpDirs: string[] = []; + + afterEach(() => { + vi.restoreAllMocks(); + for (const tmpDir of tmpDirs.splice(0)) { + rmSync(tmpDir, { recursive: true, force: true }); + } + }); + + it('constructs runWorkflow with the cwd JSONL database', async () => { + const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'run-persistence-')); + tmpDirs.push(tmpDir); + const yamlPath = path.join(tmpDir, 'relay.yaml'); + writeFileSync( + yamlPath, + [ + 'version: "1"', + 'name: run-persistence-test', + 'swarm:', + ' pattern: sequential', + 'agents: []', + 'workflows:', + ' - name: default', + ' steps: []', + ].join('\n') + ); + const dryRunSpy = vi.spyOn(WorkflowRunner.prototype, 'dryRun'); + vi.spyOn(console, 'log').mockImplementation(() => {}); + + await runWorkflow(yamlPath, { cwd: tmpDir, dryRun: true }); + + const runner = dryRunSpy.mock.instances[0] as unknown as { db: unknown }; + expect(runner.db).toBeInstanceOf(JsonFileWorkflowDb); + expect((runner.db as JsonFileWorkflowDb).getStoragePath()).toBe( + path.join(tmpDir, '.agent-relay', 'workflow-runs.jsonl') + ); + expect(runner.db).not.toBeInstanceOf(InMemoryWorkflowDb); + }); + + it('honors WorkflowRunOptions.resume before executing a new run', async () => { + const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'builder-resume-')); + tmpDirs.push(tmpDir); + const resumedRun = { id: 'resume-id', status: 'completed' } as WorkflowRunRow; + const resumeSpy = vi.spyOn(WorkflowRunner.prototype, 'resume').mockResolvedValue(resumedRun); + const executeSpy = vi.spyOn(WorkflowRunner.prototype, 'execute').mockResolvedValue(resumedRun); + + const result = await workflow('builder-resume-test') + .agent('agent-a', { cli: 'claude' }) + .step('step-a', { agent: 'agent-a', task: 'Do step A' }) + .run({ cwd: tmpDir, renderer: false, resume: 'resume-id' }); + + expect(result).toBe(resumedRun); + expect(resumeSpy).toHaveBeenCalledWith('resume-id', undefined); + expect(executeSpy).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/core/src/builder.ts b/packages/core/src/builder.ts index c560911..e907efd 100644 --- a/packages/core/src/builder.ts +++ b/packages/core/src/builder.ts @@ -152,6 +152,8 @@ export interface WorkflowRunOptions { dryRun?: boolean; /** External step executor (e.g. Daytona sandbox backend). */ executor?: RunnerStepExecutor; + /** Resume a failed run by its ID instead of starting fresh. */ + resume?: string; /** Start from a specific step, skipping all predecessors. */ startFrom?: string; /** Previous run ID whose cached outputs are used with startFrom. */ @@ -571,7 +573,7 @@ export class WorkflowBuilder { } // Auto-detect RESUME_RUN_ID env var for resuming failed runs - const resumeRunId = process.env.RESUME_RUN_ID; + const resumeRunId = options.resume ?? process.env.RESUME_RUN_ID; const startFrom = this._startFrom ?? options.startFrom ?? process.env.START_FROM; const previousRunId = this._previousRunId ?? options.previousRunId ?? process.env.PREVIOUS_RUN_ID; @@ -587,7 +589,7 @@ export class WorkflowBuilder { runner.on(renderer.onEvent); const runPromise = resumeRunId - ? runner.resume(resumeRunId, options.vars, config) + ? runner.resume(resumeRunId, options.vars) : runner.execute(config, options.workflow, options.vars, executeOptions); try { @@ -599,7 +601,7 @@ export class WorkflowBuilder { } if (resumeRunId) { - return runner.resume(resumeRunId, options.vars, config); + return runner.resume(resumeRunId, options.vars); } return runner.execute(config, options.workflow, options.vars, executeOptions); diff --git a/packages/core/src/run.ts b/packages/core/src/run.ts index b01276e..cbfdaea 100644 --- a/packages/core/src/run.ts +++ b/packages/core/src/run.ts @@ -1,5 +1,7 @@ +import path from 'node:path'; import type { RuntimeSpawnOptions } from '@agent-relay/harness-driver'; import type { DryRunReport, TrajectoryConfig, WorkflowRunRow } from './types.js'; +import { JsonFileWorkflowDb } from './file-db.js'; import { WorkflowRunner, type WorkflowEventListener } from './runner.js'; import { createDefaultEventLogger } from './default-logger.js'; import { formatDryRunReport } from './dry-run-format.js'; @@ -51,9 +53,12 @@ export async function runWorkflow( yamlPath: string, options: RunWorkflowOptions = {} ): Promise { + const dbPath = path.join(options.cwd ?? process.cwd(), '.agent-relay', 'workflow-runs.jsonl'); + const db = new JsonFileWorkflowDb(dbPath); const runner = new WorkflowRunner({ cwd: options.cwd, relay: options.relay, + db, }); const config = await runner.parseYamlFile(yamlPath); diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index ee17545..e2af8a7 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -3685,9 +3685,10 @@ export class WorkflowRunner { } } - // Reset failed steps to pending for retry + // Reset failed and stale running steps to pending for retry. No process owns + // a running step after a run is resumed. for (const [, state] of stepStates) { - if (state.row.status === 'failed') { + if (state.row.status === 'failed' || state.row.status === 'running') { state.row.status = 'pending'; state.row.error = undefined; state.row.completionReason = undefined; From e2a2e73b57d5c585c20ed740676310dbc00e1b27 Mon Sep 17 00:00:00 2001 From: AWF Verification Date: Tue, 28 Jul 2026 22:38:14 +0200 Subject: [PATCH 2/3] fix(core): restore config arg to resume() and repair regression test The initial commit dropped the third argument from builder.ts's two runner.resume() calls. That argument feeds reconstructRunFromCache(runId, config), the cached-step-output fallback used when workflow-runs.jsonl is absent -- exactly what resume-fallback.test.ts covers. Restore it. run-persistence.test.ts also could not fail: its fixture used 'steps: []', which validateWorkflow rejects, so parseYamlFile threw before any assertion ran. Give it a real deterministic step, and strengthen the builder-resume assertion to require the config argument so this regression cannot recur. --- .../core/src/__tests__/run-persistence.test.ts | 14 ++++++++++++-- packages/core/src/builder.ts | 4 ++-- 2 files changed, 14 insertions(+), 4 deletions(-) diff --git a/packages/core/src/__tests__/run-persistence.test.ts b/packages/core/src/__tests__/run-persistence.test.ts index 8f71ebf..accd565 100644 --- a/packages/core/src/__tests__/run-persistence.test.ts +++ b/packages/core/src/__tests__/run-persistence.test.ts @@ -33,7 +33,10 @@ describe('workflow run persistence', () => { 'agents: []', 'workflows:', ' - name: default', - ' steps: []', + ' steps:', + ' - name: noop', + ' type: deterministic', + ' command: "true"', ].join('\n') ); const dryRunSpy = vi.spyOn(WorkflowRunner.prototype, 'dryRun'); @@ -62,7 +65,14 @@ describe('workflow run persistence', () => { .run({ cwd: tmpDir, renderer: false, resume: 'resume-id' }); expect(result).toBe(resumedRun); - expect(resumeSpy).toHaveBeenCalledWith('resume-id', undefined); + // The third arg is the parsed config: resume() feeds it to + // reconstructRunFromCache() when workflow-runs.jsonl is absent, so dropping + // it silently disables the cached-step-output fallback. + expect(resumeSpy).toHaveBeenCalledWith( + 'resume-id', + undefined, + expect.objectContaining({ name: 'builder-resume-test' }) + ); expect(executeSpy).not.toHaveBeenCalled(); }); }); diff --git a/packages/core/src/builder.ts b/packages/core/src/builder.ts index e907efd..d82cec7 100644 --- a/packages/core/src/builder.ts +++ b/packages/core/src/builder.ts @@ -589,7 +589,7 @@ export class WorkflowBuilder { runner.on(renderer.onEvent); const runPromise = resumeRunId - ? runner.resume(resumeRunId, options.vars) + ? runner.resume(resumeRunId, options.vars, config) : runner.execute(config, options.workflow, options.vars, executeOptions); try { @@ -601,7 +601,7 @@ export class WorkflowBuilder { } if (resumeRunId) { - return runner.resume(resumeRunId, options.vars); + return runner.resume(resumeRunId, options.vars, config); } return runner.execute(config, options.workflow, options.vars, executeOptions); From e7d8a63fb298e9ae52bffa7bf1e9d357ca968c26 Mon Sep 17 00:00:00 2001 From: AWF Verification Date: Wed, 29 Jul 2026 08:14:55 +0200 Subject: [PATCH 3/3] fix(core): address review feedback on resume durability - run.ts: pass the parsed config to resume(). Without it, cache-only resume fails when workflow-runs.jsonl is unavailable, because resume() needs the config for reconstructRunFromCache(). This bug predates the PR; builder.ts was already passing it. - runner.ts: clear retryCount when resetting a step, so a step that succeeds on its first resumed attempt no longer reports stale retries. - runner.ts: gate the running-step reset behind ResumeOptions.resetRunningSteps (default false). Runs carry no lease or heartbeat, so a live owner cannot be detected; unconditionally requeueing running steps let a second resume re-execute them alongside the original process and duplicate non-idempotent side effects. The user-facing resume paths (run.ts, builder.ts) opt in, because --resume explicitly means the previous process is gone. A real ownership lease is the proper fix and is out of scope here. --- .../src/__tests__/resume-fallback.test.ts | 2 +- .../src/__tests__/run-persistence.test.ts | 5 +++- packages/core/src/builder.ts | 4 +-- packages/core/src/run.ts | 2 +- packages/core/src/runner.ts | 27 +++++++++++++++---- packages/core/src/types.ts | 14 ++++++++++ 6 files changed, 44 insertions(+), 10 deletions(-) diff --git a/packages/core/src/__tests__/resume-fallback.test.ts b/packages/core/src/__tests__/resume-fallback.test.ts index 95fc039..213288b 100644 --- a/packages/core/src/__tests__/resume-fallback.test.ts +++ b/packages/core/src/__tests__/resume-fallback.test.ts @@ -278,7 +278,7 @@ describe('resume fallback to step-output cache', () => { } }); - const run = await runner.resume(runId); + const run = await runner.resume(runId, undefined, undefined, { resetRunningSteps: true }); expect(run.status, run.error).toBe('completed'); expect(db.updateStep).toHaveBeenCalledWith( diff --git a/packages/core/src/__tests__/run-persistence.test.ts b/packages/core/src/__tests__/run-persistence.test.ts index accd565..da515f9 100644 --- a/packages/core/src/__tests__/run-persistence.test.ts +++ b/packages/core/src/__tests__/run-persistence.test.ts @@ -71,7 +71,10 @@ describe('workflow run persistence', () => { expect(resumeSpy).toHaveBeenCalledWith( 'resume-id', undefined, - expect.objectContaining({ name: 'builder-resume-test' }) + expect.objectContaining({ name: 'builder-resume-test' }), + // User-facing resume means "the previous process is gone", so the builder + // opts in to requeueing steps left running. The library default is off. + expect.objectContaining({ resetRunningSteps: true }) ); expect(executeSpy).not.toHaveBeenCalled(); }); diff --git a/packages/core/src/builder.ts b/packages/core/src/builder.ts index d82cec7..32b6c12 100644 --- a/packages/core/src/builder.ts +++ b/packages/core/src/builder.ts @@ -589,7 +589,7 @@ export class WorkflowBuilder { runner.on(renderer.onEvent); const runPromise = resumeRunId - ? runner.resume(resumeRunId, options.vars, config) + ? runner.resume(resumeRunId, options.vars, config, { resetRunningSteps: true }) : runner.execute(config, options.workflow, options.vars, executeOptions); try { @@ -601,7 +601,7 @@ export class WorkflowBuilder { } if (resumeRunId) { - return runner.resume(resumeRunId, options.vars, config); + return runner.resume(resumeRunId, options.vars, config, { resetRunningSteps: true }); } return runner.execute(config, options.workflow, options.vars, executeOptions); diff --git a/packages/core/src/run.ts b/packages/core/src/run.ts index cbfdaea..6f00805 100644 --- a/packages/core/src/run.ts +++ b/packages/core/src/run.ts @@ -88,7 +88,7 @@ export async function runWorkflow( // Resume a previous run if requested const resumeRunId = options.resume ?? process.env.RESUME_RUN_ID; if (resumeRunId) { - return runner.resume(resumeRunId, options.vars); + return runner.resume(resumeRunId, options.vars, config, { resetRunningSteps: true }); } const startFrom = options.startFrom ?? process.env.START_FROM; diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index e2af8a7..f85588e 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -116,7 +116,7 @@ import type { WorkflowStepStatus, ProcessBackend, RunnerStepExecutor, -} from './types.js'; + ResumeOptions,} from './types.js'; import { WorkflowTrajectory, type StepOutcome } from './trajectory.js'; import { runVerification, @@ -3635,7 +3635,13 @@ export class WorkflowRunner { } /** Resume a previously paused or partially completed run. */ - async resume(runId: string, vars?: VariableContext, config?: RelayYamlConfig): Promise { + async resume( + runId: string, + vars?: VariableContext, + config?: RelayYamlConfig, + options?: ResumeOptions + ): Promise { + const resetRunningSteps = options?.resetRunningSteps ?? false; // Set up abort controller early so callers can abort() even during setup this.abortController = new AbortController(); this.paused = false; @@ -3685,17 +3691,28 @@ export class WorkflowRunner { } } - // Reset failed and stale running steps to pending for retry. No process owns - // a running step after a run is resumed. + // Reset steps to pending so they are retried. + // + // `failed` is always safe to requeue. `running` is only safe when no other + // process is still executing the step: there is no lease/heartbeat on runs + // today, so we cannot detect a live owner. Requeueing blindly would let a + // second `resume` re-run steps concurrently with the original process and + // duplicate non-idempotent side effects. It is therefore opt-in via + // `resetRunningSteps`, which the user-facing resume paths set because + // `--resume` explicitly means "the previous process is gone". for (const [, state] of stepStates) { - if (state.row.status === 'failed' || state.row.status === 'running') { + const isFailed = state.row.status === 'failed'; + const isStaleRunning = state.row.status === 'running' && resetRunningSteps; + if (isFailed || isStaleRunning) { state.row.status = 'pending'; state.row.error = undefined; state.row.completionReason = undefined; + state.row.retryCount = 0; await this.db.updateStep(state.row.id, { status: 'pending', error: undefined, completionReason: undefined, + retryCount: 0, updatedAt: new Date().toISOString(), }); } diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index 0879a8e..59a5195 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -278,6 +278,20 @@ export interface PreflightCheck { description?: string; } +/** Options for {@link WorkflowRunner.resume}. */ +export interface ResumeOptions { + /** + * Requeue steps left in `running` when the run stopped. + * + * Off by default. Runs carry no lease or heartbeat, so a live owner cannot be + * detected; requeueing blindly lets a second resume re-run steps alongside the + * original process and duplicate non-idempotent side effects. The user-facing + * resume paths set this because `--resume` explicitly means the previous + * process is gone. + */ + resetRunningSteps?: boolean; +} + /** A named workflow composed of sequential or parallel steps. */ export interface WorkflowDefinition { name: string;