Skip to content

Commit dba857b

Browse files
1stvampTrigger.dev RepoOps
authored andcommitted
fix(run-engine): resolve the snapshot route when a stalled dequeue is requeued
Fix a stalled run occasionally not being requeued correctly when its snapshots are stored in Redis. Mono-RevId: 10d1c8f9d75d4cbd42475a46b68a1ec40831c7d4
1 parent d44c7e5 commit dba857b

3 files changed

Lines changed: 166 additions & 7 deletions

File tree

‎internal-packages/run-engine/src/engine/index.ts‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2912,6 +2912,13 @@ export class RunEngine {
29122912
throw new Error(`Run ${runId} not found`);
29132913
}
29142914

2915+
// The heartbeat payload carries no route, so resolve it durably for the QUEUED snapshot.
2916+
const snapshotRoute = await this.runAttemptSystem.effectiveRoute(
2917+
runId,
2918+
latestSnapshot.organizationId,
2919+
undefined
2920+
);
2921+
29152922
//it will automatically be requeued X times depending on the queue retry settings
29162923
const { wasRequeued } = await this.runAttemptSystem.tryNackAndRequeue({
29172924
run,
@@ -2929,6 +2936,7 @@ export class RunEngine {
29292936
code: "TASK_RUN_DEQUEUED_MAX_RETRIES",
29302937
message: `Trying to create an attempt failed multiple times, exceeding how many times we retry.`,
29312938
},
2939+
snapshotRoute,
29322940
tx: prisma,
29332941
});
29342942

