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
43 changes: 43 additions & 0 deletions .github/scripts/stop-session.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
#!/usr/bin/env bash
# Stops every process in a session started with `setsid`, and returns only once none is left.
#
# Usage: stop-session.sh <session-leader-pid>
#
# `next dev` runs its server and workers as child processes. Waiting on `next dev` alone returns
# while those can still be running, and the next app in the job reuses the same .next directory.
# The leader stays a zombie until the shell that started it waits on it, so zombies don't count.
set -u

leader=$1

running() {
ps -s "$leader" -o pid=,stat= 2>/dev/null | awk '$2 !~ /^Z/ { found = 1 } END { exit !found }'
}

leader_running() {
ps -p "$leader" -o stat= 2>/dev/null | grep -qv '^Z'
}

wait_for() {
for _ in $(seq 1 100); do
"$@" || return 0
sleep 0.1
done
return 1
}

kill -TERM -- "-$leader" 2>/dev/null || true
wait_for leader_running
if running; then
echo "Processes from session $leader outlived its leader:"
ps -s "$leader" -o pid,stat,etimes,args 2>/dev/null | cut -c1-200 || true
fi
wait_for running && exit 0

echo "::warning::Processes from session $leader were still running 10s after SIGTERM:"
ps -s "$leader" -o pid,stat,etimes,args 2>/dev/null || true
kill -KILL -- "-$leader" 2>/dev/null || true
wait_for running && exit 0

echo "::error::Processes from session $leader survived SIGKILL."
exit 1
39 changes: 30 additions & 9 deletions .github/workflows/test-build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -232,17 +232,22 @@ jobs:
report_dir="$RUNNER_TEMP/e2e"
server_log="$report_dir/scim-next.log"
mkdir -p "$report_dir"
node ../../node_modules/next/dist/bin/next dev --hostname 127.0.0.1 --port 3017 > "$server_log" 2>&1 &
rm -rf .next/dev
setsid node ../../node_modules/next/dist/bin/next dev --hostname 127.0.0.1 --port 3017 > "$server_log" 2>&1 &
server_pid=$!
finish() {
kill "$server_pid" 2>/dev/null || true
status=$?
bash ../../.github/scripts/stop-session.sh "$server_pid" || status=1
wait "$server_pid" 2>/dev/null || true
awk '/^ (GET|POST|PUT|PATCH|DELETE|HEAD) \/api\// { print }' "$server_log" > "$report_dir/scim-http-status.log"
if [ "$status" -ne 0 ]; then
tail -n 200 "$server_log"
exit "$status"
fi
}
trap finish EXIT
fail_startup() {
echo "::error::$1"
tail -n 200 "$server_log"
exit 1
}
started=$SECONDS
Expand All @@ -267,6 +272,12 @@ jobs:
# 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.
#
# The apps in this job share apps/sim/.next. Each starts from an empty Turbopack
# dev cache: one written under other NEXT_PUBLIC_* values, by a server `next dev`
# SIGKILLs 100ms after SIGTERM, can panic Turbopack or wedge a route compile on
# restore. Each runs in its own session, and stop-session.sh returns only once
# all of it has exited, so no app starts beside one still writing .next/dev.
- name: Verify single-block workflow runs over real HTTP
working-directory: apps/sim
env:
Expand All @@ -283,17 +294,22 @@ jobs:
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 &
rm -rf .next/dev
setsid 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
status=$?
bash ../../.github/scripts/stop-session.sh "$server_pid" || status=1
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"
if [ "$status" -ne 0 ]; then
tail -n 200 "$server_log"
exit "$status"
fi
}
trap finish EXIT
fail_startup() {
echo "::error::$1"
tail -n 200 "$server_log"
exit 1
}
started=$SECONDS
Expand Down Expand Up @@ -329,17 +345,22 @@ jobs:
report_dir="$RUNNER_TEMP/e2e"
server_log="$report_dir/desktop-inbox-next.log"
mkdir -p "$report_dir"
node ../../node_modules/next/dist/bin/next dev --hostname 127.0.0.1 --port 3019 > "$server_log" 2>&1 &
rm -rf .next/dev
setsid node ../../node_modules/next/dist/bin/next dev --hostname 127.0.0.1 --port 3019 > "$server_log" 2>&1 &
server_pid=$!
finish() {
kill "$server_pid" 2>/dev/null || true
status=$?
bash ../../.github/scripts/stop-session.sh "$server_pid" || status=1
wait "$server_pid" 2>/dev/null || true
awk '/^ (GET|POST|PUT|PATCH|DELETE|HEAD) \/api\// { print }' "$server_log" > "$report_dir/desktop-inbox-http-status.log"
if [ "$status" -ne 0 ]; then
tail -n 200 "$server_log"
exit "$status"
fi
}
trap finish EXIT
fail_startup() {
echo "::error::$1"
tail -n 200 "$server_log"
exit 1
}
started=$SECONDS
Expand Down
11 changes: 11 additions & 0 deletions apps/sim/scripts/test-desktop-inbox-e2e.ts
Original file line number Diff line number Diff line change
Expand Up @@ -413,6 +413,17 @@ async function run() {
assert(registration.reconcileMs < PICKUP_GRACE_SECONDS * 1000)
})

