Skip to content

Commit 05bf16a

Browse files
d-csTrigger.dev RepoOps
authored andcommitted
feat(run-store): implement execution-snapshot store read and write semantics
No user-facing change. The new decorator remains uninstalled and cannot be selected by production configuration in this PR. Mono-RevId: 8637f7ed56d1fc641d7c7038bae2ac4f4c36a4f2
1 parent 4f8587c commit 05bf16a

25 files changed

Lines changed: 6427 additions & 78 deletions

‎internal-packages/run-engine/src/engine/systems/completedWaitpointFreeze.test.ts‎

Lines changed: 45 additions & 47 deletions
Original file line numberDiff line numberDiff line change
@@ -6,13 +6,12 @@
66
import { isDeepStrictEqual } from "node:util";
77
import { describe, expect, it } from "vitest";
88
import type { Waitpoint } from "@trigger.dev/database";
9-
import { BatchId, RunId } from "@trigger.dev/core/v3/isomorphic";
10-
import type { CompletedWaitpoint } from "@trigger.dev/core/v3";
119
import type {
1210
CompletedWaitpointRecord,
1311
CompletedWaitpointResolver,
1412
CompletedWaitpointsPointer,
1513
ResolveCompletedWaitpointsArgs,
14+
SnapshotReadWaitpoint,
1615
} from "@internal/run-store";
1716

1817
// The frozen key sets, pinned exactly and bidirectionally. Renames, removals, widenings and
@@ -103,19 +102,15 @@ function recordOutputFor(w: Waitpoint): CompletedWaitpointRecord["output"] {
103102
return { inline: w.output };
104103
}
105104

