feat(webapp): extend SyntheticRun for replay (Phase B4)
The mollifier read-fallback's SyntheticRun previously carried just enough fields for the API retrieve/trace/spans/events/attempts/metadata endpoints. Phase C5 (replay) needs the buffered run to be passable where ReplayTaskRunService expects a TaskRun. Adds the missing fields: id, runtimeEnvironmentId, engine, workerQueue, queue, concurrencyKey, machinePreset, realtimeStreamsVersion, seedMetadata, seedMetadataType, runTags. All populated from the engine-trigger snapshot embedded in the buffer entry. Also closes a pre-existing typecheck gap in ApiRetrieveRunPresenter.synthesiseFoundRunFromBuffer — workerQueue wasn't populated and the file had been failing tsc. Now surfaces the buffered run's workerQueue, defaulting to "main" (the Prisma default).
This commit is contained in:
@@ -0,0 +1,6 @@
|
||||
---
|
||||
area: webapp
|
||||
type: improvement
|
||||
---
|
||||
|
||||
Extend `SyntheticRun` (the mollifier read-fallback synthesised TaskRun shape) with the fields `ReplayTaskRunService` reads: `id`, `runtimeEnvironmentId`, `engine`, `workerQueue`, `queue`, `concurrencyKey`, `machinePreset`, `realtimeStreamsVersion`, `seedMetadata`, `seedMetadataType`, and `runTags`. Populated from the buffered run's engine-trigger snapshot. Also closes a pre-existing typecheck gap in `ApiRetrieveRunPresenter.synthesiseFoundRunFromBuffer` by surfacing `workerQueue` (defaulting to `"main"`) on the synthesised FoundRun.
|
||||
@@ -573,6 +573,7 @@ function synthesiseFoundRunFromBuffer(buffered: SyntheticRun): FoundRun {
|
||||
attemptNumber: null,
|
||||
engine: "V2",
|
||||
taskEventStore: "taskEvent",
|
||||
workerQueue: buffered.workerQueue ?? "main",
|
||||
parentTaskRun: null,
|
||||
rootTaskRun: null,
|
||||
childRuns: [],
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import type { MollifierBuffer } from "@trigger.dev/redis-worker";
|
||||
import { RunId } from "@trigger.dev/core/v3/isomorphic";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { deserialiseMollifierSnapshot } from "./mollifierSnapshot.server";
|
||||
import { getMollifierBuffer } from "./mollifierBuffer.server";
|
||||
@@ -10,6 +11,10 @@ export type ReadFallbackInput = {
|
||||
};
|
||||
|
||||
export type SyntheticRun = {
|
||||
// Snapshot-derived TaskRun primary key. Used by ReplayTaskRunService
|
||||
// for logging and by callers passing this object where a TaskRun is
|
||||
// expected (cast). Derived deterministically from `friendlyId`.
|
||||
id: string;
|
||||
friendlyId: string;
|
||||
status: "QUEUED" | "FAILED";
|
||||
taskIdentifier: string | undefined;
|
||||
@@ -19,6 +24,12 @@ export type SyntheticRun = {
|
||||
payloadType: string | undefined;
|
||||
metadata: unknown;
|
||||
metadataType: string | undefined;
|
||||
// Seed-metadata mirrors what `triggerTask.server.ts` writes into the
|
||||
// snapshot: the original metadataPacket data preserved separately from
|
||||
// any later customer mutations. ReplayTaskRunService uses these to
|
||||
// rebuild the replay's metadata.
|
||||
seedMetadata: string | undefined;
|
||||
seedMetadataType: string | undefined;
|
||||
|
||||
idempotencyKey: string | undefined;
|
||||
idempotencyKeyOptions: string[] | undefined;
|
||||
@@ -26,6 +37,10 @@ export type SyntheticRun = {
|
||||
depth: number;
|
||||
ttl: string | undefined;
|
||||
tags: string[];
|
||||
// Mirror of `tags` under the PG field name. ReplayTaskRunService reads
|
||||
// `existingTaskRun.runTags`; both names are kept here so a synthetic
|
||||
// run can be passed wherever the PG-shape `runTags` is expected.
|
||||
runTags: string[];
|
||||
lockedToVersion: string | undefined;
|
||||
resumeParentOnCompletion: boolean;
|
||||
parentTaskRunId: string | undefined;
|
||||
@@ -36,6 +51,17 @@ export type SyntheticRun = {
|
||||
spanId: string | undefined;
|
||||
parentSpanId: string | undefined;
|
||||
|
||||
// Replay-relevant fields populated from the engine-trigger snapshot.
|
||||
// ReplayTaskRunService reads each of these from the existing TaskRun;
|
||||
// when the original lives in the buffer we synthesise them here.
|
||||
runtimeEnvironmentId: string | undefined;
|
||||
engine: "V2";
|
||||
workerQueue: string | undefined;
|
||||
queue: string | undefined;
|
||||
concurrencyKey: string | undefined;
|
||||
machinePreset: string | undefined;
|
||||
realtimeStreamsVersion: string | undefined;
|
||||
|
||||
error?: { code: string; message: string };
|
||||
};
|
||||
|
||||
@@ -77,7 +103,14 @@ export async function findRunByIdWithMollifierFallback(
|
||||
? asStringArray(idempotencyKeyOptionsRaw)
|
||||
: undefined;
|
||||
|
||||
const tags = asStringArray(snapshot.tags);
|
||||
const environment =
|
||||
snapshot.environment && typeof snapshot.environment === "object"
|
||||
? (snapshot.environment as Record<string, unknown>)
|
||||
: undefined;
|
||||
|
||||
return {
|
||||
id: RunId.fromFriendlyId(entry.runId),
|
||||
friendlyId: entry.runId,
|
||||
status: entry.status === "FAILED" ? "FAILED" : "QUEUED",
|
||||
taskIdentifier: asString(snapshot.taskIdentifier),
|
||||
@@ -87,13 +120,16 @@ export async function findRunByIdWithMollifierFallback(
|
||||
payloadType: asString(snapshot.payloadType),
|
||||
metadata: snapshot.metadata,
|
||||
metadataType: asString(snapshot.metadataType),
|
||||
seedMetadata: asString(snapshot.seedMetadata),
|
||||
seedMetadataType: asString(snapshot.seedMetadataType),
|
||||
|
||||
idempotencyKey: asString(snapshot.idempotencyKey),
|
||||
idempotencyKeyOptions,
|
||||
isTest: snapshot.isTest === true,
|
||||
depth: typeof snapshot.depth === "number" ? snapshot.depth : 0,
|
||||
ttl: asString(snapshot.ttl),
|
||||
tags: asStringArray(snapshot.tags),
|
||||
tags,
|
||||
runTags: tags,
|
||||
lockedToVersion: asString(snapshot.lockToVersion),
|
||||
resumeParentOnCompletion: snapshot.resumeParentOnCompletion === true,
|
||||
parentTaskRunId: asString(snapshot.parentTaskRunId),
|
||||
@@ -102,6 +138,15 @@ export async function findRunByIdWithMollifierFallback(
|
||||
spanId: asString(snapshot.spanId),
|
||||
parentSpanId: asString(snapshot.parentSpanId),
|
||||
|
||||
runtimeEnvironmentId:
|
||||
asString(environment?.id) ?? entry.envId,
|
||||
engine: "V2",
|
||||
workerQueue: asString(snapshot.workerQueue),
|
||||
queue: asString(snapshot.queue),
|
||||
concurrencyKey: asString(snapshot.concurrencyKey),
|
||||
machinePreset: asString(snapshot.machine),
|
||||
realtimeStreamsVersion: asString(snapshot.realtimeStreamsVersion),
|
||||
|
||||
error: entry.lastError,
|
||||
};
|
||||
} catch (err) {
|
||||
|
||||
@@ -216,4 +216,63 @@ describe("findRunByIdWithMollifierFallback", () => {
|
||||
expect(result!.traceId).toBeUndefined();
|
||||
expect(result!.spanId).toBeUndefined();
|
||||
});
|
||||
|
||||
it("populates replay-relevant fields from the snapshot", async () => {
|
||||
const entry: BufferEntry = {
|
||||
runId: "run_1",
|
||||
envId: "env_a",
|
||||
orgId: "org_1",
|
||||
payload: JSON.stringify({
|
||||
taskIdentifier: "my-task",
|
||||
environment: { id: "env_a" },
|
||||
workerQueue: "default",
|
||||
queue: "task/my-task",
|
||||
concurrencyKey: "tenant-42",
|
||||
machine: "medium-1x",
|
||||
realtimeStreamsVersion: "v2",
|
||||
seedMetadata: '{"k":"v"}',
|
||||
seedMetadataType: "application/json",
|
||||
tags: ["t1", "t2"],
|
||||
}),
|
||||
status: "QUEUED",
|
||||
attempts: 0,
|
||||
createdAt: NOW,
|
||||
};
|
||||
const result = await findRunByIdWithMollifierFallback(
|
||||
{ runId: "run_1", environmentId: "env_a", organizationId: "org_1" },
|
||||
{ getBuffer: () => fakeBuffer(entry) },
|
||||
);
|
||||
expect(result).not.toBeNull();
|
||||
expect(result!.id).toBeTypeOf("string");
|
||||
expect(result!.id.length).toBeGreaterThan(0);
|
||||
expect(result!.engine).toBe("V2");
|
||||
expect(result!.runtimeEnvironmentId).toBe("env_a");
|
||||
expect(result!.workerQueue).toBe("default");
|
||||
expect(result!.queue).toBe("task/my-task");
|
||||
expect(result!.concurrencyKey).toBe("tenant-42");
|
||||
expect(result!.machinePreset).toBe("medium-1x");
|
||||
expect(result!.realtimeStreamsVersion).toBe("v2");
|
||||
expect(result!.seedMetadata).toBe('{"k":"v"}');
|
||||
expect(result!.seedMetadataType).toBe("application/json");
|
||||
expect(result!.runTags).toEqual(["t1", "t2"]);
|
||||
});
|
||||
|
||||
it("falls back to entry.envId for runtimeEnvironmentId when snapshot lacks environment.id", async () => {
|
||||
const entry: BufferEntry = {
|
||||
runId: "run_1",
|
||||
envId: "env_a",
|
||||
orgId: "org_1",
|
||||
payload: JSON.stringify({ taskIdentifier: "t" }),
|
||||
status: "QUEUED",
|
||||
attempts: 0,
|
||||
createdAt: NOW,
|
||||
};
|
||||
const result = await findRunByIdWithMollifierFallback(
|
||||
{ runId: "run_1", environmentId: "env_a", organizationId: "org_1" },
|
||||
{ getBuffer: () => fakeBuffer(entry) },
|
||||
);
|
||||
expect(result!.runtimeEnvironmentId).toBe("env_a");
|
||||
expect(result!.workerQueue).toBeUndefined();
|
||||
expect(result!.queue).toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user