Skip to content

Commit daaaa50

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
perf: zod compile reusable server validation schemas
Mono-RevId: 749cbfed4729e03d99a710f4ce62acdc0beba8ca
1 parent 0e9d021 commit daaaa50

23 files changed

Lines changed: 306 additions & 35 deletions
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
import { WorkerApiRunLatestSnapshotResponseBody } from "@trigger.dev/core/v3/workers";
2+
import { describe, expect, it } from "vitest";
3+
import { z } from "zod";
4+
import { resolveResponseSchema } from "./responseSchemas.js";
5+
6+
describe("supervisor response schemas", () => {
7+
it("retains date normalization and validation errors for snapshots", () => {
8+
const schema = resolveResponseSchema(WorkerApiRunLatestSnapshotResponseBody);
9+
const result = schema.parse({
10+
execution: {
11+
version: "1",
12+
snapshot: {
13+
id: "snapshot_1",
14+
friendlyId: "snap_1",
15+
executionStatus: "EXECUTING",
16+
description: "Run is executing",
17+
createdAt: "2026-01-01T00:00:00.000Z",
18+
},
19+
run: { id: "run_1", friendlyId: "run_1", status: "EXECUTING" },
20+
completedWaitpoints: [],
21+
},
22+
});
23+
24+
expect(result.execution.snapshot.createdAt).toEqual(new Date("2026-01-01T00:00:00.000Z"));
25+
expect(schema.safeParse({ execution: {} }).success).toBe(false);
26+
});
27+
28+
it("does not compile unfamiliar schemas or repeat their callbacks", () => {
29+
let calls = 0;
30+
const schema = z.string().refine(() => ++calls > 1);
31+
32+
expect(resolveResponseSchema(schema).safeParse("value").success).toBe(false);
33+
expect(calls).toBe(1);
34+
});
35+
});
Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
import type { AnyZodSchema } from "@trigger.dev/core/v3";
2+
import {
3+
WorkerApiConnectResponseBody,
4+
WorkerApiContinueRunExecutionRequestBody,
5+
WorkerApiDequeueResponseBody,
6+
WorkerApiHeartbeatResponseBody,
7+
WorkerApiRunAttemptCompleteResponseBody,
8+
WorkerApiRunHeartbeatResponseBody,
9+
WorkerApiRunLatestSnapshotResponseBody,
10+
WorkerApiRunSnapshotsSinceResponseBody,
11+
WorkerApiSuspendRunResponseBody,
12+
} from "@trigger.dev/core/v3/workers";
13+
import { z } from "zod";
14+
15+
const responseSchemas = [
16+
WorkerApiConnectResponseBody,
17+
WorkerApiContinueRunExecutionRequestBody,
18+
WorkerApiDequeueResponseBody,
19+
WorkerApiHeartbeatResponseBody,
20+
WorkerApiRunAttemptCompleteResponseBody,
21+
WorkerApiRunHeartbeatResponseBody,
22+
WorkerApiRunLatestSnapshotResponseBody,
23+
WorkerApiRunSnapshotsSinceResponseBody,
24+
WorkerApiSuspendRunResponseBody,
25+
];
26+
27+
const compiledSchemas = new Map<AnyZodSchema, AnyZodSchema>(
28+
responseSchemas.map((schema) => [schema, z.compile(schema)])
29+
);
30+
31+
export function resolveResponseSchema<T extends AnyZodSchema>(schema: T): T {
32+
return (compiledSchemas.get(schema) as T | undefined) ?? schema;
33+
}

