From c82ec22ee6db19b543bcd1cbb82a728e59f4810d Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Wed, 29 Jul 2026 14:33:53 +0100 Subject: [PATCH 1/4] perf(run-engine,run-store): one execution snapshot per triggered run A non-delayed run was getting RUN_CREATED nested in the run-create transaction followed by a separate QUEUED write. It now gets a single QUEUED snapshot inside the create, with the trigger path only publishing to the queue. --- .../collapse-trigger-queued-snapshot.md | 6 + .../run-engine/src/engine/index.ts | 33 +++- .../src/engine/systems/enqueueSystem.ts | 94 ++++++---- .../engine/tests/triggerCreateRouting.test.ts | 5 +- .../tests/triggerSnapshotCollapse.test.ts | 175 ++++++++++++++++++ .../run-store/src/PostgresRunStore.ts | 1 + internal-packages/run-store/src/types.ts | 1 + 7 files changed, 273 insertions(+), 42 deletions(-) create mode 100644 .server-changes/collapse-trigger-queued-snapshot.md create mode 100644 internal-packages/run-engine/src/engine/tests/triggerSnapshotCollapse.test.ts diff --git a/.server-changes/collapse-trigger-queued-snapshot.md b/.server-changes/collapse-trigger-queued-snapshot.md new file mode 100644 index 00000000000..ccc0169c6ca --- /dev/null +++ b/.server-changes/collapse-trigger-queued-snapshot.md @@ -0,0 +1,6 @@ +--- +area: webapp +type: improvement +--- + +Triggering a task now does one less database write, so runs reach 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 b5131766559..0980b959b75 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: taskRun.createdAt, + 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.ts b/internal-packages/run-engine/src/engine/systems/enqueueSystem.ts index 2b5aac6d376..75382ce82c2 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/tests/triggerCreateRouting.test.ts b/internal-packages/run-engine/src/engine/tests/triggerCreateRouting.test.ts index 34353cb44e0..af213dc7a71 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 00000000000..32372179ce6 --- /dev/null +++ b/internal-packages/run-engine/src/engine/tests/triggerSnapshotCollapse.test.ts @@ -0,0 +1,175 @@ +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( + "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 1b0d2f89cee..676e345ad5a 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 374a550e419..ac7cc35e3c8 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; From 8743fe5030c17fcc56e440889ef3bee5e3603455 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Wed, 29 Jul 2026 15:00:50 +0100 Subject: [PATCH 2/4] test(run-engine): update snapshot-count expectations for the collapsed trigger write The store-routing counters now also count createRun, since the trigger path's QUEUED snapshot is written nested in the run create. The replica-lag test starts the run attempt so it still has a middle snapshot to exercise rather than degenerating to an empty comparison. --- .../collapse-trigger-queued-snapshot.md | 2 +- .../src/engine/systems/enqueueSystem.test.ts | 19 +++++++++++++------ .../systems/executionSnapshotSystem.test.ts | 8 ++++++++ .../engine/tests/getSnapshotsSince.test.ts | 7 ++++++- 4 files changed, 28 insertions(+), 8 deletions(-) diff --git a/.server-changes/collapse-trigger-queued-snapshot.md b/.server-changes/collapse-trigger-queued-snapshot.md index ccc0169c6ca..541c5b64505 100644 --- a/.server-changes/collapse-trigger-queued-snapshot.md +++ b/.server-changes/collapse-trigger-queued-snapshot.md @@ -3,4 +3,4 @@ area: webapp type: improvement --- -Triggering a task now does one less database write, so runs reach the queue marginally faster. +Triggering a task now reaches the queue marginally faster. 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 a5730b306f2..eca893712cc 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/executionSnapshotSystem.test.ts b/internal-packages/run-engine/src/engine/systems/executionSnapshotSystem.test.ts index aeb361019e6..61d17a6b8b2 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 61ff3bff75f..54059943fdf 100644 --- a/internal-packages/run-engine/src/engine/tests/getSnapshotsSince.test.ts +++ b/internal-packages/run-engine/src/engine/tests/getSnapshotsSince.test.ts @@ -1149,11 +1149,16 @@ describe("RunEngine getSnapshotsSince", () => { ); await setTimeout(500); - await engine.dequeueFromWorkerQueue({ + const dequeued = await engine.dequeueFromWorkerQueue({ consumerId: "test_replica_stale_tail", workerQueue: "main", }); + await engine.startRunAttempt({ + runId: dequeued[0].run.id, + snapshotId: dequeued[0].snapshot.id, + }); + const allSnapshots = await prisma.taskRunExecutionSnapshot.findMany({ where: { runId: run.id, isValid: true }, orderBy: { createdAt: "asc" }, From bebc782ddc91ef96dbdf3eef76dec85cb9c02ff0 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Wed, 29 Jul 2026 15:15:33 +0100 Subject: [PATCH 3/4] test(run-engine): build the replica-lag test's middle snapshot on the primary startRunAttempt resolves the background worker task through readOnlyPrisma, and this test configures a deliberately empty schema-only replica, so using it to add a snapshot failed on the replica read instead of exercising the lag scenario. Write the snapshot with the file's own primary-side helper. --- .../src/engine/tests/getSnapshotsSince.test.ts | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) 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 54059943fdf..4cf74753e06 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 }); @@ -1149,14 +1149,18 @@ describe("RunEngine getSnapshotsSince", () => { ); await setTimeout(500); - const dequeued = await engine.dequeueFromWorkerQueue({ + await engine.dequeueFromWorkerQueue({ consumerId: "test_replica_stale_tail", workerQueue: "main", }); - await engine.startRunAttempt({ - runId: dequeued[0].run.id, - snapshotId: dequeued[0].snapshot.id, + 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({ From a9d0749ffa4196afd06070c4ba0ce016957467e9 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Wed, 29 Jul 2026 15:33:48 +0100 Subject: [PATCH 4/4] fix(run-engine): stamp the queued timeline event at write time The run's createdAt is caller-supplied and the scheduler sets it to the exact schedule time, so using it for the emitted event backdated the queued entry on every scheduled run's timeline. Use the write moment, as the other nested-create paths do. --- .../run-engine/src/engine/index.ts | 2 +- .../tests/triggerSnapshotCollapse.test.ts | 81 +++++++++++++++++++ 2 files changed, 82 insertions(+), 1 deletion(-) diff --git a/internal-packages/run-engine/src/engine/index.ts b/internal-packages/run-engine/src/engine/index.ts index 0980b959b75..61d60695eea 100644 --- a/internal-packages/run-engine/src/engine/index.ts +++ b/internal-packages/run-engine/src/engine/index.ts @@ -1165,7 +1165,7 @@ export class RunEngine { } this.eventBus.emit("executionSnapshotCreated", { - time: taskRun.createdAt, + time: new Date(), run: { id: taskRun.id, }, diff --git a/internal-packages/run-engine/src/engine/tests/triggerSnapshotCollapse.test.ts b/internal-packages/run-engine/src/engine/tests/triggerSnapshotCollapse.test.ts index 32372179ce6..cddb2756c30 100644 --- a/internal-packages/run-engine/src/engine/tests/triggerSnapshotCollapse.test.ts +++ b/internal-packages/run-engine/src/engine/tests/triggerSnapshotCollapse.test.ts @@ -100,6 +100,87 @@ describe("RunEngine trigger() execution snapshots", () => { } ); + 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 }) => {