‎internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -697,7 +697,7 @@ export class RunAttemptSystem {
697697
// residency ONCE (forceDurable) so the transition still lands in the run's true store instead of
698698
// taking the never-enrolled Postgres shortcut. Fails closed (throws) when residency cannot be
699699
// confirmed. Callers place this AFTER their no-op early exits so a no-op does no durable read.
700-
async #effectiveRoute(
700+
public async effectiveRoute(
701701
runId: string,
702702
organizationId: string,
703703
route: SnapshotRouteWire | undefined
@@ -802,8 +802,8 @@ export class RunAttemptSystem {
802802
span.setAttribute("completionStatus", completion.ok);
803803
span.setAttribute("runId", runId);
804804

805-
// Resolve the route the terminal snapshot honors (carried, else durable). See #effectiveRoute.
806-
const effectiveRoute = await this.#effectiveRoute(
805+
// Resolve the route the terminal snapshot honors (carried, else durable). See effectiveRoute.
806+
const effectiveRoute = await this.effectiveRoute(
807807
runId,
808808
latestSnapshot.organizationId,
809809
snapshotRoute
@@ -1013,8 +1013,8 @@ export class RunAttemptSystem {
10131013
span.setAttribute("completionStatus", completion.ok);
10141014

10151015
// The route every transition this failure path writes (retry, requeue, fail, cancel) honors.
1016-
// Carried, else resolved durably; see #effectiveRoute. Resolved after the no-op exit above.
1017-
const effectiveRoute = await this.#effectiveRoute(
1016+
// Carried, else resolved durably; see effectiveRoute. Resolved after the no-op exit above.
1017+
const effectiveRoute = await this.effectiveRoute(
10181018
runId,
10191019
latestSnapshot.organizationId,
10201020
snapshotRoute
@@ -1533,7 +1533,7 @@ export class RunAttemptSystem {
15331533

15341534
switch (outcome) {
15351535
case "requeue": {
1536-
const effectiveRoute = await this.#effectiveRoute(
1536+
const effectiveRoute = await this.effectiveRoute(
15371537
runId,
15381538
latestSnapshot.organizationId,
15391539
snapshotRoute
@@ -1688,7 +1688,7 @@ export class RunAttemptSystem {
16881688
// store instead of the never-enrolled Postgres shortcut. Placed AFTER the no-transition early
16891689
// exits (already FINISHED, PENDING_CANCEL without finalize) so a no-op cancel does no durable
16901690
// read; fails closed (throws) when durable residency cannot be confirmed.
1691-
const effectiveRoute = await this.#effectiveRoute(
1691+
const effectiveRoute = await this.effectiveRoute(
16921692
runId,
16931693
latestSnapshot.organizationId,
16941694
snapshotRoute
Lines changed: 151 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,151 @@
1+
// A stalled PENDING_EXECUTING snapshot must be requeued into a redis-primary run's durable residency. The
2+
// heartbeat payload carries no route, so on an undefined-dial engine the requeue's QUEUED snapshot would
3+
// otherwise take the Postgres shortcut and leave the MemoryDB head at PENDING_EXECUTING. Real infra.
4+
import { assertNonNullable, containerTest } from "@internal/testcontainers";
5+
import {
6+
PostgresRunStore,
7+
RedisSnapshotStore,
8+
SnapshotResidencyResolver,
9+
TaskRunExecutionSnapshotStore,
10+
type SnapshotStoreDial,
11+
} from "@internal/run-store";
12+
import { trace } from "@internal/tracing";
13+
import { setTimeout } from "node:timers/promises";
14+
import { expect } from "vitest";
15+
import { RunEngine } from "../index.js";
16+
import { createCompletedWaitpointResolver } from "../systems/completedWaitpointResolver.js";
17+
import { setupAuthenticatedEnvironment, setupBackgroundWorker } from "./setup.js";
18+
19+
vi.setConfig({ testTimeout: 90_000 });
20+
21+
const ROUTE = "logical:1";
22+
const PENDING_EXECUTING_TIMEOUT_MS = 500;
23+
const machines = {
24+
defaultMachine: "small-1x",
25+
machines: {
26+
"small-1x": { name: "small-1x" as const, cpu: 0.5, memory: 0.5, centsPerMs: 0.0001 },
27+
},
28+
baseCostInCents: 0.0001,
29+
};
30+
31+
function makeEngine(
32+
prisma: any,
33+
redisOptions: any,
34+
snapshotStore: RedisSnapshotStore,
35+
dial: () => SnapshotStoreDial | undefined,
36+
workerEnabled: boolean
37+
) {
38+
const delegate = new PostgresRunStore({ prisma, readOnlyPrisma: prisma });
39+
const store = new TaskRunExecutionSnapshotStore(delegate, {
40+
store: snapshotStore,
41+
mode: "redis-only",
42+
resolveDial: dial,
43+
residencyResolver: new SnapshotResidencyResolver({
44+
store: snapshotStore,
45+
taskRunExists: async (id: string) => (await prisma.taskRun.count({ where: { id } })) > 0,
46+
}),
47+
resolveCompletedWaitpoints: createCompletedWaitpointResolver(delegate),
48+
logicalRunStoreRoute: ROUTE,
49+
});
50+
return new RunEngine({
51+
prisma,
52+
store,
53+
worker: workerEnabled
54+
? { redis: redisOptions, workers: 1, tasksPerWorker: 10, pollIntervalMs: 100 }
55+
: { redis: redisOptions, disabled: true },
56+
queue: {
57+
redis: redisOptions,
58+
retryOptions: { maxTimeoutInMs: 50 },
59+
masterQueueConsumersDisabled: true,
60+
processWorkerQueueDebounceMs: 50,
61+
},
62+
runLock: { redis: redisOptions },
63+
machines,
64+
heartbeatTimeoutsMs: { PENDING_EXECUTING: PENDING_EXECUTING_TIMEOUT_MS },
65+
tracer: trace.getTracer("test", "0.0.0"),
66+
});
67+
}
68+
69+
describe("handleStalledSnapshot honors durable residency for a stalled PENDING_EXECUTING run", () => {
70+
containerTest(
71+
"a dial=undefined engine requeues a redis-primary run into its MemoryDB head",
72+
async ({ prisma, redisOptions }) => {
73+
const env = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
74+
const snapshotStore = new RedisSnapshotStore({ redisOptions, completedTtlMs: 60_000 });
75+
const producer = makeEngine(prisma, redisOptions, snapshotStore, () => "redis-only", false);
76+
const consumer = makeEngine(prisma, redisOptions, snapshotStore, () => undefined, true);
77+
78+
try {
79+
await setupBackgroundWorker(producer, env, "test-task");
80+
81+
async function dequeueOnConsumer() {
82+
for (let i = 0; i < 25; i++) {
83+
await producer.runQueue.processMasterQueueForEnvironment(env.id, 5);
84+
const dequeued = await consumer.dequeueFromWorkerQueue({
85+
consumerId: "stall_consumer",
86+
workerQueue: "main",
87+
});
88+
if (dequeued.length > 0) return dequeued[0];
89+
await setTimeout(300);
90+
}
91+
throw new Error("run never reached the worker queue");
92+
}
93+
94+
const run = await producer.trigger(
95+
{
96+
number: 1,
97+
friendlyId: "run_stallroute",
98+
environment: env,
99+
taskIdentifier: "test-task",
100+
payload: "{}",
101+
payloadType: "application/json",
102+
context: {},
103+
traceContext: {},
104+
traceId: "t-stall",
105+
spanId: "s-stall",
106+
workerQueue: "main",
107+
queue: "task/test-task",
108+
isTest: false,
109+
tags: [],
110+
delayUntil: new Date(Date.now() + 60_000),
111+
},
112+
prisma
113+
);
114+
const runId = run.id;
115+
116+
await prisma.taskRun.update({
117+
where: { id: runId },
118+
data: { delayUntil: null, queueTimestamp: new Date() },
119+
});
120+
const runRow = await prisma.taskRun.findFirstOrThrow({ where: { id: runId } });
121+
await producer.enqueueSystem.enqueueRun({ run: runRow, env, enableFastPath: true });
122+
123+
const first = await dequeueOnConsumer();
124+
assertNonNullable(first);
125+
expect(first.snapshotRoute?.residency).toBe("redis-primary");
126+
expect(first.snapshot.executionStatus).toBe("PENDING_EXECUTING");
127+
128+
let latest = await consumer.getRunExecutionData({ runId });
129+
for (let i = 0; i < 40 && latest?.snapshot.executionStatus !== "QUEUED"; i++) {
130+
await setTimeout(250);
131+
latest = await consumer.getRunExecutionData({ runId });
132+
}
133+
assertNonNullable(latest);
134+
expect(latest.snapshot.executionStatus).toBe("QUEUED");
135+
const head = await snapshotStore.getLatest(runId);
136+
assertNonNullable(head);
137+
expect(head.id).toBe(latest.snapshot.id);
138+
expect(await prisma.taskRunExecutionSnapshot.count({ where: { runId } })).toBe(0);
139+
140+
const redelivered = await dequeueOnConsumer();
141+
assertNonNullable(redelivered);
142+
expect(redelivered.run.id).toBe(runId);
143+
expect(redelivered.snapshot.executionStatus).toBe("PENDING_EXECUTING");
144+
} finally {
145+
await producer.quit();
146+
await consumer.quit();
147+
await snapshotStore.quit();
148+
}
149+
}
150+
);
151+
});

0 commit comments

Comments
 (0)