apps/supervisor/src/index.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ import {
2323
isKubernetesEnvironment,
2424
} from "@trigger.dev/core/v3/serverOnly";
2525
import { createK8sApi, createPodCountFetcher } from "./clients/kubernetes.js";
26+
import { resolveResponseSchema } from "./clients/responseSchemas.js";
2627
import { collectDefaultMetrics, Counter, Gauge, Histogram } from "prom-client";
2728
import { register } from "./metrics.js";
2829
import { PodCleaner } from "./services/podCleaner.js";
@@ -324,6 +325,7 @@ class ManagedSupervisor {
324325
const workerToken = getWorkerToken();
325326

326327
this.workerSession = new SupervisorSession({
328+
resolveResponseSchema,
327329
workerToken,
328330
apiUrl: env.TRIGGER_API_URL,
329331
instanceName: env.TRIGGER_WORKER_INSTANCE_NAME,

apps/supervisor/src/workloadServer/index.ts

Lines changed: 14 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -66,10 +66,16 @@ const checkpointCancelRequests = new Counter({
6666
registers: [register],
6767
});
6868

69-
const WorkloadActionParams = z.object({
70-
runFriendlyId: z.string(),
71-
snapshotFriendlyId: z.string(),
72-
});
69+
const WorkloadActionParams = z.compile(
70+
z.object({
71+
runFriendlyId: z.string(),
72+
snapshotFriendlyId: z.string(),
73+
})
74+
);
75+
const WorkloadDebugParams = z.compile(WorkloadActionParams.pick({ runFriendlyId: true }));
76+
const CompiledWorkloadRunAttemptStartRequestBody = z.compile(WorkloadRunAttemptStartRequestBody);
77+
const CompiledWorkloadHeartbeatRequestBody = z.compile(WorkloadHeartbeatRequestBody);
78+
const CompiledWorkloadDebugLogRequestBody = z.compile(WorkloadDebugLogRequestBody);
7379

7480
// Workloads bundled into customer task images before CLI v4.4.4 use a strict
7581
// zod enum for checkpoint type that only allows DOCKER and KUBERNETES. The
@@ -385,7 +391,7 @@ export class WorkloadServer extends EventEmitter<WorkloadServerEvents> {
385391
"POST",
386392
{
387393
paramsSchema: WorkloadActionParams,
388-
bodySchema: WorkloadRunAttemptStartRequestBody,
394+
bodySchema: CompiledWorkloadRunAttemptStartRequestBody,
389395
handler: async (ctx) =>
390396
this.wideRoute(
391397
ctx,
@@ -490,7 +496,7 @@ export class WorkloadServer extends EventEmitter<WorkloadServerEvents> {
490496
"POST",
491497
{
492498
paramsSchema: WorkloadActionParams,
493-
bodySchema: WorkloadHeartbeatRequestBody,
499+
bodySchema: CompiledWorkloadHeartbeatRequestBody,
494500
handler: async (ctx) =>
495501
this.wideRoute(
496502
ctx,
@@ -724,8 +730,8 @@ export class WorkloadServer extends EventEmitter<WorkloadServerEvents> {
724730

725731
if (env.SEND_RUN_DEBUG_LOGS) {
726732
httpServer.route("/api/v1/workload-actions/runs/:runFriendlyId/logs/debug", "POST", {
727-
paramsSchema: WorkloadActionParams.pick({ runFriendlyId: true }),
728-
bodySchema: WorkloadDebugLogRequestBody,
733+
paramsSchema: WorkloadDebugParams,
734+
bodySchema: CompiledWorkloadDebugLogRequestBody,
729735
handler: async (ctx) =>
730736
this.wideRoute(
731737
ctx,

apps/webapp/app/routes/api.v1.tasks.$taskId.trigger.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ const { action, loader } = createActionApiRoute(
5353
{
5454
headers: HeadersSchema,
5555
params: ParamsSchema,
56-
body: TriggerTaskRequestBody,
56+
body: z.compile(TriggerTaskRequestBody),
5757
allowJWT: true,
5858
maxContentLength: env.TASK_PAYLOAD_MAXIMUM_SIZE,
5959
authorization: {

apps/webapp/app/routes/api.v2.tasks.batch.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { json } from "@remix-run/server-runtime";
2+
import { z } from "zod";
23
import type { BatchTriggerTaskV3Response } from "@trigger.dev/core/v3";
34
import { BatchTriggerTaskV3RequestBody, generateJWT } from "@trigger.dev/core/v3";
45
import { env } from "~/env.server";
@@ -28,7 +29,7 @@ const { action, loader } = createActionApiRoute(
2829
headers: HeadersSchema.extend({
2930
"batch-processing-strategy": BatchProcessingStrategy.nullish(),
3031
}),
31-
body: BatchTriggerTaskV3RequestBody,
32+
body: z.compile(BatchTriggerTaskV3RequestBody),
3233
allowJWT: true,
3334
maxContentLength: env.BATCH_TASK_PAYLOAD_MAXIMUM_SIZE,
3435
authorization: {

apps/webapp/app/routes/engine.v1.worker-actions.connect.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,11 +2,12 @@ import type { TypedResponse } from "@remix-run/server-runtime";
22
import { json } from "@remix-run/server-runtime";
33
import type { WorkerApiConnectResponseBody } from "@trigger.dev/core/v3/workers";
44
import { WorkerApiConnectRequestBody } from "@trigger.dev/core/v3/workers";
5+
import { z } from "zod";
56
import { createActionWorkerApiRoute } from "~/services/routeBuilders/apiBuilder.server";
67

78
export const action = createActionWorkerApiRoute(
89
{
9-
body: WorkerApiConnectRequestBody,
10+
body: z.compile(WorkerApiConnectRequestBody),
1011
},
1112
async ({ authenticatedWorker, body }): Promise<TypedResponse<WorkerApiConnectResponseBody>> => {
1213
await authenticatedWorker.connect(body.metadata);

apps/webapp/app/routes/engine.v1.worker-actions.dequeue.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,11 +2,12 @@ import type { TypedResponse } from "@remix-run/server-runtime";
22
import { json } from "@remix-run/server-runtime";
33
import type { WorkerApiDequeueResponseBody } from "@trigger.dev/core/v3/workers";
44
import { WorkerApiDequeueRequestBody } from "@trigger.dev/core/v3/workers";
5+
import { z } from "zod";
56
import { createActionWorkerApiRoute } from "~/services/routeBuilders/apiBuilder.server";
67

78
export const action = createActionWorkerApiRoute(
89
{
9-
body: WorkerApiDequeueRequestBody,
10+
body: z.compile(WorkerApiDequeueRequestBody),
1011
},
1112
async ({
1213
authenticatedWorker,

apps/webapp/app/routes/engine.v1.worker-actions.heartbeat.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,11 +2,12 @@ import type { TypedResponse } from "@remix-run/server-runtime";
22
import { json } from "@remix-run/server-runtime";
33
import type { WorkerApiHeartbeatResponseBody } from "@trigger.dev/core/v3/workers";
44
import { WorkerApiHeartbeatRequestBody } from "@trigger.dev/core/v3/workers";
5+
import { z } from "zod";
56
import { createActionWorkerApiRoute } from "~/services/routeBuilders/apiBuilder.server";
67

78
export const action = createActionWorkerApiRoute(
89
{
9-
body: WorkerApiHeartbeatRequestBody,
10+
body: z.compile(WorkerApiHeartbeatRequestBody),
1011
},
1112
async ({ authenticatedWorker }): Promise<TypedResponse<WorkerApiHeartbeatResponseBody>> => {
1213
await authenticatedWorker.heartbeatWorkerInstance();

apps/webapp/app/routes/engine.v1.worker-actions.runs.$runFriendlyId.logs.debug.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ export const action = createActionWorkerApiRoute(
1010
params: z.object({
1111
runFriendlyId: z.string(),
1212
}),
13-
body: WorkerApiDebugLogBody,
13+
body: z.compile(WorkerApiDebugLogBody),
1414
},
1515
async ({ body, params }): Promise<Response> => {
1616
const { runFriendlyId } = params;

0 commit comments

Comments
 (0)