Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 45 additions & 0 deletions .github/workflows/test-build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -253,6 +253,51 @@ jobs:
VERSION_COMPARE_E2E_REPORT_PATH="$report_dir/version-compare-http-report.json" \
bun run test:workflow-version-compare:e2e

# A self-hosted app: hosted billing admits a run only through a Redis usage
# reservation, and the SCIM suite above asserts PostgreSQL rate-limit storage,
# so workflow execution gets its own app rather than adding Redis to that one.
- name: Verify single-block workflow runs over real HTTP
working-directory: apps/sim
env:
NEXT_PUBLIC_APP_URL: http://127.0.0.1:3018
BETTER_AUTH_URL: http://127.0.0.1:3018
NEXT_PUBLIC_FORCE_HOSTED: 'false'
INTERNAL_API_SECRET: stop-after-http-ci-local-secret-at-least-32-characters
DB_TX_TRIPWIRE: throw
DISABLE_TELEMETRY: 'true'
NEXT_TELEMETRY_DISABLED: '1'
NEXT_PUBLIC_CHAT_DISABLED: 'true'
READY_TIMEOUT_SECONDS: 300
run: |
report_dir="$RUNNER_TEMP/e2e"
server_log="$report_dir/stop-after-next.log"
mkdir -p "$report_dir"
node ../../node_modules/next/dist/bin/next dev --hostname 127.0.0.1 --port 3018 > "$server_log" 2>&1 &
server_pid=$!
finish() {
kill "$server_pid" 2>/dev/null || true
wait "$server_pid" 2>/dev/null || true
awk '/^ (GET|POST|PUT|PATCH|DELETE|HEAD) \/api\// { print }' "$server_log" > "$report_dir/stop-after-http-status.log"
}
trap finish EXIT
fail_startup() {
echo "::error::$1"
tail -n 200 "$server_log"
exit 1
}
started=$SECONDS
until curl --fail --silent --max-time 10 http://127.0.0.1:3018/api/health > /dev/null; do
kill -0 "$server_pid" 2>/dev/null || fail_startup 'Local workflow app exited during startup.'
[ $((SECONDS - started)) -lt "$READY_TIMEOUT_SECONDS" ] ||
fail_startup "Local workflow app did not become ready within $READY_TIMEOUT_SECONDS seconds."
sleep 2
done
echo "Local workflow app ready after $((SECONDS - started))s"
STOP_AFTER_E2E_BASE_URL="$NEXT_PUBLIC_APP_URL" \
STOP_AFTER_E2E_DATABASE_URL="$DATABASE_URL" \
STOP_AFTER_E2E_REPORT_PATH="$report_dir/stop-after-http-report.json" \
bun run test:workflow-stop-after:e2e