/** `next dev` compiles a route on its first request, which must not count against the timed checks. */
await check('refuses malformed claim, lease and completion bodies', async () => {
for (const path of [
'/api/desktop/tool/claim',
'/api/desktop/tool/lease',
'/api/desktop/tool/complete',
]) {
await request(desktop, 'POST', path, { body: {}, expected: 400 })
}
})

/** Starts absent, so only the stream open below can mark the device present. */
await redis.del(`desktop:presence:${desktop.deviceId}`)
const doorbell = openDoorbell(desktop)
Expand Down
170 changes: 109 additions & 61 deletions apps/sim/scripts/test-workflow-stop-after-e2e.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,10 @@ import { readResponseTextWithLimit } from '@/lib/core/utils/stream-limits'
const logger = createLogger('WorkflowStopAfterE2E')
const execFileAsync = promisify(execFile)
const MAX_RESPONSE_BYTES = 2 * 1024 * 1024
/** The first execute request cold-compiles the route under `next dev`. */
const REQUEST_TIMEOUT_MS = 300_000
/** The first execute request cold-compiles the route's module graph under `next dev`. */
const ROUTE_COMPILE_TIMEOUT_MS = 300_000
/** Every later request hits the compiled route; the slowest fixture run waits about `SLOW_MS`. */
const REQUEST_TIMEOUT_MS = 60_000
const SLOW_SECONDS = 4
const SLOW_MS = SLOW_SECONDS * 1000
const startedAt = new Date().toISOString()
Expand Down Expand Up @@ -67,7 +69,8 @@ const personalKey = `sk-sim-fixture-${generateId()}`
const cliPath = fileURLToPath(new URL('../../../packages/sim-cli/src/index.ts', import.meta.url))
const checks: { name: string; status: 'passed' | 'failed'; durationMs: number; error?: string }[] =
[]
const requests: { method: string; path: string; status: number; durationMs: number }[] = []
/** `status` is null when the request ended without a complete response. */
const requests: { method: string; path: string; status: number | null; durationMs: number }[] = []
let directory: string | undefined

interface PipelineFixture {
Expand Down Expand Up @@ -225,37 +228,48 @@ async function seed() {
})
}

function isTimeout(error: unknown): boolean {
return error instanceof DOMException && error.name === 'TimeoutError'
}

async function execute(
workflowId: string,
body: V2ExecuteWorkflowBody,
expectedStatus = 200
{ expectedStatus = 200, timeoutMs = REQUEST_TIMEOUT_MS } = {}
): Promise<Record<string, unknown>> {
const url = new URL(`/api/v2/workflows/${workflowId}/execute`, baseUrl)
const started = performance.now()
// boundary-raw-fetch: protocol E2E exercises a separately running local app over real HTTP
const response = await fetch(url, {
method: 'POST',
redirect: 'error',
signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS),
headers: {
accept: 'application/json',
'content-type': 'application/json',
'x-api-key': personalKey,
'x-forwarded-for': '127.0.0.1',
},
body: JSON.stringify(body),
})
requests.push({
method: 'POST',
path: url.pathname,
status: response.status,
durationMs: Math.round(performance.now() - started),
})
const text = await readResponseTextWithLimit(response, {
maxBytes: MAX_RESPONSE_BYTES,
label: 'Stop-after E2E response',
})
assert.equal(response.status, expectedStatus, `${url.pathname}: ${truncate(text, 500)}`)
const elapsed = () => Math.round(performance.now() - started)
let status: number | null = null
let text: string
try {
// boundary-raw-fetch: protocol E2E exercises a separately running local app over real HTTP
const response = await fetch(url, {
method: 'POST',
redirect: 'error',
signal: AbortSignal.timeout(timeoutMs),
headers: {
accept: 'application/json',
'content-type': 'application/json',
'x-api-key': personalKey,
'x-forwarded-for': '127.0.0.1',
},
body: JSON.stringify(body),
})
text = await readResponseTextWithLimit(response, {
maxBytes: MAX_RESPONSE_BYTES,
label: 'Stop-after E2E response',
})
status = response.status
} catch (error) {
const reason = isTimeout(error)
? `no complete response within ${timeoutMs / 1000}s`
: getErrorMessage(error)
throw new Error(`POST ${url.pathname} failed after ${elapsed()} ms: ${reason}`)
} finally {
requests.push({ method: 'POST', path: url.pathname, status, durationMs: elapsed() })
}
assert.equal(status, expectedStatus, `${url.pathname}: ${truncate(text, 500)}`)
return record(JSON.parse(text))
}

