Skip to content

Commit 807523e

Browse files
waleedlatif1Waleed Latif
andauthored
feat(workflows): stop a manual v2 run after a block, and workflows run --stop-after (#8622)
* feat(workflows): stop a manual v2 run after a block, and `workflows run --stop-after` A manual v2 run can now name `run.stopAfterBlockId`; the run stops once that block completes and downstream blocks do not execute. Combined with a block entry on the same block, it re-runs exactly one block against a prior run's persisted upstream outputs, server-side: sim workflows run W --from-block X --source-run R --stop-after X --select-output X.result Agents verifying an edit no longer re-run every upstream block (often a slow LLM or API call) or toggle blocks off to skip them. - Contract: optional `stopAfterBlockId` on the manual run selection. - Application: both manual operations refuse a block missing from the saved workflow or nested in a loop/parallel (the engine would otherwise run to the end or stop after one iteration), before anything runs. - Execute service: threads the trusted value to the sync and stream paths. - CLI: `--stop-after <blockId>` implies --manual and rejects --async. - E2E: test-workflow-stop-after-e2e.ts against a running app; the http-e2e job gains a Redis service because hosted billing admits runs through a Redis usage reservation. * fix(workflows): refuse stop targets a run cannot reach, and run the E2E self-hosted - The manual operations refuse a stop block the run cannot reach from its entry (an upstream block would let the run finish everything after the entry), and look blocks up as own properties. - The executor fails a run whose stop block is absent from the workflow it executes, instead of running everything; this closes the window between validation and the executor's own draft load, for every caller. - The CLI refuses an empty --stop-after rather than dropping it. - CI: the stop-after E2E gets its own self-hosted app step; the SCIM suite asserts PostgreSQL rate-limit storage, so Redis is not added to that app. Fixture cleanup waits for run logs to finalize before deleting. * fix(workflows): refuse a disabled stop block, or one reached only through one The executor omits disabled blocks from its graph, so a disabled stop target, or one whose only path runs through a disabled block, is never reached and the run would finish everything after the entry. * fix(workflows): the executor refuses a disabled stop block too The serialized workflow keeps disabled blocks, but the DAG skips them, so a disabled stop target would never be reached. --------- Co-authored-by: Waleed Latif <waleed@sim.ai>
1 parent dc10d4c commit 807523e

16 files changed

Lines changed: 769 additions & 15 deletions

File tree

‎.github/workflows/test-build.yml‎

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -253,6 +253,51 @@ jobs:
253253
VERSION_COMPARE_E2E_REPORT_PATH="$report_dir/version-compare-http-report.json" \
254254
bun run test:workflow-version-compare:e2e
255255
256+
# A self-hosted app: hosted billing admits a run only through a Redis usage
257+
# reservation, and the SCIM suite above asserts PostgreSQL rate-limit storage,
258+
# so workflow execution gets its own app rather than adding Redis to that one.
259+
- name: Verify single-block workflow runs over real HTTP
260+
working-directory: apps/sim
261+
env:
262+
NEXT_PUBLIC_APP_URL: http://127.0.0.1:3018
263+
BETTER_AUTH_URL: http://127.0.0.1:3018
264+
NEXT_PUBLIC_FORCE_HOSTED: 'false'
265+
INTERNAL_API_SECRET: stop-after-http-ci-local-secret-at-least-32-characters
266+
DB_TX_TRIPWIRE: throw
267+
DISABLE_TELEMETRY: 'true'
268+
NEXT_TELEMETRY_DISABLED: '1'
269+
NEXT_PUBLIC_CHAT_DISABLED: 'true'
270+
READY_TIMEOUT_SECONDS: 300
271+
run: |
272+
report_dir="$RUNNER_TEMP/e2e"
273+
server_log="$report_dir/stop-after-next.log"
274+
mkdir -p "$report_dir"
275+
node ../../node_modules/next/dist/bin/next dev --hostname 127.0.0.1 --port 3018 > "$server_log" 2>&1 &
276+
server_pid=$!
277+
finish() {
278+
kill "$server_pid" 2>/dev/null || true
279+
wait "$server_pid" 2>/dev/null || true
280+
awk '/^ (GET|POST|PUT|PATCH|DELETE|HEAD) \/api\// { print }' "$server_log" > "$report_dir/stop-after-http-status.log"
281+
}
282+
trap finish EXIT
283+
fail_startup() {
284+
echo "::error::$1"
285+
tail -n 200 "$server_log"
286+
exit 1
287+
}
288+
started=$SECONDS
289+
until curl --fail --silent --max-time 10 http://127.0.0.1:3018/api/health > /dev/null; do
290+
kill -0 "$server_pid" 2>/dev/null || fail_startup 'Local workflow app exited during startup.'
291+
[ $((SECONDS - started)) -lt "$READY_TIMEOUT_SECONDS" ] ||
292+
fail_startup "Local workflow app did not become ready within $READY_TIMEOUT_SECONDS seconds."
293+
sleep 2
294+
done
295+
echo "Local workflow app ready after $((SECONDS - started))s"
296+
STOP_AFTER_E2E_BASE_URL="$NEXT_PUBLIC_APP_URL" \
297+
STOP_AFTER_E2E_DATABASE_URL="$DATABASE_URL" \
298+
STOP_AFTER_E2E_REPORT_PATH="$report_dir/stop-after-http-report.json" \
299+
bun run test:workflow-stop-after:e2e
300+
256301
- name: Upload end-to-end reports and server logs
257302
if: failure()
258303
uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6704,6 +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). |
67076708
| `--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. |
67086709
| `--include-thinking` | No | Show model reasoning while following (requires --follow). |
67096710
| `--include-tool-calls` | No | Show tool calls while following (requires --follow). |

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -642,6 +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). |
645646
| `--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. |
646647
| `--include-thinking` | No | Show model reasoning while following (requires --follow). |
647648
| `--include-tool-calls` | No | Show tool calls while following (requires --follow). |

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

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12615,6 +12615,11 @@
1261512615
"additionalProperties": false
1261612616
}
1261712617
]
12618+
},
12619+
"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.",
12621+
"type": "string",
12622+
"minLength": 1
1261812623
}
1261912624
},
1262012625
"required": ["source"],
@@ -12703,6 +12708,17 @@
1270312708
"sourceRunId": "run_123"
1270412709
}
1270512710
}
12711+
},
12712+
{
12713+
"run": {
12714+
"source": "manual",
12715+
"entry": {
12716+
"type": "block",
12717+
"blockId": "block_123",
12718+
"sourceRunId": "run_123"
12719+
},
12720+
"stopAfterBlockId": "block_123"
12721+
}
1270612722
}
1270712723
]
1270812724
},