106-
// The READ side of the freeze. Iterates `records`, never `order`.
105+
// The READ side of the freeze. Iterates `records`, never `order`, and returns ONE unenhanced row per
106+
// distinct record: no index, no nested completion objects. Index expansion and those objects belong
107+
// to the single enhancement step, which both backends then go through identically.
107108
async function referenceResolver(
108109
args: ResolveCompletedWaitpointsArgs,
109110
lookupRunOutput: (runId: string) => Promise<string | undefined>
110-
): Promise<CompletedWaitpoint[]> {
111-
const out: CompletedWaitpoint[] = [];
111+
): Promise<SnapshotReadWaitpoint[]> {
112+
const out: SnapshotReadWaitpoint[] = [];
112113
for (const record of args.records) {
113-
const indexes: (number | undefined)[] = [];
114-
for (let i = 0; i < args.order.length; i++) {
115-
if (args.order[i] === record.id) indexes.push(i);
116-
}
117-
if (indexes.length === 0) indexes.push(undefined);
118-
119114
let output: string | undefined;
120115
if (record.output === null) {
121116
output = undefined;
@@ -132,37 +127,23 @@ async function referenceResolver(
132127
throw new Error(`unknown record output variant: ${JSON.stringify(_never)}`);
133128
}
134129

135-
for (const index of indexes) {
136-
out.push({
137-
id: record.id,
138-
// Unreachable: the oracle's own loop pushes a non-negative integer or undefined.
139-
// Reproduced because the frozen index-expansion rule names it.
140-
index: index === -1 ? undefined : index,
141-
friendlyId: record.friendlyId,
142-
type: record.type,
143-
completedAt: new Date(record.completedAt),
144-
idempotencyKey: record.idempotencyKey,
145-
completedByTaskRun: record.completedByTaskRunId
146-
? {
147-
id: record.completedByTaskRunId,
148-
friendlyId: RunId.toFriendlyId(record.completedByTaskRunId),
149-
batch: args.batchId
150-
? { id: args.batchId, friendlyId: BatchId.toFriendlyId(args.batchId) }
151-
: undefined,
152-
}
153-
: undefined,
154-
completedAfter: record.completedAfter ? new Date(record.completedAfter) : undefined,
155-
completedByBatch: record.completedByBatchId
156-
? {
157-
id: record.completedByBatchId,
158-
friendlyId: BatchId.toFriendlyId(record.completedByBatchId),
159-
}
160-
: undefined,
161-
output,
162-
outputType: record.outputType,
163-
outputIsError: record.outputIsError,
164-
});
165-
}
130+
out.push({
131+
id: record.id,
132+
friendlyId: record.friendlyId,
133+
type: record.type,
134+
completedAt: new Date(record.completedAt),
135+
completedByTaskRunId: record.completedByTaskRunId ?? null,
136+
completedByBatchId: record.completedByBatchId ?? null,
137+
completedAfter: record.completedAfter ? new Date(record.completedAfter) : null,
138+
output: output ?? null,
139+
outputType: record.outputType,
140+
outputIsError: record.outputIsError,
141+
// The record already carries the resolved user-visible key, so it is re-expressed as the
142+
// triple the enhancement step reads: a key present means user-provided and active.
143+
idempotencyKey: record.idempotencyKey ?? "",
144+
userProvidedIdempotencyKey: record.idempotencyKey !== undefined,
145+
inactiveIdempotencyKey: null,
146+
});
166147
}
167148
return out;
168149
}
@@ -181,13 +162,18 @@ function makeSnapshot(batchId: string | null) {
181162
return { id: "snap_1", runId: "run_1", batchId, checkpoint: null } as never;
182163
}
183164

165+
// Parity is asserted on the FINAL runner-facing payload, not on the resolver's intermediate rows:
166+
// both sides go through the one enhancement step, the Postgres side over its own waitpoint rows and
167+
// the Redis side over the resolver's rows. Comparing the intermediate is what let a Redis read that
168+
// dropped every RUN/BATCH completion association still look correct here.
184169
async function assertParity(
185170
waitpoints: Waitpoint[],
186171
order: string[],
187172
batchId: string | null,
188173
runOutputs: Record<string, string> = {}
189174
) {
190-
const enhanced = enhanceExecutionSnapshotWithWaitpoints(makeSnapshot(batchId), waitpoints, order);
175+
const snapshot = makeSnapshot(batchId);
176+
const enhanced = enhanceExecutionSnapshotWithWaitpoints(snapshot, waitpoints, order);
191177
const args: ResolveCompletedWaitpointsArgs = {
192178
runId: "run_1",
193179
batchId: batchId ?? undefined,
@@ -197,7 +183,9 @@ async function assertParity(
197183
};
198184
// count-carried-forward behaviour (order.length, not the record count) is covered by
199185
// the run-store Redis suite, not here -- this line only constructs `args`, not asserts.
200-
const resolved = await referenceResolver(args, async (id) => runOutputs[id]);
186+
const rows = await referenceResolver(args, async (id) => runOutputs[id]);
187+
const resolvedEnhanced = enhanceExecutionSnapshotWithWaitpoints(snapshot, rows, order);
188+
const resolved = resolvedEnhanced.completedWaitpoints;
201189
expect(resolved).toEqual(enhanced.completedWaitpoints);
202190
return { enhanced, resolved };
203191
}
@@ -410,7 +398,11 @@ describe("the completed-waitpoints freeze", () => {
410398
});
411399
expect(resolved).toHaveLength(1);
412400
expect(resolved[0]!.id).toBe("wp_hook");
413-
expect(resolved[0]!.index).toBe(0);
401+
// The hook returns an UNENHANCED row even though the id sits at order position 0: no index and
402+
// no nested completion objects. Enhancement is the caller's single step.
403+
expect(resolved[0]!).not.toHaveProperty("index");
404+
expect(resolved[0]!).not.toHaveProperty("completedByTaskRun");
405+
expect(resolved[0]!).not.toHaveProperty("completedByBatch");
414406
});
415407

