Skip to content
Merged
232 changes: 232 additions & 0 deletions src/workflow/executor/dag/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2072,6 +2072,238 @@ describe("DAGExecutor", () => {
// parked run, not a dead worker, and must not restart the branch.
assertEquals(first.nodeStates["gate"]!.status, "running");
});

it("spends no recovery on a queued node that a parked wait stops from starting", async () => {
// Queueing a node for recovery is not starting it. At maxConcurrency 1
// the wait ahead of it in the queue parks and ends the pass, so the
// interrupted step never runs. Charging it there burns its only recovery
// on work that never happened, and the next pass fails the whole run as
// out of budget even though the step still has not executed once.
const executed: string[] = [];
const persistedAttempts: number[] = [];
const exec = new DAGExecutor({
stepExecutor: new MockStepExecutor(new Map(), (node) => {
executed.push(node.id);
return { success: true, output: node.id, executionTime: 1 };
}),
onRecoveryScheduled: ({ nodeId, nodeStates }) => {
persistedAttempts.push(nodeStates[nodeId]?.attempt ?? 0);
},
maxConcurrency: 1,
});
const nodes: WorkflowNode[] = [
{
id: "gate",
dependsOn: [],
config: { type: "wait", waitType: "approval", message: "m" } as any,
},
{ id: "side-effect", dependsOn: [], config: { type: "step" } as any },
];

// Crash-recovery pass: the worker died with the step in flight, and the
// wait has not been raised yet.
const first = await exec.execute(
nodes,
createTestRun({
status: "running",
nodeStates: {
"side-effect": {
nodeId: "side-effect",
status: "running",
attempt: 1,
startedAt: new Date(),
},
},
}),
);

assertEquals(first.waiting, true);
assertEquals(executed, []);
assertEquals(persistedAttempts, []);
// #715's invariant still holds: the pass parks rather than reporting
// completion, and the unexecuted step is still named as running.
assertEquals(first.completed, false);
assertEquals(first.nodeStates["side-effect"]!.status, "running");
assertEquals(
first.nodeStates["side-effect"]!.attempt,
1,
"a queued node that never started must keep its recovery",
);

// Approval-resume pass: nothing is ahead of the step now, so it runs.
const second = await exec.execute(
nodes,
createTestRun({
status: "waiting",
nodeStates: {
...first.nodeStates,
gate: { ...first.nodeStates["gate"]!, status: "completed", completedAt: new Date() },
},
context: first.context,
}),
);

assertEquals(second.error, undefined);
assertEquals(second.completed, true);
assertEquals(executed, ["side-effect"]);
assertEquals(persistedAttempts, [2]);
assertEquals(second.nodeStates["side-effect"]!.status, "completed");
});

it("spends exactly one recovery on a node that does start, and bounds the next death", async () => {
// The other direction. Deferring the charge must not stop charging it: a
// node admitted to a batch has its raised attempt persisted before it
// executes, so a worker that dies again resumes out of budget instead of
// re-running the side effect forever.
const executed: string[] = [];
let persisted: Record<string, NodeState> | undefined;
const exec = new DAGExecutor({
stepExecutor: new MockStepExecutor(new Map(), (node) => {
executed.push(node.id);
return { success: true, output: node.id, executionTime: 1 };
}),
onRecoveryScheduled: ({ nodeStates }) => {
persisted = structuredClone(nodeStates);
},
});
const nodes: WorkflowNode[] = [
{ id: "side-effect", dependsOn: [], config: { type: "step" } as any },
];
const interrupted = (attempt: number): Record<string, NodeState> => ({
"side-effect": {
nodeId: "side-effect",
status: "running",
attempt,
startedAt: new Date(),
},
});

const first = await exec.execute(
nodes,
createTestRun({ status: "running", nodeStates: interrupted(1) }),
);

assertEquals(first.completed, true);
assertEquals(executed, ["side-effect"]);
// The durable write a second worker death resumes from: the node is
// still recorded running (nothing terminal was ever written for it) and
// its attempt is already raised.
assertEquals(persisted?.["side-effect"]?.status, "running");
assertEquals(persisted?.["side-effect"]?.attempt, 2);

// Resume from that write verbatim, exactly as a second worker would.
const second = await exec.execute(
nodes,
createTestRun({ status: "running", nodeStates: persisted! }),
);

assertEquals(second.completed, false);
assertEquals(
executed,
["side-effect"],
"the recovery budget must not stretch to a second run",
);
assertStringIncludes(second.error ?? "", "retry budget exhausted");
});

