Skip to content

Commit 79cf9c4

Browse files
author
Waleed Latif
committed
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.
1 parent 1056845 commit 79cf9c4

15 files changed

Lines changed: 603 additions & 15 deletions

File tree

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

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,16 @@ jobs:
125125
runs-on: ${{ (vars.CI_PROVIDER == '' || vars.CI_PROVIDER == 'blacksmith') && 'blacksmith-8vcpu-ubuntu-2404' || 'ubuntu-latest' }}
126126
timeout-minutes: 20
127127
services:
128+
# Hosted billing admits a workflow run only through a Redis usage reservation.
129+
redis:
130+
image: redis:8.2-alpine
131+
ports:
132+
- 6379:6379
133+
options: >-
134+
--health-cmd "redis-cli ping"
135+
--health-interval 5s
136+
--health-timeout 5s
137+
--health-retries 10
128138
postgres:
129139
image: pgvector/pgvector:pg17
130140
env:
@@ -212,6 +222,7 @@ jobs:
212222
ORGANIZATIONS_ENABLED: 'true'
213223
NEXT_PUBLIC_ORGANIZATIONS_ENABLED: 'true'
214224
INTERNAL_API_SECRET: scim-http-ci-local-secret-at-least-32-characters
225+
REDIS_URL: redis://127.0.0.1:6379
215226
DB_TX_TRIPWIRE: throw
216227
DISABLE_TELEMETRY: 'true'
217228
NEXT_TELEMETRY_DISABLED: '1'
@@ -252,6 +263,10 @@ jobs:
252263
VERSION_COMPARE_E2E_AUTH_SECRET="$BETTER_AUTH_SECRET" \
253264
VERSION_COMPARE_E2E_REPORT_PATH="$report_dir/version-compare-http-report.json" \
254265
bun run test:workflow-version-compare:e2e
266+
STOP_AFTER_E2E_BASE_URL="$NEXT_PUBLIC_APP_URL" \
267+
STOP_AFTER_E2E_DATABASE_URL="$DATABASE_URL" \
268+
STOP_AFTER_E2E_REPORT_PATH="$report_dir/stop-after-http-report.json" \
269+
bun run test:workflow-stop-after:e2e
255270
256271
- name: Upload end-to-end reports and server logs
257272
if: failure()