416408
it("round-trips all four waitpoint types", async () => {
@@ -601,8 +593,9 @@ describe("the exhaustive parity grid", () => {
601593
? [id]
602594
: [id, id];
603595

596+
const snapshot = makeSnapshot(readingBatchId);
604597
const enhanced = enhanceExecutionSnapshotWithWaitpoints(
605-
makeSnapshot(readingBatchId),
598+
snapshot,
606599
[w],
607600
order
608601
);
@@ -613,10 +606,15 @@ describe("the exhaustive parity grid", () => {
613606
order,
614607
records: [toRecord(w)],
615608
};
616-
const resolved = await referenceResolver(
609+
const rows = await referenceResolver(
617610
args,
618611
async (runId) => RUN_OUTPUT_LOOKUP[runId]
619612
);
613+
const resolved = enhanceExecutionSnapshotWithWaitpoints(
614+
snapshot,
615+
rows,
616+
order
617+
).completedWaitpoints;
620618

621619
if (!isDeepStrictEqual(resolved, enhanced.completedWaitpoints)) {
622620
failures.push({

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

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,11 @@ import type {
1010
TaskRunStatus,
1111
Waitpoint,
1212
} from "@trigger.dev/database";
13-
import type { RunStore } from "@internal/run-store";
13+
import type {
14+
LatestExecutionSnapshotRead,
15+
RunStore,
16+
SnapshotReadWaitpoint,
17+
} from "@internal/run-store";
1418
import { ExecutionSnapshotNotFoundError, ServiceValidationError } from "../errors.js";
1519
import type { HeartbeatTimeouts } from "../types.js";
1620
import type { SystemResources } from "./systems.js";
@@ -31,21 +35,14 @@ export interface EnhancedExecutionSnapshot extends TaskRunExecutionSnapshot {
3135
completedWaitpoints: CompletedWaitpoint[];
3236
}
3337

34-
type ExecutionSnapshotWithCheckAndWaitpoints = Prisma.TaskRunExecutionSnapshotGetPayload<{
35-
include: {
36-
checkpoint: true;
37-
completedWaitpoints: true;
38-
};
39-
}>;
40-
4138
type ExecutionSnapshotWithCheckpoint = Prisma.TaskRunExecutionSnapshotGetPayload<{
4239
include: {
4340
checkpoint: true;
4441
};
4542
}>;
4643

4744
function enhanceExecutionSnapshot(
48-
snapshot: ExecutionSnapshotWithCheckAndWaitpoints
45+
snapshot: LatestExecutionSnapshotRead
4946
): EnhancedExecutionSnapshot {
5047
return enhanceExecutionSnapshotWithWaitpoints(
5148
snapshot,
@@ -57,10 +54,15 @@ function enhanceExecutionSnapshot(
5754
/**
5855
* Transforms a snapshot (with checkpoint but without waitpoints) into an EnhancedExecutionSnapshot
5956
* by combining it with pre-fetched waitpoints.
57+
*
58+
* This is the ONE place that expands a distinct waitpoint across its repeated positions in
59+
* `completedWaitpointOrder`, assigns `index`, and builds the nested `completedByTaskRun` (with its
60+
* batch) and `completedByBatch` objects. Every backend feeds it the same unenhanced rows, so none of
61+
* them can produce a runner payload the others would not.
6062
*/
6163
export function enhanceExecutionSnapshotWithWaitpoints(
6264
snapshot: ExecutionSnapshotWithCheckpoint,
63-
waitpoints: Waitpoint[],
65+
waitpoints: SnapshotReadWaitpoint[],
6466
completedWaitpointOrder: string[]
6567
): EnhancedExecutionSnapshot {
6668
return {

‎internal-packages/run-store/src/PostgresRunStore.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ import type {
2727
FinalizeRunData,
2828
ForWaitpointCompletionContext,
2929
IdempotencyKeyRunMatch,
30+
LatestExecutionSnapshotRead,
3031
LockRunData,
3132
PromotePendingVersionArgs,
3233
ReadClient,
@@ -1930,9 +1931,7 @@ export class PostgresRunStore implements RunStore {
19301931
runId: string,
19311932
client?: ReadClient,
19321933
environmentId?: string
1933-
): Promise<Prisma.TaskRunExecutionSnapshotGetPayload<{
1934-
include: { completedWaitpoints: true; checkpoint: true };
1935-
}> | null> {
1934+
): Promise<LatestExecutionSnapshotRead | null> {
19361935
const prisma = client ?? this.readOnlyPrisma;
19371936
const where = { runId, isValid: true, ...(environmentId ? { environmentId } : {}) };
19381937

‎internal-packages/run-store/src/delegatingRunStore.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@ import type {
3939
FinalizeRunData,
4040
ForWaitpointCompletionContext,
4141
IdempotencyKeyRunMatch,
42+
LatestExecutionSnapshotRead,
4243
LockRunData,
4344
PromotePendingVersionArgs,
4445
ReadClient,
@@ -457,9 +458,7 @@ export class DelegatingRunStore implements RunStore {
457458
// When set, scopes the read to this environment (tenant boundary); a run in another env reads as
458459
// not-found. Omit to read regardless of environment (internal callers).
459460
environmentId?: string
460-
): Promise<Prisma.TaskRunExecutionSnapshotGetPayload<{
461-
include: { completedWaitpoints: true; checkpoint: true };
462-
}> | null> {
461+
): Promise<LatestExecutionSnapshotRead | null> {
463462
return this.delegate.findLatestExecutionSnapshot(runId, client, environmentId);
464463
}
465464

‎internal-packages/run-store/src/redisSnapshotStore.ts‎

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ import {
66
type Result,
77
} from "@internal/redis";
88
import { Logger } from "@trigger.dev/core/logger";
9-
import type { CompletedWaitpoint } from "@trigger.dev/core/v3/schemas";
9+
import type { SnapshotReadWaitpoint } from "./types.js";
1010
import {
1111
SNAPSHOT_NAMESPACE,
1212
SNAPSHOT_STATE_VERSION,
@@ -117,8 +117,8 @@ export type CompletedWaitpointRecordOutput =
117117
| null;
118118

119119
/**
120-
* One completed waitpoint, one per DISTINCT id in a wait cycle. The resolver expands
121-
* this into one CompletedWaitpoint per position of the id in the cycle's order list.
120+
* One completed waitpoint, one per DISTINCT id in a wait cycle. The resolver turns each into one
121+
* unenhanced read row; run-engine's enhancement step is what expands it across the cycle's order.
122122
*/
123123
export type CompletedWaitpointRecord = {
124124
id: string;
@@ -158,10 +158,14 @@ export type ResolveCompletedWaitpointsArgs = {
158158
/**
159159
* This lane owns the signature. The waitpoint lane owns the implementation, which
160160
* lives in run-engine because a deriveFromRun record needs a Postgres read.
161+
*
162+
* It returns UNENHANCED rows, one per distinct record: the read is a store read, so it produces the
163+
* same material a Postgres read produces and leaves index expansion and the nested completion objects
164+
* to run-engine's single enhancement step.
161165
*/
162166
export type CompletedWaitpointResolver = (
163167
args: ResolveCompletedWaitpointsArgs
164-
) => Promise<CompletedWaitpoint[]>;
168+
) => Promise<SnapshotReadWaitpoint[]>;
165169

166170
export type SnapshotEntryInput = {
167171
id: string;

‎internal-packages/run-store/src/runOpsStore.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import type {
2525
FinalizeRunData,
2626
ForWaitpointCompletionContext,
2727
IdempotencyKeyRunMatch,
28+
LatestExecutionSnapshotRead,
2829
LockRunData,
2930
PromotePendingVersionArgs,
3031
ReadClient,
@@ -1243,9 +1244,7 @@ export class RoutingRunStore implements RunStore {
12431244
runId: string,
12441245
client?: ReadClient,
12451246
environmentId?: string
1246-
): Promise<Prisma.TaskRunExecutionSnapshotGetPayload<{
1247-
include: { completedWaitpoints: true; checkpoint: true };
1248-
}> | null> {
1247+
): Promise<LatestExecutionSnapshotRead | null> {
12491248
const owningStore = this.#routeOrNew(runId);
12501249
const snapshot = await owningStore.findLatestExecutionSnapshot(
12511250
runId,

0 commit comments

Comments
 (0)