Expand All @@ -268,56 +282,90 @@ async function run(
return data
}

async function expectBadRequest(workflowId: string, body: V2ExecuteWorkflowBody, code: string) {
const status = code === 'NOT_FOUND' ? 404 : 400
const error = record((await execute(workflowId, body, status)).error)
async function expectBadRequest(
workflowId: string,
body: V2ExecuteWorkflowBody,
code: string,
timeoutMs = REQUEST_TIMEOUT_MS
) {
const expectedStatus = code === 'NOT_FOUND' ? 404 : 400
const error = record((await execute(workflowId, body, { expectedStatus, timeoutMs })).error)
assert.equal(error.code, code)
}

async function runCli(args: string[]): Promise<V2ExecuteWorkflowData> {
/** Runs the CLI; a run it must fail exits non-zero and still prints the run on stdout. */
async function execCli(
args: string[]
): Promise<{ exitCode: number; stdout: string; stderr: string }> {
assert(directory, 'CLI fixture directory must exist')
const { stdout } = await execFileAsync(
'bun',
[
'--no-env-file',
cliPath,
'--endpoint',
baseUrl.origin,
'--workspace',
workspaceId,
'--output',
'json',
'workflows',
'run',
...args,
],
{
cwd: directory,
env: { ...process.env, SIM_CONFIG_DIR: directory, SIM_API_KEY: personalKey, NO_COLOR: '1' },
timeout: REQUEST_TIMEOUT_MS,
maxBuffer: MAX_RESPONSE_BYTES,
}
try {
const { stdout, stderr } = await execFileAsync(
'bun',
[
'--no-env-file',
cliPath,
'--endpoint',
baseUrl.origin,
'--workspace',
workspaceId,
'--output',
'json',
'workflows',
'run',
...args,
],
{
cwd: directory,
env: { ...process.env, SIM_CONFIG_DIR: directory, SIM_API_KEY: personalKey, NO_COLOR: '1' },
timeout: REQUEST_TIMEOUT_MS,
maxBuffer: MAX_RESPONSE_BYTES,
}
)
return { exitCode: 0, stdout, stderr }
} catch (error) {
assert(isRecordLike(error), getErrorMessage(error))
assert(!error.killed, `sim workflows run did not exit within ${REQUEST_TIMEOUT_MS / 1000}s`)
assert(
typeof error.code === 'number' &&
typeof error.stdout === 'string' &&
typeof error.stderr === 'string',
getErrorMessage(error)
)
return { exitCode: error.code, stdout: error.stdout, stderr: error.stderr }
}
}

async function runCli(args: string[]): Promise<V2ExecuteWorkflowData> {
const { exitCode, stdout, stderr } = await execCli(args)
assert.equal(
exitCode,
0,
`sim workflows run exited ${exitCode}: ${truncate(stderr || stdout, 500)}`
)
return v2ExecuteWorkflowDataSchema.parse(JSON.parse(stdout))
}

/** A CLI run the command itself must fail: exits non-zero and prints the failed run. */
async function runCliExpectingFailure(args: string[]): Promise<V2ExecuteWorkflowData> {
try {
await runCli(args)
} catch (error) {
assert(isRecordLike(error) && typeof error.stdout === 'string', getErrorMessage(error))
assert.notEqual(error.code, 0, 'a failed run must exit non-zero')
return v2ExecuteWorkflowDataSchema.parse(JSON.parse(error.stdout))
}
assert.fail('the CLI exited 0 for a run that must fail')
const { exitCode, stdout } = await execCli(args)
assert.notEqual(exitCode, 0, 'the CLI exited 0 for a run that must fail')
return v2ExecuteWorkflowDataSchema.parse(JSON.parse(stdout))
}

const selectAll = ['Slow.status', 'Check.status', 'After.status']

try {
await check('seed disposable workspace, personal key and fixtures', seed)

await check('the execute route compiles and refuses a run before it starts', () =>
expectBadRequest(
pipeline.workflowId,
{ run: { source: 'manual', stopAfterBlockId: '' } },
'BAD_REQUEST',
ROUTE_COMPILE_TIMEOUT_MS
)
)

let sourceRunId = ''
await check('a full manual run executes every block and persists its state', async () => {
const full = await run(pipeline.workflowId, {
Expand Down
Loading