From d3c593cdbf1cb422d093123b8e3c94e152fa556c Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Fri, 20 Jan 2023 08:39:58 -0800 Subject: [PATCH] =?UTF-8?q?Added=20ability=20to=20timeout=20triggers=20so?= =?UTF-8?q?=20they=20don=E2=80=99t=20run=20too=20late=20after=20the=20even?= =?UTF-8?q?t=20time?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .changeset/good-pillows-drum.md | 5 ++ apps/webapp/app/components/runs/runStatus.tsx | 13 ++++++ apps/webapp/app/models/workflowRun.server.ts | 1 + .../app/services/messageBroker.server.ts | 12 +++++ .../services/runs/runTriggerTimeout.server.ts | 37 +++++++++++++++ .../scheduler/deliverScheduledEvent.server.ts | 1 + .../workflows/registerWorkflow.server.ts | 2 + .../migration.sql | 2 + .../migration.sql | 2 + .../migration.sql | 3 ++ apps/webapp/prisma/schema.prisma | 6 +++ apps/wss/src/server.ts | 30 +++++++++++- examples/playground/src/index.ts | 1 + examples/smoke-test/src/index.ts | 1 + .../internal-bridge/src/schemas/server.ts | 1 + .../src/messages/catalogs/triggers.ts | 5 +- .../src/messages/schemas/workflowRuns.ts | 8 ++++ .../src/messages/zodPublisher.ts | 2 + .../src/messages/zodSubscriber.ts | 46 +++++++++++++++++-- .../src/schemas/workflows.ts | 1 + packages/trigger-sdk/src/client.ts | 1 + packages/trigger-sdk/src/trigger/index.ts | 7 +++ 22 files changed, 181 insertions(+), 6 deletions(-) create mode 100644 .changeset/good-pillows-drum.md create mode 100644 apps/webapp/app/services/runs/runTriggerTimeout.server.ts create mode 100644 apps/webapp/prisma/migrations/20230120010921_add_trigger_ttl_in_seconds/migration.sql create mode 100644 apps/webapp/prisma/migrations/20230120014606_add_timed_out_status_for_workflow_runs/migration.sql create mode 100644 apps/webapp/prisma/migrations/20230120014641_add_timed_out_columns/migration.sql diff --git a/.changeset/good-pillows-drum.md b/.changeset/good-pillows-drum.md new file mode 100644 index 000000000..c31981c06 --- /dev/null +++ b/.changeset/good-pillows-drum.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/sdk": patch +--- + +Added triggerTTL option that prevents old events from running a workflow diff --git a/apps/webapp/app/components/runs/runStatus.tsx b/apps/webapp/app/components/runs/runStatus.tsx index a87280da8..7f2c66de4 100644 --- a/apps/webapp/app/components/runs/runStatus.tsx +++ b/apps/webapp/app/components/runs/runStatus.tsx @@ -21,6 +21,8 @@ export function runStatusTitle(status: WorkflowRunStatus): string { return "Disconnected"; case "ERROR": return "Error"; + case "TIMED_OUT": + return "Timed out"; } } @@ -36,6 +38,8 @@ export function runStatusLabel(status: WorkflowRunStatus): ReactNode { return {runStatusTitle(status)}; case "ERROR": return {runStatusTitle(status)}; + case "TIMED_OUT": + return {runStatusTitle(status)}; } } @@ -91,5 +95,14 @@ export function runStatusIcon( )} /> ); + case "TIMED_OUT": + return ( + + ); } } diff --git a/apps/webapp/app/models/workflowRun.server.ts b/apps/webapp/app/models/workflowRun.server.ts index 663e00bac..aaa619493 100644 --- a/apps/webapp/app/models/workflowRun.server.ts +++ b/apps/webapp/app/models/workflowRun.server.ts @@ -19,6 +19,7 @@ export async function findWorklowRunById(id: string) { include: { event: true, environment: true, + workflow: true, }, }); } diff --git a/apps/webapp/app/services/messageBroker.server.ts b/apps/webapp/app/services/messageBroker.server.ts index 08f17a71e..da0b67b6b 100644 --- a/apps/webapp/app/services/messageBroker.server.ts +++ b/apps/webapp/app/services/messageBroker.server.ts @@ -43,6 +43,7 @@ import { WaitForConnection } from "./requests/waitForConnection.server"; import { WorkflowRunDisconnected } from "./runs/runDisconnected.server"; import { DeliverScheduledEvent } from "./scheduler/deliverScheduledEvent.server"; import { RegisterSchedulerSource } from "./scheduler/registerSchedulerSource.server"; +import { WorkflowRunTriggerTimeout } from "./runs/runTriggerTimeout.server"; let pulsarClient: PulsarClient; let triggerPublisher: ZodPublisher; @@ -206,6 +207,13 @@ function createCommandSubscriber() { return !!success; }, + WORKFLOW_RUN_TRIGGER_TIMEOUT: async (id, data, properties) => { + const service = new WorkflowRunTriggerTimeout(); + + await service.call(data); + + return true; + }, SEND_INTEGRATION_REQUEST: async (id, data, properties) => { const service = new CreateIntegrationRequest(); @@ -526,6 +534,10 @@ function createTaskQueue() { "x-workflow-id": run.workflowId, "x-env": run.environment.slug, "x-workflow-run-id": run.id, + "x-ttl": run.workflow.triggerTtlInSeconds, + }, + { + eventTimestamp: run.event.timestamp.getTime(), } ); diff --git a/apps/webapp/app/services/runs/runTriggerTimeout.server.ts b/apps/webapp/app/services/runs/runTriggerTimeout.server.ts new file mode 100644 index 000000000..db7464e3d --- /dev/null +++ b/apps/webapp/app/services/runs/runTriggerTimeout.server.ts @@ -0,0 +1,37 @@ +import type { PrismaClient } from "~/db.server"; +import { prisma } from "~/db.server"; + +export class WorkflowRunTriggerTimeout { + #prismaClient: PrismaClient; + + constructor(prismaClient: PrismaClient = prisma) { + this.#prismaClient = prismaClient; + } + + async call(data: { id: string; ttl: number; elapsedSeconds: number }) { + const workflowRun = await this.#prismaClient.workflowRun.findUnique({ + where: { id: data.id }, + include: { + event: true, + environment: true, + }, + }); + + if (!workflowRun) { + throw new Error("Workflow run not found"); + } + + if (workflowRun.status !== "PENDING") { + return; + } + + await this.#prismaClient.workflowRun.update({ + where: { id: data.id }, + data: { + status: "TIMED_OUT", + timedOutAt: new Date(), + timedOutReason: `Trigger timed out after ${data.elapsedSeconds}s because it exceeded the TTL of ${data.ttl}s`, + }, + }); + } +} diff --git a/apps/webapp/app/services/scheduler/deliverScheduledEvent.server.ts b/apps/webapp/app/services/scheduler/deliverScheduledEvent.server.ts index 52f420440..190400195 100644 --- a/apps/webapp/app/services/scheduler/deliverScheduledEvent.server.ts +++ b/apps/webapp/app/services/scheduler/deliverScheduledEvent.server.ts @@ -67,6 +67,7 @@ export class DeliverScheduledEvent { context: {}, organizationId: schedulerSource.organizationId, environmentId: schedulerSource.environmentId, + timestamp: payload.scheduledTime, }, }); diff --git a/apps/webapp/app/services/workflows/registerWorkflow.server.ts b/apps/webapp/app/services/workflows/registerWorkflow.server.ts index d642e434f..4232ebb65 100644 --- a/apps/webapp/app/services/workflows/registerWorkflow.server.ts +++ b/apps/webapp/app/services/workflows/registerWorkflow.server.ts @@ -110,6 +110,7 @@ export class RegisterWorkflow { type: payload.trigger.type, service: payload.trigger.service, eventNames: payload.trigger.name, + triggerTtlInSeconds: payload.triggerTTL, }, create: { organizationId: organization.id, @@ -120,6 +121,7 @@ export class RegisterWorkflow { status: payload.trigger.service === "trigger" ? "READY" : "CREATED", service: payload.trigger.service, eventNames: payload.trigger.name, + triggerTtlInSeconds: payload.triggerTTL, }, include: { externalSource: true, diff --git a/apps/webapp/prisma/migrations/20230120010921_add_trigger_ttl_in_seconds/migration.sql b/apps/webapp/prisma/migrations/20230120010921_add_trigger_ttl_in_seconds/migration.sql new file mode 100644 index 000000000..56814bce0 --- /dev/null +++ b/apps/webapp/prisma/migrations/20230120010921_add_trigger_ttl_in_seconds/migration.sql @@ -0,0 +1,2 @@ +-- AlterTable +ALTER TABLE "Workflow" ADD COLUMN "triggerTtlInSeconds" INTEGER NOT NULL DEFAULT 3600; diff --git a/apps/webapp/prisma/migrations/20230120014606_add_timed_out_status_for_workflow_runs/migration.sql b/apps/webapp/prisma/migrations/20230120014606_add_timed_out_status_for_workflow_runs/migration.sql new file mode 100644 index 000000000..e76af3ece --- /dev/null +++ b/apps/webapp/prisma/migrations/20230120014606_add_timed_out_status_for_workflow_runs/migration.sql @@ -0,0 +1,2 @@ +-- AlterEnum +ALTER TYPE "WorkflowRunStatus" ADD VALUE 'TIMED_OUT'; diff --git a/apps/webapp/prisma/migrations/20230120014641_add_timed_out_columns/migration.sql b/apps/webapp/prisma/migrations/20230120014641_add_timed_out_columns/migration.sql new file mode 100644 index 000000000..90dafe50a --- /dev/null +++ b/apps/webapp/prisma/migrations/20230120014641_add_timed_out_columns/migration.sql @@ -0,0 +1,3 @@ +-- AlterTable +ALTER TABLE "WorkflowRun" ADD COLUMN "timedOutAt" TIMESTAMP(3), +ADD COLUMN "timedOutReason" TEXT; diff --git a/apps/webapp/prisma/schema.prisma b/apps/webapp/prisma/schema.prisma index 9c9a96a4b..a2bf756e2 100644 --- a/apps/webapp/prisma/schema.prisma +++ b/apps/webapp/prisma/schema.prisma @@ -142,6 +142,8 @@ model Workflow { archivedAt DateTime? isArchived Boolean @default(false) + triggerTtlInSeconds Int @default(3600) + @@unique([organizationId, slug]) } @@ -406,6 +408,9 @@ model WorkflowRun { startedAt DateTime? finishedAt DateTime? + timedOutAt DateTime? + timedOutReason String? + isTest Boolean @default(false) requests IntegrationRequest[] delays DurableDelay[] @@ -417,6 +422,7 @@ enum WorkflowRunStatus { DISCONNECTED SUCCESS ERROR + TIMED_OUT } model WorkflowRunStep { diff --git a/apps/wss/src/server.ts b/apps/wss/src/server.ts index 72a48c664..30bff325f 100644 --- a/apps/wss/src/server.ts +++ b/apps/wss/src/server.ts @@ -324,6 +324,7 @@ export class TriggerServer { name: data.packageName, version: data.packageVersion, }, + triggerTTL: data.triggerTTL, }); this.#workflowId = response.id; @@ -340,7 +341,7 @@ export class TriggerServer { subscriptionInitialPosition: "Earliest", }, handlers: { - TRIGGER_WORKFLOW: async (id, data, properties) => { + TRIGGER_WORKFLOW: async (id, data, properties, messageAttributes) => { this.#logger.debug("Received trigger", id, data, properties); // If the API keys don't match, then we should ignore it // This ensures the workflow is triggered for the correct environment @@ -366,6 +367,33 @@ export class TriggerServer { ); } + if (properties["x-ttl"] && messageAttributes.eventTimestamp) { + const ttl = properties["x-ttl"]; + const eventTimestamp = messageAttributes.eventTimestamp; + const now = Date.now(); + const elapsedMilliseconds = now - eventTimestamp.getTime(); + const elapsedSeconds = elapsedMilliseconds / 1000; + + if (elapsedSeconds > ttl) { + this.#logger.debug("Message is expired, ignoring", { + messageAttributes, + properties, + }); + + await this.#commandPublisher.publish( + "WORKFLOW_RUN_TRIGGER_TIMEOUT", + { + id: data.id, + ttl, + elapsedSeconds, + }, + { ...properties, "x-timestamp": String(Date.now()) } + ); + + return true; + } + } + this.#logger.debug("Triggering workflow", data, properties); const runController = new WorkflowRunController({ diff --git a/examples/playground/src/index.ts b/examples/playground/src/index.ts index b4ff4d52e..de93347aa 100644 --- a/examples/playground/src/index.ts +++ b/examples/playground/src/index.ts @@ -427,6 +427,7 @@ new Trigger({ apiKey: "trigger_dev_zC25mKNn6c0q", endpoint: "ws://localhost:8889/ws", logLevel: "debug", + triggerTTL: 5, on: scheduleEvent({ rateOf: { minutes: 4 } }), run: async (event, ctx) => { await ctx.logger.info("Received the scheduled event", { diff --git a/examples/smoke-test/src/index.ts b/examples/smoke-test/src/index.ts index 16a628954..98707bff2 100644 --- a/examples/smoke-test/src/index.ts +++ b/examples/smoke-test/src/index.ts @@ -12,6 +12,7 @@ const trigger = new Trigger({ apiKey: "trigger_dev_zC25mKNn6c0q", endpoint: "ws://localhost:8889/ws", logLevel: "log", + triggerTTL: 60 * 60 * 24, on: customEvent({ name: "user.created", schema: userCreatedEvent }), run: async (event, ctx) => { await ctx.logger.info("Inside the smoke test workflow, received event", { diff --git a/packages/internal-bridge/src/schemas/server.ts b/packages/internal-bridge/src/schemas/server.ts index 7adc906ce..78e95cc2c 100644 --- a/packages/internal-bridge/src/schemas/server.ts +++ b/packages/internal-bridge/src/schemas/server.ts @@ -58,6 +58,7 @@ export const ServerRPCSchema = { trigger: TriggerMetadataSchema, packageVersion: z.string(), packageName: z.string(), + triggerTTL: z.number().optional(), }), response: z .discriminatedUnion("type", [ diff --git a/packages/internal-platform/src/messages/catalogs/triggers.ts b/packages/internal-platform/src/messages/catalogs/triggers.ts index 18f3a2802..057308d93 100644 --- a/packages/internal-platform/src/messages/catalogs/triggers.ts +++ b/packages/internal-platform/src/messages/catalogs/triggers.ts @@ -1,10 +1,13 @@ +import { z } from "zod"; import { TriggerWorkflowMessageSchema } from "../schemas/workflows"; import { WorkflowRunEventPropertiesSchema } from "../sharedSchemas"; const Catalog = { TRIGGER_WORKFLOW: { data: TriggerWorkflowMessageSchema, - properties: WorkflowRunEventPropertiesSchema, + properties: WorkflowRunEventPropertiesSchema.extend({ + "x-ttl": z.coerce.number().optional(), + }), }, }; diff --git a/packages/internal-platform/src/messages/schemas/workflowRuns.ts b/packages/internal-platform/src/messages/schemas/workflowRuns.ts index a619c73f2..e3c83f9b2 100644 --- a/packages/internal-platform/src/messages/schemas/workflowRuns.ts +++ b/packages/internal-platform/src/messages/schemas/workflowRuns.ts @@ -27,4 +27,12 @@ export const commands = { }), properties: WorkflowSendRunEventPropertiesSchema, }, + WORKFLOW_RUN_TRIGGER_TIMEOUT: { + data: z.object({ + id: z.string(), + ttl: z.number(), + elapsedSeconds: z.number(), + }), + properties: WorkflowSendRunEventPropertiesSchema, + }, }; diff --git a/packages/internal-platform/src/messages/zodPublisher.ts b/packages/internal-platform/src/messages/zodPublisher.ts index 3ba1497ba..a60c4a3d8 100644 --- a/packages/internal-platform/src/messages/zodPublisher.ts +++ b/packages/internal-platform/src/messages/zodPublisher.ts @@ -16,6 +16,7 @@ export type PublishOptions = { partitionKey?: string; orderingKey?: string; id?: string; + eventTimestamp?: number; }; type PendingMessages = Array<{ @@ -216,6 +217,7 @@ export class ZodPublisher { deliverAt: options?.deliverAt, partitionKey: options?.partitionKey, orderingKey: options?.orderingKey, + eventTimestamp: options?.eventTimestamp, }); return response.toString(); diff --git a/packages/internal-platform/src/messages/zodSubscriber.ts b/packages/internal-platform/src/messages/zodSubscriber.ts index 055d1cb19..9688aba7f 100644 --- a/packages/internal-platform/src/messages/zodSubscriber.ts +++ b/packages/internal-platform/src/messages/zodSubscriber.ts @@ -14,13 +14,21 @@ import { import { z, ZodError } from "zod"; import { ZodPubSubStatus } from "./types"; +export type SubscriberMessageAttributes = { + messageId: string; + eventTimestamp?: Date; + publishedTimestamp: Date; + redeliveryCount: number; +}; + export type ZodSubscriberHandlers< TConsumerSchema extends MessageCatalogSchema > = { [K in keyof TConsumerSchema]: ( id: string, data: z.infer, - properties: z.infer + properties: z.infer, + attributes: SubscriberMessageAttributes ) => Promise; }; @@ -109,8 +117,32 @@ export class ZodSubscriber { const properties = this.#getRawProperties(msg); + const messageId = msg.getMessageId(); + const publishedTimestamp = msg.getPublishTimestamp(); + const eventTimestamp = msg.getEventTimestamp(); + const redeliveryCount = msg.getRedeliveryCount(); + + this.#logger.debug("#onMessage", { + messageId, + publishedTimestamp, + eventTimestamp, + redeliveryCount, + }); + + const messageAttributes = { + eventTimestamp: + eventTimestamp === 0 ? undefined : new Date(eventTimestamp), + messageId: messageId.toString(), + publishedTimestamp: new Date(publishedTimestamp), + redeliveryCount, + }; + try { - const wasHandled = await this.#handleMessage(messageData, properties); + const wasHandled = await this.#handleMessage( + messageData, + properties, + messageAttributes + ); if (wasHandled) { await consumer.acknowledge(msg); @@ -134,7 +166,8 @@ export class ZodSubscriber { async #handleMessage( rawMessage: MessageData, - rawProperties: Record = {} + rawProperties: Record = {}, + messageAttributes: SubscriberMessageAttributes ): Promise { const subscriberSchema = this.#schema; type TypeKeys = keyof typeof subscriberSchema; @@ -160,7 +193,12 @@ export class ZodSubscriber { const handler = this.#handlers[typeName]; - const returnValue = await handler(rawMessage.id, message, properties); + const returnValue = await handler( + rawMessage.id, + message, + properties, + messageAttributes + ); return returnValue; } diff --git a/packages/internal-platform/src/schemas/workflows.ts b/packages/internal-platform/src/schemas/workflows.ts index bec510812..b9024b021 100644 --- a/packages/internal-platform/src/schemas/workflows.ts +++ b/packages/internal-platform/src/schemas/workflows.ts @@ -11,6 +11,7 @@ export const WorkflowMetadataSchema = z.object({ name: z.string(), trigger: TriggerMetadataSchema, package: PackageMetadataSchema, + triggerTTL: z.number().optional(), }); export type WorkflowMetadata = z.infer; diff --git a/packages/trigger-sdk/src/client.ts b/packages/trigger-sdk/src/client.ts index 720a2c065..b97093c63 100644 --- a/packages/trigger-sdk/src/client.ts +++ b/packages/trigger-sdk/src/client.ts @@ -443,6 +443,7 @@ export class TriggerClient { trigger: this.#trigger.on.metadata, packageVersion: pkg.version, packageName: pkg.name, + triggerTTL: this.#options.triggerTTL, }); if (response?.type === "error") { diff --git a/packages/trigger-sdk/src/trigger/index.ts b/packages/trigger-sdk/src/trigger/index.ts index 78b653062..1c31970bd 100644 --- a/packages/trigger-sdk/src/trigger/index.ts +++ b/packages/trigger-sdk/src/trigger/index.ts @@ -12,6 +12,13 @@ export type TriggerOptions = { apiKey?: string; endpoint?: string; logLevel?: LogLevel; + + /** + * The TTL for the trigger in seconds. If the trigger is not run within this time, it will be aborted. Defaults to 3600 seconds (1 hour). + * @type {number} + */ + triggerTTL?: number; + run: (event: z.infer, ctx: TriggerContext) => Promise; };