diff --git a/.changeset/fn-7225-parse-pause-retry.md b/.changeset/fn-7225-parse-pause-retry.md new file mode 100644 index 0000000000..af1b582edd --- /dev/null +++ b/.changeset/fn-7225-parse-pause-retry.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Auto-retry stale parse pause-resume workflow failures instead of requiring operator action. +category: fix +dev: Re-enters safe in-review parse pause-abort replays with the shared graph resume retry budget. diff --git a/packages/engine/src/__tests__/executor-paused-abort-todo-benign.test.ts b/packages/engine/src/__tests__/executor-paused-abort-todo-benign.test.ts index c726009b12..1d68a0cca2 100644 --- a/packages/engine/src/__tests__/executor-paused-abort-todo-benign.test.ts +++ b/packages/engine/src/__tests__/executor-paused-abort-todo-benign.test.ts @@ -448,6 +448,73 @@ describe("pause-abort benign requeue-to-todo (FN-6782)", () => { expect(graphSpy).not.toHaveBeenCalled(); }); + it("auto-retries a stale in-review parse pause/resume replay instead of requiring operator action", async () => { + const { store, task, executor } = makeHarness({ column: "in-review", graphResumeRetryCount: 0 }); + (executor as any).addActiveWorktree(task.id, task.worktree); + const graphSpy = vi.spyOn(executor as any, "maybeExecuteWorkflowGraph").mockResolvedValue(true); + + await invokeGraphFailure(executor, task, { + visitedNodeIds: ["parse"], + context: { "node:parse:value": "aborted" }, + }); + + const parkedFailed = store.updateTask.mock.calls.some( + (call: unknown[]) => (call[1] as { status?: string } | undefined)?.status === "failed", + ); + expect(parkedFailed).toBe(false); + expect(logText(store)).toContain("parse node pause/resume replay surfaced after task was already in-review — auto-retrying workflow graph (1/2)"); + expect(logText(store)).toContain("Auto-recovered: retrying stale in-review parse pause/resume replay"); + expect(logText(store)).not.toContain("operator action required"); + expect(store.updateTask).toHaveBeenCalledWith(task.id, { + graphResumeRetryCount: 1, + status: null, + error: null, + }, undefined); + expect((executor as any).pausedAborted.has(task.id)).toBe(false); + expect((executor as any).activeWorktrees.has(task.id)).toBe(false); + expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ + mutationType: "task:retry-stale-in-review-parse-pause-abort-replay", + metadata: expect.objectContaining({ + nodeId: "parse", + fromColumn: "in-review", + attempt: 1, + abortProvenance: "hard-cancel", + clearedStaleFailure: false, + mode: "preserved-in-review-retry-graph", + }), + })); + await flushScheduledRetry(); + expect(graphSpy).toHaveBeenCalledTimes(1); + }); + + it("clears a prior stale operator-action failure for an in-review parse replay and retries", async () => { + const { store, task, executor } = makeHarness({ + column: "in-review", + status: "failed", + error: "Workflow graph failure surfaced after paused engine abort during pause/resume in 'in-review' at node 'parse' — operator action required; retry or explicitly unpause/resume after inspecting the task", + }); + const graphSpy = vi.spyOn(executor as any, "maybeExecuteWorkflowGraph").mockResolvedValue(true); + + await invokeGraphFailure(executor, task, { + visitedNodeIds: ["parse"], + context: { "node:parse:value": "aborted" }, + }); + + expect(logText(store)).toContain("Auto-recovered: retrying stale in-review parse pause/resume replay"); + expect(store.updateTask).toHaveBeenCalledWith(task.id, { + graphResumeRetryCount: 1, + status: null, + error: null, + }, undefined); + expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ + mutationType: "task:retry-stale-in-review-parse-pause-abort-replay", + metadata: expect.objectContaining({ nodeId: "parse", clearedStaleFailure: true }), + })); + store.getTask.mockResolvedValue({ ...task, status: null, error: null, graphResumeRetryCount: 1 }); + await flushScheduledRetry(); + expect(graphSpy).toHaveBeenCalledTimes(1); + }); + it("keeps manual retry of a prior plan pause-abort park in review instead of fresh planning", async () => { const { store, task, executor } = makeHarness({ column: "in-review", diff --git a/packages/engine/src/executor.ts b/packages/engine/src/executor.ts index bb6b2c178c..dffccc39ec 100644 --- a/packages/engine/src/executor.ts +++ b/packages/engine/src/executor.ts @@ -6995,12 +6995,12 @@ export class TaskExecutor { return true; } - private isStalePauseAbortParkFailure(live: TaskDetail): boolean { + private isStalePauseAbortParkFailure(live: TaskDetail, nodeId = "plan"): boolean { return live.status === "failed" && typeof live.error === "string" && live.error.includes(PAUSE_ABORT_PARK_ERROR_MARKER) && live.error.includes("engine abort during pause/resume") - && live.error.includes("at node 'plan'"); + && live.error.includes(`at node '${nodeId}'`); } private async handleStaleInReviewPlanPauseAbortReplay( @@ -7031,7 +7031,7 @@ export class TaskExecutor { if (failureValue !== "aborted") return false; if (this.isTerminalMergeGraphFailureValue(failureValue)) return false; const cleanRow = live.status == null && live.error == null; - const staleParkedFailure = this.isStalePauseAbortParkFailure(live); + const staleParkedFailure = this.isStalePauseAbortParkFailure(live, "plan"); if (!cleanRow && !staleParkedFailure) return false; let settings: Settings; try { @@ -7075,6 +7075,112 @@ export class TaskExecutor { return true; } + private async handleStaleInReviewParsePauseAbortReplay( + live: TaskDetail, + result: WorkflowGraphTaskRunResult, + abortProvenance: "global-pause" | "merge-seam" | "hard-cancel" | "completion-finalize" | undefined, + pausedAborted: boolean, + userCanceled: boolean, + ): Promise { + /* + FNXC:WorkflowLifecycle 2026-06-29-01:18: + A stale in-review pause/resume replay at `parse` is not an operator action. Unlike `plan`, parse is a safe workflow re-entry point for review rows, so auto-retry the graph with the shared transient resume budget and suppress the parked failure notification. + */ + if (!pausedAborted) return false; + if (abortProvenance !== "hard-cancel" && abortProvenance !== "global-pause") return false; + if (userCanceled) return false; + if (live.column !== "in-review") return false; + if (live.paused || live.userPaused === true) return false; + if (live.autoMerge === false) return false; + if (live.mergeDetails?.mergeConfirmed === true) return false; + if (result.interruptedAbortKind && result.interruptedAbortKind !== WORKFLOW_NODE_ENGINE_PAUSE_ABORT_KIND) return false; + const failedNode = result.interruptedNodeId ?? result.visitedNodeIds[result.visitedNodeIds.length - 1]; + if (failedNode !== "parse") return false; + const failureValue = typeof result.context?.[`node:${failedNode}:value`] === "string" + ? result.context[`node:${failedNode}:value`] as string + : this.graphFailureValue(result); + if (failureValue !== "aborted") return false; + if (this.isTerminalMergeGraphFailureValue(failureValue)) return false; + const cleanRow = live.status == null && live.error == null; + const staleParkedFailure = this.isStalePauseAbortParkFailure(live, "parse"); + if (!cleanRow && !staleParkedFailure) return false; + const priorRetries = live.graphResumeRetryCount ?? 0; + if (priorRetries >= MAX_TRANSIENT_GRAPH_RESUME_RETRIES) return false; + let settings: Settings; + try { + settings = await this.store.getSettings(); + } catch { + return false; + } + if (settings.globalPause === true || settings.enginePaused === true) return false; + if (!allowsAutoMergeProcessing(live, settings) && !isSharedBranchGroupMemberIntegration(live)) return false; + + const nextRetries = priorRetries + 1; + this.clearPausedAborted(live.id); + this.activeWorktrees.delete(live.id); + const message = `Workflow graph parse node pause/resume replay surfaced after task was already in-review — auto-retrying workflow graph (${nextRetries}/${MAX_TRANSIENT_GRAPH_RESUME_RETRIES})`; + executorLog.log(`${live.id}: ${message}`); + await this.store.logEntry(live.id, message, undefined, this.getRunContextFor(live.id)); + await this.store.logEntry(live.id, "Auto-recovered: retrying stale in-review parse pause/resume replay — failure notification suppressed", undefined, this.getRunContextFor(live.id)); + await this.store.updateTask(live.id, { graphResumeRetryCount: nextRetries, status: null, error: null }, this.getRunContextFor(live.id)); + try { + await this.store.recordRunAuditEvent?.({ + taskId: live.id, + agentId: "executor", + runId: generateSyntheticRunId("workflow-stale-parse-retry", live.id), + domain: "database", + mutationType: "task:retry-stale-in-review-parse-pause-abort-replay", + target: live.id, + metadata: { + nodeId: failedNode, + fromColumn: live.column, + attempt: nextRetries, + maxAttempts: MAX_TRANSIENT_GRAPH_RESUME_RETRIES, + abortProvenance: abortProvenance ?? "unknown", + clearedStaleFailure: staleParkedFailure, + mode: "preserved-in-review-retry-graph", + }, + }); + } catch (error) { + executorLog.warn(`${live.id}: failed to record stale parse replay retry audit: ${error instanceof Error ? error.message : String(error)}`); + } + await this.persistTokenUsage(live.id); + + const scheduleRetry = () => { + void (async () => { + try { + const resumeTask = await this.store.getTask(live.id); + if ( + resumeTask.deletedAt + || resumeTask.paused + || resumeTask.userPaused + || resumeTask.status != null + || resumeTask.error != null + || resumeTask.column !== "in-review" + || this.activeSessions.has(live.id) + || this.activeStepExecutors.has(live.id) + || this.activeWorkflowStepSessions.has(live.id) + || this.activeWorkflowGraphAbortControllers.has(live.id) + || TaskExecutor.processWideGraphRouting.has(live.id) + ) { + executorLog.log(`${live.id}: skipping stale parse graph retry — task is no longer in a safe in-review resume state`); + return; + } + await this.maybeExecuteWorkflowGraph(resumeTask); + } catch (err) { + executorLog.error(`Failed stale parse graph retry for ${live.id}:`, err); + } + })(); + }; + if (TRANSIENT_GRAPH_RESUME_RETRY_BACKOFF_MS > 0) { + const handle = setTimeout(scheduleRetry, TRANSIENT_GRAPH_RESUME_RETRY_BACKOFF_MS); + handle.unref?.(); + } else { + setTimeout(scheduleRetry, 0).unref?.(); + } + return true; + } + private async isReentrantPausedAbortedInFlightNode( live: TaskDetail, result: WorkflowGraphTaskRunResult, @@ -7318,6 +7424,9 @@ export class TaskExecutor { await this.persistTokenUsage(task.id); return; } + if (genuinePauseAbort && await this.handleStaleInReviewParsePauseAbortReplay(live, result, abortProvenance, pausedAborted, this.userCanceledTaskIds.has(task.id))) { + return; + } if (genuinePauseAbort && await this.handleStaleInReviewPlanPauseAbortReplay(live, result, abortProvenance, pausedAborted, this.userCanceledTaskIds.has(task.id))) { return; }