Decouple zod (#500)
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled

Zod Schemas is no longer required for validating/inferring event triggers. We’ve taken inspiration from how domain-functions did it: https://github.com/seasonedcc/domain-functions/pull/114
This commit is contained in:
Eric Allam
2023-09-22 14:44:03 -07:00
committed by GitHub
parent 3703f989fa
commit 50137a6f24
21 changed files with 131 additions and 58 deletions
+6
View File
@@ -0,0 +1,6 @@
---
"@trigger.dev/sdk": patch
"@trigger.dev/core": patch
---
Decouple zod
@@ -198,6 +198,7 @@ function classForJobStatus(status: JobRunStatus) {
case "WAITING_ON_CONNECTIONS":
case "PENDING":
case "UNRESOLVED_AUTH":
case "INVALID_PAYLOAD":
return "text-rose-500";
default:
return "";
@@ -2,7 +2,7 @@ import { CodeBlock } from "~/components/code/CodeBlock";
import { DateTime } from "~/components/primitives/DateTime";
import { Paragraph } from "~/components/primitives/Paragraph";
import { RunStatusIcon, RunStatusLabel } from "~/components/runs/RunStatuses";
import { MatchedRun, useRun } from "~/hooks/useRun";
import { MatchedRun } from "~/hooks/useRun";
import { formatDuration } from "~/utils";
import {
RunPanel,
@@ -17,7 +17,8 @@ export function hasFinished(status: JobRunStatus): boolean {
status === "ABORTED" ||
status === "TIMED_OUT" ||
status === "CANCELED" ||
status === "UNRESOLVED_AUTH"
status === "UNRESOLVED_AUTH" ||
status === "INVALID_PAYLOAD"
);
}
@@ -49,6 +50,7 @@ export function RunStatusIcon({ status, className }: { status: JobRunStatus; cla
case "TIMED_OUT":
return <ExclamationTriangleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "UNRESOLVED_AUTH":
case "INVALID_PAYLOAD":
return <XCircleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "WAITING_ON_CONNECTIONS":
return <WrenchIcon className={cn(runStatusClassNameColor(status), className)} />;
@@ -77,6 +79,7 @@ export function runBasicStatus(status: JobRunStatus): RunBasicStatus {
case "UNRESOLVED_AUTH":
case "CANCELED":
case "ABORTED":
case "INVALID_PAYLOAD":
return "FAILED";
case "SUCCESS":
return "COMPLETED";
@@ -111,6 +114,8 @@ export function runStatusTitle(status: JobRunStatus): string {
return "Canceled";
case "UNRESOLVED_AUTH":
return "Unresolved auth";
case "INVALID_PAYLOAD":
return "Invalid payload";
default: {
const _exhaustiveCheck: never = status;
throw new Error(`Non-exhaustive match for value: ${status}`);
@@ -130,6 +135,7 @@ export function runStatusClassNameColor(status: JobRunStatus): string {
return "text-amber-300";
case "FAILURE":
case "UNRESOLVED_AUTH":
case "INVALID_PAYLOAD":
return "text-rose-500";
case "TIMED_OUT":
return "text-amber-300";
@@ -1,6 +1,7 @@
import {
CachedTaskSchema,
RunJobError,
RunJobInvalidPayloadError,
RunJobResumeWithTask,
RunJobRetryWithTask,
RunJobSuccess,
@@ -348,6 +349,11 @@ export class PerformRunExecutionV1Service {
break;
}
case "INVALID_PAYLOAD": {
await this.#failRunWithInvalidPayloadError(execution, safeBody.data);
break;
}
default: {
const _exhaustiveCheck: never = status;
throw new Error(`Non-exhaustive match for value: ${status}`);
@@ -453,6 +459,15 @@ export class PerformRunExecutionV1Service {
});
}
async #failRunWithInvalidPayloadError(
execution: FoundRunExecution,
data: RunJobInvalidPayloadError
) {
return await $transaction(this.#prismaClient, async (tx) => {
await this.#failRunExecution(tx, execution, data.errors, "INVALID_PAYLOAD");
});
}
async #retryRunWithTask(execution: FoundRunExecution, data: RunJobRetryWithTask) {
const { run } = execution;
@@ -572,7 +587,7 @@ export class PerformRunExecutionV1Service {
prisma: PrismaClientOrTransaction,
execution: FoundRunExecution,
output: Record<string, any>,
status: "FAILURE" | "ABORTED" | "UNRESOLVED_AUTH" = "FAILURE"
status: "FAILURE" | "ABORTED" | "UNRESOLVED_AUTH" | "INVALID_PAYLOAD" = "FAILURE"
): Promise<void> {
const { run } = execution;
@@ -1,6 +1,7 @@
import {
CachedTask,
RunJobError,
RunJobInvalidPayloadError,
RunJobResumeWithTask,
RunJobRetryWithTask,
RunJobSuccess,
@@ -360,6 +361,11 @@ export class PerformRunExecutionV2Service {
break;
}
case "INVALID_PAYLOAD": {
await this.#failRunWithInvalidPayloadError(run, safeBody.data, durationInMs);
break;
}
default: {
const _exhaustiveCheck: never = status;
throw new Error(`Non-exhaustive match for value: ${status}`);
@@ -455,6 +461,23 @@ export class PerformRunExecutionV2Service {
});
}
async #failRunWithInvalidPayloadError(
execution: FoundRun,
data: RunJobInvalidPayloadError,
durationInMs: number
) {
return await $transaction(this.#prismaClient, async (tx) => {
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
execution,
data.errors,
"INVALID_PAYLOAD",
durationInMs
);
});
}
async #retryRunWithTask(
run: FoundRun,
data: RunJobRetryWithTask,
@@ -579,7 +602,7 @@ export class PerformRunExecutionV2Service {
reason: "EXECUTE_JOB" | "PREPROCESS",
run: FoundRun,
output: Record<string, any>,
status: "FAILURE" | "ABORTED" | "TIMED_OUT" | "UNRESOLVED_AUTH" = "FAILURE",
status: "FAILURE" | "ABORTED" | "TIMED_OUT" | "UNRESOLVED_AUTH" | "INVALID_PAYLOAD" = "FAILURE",
durationInMs: number = 0
): Promise<void> {
await $transaction(prisma, async (tx) => {
+9 -1
View File
@@ -2,7 +2,7 @@ import { ulid } from "ulid";
import { z } from "zod";
import { Prettify } from "../types";
import { addMissingVersionField } from "./addMissingVersionField";
import { ErrorWithStackSchema } from "./errors";
import { ErrorWithStackSchema, SchemaErrorSchema } from "./errors";
import { EventRuleSchema } from "./eventFilter";
import { ConnectionAuthSchema, IntegrationConfigSchema } from "./integrations";
import { DeserializedJsonSchema, SerializableJsonSchema } from "./json";
@@ -437,6 +437,13 @@ export const RunJobErrorSchema = z.object({
export type RunJobError = z.infer<typeof RunJobErrorSchema>;
export const RunJobInvalidPayloadErrorSchema = z.object({
status: z.literal("INVALID_PAYLOAD"),
errors: z.array(SchemaErrorSchema),
});
export type RunJobInvalidPayloadError = z.infer<typeof RunJobInvalidPayloadErrorSchema>;
export const RunJobUnresolvedAuthErrorSchema = z.object({
status: z.literal("UNRESOLVED_AUTH_ERROR"),
issues: z.record(z.object({ id: z.string(), error: z.string() })),
@@ -477,6 +484,7 @@ export type RunJobSuccess = z.infer<typeof RunJobSuccessSchema>;
export const RunJobResponseSchema = z.discriminatedUnion("status", [
RunJobErrorSchema,
RunJobUnresolvedAuthErrorSchema,
RunJobInvalidPayloadErrorSchema,
RunJobResumeWithTaskSchema,
RunJobRetryWithTaskSchema,
RunJobCanceledWithTaskSchema,
+7
View File
@@ -7,3 +7,10 @@ export const ErrorWithStackSchema = z.object({
});
export type ErrorWithStack = z.infer<typeof ErrorWithStackSchema>;
export const SchemaErrorSchema = z.object({
path: z.array(z.string()),
message: z.string(),
});
export type SchemaError = z.infer<typeof SchemaErrorSchema>;
+1
View File
@@ -14,6 +14,7 @@ export const RunStatusSchema = z.union([
z.literal("ABORTED"),
z.literal("CANCELED"),
z.literal("UNRESOLVED_AUTH"),
z.literal("INVALID_PAYLOAD"),
]);
export const RunTaskSchema = z.object({
@@ -0,0 +1,2 @@
-- AlterEnum
ALTER TYPE "JobRunStatus" ADD VALUE 'INVALID_PAYLOAD';
+1
View File
@@ -732,6 +732,7 @@ enum JobRunStatus {
ABORTED
CANCELED
UNRESOLVED_AUTH
INVALID_PAYLOAD
}
model JobRunExecution {
+1 -2
View File
@@ -39,8 +39,7 @@
"ulid": "^2.3.0",
"uuid": "^9.0.0",
"ws": "^8.11.0",
"zod": "3.21.4",
"zod-error": "1.5.0"
"zod": "3.21.4"
},
"devDependencies": {
"@trigger.dev/tsconfig": "workspace:*",
+5 -1
View File
@@ -1,4 +1,4 @@
import { ErrorWithStack, ServerTask } from "@trigger.dev/core";
import { ErrorWithStack, SchemaError, ServerTask } from "@trigger.dev/core";
export class ResumeWithTaskError {
constructor(public task: ServerTask) {}
@@ -16,6 +16,10 @@ export class CanceledWithTaskError {
constructor(public task: ServerTask) {}
}
export class ParsedPayloadSchemaError {
constructor(public schemaErrors: SchemaError[]) {}
}
/** Use this function if you're using a `try/catch` block to catch errors.
* It checks if a thrown error is a special internal error that you should ignore.
* If this returns `true` then you must rethrow the error: `throw err;`
+10 -1
View File
@@ -31,7 +31,12 @@ import {
StatusUpdate,
} from "@trigger.dev/core";
import { ApiClient } from "./apiClient";
import { CanceledWithTaskError, ResumeWithTaskError, RetryWithTaskError } from "./errors";
import {
CanceledWithTaskError,
ParsedPayloadSchemaError,
ResumeWithTaskError,
RetryWithTaskError,
} from "./errors";
import { TriggerIntegration } from "./integrations";
import { IO } from "./io";
import { createIOWithIntegrations } from "./ioWithIntegrations";
@@ -712,6 +717,10 @@ export class TriggerClient {
return { status: "SUCCESS", output };
} catch (error) {
if (error instanceof ParsedPayloadSchemaError) {
return { status: "INVALID_PAYLOAD", errors: error.schemaErrors };
}
if (error instanceof ResumeWithTaskError) {
return { status: "RESUME_WITH_TASK", task: error.task };
}
@@ -1,8 +1,9 @@
import { EventFilter, TriggerMetadata, deepMergeFilters } from "@trigger.dev/core";
import { z } from "zod";
import { Job } from "../job";
import { TriggerClient } from "../triggerClient";
import { EventSpecification, EventSpecificationExample, Trigger } from "../types";
import { EventSpecification, EventSpecificationExample, SchemaParser, Trigger } from "../types";
import { formatSchemaErrors } from "../utils/formatSchemaErrors";
import { ParsedPayloadSchemaError } from "../errors";
type EventTriggerOptions<TEventSpecification extends EventSpecification<any>> = {
event: TEventSpecification;
@@ -50,7 +51,7 @@ type TriggerOptions<TEvent> = {
/** A [Zod](https://trigger.dev/docs/documentation/guides/zod) schema that defines the shape of the event payload.
* The default is `z.any()` which is `any`.
* */
schema?: z.Schema<TEvent>;
schema?: SchemaParser<TEvent>;
/** You can use this to filter events based on the source. */
source?: string;
/** Used to filter which events trigger the Job
@@ -94,7 +95,13 @@ export function eventTrigger<TEvent extends any = any>(
examples: options.examples,
parsePayload: (rawPayload: any) => {
if (options.schema) {
return options.schema.parse(rawPayload);
const results = options.schema.safeParse(rawPayload);
if (!results.success) {
throw new ParsedPayloadSchemaError(formatSchemaErrors(results.error.issues));
}
return results.data;
}
return rawPayload as any;
@@ -1,5 +1,3 @@
import { z } from "zod";
import {
DisplayProperty,
EventFilter,
@@ -16,7 +14,7 @@ import { IOWithIntegrations, TriggerIntegration } from "../integrations";
import { IO } from "../io";
import { Job } from "../job";
import { TriggerClient } from "../triggerClient";
import type { EventSpecification, Trigger, TriggerContext } from "../types";
import type { EventSpecification, SchemaParser, Trigger, TriggerContext } from "../types";
import { slugifyId } from "../utils";
import { SerializableJson } from "@trigger.dev/core";
import { ConnectionAuth } from "@trigger.dev/core";
@@ -154,8 +152,8 @@ type ExternalSourceOptions<
> = {
id: string;
version: string;
schema: z.Schema<TParams>;
optionSchema?: z.Schema<TTriggerOptionDefinitions>;
schema: SchemaParser<TParams>;
optionSchema?: SchemaParser<TTriggerOptionDefinitions>;
integration: TIntegration;
register: RegisterFunction<TIntegration, TParams, TChannel, TTriggerOptionDefinitions>;
filter?: FilterFunction<TParams, TTriggerOptionDefinitions>;
+13
View File
@@ -105,3 +105,16 @@ export interface EventSpecification<TEvent extends any> {
export type EventTypeFromSpecification<TEventSpec extends EventSpecification<any>> =
TEventSpec extends EventSpecification<infer TEvent> ? TEvent : never;
export type SchemaParserIssue = { path: PropertyKey[]; message: string };
export type SchemaParserResult<T> =
| {
success: true;
data: T;
}
| { success: false; error: { issues: SchemaParserIssue[] } };
export type SchemaParser<T extends unknown = unknown> = {
safeParse: (a: unknown) => SchemaParserResult<T>;
};
@@ -0,0 +1,9 @@
import type { SchemaError } from "@trigger.dev/core";
import { SchemaParserIssue } from "../types";
export function formatSchemaErrors(errors: SchemaParserIssue[]): SchemaError[] {
return errors.map((error) => {
const { path, message } = error;
return { path: path.map(String), message };
});
}
-2
View File
@@ -990,7 +990,6 @@ importers:
uuid: ^9.0.0
ws: ^8.11.0
zod: 3.21.4
zod-error: 1.5.0
dependencies:
'@trigger.dev/core': link:../core
chalk: 5.2.0
@@ -1007,7 +1006,6 @@ importers:
uuid: 9.0.0
ws: 8.12.0
zod: 3.21.4
zod-error: 1.5.0
devDependencies:
'@trigger.dev/tsconfig': link:../../config-packages/tsconfig
'@types/debug': 4.1.7
-1
View File
@@ -43,7 +43,6 @@
"@types/node": "20.4.2",
"typescript": "5.1.6",
"zod": "3.21.4",
"@trigger.dev/airtable": "workspace:*",
"@trigger.dev/linear": "workspace:*"
},
"trigger.dev": {
+4 -37
View File
@@ -58,51 +58,18 @@ client.defineJob({
});
client.defineJob({
id: "example-job",
name: "Example Job: a joke with a delay",
id: "zod-schema",
name: "Job with Zod Schema",
version: "0.0.2",
trigger: eventTrigger({
name: "shayan.event",
name: "zod.schema",
schema: z.object({
userId: z.string(),
delay: z.number(),
}),
}),
run: async (payload, io, ctx) => {
await io.wait("sleeping", payload.delay);
await io.runTask(
"init",
async () => {
console.log("init function ran", payload.userId);
},
{ name: "init" }
);
await io.runTask(
"failable",
async (task) => {
if (task.attempts > 2) {
console.log("task succeeded");
return {
ok: true,
};
}
console.log("task failed");
throw new Error(`Task failed on ${task.attempts} attempt(s)`);
},
{ name: "task-1", retry: { limit: 3 } }
);
await io.runTask(
"log",
async () => {
console.log("hello from the job", payload.userId);
},
{
name: "log",
}
);
await io.logger.info("Hello World", { ctx, payload });
},
});