diff --git a/apps/coordinator/src/server.ts b/apps/coordinator/src/server.ts index ee1e19f71..4cc6e9a7d 100644 --- a/apps/coordinator/src/server.ts +++ b/apps/coordinator/src/server.ts @@ -89,6 +89,39 @@ export class TriggerServer { sender: HostRPCSchema, receiver: ServerRPCSchema, handlers: { + INITIALIZE_DELAY: async (data) => { + if (!this.#triggerPublisher) { + // TODO: need to recover from this issue by trying to reconnect + return false; + } + + if (!this.#organizationId) { + // TODO: this should never really happen + throw new Error( + "Cannot complete workflow run without an organization ID" + ); + } + + if (!this.#workflowId) { + // TODO: this should never really happen + throw new Error("Cannot send log without a workflow ID"); + } + + const response = await this.#triggerPublisher.publish( + "INITIALIZE_DELAY", + { + id: data.waitId, + delay: data.delay, + }, + { + "x-api-key": this.#apiKey, + "x-workflow-id": this.#workflowId, + "x-workflow-run-id": data.id, + } + ); + + return !!response; + }, SEND_REQUEST: async (data) => { if (!this.#triggerPublisher) { // TODO: need to recover from this issue by trying to reconnect diff --git a/apps/webapp/app/models/workflowRun.server.ts b/apps/webapp/app/models/workflowRun.server.ts index d667590b8..d9e2e3bd2 100644 --- a/apps/webapp/app/models/workflowRun.server.ts +++ b/apps/webapp/app/models/workflowRun.server.ts @@ -134,24 +134,6 @@ export async function logMessageInRun( }); } -export async function initiateWaitInRun( - id: string, - wait: z.infer, - apiKey: string -) { - const workflowRun = await findWorkflowRunScopedToApiKey(id, apiKey); - - await prisma.workflowRunStep.create({ - data: { - runId: workflowRun.id, - type: "DURABLE_DELAY", - input: wait, - context: {}, - startedAt: new Date(), - }, - }); -} - async function findWorkflowRunScopedToApiKey(id: string, apiKey: string) { const workflowRun = await prisma.workflowRun.findFirst({ where: { id }, diff --git a/apps/webapp/app/services/delays/initiateDelay.server.ts b/apps/webapp/app/services/delays/initiateDelay.server.ts new file mode 100644 index 000000000..9319b26e4 --- /dev/null +++ b/apps/webapp/app/services/delays/initiateDelay.server.ts @@ -0,0 +1,12 @@ +import type { PrismaClient } from "~/db.server"; +import { prisma } from "~/db.server"; + +export class InitiateDelay { + #prismaClient: PrismaClient; + + constructor(prismaClient: PrismaClient = prisma) { + this.#prismaClient = prismaClient; + } + + async call(runId: string, delay: { id: string; seconds: number }) {} +} diff --git a/apps/webapp/app/services/messageBroker.server.ts b/apps/webapp/app/services/messageBroker.server.ts index 11ed535f4..5f129c842 100644 --- a/apps/webapp/app/services/messageBroker.server.ts +++ b/apps/webapp/app/services/messageBroker.server.ts @@ -16,11 +16,11 @@ import { completeWorkflowRun, failWorkflowRun, findWorklowRunById, - initiateWaitInRun, logMessageInRun, startWorkflowRun, triggerEventInRun, } from "~/models/workflowRun.server"; +import { InitiateDelay } from "./delays/initiateDelay.server"; import { DispatchEvent } from "./events/dispatch.server"; import { RegisterExternalSource } from "./externalSources/registerExternalSource.server"; import { CreateIntegrationRequest } from "./requests/createIntegrationRequest.server"; @@ -177,8 +177,13 @@ async function createTriggerSubscriber() { return true; }, - INITIATE_WAIT: async (id, data, properties) => { - await initiateWaitInRun(data.id, data.wait, properties["x-api-key"]); + INITIALIZE_DELAY: async (id, data, properties) => { + const service = new InitiateDelay(); + + await service.call(properties["x-workflow-run-id"], { + id: data.id, + seconds: data.delay, + }); return true; }, diff --git a/examples/send-to-slack/src/index.ts b/examples/send-to-slack/src/index.ts index 2b6c04182..dc4ebbf4c 100644 --- a/examples/send-to-slack/src/index.ts +++ b/examples/send-to-slack/src/index.ts @@ -17,6 +17,8 @@ const trigger = new Trigger({ }), }), run: async (event, ctx) => { + await ctx.waitFor(60); + const response = await slack.postMessage({ channel: "test-integrations", text: `New domain created: ${event.domain} by customer ${event.customerId}`, diff --git a/packages/internal-bridge/src/schemas/server.ts b/packages/internal-bridge/src/schemas/server.ts index ef3ccdf5b..51232ea99 100644 --- a/packages/internal-bridge/src/schemas/server.ts +++ b/packages/internal-bridge/src/schemas/server.ts @@ -5,6 +5,14 @@ import { import { z } from "zod"; export const ServerRPCSchema = { + INITIALIZE_DELAY: { + request: z.object({ + id: z.string(), + waitId: z.string(), + delay: z.number(), + }), + response: z.boolean(), + }, SEND_REQUEST: { request: z.object({ id: z.string(), diff --git a/packages/internal-platform/src/messages/schemas/awaits.ts b/packages/internal-platform/src/messages/schemas/awaits.ts index 4a286b4a0..d735df634 100644 --- a/packages/internal-platform/src/messages/schemas/awaits.ts +++ b/packages/internal-platform/src/messages/schemas/awaits.ts @@ -1,16 +1,14 @@ import { WaitSchema } from "@trigger.dev/common-schemas"; import { z } from "zod"; +import { WorkflowSendRunEventPropertiesSchema } from "../sharedSchemas"; const Catalog = { - INITIATE_WAIT: { + INITIALIZE_DELAY: { data: z.object({ id: z.string(), - wait: WaitSchema, - }), - properties: z.object({ - "x-workflow-id": z.string(), - "x-api-key": z.string(), + delay: z.number(), }), + properties: WorkflowSendRunEventPropertiesSchema, }, }; diff --git a/packages/trigger-sdk/src/client.ts b/packages/trigger-sdk/src/client.ts index 0f7bdae02..c828920f3 100644 --- a/packages/trigger-sdk/src/client.ts +++ b/packages/trigger-sdk/src/client.ts @@ -15,12 +15,6 @@ import { ContextLogger } from "./logger"; import { triggerRunLocalStorage } from "./localStorage"; import { ulid } from "ulid"; -type RequestResponse = { - body?: any; - headers: Record; - status: number; -}; - export class TriggerClient { #trigger: Trigger; #options: TriggerOptions; @@ -43,6 +37,14 @@ export class TriggerClient { } >(); + #waitForCallbacks = new Map< + string, + { + resolve: () => void; + reject: (err?: any) => void; + } + >(); + constructor(trigger: Trigger, options: TriggerOptions) { this.#trigger = trigger; this.#options = options; @@ -146,6 +148,26 @@ export class TriggerClient { event: JSON.parse(JSON.stringify(event)), }); }, + waitFor: async (seconds: number) => { + const waitId = ulid(); + + const result = new Promise((resolve, reject) => { + this.#waitForCallbacks.set(waitId, { + resolve, + reject, + }); + }); + + await serverRPC.send("INITIALIZE_DELAY", { + id: data.id, + waitId, + delay: seconds, + }); + + await result; + + return; + }, }; const eventData = this.#options.on.schema.parse(data.trigger.input); diff --git a/packages/trigger-sdk/src/types.ts b/packages/trigger-sdk/src/types.ts index 32ba0508d..97d79e89c 100644 --- a/packages/trigger-sdk/src/types.ts +++ b/packages/trigger-sdk/src/types.ts @@ -10,6 +10,7 @@ export interface TriggerContext { organizationId: string; logger: TriggerLogger; fireEvent(event: CustomEvent): Promise; + waitFor(seconds: number): Promise; } export interface TriggerLogger {