diff --git a/.server-changes/collapse-trigger-queued-snapshot.md b/.server-changes/collapse-trigger-queued-snapshot.md new file mode 100644 index 0000000000..541c5b6450 --- /dev/null +++ b/.server-changes/collapse-trigger-queued-snapshot.md @@ -0,0 +1,6 @@ +--- +area: webapp +type: improvement +--- + +Triggering a task now reaches the queue marginally faster. diff --git a/internal-packages/run-engine/src/engine/index.ts b/internal-packages/run-engine/src/engine/index.ts index b513176655..61d60695ee 100644 --- a/internal-packages/run-engine/src/engine/index.ts +++ b/internal-packages/run-engine/src/engine/index.ts @@ -23,6 +23,7 @@ import { } from "@trigger.dev/core/v3"; import type { TaskRunError } from "@trigger.dev/core/v3/schemas"; import { + generateInternalId, parseNaturalLanguageDurationInMs, RunId, WaitpointId, @@ -947,6 +948,7 @@ export class RunEngine { let taskRun: TaskRun & { associatedWaitpoint: Waitpoint | null }; const taskRunId = RunId.fromFriendlyId(friendlyId); + const initialSnapshotId = generateInternalId(); // App-level replacement for the dropped TaskRun env/project Cascade FKs. await this.controlPlaneResolver.assertEnvExists(environment.id); @@ -1032,9 +1034,10 @@ export class RunEngine { annotations, }, snapshot: { + id: initialSnapshotId, engine: "V2", - executionStatus: delayUntil ? "DELAYED" : "RUN_CREATED", - description: delayUntil ? "Run is delayed" : "Run was created", + executionStatus: delayUntil ? "DELAYED" : "QUEUED", + description: delayUntil ? "Run is delayed" : "Run was QUEUED", runStatus: status, environmentId: environment.id, environmentType: environment.type, @@ -1161,13 +1164,29 @@ export class RunEngine { await this.ttlSystem.scheduleExpireRun({ runId: taskRun.id, ttl: taskRun.ttl }); } - await this.enqueueSystem.enqueueRun({ + this.eventBus.emit("executionSnapshotCreated", { + time: new Date(), + run: { + id: taskRun.id, + }, + snapshot: { + id: initialSnapshotId, + executionStatus: "QUEUED", + description: "Run was QUEUED", + runStatus: taskRun.status, + attemptNumber: taskRun.attemptNumber ?? null, + checkpointId: null, + workerId: workerId ?? null, + runnerId: runnerId ?? null, + isValid: true, + error: null, + completedWaitpointIds: [], + }, + }); + + await this.enqueueSystem.publishRun({ run: taskRun, env: environment, - workerId, - runnerId, - tx: prisma, - skipRunLock: true, includeTtl: true, anchorEligibilityAtQueuePosition: true, enableFastPath, diff --git a/internal-packages/run-engine/src/engine/systems/enqueueSystem.test.ts b/internal-packages/run-engine/src/engine/systems/enqueueSystem.test.ts index a5730b306f..eca893712c 100644 --- a/internal-packages/run-engine/src/engine/systems/enqueueSystem.test.ts +++ b/internal-packages/run-engine/src/engine/systems/enqueueSystem.test.ts @@ -44,10 +44,10 @@ function createEngineOptions(redisOptions: any, prisma: any, store?: PostgresRun } /** - * A real PostgresRunStore subclass that counts the snapshot create method that enqueueRun's - * snapshot write routes through (via executionSnapshotSystem.createExecutionSnapshot). super.* - * runs the genuine store implementation, so the routing is observed over real containers without - * ever mocking prisma or the store. + * A real PostgresRunStore subclass that counts the store methods a run's QUEUED snapshot can be + * written through: nested in `createRun` on the trigger path, or standalone via + * `createExecutionSnapshot` on every re-enqueue. super.* runs the genuine store implementation, so + * the routing is observed over real containers without ever mocking prisma or the store. */ class CountingPostgresRunStore extends PostgresRunStore { public snapshotCreates = 0; @@ -59,12 +59,19 @@ class CountingPostgresRunStore extends PostgresRunStore { this.snapshotCreates++; return super.createExecutionSnapshot(input, tx); } + + override async createRun( + params: Parameters[0], + tx?: any + ): ReturnType { + this.snapshotCreates++; + return super.createRun(params, tx); + } } describe("RunEngine enqueueRun store routing", () => { - // The QUEUED snapshot written while enqueuing a run routes through the injected store. containerTest( - "enqueueRun snapshot routes through the store", + "the QUEUED snapshot routes through the store", async ({ prisma, redisOptions }) => { const countingStore = new CountingPostgresRunStore({ prisma, readOnlyPrisma: prisma }); const engine = new RunEngine(createEngineOptions(redisOptions, prisma, countingStore)); diff --git a/internal-packages/run-engine/src/engine/systems/enqueueSystem.ts b/internal-packages/run-engine/src/engine/systems/enqueueSystem.ts index 2b5aac6d37..75382ce82c 100644 --- a/internal-packages/run-engine/src/engine/systems/enqueueSystem.ts +++ b/internal-packages/run-engine/src/engine/systems/enqueueSystem.ts @@ -113,42 +113,74 @@ export class EnqueueSystem { store ); - // Force development runs to use the environment id as the worker queue. - const workerQueue = env.type === "DEVELOPMENT" ? env.id : run.workerQueue; - - const queuePositionMs = (run.queueTimestamp ?? run.createdAt).getTime(); - const timestamp = queuePositionMs - run.priorityMs; - const eligibleAtMs = anchorEligibilityAtQueuePosition ? queuePositionMs : Date.now(); - - let ttlExpiresAt: number | undefined; - if (includeTtl && run.ttl) { - const expireAt = parseNaturalLanguageDuration(run.ttl); - if (expireAt) { - ttlExpiresAt = expireAt.getTime(); - } - } - - await this.$.runQueue.enqueueMessage({ + await this.publishRun({ + run, env, - workerQueue, + includeTtl, + anchorEligibilityAtQueuePosition, enableFastPath, - message: { - runId: run.id, - taskIdentifier: run.taskIdentifier, - orgId: env.organization.id, - projectId: env.project.id, - environmentId: env.id, - environmentType: env.type, - queue: run.queue, - concurrencyKey: run.concurrencyKey ?? undefined, - timestamp, - eligibleAtMs, - attempt: 0, - ttlExpiresAt, - }, }); return newSnapshot; }); } + + /** + * Publishes the run to the RunQueue without writing an execution snapshot. Callers that already + * hold a `QUEUED` snapshot (the trigger path writes one inside the run-create transaction) use + * this so the run does not pay for a second snapshot write. + */ + public async publishRun({ + run, + env, + includeTtl = false, + anchorEligibilityAtQueuePosition = false, + enableFastPath = false, + }: { + run: TaskRun; + env: MinimalAuthenticatedEnvironment; + /** See `enqueueRun`. */ + includeTtl?: boolean; + /** See `enqueueRun`. */ + anchorEligibilityAtQueuePosition?: boolean; + /** When true, allow the queue to push directly to worker queue if concurrency is available. */ + enableFastPath?: boolean; + }) { + // Force development runs to use the environment id as the worker queue. + const workerQueue = env.type === "DEVELOPMENT" ? env.id : run.workerQueue; + + const queuePositionMs = (run.queueTimestamp ?? run.createdAt).getTime(); + const timestamp = queuePositionMs - run.priorityMs; + const eligibleAtMs = anchorEligibilityAtQueuePosition ? queuePositionMs : Date.now(); + + // Include TTL only when explicitly requested (first enqueue from trigger). + // Re-enqueues (waitpoint, checkpoint, delayed, pending version) must not add TTL. + let ttlExpiresAt: number | undefined; + if (includeTtl && run.ttl) { + const expireAt = parseNaturalLanguageDuration(run.ttl); + if (expireAt) { + ttlExpiresAt = expireAt.getTime(); + } + } + + await this.$.runQueue.enqueueMessage({ + env, + workerQueue, + enableFastPath, + message: { + runId: run.id, + taskIdentifier: run.taskIdentifier, + orgId: env.organization.id, + projectId: env.project.id, + environmentId: env.id, + environmentType: env.type, + queue: run.queue, + concurrencyKey: run.concurrencyKey ?? undefined, + timestamp, + eligibleAtMs, + attempt: 0, + ttlExpiresAt, + }, + }); + } } diff --git a/internal-packages/run-engine/src/engine/systems/executionSnapshotSystem.test.ts b/internal-packages/run-engine/src/engine/systems/executionSnapshotSystem.test.ts index aeb361019e..61d17a6b8b 100644 --- a/internal-packages/run-engine/src/engine/systems/executionSnapshotSystem.test.ts +++ b/internal-packages/run-engine/src/engine/systems/executionSnapshotSystem.test.ts @@ -65,6 +65,14 @@ class CountingPostgresRunStore extends PostgresRunStore { return super.createExecutionSnapshot(input, tx); } + override async createRun( + params: Parameters[0], + tx?: any + ): ReturnType { + this.creates++; + return super.createRun(params, tx); + } + override async findLatestExecutionSnapshot( runId: string, client?: any diff --git a/internal-packages/run-engine/src/engine/tests/getSnapshotsSince.test.ts b/internal-packages/run-engine/src/engine/tests/getSnapshotsSince.test.ts index 61ff3bff75..4cf74753e0 100644 --- a/internal-packages/run-engine/src/engine/tests/getSnapshotsSince.test.ts +++ b/internal-packages/run-engine/src/engine/tests/getSnapshotsSince.test.ts @@ -7,7 +7,7 @@ import type { PrismaClient } from "@trigger.dev/database"; import { RunEngine } from "../index.js"; import { getExecutionSnapshotsSince } from "../systems/executionSnapshotSystem.js"; import { copySnapshotsToReplica, createTestMetricsMeter } from "./helpers/replicaTestHelpers.js"; -import { setupTestScenario } from "./helpers/snapshotTestHelpers.js"; +import { createTestSnapshot, setupTestScenario } from "./helpers/snapshotTestHelpers.js"; import { setupAuthenticatedEnvironment, setupBackgroundWorker } from "./setup.js"; vi.setConfig({ testTimeout: 120_000 }); @@ -1154,6 +1154,15 @@ describe("RunEngine getSnapshotsSince", () => { workerQueue: "main", }); + await createTestSnapshot(prisma, { + runId: run.id, + status: "EXECUTING", + environmentId: authenticatedEnvironment.id, + environmentType: authenticatedEnvironment.type, + projectId: authenticatedEnvironment.project.id, + organizationId: authenticatedEnvironment.organization.id, + }); + const allSnapshots = await prisma.taskRunExecutionSnapshot.findMany({ where: { runId: run.id, isValid: true }, orderBy: { createdAt: "asc" }, diff --git a/internal-packages/run-engine/src/engine/tests/triggerCreateRouting.test.ts b/internal-packages/run-engine/src/engine/tests/triggerCreateRouting.test.ts index 34353cb44e..af213dc7a7 100644 --- a/internal-packages/run-engine/src/engine/tests/triggerCreateRouting.test.ts +++ b/internal-packages/run-engine/src/engine/tests/triggerCreateRouting.test.ts @@ -137,9 +137,6 @@ const cancelledSnapshot = (friendlyId: string, environment: any) => ({ }); describe("RunEngine trigger/create routing", () => { - // trigger create routes through runStore.createRun with the structured - // DTO, and the persisted run + its nested first RUN_CREATED snapshot land via - // the single create call. containerTest( "trigger routes createRun and lands run + first snapshot", async ({ prisma, redisOptions }) => { @@ -169,7 +166,7 @@ describe("RunEngine trigger/create routing", () => { orderBy: { createdAt: "asc" }, }); expect(snapshot).not.toBeNull(); - expect(snapshot!.executionStatus).toBe("RUN_CREATED"); + expect(snapshot!.executionStatus).toBe("QUEUED"); } finally { await engine.quit(); } diff --git a/internal-packages/run-engine/src/engine/tests/triggerSnapshotCollapse.test.ts b/internal-packages/run-engine/src/engine/tests/triggerSnapshotCollapse.test.ts new file mode 100644 index 0000000000..cddb2756c3 --- /dev/null +++ b/internal-packages/run-engine/src/engine/tests/triggerSnapshotCollapse.test.ts @@ -0,0 +1,256 @@ +import { assertNonNullable, containerTest } from "@internal/testcontainers"; +import { trace } from "@internal/tracing"; +import type { PrismaClient, TaskRunExecutionStatus } from "@trigger.dev/database"; +import { setTimeout } from "node:timers/promises"; +import { expect } from "vitest"; +import { RunEngine } from "../index.js"; +import { setupAuthenticatedEnvironment, setupBackgroundWorker } from "./setup.js"; + +vi.setConfig({ testTimeout: 60_000 }); + +async function snapshotStatuses( + prisma: PrismaClient, + runId: string +): Promise { + const snapshots = await prisma.taskRunExecutionSnapshot.findMany({ + where: { runId }, + orderBy: { createdAt: "asc" }, + select: { executionStatus: true }, + }); + + return snapshots.map((snapshot) => snapshot.executionStatus); +} + +describe("RunEngine trigger() execution snapshots", () => { + containerTest( + "a non-delayed run is created with a single QUEUED snapshot", + async ({ prisma, redisOptions }) => { + const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION"); + + const engine = new RunEngine({ + prisma, + worker: { + redis: redisOptions, + workers: 1, + tasksPerWorker: 10, + pollIntervalMs: 100, + }, + queue: { + redis: redisOptions, + 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.0001, + }, + tracer: trace.getTracer("test", "0.0.0"), + }); + + try { + const taskIdentifier = "test-task"; + + await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier); + + const run = await engine.trigger( + { + number: 1, + friendlyId: "run_1234", + environment: authenticatedEnvironment, + taskIdentifier, + payload: "{}", + payloadType: "application/json", + context: {}, + traceContext: {}, + traceId: "t_collapse_1", + spanId: "s_collapse_1", + workerQueue: "main", + queue: "task/test-task", + isTest: false, + tags: [], + }, + prisma + ); + + expect(await snapshotStatuses(prisma, run.id)).toEqual(["QUEUED"]); + + const queueLength = await engine.runQueue.lengthOfQueue( + authenticatedEnvironment, + run.queue + ); + expect(queueLength).toBe(1); + + const executionData = await engine.getRunExecutionData({ runId: run.id }); + assertNonNullable(executionData); + expect(executionData.snapshot.executionStatus).toBe("QUEUED"); + } finally { + await engine.quit(); + } + } + ); + + containerTest( + "the QUEUED snapshot event is stamped at write time, not at an overridden run createdAt", + async ({ prisma, redisOptions }) => { + const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION"); + + const engine = new RunEngine({ + prisma, + worker: { + redis: redisOptions, + workers: 1, + tasksPerWorker: 10, + pollIntervalMs: 100, + }, + queue: { + redis: redisOptions, + 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.0001, + }, + tracer: trace.getTracer("test", "0.0.0"), + }); + + try { + const taskIdentifier = "test-task"; + + await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier); + + const stampedTimes: Date[] = []; + engine.eventBus.on("executionSnapshotCreated", ({ time }) => { + stampedTimes.push(time); + }); + + const backdatedCreatedAt = new Date(Date.now() - 60 * 60 * 1000); + const triggeredAt = Date.now(); + + const run = await engine.trigger( + { + number: 1, + friendlyId: "run_1236", + environment: authenticatedEnvironment, + taskIdentifier, + payload: "{}", + payloadType: "application/json", + context: {}, + traceContext: {}, + traceId: "t_collapse_3", + spanId: "s_collapse_3", + workerQueue: "main", + queue: "task/test-task", + isTest: false, + tags: [], + createdAt: backdatedCreatedAt, + }, + prisma + ); + + const storedRun = await prisma.taskRun.findUnique({ where: { id: run.id } }); + expect(storedRun?.createdAt.getTime()).toBe(backdatedCreatedAt.getTime()); + + expect(stampedTimes.length).toBe(1); + expect(stampedTimes[0].getTime()).toBeGreaterThanOrEqual(triggeredAt); + } finally { + await engine.quit(); + } + } + ); + + containerTest( + "a delayed run keeps DELAYED and QUEUED as separate snapshots", + async ({ prisma, redisOptions }) => { + const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION"); + + const engine = new RunEngine({ + prisma, + worker: { + redis: redisOptions, + workers: 1, + tasksPerWorker: 10, + pollIntervalMs: 100, + }, + queue: { + redis: redisOptions, + 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.0001, + }, + tracer: trace.getTracer("test", "0.0.0"), + }); + + try { + const taskIdentifier = "test-task"; + + await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier); + + const run = await engine.trigger( + { + number: 1, + friendlyId: "run_1235", + environment: authenticatedEnvironment, + taskIdentifier, + payload: "{}", + payloadType: "application/json", + context: {}, + traceContext: {}, + traceId: "t_collapse_2", + spanId: "s_collapse_2", + workerQueue: "main", + queue: "task/test-task", + isTest: false, + tags: [], + delayUntil: new Date(Date.now() + 500), + }, + prisma + ); + + expect(await snapshotStatuses(prisma, run.id)).toEqual(["DELAYED"]); + + await setTimeout(1_500); + + expect(await snapshotStatuses(prisma, run.id)).toEqual(["DELAYED", "QUEUED"]); + } finally { + await engine.quit(); + } + } + ); +}); diff --git a/internal-packages/run-store/src/PostgresRunStore.ts b/internal-packages/run-store/src/PostgresRunStore.ts index 1b0d2f89ce..676e345ad5 100644 --- a/internal-packages/run-store/src/PostgresRunStore.ts +++ b/internal-packages/run-store/src/PostgresRunStore.ts @@ -680,6 +680,7 @@ export class PostgresRunStore implements RunStore { const client = tx ?? this.prisma; const snapshotCreate = { + id: params.snapshot.id, engine: params.snapshot.engine, executionStatus: params.snapshot.executionStatus, description: params.snapshot.description, diff --git a/internal-packages/run-store/src/types.ts b/internal-packages/run-store/src/types.ts index 374a550e41..ac7cc35e3c 100644 --- a/internal-packages/run-store/src/types.ts +++ b/internal-packages/run-store/src/types.ts @@ -29,6 +29,7 @@ export type IdempotencyKeyRunMatch = { }; export type CreateRunSnapshotInput = { + id?: string; engine: "V2"; executionStatus: TaskRunExecutionStatus; description: string;