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: 32 additions & 13 deletions .github/scripts/stop-session.sh
Original file line number Diff line number Diff line change
Expand Up @@ -5,39 +5,58 @@
#
# `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.
# Every member of the session is signalled, including one that moved to its own process group.
# The leader stays a zombie until the shell that started it waits on it, so zombies don't count.
set -u

leader=$1
grace_seconds=10

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

signal_session() {
local pids
pids=$(session_pids)
[ -z "$pids" ] || kill "-$1" $pids 2>/dev/null || true
}

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

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

wait_for() {
for _ in $(seq 1 100); do
now_us() {
echo "${EPOCHREALTIME/./}"
}

# Waits while the given check holds, until the shared deadline. Succeeds once it stops holding.
wait_while() {
while (($(now_us) < deadline)); do
"$@" || return 0
sleep 0.1
done
return 1
! "$@"
}

kill -TERM -- "-$leader" 2>/dev/null || true
wait_for leader_running
if running; then
deadline=$(($(now_us) + grace_seconds * 1000000))
signal_session TERM
wait_while leader_running
if ! leader_running && 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
wait_while 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 "::warning::Processes from session $leader were still running ${grace_seconds}s after SIGTERM:"
ps -s "$leader" -o pid,stat,etimes,args 2>/dev/null | cut -c1-200 || true
deadline=$(($(now_us) + grace_seconds * 1000000))
signal_session KILL
wait_while running && exit 0

echo "::error::Processes from session $leader survived SIGKILL."
echo "::error::Processes from session $leader survived SIGKILL for ${grace_seconds}s."
exit 1
24 changes: 15 additions & 9 deletions .github/workflows/test-build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,13 @@ jobs:
# immediately, and either way the server log tail lands in the job log. No
# step timeout: the job's bound covers a hang without cutting a slow but
# healthy suite short of writing its report.
#
# Every `next dev` app in this job shares apps/sim/.next, so each one follows the
# same lifecycle. It starts from an empty Turbopack dev cache: a cache written under
# other NEXT_PUBLIC_* values, by a server that `next dev` SIGKILLs 100ms after
# SIGTERM, can panic Turbopack or wedge a route compile on restore. It 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 SCIM, administration and workflow comparisons over real HTTP
working-directory: apps/sim
env:
Expand Down Expand Up @@ -285,16 +292,21 @@ jobs:
report_dir="$RUNNER_TEMP/e2e"
server_log="$report_dir/cli-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
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 @@ -315,12 +327,6 @@ 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 Down
51 changes: 32 additions & 19 deletions apps/sim/scripts/test-desktop-inbox-e2e.ts
Original file line number Diff line number Diff line change
Expand Up @@ -284,20 +284,16 @@ function openDoorbell(desktop: Desktop) {
function startExecutor(desktop: Desktop, doorbell: ReturnType<typeof openDoorbell>) {
const ran = new Map<string, { toolName: string; chatId: string; executionToken: string }>()
const completed = new Set<string>()
let running = false
let inFlight: Promise<void> | null = null
let rerun = false
let stopped = false
const pull = async () => {
if (running) {
rerun = true
return
}
running = true
const drain = async () => {
try {
do {
rerun = false
const inbox = await pullInbox(desktop)
for (const item of inbox.items) {
if (stopped) return
if (item.kind !== 'call' || ran.has(item.toolCallId)) continue
const claimBody: ClaimDesktopToolBody = {
deviceId: desktop.deviceId,
Expand Down Expand Up @@ -351,8 +347,16 @@ function startExecutor(desktop: Desktop, doorbell: ReturnType<typeof openDoorbel
}
} while (rerun && !stopped)
} finally {
running = false
inFlight = null
}
}
const pull = () => {
if (inFlight) {
rerun = true
return inFlight
}
inFlight = drain()
return inFlight
}
const offDoorbell = doorbell.onEvent(
() => void pull().catch((error) => logger.error('pull failed', error))
Expand All @@ -361,14 +365,19 @@ function startExecutor(desktop: Desktop, doorbell: ReturnType<typeof openDoorbel
() => void pull().catch((error) => logger.error('reconcile failed', error)),
RECONCILE_MS
)
void pull()
void pull().catch((error) => logger.error('pull failed', error))
return {
ran,
completed,
stop() {
/**
* Resolves once the pull in flight has settled, so the fixture is never torn down under a
* request the executor already started.
*/
async stop() {
stopped = true
offDoorbell()
clearInterval(timer)
await inFlight?.catch(() => {})
},
}
}
Expand Down Expand Up @@ -414,15 +423,19 @@ async function run() {
})

/** `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 })
await check(
'reads the inbox and refuses malformed claim, lease and completion bodies',
async () => {
await pullInbox(desktop)
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}`)
Expand Down Expand Up @@ -544,7 +557,7 @@ async function run() {
})
})
} finally {
executor.stop()
await executor.stop()
}

await check('tells the device to cancel a call whose chat was stopped', async () => {
Expand Down
4 changes: 4 additions & 0 deletions apps/sim/scripts/test-workflow-stop-after-e2e.ts
Original file line number Diff line number Diff line change
Expand Up @@ -324,6 +324,10 @@ async function execCli(
return { exitCode: 0, stdout, stderr }
} catch (error) {
assert(isRecordLike(error), getErrorMessage(error))
assert(
error.code !== 'ERR_CHILD_PROCESS_STDIO_MAXBUFFER',
`sim workflows run printed more than ${MAX_RESPONSE_BYTES} bytes`
)
assert(!error.killed, `sim workflows run did not exit within ${REQUEST_TIMEOUT_MS / 1000}s`)
assert(
typeof error.code === 'number' &&
Expand Down
60 changes: 44 additions & 16 deletions apps/sim/scripts/test-workflow-version-compare-e2e.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,9 @@ import type { WorkflowState } from '@/stores/workflows/workflow/types'
const logger = createLogger('WorkflowVersionCompareE2E')
const execFileAsync = promisify(execFile)
const MAX_RESPONSE_BYTES = 2 * 1024 * 1024
/** The first request to each route cold-compiles it under `next dev`. */
const ROUTE_COMPILE_TIMEOUT_MS = 300_000
/** Every later request hits a compiled route. */
const REQUEST_TIMEOUT_MS = 60_000
const startedAt = new Date().toISOString()

Expand Down Expand Up @@ -309,30 +312,39 @@ async function seedBrowserWorkflows() {
}
}

async function observedFetch(input: string | URL | Request, init?: RequestInit): Promise<Response> {
async function observedFetch(
input: string | URL | Request,
init?: RequestInit,
timeoutMs = REQUEST_TIMEOUT_MS
): Promise<Response> {
const url = new URL(input instanceof Request ? input.url : input)
assert.equal(
url.origin,
baseUrl.origin,
'E2E requests must remain on the configured loopback app'
)
const method = init?.method ?? (input instanceof Request ? input.method : 'GET')
const started = performance.now()
const signal = init?.signal ?? (input instanceof Request ? input.signal : undefined)
// boundary-raw-fetch: protocol E2E exercises a separately running local app over real HTTP
const response = await fetch(input, {
...init,
redirect: 'error',
signal: signal
? AbortSignal.any([signal, AbortSignal.timeout(REQUEST_TIMEOUT_MS)])
: AbortSignal.timeout(REQUEST_TIMEOUT_MS),
})
requests.push({
method: init?.method ?? (input instanceof Request ? input.method : 'GET'),
path: url.pathname,
status: response.status,
durationMs: Math.round(performance.now() - started),
})
return response
const timeout = AbortSignal.timeout(timeoutMs)
try {
// boundary-raw-fetch: protocol E2E exercises a separately running local app over real HTTP
const response = await fetch(input, {
...init,
redirect: 'error',
signal: signal ? AbortSignal.any([signal, timeout]) : timeout,
})
requests.push({
method,
path: url.pathname,
status: response.status,
durationMs: Math.round(performance.now() - started),
})
return response
} catch (error) {
if (!timeout.aborted) throw error
throw new Error(`${method} ${url.pathname}: no response within ${timeoutMs / 1000}s`)
}
}

async function get(path: string, auth: { key?: string; session?: string } = {}, expected = 200) {
Expand Down Expand Up @@ -432,6 +444,22 @@ async function runMcp() {
try {
await check('seed disposable workspace, versions and credentials', seed)
if (browserFixturesPath) await check('seed browser scenarios', seedBrowserWorkflows)
await check('every route under test compiles and refuses an anonymous request', async () => {
for (const [method, path] of [
['GET', `/api/workflows/${workflowId}/deployments/1`],
['GET', `/api/v2/workflows/${workflowId}/versions/1`],
['GET', comparisonPath(1, 2)],
['POST', '/api/mcp'],
] as const) {
const response = await observedFetch(
new URL(path, baseUrl),
{ method, headers: { 'x-forwarded-for': '127.0.0.1' } },
ROUTE_COMPILE_TIMEOUT_MS
)
await response.body?.cancel()
assert.equal(response.status, 401, `${method} ${path}: unexpected HTTP status`)
}
})
await check(
'session preview migrates both versions while the public archive stays pinned',
async () => {
Expand Down
Loading