‎apps/sim/app/api/v2/workflows/[workflowId]/execute/route.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -448,6 +448,7 @@ export const POST = withRouteHandler(
448448
mode: body.stream ? 'stream' : resultStream ? 'sync-result-stream' : 'sync',
449449
blockId: manualRun.entry.blockId,
450450
sourceRunId: manualRun.entry.sourceRunId,
451+
stopAfterBlockId: manualRun.stopAfterBlockId,
451452
},
452453
request: req,
453454
})
@@ -460,6 +461,7 @@ export const POST = withRouteHandler(
460461
mode: body.stream ? 'stream' : resultStream ? 'sync-result-stream' : 'sync',
461462
triggerBlockId: manualRun.entry?.blockId,
462463
useMockPayload: manualRun.entry?.useMockPayload === true,
464+
stopAfterBlockId: manualRun.stopAfterBlockId,
463465
},
464466
request: req,
465467
})

‎apps/sim/lib/api/contracts/v2/workflows.ts‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1297,6 +1297,13 @@ export const v2WorkflowRunSelectionSchema = z.discriminatedUnion('source', [
12971297
.describe(
12981298
'Manual entry mode. Omit to enter through the workflow trigger; a block entry requires an exact source run.'
12991299
),
1300+
stopAfterBlockId: z
1301+
.string()
1302+
.min(1, 'run.stopAfterBlockId cannot be empty')
1303+
.optional()
1304+
.describe(
1305+
'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.'
1306+
),
13001307
})
13011308
.strict(),
13021309
])
@@ -1412,6 +1419,13 @@ export const v2ExecuteWorkflowBodySchema = z
14121419
entry: { type: 'block', blockId: 'block_123', sourceRunId: 'run_123' },
14131420
},
14141421
},
1422+
{
1423+
run: {
1424+
source: 'manual',
1425+
entry: { type: 'block', blockId: 'block_123', sourceRunId: 'run_123' },
1426+
stopAfterBlockId: 'block_123',
1427+
},
1428+
},
14151429
],
14161430
})
14171431
export type V2ExecuteWorkflowBody = z.input<typeof v2ExecuteWorkflowBodySchema>

‎apps/sim/lib/workflows/application/execute-manual-workflow.test.ts‎

Lines changed: 118 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -300,6 +300,124 @@ describe('manual workflow execution application operations', () => {
300300
expect(mocks.loadSourceState).not.toHaveBeenCalled()
301301
})
302302

