fix(engine): reject out-of-order step starts (#2403)
## Summary Ordered task steps can no longer appear active ahead of unfinished predecessors. Step starts now use the same dependency-aware ordering guard as completions, while steps explicitly declared independent remain parallelizable. Rejected executor updates explain that the lifecycle transition was suppressed instead of implying completed work was overwritten. ## Validation - Reproduced the FN-8490 concurrent update sequence and verified later steps remain pending. - Passed 15 PostgreSQL step-order tests, the focused executor response test, core and engine typechecks, changeset validation, and `pnpm verify:fast` including boot smoke. - The full `executor-prompt.test.ts` run retains five pause-behavior expectation failures that reproduce unchanged on `origin/main`. <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **New Features** * Enhanced the step start hook to support an awaited “pre-start projection” that can reject startup via `false` (sync or async), preventing step-session creation/completion. * Added a step-start “verdict” so steps can be started or blocked deterministically (including “resumed” behavior). * **Bug Fixes** * Prevented ordered/dependency steps from transitioning out-of-order by enforcing guards for both in-progress and done transitions, including concurrent update attempts. * Improved integrity/out-of-order warning behavior and suppression details when persisted status doesn’t match expectations. * **Tests** * Added/updated PostgreSQL and engine regression coverage for blocked/resumed start and start-rejection control flow. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
7
.changeset/ordered-step-starts.md
Normal file
7
.changeset/ordered-step-starts.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Prevent executors from starting ordered task steps before their required predecessors finish.
|
||||||
|
category: fix
|
||||||
|
dev: Applies dependency-aware ordering to both in-progress and done step transitions.
|
||||||
@@ -35,6 +35,68 @@ pgTest("TaskStore.updateStep step-order guard (PostgreSQL)", () => {
|
|||||||
expect(updated.steps[2].status).toBe("pending");
|
expect(updated.steps[2].status).toBe("pending");
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("does not start a later step when an earlier ordered step is still in progress", async () => {
|
||||||
|
const store = h.store();
|
||||||
|
const task = await h.createTaskWithSteps();
|
||||||
|
|
||||||
|
await store.updateStep(task.id, 0, "in-progress");
|
||||||
|
await Promise.all([
|
||||||
|
store.updateStep(task.id, 1, "done"),
|
||||||
|
store.updateStep(task.id, 2, "in-progress"),
|
||||||
|
]);
|
||||||
|
|
||||||
|
const updated = await store.getTask(task.id);
|
||||||
|
expect(updated.steps.map((step) => step.status)).toEqual([
|
||||||
|
"in-progress",
|
||||||
|
"pending",
|
||||||
|
"pending",
|
||||||
|
]);
|
||||||
|
expect(updated.currentStep).toBe(0);
|
||||||
|
expect(updated.log.map((entry) => entry.action)).toContainEqual(
|
||||||
|
expect.stringContaining("Ignored out-of-order in-progress for step 2"),
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
* FNXC:StepLifecycle 2026-07-22-10:30:
|
||||||
|
* Legacy state can already contain the inversion being repaired. The atomic start verdict
|
||||||
|
* must report the dependency rejection even though the target remains in-progress.
|
||||||
|
*/
|
||||||
|
it("reports blocked for a corrupted in-progress step with an unfinished predecessor", async () => {
|
||||||
|
const store = h.store();
|
||||||
|
const task = await h.createTaskWithSteps();
|
||||||
|
await store.updateTask(task.id, {
|
||||||
|
steps: task.steps.map((step, index) => ({
|
||||||
|
...step,
|
||||||
|
status: index < 2 ? "in-progress" as const : step.status,
|
||||||
|
})),
|
||||||
|
});
|
||||||
|
|
||||||
|
const result = await store.startStep(task.id, 1);
|
||||||
|
|
||||||
|
expect(result).toMatchObject({
|
||||||
|
accepted: false,
|
||||||
|
disposition: "blocked",
|
||||||
|
blockingStepIndex: 0,
|
||||||
|
});
|
||||||
|
expect(result.task.steps[1].status).toBe("in-progress");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("reports resumed for a valid in-progress step whose predecessors are terminal", async () => {
|
||||||
|
const store = h.store();
|
||||||
|
const task = await h.createTaskWithSteps();
|
||||||
|
await store.updateStep(task.id, 0, "done");
|
||||||
|
await store.updateStep(task.id, 1, "in-progress");
|
||||||
|
|
||||||
|
const result = await store.startStep(task.id, 1);
|
||||||
|
|
||||||
|
expect(result).toMatchObject({
|
||||||
|
accepted: true,
|
||||||
|
disposition: "resumed",
|
||||||
|
});
|
||||||
|
expect(result.task.steps[1].status).toBe("in-progress");
|
||||||
|
});
|
||||||
|
|
||||||
it("allows an explicitly independent step to finish out of index order", async () => {
|
it("allows an explicitly independent step to finish out of index order", async () => {
|
||||||
const store = h.store();
|
const store = h.store();
|
||||||
const task = await h.createTaskWithSteps();
|
const task = await h.createTaskWithSteps();
|
||||||
@@ -77,6 +139,25 @@ pgTest("TaskStore.updateStep step-order guard (PostgreSQL)", () => {
|
|||||||
expect(updated.log.at(-1)?.action).toContain("dependency step 1");
|
expect(updated.log.at(-1)?.action).toContain("dependency step 1");
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("blocks a start when an explicit dependency is still in progress", async () => {
|
||||||
|
const store = h.store();
|
||||||
|
const task = await h.createTaskWithSteps();
|
||||||
|
await store.updateTask(task.id, {
|
||||||
|
steps: task.steps.map((step, index) =>
|
||||||
|
index === 2 ? { ...step, dependsOn: [1] } : step,
|
||||||
|
),
|
||||||
|
});
|
||||||
|
|
||||||
|
await store.updateStep(task.id, 0, "done");
|
||||||
|
const started = await store.updateStep(task.id, 1, "in-progress");
|
||||||
|
expect(started.steps[1].status).toBe("in-progress");
|
||||||
|
|
||||||
|
const updated = await store.updateStep(task.id, 2, "in-progress");
|
||||||
|
|
||||||
|
expect(updated.steps[2].status).toBe("pending");
|
||||||
|
expect(updated.log.at(-1)?.action).toContain("dependency step 1");
|
||||||
|
});
|
||||||
|
|
||||||
it("allows a satisfied explicit dependency while an unrelated earlier step is pending", async () => {
|
it("allows a satisfied explicit dependency while an unrelated earlier step is pending", async () => {
|
||||||
const store = h.store();
|
const store = h.store();
|
||||||
const task = await h.createTaskWithSteps();
|
const task = await h.createTaskWithSteps();
|
||||||
@@ -112,6 +193,21 @@ pgTest("TaskStore.updateStep step-order guard (PostgreSQL)", () => {
|
|||||||
expect(updated.steps[2].status).toBe("done");
|
expect(updated.steps[2].status).toBe("done");
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("blocks a graph-source start while its implicit predecessor is unfinished", async () => {
|
||||||
|
const store = h.store();
|
||||||
|
const task = await h.createTaskWithSteps();
|
||||||
|
|
||||||
|
const updated = await store.updateStep(task.id, 2, "in-progress", { source: "graph" });
|
||||||
|
|
||||||
|
expect(updated.steps[2].status).toBe("pending");
|
||||||
|
expect(updated.log.map((entry) => entry.action)).toContainEqual(
|
||||||
|
expect.stringContaining("Ignored dependency-order in-progress for step 2"),
|
||||||
|
);
|
||||||
|
expect(updated.log.map((entry) => entry.action)).toContainEqual(
|
||||||
|
expect.stringContaining("[integrity-warning] graph-source updateStep suppressed"),
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
it.each([{ dependsOn: [2] }, { dependsOn: [99] }, { dependsOn: [-1] }, { dependsOn: [1.5] }])(
|
it.each([{ dependsOn: [2] }, { dependsOn: [99] }, { dependsOn: [-1] }, { dependsOn: [1.5] }])(
|
||||||
"falls back to strict ordering for malformed dependsOn %j",
|
"falls back to strict ordering for malformed dependsOn %j",
|
||||||
async ({ dependsOn }) => {
|
async ({ dependsOn }) => {
|
||||||
|
|||||||
@@ -126,7 +126,7 @@ import { flushAgentLogBufferImpl, appendAgentLogBatchImpl } from "./task-store/a
|
|||||||
import { refineTaskImpl, updateTaskDependenciesImpl } from "./task-store/update-task-deps.js";
|
import { refineTaskImpl, updateTaskDependenciesImpl } from "./task-store/update-task-deps.js";
|
||||||
import { createWorkflowStepImpl, updateWorkflowStepImpl, updateWorkflowDefinitionImpl, deleteWorkflowDefinitionImpl, setDefaultWorkflowIdImpl, selectTaskWorkflowImpl } from "./task-store/workflow-ops.js";
|
import { createWorkflowStepImpl, updateWorkflowStepImpl, updateWorkflowDefinitionImpl, deleteWorkflowDefinitionImpl, setDefaultWorkflowIdImpl, selectTaskWorkflowImpl } from "./task-store/workflow-ops.js";
|
||||||
import { initImpl, setupActivityLogListenersImpl, reconcileOrphanedTaskDirsImpl, watchImpl, checkForChangesImpl, migrateAgentLogEntriesImpl, migrateMovedSettingsImpl, recoverStaleTransitionPendingImpl, migrateLegacyWorkflowStepsImpl, emitTaskLifecycleEventSafelyImpl } from "./task-store/lifecycle-ops.js";
|
import { initImpl, setupActivityLogListenersImpl, reconcileOrphanedTaskDirsImpl, watchImpl, checkForChangesImpl, migrateAgentLogEntriesImpl, migrateMovedSettingsImpl, recoverStaleTransitionPendingImpl, migrateLegacyWorkflowStepsImpl, emitTaskLifecycleEventSafelyImpl } from "./task-store/lifecycle-ops.js";
|
||||||
import { updateStepImpl, acquireMergeQueueLeaseImpl, mergeTaskImpl } from "./task-store/merge-queue-ops.js";
|
import { updateStepImpl, startStepImpl, acquireMergeQueueLeaseImpl, mergeTaskImpl } from "./task-store/merge-queue-ops.js";
|
||||||
import { addCommentImpl, publishArchivedTaskDocumentAdditionImpl, upsertTaskDocumentImpl } from "./task-store/comments-ops.js";
|
import { addCommentImpl, publishArchivedTaskDocumentAdditionImpl, upsertTaskDocumentImpl } from "./task-store/comments-ops.js";
|
||||||
import { deleteTaskImpl, deleteTaskIfImpl, archiveTaskImpl, type DeleteTaskIfResult } from "./task-store/archive-lifecycle.js";
|
import { deleteTaskImpl, deleteTaskIfImpl, archiveTaskImpl, type DeleteTaskIfResult } from "./task-store/archive-lifecycle.js";
|
||||||
import { updateSettingsImpl, updateGlobalSettingsImpl } from "./task-store/settings-ops.js";
|
import { updateSettingsImpl, updateGlobalSettingsImpl } from "./task-store/settings-ops.js";
|
||||||
@@ -1426,6 +1426,10 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
|||||||
async updateStep( id: string, stepIndex: number, status: import("./types.js").StepStatus, options?: { source?: "graph" }, ): Promise<Task> {
|
async updateStep( id: string, stepIndex: number, status: import("./types.js").StepStatus, options?: { source?: "graph" }, ): Promise<Task> {
|
||||||
return updateStepImpl(this, id, stepIndex, status, options);
|
return updateStepImpl(this, id, stepIndex, status, options);
|
||||||
}
|
}
|
||||||
|
// FNXC:StepLifecycle 2026-07-22-10:30: Execution callers need the locked start verdict; updateStep retains its legacy Task-only contract.
|
||||||
|
async startStep( id: string, stepIndex: number, options?: { source?: "graph" }, ): Promise<import("./task-store/merge-queue-ops.js").StepStartResult> {
|
||||||
|
return startStepImpl(this, id, stepIndex, options);
|
||||||
|
}
|
||||||
async logEntry(id: string, action: string, outcome?: string, runContext?: RunMutationContext): Promise<Task> {
|
async logEntry(id: string, action: string, outcome?: string, runContext?: RunMutationContext): Promise<Task> {
|
||||||
return logEntryImpl(this, id, action, outcome, runContext);
|
return logEntryImpl(this, id, action, outcome, runContext);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -18,6 +18,15 @@ import {assertSafeGitBranchName, assertSafeAbsolutePath} from "../task-store/she
|
|||||||
import {acquireMergeQueueLease as acquireMergeQueueLeaseAsync} from "../task-store/async-merge-coordination.js";
|
import {acquireMergeQueueLease as acquireMergeQueueLeaseAsync} from "../task-store/async-merge-coordination.js";
|
||||||
import type {MergeQueueRow} from "../task-store/row-types.js";
|
import type {MergeQueueRow} from "../task-store/row-types.js";
|
||||||
|
|
||||||
|
export type StepStartDisposition = "started" | "resumed" | "blocked" | "terminal";
|
||||||
|
|
||||||
|
export interface StepStartResult {
|
||||||
|
task: Task;
|
||||||
|
accepted: boolean;
|
||||||
|
disposition: StepStartDisposition;
|
||||||
|
blockingStepIndex?: number;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Step state is written from more places than an agent's explicit
|
* Step state is written from more places than an agent's explicit
|
||||||
* `fn_task_update` call: workflow projection, review auto-approval, restart
|
* `fn_task_update` call: workflow projection, review auto-approval, restart
|
||||||
@@ -53,15 +62,16 @@ async function appendProactiveStepStatus(store: TaskStore, taskId: string, messa
|
|||||||
await store.appendAgentLog(taskId, message, "status", undefined, "executor");
|
await store.appendAgentLog(taskId, message, "status", undefined, "executor");
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function updateStepImpl(store: TaskStore, id: string, stepIndex: number, status: import("../types.js").StepStatus, options?: { source?: "graph" },): Promise<Task> {
|
async function mutateStepImpl(store: TaskStore, id: string, stepIndex: number, status: import("../types.js").StepStatus, options?: { source?: "graph" },): Promise<{ task: Task; startResult?: Omit<StepStartResult, "task"> }> {
|
||||||
// FNXC:WorkflowStepOrdering 2026-07-20-20:05:
|
// FNXC:WorkflowStepOrdering 2026-07-20-20:05:
|
||||||
// Step-inversion projection discipline (U6/KTD-7). A `source: "graph"` write
|
// Step-inversion projection discipline (U6/KTD-7). A `source: "graph"` write
|
||||||
// is the workflow-graph executor projecting a foreach instance's lifecycle
|
// is the workflow-graph executor projecting a foreach instance's lifecycle
|
||||||
// (in-progress / done / pending) onto Task.steps[] with EXPLICIT indices. Three
|
// (in-progress / done / pending) onto Task.steps[] with EXPLICIT indices. Three
|
||||||
// behaviors diverge from a legacy write with no explicit dependency metadata:
|
// behaviors diverge from a legacy write with no explicit dependency metadata:
|
||||||
// (a) the out-of-order-done guard uses DEPENDENCY order (a done write is
|
// (a) the out-of-order start/completion guard uses DEPENDENCY order (an
|
||||||
// legal when every dependsOn step — default: the immediately-preceding
|
// in-progress or done write is legal when every dependsOn step —
|
||||||
// step — is done/skipped, KTD-11). Explicit dependsOn metadata is
|
// default: the immediately-preceding step — is done/skipped, KTD-11).
|
||||||
|
// Explicit dependsOn metadata is
|
||||||
// authoritative for every writer, not only graph-tagged writes;
|
// authoritative for every writer, not only graph-tagged writes;
|
||||||
// (b) a guard that DOES suppress a graph write logs an audit warning loudly
|
// (b) a guard that DOES suppress a graph write logs an audit warning loudly
|
||||||
// (legacy stays silent — a graph suppression is a projection bug);
|
// (legacy stays silent — a graph suppression is a projection bug);
|
||||||
@@ -127,16 +137,23 @@ export async function updateStepImpl(store: TaskStore, id: string, stepIndex: nu
|
|||||||
await store.atomicWriteTaskJson(dir, task);
|
await store.atomicWriteTaskJson(dir, task);
|
||||||
if (store.isWatching) store.taskCache.set(id, { ...task });
|
if (store.isWatching) store.taskCache.set(id, { ...task });
|
||||||
store.emit("task:updated", task);
|
store.emit("task:updated", task);
|
||||||
return task;
|
return {
|
||||||
|
task,
|
||||||
|
startResult: {
|
||||||
|
accepted: false,
|
||||||
|
disposition: "terminal",
|
||||||
|
},
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
if (status === "done") {
|
if (status === "done" || status === "in-progress") {
|
||||||
// The set of predecessor steps that must be done/skipped before this step
|
// The set of predecessor steps that must be done/skipped before this step
|
||||||
// may go done. Explicit dependency metadata is authoritative regardless
|
// may start or finish. Explicit dependency metadata is authoritative
|
||||||
// of which execution surface performs the write: parallel step sessions
|
// regardless of which execution surface performs the write: parallel
|
||||||
// and graph foreach instances share the same Task.steps[] contract. When
|
// step sessions and graph foreach instances share the same Task.steps[]
|
||||||
// metadata is absent, legacy callers retain strict index order while a
|
// contract. When metadata is absent, legacy callers retain strict index
|
||||||
// graph-source write defaults to the immediately preceding step (KTD-11).
|
// order while a graph-source write defaults to the immediately preceding
|
||||||
|
// step (KTD-11).
|
||||||
const explicitDependencies = task.steps[stepIndex]?.dependsOn;
|
const explicitDependencies = task.steps[stepIndex]?.dependsOn;
|
||||||
const hasExplicitDependencies = Array.isArray(explicitDependencies);
|
const hasExplicitDependencies = Array.isArray(explicitDependencies);
|
||||||
const validExplicitDependencies =
|
const validExplicitDependencies =
|
||||||
@@ -203,14 +220,25 @@ export async function updateStepImpl(store: TaskStore, id: string, stepIndex: nu
|
|||||||
timestamp: ts,
|
timestamp: ts,
|
||||||
action:
|
action:
|
||||||
`[integrity-warning] graph-source updateStep suppressed: step ${stepIndex} ` +
|
`[integrity-warning] graph-source updateStep suppressed: step ${stepIndex} ` +
|
||||||
`(${task.steps[stepIndex].name}) → done blocked by unmet dependency ` +
|
`(${task.steps[stepIndex].name}) → ${status} blocked by unmet dependency ` +
|
||||||
`step ${blockingIndex} (${blockingStatus})`,
|
`step ${blockingIndex} (${blockingStatus})`,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
await store.atomicWriteTaskJson(dir, task);
|
await store.atomicWriteTaskJson(dir, task);
|
||||||
if (store.isWatching) store.taskCache.set(id, { ...task });
|
if (store.isWatching) store.taskCache.set(id, { ...task });
|
||||||
store.emit("task:updated", task);
|
store.emit("task:updated", task);
|
||||||
return task;
|
return {
|
||||||
|
task,
|
||||||
|
...(status === "in-progress"
|
||||||
|
? {
|
||||||
|
startResult: {
|
||||||
|
accepted: false,
|
||||||
|
disposition: "blocked" as const,
|
||||||
|
blockingStepIndex: blockingIndex,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
: {}),
|
||||||
|
};
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -264,9 +292,38 @@ export async function updateStepImpl(store: TaskStore, id: string, stepIndex: nu
|
|||||||
id,
|
id,
|
||||||
proactiveStepStatusMessage(stepIndex, task.steps[stepIndex].name, currentStatus, status),
|
proactiveStepStatusMessage(stepIndex, task.steps[stepIndex].name, currentStatus, status),
|
||||||
).catch(() => undefined);
|
).catch(() => undefined);
|
||||||
return task;
|
return {
|
||||||
|
task,
|
||||||
|
...(status === "in-progress"
|
||||||
|
? {
|
||||||
|
startResult: {
|
||||||
|
accepted: true,
|
||||||
|
disposition: currentStatus === "in-progress" ? "resumed" as const : "started" as const,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
: {}),
|
||||||
|
};
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export async function updateStepImpl(store: TaskStore, id: string, stepIndex: number, status: import("../types.js").StepStatus, options?: { source?: "graph" },): Promise<Task> {
|
||||||
|
return (await mutateStepImpl(store, id, stepIndex, status, options)).task;
|
||||||
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:StepLifecycle 2026-07-22-10:30:
|
||||||
|
Execution must distinguish a dependency-blocked projection from a valid restart resume even when both return a task whose target step is already in-progress. Keep that verdict inside the same task lock as the dependency check; a caller-side pre-read would race and duplicate ordering policy.
|
||||||
|
*/
|
||||||
|
export async function startStepImpl(store: TaskStore, id: string, stepIndex: number, options?: { source?: "graph" },): Promise<StepStartResult> {
|
||||||
|
const result = await mutateStepImpl(store, id, stepIndex, "in-progress", options);
|
||||||
|
return {
|
||||||
|
task: result.task,
|
||||||
|
...(result.startResult ?? {
|
||||||
|
accepted: false,
|
||||||
|
disposition: "terminal" as const,
|
||||||
|
}),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
export async function acquireMergeQueueLeaseImpl(store: TaskStore, workerId: string, opts: MergeQueueAcquireOptions): Promise<MergeQueueEntry | null> {
|
export async function acquireMergeQueueLeaseImpl(store: TaskStore, workerId: string, opts: MergeQueueAcquireOptions): Promise<MergeQueueEntry | null> {
|
||||||
if (store.backendMode) {
|
if (store.backendMode) {
|
||||||
|
|||||||
@@ -1502,6 +1502,86 @@ describe("swallowed async store failure observability", () => {
|
|||||||
mockedWithRateLimitRetry.mockImplementation((fn: () => Promise<unknown>) => fn());
|
mockedWithRateLimitRetry.mockImplementation((fn: () => Promise<unknown>) => fn());
|
||||||
});
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
* FNXC:StepLifecycle 2026-07-22-10:30:
|
||||||
|
* A legacy inversion leaves the target in-progress even when the predecessor guard rejects
|
||||||
|
* its restart. The executor must consume the atomic verdict instead of inferring acceptance
|
||||||
|
* from that unchanged target status.
|
||||||
|
*/
|
||||||
|
it("turns a blocked corrupted in-progress start into a failed step-session result", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = {
|
||||||
|
id: "FN-8490",
|
||||||
|
title: "Ordered step start",
|
||||||
|
description: "Do not execute a rejected later step",
|
||||||
|
column: "in-progress" as const,
|
||||||
|
dependencies: [] as string[],
|
||||||
|
steps: [
|
||||||
|
{ name: "Step 0", status: "in-progress" as const },
|
||||||
|
{ name: "Step 1", status: "in-progress" as const },
|
||||||
|
],
|
||||||
|
currentStep: 0,
|
||||||
|
log: [] as any[],
|
||||||
|
prompt: "# test\n## Steps\n### Step 0: Preflight\n- [ ] check\n### Step 1: Implement\n- [ ] build",
|
||||||
|
worktree: "/tmp/test/.worktrees/fn-8490",
|
||||||
|
baseCommitSha: "abc123",
|
||||||
|
enabledWorkflowSteps: [],
|
||||||
|
createdAt: new Date().toISOString(),
|
||||||
|
updatedAt: new Date().toISOString(),
|
||||||
|
};
|
||||||
|
store.getSettings.mockResolvedValue({
|
||||||
|
maxConcurrent: 2,
|
||||||
|
maxWorktrees: 4,
|
||||||
|
pollIntervalMs: 15_000,
|
||||||
|
groupOverlappingFiles: false,
|
||||||
|
autoMerge: false,
|
||||||
|
runStepsInNewSessions: true,
|
||||||
|
maxParallelSteps: 1,
|
||||||
|
});
|
||||||
|
store.getTask.mockResolvedValue(task);
|
||||||
|
store.startStep
|
||||||
|
.mockResolvedValueOnce({
|
||||||
|
task,
|
||||||
|
accepted: true,
|
||||||
|
disposition: "resumed",
|
||||||
|
})
|
||||||
|
.mockResolvedValueOnce({
|
||||||
|
task,
|
||||||
|
accepted: false,
|
||||||
|
disposition: "blocked",
|
||||||
|
blockingStepIndex: 0,
|
||||||
|
});
|
||||||
|
mockExecuteAll.mockImplementation(async () => {
|
||||||
|
const options = mockedStepSessionExecutor.mock.calls.at(-1)?.[0] as {
|
||||||
|
onStepStart?: (stepIndex: number) => Promise<void | boolean>;
|
||||||
|
};
|
||||||
|
const accepted = await options.onStepStart?.(1);
|
||||||
|
return accepted === false
|
||||||
|
? [{ stepIndex: 1, success: false, error: "start rejected", retries: 0 }]
|
||||||
|
: [{ stepIndex: 1, success: true, retries: 0 }];
|
||||||
|
});
|
||||||
|
|
||||||
|
const onError = vi.fn();
|
||||||
|
const executor = new TaskExecutor(store, "/tmp/test", { onError });
|
||||||
|
await executor.execute(task);
|
||||||
|
|
||||||
|
expect(store.startStep).toHaveBeenLastCalledWith("FN-8490", 1, undefined);
|
||||||
|
expect(
|
||||||
|
store.updateStep.mock.calls.some(
|
||||||
|
([taskId, stepIndex, status]) => taskId === "FN-8490" && stepIndex === 0 && status === "done",
|
||||||
|
),
|
||||||
|
).toBe(false);
|
||||||
|
expect(store.moveTask).toHaveBeenCalledWith(
|
||||||
|
"FN-8490",
|
||||||
|
"todo",
|
||||||
|
expect.objectContaining({ preserveProgress: true, recoveryRehome: true }),
|
||||||
|
);
|
||||||
|
expect(onError).toHaveBeenCalledWith(
|
||||||
|
expect.objectContaining({ id: "FN-8490" }),
|
||||||
|
expect.objectContaining({ message: "Step 1: start rejected" }),
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
it("logs warning when rate-limit retry logEntry fails in step-session mode", async () => {
|
it("logs warning when rate-limit retry logEntry fails in step-session mode", async () => {
|
||||||
const warnSpy = vi.spyOn(executorLog, "warn");
|
const warnSpy = vi.spyOn(executorLog, "warn");
|
||||||
const store = createMockStore();
|
const store = createMockStore();
|
||||||
@@ -2790,6 +2870,30 @@ describe("fn_task_update bare-call guard (P1 api-contract)", () => {
|
|||||||
expect(text).toContain("skipped");
|
expect(text).toContain("skipped");
|
||||||
expect(store.appendAgentLog).not.toHaveBeenCalled();
|
expect(store.appendAgentLog).not.toHaveBeenCalled();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// FNXC:StepLifecycle 2026-07-22-09:50: Rejected starts must clearly preserve
|
||||||
|
// lifecycle invariants so agents do not execute work for a pending step.
|
||||||
|
it("explains that a rejected out-of-order start preserves lifecycle invariants", async () => {
|
||||||
|
const { store, tool } = makeTool();
|
||||||
|
store.getTask.mockResolvedValue(createMockTaskDetail({
|
||||||
|
steps: [
|
||||||
|
{ name: "Preflight", status: "in-progress" },
|
||||||
|
{ name: "Implement", status: "pending" },
|
||||||
|
],
|
||||||
|
}));
|
||||||
|
store.updateStep.mockResolvedValue(createMockTaskDetail({
|
||||||
|
steps: [
|
||||||
|
{ name: "Preflight", status: "in-progress" },
|
||||||
|
{ name: "Implement", status: "pending" },
|
||||||
|
],
|
||||||
|
}));
|
||||||
|
|
||||||
|
const result = await tool.execute("call-1", { step: 1, status: "in-progress" });
|
||||||
|
|
||||||
|
const text = result.content[0]?.type === "text" ? result.content[0].text : "";
|
||||||
|
expect(text).toContain("remains pending");
|
||||||
|
expect(text).toContain("ignored to preserve step lifecycle invariants");
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|||||||
@@ -654,6 +654,15 @@ export function createMockStore() {
|
|||||||
applyPatch(id, { steps });
|
applyPatch(id, { steps });
|
||||||
return { ...current, steps };
|
return { ...current, steps };
|
||||||
}),
|
}),
|
||||||
|
// FNXC:EngineTests 2026-07-22-10:30: Mirror the production atomic-start surface while retaining this helper's write-through projection behavior.
|
||||||
|
startStep: vi.fn(async (id: string, stepIndex: number, options?: { source?: "graph" }) => {
|
||||||
|
const task = await store.updateStep(id, stepIndex, "in-progress", options);
|
||||||
|
return {
|
||||||
|
task,
|
||||||
|
accepted: task.steps?.[stepIndex]?.status === "in-progress",
|
||||||
|
disposition: task.steps?.[stepIndex]?.status === "in-progress" ? "started" as const : "blocked" as const,
|
||||||
|
};
|
||||||
|
}),
|
||||||
getWorkflowStep: vi.fn().mockResolvedValue(undefined),
|
getWorkflowStep: vi.fn().mockResolvedValue(undefined),
|
||||||
listWorkflowSteps: vi.fn().mockResolvedValue([]),
|
listWorkflowSteps: vi.fn().mockResolvedValue([]),
|
||||||
setPluginWorkflowStepTemplates: vi.fn(),
|
setPluginWorkflowStepTemplates: vi.fn(),
|
||||||
|
|||||||
@@ -25,6 +25,11 @@ import {
|
|||||||
|
|
||||||
function makeStore() {
|
function makeStore() {
|
||||||
return {
|
return {
|
||||||
|
startStep: vi.fn().mockResolvedValue({
|
||||||
|
accepted: true,
|
||||||
|
disposition: "started",
|
||||||
|
task: makeTask([{ name: "Implement", status: "in-progress" }]),
|
||||||
|
}),
|
||||||
updateStep: vi.fn().mockResolvedValue(undefined),
|
updateStep: vi.fn().mockResolvedValue(undefined),
|
||||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||||
};
|
};
|
||||||
@@ -72,11 +77,9 @@ describe("runTaskStep", () => {
|
|||||||
// Baseline is captured BEFORE the step runs.
|
// Baseline is captured BEFORE the step runs.
|
||||||
expect(gitRevParse).toHaveBeenCalledWith("/wt");
|
expect(gitRevParse).toHaveBeenCalledWith("/wt");
|
||||||
expect(runStep).toHaveBeenCalledWith(0);
|
expect(runStep).toHaveBeenCalledWith(0);
|
||||||
// Projection ordering: in-progress before done.
|
// FNXC:StepLifecycle 2026-07-22-10:30: The accepted start projection must precede terminal completion.
|
||||||
expect(store.updateStep.mock.calls).toEqual([
|
expect(store.startStep).toHaveBeenCalledWith("FN-001", 0);
|
||||||
["FN-001", 0, "in-progress"],
|
expect(store.updateStep.mock.calls).toEqual([["FN-001", 0, "done"]]);
|
||||||
["FN-001", 0, "done"],
|
|
||||||
]);
|
|
||||||
});
|
});
|
||||||
|
|
||||||
it("passes graph projection source through in-progress and done writes", async () => {
|
it("passes graph projection source through in-progress and done writes", async () => {
|
||||||
@@ -96,8 +99,8 @@ describe("runTaskStep", () => {
|
|||||||
{ projectionSource: "graph" },
|
{ projectionSource: "graph" },
|
||||||
);
|
);
|
||||||
|
|
||||||
|
expect(store.startStep).toHaveBeenCalledWith("FN-001", 0, { source: "graph" });
|
||||||
expect(store.updateStep.mock.calls).toEqual([
|
expect(store.updateStep.mock.calls).toEqual([
|
||||||
["FN-001", 0, "in-progress", { source: "graph" }],
|
|
||||||
["FN-001", 0, "done", { source: "graph" }],
|
["FN-001", 0, "done", { source: "graph" }],
|
||||||
]);
|
]);
|
||||||
});
|
});
|
||||||
@@ -140,12 +143,74 @@ describe("runTaskStep", () => {
|
|||||||
);
|
);
|
||||||
|
|
||||||
expect(result).toEqual({ outcome: "failure", baselineSha: "baseSHA", checkpointId: "leaf" });
|
expect(result).toEqual({ outcome: "failure", baselineSha: "baseSHA", checkpointId: "leaf" });
|
||||||
// Only the in-progress write happened — the failed step is left non-done.
|
// FNXC:StepLifecycle 2026-07-22-10:30: A failed run preserves the accepted non-terminal start.
|
||||||
expect(store.updateStep.mock.calls).toEqual([["FN-001", 0, "in-progress"]]);
|
expect(store.startStep).toHaveBeenCalledWith("FN-001", 0);
|
||||||
|
expect(store.updateStep).not.toHaveBeenCalled();
|
||||||
expect(store.updateStep).not.toHaveBeenCalledWith("FN-001", 0, "done");
|
expect(store.updateStep).not.toHaveBeenCalledWith("FN-001", 0, "done");
|
||||||
expect(store.updateStep).not.toHaveBeenCalledWith("FN-001", 0, "skipped");
|
expect(store.updateStep).not.toHaveBeenCalledWith("FN-001", 0, "skipped");
|
||||||
});
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
* FNXC:StepLifecycle 2026-07-22-10:30:
|
||||||
|
* Exercise both sides of the atomic verdict: corrupted in-progress state must not run,
|
||||||
|
* while a dependency-valid in-progress restart remains resumable.
|
||||||
|
*/
|
||||||
|
it("does not execute a step whose atomic start is blocked despite an in-progress status", async () => {
|
||||||
|
const store = makeStore();
|
||||||
|
store.startStep.mockResolvedValue({
|
||||||
|
accepted: false,
|
||||||
|
disposition: "blocked",
|
||||||
|
blockingStepIndex: 0,
|
||||||
|
task: makeTask([
|
||||||
|
{ name: "Preflight", status: "in-progress" },
|
||||||
|
{ name: "Implement", status: "in-progress" },
|
||||||
|
]),
|
||||||
|
});
|
||||||
|
const runStep = vi.fn().mockResolvedValue({ success: true });
|
||||||
|
const gitRevParse = vi.fn().mockResolvedValue("baseSHA");
|
||||||
|
|
||||||
|
const result = await runTaskStep(
|
||||||
|
{ store, worktreePath: "/wt", runStep, gitRevParse },
|
||||||
|
makeTask([
|
||||||
|
{ name: "Preflight", status: "in-progress" },
|
||||||
|
{ name: "Implement", status: "in-progress" },
|
||||||
|
]),
|
||||||
|
1,
|
||||||
|
{ projectionSource: "graph" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(result).toEqual({ outcome: "failure" });
|
||||||
|
expect(runStep).not.toHaveBeenCalled();
|
||||||
|
expect(gitRevParse).not.toHaveBeenCalled();
|
||||||
|
expect(store.updateStep).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("executes a valid in-progress resume accepted by the atomic start", async () => {
|
||||||
|
const store = makeStore();
|
||||||
|
store.startStep.mockResolvedValue({
|
||||||
|
accepted: true,
|
||||||
|
disposition: "resumed",
|
||||||
|
task: makeTask([
|
||||||
|
{ name: "Preflight", status: "done" },
|
||||||
|
{ name: "Implement", status: "in-progress" },
|
||||||
|
]),
|
||||||
|
});
|
||||||
|
const runStep = vi.fn().mockResolvedValue({ success: false });
|
||||||
|
|
||||||
|
const result = await runTaskStep(
|
||||||
|
{ store, worktreePath: "/wt", runStep, gitRevParse: async () => "baseSHA" },
|
||||||
|
makeTask([
|
||||||
|
{ name: "Preflight", status: "done" },
|
||||||
|
{ name: "Implement", status: "in-progress" },
|
||||||
|
]),
|
||||||
|
1,
|
||||||
|
{ projectionSource: "graph" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(result.outcome).toBe("failure");
|
||||||
|
expect(runStep).toHaveBeenCalledWith(1);
|
||||||
|
});
|
||||||
|
|
||||||
it("still returns a result when baseline capture fails (best-effort)", async () => {
|
it("still returns a result when baseline capture fails (best-effort)", async () => {
|
||||||
const store = makeStore();
|
const store = makeStore();
|
||||||
const runStep = vi.fn().mockResolvedValue({ success: true });
|
const runStep = vi.fn().mockResolvedValue({ success: true });
|
||||||
|
|||||||
@@ -1525,6 +1525,63 @@ describe("StepSessionExecutor", () => {
|
|||||||
expect(onStepStart).toHaveBeenNthCalledWith(3, 2);
|
expect(onStepStart).toHaveBeenNthCalledWith(3, 2);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("does not create or complete a step session when the persisted start is rejected", async () => {
|
||||||
|
const prompt = makeStepPrompt("FN-8490", 1);
|
||||||
|
const task = makeTaskDetail({
|
||||||
|
id: "FN-8490",
|
||||||
|
prompt,
|
||||||
|
steps: [{ name: "Ordered step", status: "pending" }],
|
||||||
|
});
|
||||||
|
const onStepStart = vi.fn().mockResolvedValue(false);
|
||||||
|
const onStepComplete = vi.fn();
|
||||||
|
|
||||||
|
const executor = new StepSessionExecutor({
|
||||||
|
taskDetail: task,
|
||||||
|
worktreePath: "/project/.worktrees/main",
|
||||||
|
rootDir: "/project",
|
||||||
|
settings: makeSettings({ maxParallelSteps: 1 }),
|
||||||
|
onStepStart,
|
||||||
|
onStepComplete,
|
||||||
|
});
|
||||||
|
|
||||||
|
const results = await executor.executeAll();
|
||||||
|
|
||||||
|
expect(results).toEqual([
|
||||||
|
expect.objectContaining({
|
||||||
|
stepIndex: 0,
|
||||||
|
success: false,
|
||||||
|
error: expect.stringContaining("start was rejected"),
|
||||||
|
retries: 0,
|
||||||
|
}),
|
||||||
|
]);
|
||||||
|
expect(onStepStart).toHaveBeenCalledWith(0);
|
||||||
|
expect(mockedCreateFnAgent).not.toHaveBeenCalled();
|
||||||
|
expect(onStepComplete).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("preserves notification-only start callbacks that return void", async () => {
|
||||||
|
const task = makeTaskDetail({
|
||||||
|
prompt: makeStepPrompt("FN-8490-LEGACY", 1),
|
||||||
|
steps: [{ name: "Legacy callback step", status: "pending" }],
|
||||||
|
});
|
||||||
|
const session = makeMockSession();
|
||||||
|
mockedCreateFnAgent.mockResolvedValue({ session } as any);
|
||||||
|
const onStepStart = vi.fn(() => undefined);
|
||||||
|
|
||||||
|
const executor = new StepSessionExecutor({
|
||||||
|
taskDetail: task,
|
||||||
|
worktreePath: "/project/.worktrees/main",
|
||||||
|
rootDir: "/project",
|
||||||
|
settings: makeSettings({ maxParallelSteps: 1 }),
|
||||||
|
onStepStart,
|
||||||
|
});
|
||||||
|
|
||||||
|
const results = await executor.executeAll();
|
||||||
|
|
||||||
|
expect(results[0]?.success).toBe(true);
|
||||||
|
expect(mockedCreateFnAgent).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
|
|
||||||
it("skips live-terminal steps before starting resumed sessions", async () => {
|
it("skips live-terminal steps before starting resumed sessions", async () => {
|
||||||
const prompt = makeStepPrompt("FN-7248", 2);
|
const prompt = makeStepPrompt("FN-7248", 2);
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
|
|||||||
@@ -67,6 +67,20 @@ function makeFakeStore(steps: TaskStep[]) {
|
|||||||
return {
|
return {
|
||||||
trajectory,
|
trajectory,
|
||||||
steps,
|
steps,
|
||||||
|
startStep: async (
|
||||||
|
_id: string,
|
||||||
|
stepIndex: number,
|
||||||
|
options?: { source?: "graph" },
|
||||||
|
) => {
|
||||||
|
const priorStatus = steps[stepIndex]?.status;
|
||||||
|
trajectory.push({ step: stepIndex, status: "in-progress", ...(options?.source ? { source: options.source } : {}) });
|
||||||
|
if (steps[stepIndex]) steps[stepIndex] = { ...steps[stepIndex], status: "in-progress" };
|
||||||
|
return {
|
||||||
|
task: { steps } as never,
|
||||||
|
accepted: true,
|
||||||
|
disposition: priorStatus === "in-progress" ? "resumed" as const : "started" as const,
|
||||||
|
};
|
||||||
|
},
|
||||||
updateStep: async (
|
updateStep: async (
|
||||||
_id: string,
|
_id: string,
|
||||||
stepIndex: number,
|
stepIndex: number,
|
||||||
|
|||||||
@@ -11681,14 +11681,25 @@ export class TaskExecutor {
|
|||||||
sourceTaskId: task.id,
|
sourceTaskId: task.id,
|
||||||
sourceAgentId: stepIdentityAgent?.id,
|
sourceAgentId: stepIdentityAgent?.id,
|
||||||
taskEnv,
|
taskEnv,
|
||||||
onStepStart: (stepIndex) => {
|
// FNXC:StepLifecycle 2026-07-22-09:53: Await the dependency-aware store projection before session allocation so a rejected out-of-order start cannot execute while its persisted step remains pending.
|
||||||
this.options.stuckTaskDetector?.recordProgress(task.id);
|
onStepStart: async (stepIndex) => {
|
||||||
try {
|
try {
|
||||||
this.store.updateStep(task.id, stepIndex, "in-progress", stepProjectionOptions).catch((err) => {
|
const startResult = await this.store.startStep(
|
||||||
executorLog.warn(`${task.id}: failed to update step ${stepIndex} status to in-progress: ${err}`);
|
task.id,
|
||||||
});
|
stepIndex,
|
||||||
|
stepProjectionOptions,
|
||||||
|
);
|
||||||
|
if (!startResult.accepted) {
|
||||||
|
executorLog.warn(
|
||||||
|
`${task.id}: step ${stepIndex} start was rejected (${startResult.disposition}); persisted status is ` +
|
||||||
|
`${startResult.task.steps?.[stepIndex]?.status ?? "missing"}`,
|
||||||
|
);
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
this.options.stuckTaskDetector?.recordProgress(task.id);
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
executorLog.warn(`${task.id}: failed to update step ${stepIndex} status to in-progress: ${err}`);
|
executorLog.warn(`${task.id}: failed to update step ${stepIndex} status to in-progress: ${err}`);
|
||||||
|
return false;
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
onStepComplete: (stepIndex, result) => {
|
onStepComplete: (stepIndex, result) => {
|
||||||
@@ -14373,11 +14384,10 @@ export class TaskExecutor {
|
|||||||
RETHINK re-enters the implementation node instead of rewinding the live conversation.
|
RETHINK re-enters the implementation node instead of rewinding the live conversation.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
// If the persisted status doesn't match the requested status, the
|
// FNXC:StepLifecycle 2026-07-22-09:50: A persisted-status mismatch means
|
||||||
// store rejected the transition (currently: in-progress regression
|
// the store rejected the transition (for example, a completed-step
|
||||||
// on a done/skipped step). FN-5168 treats repeated rebuffs after loop
|
// regression or an out-of-order start/completion). FN-5168 treats
|
||||||
// recovery as a deterministic churn signal, but the agent-facing text
|
// repeated rebuffs after loop recovery as a deterministic churn signal.
|
||||||
// stays unchanged so the tool contract is preserved.
|
|
||||||
if (persistedStatus !== status) {
|
if (persistedStatus !== status) {
|
||||||
stuckDetector?.recordIgnoredStepUpdate(taskId);
|
stuckDetector?.recordIgnoredStepUpdate(taskId);
|
||||||
|
|
||||||
@@ -14394,7 +14404,7 @@ export class TaskExecutor {
|
|||||||
return {
|
return {
|
||||||
content: [{
|
content: [{
|
||||||
type: "text" as const,
|
type: "text" as const,
|
||||||
text: `Step ${step} (${stepInfo.name}) is already ${persistedStatus} — ${status} request ignored to preserve completed work. Progress: ${progress}/${task.steps.length} done.`,
|
text: `Step ${step} (${stepInfo.name}) remains ${persistedStatus} — ${status} request ignored to preserve step lifecycle invariants. Progress: ${progress}/${task.steps.length} done.`,
|
||||||
}],
|
}],
|
||||||
details: {},
|
details: {},
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -63,7 +63,7 @@ export type RunSingleStep = (stepIndex: number) => Promise<{ success: boolean; e
|
|||||||
/** Dependencies for {@link runTaskStep}. */
|
/** Dependencies for {@link runTaskStep}. */
|
||||||
export interface RunTaskStepDeps {
|
export interface RunTaskStepDeps {
|
||||||
/** Step-state projection sink (KTD-7). */
|
/** Step-state projection sink (KTD-7). */
|
||||||
store: Pick<TaskStore, "updateStep" | "logEntry">;
|
store: Pick<TaskStore, "startStep" | "updateStep" | "logEntry">;
|
||||||
/** Absolute path to the task's worktree (where `git rev-parse HEAD` runs). */
|
/** Absolute path to the task's worktree (where `git rev-parse HEAD` runs). */
|
||||||
worktreePath: string;
|
worktreePath: string;
|
||||||
/** Run exactly step `i` (step-session physics). */
|
/** Run exactly step `i` (step-session physics). */
|
||||||
@@ -135,17 +135,27 @@ export async function runTaskStep(
|
|||||||
const captureCheckpointId =
|
const captureCheckpointId =
|
||||||
deps.captureCheckpointId ?? (() => defaultCaptureCheckpointId(opts.sessionRef));
|
deps.captureCheckpointId ?? (() => defaultCaptureCheckpointId(opts.sessionRef));
|
||||||
|
|
||||||
// 1. Projection: step → in-progress (KTD-7). updateStep's own guards apply.
|
/*
|
||||||
|
* FNXC:StepLifecycle 2026-07-22-10:30:
|
||||||
|
* The atomic start verdict, rather than the returned status alone, distinguishes a valid
|
||||||
|
* in-progress restart resume from a blocked legacy-corruption state. Never run step work
|
||||||
|
* after the authoritative dependency guard rejects its projection.
|
||||||
|
*/
|
||||||
try {
|
try {
|
||||||
if (opts.projectionSource) {
|
const startResult = opts.projectionSource
|
||||||
await store.updateStep(task.id, stepIndex, "in-progress", { source: opts.projectionSource });
|
? await store.startStep(task.id, stepIndex, { source: opts.projectionSource })
|
||||||
} else {
|
: await store.startStep(task.id, stepIndex);
|
||||||
await store.updateStep(task.id, stepIndex, "in-progress");
|
if (!startResult.accepted) {
|
||||||
|
executorLog.warn(
|
||||||
|
`${task.id}: runTaskStep rejected step ${stepIndex} start (${startResult.disposition})`,
|
||||||
|
);
|
||||||
|
return { outcome: "failure" };
|
||||||
}
|
}
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
executorLog.warn(
|
executorLog.warn(
|
||||||
`${task.id}: runTaskStep failed to mark step ${stepIndex} in-progress: ${errMsg(err)}`,
|
`${task.id}: runTaskStep failed to mark step ${stepIndex} in-progress: ${errMsg(err)}`,
|
||||||
);
|
);
|
||||||
|
return { outcome: "failure" };
|
||||||
}
|
}
|
||||||
|
|
||||||
// 2. Baseline capture at instance start, before step work (KTD-2).
|
// 2. Baseline capture at instance start, before step work (KTD-2).
|
||||||
|
|||||||
@@ -122,8 +122,12 @@ export interface StepSessionExecutorOptions {
|
|||||||
runtimeHint?: string;
|
runtimeHint?: string;
|
||||||
/** Optional assigned-agent runtime config for model override precedence. */
|
/** Optional assigned-agent runtime config for model override precedence. */
|
||||||
assignedAgentRuntimeConfig?: Record<string, unknown>;
|
assignedAgentRuntimeConfig?: Record<string, unknown>;
|
||||||
/** Callback invoked when a step starts executing. */
|
/**
|
||||||
onStepStart?: (stepIndex: number) => void;
|
* FNXC:StepLifecycle 2026-07-22-09:53: This awaitable pre-start contract lets the
|
||||||
|
* authoritative lifecycle projection reject execution before session allocation;
|
||||||
|
* `void` preserves the legacy notification-only contract.
|
||||||
|
*/
|
||||||
|
onStepStart?: (stepIndex: number) => void | boolean | Promise<void | boolean>;
|
||||||
/** Callback invoked when a step completes (success or failure). */
|
/** Callback invoked when a step completes (success or failure). */
|
||||||
onStepComplete?: (stepIndex: number, result: StepResult) => void;
|
onStepComplete?: (stepIndex: number, result: StepResult) => void;
|
||||||
/** Optional skill selection context for session creation. */
|
/** Optional skill selection context for session creation. */
|
||||||
@@ -1227,8 +1231,23 @@ export class StepSessionExecutor {
|
|||||||
return { stepIndex, success: false, error: "Execution aborted", retries: 0 };
|
return { stepIndex, success: false, error: "Execution aborted", retries: 0 };
|
||||||
}
|
}
|
||||||
|
|
||||||
// Notify caller that this step is starting
|
// FNXC:StepLifecycle 2026-07-22-09:53: Project the start before allocating session
|
||||||
this.options.onStepStart?.(stepIndex);
|
// resources; a store-rejected out-of-order transition must not execute while pending.
|
||||||
|
if (this.options.onStepStart) {
|
||||||
|
try {
|
||||||
|
const startAccepted = await this.options.onStepStart(stepIndex);
|
||||||
|
if (startAccepted === false) {
|
||||||
|
const error = `Step ${stepIndex} start was rejected by the task lifecycle projection`;
|
||||||
|
stepExecLog.warn(`${error} for task ${taskDetail.id}`);
|
||||||
|
return { stepIndex, success: false, error, retries: 0 };
|
||||||
|
}
|
||||||
|
} catch (err) {
|
||||||
|
const reason = err instanceof Error ? err.message : String(err);
|
||||||
|
const error = `Step ${stepIndex} start projection failed: ${reason}`;
|
||||||
|
stepExecLog.warn(`${error} for task ${taskDetail.id}`);
|
||||||
|
return { stepIndex, success: false, error, retries: 0 };
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Build step prompt
|
// Build step prompt
|
||||||
const promptTaskDetail = this.consumeTaskDetailForStepPrompt();
|
const promptTaskDetail = this.consumeTaskDetailForStepPrompt();
|
||||||
|
|||||||
Reference in New Issue
Block a user