it("still spends only one recovery when the interrupted state has no startedAt", async () => {
// `startedAt` is optional on NodeState, and `WorkflowBackend` is an
// exported interface a project can implement, so a backend that does not
// round-trip the timestamp hands back a running node without it. The
// recovery budget must come from `attempt` alone: inferring "never
// started" from a missing timestamp gives such a node a second recovery
// and duplicates its side effect.
const executed: string[] = [];
let lastDurable: Record<string, NodeState> | undefined;
let atExecution: Record<string, NodeState> | undefined;
const build = () =>
new DAGExecutor({
stepExecutor: new MockStepExecutor(new Map(), (node) => {
executed.push(node.id);
atExecution = lastDurable;
return { success: true, output: node.id, executionTime: 1 };
}),
onRecoveryScheduled: ({ nodeStates }) => {
lastDurable = structuredClone(nodeStates);
},
onNodeStatesChanged: ({ nodeStates }) => {
lastDurable = structuredClone(nodeStates);
},
});
const nodes: WorkflowNode[] = [
{ id: "side-effect", dependsOn: [], config: { type: "step" } as any },
];

const first = await build().execute(
nodes,
createTestRun({
status: "running",
nodeStates: {
"side-effect": { nodeId: "side-effect", status: "running", attempt: 1 },
},
}),
);
assertEquals(first.completed, true);
assertEquals(executed, ["side-effect"], "the one recovery runs it once");

// Resume from the durable write in force while it was executing, exactly
// as a second worker would after this one died.
const second = await build().execute(
nodes,
createTestRun({ status: "running", nodeStates: atExecution! }),
);
assertEquals(
executed,
["side-effect"],
"a missing startedAt must not buy a second recovery",
);
assertEquals(second.completed, false);
assertStringIncludes(second.error ?? "", "retry budget exhausted");
});

it("bounds a child graph's recovery too, not just the top-level run", async () => {
// A composite runs its children against a synthetic run, and only the
// top-level run persists. The budget check must not ride on that: a
// child node out of budget has to be refused in the child graph, where
// nothing is written, or the composite re-runs a side effect that
// already exhausted its recoveries.
const executed: string[] = [];
const exec = new DAGExecutor({
stepExecutor: new MockStepExecutor(new Map(), (node) => {
executed.push(node.id);
return { success: true, output: node.id, executionTime: 1 };
}),
});
const nodes: WorkflowNode[] = [
{
id: "outer",
dependsOn: [],
config: {
type: "parallel",
nodes: [{ id: "inner", dependsOn: [], config: { type: "step" } }],
} as any,
},
];

const result = await exec.execute(
nodes,
createTestRun({
status: "running",
nodeStates: {
inner: {
nodeId: "inner",
status: "running",
attempt: 2,
startedAt: new Date(),
},
},
}),
);

assertEquals(executed, [], "an out-of-budget child must not be re-run");
assertEquals(result.completed, false);
assertStringIncludes(result.error ?? "", "retry budget exhausted");
});
});

