fix(workflows): address lifecycle review follow-ups (#2380)
## Summary - preserve workflow IR hashes in production column-transition audit metadata - centralize active workflow-continuation states across release, runtime, and executor paths - extract and test actionable planning-continuation selection - expand Coding (Ideas) remapping/removal coverage and add required lifecycle decision records Follow-up to the review body on #2378 after that PR was merged. ## Validation - `pnpm lint` - 123 focused core/engine tests - `pnpm verify:fast` - `pnpm test:gate` (487 tests) <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit - **Bug Fixes** - Improved workflow continuation handling by centralizing “active/continuation-eligible” state selection across executor, hold/release logic, and in-process runtime. - Persisted richer task column-transition metadata (including `irHash`) to preserve workflow provenance. - Ensured planning continuations exclude paused/missing/invalid tasks and that task resolution failures surface instead of being ignored. - Corrected fresh-worktree step execution ordering to return expected `baselineSha`/`checkpointId` behavior. - **New Features** - Added and exposed `ACTIVE_WORKFLOW_WORK_ITEM_STATES` for consistent work-item “active” semantics. - Introduced a shared planning-continuation candidate selector to standardize dispatchable planning work filtering. - **Documentation** - Clarified the small coding-ideas workflow preset omits verification while preserving a continuous executable path. - **Tests** - Added coverage for planning continuation filtering and fresh-worktree ordering behavior. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
7
.changeset/workflow-review-followups.md
Normal file
7
.changeset/workflow-review-followups.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Preserve workflow lifecycle state and start execution steps only after worktree creation.
|
||||||
|
category: fix
|
||||||
|
dev: Adds shared active-state semantics, lifecycle records, and worktree-first graph step projection.
|
||||||
@@ -100,9 +100,13 @@ describe("builtin coding-ideas workflow ir", () => {
|
|||||||
expect(nodeColumn("parse")).toBe("in-progress");
|
expect(nodeColumn("parse")).toBe("in-progress");
|
||||||
expect(nodeColumn("steps")).toBe("in-progress");
|
expect(nodeColumn("steps")).toBe("in-progress");
|
||||||
expect(nodeColumn("code-review")).toBe("in-review");
|
expect(nodeColumn("code-review")).toBe("in-review");
|
||||||
|
expect(nodeColumn("code-review-remediation")).toBe("in-progress");
|
||||||
expect(nodeColumn("completion-summary")).toBe("in-review");
|
expect(nodeColumn("completion-summary")).toBe("in-review");
|
||||||
expect(nodeColumn("merge-gate")).toBe("in-review");
|
for (const node of ir.nodes.filter((candidate) => candidate.id.startsWith("merge-"))) {
|
||||||
|
expect(node.column, `${node.id} should remain in review`).toBe("in-review");
|
||||||
|
}
|
||||||
expect(ir.nodes.some((node) => node.id === "browser-verification")).toBe(false);
|
expect(ir.nodes.some((node) => node.id === "browser-verification")).toBe(false);
|
||||||
|
expect(ir.nodes.some((node) => node.id === "browser-verification-remediation")).toBe(false);
|
||||||
expect(ir.nodes.some((node) => node.id === "post-merge-verification")).toBe(false);
|
expect(ir.nodes.some((node) => node.id === "post-merge-verification")).toBe(false);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -94,9 +94,12 @@ const RAW_BUILTIN_CODING_IDEAS_WORKFLOW_IR: WorkflowIr = (() => {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// FNXC:CodingIdeasWorkflow 2026-07-21-12:20:
|
||||||
// Keep this preset intentionally small: planning + plan review in Todo,
|
// Keep this preset intentionally small: planning + plan review in Todo,
|
||||||
// implementation in In progress, then code review and merge in In review.
|
// implementation in In progress, then code review and merge in In review.
|
||||||
// Browser and post-merge verification remain available in richer workflows.
|
// Browser and post-merge verification remain available in richer workflows;
|
||||||
|
// direct success edges replace the removed nodes so the reduced preset keeps
|
||||||
|
// one continuous executable path.
|
||||||
const removedNodeIds = new Set([
|
const removedNodeIds = new Set([
|
||||||
"browser-verification",
|
"browser-verification",
|
||||||
"browser-verification-remediation",
|
"browser-verification-remediation",
|
||||||
|
|||||||
@@ -2247,3 +2247,4 @@ export { evaluateTransitionInvariants, evaluateMergeBlockerPostcondition, evalua
|
|||||||
export { StaleBinarySchemaError, assertBinaryNotOlderThanDatabase } from "./postgres/schema-applier.js";
|
export { StaleBinarySchemaError, assertBinaryNotOlderThanDatabase } from "./postgres/schema-applier.js";
|
||||||
export { promoteResearchFinding } from "./research-feature-promotion.js";
|
export { promoteResearchFinding } from "./research-feature-promotion.js";
|
||||||
export type { ResearchFeaturePromotionInput } from "./research-feature-promotion.js";
|
export type { ResearchFeaturePromotionInput } from "./research-feature-promotion.js";
|
||||||
|
export { ACTIVE_WORKFLOW_WORK_ITEM_STATES } from "./types.js";
|
||||||
|
|||||||
@@ -2609,3 +2609,4 @@ export type { LanguageFamily, DetectedContentLanguage } from "./detect-content-l
|
|||||||
export { promoteResearchFinding } from "./research-feature-promotion.js";
|
export { promoteResearchFinding } from "./research-feature-promotion.js";
|
||||||
export type { ResearchFeaturePromotionInput } from "./research-feature-promotion.js";
|
export type { ResearchFeaturePromotionInput } from "./research-feature-promotion.js";
|
||||||
export { getTotalAgentActiveMs, startPlanningSegment, finalizePlanningSegment } from "./task-timing.js";
|
export { getTotalAgentActiveMs, startPlanningSegment, finalizePlanningSegment } from "./task-timing.js";
|
||||||
|
export { ACTIVE_WORKFLOW_WORK_ITEM_STATES } from "./types.js";
|
||||||
|
|||||||
@@ -78,6 +78,7 @@ export type { ThinkingLevel, Column, ColumnId, TaskPriority };
|
|||||||
|
|
||||||
import {
|
import {
|
||||||
MERGE_REQUEST_STATES,
|
MERGE_REQUEST_STATES,
|
||||||
|
ACTIVE_WORKFLOW_WORK_ITEM_STATES,
|
||||||
WORKFLOW_WORK_ITEM_KINDS,
|
WORKFLOW_WORK_ITEM_KINDS,
|
||||||
WORKFLOW_WORK_ITEM_STATES,
|
WORKFLOW_WORK_ITEM_STATES,
|
||||||
} from "./types/merge-queue.js";
|
} from "./types/merge-queue.js";
|
||||||
@@ -101,6 +102,7 @@ import type {
|
|||||||
} from "./types/merge-queue.js";
|
} from "./types/merge-queue.js";
|
||||||
export {
|
export {
|
||||||
MERGE_REQUEST_STATES,
|
MERGE_REQUEST_STATES,
|
||||||
|
ACTIVE_WORKFLOW_WORK_ITEM_STATES,
|
||||||
WORKFLOW_WORK_ITEM_KINDS,
|
WORKFLOW_WORK_ITEM_KINDS,
|
||||||
WORKFLOW_WORK_ITEM_STATES,
|
WORKFLOW_WORK_ITEM_STATES,
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -43,6 +43,15 @@ export const WORKFLOW_WORK_ITEM_STATES = [
|
|||||||
|
|
||||||
export type WorkflowWorkItemState = (typeof WORKFLOW_WORK_ITEM_STATES)[number];
|
export type WorkflowWorkItemState = (typeof WORKFLOW_WORK_ITEM_STATES)[number];
|
||||||
|
|
||||||
|
/** FNXC:WorkflowContinuations 2026-07-21-12:30:
|
||||||
|
* States that keep a workflow work item eligible for continuation ownership. */
|
||||||
|
export const ACTIVE_WORKFLOW_WORK_ITEM_STATES: readonly WorkflowWorkItemState[] = [
|
||||||
|
"runnable",
|
||||||
|
"running",
|
||||||
|
"held",
|
||||||
|
"retrying",
|
||||||
|
];
|
||||||
|
|
||||||
export interface WorkflowWorkItem {
|
export interface WorkflowWorkItem {
|
||||||
id: string;
|
id: string;
|
||||||
runId: string;
|
runId: string;
|
||||||
|
|||||||
@@ -1553,6 +1553,7 @@ function validateColumnAgent(column: WorkflowIrColumn): void {
|
|||||||
function validateV2(ir: WorkflowIrV2): void {
|
function validateV2(ir: WorkflowIrV2): void {
|
||||||
validateColumns(ir);
|
validateColumns(ir);
|
||||||
|
|
||||||
|
// FNXC:WorkflowValidation 2026-07-21-12:20:
|
||||||
// Capacity holds must have somewhere the scheduler can actually release
|
// Capacity holds must have somewhere the scheduler can actually release
|
||||||
// them. Failing authoring here avoids durable continuations that can never
|
// them. Failing authoring here avoids durable continuations that can never
|
||||||
// become runnable.
|
// become runnable.
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import "./executor-test-helpers.js";
|
|||||||
import { getBuiltinWorkflow } from "@fusion/core";
|
import { getBuiltinWorkflow } from "@fusion/core";
|
||||||
import { TaskExecutor } from "../executor.js";
|
import { TaskExecutor } from "../executor.js";
|
||||||
import { WorkflowGraphTaskRunner } from "../workflow-graph-task-runner.js";
|
import { WorkflowGraphTaskRunner } from "../workflow-graph-task-runner.js";
|
||||||
|
import { FOREACH_ACTIVE_CONTEXT_KEY } from "../workflow-node-handlers.js";
|
||||||
import {
|
import {
|
||||||
createMockStore,
|
createMockStore,
|
||||||
mockedCreateFnAgent,
|
mockedCreateFnAgent,
|
||||||
@@ -187,6 +188,99 @@ describe("fast mode workflow/runtime invariants", () => {
|
|||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("does not project a fresh graph step or capture its baseline before the executor creates its worktree", async () => {
|
||||||
|
let liveTask = task({
|
||||||
|
steps: [{ name: "Preflight", status: "pending" }],
|
||||||
|
worktree: undefined,
|
||||||
|
branch: undefined,
|
||||||
|
baseCommitSha: undefined,
|
||||||
|
});
|
||||||
|
const store = createMockStore();
|
||||||
|
store.getTask.mockImplementation(async () => liveTask);
|
||||||
|
const executor = new TaskExecutor(store, "/tmp/project-root");
|
||||||
|
const runGraphTaskStep = vi.spyOn(executor as any, "runGraphTaskStep").mockImplementation(async () => {
|
||||||
|
expect(store.updateStep).not.toHaveBeenCalled();
|
||||||
|
liveTask = {
|
||||||
|
...liveTask,
|
||||||
|
worktree: "/tmp/project-root/.worktrees/fresh-step",
|
||||||
|
branch: "fusion/fn-6226",
|
||||||
|
baseCommitSha: "fresh-worktree-base",
|
||||||
|
steps: [{ name: "Preflight", status: "done" }],
|
||||||
|
};
|
||||||
|
return { success: true };
|
||||||
|
});
|
||||||
|
|
||||||
|
const result = await (executor as any)
|
||||||
|
.createAuthoritativeWorkflowPrimitives({ experimentalFeatures: { workflowGraphExecutor: true } })
|
||||||
|
.runTaskStep(
|
||||||
|
{
|
||||||
|
run: { taskId: liveTask.id },
|
||||||
|
node: {
|
||||||
|
node: { id: "steps#0:step-execute" },
|
||||||
|
context: {
|
||||||
|
[FOREACH_ACTIVE_CONTEXT_KEY]: {
|
||||||
|
foreachNodeId: "steps",
|
||||||
|
stepIndex: 0,
|
||||||
|
instanceId: "steps#0",
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
liveTask,
|
||||||
|
0,
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(runGraphTaskStep).toHaveBeenCalledTimes(1);
|
||||||
|
expect(store.updateStep).not.toHaveBeenCalled();
|
||||||
|
expect(result).toEqual({
|
||||||
|
outcome: "success",
|
||||||
|
baselineSha: "fresh-worktree-base",
|
||||||
|
checkpointId: undefined,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
it("applies fresh-worktree step ordering through the legacy graph seam", async () => {
|
||||||
|
let liveTask = task({
|
||||||
|
steps: [{ name: "Preflight", status: "pending" }],
|
||||||
|
worktree: undefined,
|
||||||
|
baseCommitSha: undefined,
|
||||||
|
});
|
||||||
|
const store = createMockStore();
|
||||||
|
store.getTask.mockImplementation(async () => liveTask);
|
||||||
|
const executor = new TaskExecutor(store, "/tmp/project-root");
|
||||||
|
vi.spyOn(executor as any, "runGraphTaskStep").mockImplementation(async () => {
|
||||||
|
expect(store.updateStep).not.toHaveBeenCalled();
|
||||||
|
liveTask = {
|
||||||
|
...liveTask,
|
||||||
|
worktree: "/tmp/project-root/.worktrees/fresh-step",
|
||||||
|
baseCommitSha: "fresh-worktree-base",
|
||||||
|
steps: [{ name: "Preflight", status: "done" }],
|
||||||
|
};
|
||||||
|
return { success: true };
|
||||||
|
});
|
||||||
|
const active = {
|
||||||
|
foreachNodeId: "steps",
|
||||||
|
stepIndex: 0,
|
||||||
|
instanceId: "steps#0",
|
||||||
|
};
|
||||||
|
|
||||||
|
const result = await executor.createAuthoritativeWorkflowSeams({} as any).stepExecute?.(
|
||||||
|
liveTask,
|
||||||
|
{ [FOREACH_ACTIVE_CONTEXT_KEY]: active },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(store.updateStep).not.toHaveBeenCalled();
|
||||||
|
expect(result).toMatchObject({
|
||||||
|
outcome: "success",
|
||||||
|
contextPatch: {
|
||||||
|
[FOREACH_ACTIVE_CONTEXT_KEY]: {
|
||||||
|
baselineSha: "fresh-worktree-base",
|
||||||
|
checkpointId: undefined,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
it("fast builtin:coding still parses and executes steps while disabled optional groups stay inert", async () => {
|
it("fast builtin:coding still parses and executes steps while disabled optional groups stay inert", async () => {
|
||||||
const calls: string[] = [];
|
const calls: string[] = [];
|
||||||
const prompt = "# Task\n\n## Steps\n\n### Step 1: Do the work\n- [ ] edit files";
|
const prompt = "# Task\n\n## Steps\n\n### Step 1: Do the work\n- [ ] edit files";
|
||||||
|
|||||||
@@ -0,0 +1,29 @@
|
|||||||
|
import { describe, expect, it } from "vitest";
|
||||||
|
import type { Task, WorkflowWorkItem } from "@fusion/core";
|
||||||
|
import { selectActionablePlanningContinuations } from "../runtimes/in-process-runtime.js";
|
||||||
|
|
||||||
|
function workItem(id: string, waitReason: WorkflowWorkItem["waitReason"]): WorkflowWorkItem {
|
||||||
|
return { id, waitReason } as WorkflowWorkItem;
|
||||||
|
}
|
||||||
|
|
||||||
|
function task(id: string, patch: Partial<Task> = {}): Task {
|
||||||
|
return { id, paused: false, userPaused: false, ...patch } as Task;
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("selectActionablePlanningContinuations", () => {
|
||||||
|
it("retains only planning items whose tasks are present and unpaused", () => {
|
||||||
|
const selected = selectActionablePlanningContinuations([
|
||||||
|
{ item: workItem("eligible", "planning"), task: task("T-1") },
|
||||||
|
{ item: workItem("capacity", "capacity"), task: task("T-2") },
|
||||||
|
{ item: workItem("missing", "planning"), task: undefined },
|
||||||
|
{ item: workItem("null-task", "planning"), task: null },
|
||||||
|
{ item: workItem("no-wait-reason", null), task: task("T-5") },
|
||||||
|
{ item: workItem("paused", "planning"), task: task("T-3", { paused: true }) },
|
||||||
|
{ item: workItem("user-paused", "planning"), task: task("T-4", { userPaused: true }) },
|
||||||
|
]);
|
||||||
|
|
||||||
|
expect(selected.map(({ item, task: selectedTask }) => [item.id, selectedTask.id])).toEqual([
|
||||||
|
["eligible", "T-1"],
|
||||||
|
]);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -14,7 +14,7 @@ import { existsSync, lstatSync, realpathSync } from "node:fs";
|
|||||||
import { readFile, rm, writeFile } from "node:fs/promises";
|
import { readFile, rm, writeFile } from "node:fs/promises";
|
||||||
import type { TaskStore, Task, TaskDetail, TaskTokenUsage, StepStatus, Settings, WorkflowStep, MissionStore, AsyncMissionStore, Slice, AgentState, AgentCapability, RunMutationContext, AgentHeartbeatConfig, Agent, AgentMemoryInclusionMode, ProjectSettings, MergeResult, WorkflowIrNode, WorkflowIrNodeKind, WorkflowStepResult as CoreWorkflowStepResult, ThinkingLevel } from "@fusion/core";
|
import type { TaskStore, Task, TaskDetail, TaskTokenUsage, StepStatus, Settings, WorkflowStep, MissionStore, AsyncMissionStore, Slice, AgentState, AgentCapability, RunMutationContext, AgentHeartbeatConfig, Agent, AgentMemoryInclusionMode, ProjectSettings, MergeResult, WorkflowIrNode, WorkflowIrNodeKind, WorkflowStepResult as CoreWorkflowStepResult, ThinkingLevel } from "@fusion/core";
|
||||||
import { getUnmetSchedulingDependencies } from "./scheduler.js";
|
import { getUnmetSchedulingDependencies } from "./scheduler.js";
|
||||||
import { RetryStormError, serializeRetryStormError, evaluateCompletedPromotionFailureProvenance, evaluateSkipBypassTaint, resolveWorkflowIrForTask, resolveCompleteColumn, resolveMergeOrchestrationColumn, resolveReboundTarget, resolveColumnAgentBinding, resolveEffectiveAgent, instanceNodeId, getWorkflowExtensionRegistry, getBuiltinWorkflow, parseNoOpCompletionMarker, allowsAutoMergeProcessing, resolveEffectiveAutoMerge, isLiveSharedBranchGroupMemberIntegration, resolveMaxAutoMergeRetries, resolveMaxConsecutiveToolFailureRetries, resolveConsecutiveToolFailureRetryBackoffMs, resolveConsecutiveToolFailureThreshold, resolveExecutorEscalationTarget, resolveOptionalStepRevisionBudget, resolveOptionalReviewRevisionBudget, COMPLETION_SUMMARY_NODE_ID, upsertWorkflowStepResult, AWAITING_APPROVAL_PAUSE_REASON, THINKING_LEVELS, AgentStore, resolveExecutorFallbackModel } from "@fusion/core";
|
import { RetryStormError, serializeRetryStormError, evaluateCompletedPromotionFailureProvenance, evaluateSkipBypassTaint, resolveWorkflowIrForTask, resolveCompleteColumn, resolveMergeOrchestrationColumn, resolveReboundTarget, resolveColumnAgentBinding, resolveEffectiveAgent, instanceNodeId, getWorkflowExtensionRegistry, getBuiltinWorkflow, parseNoOpCompletionMarker, allowsAutoMergeProcessing, resolveEffectiveAutoMerge, isLiveSharedBranchGroupMemberIntegration, resolveMaxAutoMergeRetries, resolveMaxConsecutiveToolFailureRetries, resolveConsecutiveToolFailureRetryBackoffMs, resolveConsecutiveToolFailureThreshold, resolveExecutorEscalationTarget, resolveOptionalStepRevisionBudget, resolveOptionalReviewRevisionBudget, COMPLETION_SUMMARY_NODE_ID, upsertWorkflowStepResult, AWAITING_APPROVAL_PAUSE_REASON, THINKING_LEVELS, ACTIVE_WORKFLOW_WORK_ITEM_STATES, AgentStore, resolveExecutorFallbackModel } from "@fusion/core";
|
||||||
import { finalizeProvenAutoMergeTask } from "./auto-merge-finalization.js";
|
import { finalizeProvenAutoMergeTask } from "./auto-merge-finalization.js";
|
||||||
import { mergeEffectiveSettings } from "./effective-settings.js";
|
import { mergeEffectiveSettings } from "./effective-settings.js";
|
||||||
import { generateFeatureVideo, type GenerateFeatureVideoOptions } from "./review-artifacts/feature-video.js";
|
import { generateFeatureVideo, type GenerateFeatureVideoOptions } from "./review-artifacts/feature-video.js";
|
||||||
@@ -191,7 +191,7 @@ import type { StuckTaskDetector, StuckTaskEvent } from "./stuck-task-detector.js
|
|||||||
import type { PluginRunner } from "./plugin-runner.js";
|
import type { PluginRunner } from "./plugin-runner.js";
|
||||||
import { isContextLimitError } from "./context-limit-detector.js";
|
import { isContextLimitError } from "./context-limit-detector.js";
|
||||||
import { StepSessionExecutor } from "./step-session-executor.js";
|
import { StepSessionExecutor } from "./step-session-executor.js";
|
||||||
import { makeAncestryBlastRadiusGuard, resetStepToBaseline, runTaskStep } from "./step-runner.js";
|
import { makeAncestryBlastRadiusGuard, resetStepToBaseline, runTaskStep, type RunTaskStepResult } from "./step-runner.js";
|
||||||
// FNXC:MergerUnification 2026-06-21-19:05: the foundation branch imported `acquireWorkspaceRepoWorktree` here but never used it in executor.ts (the agent tool wraps it via agent-tools.ts), which fails lint on the inherited base. Removed until master-plan U1 re-adds it together with its per-repo acquisition usage.
|
// FNXC:MergerUnification 2026-06-21-19:05: the foundation branch imported `acquireWorkspaceRepoWorktree` here but never used it in executor.ts (the agent tool wraps it via agent-tools.ts), which fails lint on the inherited base. Removed until master-plan U1 re-adds it together with its per-repo acquisition usage.
|
||||||
import { acquireTaskWorktree, type AcquireTaskWorktreeResult } from "./worktree-acquisition.js";
|
import { acquireTaskWorktree, type AcquireTaskWorktreeResult } from "./worktree-acquisition.js";
|
||||||
import { resolveCapturedBaseCommitSha } from "./base-commit-capture.js";
|
import { resolveCapturedBaseCommitSha } from "./base-commit-capture.js";
|
||||||
@@ -5659,7 +5659,7 @@ export class TaskExecutor {
|
|||||||
const workItems = await this.store.listWorkflowWorkItemsForTask?.(task.id, { kinds: ["task"] }) ?? [];
|
const workItems = await this.store.listWorkflowWorkItemsForTask?.(task.id, { kinds: ["task"] }) ?? [];
|
||||||
for (let index = workItems.length - 1; index >= 0; index -= 1) {
|
for (let index = workItems.length - 1; index >= 0; index -= 1) {
|
||||||
const candidate = workItems[index];
|
const candidate = workItems[index];
|
||||||
if (["held", "runnable", "running", "retrying"].includes(candidate.state)) {
|
if (ACTIVE_WORKFLOW_WORK_ITEM_STATES.includes(candidate.state)) {
|
||||||
continuation = candidate;
|
continuation = candidate;
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
@@ -5956,7 +5956,7 @@ export class TaskExecutor {
|
|||||||
clearPin: pinPersistence.clearPin,
|
clearPin: pinPersistence.clearPin,
|
||||||
onSuspend: async (suspension) => {
|
onSuspend: async (suspension) => {
|
||||||
const items = await this.store.listWorkflowWorkItemsForTask(task.id, { kinds: ["task"] });
|
const items = await this.store.listWorkflowWorkItemsForTask(task.id, { kinds: ["task"] });
|
||||||
const live = items.filter((item) => ["held", "runnable", "running", "retrying"].includes(item.state));
|
const live = items.filter((item) => ACTIVE_WORKFLOW_WORK_ITEM_STATES.includes(item.state));
|
||||||
if (live.some((item) => item.nodeId === suspension.nodeId)) return;
|
if (live.some((item) => item.nodeId === suspension.nodeId)) return;
|
||||||
await this.store.replaceActiveTaskWorkflowContinuation({
|
await this.store.replaceActiveTaskWorkflowContinuation({
|
||||||
runId: `${workflowRunId ?? `${task.id}:workflow`}:continuation:${suspension.nodeId}:${items.length}`,
|
runId: `${workflowRunId ?? `${task.id}:workflow`}:continuation:${suspension.nodeId}:${items.length}`,
|
||||||
@@ -5996,7 +5996,7 @@ export class TaskExecutor {
|
|||||||
target: event.taskId,
|
target: event.taskId,
|
||||||
metadata:
|
metadata:
|
||||||
event.type === "task:column-transition"
|
event.type === "task:column-transition"
|
||||||
? { taskId: event.taskId, workflowId: event.workflowId, fromColumn: event.fromColumn, toColumn: event.toColumn, nodeId: event.nodeId }
|
? { taskId: event.taskId, workflowId: event.workflowId, fromColumn: event.fromColumn, toColumn: event.toColumn, nodeId: event.nodeId, irHash: event.irHash }
|
||||||
: { taskId: event.taskId, workflowId: event.workflowId, pinnedNodeId: event.pinnedNodeId, reason: event.reason },
|
: { taskId: event.taskId, workflowId: event.workflowId, pinnedNodeId: event.pinnedNodeId, reason: event.reason },
|
||||||
});
|
});
|
||||||
},
|
},
|
||||||
@@ -6688,6 +6688,57 @@ export class TaskExecutor {
|
|||||||
return only;
|
return only;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Project a graph-owned step only after it has a real worktree.
|
||||||
|
*
|
||||||
|
* A fresh task has no worktree until the authoritative implementation pass
|
||||||
|
* acquires one. Projecting before that pass produces a false "step started"
|
||||||
|
* event and captures the baseline from the project root. In that fresh path,
|
||||||
|
* let the implementation pass own the first projection and reuse the base SHA
|
||||||
|
* it captures during worktree acquisition. Resumed and isolated-step runs
|
||||||
|
* already have a worktree, so they keep the normal per-step projection and
|
||||||
|
* pre-work baseline behavior.
|
||||||
|
*/
|
||||||
|
private async runProjectedGraphTaskStep(
|
||||||
|
task: Task,
|
||||||
|
live: TaskDetail,
|
||||||
|
stepIndex: number,
|
||||||
|
active: ForeachActiveContext,
|
||||||
|
governingNodeId?: string,
|
||||||
|
thinkingLevel?: ThinkingLevel,
|
||||||
|
): Promise<RunTaskStepResult> {
|
||||||
|
const worktreePath = active.worktreePath || live.worktree;
|
||||||
|
const runStep = (idx: number) =>
|
||||||
|
this.runGraphTaskStep(
|
||||||
|
task,
|
||||||
|
idx,
|
||||||
|
active.instanceId,
|
||||||
|
governingNodeId,
|
||||||
|
thinkingLevel,
|
||||||
|
);
|
||||||
|
|
||||||
|
if (!worktreePath) {
|
||||||
|
const result = await runStep(stepIndex);
|
||||||
|
const refreshed = await this.store.getTask(task.id).catch(() => live);
|
||||||
|
return {
|
||||||
|
outcome: result.success ? "success" : "failure",
|
||||||
|
baselineSha: refreshed.baseCommitSha,
|
||||||
|
checkpointId: undefined,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
return runTaskStep(
|
||||||
|
{
|
||||||
|
store: this.store,
|
||||||
|
worktreePath,
|
||||||
|
runStep,
|
||||||
|
},
|
||||||
|
{ id: task.id, steps: live.steps },
|
||||||
|
stepIndex,
|
||||||
|
{ markDoneOnSuccess: active.deferDoneToReview !== true, projectionSource: "graph" },
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
/** Public authoritative-driver seam factory: exposes the same real lifecycle
|
/** Public authoritative-driver seam factory: exposes the same real lifecycle
|
||||||
* seams the internal graph runner uses, without changing legacy behavior. */
|
* seams the internal graph runner uses, without changing legacy behavior. */
|
||||||
public createAuthoritativeWorkflowPrimitives(settings: Settings): WorkflowRuntimePrimitives {
|
public createAuthoritativeWorkflowPrimitives(settings: Settings): WorkflowRuntimePrimitives {
|
||||||
@@ -6790,24 +6841,14 @@ export class TaskExecutor {
|
|||||||
data: { status: liveStatus },
|
data: { status: liveStatus },
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
const worktreePath = active.worktreePath || live.worktree || this.rootDir;
|
|
||||||
this.graphStepActiveContext.set(this.graphActiveContextKey(task.id, active.instanceId), active);
|
this.graphStepActiveContext.set(this.graphActiveContextKey(task.id, active.instanceId), active);
|
||||||
const stepGoverningNodeId = context[SEAM_GOVERNING_NODE_CONTEXT_KEY];
|
const stepGoverningNodeId = context[SEAM_GOVERNING_NODE_CONTEXT_KEY];
|
||||||
return await runTaskStep(
|
return await this.runProjectedGraphTaskStep(
|
||||||
{
|
task,
|
||||||
store: this.store,
|
live,
|
||||||
worktreePath,
|
|
||||||
runStep: (idx) =>
|
|
||||||
this.runGraphTaskStep(
|
|
||||||
task,
|
|
||||||
idx,
|
|
||||||
active.instanceId,
|
|
||||||
typeof stepGoverningNodeId === "string" ? stepGoverningNodeId : undefined,
|
|
||||||
),
|
|
||||||
},
|
|
||||||
{ id: task.id, steps: live.steps },
|
|
||||||
stepIndex,
|
stepIndex,
|
||||||
{ markDoneOnSuccess: active.deferDoneToReview !== true, projectionSource: "graph" },
|
active,
|
||||||
|
typeof stepGoverningNodeId === "string" ? stepGoverningNodeId : undefined,
|
||||||
);
|
);
|
||||||
},
|
},
|
||||||
resetTaskStep: async (ctx, task, stepIndex, baselineSha, checkpointId) => {
|
resetTaskStep: async (ctx, task, stepIndex, baselineSha, checkpointId) => {
|
||||||
@@ -7390,7 +7431,6 @@ export class TaskExecutor {
|
|||||||
// worktree when the foreach allocated one; otherwise the task's main
|
// worktree when the foreach allocated one; otherwise the task's main
|
||||||
// worktree (shared isolation — unchanged). The file-scope guard the session
|
// worktree (shared isolation — unchanged). The file-scope guard the session
|
||||||
// machinery installs applies to either worktree unchanged (not bypassed).
|
// machinery installs applies to either worktree unchanged (not bypassed).
|
||||||
const worktreePath = active.worktreePath || live.worktree || this.rootDir;
|
|
||||||
// Stamp the active instance so `runGraphTaskStep` can honor
|
// Stamp the active instance so `runGraphTaskStep` can honor
|
||||||
// `deferDoneToReview` when judging a non-terminal step (FIX 3).
|
// `deferDoneToReview` when judging a non-terminal step (FIX 3).
|
||||||
this.graphStepActiveContext.set(this.graphActiveContextKey(seamTask.id, active.instanceId), active);
|
this.graphStepActiveContext.set(this.graphActiveContextKey(seamTask.id, active.instanceId), active);
|
||||||
@@ -7405,35 +7445,15 @@ export class TaskExecutor {
|
|||||||
// foreach (overwrite mid-build, or clear while the shared pass is live).
|
// foreach (overwrite mid-build, or clear while the shared pass is live).
|
||||||
const stepGoverningNodeId = context[SEAM_GOVERNING_NODE_CONTEXT_KEY];
|
const stepGoverningNodeId = context[SEAM_GOVERNING_NODE_CONTEXT_KEY];
|
||||||
const seamThinkingLevel = context[SEAM_THINKING_LEVEL_CONTEXT_KEY];
|
const seamThinkingLevel = context[SEAM_THINKING_LEVEL_CONTEXT_KEY];
|
||||||
const result: Awaited<ReturnType<typeof runTaskStep>> = await runTaskStep(
|
const result = await this.runProjectedGraphTaskStep(
|
||||||
{
|
seamTask,
|
||||||
store: this.store,
|
live,
|
||||||
worktreePath,
|
|
||||||
// U6/U8: graph-owned per-step physics. Per-step-review workflows
|
|
||||||
// pin StepSessionExecutor inside runGraphTaskStep; final-review coding
|
|
||||||
// honors runStepsInNewSessions and may reuse one executor session.
|
|
||||||
// Thread the instanceId so the active-context read is per-instance
|
|
||||||
// (parallel-foreach safe).
|
|
||||||
runStep: (stepIndex) =>
|
|
||||||
this.runGraphTaskStep(
|
|
||||||
seamTask,
|
|
||||||
stepIndex,
|
|
||||||
active.instanceId,
|
|
||||||
typeof stepGoverningNodeId === "string" ? stepGoverningNodeId : undefined,
|
|
||||||
typeof seamThinkingLevel === "string" && WORKFLOW_THINKING_LEVEL_SET.has(seamThinkingLevel)
|
|
||||||
? (seamThinkingLevel as ThinkingLevel)
|
|
||||||
: undefined,
|
|
||||||
),
|
|
||||||
},
|
|
||||||
{ id: seamTask.id, steps: live.steps },
|
|
||||||
active.stepIndex,
|
active.stepIndex,
|
||||||
{
|
active,
|
||||||
// Single-authority done-marking (U6/KTD-4): when the foreach template
|
typeof stepGoverningNodeId === "string" ? stepGoverningNodeId : undefined,
|
||||||
// has a step-review node, leave the step in-progress so the review's
|
typeof seamThinkingLevel === "string" && WORKFLOW_THINKING_LEVEL_SET.has(seamThinkingLevel)
|
||||||
// APPROVE marks it done (the review is the single done authority).
|
? (seamThinkingLevel as ThinkingLevel)
|
||||||
markDoneOnSuccess: active.deferDoneToReview !== true,
|
: undefined,
|
||||||
projectionSource: "graph",
|
|
||||||
},
|
|
||||||
);
|
);
|
||||||
// Capture baseline/checkpoint back into the reserved active context so the
|
// Capture baseline/checkpoint back into the reserved active context so the
|
||||||
// foreach sub-walk threads them to later template nodes (step-review/reset).
|
// foreach sub-walk threads them to later template nodes (step-review/reset).
|
||||||
|
|||||||
@@ -42,6 +42,7 @@ import {
|
|||||||
resolveColumnFlags,
|
resolveColumnFlags,
|
||||||
resolveColumnAdjacency,
|
resolveColumnAdjacency,
|
||||||
PLAN_REVIEW_GROUP_ID,
|
PLAN_REVIEW_GROUP_ID,
|
||||||
|
ACTIVE_WORKFLOW_WORK_ITEM_STATES,
|
||||||
DEFAULT_WORKFLOW_POOL_ID,
|
DEFAULT_WORKFLOW_POOL_ID,
|
||||||
TransitionRejectionError,
|
TransitionRejectionError,
|
||||||
resolveWorkflowIrForTask,
|
resolveWorkflowIrForTask,
|
||||||
@@ -142,6 +143,7 @@ function isHeldTask(ir: WorkflowIr, task: Task): boolean {
|
|||||||
* `promoteHeldTask`, `releaseHeldTaskByEvent`) enforces the same invariant.
|
* `promoteHeldTask`, `releaseHeldTaskByEvent`) enforces the same invariant.
|
||||||
*/
|
*/
|
||||||
/**
|
/**
|
||||||
|
* FNXC:PlanReview 2026-07-21-12:20:
|
||||||
* Locate a Plan Review node placed before WIP. Disabled optional groups still
|
* Locate a Plan Review node placed before WIP. Disabled optional groups still
|
||||||
* traverse this node, allowing the graph to persist the same generic capacity
|
* traverse this node, allowing the graph to persist the same generic capacity
|
||||||
* continuation without invoking a reviewer.
|
* continuation without invoking a reviewer.
|
||||||
@@ -176,9 +178,7 @@ export async function isUnplannedForExecution(store: TaskStore, task: Task, ir:
|
|||||||
if (!legacyPassed) {
|
if (!legacyPassed) {
|
||||||
if (typeof store.listWorkflowWorkItemsForTask !== "function") return true;
|
if (typeof store.listWorkflowWorkItemsForTask !== "function") return true;
|
||||||
const continuations = await store.listWorkflowWorkItemsForTask(task.id, { kinds: ["task"] });
|
const continuations = await store.listWorkflowWorkItemsForTask(task.id, { kinds: ["task"] });
|
||||||
const active = continuations.filter((item) =>
|
const active = continuations.filter((item) => ACTIVE_WORKFLOW_WORK_ITEM_STATES.includes(item.state));
|
||||||
["held", "runnable", "running", "retrying"].includes(item.state),
|
|
||||||
);
|
|
||||||
// Readiness is represented by the graph's durable boundary continuation,
|
// Readiness is represented by the graph's durable boundary continuation,
|
||||||
// not by a special-case review result. Optional groups that are disabled
|
// not by a special-case review result. Optional groups that are disabled
|
||||||
// are still traversed and therefore reach the same capacity boundary.
|
// are still traversed and therefore reach the same capacity boundary.
|
||||||
|
|||||||
@@ -14,11 +14,13 @@ import type {
|
|||||||
GithubIssueAction,
|
GithubIssueAction,
|
||||||
CliSession,
|
CliSession,
|
||||||
NotificationPayload,
|
NotificationPayload,
|
||||||
|
WorkflowWorkItem,
|
||||||
} from "@fusion/core";
|
} from "@fusion/core";
|
||||||
import {
|
import {
|
||||||
AsyncCentralClaimStore,
|
AsyncCentralClaimStore,
|
||||||
ChatStore,
|
ChatStore,
|
||||||
computeWorkflowIrPin,
|
computeWorkflowIrPin,
|
||||||
|
ACTIVE_WORKFLOW_WORK_ITEM_STATES,
|
||||||
isEphemeralAgent,
|
isEphemeralAgent,
|
||||||
resolveWorkflowIrForTask,
|
resolveWorkflowIrForTask,
|
||||||
} from "@fusion/core";
|
} from "@fusion/core";
|
||||||
@@ -66,6 +68,26 @@ const yieldEventLoop = (): Promise<void> => new Promise((resolve) => setImmediat
|
|||||||
export const CLI_AGENT_AWAITING_INPUT_EVENT = "cli-agent-awaiting-input" as const;
|
export const CLI_AGENT_AWAITING_INPUT_EVENT = "cli-agent-awaiting-input" as const;
|
||||||
const TASK_PLANNER_CHAT_AGENT_ID_PREFIX = "task-planner:";
|
const TASK_PLANNER_CHAT_AGENT_ID_PREFIX = "task-planner:";
|
||||||
|
|
||||||
|
export interface PlanningContinuationCandidate {
|
||||||
|
item: WorkflowWorkItem;
|
||||||
|
task: Task | null | undefined;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** FNXC:WorkflowScheduling 2026-07-21-12:30:
|
||||||
|
* Select due planning continuations whose task remains dispatchable. */
|
||||||
|
export function selectActionablePlanningContinuations(
|
||||||
|
candidates: readonly PlanningContinuationCandidate[],
|
||||||
|
): Array<{ item: WorkflowWorkItem; task: Task }> {
|
||||||
|
return candidates.filter(
|
||||||
|
(candidate): candidate is { item: WorkflowWorkItem; task: Task } =>
|
||||||
|
candidate.item.waitReason === "planning"
|
||||||
|
&& candidate.task !== null
|
||||||
|
&& candidate.task !== undefined
|
||||||
|
&& candidate.task.paused !== true
|
||||||
|
&& candidate.task.userPaused !== true,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
export interface CliAgentAwaitingInputNotificationInfo {
|
export interface CliAgentAwaitingInputNotificationInfo {
|
||||||
sessionId: string;
|
sessionId: string;
|
||||||
notification: Record<string, unknown> | undefined;
|
notification: Record<string, unknown> | undefined;
|
||||||
@@ -996,8 +1018,14 @@ export class InProcessRuntime
|
|||||||
const planReview = resolvePreReleasePlanReviewNode(ir);
|
const planReview = resolvePreReleasePlanReviewNode(ir);
|
||||||
if (!planReview || planReview.column !== live.column) return;
|
if (!planReview || planReview.column !== live.column) return;
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:PlanReview 2026-07-21-12:20:
|
||||||
|
Specification completion creates a planning continuation only
|
||||||
|
when the review node belongs to the card's current column and no
|
||||||
|
active continuation already owns the task.
|
||||||
|
*/
|
||||||
const active = await this.taskStore.listWorkflowWorkItemsForTask(live.id, { kinds: ["task"] });
|
const active = await this.taskStore.listWorkflowWorkItemsForTask(live.id, { kinds: ["task"] });
|
||||||
if (!active.some((item) => ["held", "runnable", "running", "retrying"].includes(item.state))) {
|
if (!active.some((item) => ACTIVE_WORKFLOW_WORK_ITEM_STATES.includes(item.state))) {
|
||||||
await this.taskStore.replaceActiveTaskWorkflowContinuation({
|
await this.taskStore.replaceActiveTaskWorkflowContinuation({
|
||||||
runId: `${live.id}:planning-continuation:${planReview.id}:${active.length}`,
|
runId: `${live.id}:planning-continuation:${planReview.id}:${active.length}`,
|
||||||
taskId: live.id,
|
taskId: live.id,
|
||||||
@@ -1840,7 +1868,11 @@ export class InProcessRuntime
|
|||||||
return this.missionExecutionLoop;
|
return this.missionExecutionLoop;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Wake the durable task-continuation consumer without nesting execution in triage. */
|
/**
|
||||||
|
* FNXC:WorkflowScheduling 2026-07-21-12:20:
|
||||||
|
* Wake the durable task-continuation consumer in a microtask so triage can
|
||||||
|
* release its own execution slot before continuation dispatch begins.
|
||||||
|
*/
|
||||||
private kickWorkflowContinuationProcessor(): void {
|
private kickWorkflowContinuationProcessor(): void {
|
||||||
queueMicrotask(() => {
|
queueMicrotask(() => {
|
||||||
void this.drainWorkflowContinuations().catch((error) => {
|
void this.drainWorkflowContinuations().catch((error) => {
|
||||||
@@ -1850,6 +1882,11 @@ export class InProcessRuntime
|
|||||||
}
|
}
|
||||||
|
|
||||||
private async drainWorkflowContinuations(): Promise<void> {
|
private async drainWorkflowContinuations(): Promise<void> {
|
||||||
|
/*
|
||||||
|
FNXC:WorkflowScheduling 2026-07-21-12:20:
|
||||||
|
A single runtime drain owns selection at a time. Concurrent wakeups collapse
|
||||||
|
behind this guard and the recurring processor supplies the next bounded pass.
|
||||||
|
*/
|
||||||
if (this.workflowContinuationDrainActive || this.status !== "active") return;
|
if (this.workflowContinuationDrainActive || this.status !== "active") return;
|
||||||
this.workflowContinuationDrainActive = true;
|
this.workflowContinuationDrainActive = true;
|
||||||
try {
|
try {
|
||||||
@@ -1858,10 +1895,12 @@ export class InProcessRuntime
|
|||||||
states: ["runnable", "retrying"],
|
states: ["runnable", "retrying"],
|
||||||
limit: 20,
|
limit: 20,
|
||||||
});
|
});
|
||||||
|
const candidates: PlanningContinuationCandidate[] = [];
|
||||||
for (const item of items) {
|
for (const item of items) {
|
||||||
if (item.waitReason !== "planning") continue;
|
const task = await this.taskStore.getTask(item.taskId);
|
||||||
const task = await this.taskStore.getTask(item.taskId).catch(() => undefined);
|
candidates.push({ item, task });
|
||||||
if (!task || task.paused || task.userPaused) continue;
|
}
|
||||||
|
for (const { item, task } of selectActionablePlanningContinuations(candidates)) {
|
||||||
void this.executor.execute(task).catch((error) => {
|
void this.executor.execute(task).catch((error) => {
|
||||||
runtimeLog.error(`Workflow continuation ${item.id} failed:`, error);
|
runtimeLog.error(`Workflow continuation ${item.id} failed:`, error);
|
||||||
});
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user