From d660c581f49620db0832a6c4d17e2f547588f457 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Mon, 2 Jan 2023 16:17:28 +0000 Subject: [PATCH] Add support for durable delays (either durations or scheduled times) --- apps/coordinator/src/server.ts | 34 ++++++++++- .../services/delays/initiateDelay.server.ts | 56 ++++++++++++++++++- .../services/delays/resolveDelay.server.ts | 33 +++++++++++ .../app/services/messageBroker.server.ts | 25 ++++++++- apps/webapp/app/utils/delays.ts | 13 +++++ .../migration.sql | 20 +++++++ apps/webapp/prisma/schema.prisma | 23 +++++++- examples/send-to-slack/src/index.ts | 10 +++- packages/common-schemas/src/waits.ts | 5 +- packages/internal-bridge/src/schemas/host.ts | 13 +++++ .../internal-bridge/src/schemas/server.ts | 3 +- .../src/messages/catalogs/coordinator.ts | 24 ++++---- .../src/messages/catalogs/platform.ts | 10 ++-- .../src/messages/schemas/awaits.ts | 15 ----- .../messages/schemas/completeWorkflowRun.ts | 16 ------ ...{triggerCustomEvent.ts => customEvents.ts} | 4 +- .../src/messages/schemas/delays.ts | 25 +++++++++ .../src/messages/schemas/failWorkflowRun.ts | 17 ------ .../schemas/finishIntegrationRequest.ts | 15 ----- .../messages/schemas/integrationRequests.ts | 28 ++++++++++ .../schemas/{logMessage.ts => logs.ts} | 4 +- .../schemas/sendIntegrationRequest.ts | 16 ------ .../src/messages/schemas/startWorkflowRun.ts | 15 ----- .../src/messages/schemas/workflowRuns.ts | 34 +++++++++++ .../{triggerWorkflow.ts => workflows.ts} | 4 +- packages/trigger-sdk/src/client.ts | 52 ++++++++++++++++- packages/trigger-sdk/src/types.ts | 10 +++- 27 files changed, 390 insertions(+), 134 deletions(-) create mode 100644 apps/webapp/app/services/delays/resolveDelay.server.ts create mode 100644 apps/webapp/app/utils/delays.ts create mode 100644 apps/webapp/prisma/migrations/20230102142756_add_durable_delay_model/migration.sql delete mode 100644 packages/internal-platform/src/messages/schemas/awaits.ts delete mode 100644 packages/internal-platform/src/messages/schemas/completeWorkflowRun.ts rename packages/internal-platform/src/messages/schemas/{triggerCustomEvent.ts => customEvents.ts} (87%) create mode 100644 packages/internal-platform/src/messages/schemas/delays.ts delete mode 100644 packages/internal-platform/src/messages/schemas/failWorkflowRun.ts delete mode 100644 packages/internal-platform/src/messages/schemas/finishIntegrationRequest.ts create mode 100644 packages/internal-platform/src/messages/schemas/integrationRequests.ts rename packages/internal-platform/src/messages/schemas/{logMessage.ts => logs.ts} (87%) delete mode 100644 packages/internal-platform/src/messages/schemas/sendIntegrationRequest.ts delete mode 100644 packages/internal-platform/src/messages/schemas/startWorkflowRun.ts create mode 100644 packages/internal-platform/src/messages/schemas/workflowRuns.ts rename packages/internal-platform/src/messages/schemas/{triggerWorkflow.ts => workflows.ts} (90%) diff --git a/apps/coordinator/src/server.ts b/apps/coordinator/src/server.ts index 82372736e..e2fa9cf4a 100644 --- a/apps/coordinator/src/server.ts +++ b/apps/coordinator/src/server.ts @@ -111,7 +111,7 @@ export class TriggerServer { "INITIALIZE_DELAY", { id: data.waitId, - delay: data.delay, + config: data.config, }, { "x-api-key": this.#apiKey, @@ -373,6 +373,38 @@ export class TriggerServer { subscriptionInitialPosition: "Earliest", }, handlers: { + RESOLVE_DELAY: async (id, data, properties) => { + this.#logger.debug("Received resolve delay", id, data, properties); + + if (!this.#serverRPC) { + throw new Error("Cannot resolve delay without an RPC connection"); + } + + // If the API keys don't match, then we should ignore it + // This ensures the workflow is triggered for the correct environment + if (properties["x-api-key"] !== this.#apiKey) { + return true; + } + + // If the workflow id is not the same as the workflow id + // that we are listening for, then we should ignore it + if (properties["x-workflow-id"] !== this.#workflowId) { + return true; + } + + const success = await this.#serverRPC.send("RESOLVE_DELAY", { + id: data.id, + meta: { + workflowId: properties["x-workflow-id"], + organizationId: properties["x-org-id"], + environment: properties["x-env"], + apiKey: properties["x-api-key"], + runId: properties["x-workflow-run-id"], + }, + }); + + return success; + }, RESOLVE_INTEGRATION_REQUEST: async (id, data, properties) => { this.#logger.debug( "Received finish integration request", diff --git a/apps/webapp/app/services/delays/initiateDelay.server.ts b/apps/webapp/app/services/delays/initiateDelay.server.ts index 9319b26e4..e7484ed0e 100644 --- a/apps/webapp/app/services/delays/initiateDelay.server.ts +++ b/apps/webapp/app/services/delays/initiateDelay.server.ts @@ -1,5 +1,11 @@ +import type { WaitSchema } from "@trigger.dev/common-schemas"; +import type { z } from "zod"; import type { PrismaClient } from "~/db.server"; import { prisma } from "~/db.server"; +import { calculateDurationInMs } from "~/utils/delays"; +import { internalPubSub } from "../messageBroker.server"; + +type DelayConfig = z.infer; export class InitiateDelay { #prismaClient: PrismaClient; @@ -8,5 +14,53 @@ export class InitiateDelay { this.#prismaClient = prismaClient; } - async call(runId: string, delay: { id: string; seconds: number }) {} + async call(runId: string, delay: { id: string; config: DelayConfig }) { + const delayUntil = this.#calculateDelayUntil(delay.config); + + // Make sure the delay is not more than 1 year in the future + if (delayUntil.getTime() > Date.now() + 365 * 24 * 60 * 60 * 1000) { + throw new Error( + `Delay is more than 1 year in the future, which is the maximum allowed by trigger.dev` + ); + } + + const workflowStep = await this.#prismaClient.workflowRunStep.create({ + data: { + runId, + type: "DURABLE_DELAY", + input: delay.config, + context: { id: delay.id, delayUntil: delayUntil.toISOString() }, + status: "RUNNING", + startedAt: new Date(), + }, + }); + + // Create the durable delay + const durableDelay = await this.#prismaClient.durableDelay.create({ + data: { + id: delay.id, + runId, + stepId: workflowStep.id, + delayUntil: this.#calculateDelayUntil(delay.config), + }, + }); + + await internalPubSub.publish( + "RESOLVE_DELAY", + { + id: delay.id, + }, + {}, + { deliverAt: durableDelay.delayUntil.getTime() } + ); + } + + #calculateDelayUntil(config: DelayConfig): Date { + switch (config.type) { + case "DELAY": + return new Date(Date.now() + calculateDurationInMs(config)); + case "SCHEDULE_FOR": + return new Date(config.scheduledFor); + } + } } diff --git a/apps/webapp/app/services/delays/resolveDelay.server.ts b/apps/webapp/app/services/delays/resolveDelay.server.ts new file mode 100644 index 000000000..a96893891 --- /dev/null +++ b/apps/webapp/app/services/delays/resolveDelay.server.ts @@ -0,0 +1,33 @@ +import type { PrismaClient } from "~/db.server"; +import { prisma } from "~/db.server"; + +export class ResolveDelay { + #prismaClient: PrismaClient; + + constructor(prismaClient: PrismaClient = prisma) { + this.#prismaClient = prismaClient; + } + + async call(id: string) { + const delay = await this.#prismaClient.durableDelay.update({ + where: { id }, + data: { resolvedAt: new Date() }, + include: { + step: { + include: { + run: { + include: { environment: true }, + }, + }, + }, + }, + }); + + await this.#prismaClient.workflowRunStep.update({ + where: { id: delay.step.id }, + data: { status: "SUCCESS", finishedAt: delay.resolvedAt }, + }); + + return delay.step.run; + } +} diff --git a/apps/webapp/app/services/messageBroker.server.ts b/apps/webapp/app/services/messageBroker.server.ts index bb901b94f..b9af35cb9 100644 --- a/apps/webapp/app/services/messageBroker.server.ts +++ b/apps/webapp/app/services/messageBroker.server.ts @@ -28,6 +28,7 @@ import { PerformIntegrationRequest } from "./requests/performIntegrationRequest. import { StartIntegrationRequest } from "./requests/startIntegrationRequest.server"; import { WaitForConnection } from "./requests/waitForConnection.server"; import { HandleNewServiceConnection } from "./externalServices/handleNewConnection.server"; +import { ResolveDelay } from "./delays/resolveDelay.server"; let pulsarClient: PulsarClient; let triggerPublisher: ZodPublisher; @@ -183,7 +184,7 @@ async function createTriggerSubscriber() { await service.call(properties["x-workflow-run-id"], { id: data.id, - seconds: data.delay, + config: data.config, }); return true; @@ -273,6 +274,10 @@ const InternalCatalog = { data: z.object({ id: z.string() }), properties: z.object({}), }, + RESOLVE_DELAY: { + data: z.object({ id: z.string() }), + properties: z.object({}), + }, }; async function createInternalPubSub() { @@ -288,6 +293,24 @@ async function createInternalPubSub() { }, schema: InternalCatalog, handlers: { + RESOLVE_DELAY: async (id, data, properties) => { + const service = new ResolveDelay(); + + const run = await service.call(data.id); + + triggerPublisher.publish( + "RESOLVE_DELAY", + { id: data.id }, + { + "x-workflow-run-id": run.id, + "x-api-key": run.environment.apiKey, + "x-org-id": run.environment.organizationId, + "x-workflow-id": run.workflowId, + "x-env": run.environment.slug, + } + ); + return true; + }, INTEGRATION_REQUEST_CREATED: async (id, data, properties) => { const integrationRequest = await findIntegrationRequestById(data.id); diff --git a/apps/webapp/app/utils/delays.ts b/apps/webapp/app/utils/delays.ts new file mode 100644 index 000000000..6faa67c67 --- /dev/null +++ b/apps/webapp/app/utils/delays.ts @@ -0,0 +1,13 @@ +export const calculateDurationInMs = (options: { + seconds?: number; + minutes?: number; + hours?: number; + days?: number; +}) => { + return ( + (options?.seconds ?? 0) * 1000 + + (options?.minutes ?? 0) * 60 * 1000 + + (options?.hours ?? 0) * 60 * 60 * 1000 + + (options?.days ?? 0) * 24 * 60 * 60 * 1000 + ); +}; diff --git a/apps/webapp/prisma/migrations/20230102142756_add_durable_delay_model/migration.sql b/apps/webapp/prisma/migrations/20230102142756_add_durable_delay_model/migration.sql new file mode 100644 index 000000000..dedca8fe5 --- /dev/null +++ b/apps/webapp/prisma/migrations/20230102142756_add_durable_delay_model/migration.sql @@ -0,0 +1,20 @@ +-- CreateTable +CREATE TABLE "DurableDelay" ( + "id" TEXT NOT NULL, + "runId" TEXT NOT NULL, + "stepId" TEXT NOT NULL, + "delayUntil" TIMESTAMP(3) NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "resolvedAt" TIMESTAMP(3), + + CONSTRAINT "DurableDelay_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE UNIQUE INDEX "DurableDelay_stepId_key" ON "DurableDelay"("stepId"); + +-- AddForeignKey +ALTER TABLE "DurableDelay" ADD CONSTRAINT "DurableDelay_runId_fkey" FOREIGN KEY ("runId") REFERENCES "WorkflowRun"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "DurableDelay" ADD CONSTRAINT "DurableDelay_stepId_fkey" FOREIGN KEY ("stepId") REFERENCES "WorkflowRunStep"("id") ON DELETE CASCADE ON UPDATE CASCADE; diff --git a/apps/webapp/prisma/schema.prisma b/apps/webapp/prisma/schema.prisma index 4aa35e956..5b29a9b3b 100644 --- a/apps/webapp/prisma/schema.prisma +++ b/apps/webapp/prisma/schema.prisma @@ -59,11 +59,11 @@ model APIConnection { apiIdentifier String status APIConnectionStatus @default(CREATED) scopes String[] - + authenticationMethod APIAuthenticationMethod @default(OAUTH) authenticationConfig Json? - type APIConnectionType + type APIConnectionType createdAt DateTime @default(now()) updatedAt DateTime @updatedAt @@ -293,6 +293,21 @@ model IntegrationResponse { createdAt DateTime @default(now()) } +model DurableDelay { + id String @id + + run WorkflowRun @relation(fields: [runId], references: [id], onDelete: Cascade, onUpdate: Cascade) + runId String + + step WorkflowRunStep @relation(fields: [stepId], references: [id], onDelete: Cascade, onUpdate: Cascade) + stepId String @unique + + delayUntil DateTime + + createdAt DateTime @default(now()) + resolvedAt DateTime? +} + model TriggerEvent { id String @id @default(cuid()) service String @@ -351,6 +366,7 @@ model WorkflowRun { isTest Boolean @default(false) requests IntegrationRequest[] + delays DurableDelay[] } enum WorkflowRunStatus { @@ -380,6 +396,7 @@ model WorkflowRunStep { status WorkflowRunStepStatus @default(PENDING) integrationRequest IntegrationRequest? + delay DurableDelay? } enum WorkflowRunStepStatus { @@ -395,4 +412,4 @@ enum WorkflowRunStepType { DURABLE_DELAY CUSTOM_EVENT INTEGRATION_REQUEST -} \ No newline at end of file +} diff --git a/examples/send-to-slack/src/index.ts b/examples/send-to-slack/src/index.ts index 604b85066..17cf4d62a 100644 --- a/examples/send-to-slack/src/index.ts +++ b/examples/send-to-slack/src/index.ts @@ -17,13 +17,21 @@ const trigger = new Trigger({ }), }), run: async (event, ctx) => { - // await ctx.waitFor(60); + await ctx.logger.info( + "Received domain.created event, waiting for 60 seconds..." + ); + + await ctx.waitFor({ seconds: 60 }); + + await ctx.logger.info("Posting to Slack..."); const response = await slack.postMessage({ channel: "test-integrations", text: `New domain created: ${event.domain} by customer ${event.customerId}`, }); + await ctx.logger.info("Posted to Slack!"); + return response.message; }, }); diff --git a/packages/common-schemas/src/waits.ts b/packages/common-schemas/src/waits.ts index 40203f57f..d2724fda7 100644 --- a/packages/common-schemas/src/waits.ts +++ b/packages/common-schemas/src/waits.ts @@ -2,7 +2,10 @@ import { z } from "zod"; export const DelaySchema = z.object({ type: z.literal("DELAY"), - durationInMs: z.number(), + seconds: z.number().optional(), + minutes: z.number().optional(), + hours: z.number().optional(), + days: z.number().optional(), }); export const ScheduledForSchema = z.object({ diff --git a/packages/internal-bridge/src/schemas/host.ts b/packages/internal-bridge/src/schemas/host.ts index f173de5a4..1803b26fb 100644 --- a/packages/internal-bridge/src/schemas/host.ts +++ b/packages/internal-bridge/src/schemas/host.ts @@ -32,6 +32,19 @@ export const HostRPCSchema = { }), response: z.boolean(), }, + RESOLVE_DELAY: { + request: z.object({ + id: z.string(), + meta: z.object({ + environment: z.string(), + workflowId: z.string(), + organizationId: z.string(), + apiKey: z.string(), + runId: z.string(), + }), + }), + response: z.boolean(), + }, }; export type HostRPC = typeof HostRPCSchema; diff --git a/packages/internal-bridge/src/schemas/server.ts b/packages/internal-bridge/src/schemas/server.ts index 51232ea99..56c3b9838 100644 --- a/packages/internal-bridge/src/schemas/server.ts +++ b/packages/internal-bridge/src/schemas/server.ts @@ -1,6 +1,7 @@ import { CustomEventSchema, TriggerMetadataSchema, + WaitSchema, } from "@trigger.dev/common-schemas"; import { z } from "zod"; @@ -9,7 +10,7 @@ export const ServerRPCSchema = { request: z.object({ id: z.string(), waitId: z.string(), - delay: z.number(), + config: WaitSchema, }), response: z.boolean(), }, diff --git a/packages/internal-platform/src/messages/catalogs/coordinator.ts b/packages/internal-platform/src/messages/catalogs/coordinator.ts index bc7e1c4c0..6517ce00f 100644 --- a/packages/internal-platform/src/messages/catalogs/coordinator.ts +++ b/packages/internal-platform/src/messages/catalogs/coordinator.ts @@ -1,19 +1,15 @@ -import sendIntegrationRequest from "../schemas/sendIntegrationRequest"; -import startWorklowRun from "../schemas/startWorkflowRun"; -import failWorkflowRun from "../schemas/failWorkflowRun"; -import completeWorkflowRun from "../schemas/completeWorkflowRun"; -import logMessage from "../schemas/logMessage"; -import triggerCustomEvent from "../schemas/triggerCustomEvent"; -import awaits from "../schemas/awaits"; +import { coordinator as integrationRequests } from "../schemas/integrationRequests"; +import { coordinator as workflowRuns } from "../schemas/workflowRuns"; +import { coordinator as logs } from "../schemas/logs"; +import { coordinator as customEvents } from "../schemas/customEvents"; +import { coordinator as delays } from "../schemas/delays"; const Catalog = { - ...sendIntegrationRequest, - ...startWorklowRun, - ...failWorkflowRun, - ...completeWorkflowRun, - ...logMessage, - ...triggerCustomEvent, - ...awaits, + ...integrationRequests, + ...workflowRuns, + ...logs, + ...customEvents, + ...delays, }; export default Catalog; diff --git a/packages/internal-platform/src/messages/catalogs/platform.ts b/packages/internal-platform/src/messages/catalogs/platform.ts index f3232a9ac..d687056e8 100644 --- a/packages/internal-platform/src/messages/catalogs/platform.ts +++ b/packages/internal-platform/src/messages/catalogs/platform.ts @@ -1,9 +1,11 @@ -import triggerWorkflow from "../schemas/triggerWorkflow"; -import finishIntegrationRequest from "../schemas/finishIntegrationRequest"; +import { platform as workflows } from "../schemas/workflows"; +import { platform as integrationRequests } from "../schemas/integrationRequests"; +import { platform as delays } from "../schemas/delays"; const Catalog = { - ...triggerWorkflow, - ...finishIntegrationRequest, + ...workflows, + ...integrationRequests, + ...delays, }; export default Catalog; diff --git a/packages/internal-platform/src/messages/schemas/awaits.ts b/packages/internal-platform/src/messages/schemas/awaits.ts deleted file mode 100644 index d735df634..000000000 --- a/packages/internal-platform/src/messages/schemas/awaits.ts +++ /dev/null @@ -1,15 +0,0 @@ -import { WaitSchema } from "@trigger.dev/common-schemas"; -import { z } from "zod"; -import { WorkflowSendRunEventPropertiesSchema } from "../sharedSchemas"; - -const Catalog = { - INITIALIZE_DELAY: { - data: z.object({ - id: z.string(), - delay: z.number(), - }), - properties: WorkflowSendRunEventPropertiesSchema, - }, -}; - -export default Catalog; diff --git a/packages/internal-platform/src/messages/schemas/completeWorkflowRun.ts b/packages/internal-platform/src/messages/schemas/completeWorkflowRun.ts deleted file mode 100644 index c3e24de1d..000000000 --- a/packages/internal-platform/src/messages/schemas/completeWorkflowRun.ts +++ /dev/null @@ -1,16 +0,0 @@ -import { z } from "zod"; - -const Catalog = { - COMPLETE_WORKFLOW_RUN: { - data: z.object({ - id: z.string(), - output: z.string(), - }), - properties: z.object({ - "x-workflow-id": z.string(), - "x-api-key": z.string(), - }), - }, -}; - -export default Catalog; diff --git a/packages/internal-platform/src/messages/schemas/triggerCustomEvent.ts b/packages/internal-platform/src/messages/schemas/customEvents.ts similarity index 87% rename from packages/internal-platform/src/messages/schemas/triggerCustomEvent.ts rename to packages/internal-platform/src/messages/schemas/customEvents.ts index d5e0e1de3..59c920914 100644 --- a/packages/internal-platform/src/messages/schemas/triggerCustomEvent.ts +++ b/packages/internal-platform/src/messages/schemas/customEvents.ts @@ -1,7 +1,7 @@ import { z } from "zod"; import { CustomEventSchema } from "@trigger.dev/common-schemas"; -const Catalog = { +export const coordinator = { TRIGGER_CUSTOM_EVENT: { data: z.object({ id: z.string(), @@ -13,5 +13,3 @@ const Catalog = { }), }, }; - -export default Catalog; diff --git a/packages/internal-platform/src/messages/schemas/delays.ts b/packages/internal-platform/src/messages/schemas/delays.ts new file mode 100644 index 000000000..f4aac921d --- /dev/null +++ b/packages/internal-platform/src/messages/schemas/delays.ts @@ -0,0 +1,25 @@ +import { WaitSchema } from "@trigger.dev/common-schemas"; +import { z } from "zod"; +import { + WorkflowRunEventPropertiesSchema, + WorkflowSendRunEventPropertiesSchema, +} from "../sharedSchemas"; + +export const coordinator = { + INITIALIZE_DELAY: { + data: z.object({ + id: z.string(), + config: WaitSchema, + }), + properties: WorkflowSendRunEventPropertiesSchema, + }, +}; + +export const platform = { + RESOLVE_DELAY: { + data: z.object({ + id: z.string(), + }), + properties: WorkflowRunEventPropertiesSchema, + }, +}; diff --git a/packages/internal-platform/src/messages/schemas/failWorkflowRun.ts b/packages/internal-platform/src/messages/schemas/failWorkflowRun.ts deleted file mode 100644 index 826bfc8d1..000000000 --- a/packages/internal-platform/src/messages/schemas/failWorkflowRun.ts +++ /dev/null @@ -1,17 +0,0 @@ -import { z } from "zod"; -import { ErrorSchema } from "@trigger.dev/common-schemas"; - -const Catalog = { - FAIL_WORKFLOW_RUN: { - data: z.object({ - id: z.string(), - error: ErrorSchema, - }), - properties: z.object({ - "x-workflow-id": z.string(), - "x-api-key": z.string(), - }), - }, -}; - -export default Catalog; diff --git a/packages/internal-platform/src/messages/schemas/finishIntegrationRequest.ts b/packages/internal-platform/src/messages/schemas/finishIntegrationRequest.ts deleted file mode 100644 index a787dab75..000000000 --- a/packages/internal-platform/src/messages/schemas/finishIntegrationRequest.ts +++ /dev/null @@ -1,15 +0,0 @@ -import { JsonSchema } from "@trigger.dev/common-schemas"; -import { z } from "zod"; -import { WorkflowRunEventPropertiesSchema } from "../sharedSchemas"; - -const Catalog = { - RESOLVE_INTEGRATION_REQUEST: { - data: z.object({ - id: z.string(), - output: JsonSchema.default({}), - }), - properties: WorkflowRunEventPropertiesSchema, - }, -}; - -export default Catalog; diff --git a/packages/internal-platform/src/messages/schemas/integrationRequests.ts b/packages/internal-platform/src/messages/schemas/integrationRequests.ts new file mode 100644 index 000000000..57c8199f9 --- /dev/null +++ b/packages/internal-platform/src/messages/schemas/integrationRequests.ts @@ -0,0 +1,28 @@ +import { JsonSchema } from "@trigger.dev/common-schemas"; +import { z } from "zod"; +import { + WorkflowRunEventPropertiesSchema, + WorkflowSendRunEventPropertiesSchema, +} from "../sharedSchemas"; + +export const platform = { + RESOLVE_INTEGRATION_REQUEST: { + data: z.object({ + id: z.string(), + output: JsonSchema.default({}), + }), + properties: WorkflowRunEventPropertiesSchema, + }, +}; + +export const coordinator = { + SEND_INTEGRATION_REQUEST: { + data: z.object({ + id: z.string(), + service: z.string(), + endpoint: z.string(), + params: z.any(), + }), + properties: WorkflowSendRunEventPropertiesSchema, + }, +}; diff --git a/packages/internal-platform/src/messages/schemas/logMessage.ts b/packages/internal-platform/src/messages/schemas/logs.ts similarity index 87% rename from packages/internal-platform/src/messages/schemas/logMessage.ts rename to packages/internal-platform/src/messages/schemas/logs.ts index d450ebca4..b070b68db 100644 --- a/packages/internal-platform/src/messages/schemas/logMessage.ts +++ b/packages/internal-platform/src/messages/schemas/logs.ts @@ -1,7 +1,7 @@ import { LogMessageSchema } from "@trigger.dev/common-schemas"; import { z } from "zod"; -const Catalog = { +export const coordinator = { LOG_MESSAGE: { data: z.object({ id: z.string(), @@ -13,5 +13,3 @@ const Catalog = { }), }, }; - -export default Catalog; diff --git a/packages/internal-platform/src/messages/schemas/sendIntegrationRequest.ts b/packages/internal-platform/src/messages/schemas/sendIntegrationRequest.ts deleted file mode 100644 index ea8c72b85..000000000 --- a/packages/internal-platform/src/messages/schemas/sendIntegrationRequest.ts +++ /dev/null @@ -1,16 +0,0 @@ -import { z } from "zod"; -import { WorkflowSendRunEventPropertiesSchema } from "../sharedSchemas"; - -const Catalog = { - SEND_INTEGRATION_REQUEST: { - data: z.object({ - id: z.string(), - service: z.string(), - endpoint: z.string(), - params: z.any(), - }), - properties: WorkflowSendRunEventPropertiesSchema, - }, -}; - -export default Catalog; diff --git a/packages/internal-platform/src/messages/schemas/startWorkflowRun.ts b/packages/internal-platform/src/messages/schemas/startWorkflowRun.ts deleted file mode 100644 index 5bfc44247..000000000 --- a/packages/internal-platform/src/messages/schemas/startWorkflowRun.ts +++ /dev/null @@ -1,15 +0,0 @@ -import { z } from "zod"; - -const Catalog = { - START_WORKFLOW_RUN: { - data: z.object({ - id: z.string(), - }), - properties: z.object({ - "x-workflow-id": z.string(), - "x-api-key": z.string(), - }), - }, -}; - -export default Catalog; diff --git a/packages/internal-platform/src/messages/schemas/workflowRuns.ts b/packages/internal-platform/src/messages/schemas/workflowRuns.ts new file mode 100644 index 000000000..ac912127b --- /dev/null +++ b/packages/internal-platform/src/messages/schemas/workflowRuns.ts @@ -0,0 +1,34 @@ +import { ErrorSchema } from "@trigger.dev/common-schemas"; +import { z } from "zod"; + +export const coordinator = { + COMPLETE_WORKFLOW_RUN: { + data: z.object({ + id: z.string(), + output: z.string(), + }), + properties: z.object({ + "x-workflow-id": z.string(), + "x-api-key": z.string(), + }), + }, + FAIL_WORKFLOW_RUN: { + data: z.object({ + id: z.string(), + error: ErrorSchema, + }), + properties: z.object({ + "x-workflow-id": z.string(), + "x-api-key": z.string(), + }), + }, + START_WORKFLOW_RUN: { + data: z.object({ + id: z.string(), + }), + properties: z.object({ + "x-workflow-id": z.string(), + "x-api-key": z.string(), + }), + }, +}; diff --git a/packages/internal-platform/src/messages/schemas/triggerWorkflow.ts b/packages/internal-platform/src/messages/schemas/workflows.ts similarity index 90% rename from packages/internal-platform/src/messages/schemas/triggerWorkflow.ts rename to packages/internal-platform/src/messages/schemas/workflows.ts index 34b72968e..e9a9d7f70 100644 --- a/packages/internal-platform/src/messages/schemas/triggerWorkflow.ts +++ b/packages/internal-platform/src/messages/schemas/workflows.ts @@ -8,11 +8,9 @@ export const TriggerWorkflowMessageSchema = z.object({ context: JsonSchema.default({}), }); -const Catalog = { +export const platform = { TRIGGER_WORKFLOW: { data: TriggerWorkflowMessageSchema, properties: WorkflowEventPropertiesSchema, }, }; - -export default Catalog; diff --git a/packages/trigger-sdk/src/client.ts b/packages/trigger-sdk/src/client.ts index 1a567ad16..d5e4c6f89 100644 --- a/packages/trigger-sdk/src/client.ts +++ b/packages/trigger-sdk/src/client.ts @@ -10,7 +10,7 @@ import { } from "internal-bridge"; import * as pkg from "../package.json"; import { Trigger, TriggerOptions } from "./trigger"; -import { TriggerContext } from "./types"; +import { TriggerContext, WaitForOptions } from "./types"; import { ContextLogger } from "./logger"; import { triggerRunLocalStorage } from "./localStorage"; import { ulid } from "ulid"; @@ -109,6 +109,23 @@ export class TriggerClient { sender: ServerRPCSchema, receiver: HostRPCSchema, handlers: { + RESOLVE_DELAY: async (data) => { + console.log(`RESOLVE_DELAY(${data.id})`); + + const waitCallbacks = this.#waitForCallbacks.get(data.id); + + if (!waitCallbacks) { + throw new Error( + `Could not find wait callbacks for wait ID ${data.id}` + ); + } + + const { resolve, reject } = waitCallbacks; + + resolve(); + + return true; + }, RESOLVE_REQUEST: async (data) => { const requestCallbacks = this.#responseCompleteCallbacks.get(data.id); @@ -148,7 +165,7 @@ export class TriggerClient { event: JSON.parse(JSON.stringify(event)), }); }, - waitFor: async (seconds: number) => { + waitFor: async (options: WaitForOptions) => { const waitId = ulid(); const result = new Promise((resolve, reject) => { @@ -161,7 +178,36 @@ export class TriggerClient { await serverRPC.send("INITIALIZE_DELAY", { id: data.id, waitId, - delay: seconds, + config: { + type: "DELAY", + seconds: options.seconds, + minutes: options.minutes, + hours: options.hours, + days: options.days, + }, + }); + + await result; + + return; + }, + waitUntil: async (date: Date) => { + const waitId = ulid(); + + const result = new Promise((resolve, reject) => { + this.#waitForCallbacks.set(waitId, { + resolve, + reject, + }); + }); + + await serverRPC.send("INITIALIZE_DELAY", { + id: data.id, + waitId, + config: { + type: "SCHEDULE_FOR", + scheduledFor: date.toISOString(), + }, }); await result; diff --git a/packages/trigger-sdk/src/types.ts b/packages/trigger-sdk/src/types.ts index 97d79e89c..8894f6c38 100644 --- a/packages/trigger-sdk/src/types.ts +++ b/packages/trigger-sdk/src/types.ts @@ -3,6 +3,13 @@ import { z } from "zod"; type CustomEvent = z.infer; +export type WaitForOptions = { + seconds?: number; + minutes?: number; + hours?: number; + days?: number; +}; + export interface TriggerContext { id: string; environment: string; @@ -10,7 +17,8 @@ export interface TriggerContext { organizationId: string; logger: TriggerLogger; fireEvent(event: CustomEvent): Promise; - waitFor(seconds: number): Promise; + waitFor(options: WaitForOptions): Promise; + waitUntil(date: Date): Promise; } export interface TriggerLogger {