describe("wait resume inside a nested composite", () => {
Expand Down
79 changes: 49 additions & 30 deletions src/workflow/executor/dag/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,23 @@ export class DAGExecutor {
};
}

// Only the top-level run has a row in the backend to write. Composites
// execute their children against synthetic runs (`${node.id}_parallel`,
// `_branch`, `_iter_N`) whose node states are a different keyspace: a loop
// iteration's run carries only that iteration's children. Persisting one
// of those under the real run id would replace the run's whole node-state
// map with an iteration-local fragment, stranding every completed
// top-level node as pending and re-running the workflow from the start --
// the duplicate side effect this recovery path exists to prevent.
// Child recoveries are persisted by the parent when it returns.
const isDurableRun = run.id === scope.rootRunId;
// Nodes queued for crash recovery that have not been admitted to a batch
// yet. Being queued spends nothing: an earlier node in the queue can park
// on a wait and end the pass before these ever start, and a node that
// never started must keep its one recovery for the next pass. The budget
// is charged at admission instead, below.
const recoveryQueued = new Set<string>();
Comment thread
kwakayama marked this conversation as resolved.

let ready = startFromNode ? [startFromNode] : getReadyNodes(inDegree, nodeStates);
if (!startFromNode) {
// A node recorded as running means different things depending on why the
Expand All @@ -151,16 +168,6 @@ export class DAGExecutor {
// charge every nested composite re-entry to the crash budget on an
// ordinary approval.
const { resumingWait } = scope;
// Only the top-level run has a row in the backend to write. Composites
// execute their children against synthetic runs (`${node.id}_parallel`,
// `_branch`, `_iter_N`) whose node states are a different keyspace: a loop
// iteration's run carries only that iteration's children. Persisting one
// of those under the real run id would replace the run's whole node-state
// map with an iteration-local fragment, stranding every completed
// top-level node as pending and re-running the workflow from the start --
// the duplicate side effect this recovery path exists to prevent.
// Child recoveries are persisted by the parent when it returns.
const isDurableRun = run.id === scope.rootRunId;
const exhausted: Array<{ nodeId: string; attempts: number; maxAttempts: number }> = [];
for (const [nodeId, degree] of inDegree) {
if (degree !== 0 || ready.includes(nodeId)) continue;
Expand Down Expand Up @@ -197,30 +204,17 @@ export class DAGExecutor {
//
// An interrupted attempt never finished, so it does not consume the
// node's retry budget outright -- a default node still gets recovered
// once. What it does consume is one recovery: the count below is
// written back before scheduling, and nothing overwrites it if the
// worker dies again, so the next resume sees a higher number.
// once. What it does consume is one recovery, charged when the node is
// admitted to a batch: the raised count is written back before the node
// executes, and nothing overwrites it if the worker dies again, so the
// next resume sees a higher number.
const maxAttempts = node.config.retry?.maxAttempts ?? 1;
const attempts = state.attempt ?? 0;
if (attempts > maxAttempts) {
exhausted.push({ nodeId, attempts, maxAttempts });
continue;
}
nodeStates[nodeId] = { ...state, attempt: attempts + 1 };
if (isDurableRun) {
const recovered = await this.config.onRecoveryScheduled?.({
runId: scope.rootRunId,
nodeId,
nodeStates: structuredClone(nodeStates),
ownership: scope.ownership,
});
if (recovered === false) {
throw ORCHESTRATION_ERROR.create({
detail:
`Cannot recover workflow node "${nodeId}" because execution ownership changed`,
});
}
}
recoveryQueued.add(nodeId);
ready.push(nodeId);
}

Expand Down Expand Up @@ -252,15 +246,26 @@ export class DAGExecutor {
const batch = ready.slice(0, this.config.maxConcurrency);
ready = ready.slice(this.config.maxConcurrency);

const isDurableRun = run.id === scope.rootRunId;
const batchStartedAt = new Date();
// Recoveries this batch charges. A node only reaches here once it is
// actually starting, so the recovery it spends is one it gets to use.
const recovered: string[] = [];
for (const nodeId of batch) {
const existing = nodeStates[nodeId];
// A node already recorded running keeps its attempt: it is being
// re-entered, not restarted (a composite resuming a parked wait). A
// node admitted off the recovery queue is the exception -- an
// interrupted attempt is being replaced, and raising the count here is
// what bounds repeated worker deaths.
const isRecovery = recoveryQueued.delete(nodeId);
if (isRecovery) recovered.push(nodeId);
const runningState: NodeState = {
...existing,
nodeId,
status: "running",
attempt: existing?.status === "running" ? existing.attempt : (existing?.attempt ?? 0) + 1,
attempt: existing?.status === "running" && !isRecovery
? existing.attempt
: (existing?.attempt ?? 0) + 1,
startedAt: existing?.status === "running" && existing.startedAt
? existing.startedAt
: batchStartedAt,
Expand All @@ -271,6 +276,20 @@ export class DAGExecutor {
nodeStates[nodeId] = runningState;
}
if (isDurableRun) {
for (const nodeId of recovered) {
const persisted = await this.config.onRecoveryScheduled?.({
runId: scope.rootRunId,
nodeId,
nodeStates: structuredClone(nodeStates),
ownership: scope.ownership,
});
if (persisted === false) {
throw ORCHESTRATION_ERROR.create({
detail:
`Cannot recover workflow node "${nodeId}" because execution ownership changed`,
});
}
}
await this.publishNodeStates(scope, nodeStates, context, batch);
abortSignal?.throwIfAborted();
Comment thread
kwakayama marked this conversation as resolved.
}
Expand Down
4 changes: 4 additions & 0 deletions src/workflow/executor/workflow-executor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -462,6 +462,10 @@ describe("workflow/executor/workflow-executor", () => {
assertEquals(persistedWhileRunning.workerId, "run-execution:old-owner");
assertEquals(persistedWhileRunning.nodeStates["side-effect"]?.attempt, 2);
assertEquals(persistedWhileRunning.nodeStates["side-effect"]?.status, "running");
assertExists(
persistedWhileRunning.nodeStates["side-effect"]?.startedAt,
"the recovery is persisted at admission, so the durable write carries its start",
);

release.resolve();
await execution;
Expand Down