ae46e3f7c8
This PR implements a new run TTL system and queue size limits to prevent unbounded queue growth which should help prevent situations where queues enter a "death spiral" where the queue will never be able to catch up. The main/correct way to battle this situation is to enforce a maximum TTL on all runs (e.g. up to 14 days) where runs that have been queued for that maximum TTL will get auto-expired, making room for newer runs to execute. This required creating a new TTL system that can handle higher workloads and is now deeply integrated into the RunQueue. When runs are enqueued with a TTL, they are added to their normal queue as well as to the TTL queue. When runs are dequeued, they are removed from both their normal queue and the TTL queue. If runs are dequeued by the TTL system, they are removed from their normal queue. Both these dequeues happen automatically so there is no race condition. The TTL expiration system is also made reliable by expiring runs via a Redis worker, which is enqueued to atomically inside the TTL dequeue lua script. ### Optional associated waitpoints Additionally, this PR implements an optimization where runs that aren't triggered with a dependent parent run will no longer create an associated waitpoint. Associated waitpoints are then lazily created if a dependent run wants to wait for the child run post-facto (via debounce or idempotency), which is a rare situation but is possible. This means fewer waitpoint creations but also fewer waitpoint completions for runs with no dependencies. ### Environment Queue Limits Prevents any single queue growing too large by enforcing queue size limits at trigger time. - Queue size checks happen at trigger time - runs are rejected if queue would exceed limit - Dashboard UI shows queue limits on both the Queues page and a new Limits page - In-memory caching for queue size checks to reduce Redis load ### Batch trigger fixes Currently when a batch item cannot be created for whatever reason (e.g. queue limits) the run will never get created, which means a stalled run if using `batchTriggerAndWait`. We've updated the system to handle this differently: now when a batch item cannot be triggered and converted into a run, we will eventually (after retrying 8 times up to 30s) we will create a "pre-failed" run with the error details, correctly resolving the batchTriggerAndWait.
1176 lines
34 KiB
TypeScript
1176 lines
34 KiB
TypeScript
import { describe, expect, vi } from "vitest";
|
|
|
|
// Mock the db prisma client
|
|
vi.mock("~/db.server", () => ({
|
|
prisma: {},
|
|
$replica: {},
|
|
}));
|
|
|
|
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 { IOPacket } from "@trigger.dev/core/v3";
|
|
import { TaskRun } from "@trigger.dev/database";
|
|
import { IdempotencyKeyConcern } from "~/runEngine/concerns/idempotencyKeys.server";
|
|
import { DefaultQueueManager } from "~/runEngine/concerns/queues.server";
|
|
import {
|
|
EntitlementValidationParams,
|
|
MaxAttemptsValidationParams,
|
|
ParentRunValidationParams,
|
|
PayloadProcessor,
|
|
TagValidationParams,
|
|
TracedEventSpan,
|
|
TraceEventConcern,
|
|
TriggerRacepoints,
|
|
TriggerRacepointSystem,
|
|
TriggerTaskRequest,
|
|
TriggerTaskValidator,
|
|
ValidationResult,
|
|
} from "~/runEngine/types";
|
|
import { RunEngineTriggerTaskService } from "../../app/runEngine/services/triggerTask.server";
|
|
import { promiseWithResolvers } from "@trigger.dev/core";
|
|
import { setTimeout } from "node:timers/promises";
|
|
|
|
vi.setConfig({ testTimeout: 60_000 }); // 60 seconds timeout
|
|
|
|
class MockPayloadProcessor implements PayloadProcessor {
|
|
async process(request: TriggerTaskRequest): Promise<IOPacket> {
|
|
return {
|
|
data: JSON.stringify(request.body.payload),
|
|
dataType: "application/json",
|
|
};
|
|
}
|
|
}
|
|
|
|
class MockTriggerTaskValidator implements TriggerTaskValidator {
|
|
validateTags(params: TagValidationParams): ValidationResult {
|
|
return { ok: true };
|
|
}
|
|
validateEntitlement(params: EntitlementValidationParams): Promise<ValidationResult> {
|
|
return Promise.resolve({ ok: true });
|
|
}
|
|
validateMaxAttempts(params: MaxAttemptsValidationParams): ValidationResult {
|
|
return { ok: true };
|
|
}
|
|
validateParentRun(params: ParentRunValidationParams): ValidationResult {
|
|
return { ok: true };
|
|
}
|
|
}
|
|
|
|
class MockTraceEventConcern implements TraceEventConcern {
|
|
async traceRun<T>(
|
|
request: TriggerTaskRequest,
|
|
parentStore: string | undefined,
|
|
callback: (span: TracedEventSpan, store: string) => Promise<T>
|
|
): Promise<T> {
|
|
return await callback(
|
|
{
|
|
traceId: "test",
|
|
spanId: "test",
|
|
traceContext: {},
|
|
traceparent: undefined,
|
|
setAttribute: () => { },
|
|
failWithError: () => { },
|
|
stop: () => { },
|
|
},
|
|
"test"
|
|
);
|
|
}
|
|
|
|
async traceIdempotentRun<T>(
|
|
request: TriggerTaskRequest,
|
|
parentStore: string | undefined,
|
|
options: {
|
|
existingRun: TaskRun;
|
|
idempotencyKey: string;
|
|
incomplete: boolean;
|
|
isError: boolean;
|
|
},
|
|
callback: (span: TracedEventSpan, store: string) => Promise<T>
|
|
): Promise<T> {
|
|
return await callback(
|
|
{
|
|
traceId: "test",
|
|
spanId: "test",
|
|
traceContext: {},
|
|
traceparent: undefined,
|
|
setAttribute: () => { },
|
|
failWithError: () => { },
|
|
stop: () => { },
|
|
},
|
|
"test"
|
|
);
|
|
}
|
|
|
|
async traceDebouncedRun<T>(
|
|
request: TriggerTaskRequest,
|
|
parentStore: string | undefined,
|
|
options: {
|
|
existingRun: TaskRun;
|
|
debounceKey: string;
|
|
incomplete: boolean;
|
|
isError: boolean;
|
|
},
|
|
callback: (span: TracedEventSpan, store: string) => Promise<T>
|
|
): Promise<T> {
|
|
return await callback(
|
|
{
|
|
traceId: "test",
|
|
spanId: "test",
|
|
traceContext: {},
|
|
traceparent: undefined,
|
|
setAttribute: () => { },
|
|
failWithError: () => { },
|
|
stop: () => { },
|
|
},
|
|
"test"
|
|
);
|
|
}
|
|
}
|
|
|
|
type TriggerRacepoint = { promise: Promise<void>; resolve: (value: void) => void };
|
|
|
|
class MockTriggerRacepointSystem implements TriggerRacepointSystem {
|
|
private racepoints: Record<string, TriggerRacepoint | undefined> = {};
|
|
|
|
async waitForRacepoint({ id }: { racepoint: TriggerRacepoints; id: string }): Promise<void> {
|
|
const racepoint = this.racepoints[id];
|
|
|
|
if (racepoint) {
|
|
return racepoint.promise;
|
|
}
|
|
|
|
return Promise.resolve();
|
|
}
|
|
|
|
registerRacepoint(racepoint: TriggerRacepoints, id: string): TriggerRacepoint {
|
|
const { promise, resolve } = promiseWithResolvers<void>();
|
|
this.racepoints[id] = { promise, resolve };
|
|
|
|
return { promise, resolve };
|
|
}
|
|
}
|
|
|
|
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"),
|
|
});
|
|
|
|
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
|
|
|
|
const taskIdentifier = "test-task";
|
|
|
|
//create background worker
|
|
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.findUnique({
|
|
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);
|
|
|
|
await engine.quit();
|
|
});
|
|
|
|
containerTest("should handle idempotency keys correctly", 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"),
|
|
});
|
|
|
|
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
|
|
|
|
const taskIdentifier = "test-task";
|
|
|
|
//create background worker
|
|
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" },
|
|
options: {
|
|
idempotencyKey: "test-idempotency-key",
|
|
},
|
|
},
|
|
});
|
|
|
|
expect(result).toBeDefined();
|
|
expect(result?.run.friendlyId).toBeDefined();
|
|
expect(result?.run.status).toBe("PENDING");
|
|
expect(result?.isCached).toBe(false);
|
|
|
|
const run = await prisma.taskRun.findUnique({
|
|
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);
|
|
|
|
// Now lets try to trigger the same task with the same idempotency key
|
|
const cachedResult = await triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
idempotencyKey: "test-idempotency-key",
|
|
},
|
|
},
|
|
});
|
|
|
|
expect(cachedResult).toBeDefined();
|
|
expect(cachedResult?.run.friendlyId).toBe(result?.run.friendlyId);
|
|
expect(cachedResult?.isCached).toBe(true);
|
|
|
|
await engine.quit();
|
|
});
|
|
|
|
containerTest(
|
|
"should handle idempotency keys when the engine throws an RunDuplicateIdempotencyKeyError",
|
|
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"),
|
|
logLevel: "debug",
|
|
});
|
|
|
|
const parentTask = "parent-task";
|
|
|
|
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
|
|
|
|
const taskIdentifier = "test-task";
|
|
|
|
//create background worker
|
|
await setupBackgroundWorker(engine, authenticatedEnvironment, [parentTask, taskIdentifier]);
|
|
|
|
const parentRun1 = await engine.trigger(
|
|
{
|
|
number: 1,
|
|
friendlyId: "run_p1",
|
|
environment: authenticatedEnvironment,
|
|
taskIdentifier: parentTask,
|
|
payload: "{}",
|
|
payloadType: "application/json",
|
|
context: {},
|
|
traceContext: {},
|
|
traceId: "t12345",
|
|
spanId: "s12345",
|
|
queue: `task/${parentTask}`,
|
|
isTest: false,
|
|
tags: [],
|
|
workerQueue: "main",
|
|
},
|
|
prisma
|
|
);
|
|
|
|
//dequeue parent and create the attempt
|
|
await setTimeout(500);
|
|
const dequeued = await engine.dequeueFromWorkerQueue({
|
|
consumerId: "test_12345",
|
|
workerQueue: "main",
|
|
});
|
|
await engine.startRunAttempt({
|
|
runId: parentRun1.id,
|
|
snapshotId: dequeued[0].snapshot.id,
|
|
});
|
|
|
|
const parentRun2 = await engine.trigger(
|
|
{
|
|
number: 2,
|
|
friendlyId: "run_p2",
|
|
environment: authenticatedEnvironment,
|
|
taskIdentifier: parentTask,
|
|
payload: "{}",
|
|
payloadType: "application/json",
|
|
context: {},
|
|
traceContext: {},
|
|
traceId: "t12346",
|
|
spanId: "s12346",
|
|
queue: `task/${parentTask}`,
|
|
isTest: false,
|
|
tags: [],
|
|
workerQueue: "main",
|
|
},
|
|
prisma
|
|
);
|
|
|
|
await setTimeout(500);
|
|
const dequeued2 = await engine.dequeueFromWorkerQueue({
|
|
consumerId: "test_12345",
|
|
workerQueue: "main",
|
|
});
|
|
await engine.startRunAttempt({
|
|
runId: parentRun2.id,
|
|
snapshotId: dequeued2[0].snapshot.id,
|
|
});
|
|
|
|
const queuesManager = new DefaultQueueManager(prisma, engine);
|
|
|
|
const idempotencyKeyConcern = new IdempotencyKeyConcern(
|
|
prisma,
|
|
engine,
|
|
new MockTraceEventConcern()
|
|
);
|
|
|
|
const triggerRacepointSystem = new MockTriggerRacepointSystem();
|
|
|
|
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
|
|
triggerRacepointSystem,
|
|
});
|
|
|
|
const idempotencyKey = "test-idempotency-key";
|
|
|
|
const racepoint = triggerRacepointSystem.registerRacepoint("idempotencyKey", idempotencyKey);
|
|
|
|
const childTriggerPromise1 = triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
idempotencyKey,
|
|
parentRunId: parentRun1.friendlyId,
|
|
resumeParentOnCompletion: true,
|
|
},
|
|
},
|
|
});
|
|
|
|
const childTriggerPromise2 = triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
idempotencyKey,
|
|
parentRunId: parentRun2.friendlyId,
|
|
resumeParentOnCompletion: true,
|
|
},
|
|
},
|
|
});
|
|
|
|
await setTimeout(500);
|
|
|
|
// Now we can resolve the racepoint
|
|
racepoint.resolve();
|
|
|
|
const result = await childTriggerPromise1;
|
|
const result2 = await childTriggerPromise2;
|
|
|
|
expect(result).toBeDefined();
|
|
expect(result?.run.friendlyId).toBeDefined();
|
|
expect(result?.run.status).toBe("PENDING");
|
|
|
|
const run = await prisma.taskRun.findUnique({
|
|
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}`);
|
|
|
|
expect(result2).toBeDefined();
|
|
expect(result2?.run.friendlyId).toBe(result?.run.friendlyId);
|
|
|
|
const parent1ExecutionData = await engine.getRunExecutionData({ runId: parentRun1.id });
|
|
assertNonNullable(parent1ExecutionData);
|
|
expect(parent1ExecutionData.snapshot.executionStatus).toBe("EXECUTING_WITH_WAITPOINTS");
|
|
|
|
const parent2ExecutionData = await engine.getRunExecutionData({ runId: parentRun2.id });
|
|
assertNonNullable(parent2ExecutionData);
|
|
expect(parent2ExecutionData.snapshot.executionStatus).toBe("EXECUTING_WITH_WAITPOINTS");
|
|
|
|
const parent1RunWaitpoint = await prisma.taskRunWaitpoint.findFirst({
|
|
where: {
|
|
taskRunId: parentRun1.id,
|
|
},
|
|
include: {
|
|
waitpoint: true,
|
|
},
|
|
});
|
|
|
|
assertNonNullable(parent1RunWaitpoint);
|
|
expect(parent1RunWaitpoint.waitpoint.type).toBe("RUN");
|
|
expect(parent1RunWaitpoint.waitpoint.completedByTaskRunId).toBe(result?.run.id);
|
|
|
|
const parent2RunWaitpoint = await prisma.taskRunWaitpoint.findFirst({
|
|
where: {
|
|
taskRunId: parentRun2.id,
|
|
},
|
|
include: {
|
|
waitpoint: true,
|
|
},
|
|
});
|
|
|
|
assertNonNullable(parent2RunWaitpoint);
|
|
expect(parent2RunWaitpoint.waitpoint.type).toBe("RUN");
|
|
expect(parent2RunWaitpoint.waitpoint.completedByTaskRunId).toBe(result2?.run.id);
|
|
|
|
await engine.quit();
|
|
}
|
|
);
|
|
|
|
containerTest(
|
|
"should resolve queue names correctly when locked to version",
|
|
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"),
|
|
});
|
|
|
|
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
|
|
const taskIdentifier = "test-task";
|
|
|
|
// Create a background worker with a specific version
|
|
const worker = await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier, {
|
|
preset: "small-1x",
|
|
});
|
|
|
|
// Create a specific queue for this worker
|
|
const specificQueue = await prisma.taskQueue.create({
|
|
data: {
|
|
name: "specific-queue",
|
|
friendlyId: "specific-queue",
|
|
projectId: authenticatedEnvironment.projectId,
|
|
runtimeEnvironmentId: authenticatedEnvironment.id,
|
|
workers: {
|
|
connect: {
|
|
id: worker.worker.id,
|
|
},
|
|
},
|
|
},
|
|
});
|
|
|
|
// Associate the task with the queue
|
|
await prisma.backgroundWorkerTask.update({
|
|
where: {
|
|
workerId_slug: {
|
|
workerId: worker.worker.id,
|
|
slug: taskIdentifier,
|
|
},
|
|
},
|
|
data: {
|
|
queueId: specificQueue.id,
|
|
},
|
|
});
|
|
|
|
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
|
|
});
|
|
|
|
// Test case 1: Trigger with lockToVersion but no specific queue
|
|
const result1 = await triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
lockToVersion: worker.worker.version,
|
|
},
|
|
},
|
|
});
|
|
|
|
expect(result1).toBeDefined();
|
|
expect(result1?.run.queue).toBe("specific-queue");
|
|
|
|
// Test case 2: Trigger with lockToVersion and specific queue
|
|
const result2 = await triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
lockToVersion: worker.worker.version,
|
|
queue: {
|
|
name: "specific-queue",
|
|
},
|
|
},
|
|
},
|
|
});
|
|
|
|
expect(result2).toBeDefined();
|
|
expect(result2?.run.queue).toBe("specific-queue");
|
|
expect(result2?.run.lockedQueueId).toBe(specificQueue.id);
|
|
|
|
// Test case 3: Try to use non-existent queue with locked version (should throw)
|
|
await expect(
|
|
triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
lockToVersion: worker.worker.version,
|
|
queue: {
|
|
name: "non-existent-queue",
|
|
},
|
|
},
|
|
},
|
|
})
|
|
).rejects.toThrow(
|
|
`Specified queue 'non-existent-queue' not found or not associated with locked version '${worker.worker.version}'`
|
|
);
|
|
|
|
// Test case 4: Trigger with a non-existent queue without a locked version
|
|
const result4 = await triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
queue: {
|
|
name: "non-existent-queue",
|
|
},
|
|
},
|
|
},
|
|
});
|
|
|
|
expect(result4).toBeDefined();
|
|
expect(result4?.run.queue).toBe("non-existent-queue");
|
|
expect(result4?.run.status).toBe("PENDING");
|
|
|
|
await engine.quit();
|
|
}
|
|
);
|
|
|
|
containerTest(
|
|
"should preserve runFriendlyId across retries when RunDuplicateIdempotencyKeyError is thrown",
|
|
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"),
|
|
logLevel: "debug",
|
|
});
|
|
|
|
const parentTask = "parent-task";
|
|
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
|
|
const taskIdentifier = "test-task";
|
|
|
|
// Create background worker
|
|
await setupBackgroundWorker(engine, authenticatedEnvironment, [parentTask, taskIdentifier]);
|
|
|
|
// Create parent runs and start their attempts (required for resumeParentOnCompletion)
|
|
const parentRun1 = await engine.trigger(
|
|
{
|
|
number: 1,
|
|
friendlyId: "run_p1",
|
|
environment: authenticatedEnvironment,
|
|
taskIdentifier: parentTask,
|
|
payload: "{}",
|
|
payloadType: "application/json",
|
|
context: {},
|
|
traceContext: {},
|
|
traceId: "t12345",
|
|
spanId: "s12345",
|
|
queue: `task/${parentTask}`,
|
|
isTest: false,
|
|
tags: [],
|
|
workerQueue: "main",
|
|
},
|
|
prisma
|
|
);
|
|
|
|
await setTimeout(500);
|
|
const dequeued = await engine.dequeueFromWorkerQueue({
|
|
consumerId: "test_12345",
|
|
workerQueue: "main",
|
|
});
|
|
await engine.startRunAttempt({
|
|
runId: parentRun1.id,
|
|
snapshotId: dequeued[0].snapshot.id,
|
|
});
|
|
|
|
const parentRun2 = await engine.trigger(
|
|
{
|
|
number: 2,
|
|
friendlyId: "run_p2",
|
|
environment: authenticatedEnvironment,
|
|
taskIdentifier: parentTask,
|
|
payload: "{}",
|
|
payloadType: "application/json",
|
|
context: {},
|
|
traceContext: {},
|
|
traceId: "t12346",
|
|
spanId: "s12346",
|
|
queue: `task/${parentTask}`,
|
|
isTest: false,
|
|
tags: [],
|
|
workerQueue: "main",
|
|
},
|
|
prisma
|
|
);
|
|
|
|
await setTimeout(500);
|
|
const dequeued2 = await engine.dequeueFromWorkerQueue({
|
|
consumerId: "test_12345",
|
|
workerQueue: "main",
|
|
});
|
|
await engine.startRunAttempt({
|
|
runId: parentRun2.id,
|
|
snapshotId: dequeued2[0].snapshot.id,
|
|
});
|
|
|
|
const queuesManager = new DefaultQueueManager(prisma, engine);
|
|
const idempotencyKeyConcern = new IdempotencyKeyConcern(
|
|
prisma,
|
|
engine,
|
|
new MockTraceEventConcern()
|
|
);
|
|
|
|
const triggerRacepointSystem = new MockTriggerRacepointSystem();
|
|
|
|
// Track all friendlyIds passed to the payload processor
|
|
const processedFriendlyIds: string[] = [];
|
|
class TrackingPayloadProcessor implements PayloadProcessor {
|
|
async process(request: TriggerTaskRequest): Promise<IOPacket> {
|
|
processedFriendlyIds.push(request.friendlyId);
|
|
return {
|
|
data: JSON.stringify(request.body.payload),
|
|
dataType: "application/json",
|
|
};
|
|
}
|
|
}
|
|
|
|
const triggerTaskService = new RunEngineTriggerTaskService({
|
|
engine,
|
|
prisma,
|
|
payloadProcessor: new TrackingPayloadProcessor(),
|
|
queueConcern: queuesManager,
|
|
idempotencyKeyConcern,
|
|
validator: new MockTriggerTaskValidator(),
|
|
traceEventConcern: new MockTraceEventConcern(),
|
|
tracer: trace.getTracer("test", "0.0.0"),
|
|
metadataMaximumSize: 1024 * 1024 * 1, // 1MB
|
|
triggerRacepointSystem,
|
|
});
|
|
|
|
const idempotencyKey = "test-preserve-friendly-id";
|
|
const racepoint = triggerRacepointSystem.registerRacepoint("idempotencyKey", idempotencyKey);
|
|
|
|
// Trigger two concurrent requests with same idempotency key
|
|
// One will succeed, one will fail with RunDuplicateIdempotencyKeyError and retry
|
|
const childTriggerPromise1 = triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test1" },
|
|
options: {
|
|
idempotencyKey,
|
|
parentRunId: parentRun1.friendlyId,
|
|
resumeParentOnCompletion: true,
|
|
},
|
|
},
|
|
});
|
|
|
|
const childTriggerPromise2 = triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test2" },
|
|
options: {
|
|
idempotencyKey,
|
|
parentRunId: parentRun2.friendlyId,
|
|
resumeParentOnCompletion: true,
|
|
},
|
|
},
|
|
});
|
|
|
|
await setTimeout(500);
|
|
|
|
// Resolve the racepoint to allow both requests to proceed
|
|
racepoint.resolve();
|
|
|
|
const result1 = await childTriggerPromise1;
|
|
const result2 = await childTriggerPromise2;
|
|
|
|
// Both should return the same run (one created, one cached)
|
|
expect(result1).toBeDefined();
|
|
expect(result2).toBeDefined();
|
|
expect(result1?.run.friendlyId).toBe(result2?.run.friendlyId);
|
|
|
|
// The key assertion: When a retry happens due to RunDuplicateIdempotencyKeyError,
|
|
// the same friendlyId should be used. We expect exactly 2 calls to payloadProcessor
|
|
// (one for each concurrent request), not 3 (which would indicate a new friendlyId on retry)
|
|
// Since the retry returns early from the idempotency cache, payloadProcessor is not called again.
|
|
expect(processedFriendlyIds.length).toBe(2);
|
|
|
|
// Verify that we have exactly 2 unique friendlyIds (one per original request)
|
|
const uniqueFriendlyIds = new Set(processedFriendlyIds);
|
|
expect(uniqueFriendlyIds.size).toBe(2);
|
|
|
|
await engine.quit();
|
|
}
|
|
);
|
|
|
|
containerTest(
|
|
"should reject invalid debounce.delay when no explicit delay is provided",
|
|
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"),
|
|
});
|
|
|
|
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,
|
|
});
|
|
|
|
// Invalid debounce.delay format (ms not supported)
|
|
await expect(
|
|
triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
debounce: {
|
|
key: "test-key",
|
|
delay: "300ms", // Invalid - ms not supported
|
|
},
|
|
},
|
|
},
|
|
})
|
|
).rejects.toThrow("Debounce requires a valid delay duration");
|
|
|
|
await engine.quit();
|
|
}
|
|
);
|
|
|
|
containerTest(
|
|
"should reject invalid debounce.delay even when explicit delay is valid",
|
|
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"),
|
|
});
|
|
|
|
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,
|
|
});
|
|
|
|
// Valid explicit delay but invalid debounce.delay
|
|
// This is the bug case: the explicit delay passes validation,
|
|
// but debounce.delay would fail later when rescheduling
|
|
await expect(
|
|
triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
delay: "5m", // Valid explicit delay
|
|
debounce: {
|
|
key: "test-key",
|
|
delay: "invalid-delay", // Invalid debounce delay
|
|
},
|
|
},
|
|
},
|
|
})
|
|
).rejects.toThrow("Invalid debounce delay");
|
|
|
|
await engine.quit();
|
|
}
|
|
);
|
|
|
|
containerTest(
|
|
"should accept valid debounce.delay formats",
|
|
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"),
|
|
});
|
|
|
|
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,
|
|
});
|
|
|
|
// Valid debounce.delay format
|
|
const result = await triggerTaskService.call({
|
|
taskId: taskIdentifier,
|
|
environment: authenticatedEnvironment,
|
|
body: {
|
|
payload: { test: "test" },
|
|
options: {
|
|
debounce: {
|
|
key: "test-key",
|
|
delay: "5s", // Valid format
|
|
},
|
|
},
|
|
},
|
|
});
|
|
|
|
expect(result).toBeDefined();
|
|
expect(result?.run.friendlyId).toBeDefined();
|
|
|
|
await engine.quit();
|
|
}
|
|
);
|
|
});
|