- name: Upload end-to-end reports and server logs
if: failure()
uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4
Expand Down
1 change: 1 addition & 0 deletions apps/docs/content/docs/cli/reference.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -6704,6 +6704,7 @@ sim workflows run <workflowId> [options]
| `--mock-payload` | No | Use the selected trigger's server-derived mock payload; runs the current saved workflow state (implies --manual). |
| `--from-block <blockId>` | No | Run manually from this saved workflow block. |
| `--source-run <runId>` | No | Prior run whose persisted state supplies upstream outputs (requires --from-block). |
| `--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). |
| `--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. |
| `--include-thinking` | No | Show model reasoning while following (requires --follow). |
| `--include-tool-calls` | No | Show tool calls while following (requires --follow). |
Expand Down
1 change: 1 addition & 0 deletions apps/docs/content/docs/cli/workflows.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -642,6 +642,7 @@ sim workflows run <workflowId> [options]
| `--mock-payload` | No | Use the selected trigger's server-derived mock payload; runs the current saved workflow state (implies --manual). |
| `--from-block <blockId>` | No | Run manually from this saved workflow block. |
| `--source-run <runId>` | No | Prior run whose persisted state supplies upstream outputs (requires --from-block). |
| `--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). |
| `--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. |
| `--include-thinking` | No | Show model reasoning while following (requires --follow). |
| `--include-tool-calls` | No | Show tool calls while following (requires --follow). |
Expand Down
16 changes: 16 additions & 0 deletions apps/docs/openapi-v2-workflows.json
Original file line number Diff line number Diff line change
Expand Up @@ -12615,6 +12615,11 @@
"additionalProperties": false
}
]
},
"stopAfterBlockId": {
"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.",
"type": "string",
"minLength": 1
}
},
"required": ["source"],
Expand Down Expand Up @@ -12703,6 +12708,17 @@
"sourceRunId": "run_123"
}
}
},
{
"run": {
"source": "manual",
"entry": {
"type": "block",
"blockId": "block_123",
"sourceRunId": "run_123"
},
"stopAfterBlockId": "block_123"
}
}
]
},
Expand Down
2 changes: 2 additions & 0 deletions apps/sim/app/api/v2/workflows/[workflowId]/execute/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -448,6 +448,7 @@ export const POST = withRouteHandler(
mode: body.stream ? 'stream' : resultStream ? 'sync-result-stream' : 'sync',
blockId: manualRun.entry.blockId,
sourceRunId: manualRun.entry.sourceRunId,
stopAfterBlockId: manualRun.stopAfterBlockId,
},
request: req,
})
Expand All @@ -460,6 +461,7 @@ export const POST = withRouteHandler(
mode: body.stream ? 'stream' : resultStream ? 'sync-result-stream' : 'sync',
triggerBlockId: manualRun.entry?.blockId,
useMockPayload: manualRun.entry?.useMockPayload === true,
stopAfterBlockId: manualRun.stopAfterBlockId,
},
request: req,
})
Expand Down
14 changes: 14 additions & 0 deletions apps/sim/lib/api/contracts/v2/workflows.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1297,6 +1297,13 @@ export const v2WorkflowRunSelectionSchema = z.discriminatedUnion('source', [
.describe(
'Manual entry mode. Omit to enter through the workflow trigger; a block entry requires an exact source run.'
),
stopAfterBlockId: z
.string()
.min(1, 'run.stopAfterBlockId cannot be empty')
.optional()
.describe(
'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.'
),
})
.strict(),
])
Expand Down Expand Up @@ -1412,6 +1419,13 @@ export const v2ExecuteWorkflowBodySchema = z
entry: { type: 'block', blockId: 'block_123', sourceRunId: 'run_123' },
},
},
{
run: {
source: 'manual',
entry: { type: 'block', blockId: 'block_123', sourceRunId: 'run_123' },
stopAfterBlockId: 'block_123',
},
},
],
})
export type V2ExecuteWorkflowBody = z.input<typeof v2ExecuteWorkflowBodySchema>
Expand Down
118 changes: 118 additions & 0 deletions apps/sim/lib/workflows/application/execute-manual-workflow.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -300,6 +300,124 @@ describe('manual workflow execution application operations', () => {
expect(mocks.loadSourceState).not.toHaveBeenCalled()
})

it('rejects a stop block missing from the saved workflow before anything runs', async () => {
await expect(
executeManualWorkflowOperation.execute({
principal,
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'missing' },
})
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') })
await expect(
executeManualWorkflowFromBlockOperation.execute({
principal,
input: {
...baseInput,
blockId: 'agent-1',
sourceRunId: 'source-run-1',
stopAfterBlockId: 'missing',
},
})
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') })
expect(mocks.loadSourceState).not.toHaveBeenCalled()
expect(mocks.executeService).not.toHaveBeenCalled()
})

it('rejects a stop block nested in a loop or parallel, which the engine cannot stop on', async () => {
mockLoadManualState.mockResolvedValue({
blocks: {
'trigger-1': {},
'loop-1': { type: 'loop' },
'agent-1': { data: { parentId: 'loop-1' } },
},
edges: [],
})

await expect(
executeManualWorkflowOperation.execute({
principal,
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-1' },
})
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('inside loop') })
expect(mocks.executeService).not.toHaveBeenCalled()
})

it('rejects a stop block the run cannot reach from its entry, which would run everything', async () => {
mockLoadManualState.mockResolvedValue({
blocks: { 'trigger-1': {}, 'agent-1': {}, 'agent-2': {}, unconnected: {} },
edges: [
{ source: 'trigger-1', target: 'agent-1' },
{ source: 'agent-1', target: 'agent-2' },
],
})

await expect(
executeManualWorkflowFromBlockOperation.execute({
principal,
input: {
...baseInput,
blockId: 'agent-2',
sourceRunId: 'source-run-1',
stopAfterBlockId: 'agent-1',
},
})
).rejects.toMatchObject({
code: 'validation',
message: expect.stringContaining('not reachable'),
})
await expect(
executeManualWorkflowOperation.execute({
principal,
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'unconnected' },
})
).rejects.toMatchObject({
code: 'validation',
message: expect.stringContaining('not reachable'),
})
expect(mocks.loadSourceState).not.toHaveBeenCalled()
expect(mocks.executeService).not.toHaveBeenCalled()
})

