From e37a2001e1767340b0addfb8f1e96ca4420a000c Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Thu, 26 Jan 2023 10:12:57 +0000 Subject: [PATCH] Added lastRunAt to the scheduleEvent payload --- .changeset/calm-dingos-develop.md | 5 ++ apps/docs/mint.json | 19 ++++-- apps/docs/reference/custom-event.mdx | 2 + apps/docs/reference/schedule-event.mdx | 68 +++++++++++++++++++ apps/docs/reference/trigger.mdx | 7 +- apps/docs/reference/webhook-event.mdx | 2 + apps/docs/triggers/scheduled.mdx | 32 +++++++++ .../scheduler/scheduleNextEvent.server.ts | 8 ++- apps/webapp/app/triggers/monitoring.server.ts | 13 ++-- examples/schedule-to-slack/src/index.ts | 15 ++-- packages/common-schemas/src/events.ts | 1 + 11 files changed, 148 insertions(+), 24 deletions(-) create mode 100644 .changeset/calm-dingos-develop.md create mode 100644 apps/docs/reference/schedule-event.mdx diff --git a/.changeset/calm-dingos-develop.md b/.changeset/calm-dingos-develop.md new file mode 100644 index 000000000..df988a40f --- /dev/null +++ b/.changeset/calm-dingos-develop.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/sdk": patch +--- + +Added lastRunAt to the scheduleEvent payload diff --git a/apps/docs/mint.json b/apps/docs/mint.json index 7dd0e81b0..41e403a37 100644 --- a/apps/docs/mint.json +++ b/apps/docs/mint.json @@ -39,7 +39,11 @@ "navigation": [ { "group": "Getting Started", - "pages": ["welcome", "getting-started", "get-help"] + "pages": [ + "welcome", + "getting-started", + "get-help" + ] }, { "group": "Examples", @@ -67,11 +71,15 @@ "pages": [ { "group": "Slack", - "pages": ["integrations/apis/slack/actions/post-message"] + "pages": [ + "integrations/apis/slack/actions/post-message" + ] }, { "group": "Resend.com", - "pages": ["integrations/apis/resend/actions/send-email"] + "pages": [ + "integrations/apis/resend/actions/send-email" + ] } ] }, @@ -102,7 +110,8 @@ "pages": [ "reference/trigger", "reference/custom-event", - "reference/webhook-event" + "reference/webhook-event", + "reference/schedule-event" ] }, { @@ -128,4 +137,4 @@ "github": "https://github.com/triggerdotdev/trigger.dev", "discord": "https://discord.gg/nkqV9xBYWy" } -} +} \ No newline at end of file diff --git a/apps/docs/reference/custom-event.mdx b/apps/docs/reference/custom-event.mdx index a12441898..a6a4f892b 100644 --- a/apps/docs/reference/custom-event.mdx +++ b/apps/docs/reference/custom-event.mdx @@ -7,6 +7,8 @@ description: "Trigger a workflow when a custom event is received." ## Usage ```ts +import { customEvent, Trigger } from "@trigger.dev/sdk"; + new Trigger({ id: "user-created-notify-slack", name: "User Created - Notify Slack", diff --git a/apps/docs/reference/schedule-event.mdx b/apps/docs/reference/schedule-event.mdx new file mode 100644 index 000000000..9cce0351c --- /dev/null +++ b/apps/docs/reference/schedule-event.mdx @@ -0,0 +1,68 @@ +--- +title: "scheduleEvent Trigger" +sidebarTitle: "scheduleEvent" +description: "Run a workflow on a recurring schedule." +--- + +## Usage + +```ts +import { scheduleEvent } from "@trigger.dev/sdk"; + +new Trigger({ + id: "usage", + name: "usage", + on: scheduleEvent({ rateof: { minutes: 10 } }), + run: async (event, ctx) => {}, +}).listen(); +``` + +## Options + +You must use one of the following options, but not both: + + + The rate of the schedule. This can be a number of minutes, hours, + or days. For example, `{ rateOf: { minutes: 10 } }` will run + every 10 minutes. + + + + A cron expression to run the workflow on. For example, `0 0 * * *` will run + the workflow every hour at the top of the hour. See + [crontab.guru](https://crontab.guru/) for more information. + + +## Event Payload + + + The time the event was scheduled to run. + + + + The time the event last run. This will be `undefined` if the event has never run. Use this parameter to run window queries: + +```ts +import { scheduleEvent, Trigger } from "@trigger.dev/sdk"; + +new Trigger({ + id: "usage", + name: "usage", + on: scheduleEvent({ rateof: { minutes: 10 } }), + run: async (event, ctx) => { + const { lastRunAt, scheduledTime } = event; + + const query = `SELECT * FROM users WHERE created_at < ${scheduledTime}`; + + if (lastRunAt) { + query += ` AND created_at > ${lastRunAt}`; + } + + const latestUsers = await db.query(query); + + // ... + }, +}).listen(); +``` + + diff --git a/apps/docs/reference/trigger.mdx b/apps/docs/reference/trigger.mdx index c8bd6a9b4..d06f51cbf 100644 --- a/apps/docs/reference/trigger.mdx +++ b/apps/docs/reference/trigger.mdx @@ -9,6 +9,8 @@ description: "The Trigger class let's you define a workflow that is triggered by ### Usage ```ts +import { customEvent, Trigger } from "@trigger.dev/sdk"; + const trigger = new Trigger({ id: "user-created-notify-slack", name: "User Created - Notify Slack", @@ -39,11 +41,6 @@ const trigger = new Trigger({ `TRIGGER_API_KEY` environment variable. - - Your Trigger.dev API key. If not provided, the API key will be read from the - `TRIGGER_API_KEY` environment variable. - - The URL of the Trigger.dev WebSocket server. If not provided, the endpoint will point to the production server. diff --git a/apps/docs/reference/webhook-event.mdx b/apps/docs/reference/webhook-event.mdx index 1785cbac6..4fce5d782 100644 --- a/apps/docs/reference/webhook-event.mdx +++ b/apps/docs/reference/webhook-event.mdx @@ -7,6 +7,8 @@ description: "Trigger a workflow when a webhook event is received" ## Usage ```ts +import { webhookEvent, Trigger } from "@trigger.dev/sdk"; + new Trigger({ id: "caldotcom-to-slack", name: "Cal.com To Slack", diff --git a/apps/docs/triggers/scheduled.mdx b/apps/docs/triggers/scheduled.mdx index 5bea86e99..5a88f4b6e 100644 --- a/apps/docs/triggers/scheduled.mdx +++ b/apps/docs/triggers/scheduled.mdx @@ -4,6 +4,8 @@ sidebarTitle: "Scheduled" description: "Run a workflow on a recurring schedule" --- +See the [reference](/reference/schedule-event) for more details. + ## Examples ### Every 5 minutes @@ -51,3 +53,33 @@ new Trigger({ }, }).listen(); ``` + +## Preventing late runs + +To prevent a scheduled trigger from running late, you can set a `triggerTTL` option when creating the `Trigger`, like so: + +```ts +new Trigger({ + id: "scheduled-workflow", + name: "Scheduled Workflow", + apiKey: "", + on: scheduleEvent({ rateOf: { minutes: 5 } }), + triggerTTL: 300, + run: async (event, ctx) => { + await ctx.logger.info("Received the scheduled event", { + event, + wallTime: new Date(), + }); + + return { foo: "bar" }; + }, +}).listen(); +``` + +This will prevent the trigger from running if it is running more than `300` seconds behind, which can happen if the server running your `Trigger` code goes down or is otherwise unavailable. + +This is especially useful for scheduled triggers that run on a very short interval, like every minute, so you don't get a backlog of runs that all run at once when the server comes back online. + + + Set your `triggerTTL` to the same time (or double) as the rateOf the trigger. + diff --git a/apps/webapp/app/services/scheduler/scheduleNextEvent.server.ts b/apps/webapp/app/services/scheduler/scheduleNextEvent.server.ts index 39284533b..7f21d0b26 100644 --- a/apps/webapp/app/services/scheduler/scheduleNextEvent.server.ts +++ b/apps/webapp/app/services/scheduler/scheduleNextEvent.server.ts @@ -29,7 +29,13 @@ export class ScheduleNextEvent { const messageId = await taskQueue.publish( "DELIVER_SCHEDULED_EVENT", - { externalSourceId: schedulerSource.id, payload: { scheduledTime } }, + { + externalSourceId: schedulerSource.id, + payload: { + scheduledTime, + lastRunAt: fromEvent ? fromEvent.createdAt : undefined, + }, + }, {}, { deliverAt: scheduledTime.getTime() } ); diff --git a/apps/webapp/app/triggers/monitoring.server.ts b/apps/webapp/app/triggers/monitoring.server.ts index 01ff77c68..d4ea8637c 100644 --- a/apps/webapp/app/triggers/monitoring.server.ts +++ b/apps/webapp/app/triggers/monitoring.server.ts @@ -5,19 +5,20 @@ import { prisma } from "~/db.server"; export const uptimeCheck = new Trigger({ id: "uptime-check", name: "Uptime Check", - on: scheduleEvent({ rateOf: { minutes: 1 } }), - triggerTTL: 300, + on: scheduleEvent({ rateOf: { minutes: 5 } }), + logLevel: "info", + triggerTTL: 60, run: async (event, context) => { + if (context.environment === "development" && !context.isTest) { + return; + } + // Grab counts of workflows, runs, and steps const userCount = await prisma.user.count(); const workflowCount = await prisma.workflow.count(); const runCount = await prisma.workflowRun.count(); const stepCount = await prisma.workflowRunStep.count(); - if (context.environment === "development") { - return; - } - await slack.postMessage("Uptime Notification", { channelName: "monitoring", text: `[${context.environment}] Uptime Check: ${userCount} users, ${workflowCount} workflows, ${runCount} runs, ${stepCount} steps.`, diff --git a/examples/schedule-to-slack/src/index.ts b/examples/schedule-to-slack/src/index.ts index 4f6eadbe1..6ca2262e2 100644 --- a/examples/schedule-to-slack/src/index.ts +++ b/examples/schedule-to-slack/src/index.ts @@ -2,20 +2,21 @@ import { Trigger, scheduleEvent } from "@trigger.dev/sdk"; import { slack } from "@trigger.dev/integrations"; const trigger = new Trigger({ - id: "schedule-to-slack", + id: "schedule-to-slack-2", name: "Send to Slack every minute", - apiKey: "trigger_development_vzNnO2DGBGcG", + apiKey: "trigger_dev_zC25mKNn6c0q", + endpoint: "ws://localhost:8889/ws", logLevel: "debug", on: scheduleEvent({ rateOf: { minutes: 1 } }), run: async (event, ctx) => { await ctx.logger.info("It's me, the annoying slack bot!"); - const response = await slack.postMessage("slaaaaaack", { - channel: "test-integrations", - text: `Hello, the time is ${event.scheduledTime}`, - }); + // const response = await slack.postMessage("slaaaaaack", { + // channelName: "test-integrations", + // text: `Hello, the time is ${event.scheduledTime}, and I was last run at ${event.lastRunAt}!`, + // }); - return response.message; + return event; }, }); diff --git a/packages/common-schemas/src/events.ts b/packages/common-schemas/src/events.ts index 26805d544..5ece38727 100644 --- a/packages/common-schemas/src/events.ts +++ b/packages/common-schemas/src/events.ts @@ -29,6 +29,7 @@ export const EventFilterSchema: z.ZodType = z.lazy(() => ); export const ScheduledEventPayloadSchema = z.object({ + lastRunAt: z.coerce.date().optional(), scheduledTime: z.coerce.date(), });