Context on job run

This commit is contained in:
nicktrn
2023-12-07 08:17:54 +00:00
parent d5ac3423cd
commit 9ea24f8815
15 changed files with 86 additions and 6 deletions
@@ -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 (
<RunPanel selected={false}>
@@ -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: {
@@ -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 (
<TriggerDetail
trigger={trigger}
@@ -39,6 +42,7 @@ export default function Page() {
event={job.event}
eventIds={run.eventIds}
payload={payload}
context={context}
properties={run.properties}
/>
);
@@ -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={[]}
/>
);
@@ -29,6 +29,7 @@ export default function Page() {
trigger={trigger}
event={{ icon: "webhook", title: "Register Webhook" }}
payload={trigger.payload}
context={trigger.context}
properties={[]}
/>
);
@@ -29,6 +29,7 @@ export default function Page() {
trigger={trigger}
event={{ icon: "mail-fast", title: "Deliver Webhook" }}
payload={trigger.payload}
context={trigger.context}
properties={[]}
/>
);
@@ -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,
@@ -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,
+11
View File
@@ -590,9 +590,20 @@ export const AutoYieldConfigSchema = z.object({
export type AutoYieldConfig = z.infer<typeof AutoYieldConfigSchema>;
export const EventMetadataSchema = ApiEventLogSchema.pick({
id: true,
name: true,
timestamp: true,
context: true,
});
export type EventMetadata = z.infer<typeof EventMetadataSchema>;
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(),
+4
View File
@@ -9,3 +9,7 @@ export interface AsyncMap {
has: (key: string) => Promise<boolean>;
set: (key: string, value: any) => Promise<any>;
}
export type SomeRequired<T extends Record<any, any>, TSome extends keyof T> = T & {
[K in TSome]-?: T[K];
};
@@ -0,0 +1,2 @@
-- AlterTable
ALTER TABLE "JobRun" ADD COLUMN "context" TEXT;
+1
View File
@@ -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)
+8 -2
View File
@@ -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<TTrigger, TriggerEventType<TTrigger>>,
io: IOWithIntegrations<TIntegrations>,
context: TriggerContext
context: GetTriggerContext<TTrigger>
) => Promise<TOutput>;
onSuccess?: (
@@ -112,6 +112,12 @@ type ArrayIfBatched<TTrigger extends Trigger<any>, T> = TTrigger extends EventTr
? ArrayIfTrue<HasBatchingEnabled<U>, T>
: T;
type GetTriggerContext<TTrigger extends Trigger<any>> = TTrigger extends EventTrigger<any, infer U>
? GenericTriggerContext<HasBatchingEnabled<U>>
: TTrigger extends WebhookTrigger<any, any, any, infer U>
? GenericTriggerContext<HasBatchingEnabled<U>>
: GenericTriggerContext<false>;
export type JobPayload<TJob> = TJob extends Job<Trigger<EventSpecification<infer TEvent>>, any>
? TEvent
: never;
+29 -2
View File
@@ -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<any> {
const { event, organization, project, environment, job, run, source } = execution;
let eventMetadata: SomeRequired<EventMetadata, "context">[] | 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,
+12 -1
View File
@@ -12,6 +12,7 @@ import type {
RunTaskOptions,
RuntimeEnvironmentType,
ServerTask,
SomeRequired,
SourceEventOption,
SuccessfulRunNotification,
TriggerMetadata,
@@ -32,6 +33,14 @@ export type {
SourceEventOption,
};
type BatchedContext = SomeRequired<Omit<TriggerContext, "event">, "events">;
type NonBatchedContext = SomeRequired<Omit<TriggerContext, "events">, "event">;
export type GenericTriggerContext<TBatched extends boolean> = 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 */