‎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: 71 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -300,6 +300,77 @@ 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' })
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' })
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' })
341+
expect(mocks.executeService).not.toHaveBeenCalled()
342+
})
343+
344+
it('stops the trigger and block entries after the requested block', async () => {
345+
mocks.loadSourceState.mockResolvedValueOnce({ blockStates: {}, executedBlocks: [] })
346+
347+
await executeManualWorkflowOperation.execute({
348+
principal,
349+
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-1' },
350+
})
351+
await executeManualWorkflowFromBlockOperation.execute({
352+
principal,
353+
input: {
354+
...baseInput,
355+
blockId: 'agent-1',
356+
sourceRunId: 'source-run-1',
357+
stopAfterBlockId: 'agent-1',
358+
},
359+
})
360+
361+
expect(mocks.executeService).toHaveBeenNthCalledWith(
362+
1,
363+
expect.objectContaining({ triggerBlockId: 'trigger-1', stopAfterBlockId: 'agent-1' })
364+
)
365+
expect(mocks.executeService).toHaveBeenNthCalledWith(
366+
2,
367+
expect.objectContaining({
368+
runFromBlock: expect.objectContaining({ startBlockId: 'agent-1' }),
369+
stopAfterBlockId: 'agent-1',
370+
})
371+
)
372+
})
373+
303374
it('rejects a source run without persisted state for this workflow', async () => {
304375
mocks.loadSourceState.mockResolvedValueOnce(null)
305376

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

Lines changed: 29 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,31 @@ async function loadManualState(workflowId: string) {
4748
return state
4849
}
4950

51+
type ManualWorkflowState = Awaited<ReturnType<typeof loadManualState>>
52+
53+
/**
54+
* The engine stops only when it completes a node whose id equals the target, so
55+
* an unknown id would silently run the whole workflow, and a block inside a loop
56+
* or parallel would stop after its first iteration or never. Both are refused,
57+
* matching the editor, which offers "Run until block" only outside subflows.
58+
*/
59+
function assertStopAfterBlock(state: ManualWorkflowState, blockId: string | undefined): void {
60+
if (blockId === undefined) return
61+
const block = state.blocks[blockId]
62+
if (!block) {
63+
throw new OrchestrationError(
64+
'validation',
65+
`run.stopAfterBlockId "${blockId}" is not a block in the current saved workflow.`
66+
)
67+
}
68+
if (block.data?.parentId) {
69+
throw new OrchestrationError(
70+
'validation',
71+
`run.stopAfterBlockId "${blockId}" is inside loop or parallel "${block.data.parentId}"; stop after that container instead.`
72+
)
73+
}
74+
}
75+
5076
function listTriggers(options: ReturnType<typeof resolveTriggerRunOptions>): string {
5177
return options.map((option) => `${option.triggerBlockId} (${option.blockName})`).join(', ')
5278
}
@@ -68,6 +94,7 @@ function executionServiceInput(params: {
6894
includeFileBase64: params.input.includeFileBase64,
6995
base64MaxBytes: params.input.base64MaxBytes,
7096
selectedOutputs: params.input.selectedOutputs,
97+
stopAfterBlockId: params.input.stopAfterBlockId,
7198
rateLimitCounter: 'sync' as const,
7299
abortSignal: params.input.abortSignal,
73100
mode: params.input.mode,
@@ -90,6 +117,7 @@ export const executeManualWorkflowOperation = defineAuthorizedWorkflowUseCase({
90117
)
91118
}
92119
const state = await loadManualState(context.workflowId)
120+
assertStopAfterBlock(state, input.stopAfterBlockId)
93121
const options = resolveTriggerRunOptions(mergeSubblockStateWithValues(state.blocks))
94122
if (options.length === 0) {
95123
throw new OrchestrationError(
@@ -141,6 +169,7 @@ export const executeManualWorkflowFromBlockOperation = defineAuthorizedWorkflowU
141169
`run.entry.blockId "${input.blockId}" is not a block in the current saved workflow.`
142170
)
143171
}
172+
assertStopAfterBlock(state, input.stopAfterBlockId)
144173

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

‎apps/sim/lib/workflows/executor/execute-service.ts‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,6 +120,8 @@ export interface ExecuteWorkflowServiceParams {
120120
/** Mocked upstream outputs (block name/id → output object) overlaid on the snapshot. */
121121
variableInputs?: Record<string, unknown>
122122
}
123+
/** Saved block after which execution stops, validated by the application use case. */
124+
stopAfterBlockId?: string
123125
}
124126

125127
export interface ExecuteWorkflowServiceFailure {
@@ -272,6 +274,7 @@ export async function executeWorkflowService(
272274
useDraftState = false,
273275
triggerBlockId,
274276
runFromBlock,
277+
stopAfterBlockId,
275278
} = params
276279

277280
let reqLogger = logger.withMetadata({ requestId, workflowId, userId })
@@ -289,6 +292,9 @@ export async function executeWorkflowService(
289292
if (runFromBlock && !useDraftState) {
290293
throw new Error('Run-from-block requires manual execution state')
291294
}
295+
if (stopAfterBlockId && !useDraftState) {
296+
throw new Error('Stop-after-block requires manual execution state')
297+
}
292298

293299
if (callChain) {
294300
const chainError = validateCallChain(callChain)
@@ -589,6 +595,7 @@ export async function executeWorkflowService(
589595
triggerBlockId,
590596
useDraftState,
591597
runFromBlock,
598+
stopAfterBlockId,
592599
onStream,
593600
onBlockComplete: (blockId, data) =>
594601
onBlockComplete(blockId, data.output, data.outputBlockId),
@@ -689,6 +696,7 @@ export async function executeWorkflowService(
689696
base64MaxBytes,
690697
abortSignal: timeoutController.signal,
691698
runFromBlock,
699+
stopAfterBlockId,
692700
})
693701

694702
await handlePostExecutionPauseState({ result, workflowId, executionId, loggingSession })

‎apps/sim/package.json‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
"test": "vitest run",
2424
"test:scim:e2e": "bun run scripts/test-scim-e2e.ts",
2525
"test:workflow-version-compare:e2e": "bun --no-env-file scripts/test-workflow-version-compare-e2e.ts",
26+
"test:workflow-stop-after:e2e": "bun --no-env-file scripts/test-workflow-stop-after-e2e.ts",
2627
"test:watch": "vitest",
2728
"test:coverage": "vitest run --coverage",
2829
"email:dev": "email dev --dir components/emails",

0 commit comments

Comments
 (0)