Skip to content

Commit 844d181

Browse files
authored
fix(executor): fail a stop-after run whose routing skips the stop block (#8635)
* fix(executor): fail a stop-after run whose routing skips the stop block A run with stopAfterBlockId only stopped when the stop block completed. When a router, condition, or untaken error path routed the run away from it, the stop never triggered and the run finished every other branch, reporting success as if it had stopped there. A static check before the run cannot see this. - The engine ends the run as soon as every path into the stop block has been deactivated, before any further block starts, and fails it with `Stop block "<name>" (<id>) was not reached: no path this run took leads to it`. - Any run that ends without completing its stop block fails the same way: a stop block missing from the executed graph, or a Response block that ended the run first. - A loop or parallel stop with nothing to run completes at its start sentinel, whose end sentinel never runs, so that exit now counts as reaching it. - The v2 contract and the CLI `--stop-after` help describe the failure. - E2E: a condition fixture checks the stop on the taken branch still stops there, a stop on the skipped branch fails the run before the other branch's slow block finishes, and the CLI exits non-zero. * fix(executor): a skipped stop block fails a run another branch paused, and names a Response ending - A run whose stop block was proven unreachable fails even when another branch paused, instead of returning a paused run that would resume past it. - When a Response block ended the run first, the error says so rather than claiming no path leads to the stop block.
1 parent 6b3cd82 commit 844d181

9 files changed

Lines changed: 533 additions & 30 deletions

File tree

‎apps/docs/content/docs/cli/reference.mdx‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6704,7 +6704,7 @@ sim workflows run <workflowId> [options]
67046704
| `--mock-payload` | No | Use the selected trigger's server-derived mock payload; runs the current saved workflow state (implies --manual). |
67056705
| `--from-block <blockId>` | No | Run manually from this saved workflow block. |
67066706
| `--source-run <runId>` | No | Prior run whose persisted state supplies upstream outputs (requires --from-block). |
6707-
| `--stop-after <blockId>` | No | Stop the run after this saved block; with --from-block on the same block, re-runs only that block (implies --manual). |
6707+
| `--stop-after <blockId>` | No | Stop the run after this saved block, failing it if the run takes a path that skips the block; with --from-block on the same block, re-runs only that block (implies --manual). |
67086708
| `--follow` | No | Stream the run as it happens; progress on stderr, result on stdout. The stream reports only success and output, so the result omits the run id and timings a non-streaming run returns. |
67096709
| `--include-thinking` | No | Show model reasoning while following (requires --follow). |
67106710
| `--include-tool-calls` | No | Show tool calls while following (requires --follow). |

‎apps/docs/content/docs/cli/workflows.mdx‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -642,7 +642,7 @@ sim workflows run <workflowId> [options]
642642
| `--mock-payload` | No | Use the selected trigger's server-derived mock payload; runs the current saved workflow state (implies --manual). |
643643
| `--from-block <blockId>` | No | Run manually from this saved workflow block. |
644644
| `--source-run <runId>` | No | Prior run whose persisted state supplies upstream outputs (requires --from-block). |
645-
| `--stop-after <blockId>` | No | Stop the run after this saved block; with --from-block on the same block, re-runs only that block (implies --manual). |
645+
| `--stop-after <blockId>` | No | Stop the run after this saved block, failing it if the run takes a path that skips the block; with --from-block on the same block, re-runs only that block (implies --manual). |
646646
| `--follow` | No | Stream the run as it happens; progress on stderr, result on stdout. The stream reports only success and output, so the result omits the run id and timings a non-streaming run returns. |
647647
| `--include-thinking` | No | Show model reasoning while following (requires --follow). |
648648
| `--include-tool-calls` | No | Show tool calls while following (requires --follow). |

‎apps/docs/openapi-v2-workflows.json‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12617,7 +12617,7 @@
1261712617
]
1261812618
},
1261912619
"stopAfterBlockId": {
12620-
"description": "Saved workflow block after which the run stops; downstream blocks do not execute. Must not be inside a loop or parallel. With a block entry naming the same block, re-runs only that block against the source run.",
12620+
"description": "Saved workflow block after which the run stops; downstream blocks do not execute. Must not be inside a loop or parallel. If a router, condition, or untaken error path routes the run away from the block, the run fails as soon as that is decided, without running the other branches. With a block entry naming the same block, re-runs only that block against the source run.",
1262112621
"type": "string",
1262212622
"minLength": 1
1262312623
}

