Skip to content

Commit cde4976

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

7 files changed

Lines changed: 171 additions & 48 deletions

File tree

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

Lines changed: 41 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -125,16 +125,6 @@ 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
138128
postgres:
139129
image: pgvector/pgvector:pg17
140130
env:
@@ -222,7 +212,6 @@ jobs:
222212
ORGANIZATIONS_ENABLED: 'true'
223213
NEXT_PUBLIC_ORGANIZATIONS_ENABLED: 'true'
224214
INTERNAL_API_SECRET: scim-http-ci-local-secret-at-least-32-characters
225-
REDIS_URL: redis://127.0.0.1:6379
226215
DB_TX_TRIPWIRE: throw
227216
DISABLE_TELEMETRY: 'true'
228217
NEXT_TELEMETRY_DISABLED: '1'
@@ -263,6 +252,47 @@ jobs:
263252
VERSION_COMPARE_E2E_AUTH_SECRET="$BETTER_AUTH_SECRET" \
264253
VERSION_COMPARE_E2E_REPORT_PATH="$report_dir/version-compare-http-report.json" \
265254
bun run test:workflow-version-compare:e2e
255+
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"
266296
STOP_AFTER_E2E_BASE_URL="$NEXT_PUBLIC_APP_URL" \
267297
STOP_AFTER_E2E_DATABASE_URL="$DATABASE_URL" \
268298
STOP_AFTER_E2E_REPORT_PATH="$report_dir/stop-after-http-report.json" \

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

Lines changed: 42 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -306,7 +306,7 @@ describe('manual workflow execution application operations', () => {
306306
principal,
307307
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'missing' },
308308
})
309-
).rejects.toMatchObject({ code: 'validation' })
309+
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') })
310310
await expect(
311311
executeManualWorkflowFromBlockOperation.execute({
312312
principal,
@@ -317,7 +317,7 @@ describe('manual workflow execution application operations', () => {
317317
stopAfterBlockId: 'missing',
318318
},
319319
})
320-
).rejects.toMatchObject({ code: 'validation' })
320+
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') })
321321
expect(mocks.loadSourceState).not.toHaveBeenCalled()
322322
expect(mocks.executeService).not.toHaveBeenCalled()
323323
})
@@ -337,38 +337,54 @@ describe('manual workflow execution application operations', () => {
337337
principal,
338338
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-1' },
339339
})
340-
).rejects.toMatchObject({ code: 'validation' })
340+
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('inside loop') })
341341
expect(mocks.executeService).not.toHaveBeenCalled()
342342
})
343343

344-
it('stops the trigger and block entries after the requested block', async () => {
345-
mocks.loadSourceState.mockResolvedValueOnce({ blockStates: {}, executedBlocks: [] })
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+
})
346352

347-
await executeManualWorkflowOperation.execute({
348-
principal,
349-
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-1' },
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'),
350366
})
351-
await executeManualWorkflowFromBlockOperation.execute({
352-
principal,
353-
input: {
354-
...baseInput,
355-
blockId: 'agent-1',
356-
sourceRunId: 'source-run-1',
357-
stopAfterBlockId: 'agent-1',
358-
},
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'),
359375
})
376+
expect(mocks.loadSourceState).not.toHaveBeenCalled()
377+
expect(mocks.executeService).not.toHaveBeenCalled()
378+
})
360379

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',
380+
it('rejects a stop block named by an inherited object key', async () => {
381+
await expect(
382+
executeManualWorkflowOperation.execute({
383+
principal,
384+
input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'toString' },
370385
})
371-
)
386+
).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') })
387+
expect(mocks.executeService).not.toHaveBeenCalled()
372388
})
373389