it('rejects a disabled stop block, or one reached only through a disabled block', async () => {
mockLoadManualState.mockResolvedValue({
blocks: {
'trigger-1': {},
'agent-1': { enabled: false },
'agent-2': {},
},
edges: [
{ source: 'trigger-1', target: 'agent-1' },
{ source: 'agent-1', target: 'agent-2' },
],
})

await expect(
executeManualWorkflowOperation.execute({
principal,
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-1' },
})
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('is disabled') })
await expect(
executeManualWorkflowOperation.execute({
principal,
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-2' },
})
).rejects.toMatchObject({
code: 'validation',
message: expect.stringContaining('not reachable'),
})
expect(mocks.executeService).not.toHaveBeenCalled()
})

it('rejects a stop block named by an inherited object key', async () => {
await expect(
executeManualWorkflowOperation.execute({
principal,
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'toString' },
})
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') })
expect(mocks.executeService).not.toHaveBeenCalled()
})

it('rejects a source run without persisted state for this workflow', async () => {
mocks.loadSourceState.mockResolvedValueOnce(null)

Expand Down
67 changes: 67 additions & 0 deletions apps/sim/lib/workflows/application/execute-manual-workflow.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ interface ManualExecutionInput
extends Omit<ExecuteWorkflowInput, 'input' | 'mode' | 'requestedTimeoutSeconds'> {
input?: unknown
mode: 'sync' | 'stream' | 'sync-result-stream'
stopAfterBlockId?: string
}

export interface ExecuteManualWorkflowInput extends ManualExecutionInput {
Expand Down Expand Up @@ -47,6 +48,69 @@ async function loadManualState(workflowId: string) {
return state
}

type ManualWorkflowState = Awaited<ReturnType<typeof loadManualState>>

/** Blocks a run entering at `entryBlockId` can reach; the executor skips disabled blocks. */
function reachableFrom(state: ManualWorkflowState, entryBlockId: string): Set<string> {
const targetsBySource = new Map<string, string[]>()
for (const edge of state.edges) {
const targets = targetsBySource.get(edge.source) ?? []
targets.push(edge.target)
targetsBySource.set(edge.source, targets)
}
const reached = new Set([entryBlockId])
const queue = [entryBlockId]
for (let next = queue.pop(); next !== undefined; next = queue.pop()) {
for (const target of targetsBySource.get(next) ?? []) {
if (reached.has(target) || state.blocks[target]?.enabled === false) continue
reached.add(target)
queue.push(target)
}
}
return reached
}

/**
* The engine stops only when it completes a node whose id equals the target, so
* a target the run cannot reach would silently run everything after the entry:
* an unknown or disabled block, a block upstream of the entry or only behind a
* disabled one, or a block inside a loop or parallel (which would stop after its
* first iteration or never). All are refused, matching the editor, which offers
* "Run until block" only outside subflows.
*/
function assertStopAfterBlock(
state: ManualWorkflowState,
blockId: string | undefined,
entryBlockId: string
): void {
if (blockId === undefined) return
const block = Object.hasOwn(state.blocks, blockId) ? state.blocks[blockId] : undefined
if (!block) {
throw new OrchestrationError(
'validation',
`run.stopAfterBlockId "${blockId}" is not a block in the current saved workflow.`
)
Comment thread
waleedlatif1 marked this conversation as resolved.
}
if (block.enabled === false) {
throw new OrchestrationError(
'validation',
`run.stopAfterBlockId "${blockId}" is disabled, so the run never executes it.`
)
}
if (block.data?.parentId) {
throw new OrchestrationError(
'validation',
`run.stopAfterBlockId "${blockId}" is inside loop or parallel "${block.data.parentId}"; stop after that container instead.`
)
}
if (!reachableFrom(state, entryBlockId).has(blockId)) {
throw new OrchestrationError(
'validation',
`run.stopAfterBlockId "${blockId}" is not reachable from entry block "${entryBlockId}"; stop after the entry block or one downstream of it.`
)
}
}

function listTriggers(options: ReturnType<typeof resolveTriggerRunOptions>): string {
return options.map((option) => `${option.triggerBlockId} (${option.blockName})`).join(', ')
}
Expand All @@ -68,6 +132,7 @@ function executionServiceInput(params: {
includeFileBase64: params.input.includeFileBase64,
base64MaxBytes: params.input.base64MaxBytes,
selectedOutputs: params.input.selectedOutputs,
stopAfterBlockId: params.input.stopAfterBlockId,
rateLimitCounter: 'sync' as const,
abortSignal: params.input.abortSignal,
mode: params.input.mode,
Expand Down Expand Up @@ -115,6 +180,7 @@ export const executeManualWorkflowOperation = defineAuthorizedWorkflowUseCase({
)
}

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

const sourceSnapshot = await getExecutionStateForWorkflow(input.sourceRunId, context.workflowId)
if (!sourceSnapshot) {
Expand Down
Loading
Loading