Files
triggerdotdev--trigger.dev/apps/webapp/test/engine/triggerTask.test.ts
Daniel Sutton 4c2c25511b test(webapp): split triggerTask engine test into per-concern files (#4167)
The engine `triggerTask` suite was a single 2447-line file with 23
`containerTest` cases, each spinning its own Postgres + Redis. vitest
shards by whole file, so all 23 container setups landed on one shard and
dominated its wall-clock. The recorded entry in `test-timings.json`
badly under-counts the real cost (it does not capture the
per-`containerTest` container startup that dominates on CI), so the
duration-sharding sequencer treated the file as light and stacked it,
producing one ~21 minute shard.

Splitting does not reduce the number of container setups; it lets those
23 cases distribute across shards instead of stacking on one. The webapp
unit-test stage is gated by its slowest shard, so this cuts the stage's
wall-clock roughly in half.

## CI timing (before vs after)

Real CI wall-clock of the `Unit Tests: Webapp` shards (`--shard=i/10`).
"Before" is sampled from recent runs on other branches (unsplit file,
from `main`); "after" is this PR.

| Shard | Before (s) | After (s) |
|------:|-----------:|----------:|
| 1  | 250 | 359 |
| 2  | 444 | 411 |
| 3  | 497 | 659 |
| 4  | **1257** | 284 |
| 5  | 545 | 641 |
| 6  | 284 | 644 |
| 7  | 244 | 214 |
| 8  | 340 | 445 |
| 9  | 188 | 395 |
| 10 | 234 | 567 |
| **Slowest shard (gates the stage)** | **~1247s (≈21m)** | **659s
(≈11m)** |
| Sum of all shards | 4283 | 4619 |

Before: shard 4 is the long pole at 1237s / 1247s / 1257s across three
sampled runs (the `triggerTask` file plus whatever else the packer put
with it). After: the six pieces spread across shards, the slowest drops
to 659s. The small rise in summed time is the extra per-file container
startup, paid in parallel across shards, so the gating number still
falls by about 10 minutes.

## Change

Split into six per-concern files that share a `triggerTaskTestHelpers`
module (the `vi.mock` calls stay per-file, since vitest hoists them):

- `triggerTask.test.ts` (3): trigger + concurrencyKey coercion
- `triggerTask.idempotency.test.ts` (4): idempotency + queue resolution
- `triggerTask.debounce.test.ts` (4): retries + debounce validation
- `triggerTask.mollifier.test.ts` (4): mollifier call-site behaviour
- `triggerTask.metadataCache.test.ts` (4): DefaultQueueManager task
metadata cache
- `triggerTask.residency.test.ts` (4): child run residency inheritance

All 23 cases are preserved. The file's `test-timings.json` entry is
split across the new files so bin-packing stays balanced.

While rewriting these files, cleanup was moved to `onTestFinished(() =>
engine.quit())` so an `engine`/`Redis` leaked on a failing assertion no
longer persists on the worker-scoped Redis and cascades into later cases
(`hookTimeout` raised to 60s so the after-cleanup gets the full budget).
Prisma lookups switched from `findUnique` to `findFirst` to match the
repo convention.

Verified: all six files run green locally (23/23), oxlint and oxfmt
clean.
2026-07-06 13:03:27 +01:00

314 lines
11 KiB
TypeScript

import { describe, expect, onTestFinished, vi } from "vitest";
// db.server + splitMode are mocked so the idempotency dedup client resolves to
// the container prisma passed into the concern (split stays off).
vi.mock("~/db.server", () => ({
prisma: {},
$replica: {},
runOpsNewPrisma: {},
runOpsLegacyPrisma: {},
}));
vi.mock("~/v3/runOpsMigration/splitMode.server", () => ({ isSplitEnabled: async () => false }));
vi.mock("~/services/platform.v3.server", async (importOriginal) => {
const actual = (await importOriginal()) as Record<string, unknown>;
return {
...actual,
getEntitlement: vi.fn(),
};
});
import { RunEngine } from "@internal/run-engine";
import { setupAuthenticatedEnvironment, setupBackgroundWorker } from "@internal/run-engine/tests";
import { assertNonNullable, containerTest } from "@internal/testcontainers";
import { trace } from "@opentelemetry/api";
import { IdempotencyKeyConcern } from "~/runEngine/concerns/idempotencyKeys.server";
import { DefaultQueueManager } from "~/runEngine/concerns/queues.server";
import { RunEngineTriggerTaskService } from "../../app/runEngine/services/triggerTask.server";
import { setTimeout } from "node:timers/promises";
import {
MockPayloadProcessor,
MockTraceEventConcern,
MockTriggerTaskValidator,
} from "./triggerTaskTestHelpers";
vi.setConfig({ testTimeout: 60_000, hookTimeout: 60_000 });
describe("RunEngineTriggerTaskService", () => {
containerTest("should trigger a task with minimal options", async ({ prisma, redisOptions }) => {
const engine = new RunEngine({
prisma,
worker: {
redis: redisOptions,
workers: 1,
tasksPerWorker: 10,
pollIntervalMs: 100,
},
queue: {
redis: redisOptions,
},
runLock: {
redis: redisOptions,
},
machines: {
defaultMachine: "small-1x",
machines: {
"small-1x": {
name: "small-1x" as const,
cpu: 0.5,
memory: 0.5,
centsPerMs: 0.0001,
},
},
baseCostInCents: 0.0005,
},
tracer: trace.getTracer("test", "0.0.0"),
});
onTestFinished(() => engine.quit());
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
const taskIdentifier = "test-task";
await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier);
const queuesManager = new DefaultQueueManager(prisma, engine);
const idempotencyKeyConcern = new IdempotencyKeyConcern(
prisma,
engine,
new MockTraceEventConcern()
);
const triggerTaskService = new RunEngineTriggerTaskService({
engine,
prisma,
payloadProcessor: new MockPayloadProcessor(),
queueConcern: queuesManager,
idempotencyKeyConcern,
validator: new MockTriggerTaskValidator(),
traceEventConcern: new MockTraceEventConcern(),
tracer: trace.getTracer("test", "0.0.0"),
metadataMaximumSize: 1024 * 1024 * 1, // 1MB
});
const result = await triggerTaskService.call({
taskId: taskIdentifier,
environment: authenticatedEnvironment,
body: { payload: { test: "test" } },
});
expect(result).toBeDefined();
expect(result?.run.friendlyId).toBeDefined();
expect(result?.run.status).toBe("PENDING");
expect(result?.isCached).toBe(false);
const run = await prisma.taskRun.findFirst({
where: {
id: result?.run.id,
},
});
expect(run).toBeDefined();
expect(run?.friendlyId).toBe(result?.run.friendlyId);
expect(run?.engine).toBe("V2");
expect(run?.queuedAt).toBeDefined();
expect(run?.queue).toBe(`task/${taskIdentifier}`);
// Lets make sure the task is in the queue
const queueLength = await engine.runQueue.lengthOfQueue(
authenticatedEnvironment,
`task/${taskIdentifier}`
);
expect(queueLength).toBe(1);
});
containerTest(
"routes scheduled-lineage runs to a separate worker queue that dequeues independently",
async ({ prisma, redisOptions }) => {
const engine = new RunEngine({
prisma,
worker: {
redis: redisOptions,
workers: 1,
tasksPerWorker: 10,
pollIntervalMs: 100,
},
queue: {
redis: redisOptions,
// Disable the background master-queue consumers so our manual
// processMasterQueueForEnvironment + dequeue calls are deterministic.
masterQueueConsumersDisabled: true,
processWorkerQueueDebounceMs: 50,
},
runLock: { redis: redisOptions },
machines: {
defaultMachine: "small-1x",
machines: {
"small-1x": {
name: "small-1x" as const,
cpu: 0.5,
memory: 0.5,
centsPerMs: 0.0001,
},
},
baseCostInCents: 0.0005,
},
tracer: trace.getTracer("test", "0.0.0"),
});
try {
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
// Turn the per-org split flag on in-memory — the resolver reads this
// object directly (no DB round-trip on the trigger hot path).
(authenticatedEnvironment.organization as { featureFlags?: unknown }).featureFlags = {
workerQueueScheduledSplitEnabled: true,
};
const taskIdentifier = "test-task";
await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier);
const triggerTaskService = new RunEngineTriggerTaskService({
engine,
prisma,
payloadProcessor: new MockPayloadProcessor(),
queueConcern: new DefaultQueueManager(prisma, engine),
idempotencyKeyConcern: new IdempotencyKeyConcern(
prisma,
engine,
new MockTraceEventConcern()
),
validator: new MockTriggerTaskValidator(),
traceEventConcern: new MockTraceEventConcern(),
tracer: trace.getTracer("test", "0.0.0"),
metadataMaximumSize: 1024 * 1024 * 1,
});
// A standard run (default triggerSource) stays on the region queue.
const standardResult = await triggerTaskService.call({
taskId: taskIdentifier,
environment: authenticatedEnvironment,
body: { payload: { kind: "standard" } },
});
assertNonNullable(standardResult);
// A scheduled run routes to the `<region>:scheduled` queue. Descendants
// would too, via rootTriggerSource propagation.
const scheduledResult = await triggerTaskService.call({
taskId: taskIdentifier,
environment: authenticatedEnvironment,
body: { payload: { kind: "scheduled" } },
options: { triggerSource: "schedule" },
});
assertNonNullable(scheduledResult);
const standardRun = await prisma.taskRun.findUniqueOrThrow({
where: { id: standardResult.run.id },
});
const scheduledRun = await prisma.taskRun.findUniqueOrThrow({
where: { id: scheduledResult.run.id },
});
// Producer routing: the persisted worker queue carries the class.
const baseWorkerQueue = standardRun.workerQueue;
expect(scheduledRun.workerQueue).toBe(`${baseWorkerQueue}:scheduled`);
// Move both runs from the env queue onto their respective worker queues.
await engine.runQueue.processMasterQueueForEnvironment(authenticatedEnvironment.id, 10);
await setTimeout(500);
// Dequeue isolation: the scheduled queue yields only the scheduled run...
const dequeuedScheduled = await engine.dequeueFromWorkerQueue({
consumerId: "test-scheduled-consumer",
workerQueue: `${baseWorkerQueue}:scheduled`,
});
expect(dequeuedScheduled.length).toBe(1);
assertNonNullable(dequeuedScheduled[0]);
expect(dequeuedScheduled[0].run.id).toBe(scheduledResult.run.id);
// ...and the base queue yields only the standard run.
const dequeuedStandard = await engine.dequeueFromWorkerQueue({
consumerId: "test-standard-consumer",
workerQueue: baseWorkerQueue,
});
expect(dequeuedStandard.length).toBe(1);
assertNonNullable(dequeuedStandard[0]);
expect(dequeuedStandard[0].run.id).toBe(standardResult.run.id);
} finally {
await engine.quit();
}
}
);
// The BatchQueue worker rebuilds body.options from Redis-stored items
// (Record<string, unknown>), so the Phase-2 schema coercion doesn't apply
// to in-flight items enqueued before the schema fix. The defensive
// `typeof === "number"` coercion at the engine.trigger call site is what
// prevents these from failing at prisma.taskRun.create with
// "Argument concurrencyKey: Expected String or Null, provided Int".
containerTest(
"coerces a numeric concurrencyKey to a string at the engine.trigger boundary",
async ({ prisma, redisOptions }) => {
const engine = new RunEngine({
prisma,
worker: {
redis: redisOptions,
workers: 1,
tasksPerWorker: 10,
pollIntervalMs: 100,
},
queue: { redis: redisOptions },
runLock: { redis: redisOptions },
machines: {
defaultMachine: "small-1x",
machines: {
"small-1x": {
name: "small-1x" as const,
cpu: 0.5,
memory: 0.5,
centsPerMs: 0.0001,
},
},
baseCostInCents: 0.0005,
},
tracer: trace.getTracer("test", "0.0.0"),
});
onTestFinished(() => engine.quit());
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
const taskIdentifier = "test-task";
await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier);
const triggerTaskService = new RunEngineTriggerTaskService({
engine,
prisma,
payloadProcessor: new MockPayloadProcessor(),
queueConcern: new DefaultQueueManager(prisma, engine),
idempotencyKeyConcern: new IdempotencyKeyConcern(
prisma,
engine,
new MockTraceEventConcern()
),
validator: new MockTriggerTaskValidator(),
traceEventConcern: new MockTraceEventConcern(),
tracer: trace.getTracer("test", "0.0.0"),
metadataMaximumSize: 1024 * 1024 * 1,
});
const result = await triggerTaskService.call({
taskId: taskIdentifier,
environment: authenticatedEnvironment,
// Cast through `any` to simulate the in-flight Redis batch-item shape
// (Record<string, unknown>) that bypasses the BatchItemNDJSON schema.
body: { payload: { userId: 51262 }, options: { concurrencyKey: 51262 as any } },
});
expect(result).toBeDefined();
const run = await prisma.taskRun.findFirst({ where: { id: result!.run.id } });
expect(run?.concurrencyKey).toBe("51262");
}
);
});