303+
it('rejects a stop block missing from the saved workflow before anything runs', async () => {
304+
await expect(
305+
executeManualWorkflowOperation.execute({
306+
principal,
307+
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'missing' },
308+
})
309+
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') })
310+
await expect(
311+
executeManualWorkflowFromBlockOperation.execute({
312+
principal,
313+
input: {
314+
...baseInput,
315+
blockId: 'agent-1',
316+
sourceRunId: 'source-run-1',
317+
stopAfterBlockId: 'missing',
318+
},
319+
})
320+
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') })
321+
expect(mocks.loadSourceState).not.toHaveBeenCalled()
322+
expect(mocks.executeService).not.toHaveBeenCalled()
323+
})
324+
325+
it('rejects a stop block nested in a loop or parallel, which the engine cannot stop on', async () => {
326+
mockLoadManualState.mockResolvedValue({
327+
blocks: {
328+
'trigger-1': {},
329+
'loop-1': { type: 'loop' },
330+
'agent-1': { data: { parentId: 'loop-1' } },
331+
},
332+
edges: [],
333+
})
334+
335+
await expect(
336+
executeManualWorkflowOperation.execute({
337+
principal,
338+
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-1' },
339+
})
340+
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('inside loop') })
341+
expect(mocks.executeService).not.toHaveBeenCalled()
342+
})
343+
344+
it('rejects a stop block the run cannot reach from its entry, which would run everything', async () => {
345+
mockLoadManualState.mockResolvedValue({
346+
blocks: { 'trigger-1': {}, 'agent-1': {}, 'agent-2': {}, unconnected: {} },
347+
edges: [
348+
{ source: 'trigger-1', target: 'agent-1' },
349+
{ source: 'agent-1', target: 'agent-2' },
350+
],
351+
})
352+
353+
await expect(
354+
executeManualWorkflowFromBlockOperation.execute({
355+
principal,
356+
input: {
357+
...baseInput,
358+
blockId: 'agent-2',
359+
sourceRunId: 'source-run-1',
360+
stopAfterBlockId: 'agent-1',
361+
},
362+
})
363+
).rejects.toMatchObject({
364+
code: 'validation',
365+
message: expect.stringContaining('not reachable'),
366+
})
367+
await expect(
368+
executeManualWorkflowOperation.execute({
369+
principal,
370+
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'unconnected' },
371+
})
372+
).rejects.toMatchObject({
373+
code: 'validation',
374+
message: expect.stringContaining('not reachable'),
375+
})
376+
expect(mocks.loadSourceState).not.toHaveBeenCalled()
377+
expect(mocks.executeService).not.toHaveBeenCalled()
378+
})
379+
380+
it('rejects a disabled stop block, or one reached only through a disabled block', async () => {
381+
mockLoadManualState.mockResolvedValue({
382+
blocks: {
383+
'trigger-1': {},
384+
'agent-1': { enabled: false },
385+
'agent-2': {},
386+
},
387+
edges: [
388+
{ source: 'trigger-1', target: 'agent-1' },
389+
{ source: 'agent-1', target: 'agent-2' },
390+
],
391+
})
392+
393+
await expect(
394+
executeManualWorkflowOperation.execute({
395+
principal,
396+
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-1' },
397+
})
398+
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('is disabled') })
399+
await expect(
400+
executeManualWorkflowOperation.execute({
401+
principal,
402+
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-2' },
403+
})
404+
).rejects.toMatchObject({
405+
code: 'validation',
406+
message: expect.stringContaining('not reachable'),
407+
})
408+
expect(mocks.executeService).not.toHaveBeenCalled()
409+
})
410+
411+
it('rejects a stop block named by an inherited object key', async () => {
412+
await expect(
413+
executeManualWorkflowOperation.execute({
414+
principal,
415+
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'toString' },
416+
})
417+
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') })
418+
expect(mocks.executeService).not.toHaveBeenCalled()
419+
})
420+
303421
it('rejects a source run without persisted state for this workflow', async () => {
304422
mocks.loadSourceState.mockResolvedValueOnce(null)
305423

‎apps/sim/lib/workflows/application/execute-manual-workflow.ts‎

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ interface ManualExecutionInput
2020
extends Omit<ExecuteWorkflowInput, 'input' | 'mode' | 'requestedTimeoutSeconds'> {
2121
input?: unknown
2222
mode: 'sync' | 'stream' | 'sync-result-stream'
23+
stopAfterBlockId?: string
2324
}
2425

2526
export interface ExecuteManualWorkflowInput extends ManualExecutionInput {
@@ -47,6 +48,69 @@ async function loadManualState(workflowId: string) {
4748
return state
4849
}
4950

51+
type ManualWorkflowState = Awaited<ReturnType<typeof loadManualState>>
52+
53+
/** Blocks a run entering at `entryBlockId` can reach; the executor skips disabled blocks. */
54+
function reachableFrom(state: ManualWorkflowState, entryBlockId: string): Set<string> {
55+
const targetsBySource = new Map<string, string[]>()
56+
for (const edge of state.edges) {
57+
const targets = targetsBySource.get(edge.source) ?? []
58+
targets.push(edge.target)
59+
targetsBySource.set(edge.source, targets)
60+
}
61+
const reached = new Set([entryBlockId])
62+
const queue = [entryBlockId]
63+
for (let next = queue.pop(); next !== undefined; next = queue.pop()) {
64+
for (const target of targetsBySource.get(next) ?? []) {
65+
if (reached.has(target) || state.blocks[target]?.enabled === false) continue
66+
reached.add(target)
67+
queue.push(target)
68+
}
69+
}
70+
return reached
71+
}
72+
73+
/**
74+
* The engine stops only when it completes a node whose id equals the target, so
75+
* a target the run cannot reach would silently run everything after the entry:
76+
* an unknown or disabled block, a block upstream of the entry or only behind a
77+
* disabled one, or a block inside a loop or parallel (which would stop after its
78+
* first iteration or never). All are refused, matching the editor, which offers
79+
* "Run until block" only outside subflows.
80+
*/
81+
function assertStopAfterBlock(
82+
state: ManualWorkflowState,
83+
blockId: string | undefined,
84+
entryBlockId: string
85+
): void {
86+
if (blockId === undefined) return
87+
const block = Object.hasOwn(state.blocks, blockId) ? state.blocks[blockId] : undefined
88+
if (!block) {
89+
throw new OrchestrationError(
90+
'validation',
91+
`run.stopAfterBlockId "${blockId}" is not a block in the current saved workflow.`
92+
)
93+
}
94+
if (block.enabled === false) {
95+
throw new OrchestrationError(
96+
'validation',
97+
`run.stopAfterBlockId "${blockId}" is disabled, so the run never executes it.`
98+
)
99+
}
100+
if (block.data?.parentId) {
101+
throw new OrchestrationError(
102+
'validation',
103+
`run.stopAfterBlockId "${blockId}" is inside loop or parallel "${block.data.parentId}"; stop after that container instead.`
104+
)
105+
}
106+
if (!reachableFrom(state, entryBlockId).has(blockId)) {
107+
throw new OrchestrationError(
108+
'validation',
109+
`run.stopAfterBlockId "${blockId}" is not reachable from entry block "${entryBlockId}"; stop after the entry block or one downstream of it.`
110+
)
111+
}
112+
}
113+
50114
function listTriggers(options: ReturnType<typeof resolveTriggerRunOptions>): string {
51115
return options.map((option) => `${option.triggerBlockId} (${option.blockName})`).join(', ')
52116
}
@@ -68,6 +132,7 @@ function executionServiceInput(params: {
68132
includeFileBase64: params.input.includeFileBase64,
69133
base64MaxBytes: params.input.base64MaxBytes,
70134
selectedOutputs: params.input.selectedOutputs,
135+
stopAfterBlockId: params.input.stopAfterBlockId,
71136
rateLimitCounter: 'sync' as const,
72137
abortSignal: params.input.abortSignal,
73138
mode: params.input.mode,
@@ -115,6 +180,7 @@ export const executeManualWorkflowOperation = defineAuthorizedWorkflowUseCase({
115180
)
116181
}
117182

183+
assertStopAfterBlock(state, input.stopAfterBlockId, selected.triggerBlockId)
118184
const executionInput = input.useMockPayload ? selected.mockPayload : input.input
119185
const validation = validateTriggerInput(selected, executionInput)
120186
if (!validation.ok) {
@@ -141,6 +207,7 @@ export const executeManualWorkflowFromBlockOperation = defineAuthorizedWorkflowU
141207
`run.entry.blockId "${input.blockId}" is not a block in the current saved workflow.`
142208
)
143209
}
210+
assertStopAfterBlock(state, input.stopAfterBlockId, input.blockId)
144211

145212
const sourceSnapshot = await getExecutionStateForWorkflow(input.sourceRunId, context.workflowId)
146213
if (!sourceSnapshot) {

0 commit comments

Comments
 (0)