‎apps/sim/executor/execution/edge-manager.ts‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,17 @@ export class EdgeManager {
122122
return node.incomingEdges.size === 0 || this.countActiveIncomingEdges(node) === 0
123123
}
124124

125+
/**
126+
* Whether a node that has not been queued can still run: it has received an activated edge, or
127+
* an incoming edge is still undecided. False means every path into it was deactivated (a router
128+
* or condition chose another route, or an error path was not taken), so nothing will queue it.
129+
*/
130+
canNodeStillRun(nodeId: string): boolean {
131+
if (this.nodesWithActivatedEdge.has(nodeId)) return true
132+
const node = this.dag.nodes.get(nodeId)
133+
return node !== undefined && this.countActiveIncomingEdges(node) > 0
134+
}
135+
125136
restoreIncomingEdge(targetNodeId: string, sourceNodeId: string): void {
126137
const targetNode = this.dag.nodes.get(targetNodeId)
127138
if (!targetNode) {

‎apps/sim/executor/execution/engine.test.ts‎

Lines changed: 313 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -21,13 +21,17 @@ vi.mock('@/lib/execution/cancellation', () => ({
2121
},
2222
}))
2323

24-
import { EDGE } from '@/executor/constants'
25-
import type { DAG, DAGNode } from '@/executor/dag/builder'
26-
import type { EdgeManager } from '@/executor/execution/edge-manager'
24+
import { BlockType, EDGE } from '@/executor/constants'
25+
import { type DAG, DAGBuilder, type DAGNode } from '@/executor/dag/builder'
26+
import { EdgeManager } from '@/executor/execution/edge-manager'
2727
import type { NodeExecutionOrchestrator } from '@/executor/orchestrators/node'
2828
import type { ExecutionContext, ExecutionResult } from '@/executor/types'
2929
import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
30-
import type { SerializedBlock } from '@/serializer/types'
30+
import {
31+
buildLoopSentinelEndId,
32+
buildLoopSentinelStartId,
33+
} from '@/executor/utils/subflow-node-id-codec'
34+
import type { SerializedBlock, SerializedWorkflow } from '@/serializer/types'
3135
import { ExecutionEngine } from './engine'
3236

3337
const executionEngineLoggerCallIndex = loggerMock.createLogger.mock.calls.findIndex(
@@ -1088,4 +1092,309 @@ describe('ExecutionEngine', () => {
10881092
expect(result.output).toEqual({ data: { response: true }, status: 200, headers: {} })
10891093
})
10901094
})
1095+
/**
1096+
* A stop-after run is a promise that the run ends with the stop block. These run the real DAG
1097+
* builder and edge manager, so a router or condition deciding a path is the same decision the
1098+
* engine makes in production.
1099+
*/
1100+
describe('Stop-after block', () => {
1101+
function block(id: string, type = BlockType.FUNCTION): SerializedBlock {
1102+
return { ...createMockBlock(id), metadata: { id: type, name: id } }
1103+
}
1104+
1105+
function buildRun(
1106+
workflow: SerializedWorkflow,
1107+
stopAfterBlockId: string,
1108+
outputs: Record<string, ExecutionResult['output']> = {}
1109+
) {
1110+
const dag = new DAGBuilder().build(workflow, { triggerBlockId: 'start' })
1111+
const executed: string[] = []
1112+
const nodeOrchestrator = {
1113+
executeNode: vi.fn(async (_ctx: ExecutionContext, nodeId: string) => {
1114+
executed.push(nodeId)
1115+
return { nodeId, output: outputs[nodeId] ?? {}, isFinalOutput: false }
1116+
}),
1117+
handleNodeCompletion: vi.fn(),
1118+
} as unknown as NodeExecutionOrchestrator
1119+
const engine = new ExecutionEngine(
1120+
createMockContext({
1121+
stopAfterBlockId,
1122+
decisions: { router: new Map(), condition: new Map() },
1123+
}),
1124+
dag,
1125+
new EdgeManager(dag),
1126+
nodeOrchestrator
1127+
)
1128+
return { engine, executed }
1129+
}
1130+
1131+
/** start → condition; `if` → taken → takenTail, `else` → stop → after. */
1132+
const conditionWorkflow: SerializedWorkflow = {
1133+
version: '1',
1134+
blocks: [
1135+
block('start', BlockType.STARTER),
1136+
block('condition', BlockType.CONDITION),
1137+
block('taken'),
1138+
block('takenTail'),
1139+
block('stop'),
1140+
block('after'),
1141+
],
1142+
connections: [
1143+
{ source: 'start', target: 'condition' },
1144+
{ source: 'condition', target: 'taken', sourceHandle: 'condition-if' },
1145+
{ source: 'condition', target: 'stop', sourceHandle: 'condition-else' },
1146+
{ source: 'taken', target: 'takenTail' },
1147+
{ source: 'stop', target: 'after' },
1148+
],
1149+
loops: {},
1150+
parallels: {},
1151+
}
1152+
1153+
it('fails as soon as a condition routes away from the stop block', async () => {
1154+
const { engine, executed } = buildRun(conditionWorkflow, 'stop', {
1155+
condition: { selectedOption: 'if' },
1156+
})
1157+
1158+
await expect(engine.run('start')).rejects.toThrow('Stop block "stop" (stop) was not reached')
1159+
expect(executed).toEqual(['start', 'condition'])
1160+
})
1161+
1162+
it('fails as soon as a router routes away from the stop block', async () => {
1163+
const { engine, executed } = buildRun(
1164+
{
1165+
version: '1',
1166+
blocks: [
1167+
block('start', BlockType.STARTER),
1168+
block('router', BlockType.ROUTER_V2),
1169+
block('taken'),
1170+
block('takenTail'),
1171+
block('stop'),
1172+
],
1173+
connections: [
1174+
{ source: 'start', target: 'router' },
1175+
{ source: 'router', target: 'taken', sourceHandle: 'router-route-a' },
1176+
{ source: 'router', target: 'stop', sourceHandle: 'router-route-b' },
1177+
{ source: 'taken', target: 'takenTail' },
1178+
],
1179+
loops: {},
1180+
parallels: {},
1181+
},
1182+
'stop',
1183+
{ router: { selectedRoute: 'route-a' } }
1184+
)
1185+
1186+
await expect(engine.run('start')).rejects.toThrow('Stop block "stop" (stop) was not reached')
1187+
expect(executed).toEqual(['start', 'router'])
1188+
})
1189+
1190+
it('fails rather than pausing when another branch pauses after the stop block is skipped', async () => {
1191+
const { engine, executed } = buildRun(
1192+
{
1193+
version: '1',
1194+
blocks: [
1195+
block('start', BlockType.STARTER),
1196+
block('approval'),
1197+
block('condition', BlockType.CONDITION),
1198+
block('taken'),
1199+
block('stop'),
1200+
],
1201+
connections: [
1202+
{ source: 'start', target: 'approval' },
1203+
{ source: 'start', target: 'condition' },
1204+
{ source: 'condition', target: 'taken', sourceHandle: 'condition-if' },
1205+
{ source: 'condition', target: 'stop', sourceHandle: 'condition-else' },
1206+
],
1207+
loops: {},
1208+
parallels: {},
1209+
},
1210+
'stop',
1211+
{
1212+
approval: {
1213+
response: { status: 'paused' },
1214+
_pauseMetadata: {
1215+
contextId: 'pause-1',
1216+
blockId: 'approval',
1217+
response: { status: 'paused' },
1218+
timestamp: new Date().toISOString(),
1219+
pauseKind: 'hitl',
1220+
},
1221+
},
1222+
condition: { selectedOption: 'if' },
1223+
}
1224+
)
1225+
1226+
await expect(engine.run('start')).rejects.toThrow('Stop block "stop" (stop) was not reached')
1227+
expect(executed).toEqual(['start', 'approval', 'condition'])
1228+
})
1229+
1230+
it('fails when the stop block sits on an error path the run never takes', async () => {
1231+
const { engine, executed } = buildRun(
1232+
{
1233+
version: '1',
1234+
blocks: [block('start', BlockType.STARTER), block('work'), block('next'), block('stop')],
1235+
connections: [
1236+
{ source: 'start', target: 'work' },
1237+
{ source: 'work', target: 'next', sourceHandle: 'source' },
1238+
{ source: 'work', target: 'stop', sourceHandle: 'error' },
1239+
],
1240+
loops: {},
1241+
parallels: {},
1242+
},
1243+
'stop'
1244+
)
1245+
1246+
await expect(engine.run('start')).rejects.toThrow('Stop block "stop" (stop) was not reached')
1247+
expect(executed).toEqual(['start', 'work'])
1248+
})
1249+
1250+
it('stops after the stop block when the condition routes to it', async () => {
1251+
const { engine, executed } = buildRun(conditionWorkflow, 'stop', {
1252+
condition: { selectedOption: 'else' },
1253+
})
1254+
1255+
const result = await engine.run('start')
1256+
1257+
expect(result.success).toBe(true)
1258+
expect(executed).toEqual(['start', 'condition', 'stop'])
1259+
})
1260+
1261+
it('runs a join reached through one taken and one skipped branch', async () => {
1262+
const { engine, executed } = buildRun(
1263+
{
1264+
version: '1',
1265+
blocks: [
1266+
block('start', BlockType.STARTER),
1267+
block('condition', BlockType.CONDITION),
1268+
block('taken'),
1269+
block('skipped'),
1270+
block('stop'),
1271+
block('after'),
1272+
],
1273+
connections: [
1274+
{ source: 'start', target: 'condition' },
1275+
{ source: 'condition', target: 'taken', sourceHandle: 'condition-if' },
1276+
{ source: 'condition', target: 'skipped', sourceHandle: 'condition-else' },
1277+
{ source: 'taken', target: 'stop' },
1278+
{ source: 'skipped', target: 'stop' },
1279+
{ source: 'stop', target: 'after' },
1280+
],
1281+
loops: {},
1282+
parallels: {},
1283+
},
1284+
'stop',
1285+
{ condition: { selectedOption: 'if' } }
1286+
)
1287+
1288+
const result = await engine.run('start')
1289+
1290+
expect(result.success).toBe(true)
1291+
expect(executed).toEqual(['start', 'condition', 'taken', 'stop'])
1292+
})
1293+
1294+
it('runs a loop stop block queued by a dead end inside the loop', async () => {
1295+
const sentinelStart = buildLoopSentinelStartId('loop')
1296+
const sentinelEnd = buildLoopSentinelEndId('loop')
1297+
const { engine, executed } = buildRun(
1298+
{
1299+
version: '1',
1300+
blocks: [
1301+
block('start', BlockType.STARTER),
1302+
block('loop', BlockType.LOOP),
1303+
block('condition', BlockType.CONDITION),
1304+
block('inner'),
1305+
block('after'),
1306+
],
1307+
connections: [
1308+
{ source: 'start', target: 'loop' },
1309+
{ source: 'loop', target: 'condition', sourceHandle: 'loop-start-source' },
1310+
{ source: 'condition', target: 'inner', sourceHandle: 'condition-if' },
1311+
{ source: 'loop', target: 'after', sourceHandle: 'loop-end-source' },
1312+
],
1313+
loops: { loop: { id: 'loop', nodes: ['condition', 'inner'], iterations: 1 } },
1314+
parallels: {},
1315+
},
1316+
sentinelEnd,
1317+
{ condition: { selectedOption: 'else' } }
1318+
)
1319+
1320+
const result = await engine.run('start')
1321+
1322+
expect(result.success).toBe(true)
1323+
expect(executed).toEqual(['start', sentinelStart, 'condition', sentinelEnd])
1324+
})
1325+
1326+
it('stops after a loop stop block whose loop has nothing to run', async () => {
1327+
const sentinelStart = buildLoopSentinelStartId('loop')
1328+
const { engine, executed } = buildRun(
1329+
{
1330+
version: '1',
1331+
blocks: [
1332+
block('start', BlockType.STARTER),
1333+
block('loop', BlockType.LOOP),
1334+
block('inner'),
1335+
block('after'),
1336+
],
1337+
connections: [
1338+
{ source: 'start', target: 'loop' },
1339+
{ source: 'loop', target: 'inner', sourceHandle: 'loop-start-source' },
1340+
{ source: 'loop', target: 'after', sourceHandle: 'loop-end-source' },
1341+
],
1342+
loops: { loop: { id: 'loop', nodes: ['inner'], iterations: 0 } },
1343+
parallels: {},
1344+
},
1345+
buildLoopSentinelEndId('loop'),
1346+
{
1347+
[sentinelStart]: { sentinelStart: true, shouldExit: true, selectedRoute: EDGE.LOOP_EXIT },
1348+
}
1349+
)
1350+
1351+
const result = await engine.run('start')
1352+
1353+
expect(result.success).toBe(true)
1354+
expect(executed).toEqual(['start', sentinelStart])
1355+
})
1356+
1357+
it('succeeds when the stop block is a Response block, and fails when one ends the run first', async () => {
1358+
const workflow: SerializedWorkflow = {
1359+
version: '1',
1360+
blocks: [
1361+
block('start', BlockType.STARTER),
1362+
block('respond', BlockType.RESPONSE),
1363+
block('stop'),
1364+
],
1365+
connections: [
1366+
{ source: 'start', target: 'respond' },
1367+
{ source: 'respond', target: 'stop' },
1368+
],
1369+
loops: {},
1370+
parallels: {},
1371+
}
1372+
1373+
const atResponse = buildRun(workflow, 'respond')
1374+
expect((await atResponse.engine.run('start')).success).toBe(true)
1375+
1376+
const pastResponse = buildRun(workflow, 'stop')
1377+
await expect(pastResponse.engine.run('start')).rejects.toThrow(
1378+
'Stop block "stop" (stop) was not reached: a Response block ended the run first'
1379+
)
1380+
expect(pastResponse.executed).toEqual(['start', 'respond'])
1381+
})
1382+
1383+
it('stops after the entry block when the stop block is the entry', async () => {
1384+
const { engine, executed } = buildRun(conditionWorkflow, 'start')
1385+
1386+
const result = await engine.run('start')
1387+
1388+
expect(result.success).toBe(true)
1389+
expect(executed).toEqual(['start'])
1390+
})
1391+
1392+
it('fails when the stop block is not in the graph the run executes', async () => {
1393+
const { engine } = buildRun(conditionWorkflow, 'missing', {
1394+
condition: { selectedOption: 'if' },
1395+
})
1396+
1397+
await expect(engine.run('start')).rejects.toThrow('Stop block missing was not reached')
1398+
})
1399+
})
10911400
})

0 commit comments

Comments
 (0)