From 9ea24f8815e53f9893082401ef88e0a4e164ab18 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 7 Dec 2023 08:17:54 +0000 Subject: [PATCH] Context on job run --- .../app/components/run/TriggerDetail.tsx | 4 ++- .../app/presenters/RunPresenter.server.ts | 2 ++ .../route.tsx | 4 +++ .../route.tsx | 1 + .../route.tsx | 1 + .../route.tsx | 1 + .../app/services/runs/createRun.server.ts | 3 ++ .../runs/performRunExecutionV3.server.ts | 4 +++ packages/core/src/schemas/api.ts | 11 +++++++ packages/core/src/types.ts | 4 +++ .../migration.sql | 2 ++ packages/database/prisma/schema.prisma | 1 + packages/trigger-sdk/src/job.ts | 10 ++++-- packages/trigger-sdk/src/triggerClient.ts | 31 +++++++++++++++++-- packages/trigger-sdk/src/types.ts | 13 +++++++- 15 files changed, 86 insertions(+), 6 deletions(-) create mode 100644 packages/database/prisma/migrations/20231206153116_add_context_to_job_run/migration.sql diff --git a/apps/webapp/app/components/run/TriggerDetail.tsx b/apps/webapp/app/components/run/TriggerDetail.tsx index 8f5c7e3a6..5d87f35c8 100644 --- a/apps/webapp/app/components/run/TriggerDetail.tsx +++ b/apps/webapp/app/components/run/TriggerDetail.tsx @@ -16,6 +16,7 @@ import { DisplayProperty } from "@trigger.dev/core"; export function TriggerDetail({ trigger, payload, + context, event, properties, batched = false, @@ -23,6 +24,7 @@ export function TriggerDetail({ }: { trigger: DetailedEvent; payload: string; + context: string; event: { title: string; icon: string; @@ -31,7 +33,7 @@ export function TriggerDetail({ batched?: boolean; eventIds?: string[]; }) { - const { id, name, context, timestamp, deliveredAt } = trigger; + const { id, name, timestamp, deliveredAt } = trigger; return ( diff --git a/apps/webapp/app/presenters/RunPresenter.server.ts b/apps/webapp/app/presenters/RunPresenter.server.ts index 7856fb821..a932fab3b 100644 --- a/apps/webapp/app/presenters/RunPresenter.server.ts +++ b/apps/webapp/app/presenters/RunPresenter.server.ts @@ -84,6 +84,7 @@ export class RunPresenter { eventIds: run.eventIds, batched: run.batched, payload: run.payload, + context: run.context, tasks, runConnections: run.runConnections, missingConnections: run.missingConnections, @@ -123,6 +124,7 @@ export class RunPresenter { batched: true, eventIds: true, payload: true, + context: true, executionCount: true, executionDuration: true, version: { diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.runs.$runParam.trigger/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.runs.$runParam.trigger/route.tsx index e7caf9ef5..7364b4925 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.runs.$runParam.trigger/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.runs.$runParam.trigger/route.tsx @@ -32,6 +32,9 @@ export default function Page() { const payload = run.payload !== null ? JSON.stringify(JSON.parse(run.payload), null, 2) : trigger.payload; + const context = + run.context !== null ? JSON.stringify(JSON.parse(run.context), null, 2) : trigger.context; + return ( ); diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.external.$triggerParam_.runs.$runParam.trigger/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.external.$triggerParam_.runs.$runParam.trigger/route.tsx index 9715800be..c55882c59 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.external.$triggerParam_.runs.$runParam.trigger/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.external.$triggerParam_.runs.$runParam.trigger/route.tsx @@ -29,6 +29,7 @@ export default function Page() { trigger={trigger} event={{ icon: "register-source", title: "Register external source" }} payload={trigger.payload} + context={trigger.context} properties={[]} /> ); diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.webhooks.$triggerParam_.runs.$runParam.trigger/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.webhooks.$triggerParam_.runs.$runParam.trigger/route.tsx index decdd162b..e917082b4 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.webhooks.$triggerParam_.runs.$runParam.trigger/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.webhooks.$triggerParam_.runs.$runParam.trigger/route.tsx @@ -29,6 +29,7 @@ export default function Page() { trigger={trigger} event={{ icon: "webhook", title: "Register Webhook" }} payload={trigger.payload} + context={trigger.context} properties={[]} /> ); diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.webhooks.$triggerParam_.runs.delivery.$runParam.trigger/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.webhooks.$triggerParam_.runs.delivery.$runParam.trigger/route.tsx index 99a1d239f..dedfb2901 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.webhooks.$triggerParam_.runs.delivery.$runParam.trigger/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.triggers_.webhooks.$triggerParam_.runs.delivery.$runParam.trigger/route.tsx @@ -29,6 +29,7 @@ export default function Page() { trigger={trigger} event={{ icon: "mail-fast", title: "Deliver Webhook" }} payload={trigger.payload} + context={trigger.context} properties={[]} /> ); diff --git a/apps/webapp/app/services/runs/createRun.server.ts b/apps/webapp/app/services/runs/createRun.server.ts index a83870ef7..6d70c99fa 100644 --- a/apps/webapp/app/services/runs/createRun.server.ts +++ b/apps/webapp/app/services/runs/createRun.server.ts @@ -65,6 +65,9 @@ export class CreateRunService { payload: JSON.stringify( batched ? eventRecords.map((event) => event.payload) ?? [{}] : firstEvent.payload ?? {} ), + context: JSON.stringify( + batched ? eventRecords.map((event) => event.context) ?? [{}] : firstEvent.context ?? {} + ), externalAccountId: firstEvent.externalAccountId ? firstEvent.externalAccountId : undefined, diff --git a/apps/webapp/app/services/runs/performRunExecutionV3.server.ts b/apps/webapp/app/services/runs/performRunExecutionV3.server.ts index c52cebae0..f45e7dd34 100644 --- a/apps/webapp/app/services/runs/performRunExecutionV3.server.ts +++ b/apps/webapp/app/services/runs/performRunExecutionV3.server.ts @@ -451,8 +451,10 @@ export class PerformRunExecutionV3Service { return { event, + eventIds: run.eventIds, batched: run.batched, payload: run.payload, + context: run.context, job: { id: run.version.job.slug, version: run.version.version, @@ -504,8 +506,10 @@ export class PerformRunExecutionV3Service { return { event, + eventIds: run.eventIds, batched: run.batched, payload: run.payload, + context: run.context, job: { id: run.version.job.slug, version: run.version.version, diff --git a/packages/core/src/schemas/api.ts b/packages/core/src/schemas/api.ts index 28428140c..7314d614a 100644 --- a/packages/core/src/schemas/api.ts +++ b/packages/core/src/schemas/api.ts @@ -590,9 +590,20 @@ export const AutoYieldConfigSchema = z.object({ export type AutoYieldConfig = z.infer; +export const EventMetadataSchema = ApiEventLogSchema.pick({ + id: true, + name: true, + timestamp: true, + context: true, +}); + +export type EventMetadata = z.infer; + export const RunJobBodySchema = z.object({ event: ApiEventLogSchema, + eventIds: z.string().array(), payload: z.string().nullable(), + context: z.string().nullable(), batched: z.boolean(), job: z.object({ id: z.string(), diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index ae413ce93..407f9d9ec 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -9,3 +9,7 @@ export interface AsyncMap { has: (key: string) => Promise; set: (key: string, value: any) => Promise; } + +export type SomeRequired, TSome extends keyof T> = T & { + [K in TSome]-?: T[K]; +}; diff --git a/packages/database/prisma/migrations/20231206153116_add_context_to_job_run/migration.sql b/packages/database/prisma/migrations/20231206153116_add_context_to_job_run/migration.sql new file mode 100644 index 000000000..7eabef22a --- /dev/null +++ b/packages/database/prisma/migrations/20231206153116_add_context_to_job_run/migration.sql @@ -0,0 +1,2 @@ +-- AlterTable +ALTER TABLE "JobRun" ADD COLUMN "context" TEXT; diff --git a/packages/database/prisma/schema.prisma b/packages/database/prisma/schema.prisma index c0b593963..9acc56c9f 100644 --- a/packages/database/prisma/schema.prisma +++ b/packages/database/prisma/schema.prisma @@ -764,6 +764,7 @@ model JobRun { internal Boolean @default(false) payload String? + context String? batched Boolean @default(false) job Job @relation(fields: [jobId], references: [id], onDelete: Cascade, onUpdate: Cascade) diff --git a/packages/trigger-sdk/src/job.ts b/packages/trigger-sdk/src/job.ts index 32aa47b20..e00db5c61 100644 --- a/packages/trigger-sdk/src/job.ts +++ b/packages/trigger-sdk/src/job.ts @@ -15,8 +15,8 @@ import { runLocalStorage } from "./runLocalStorage"; import { TriggerClient } from "./triggerClient"; import type { EventSpecification, + GenericTriggerContext, Trigger, - TriggerContext, TriggerEventType, TriggerInvokeType, } from "./types"; @@ -84,7 +84,7 @@ export type JobOptions< run: ( payload: ArrayIfBatched>, io: IOWithIntegrations, - context: TriggerContext + context: GetTriggerContext ) => Promise; onSuccess?: ( @@ -112,6 +112,12 @@ type ArrayIfBatched, T> = TTrigger extends EventTr ? ArrayIfTrue, T> : T; +type GetTriggerContext> = TTrigger extends EventTrigger + ? GenericTriggerContext> + : TTrigger extends WebhookTrigger + ? GenericTriggerContext> + : GenericTriggerContext; + export type JobPayload = TJob extends Job>, any> ? TEvent : never; diff --git a/packages/trigger-sdk/src/triggerClient.ts b/packages/trigger-sdk/src/triggerClient.ts index 8aca65f72..83405d344 100644 --- a/packages/trigger-sdk/src/triggerClient.ts +++ b/packages/trigger-sdk/src/triggerClient.ts @@ -6,6 +6,8 @@ import { DeserializedJson, EphemeralEventDispatcherRequestBody, ErrorWithStackSchema, + EventMetadata, + EventMetadataSchema, FailedRunNotification, GetRunOptionsWithTaskDetails, GetRunsOptions, @@ -40,6 +42,7 @@ import { ScheduleMetadata, SendEvent, SendEventOptions, + SomeRequired, SourceMetadataV2, StatusUpdate, SuccessfulRunNotification, @@ -74,6 +77,7 @@ import { ExternalSource } from "./triggers/externalSource"; import { DynamicIntervalOptions, DynamicSchedule } from "./triggers/scheduled"; import type { EventSpecification, + GenericTriggerContext, NotificationsEventEmitter, Trigger, TriggerContext, @@ -125,6 +129,7 @@ import { ConcurrencyLimit, ConcurrencyLimitOptions } from "./concurrencyLimit"; import { KeyValueStore } from "./store/keyValueStore"; import { WebhookDeliveryContext, WebhookSource } from "./triggers/webhook"; import { formatSchemaErrors } from "./utils/formatSchemaErrors"; +import { z } from "zod"; export type TriggerClientOptions = { /** The `id` property is used to uniquely identify the client. @@ -1235,7 +1240,7 @@ export class TriggerClient { } const output = await runLocalStorage.runWith({ io, ctx: context }, () => { - return job.options.run(parsedPayload, ioWithConnections, context); + return job.options.run(parsedPayload, ioWithConnections, context as any); }); if (this.#options.verbose) { @@ -1364,9 +1369,30 @@ export class TriggerClient { }; } - #createRunContext(execution: RunJobBody): TriggerContext { + #createRunContext(execution: RunJobBody): GenericTriggerContext { const { event, organization, project, environment, job, run, source } = execution; + let eventMetadata: SomeRequired[] | undefined; + + if (execution.batched) { + if (!execution.context) { + throw new Error("Context needs to be defined for batched events."); + } + + const parsedContext = z.array(z.any()).parse(JSON.parse(execution.context)); + + if (parsedContext.length !== execution.eventIds.length) { + throw new Error("Context length needs to match number of event IDs."); + } + + eventMetadata = execution.eventIds.map((id, i) => ({ + id, + name: execution.event.name, + context: parsedContext[i], + timestamp: execution.event.timestamp, + })); + } + return { event: { id: event.id, @@ -1374,6 +1400,7 @@ export class TriggerClient { context: event.context, timestamp: event.timestamp, }, + events: eventMetadata, organization, project: project ?? { id: "unknown", name: "unknown", slug: "unknown" }, // backwards compat with old servers environment, diff --git a/packages/trigger-sdk/src/types.ts b/packages/trigger-sdk/src/types.ts index 2b38f5a97..6b96e1b57 100644 --- a/packages/trigger-sdk/src/types.ts +++ b/packages/trigger-sdk/src/types.ts @@ -12,6 +12,7 @@ import type { RunTaskOptions, RuntimeEnvironmentType, ServerTask, + SomeRequired, SourceEventOption, SuccessfulRunNotification, TriggerMetadata, @@ -32,6 +33,14 @@ export type { SourceEventOption, }; +type BatchedContext = SomeRequired, "events">; + +type NonBatchedContext = SomeRequired, "event">; + +export type GenericTriggerContext = TBatched extends true + ? BatchedContext + : NonBatchedContext; + export interface TriggerContext { /** Job metadata */ job: { id: string; version: string }; @@ -44,7 +53,9 @@ export interface TriggerContext { /** Run metadata */ run: { id: string; isTest: boolean; startedAt: Date; isRetry: boolean }; /** Event metadata */ - event: { id: string; name: string; context: any; timestamp: Date }; + event?: { id: string; name: string; context: any; timestamp: Date }; + /** Batched event metadata */ + events?: { id: string; name: string; context: any; timestamp: Date }[]; /** Source metadata */ source?: { id: string; metadata?: any }; /** Account metadata */