374390
it('rejects a source run without persisted state for this workflow', async () => {

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

Lines changed: 38 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -50,15 +50,40 @@ async function loadManualState(workflowId: string) {
5050

5151
type ManualWorkflowState = Awaited<ReturnType<typeof loadManualState>>
5252

53+
function reachableFrom(state: ManualWorkflowState, entryBlockId: string): Set<string> {
54+
const targetsBySource = new Map<string, string[]>()
55+
for (const edge of state.edges) {
56+
const targets = targetsBySource.get(edge.source) ?? []
57+
targets.push(edge.target)
58+
targetsBySource.set(edge.source, targets)
59+
}
60+
const reached = new Set([entryBlockId])
61+
const queue = [entryBlockId]
62+
for (let next = queue.pop(); next !== undefined; next = queue.pop()) {
63+
for (const target of targetsBySource.get(next) ?? []) {
64+
if (reached.has(target)) continue
65+
reached.add(target)
66+
queue.push(target)
67+
}
68+
}
69+
return reached
70+
}
71+
5372
/**
5473
* 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.
74+
* a target the run cannot reach would silently run everything after the entry:
75+
* an unknown id, a block upstream of the entry, or a block inside a loop or
76+
* parallel (which would stop after its first iteration or never). All are
77+
* refused, matching the editor, which offers "Run until block" only outside
78+
* subflows.
5879
*/
59-
function assertStopAfterBlock(state: ManualWorkflowState, blockId: string | undefined): void {
80+
function assertStopAfterBlock(
81+
state: ManualWorkflowState,
82+
blockId: string | undefined,
83+
entryBlockId: string
84+
): void {
6085
if (blockId === undefined) return
61-
const block = state.blocks[blockId]
86+
const block = Object.hasOwn(state.blocks, blockId) ? state.blocks[blockId] : undefined
6287
if (!block) {
6388
throw new OrchestrationError(
6489
'validation',
@@ -71,6 +96,12 @@ function assertStopAfterBlock(state: ManualWorkflowState, blockId: string | unde
7196
`run.stopAfterBlockId "${blockId}" is inside loop or parallel "${block.data.parentId}"; stop after that container instead.`
7297
)
7398
}
99+
if (!reachableFrom(state, entryBlockId).has(blockId)) {
100+
throw new OrchestrationError(
101+
'validation',
102+
`run.stopAfterBlockId "${blockId}" is not reachable from entry block "${entryBlockId}"; stop after the entry block or one downstream of it.`
103+
)
104+
}
74105
}
75106

76107
function listTriggers(options: ReturnType<typeof resolveTriggerRunOptions>): string {
@@ -117,7 +148,6 @@ export const executeManualWorkflowOperation = defineAuthorizedWorkflowUseCase({
117148
)
118149
}
119150
const state = await loadManualState(context.workflowId)
120-
assertStopAfterBlock(state, input.stopAfterBlockId)
121151
const options = resolveTriggerRunOptions(mergeSubblockStateWithValues(state.blocks))
122152
if (options.length === 0) {
123153
throw new OrchestrationError(
@@ -143,6 +173,7 @@ export const executeManualWorkflowOperation = defineAuthorizedWorkflowUseCase({
143173
)
144174
}
145175

176+
assertStopAfterBlock(state, input.stopAfterBlockId, selected.triggerBlockId)
146177
const executionInput = input.useMockPayload ? selected.mockPayload : input.input
147178
const validation = validateTriggerInput(selected, executionInput)
148179
if (!validation.ok) {
@@ -169,7 +200,7 @@ export const executeManualWorkflowFromBlockOperation = defineAuthorizedWorkflowU
169200
`run.entry.blockId "${input.blockId}" is not a block in the current saved workflow.`
170201
)
171202
}
172-
assertStopAfterBlock(state, input.stopAfterBlockId)
203+
assertStopAfterBlock(state, input.stopAfterBlockId, input.blockId)
173204

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

‎apps/sim/lib/workflows/executor/execution-core.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -907,6 +907,9 @@ async function executeWorkflowCoreImpl(
907907
resolvedStopAfterBlockId = buildLoopSentinelEndId(stopAfterBlockId)
908908
} else if (serializedWorkflow.parallels?.[stopAfterBlockId]) {
909909
resolvedStopAfterBlockId = buildParallelSentinelEndId(stopAfterBlockId)
910+
} else if (!serializedWorkflow.blocks.some((block) => block.id === stopAfterBlockId)) {
911+
// The engine stops on an exact node id, so an absent target would run everything.
912+
throw new Error(`Stop block ${stopAfterBlockId} is not in the workflow being executed`)
910913
}
911914
}
912915

‎apps/sim/scripts/test-workflow-stop-after-e2e.ts‎

Lines changed: 36 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import { assertDisposableTestDatabaseUrl } from '@sim/db/testing/test-infrastruc
99
import { createLogger } from '@sim/logger'
1010
import { sha256Hex } from '@sim/security/hash'
1111
import { getErrorMessage } from '@sim/utils/errors'
12+
import { sleep } from '@sim/utils/helpers'
1213
import { generateId } from '@sim/utils/id'
1314
import { isRecordLike } from '@sim/utils/object'
1415
import { truncate } from '@sim/utils/string'
@@ -32,7 +33,8 @@ import { readResponseTextWithLimit } from '@/lib/core/utils/stream-limits'
3233
const logger = createLogger('WorkflowStopAfterE2E')
3334
const execFileAsync = promisify(execFile)
3435
const MAX_RESPONSE_BYTES = 2 * 1024 * 1024
35-
const REQUEST_TIMEOUT_MS = 120_000
36+
/** The first execute request cold-compiles the route under `next dev`. */
37+
const REQUEST_TIMEOUT_MS = 300_000
3638
const SLOW_SECONDS = 4
3739
const SLOW_MS = SLOW_SECONDS * 1000
3840
const startedAt = new Date().toISOString()
@@ -319,6 +321,17 @@ try {
319321
{ run: { source: 'manual', stopAfterBlockId: '' } },
320322
'BAD_REQUEST'
321323
)
324+
await expectBadRequest(
325+
pipeline.workflowId,
326+
{
327+
run: {
328+
source: 'manual',
329+
entry: { type: 'block', blockId: pipeline.check, sourceRunId },
330+
stopAfterBlockId: pipeline.slow,
331+
},
332+
},
333+
'BAD_REQUEST'
334+
)
322335
await expectBadRequest(
323336
pipeline.workflowId,
324337
{ run: { source: 'manual', stopAfterBlockId: pipeline.check }, async: true },
@@ -360,8 +373,28 @@ try {
360373
} finally {
361374
try {
362375
await check('remove disposable fixtures', async () => {
363-
await sql`delete from workspace where id = ${workspaceId}`
364-
await sql`delete from "user" where id = ${ownerId}`
376+
// A response returns before its run finishes persisting logs and large-value
377+
// references; a cascade delete racing those writes can be chosen as a deadlock victim.
378+
const workflowIds = [pipeline.workflowId, otherPipeline.workflowId]
379+
for (let attempt = 0; attempt < 120; attempt++) {
380+
const [{ open }] =
381+
await sql`select count(*)::int as open from workflow_execution_logs where workflow_id in ${sql(workflowIds)} and ended_at is null`
382+
if (open === 0) break
383+
await sleep(500)
384+
}
385+
for (let attempt = 1; ; attempt++) {
386+
try {
387+
await sql.begin(async (tx) => {
388+
await tx`delete from workspace where id = ${workspaceId}`
389+
await tx`delete from "user" where id = ${ownerId}`
390+
})
391+
break
392+
} catch (error) {
393+
const deadlocked = isRecordLike(error) && error.code === '40P01'
394+
if (!deadlocked || attempt === 5) throw error
395+
await sleep(1000)
396+
}
397+
}
365398
if (directory) await rm(directory, { recursive: true, force: true })
366399
})
367400
} catch (error) {

‎packages/sim-cli/src/commands/protocol/workflow-run-follow.test.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -243,6 +243,13 @@ describe('sim workflows run --follow', () => {
243243
})
244244
})
245245

246+
it('refuses an empty --stop-after rather than running the whole draft', async () => {
247+
await expect(run(WORKFLOW_ID, '--stop-after', '')).rejects.toThrow(
248+
'--stop-after requires a block ID'
249+
)
250+
expect(requestRaw).not.toHaveBeenCalled()
251+
})
252+
246253
it('refuses --stop-after with --async before sending anything', async () => {
247254
await expect(run(WORKFLOW_ID, '--stop-after', 'agent-1', '--async')).rejects.toThrow(
248255
'Manual execution does not support --async'

‎packages/sim-cli/src/commands/protocol/workflow-run-follow.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -79,8 +79,11 @@ export function resolveWorkflowRunSelection(
7979
if (useMockPayload && flags.input !== undefined) {
8080
throw new SimApiError('--mock-payload cannot be combined with --input', 0)
8181
}
82+
if (stopAfter !== undefined && stopAfter.trim() === '') {
83+
throw new SimApiError('--stop-after requires a block ID', 0)
84+
}
8285

83-
const stop = stopAfter ? { stopAfterBlockId: stopAfter } : {}
86+
const stop = stopAfter !== undefined ? { stopAfterBlockId: stopAfter } : {}
8487
if (fromBlock && sourceRun) {
8588
return {
8689
source: 'manual',

0 commit comments

Comments
 (0)