From f08eefc09becfdeb07d6a4328df6696fa34e2ff3 Mon Sep 17 00:00:00 2001 From: Dan Sutton Date: Mon, 11 May 2026 11:52:19 +0100 Subject: [PATCH] feat(references): add stress-tasks reference project for trigger fan-out repro --- apps/webapp/seed.mts | 26 +- pnpm-lock.yaml | 16 + references/stress-tasks/EXAMPLES.md | 119 ++++++++ references/stress-tasks/package.json | 17 ++ references/stress-tasks/src/trigger/fanout.ts | 273 ++++++++++++++++++ references/stress-tasks/trigger.config.ts | 15 + references/stress-tasks/tsconfig.json | 15 + 7 files changed, 480 insertions(+), 1 deletion(-) create mode 100644 references/stress-tasks/EXAMPLES.md create mode 100644 references/stress-tasks/package.json create mode 100644 references/stress-tasks/src/trigger/fanout.ts create mode 100644 references/stress-tasks/trigger.config.ts create mode 100644 references/stress-tasks/tsconfig.json diff --git a/apps/webapp/seed.mts b/apps/webapp/seed.mts index 9eb30cd25..7f364595f 100644 --- a/apps/webapp/seed.mts +++ b/apps/webapp/seed.mts @@ -67,11 +67,35 @@ async function seed() { name: "realtime-streams", externalRef: "proj_klxlzjnzxmbgiwuuwhvb", }, + { + name: "stress-tasks", + externalRef: "proj_stresstaskslocaldevx", + // Stress-tasks fan-outs need a much higher concurrency ceiling than the + // default 300 — at 1000+ children per parent, runs would otherwise queue + // and the local repro wouldn't track the production fan-out signature. + environmentConcurrencyLimit: 25000, + }, ]; // Create or find each project for (const projectConfig of referenceProjects) { - await findOrCreateProject(projectConfig.name, organization, user.id, projectConfig.externalRef); + const result = await findOrCreateProject( + projectConfig.name, + organization, + user.id, + projectConfig.externalRef, + ); + + if (projectConfig.environmentConcurrencyLimit) { + const updated = await prisma.runtimeEnvironment.updateMany({ + where: { projectId: result.project.id }, + data: { maximumConcurrencyLimit: projectConfig.environmentConcurrencyLimit }, + }); + console.log( + ` Updated ${updated.count} environment(s) on ${projectConfig.name} ` + + `to maximumConcurrencyLimit=${projectConfig.environmentConcurrencyLimit}`, + ); + } } await createBatchLimitOrgs(user); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 17f73d9a2..f3e61e660 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -3028,6 +3028,22 @@ importers: specifier: workspace:* version: link:../../packages/cli-v3 + references/stress-tasks: + dependencies: + '@trigger.dev/build': + specifier: workspace:* + version: link:../../packages/build + '@trigger.dev/sdk': + specifier: workspace:* + version: link:../../packages/trigger-sdk + zod: + specifier: 3.25.76 + version: 3.25.76 + devDependencies: + trigger.dev: + specifier: workspace:* + version: link:../../packages/cli-v3 + references/telemetry: dependencies: '@opentelemetry/resources': diff --git a/references/stress-tasks/EXAMPLES.md b/references/stress-tasks/EXAMPLES.md new file mode 100644 index 000000000..08b626bf2 --- /dev/null +++ b/references/stress-tasks/EXAMPLES.md @@ -0,0 +1,119 @@ +# Stress-tasks — example payloads + +Copy any of these into the dashboard test UI (Tasks → pick the task → Test). +The trigger.dev test UI defaults to the most recent run's payload, so once +you've fired a particular shape once, it'll be remembered. + +## `stress-fan-out-trigger` — N individual `.trigger()` calls in a single trace + +Mirrors the production failure mode (events 1–10 in +`prisma-connection-investigation-results.md`) where one trace fans out N HTTP +triggers and exhausts the api-prod Prisma connection pool. + +### Smoke test (use this first to confirm wiring) + +```json +{ "count": 10 } +``` + +### Reproduce the prod fan-out — 1,000 all at once + +```json +{ "count": 1000 } +``` + +### Bounded producer — only 100 in-flight at a time + +```json +{ "count": 1000, "concurrency": 100 } +``` + +### Exercise the `runTags ||` row-lock contention path (events 3, 4, 5, 7) + +```json +{ "count": 1000, "tags": ["stress-test", "burst-2026-05-08"] } +``` + +### Children doing real work — 500 triggers, 2 s child sleep, 200 in flight + +```json +{ "count": 500, "concurrency": 200, "childSleepMs": 2000 } +``` + +### Large payloads — 200 triggers, 50 KB pad each + +```json +{ "count": 200, "childPayloadBytes": 50000 } +``` + +### Combined contention — fan-out + tags + child work + +```json +{ "count": 1000, "concurrency": 250, "childSleepMs": 500, "tags": ["combined"] } +``` + +--- + +## `stress-fan-out-batch` — N triggers via chunked `batchTrigger` + +Different server-side code path: one HTTP request per chunk, server-side +bulk insert. Useful contrast for understanding whether pool pressure is +specific to the N-trigger path or surfaces here too. + +### Smoke test + +```json +{ "count": 10, "batchSize": 10 } +``` + +### Default — 1,000 across two sequential 500-payload batches + +```json +{ "count": 1000 } +``` + +### Parallel batches — same volume, two batchTrigger calls in flight + +```json +{ "count": 1000, "chunkConcurrency": 2 } +``` + +### Many small batches — 100 chunks of 10, sequential + +```json +{ "count": 1000, "batchSize": 10 } +``` + +### Many small batches in parallel — 100 chunks of 10, 8 in flight + +```json +{ "count": 1000, "batchSize": 10, "chunkConcurrency": 8 } +``` + +### With tags — exercise `runTags ||` contention via the batch path + +```json +{ "count": 1000, "tags": ["stress-batch"] } +``` + +### Children doing real work + +```json +{ "count": 500, "batchSize": 100, "chunkConcurrency": 5, "childSleepMs": 2000 } +``` + +--- + +## What to watch while these run + +- **Axiom** (`['trigger-cloud-prod']` equivalent locally — wherever your local + OTel goes): `prisma:engine:connection` span durations on `trigger-api-prod` + / engine. Baseline is sub-millisecond; > 100 ms is the early signal. +- **Webapp logs**: P2024 ("Timed out fetching a new connection from the + connection pool") and P1001 ("Can't reach database server") surfaces during + the burst. +- **Postgres** (`docker exec database psql -U postgres -d postgres`): + `SELECT count(*) FROM pg_stat_activity;` — connection count under load. +- **Run dashboard**: how many runs queued vs. executing vs. failed; the spread + is what tells you whether the producer-side bottleneck (trigger plumbing) + or the consumer-side bottleneck (worker concurrency) was hit first. diff --git a/references/stress-tasks/package.json b/references/stress-tasks/package.json new file mode 100644 index 000000000..9a3ad1db8 --- /dev/null +++ b/references/stress-tasks/package.json @@ -0,0 +1,17 @@ +{ + "name": "references-stress-tasks", + "private": true, + "type": "module", + "devDependencies": { + "trigger.dev": "workspace:*" + }, + "dependencies": { + "@trigger.dev/build": "workspace:*", + "@trigger.dev/sdk": "workspace:*", + "zod": "3.25.76" + }, + "scripts": { + "dev": "trigger dev", + "deploy": "trigger deploy" + } +} diff --git a/references/stress-tasks/src/trigger/fanout.ts b/references/stress-tasks/src/trigger/fanout.ts new file mode 100644 index 000000000..a8bbaba37 --- /dev/null +++ b/references/stress-tasks/src/trigger/fanout.ts @@ -0,0 +1,273 @@ +import { logger, task } from "@trigger.dev/sdk"; +import { setTimeout as sleep } from "node:timers/promises"; + +/** + * Minimal child task — the fan-out target. The body does nothing meaningful; + * the cost we want to exercise lives in the trigger plumbing on the server. + * + * Optional `sleepMs` lets you keep the child run busy for a while (so concurrent + * children pile up against worker concurrency limits). Optional `pad` is opaque + * data — used by the parent tasks to inflate payload size. + */ +export const noopChildTask = task({ + id: "stress-noop-child", + retry: { maxAttempts: 1 }, + run: async (payload: { index: number; sleepMs?: number; pad?: string }) => { + if (payload.sleepMs && payload.sleepMs > 0) { + await sleep(payload.sleepMs); + } + return { ok: true, index: payload.index }; + }, +}); + +type TriggerOutcome = + | { success: true } + | { success: false; errorName: string; errorMessage: string }; + +/** + * Run an async-task pool. Up to `concurrency` workers pull from a shared cursor. + * Returns results in submission order. Used to cap simultaneous in-flight + * triggers without sequentialising — closer to a real producer with a connection + * pool than `Promise.all` over the full list (which fires everything immediately + * and lets the runtime decide how to interleave). + */ +async function asyncPool( + concurrency: number, + total: number, + produce: (index: number) => Promise, +): Promise { + const results = new Array(total); + let cursor = 0; + const workerCount = Math.max(1, Math.min(concurrency, total)); + const workers = Array.from({ length: workerCount }, async () => { + while (true) { + const i = cursor++; + if (i >= total) return; + results[i] = await produce(i); + } + }); + await Promise.all(workers); + return results; +} + +/** + * Fan-out via N concurrent `.trigger()` calls in a single trace. + * + * This mirrors the production failure mode catalogued in + * `prisma-connection-investigation-results.md` — a single trace fans out + * N HTTP triggers against the webapp api. Run against a local + * `pnpm run dev --filter webapp` to reproduce `prisma:engine:connection` + * acquire-wait spikes and the P2024 / "Can't reach database server" surface. + * + * Parameters: + * count total triggers to fire (default 1000) + * concurrency max simultaneous in-flight triggers (default = count, i.e. all at once) + * childSleepMs sleep duration the child should observe in its body (default 0) + * childPayloadBytes pad each child payload with this many bytes of opaque data (default 0) + * tags tags applied to every child trigger (default []) + * + * Example payloads (copy-paste into the test UI): + * + * @example Smoke test — 10 triggers, all defaults + * { "count": 10 } + * + * @example Reproduce the prod fan-out — 1,000 all at once, single trace + * { "count": 1000 } + * + * @example Bounded producer — 1,000 triggers but only 100 in-flight at any time + * { "count": 1000, "concurrency": 100 } + * + * @example Exercise the `runTags ||` row-lock contention path (events 3, 4, 5, 7) + * { "count": 1000, "tags": ["stress-test", "burst-2026-05-08"] } + * + * @example Children doing real work — 500 triggers, 2s child sleep, 200 in-flight + * { "count": 500, "concurrency": 200, "childSleepMs": 2000 } + * + * @example Large payloads — 200 triggers, 50KB pad each (marshalling pressure) + * { "count": 200, "childPayloadBytes": 50000 } + * + * @example Combined contention — fan-out + tags + child work + * { "count": 1000, "concurrency": 250, "childSleepMs": 500, "tags": ["combined"] } + */ +export const fanOutTriggerTask = task({ + id: "stress-fan-out-trigger", + maxDuration: 600, + retry: { maxAttempts: 1 }, + run: async (payload: { + count?: number; + concurrency?: number; + childSleepMs?: number; + childPayloadBytes?: number; + tags?: string[]; + }) => { + const count = payload.count ?? 1000; + const concurrency = payload.concurrency ?? count; + const childSleepMs = payload.childSleepMs ?? 0; + const childPayloadBytes = payload.childPayloadBytes ?? 0; + const tags = payload.tags ?? []; + + const pad = childPayloadBytes > 0 ? "x".repeat(childPayloadBytes) : undefined; + const triggerOptions = tags.length > 0 ? { tags } : undefined; + + logger.info("Starting fan-out via individual triggers", { + count, + concurrency, + childSleepMs, + childPayloadBytes, + tags, + }); + const start = Date.now(); + + const results = await asyncPool(concurrency, count, async (index) => { + try { + await noopChildTask.trigger( + { index, sleepMs: childSleepMs, pad }, + triggerOptions, + ); + return { success: true }; + } catch (err) { + const e = err as Error; + return { + success: false, + errorName: e?.constructor?.name ?? "Unknown", + errorMessage: e?.message ?? String(err), + }; + } + }); + + const fulfilled = results.filter((r) => r.success).length; + const failures = results.filter( + (r): r is Extract => !r.success, + ); + + const errorCounts: Record = {}; + for (const f of failures) { + errorCounts[f.errorName] = (errorCounts[f.errorName] ?? 0) + 1; + } + + const durationMs = Date.now() - start; + const summary = { + count, + concurrency, + childSleepMs, + childPayloadBytes, + fulfilled, + rejected: failures.length, + durationMs, + triggersPerSecond: + durationMs > 0 ? Math.round((fulfilled / durationMs) * 1000) : 0, + errorCounts, + sampleErrors: failures.slice(0, 5).map((f) => ({ + name: f.errorName, + message: f.errorMessage, + })), + }; + + logger.info("Fan-out complete", summary); + return summary; + }, +}); + +/** + * Fan-out via `batchTrigger`, chunked into `batchSize`-payload calls. + * + * Different server-side code path from `fanOutTriggerTask`: one HTTP + * request per chunk and a server-side bulk insert, vs. N individual API + * round-trips. Useful contrast for understanding whether pool pressure + * is specific to the N-trigger path or shows up here too. + * + * Parameters: + * count total triggers to fire (default 1000) + * batchSize payloads per batchTrigger call (default 500, the SDK default cap) + * chunkConcurrency max simultaneous in-flight batchTrigger calls (default 1, sequential) + * childSleepMs sleep duration the child should observe in its body (default 0) + * childPayloadBytes pad each child payload with this many bytes of opaque data (default 0) + * tags tags applied to every child trigger (default []) + * + * Example payloads (copy-paste into the test UI): + * + * @example Smoke test — single small batch + * { "count": 10, "batchSize": 10 } + * + * @example Default — 1,000 triggers across two sequential 500-payload batches + * { "count": 1000 } + * + * @example Parallel batches — same volume, two batchTrigger calls in flight + * { "count": 1000, "chunkConcurrency": 2 } + * + * @example Many small batches — 100 chunks of 10, sequential + * { "count": 1000, "batchSize": 10 } + * + * @example Many small batches in parallel — 100 chunks of 10, 8 in flight + * { "count": 1000, "batchSize": 10, "chunkConcurrency": 8 } + * + * @example With tags — exercise `runTags ||` contention via the batch path + * { "count": 1000, "tags": ["stress-batch"] } + * + * @example Children doing real work + * { "count": 500, "batchSize": 100, "chunkConcurrency": 5, "childSleepMs": 2000 } + */ +export const fanOutBatchTask = task({ + id: "stress-fan-out-batch", + maxDuration: 600, + retry: { maxAttempts: 1 }, + run: async (payload: { + count?: number; + batchSize?: number; + chunkConcurrency?: number; + childSleepMs?: number; + childPayloadBytes?: number; + tags?: string[]; + }) => { + const count = payload.count ?? 1000; + const batchSize = payload.batchSize ?? 500; + const chunkConcurrency = payload.chunkConcurrency ?? 1; + const childSleepMs = payload.childSleepMs ?? 0; + const childPayloadBytes = payload.childPayloadBytes ?? 0; + const tags = payload.tags ?? []; + + const pad = childPayloadBytes > 0 ? "x".repeat(childPayloadBytes) : undefined; + const itemOptions = tags.length > 0 ? { tags } : undefined; + + logger.info("Starting fan-out via batchTrigger", { + count, + batchSize, + chunkConcurrency, + childSleepMs, + childPayloadBytes, + tags, + }); + const start = Date.now(); + + const chunkCount = Math.ceil(count / batchSize); + const chunks = Array.from({ length: chunkCount }, (_, chunkIndex) => { + const startIdx = chunkIndex * batchSize; + const endIdx = Math.min(startIdx + batchSize, count); + return Array.from({ length: endIdx - startIdx }, (_, k) => ({ + payload: { index: startIdx + k, sleepMs: childSleepMs, pad }, + ...(itemOptions ? { options: itemOptions } : {}), + })); + }); + + const chunkResults = await asyncPool( + chunkConcurrency, + chunkCount, + async (i) => noopChildTask.batchTrigger(chunks[i]), + ); + + const totalCreated = chunkResults.reduce((sum, r) => sum + r.runCount, 0); + const durationMs = Date.now() - start; + const summary = { + count, + batchSize, + chunkConcurrency, + chunkCount, + totalCreated, + durationMs, + triggersPerSecond: + durationMs > 0 ? Math.round((totalCreated / durationMs) * 1000) : 0, + }; + logger.info("Batch fan-out complete", summary); + return summary; + }, +}); diff --git a/references/stress-tasks/trigger.config.ts b/references/stress-tasks/trigger.config.ts new file mode 100644 index 000000000..333c66dd2 --- /dev/null +++ b/references/stress-tasks/trigger.config.ts @@ -0,0 +1,15 @@ +import { defineConfig } from "@trigger.dev/sdk/v3"; + +export default defineConfig({ + compatibilityFlags: ["run_engine_v2"], + project: "proj_stresstaskslocaldevx", + logLevel: "debug", + maxDuration: 3600, + retries: { + enabledInDev: false, + default: { + maxAttempts: 1, + }, + }, + machine: "small-2x", +}); diff --git a/references/stress-tasks/tsconfig.json b/references/stress-tasks/tsconfig.json new file mode 100644 index 000000000..9a5ee0b9d --- /dev/null +++ b/references/stress-tasks/tsconfig.json @@ -0,0 +1,15 @@ +{ + "compilerOptions": { + "target": "ES2023", + "module": "Node16", + "moduleResolution": "Node16", + "esModuleInterop": true, + "strict": true, + "skipLibCheck": true, + "customConditions": ["@triggerdotdev/source"], + "jsx": "preserve", + "lib": ["DOM", "DOM.Iterable"], + "noEmit": true + }, + "include": ["./src/**/*.ts", "trigger.config.ts"] +}