pg-workflows

Recurring schedules

Start runs on a cron expression or a fixed interval.

Set schedule to start a run on a recurring basis. It accepts a cron expression, a duration string, or a duration object.

workflow('weekday-report', handler, { schedule: '0 9 * * 1-5', timezone: 'America/New_York' })
workflow('every-5-minutes', handler, { schedule: '5m' })
workflow('hourly', handler, { schedule: '1 hour' })
workflow('daily', handler, { schedule: { days: 1 } })
  • Cron vs. duration. A string of 5 or 6 space-separated fields that uses only cron characters (0-9 * / , - ? L W #) is parsed as cron. Anything else is parsed as a duration.
  • Durations must divide evenly. A duration is converted to cron, so it must be whole minutes that divide 60, whole hours that divide 24, or exactly one day. '23m' and '7h' throw at registration. Use a cron expression for those.
  • timezone applies to cron only. The default is UTC.
  • A worker must be running. Scheduled runs are started by an engine that has the workflow registered and has called engine.start().

A scheduled run has ctx.schedule.timestamp, the time the schedule fired. A run started with startWorkflow has ctx.schedule === undefined. For incremental syncs, use the last completed run as a cursor. Read it inside a step so the value stays fixed when the run resumes or retries:

import { WorkflowStatus, workflow } from 'pg-workflows'

const syncOrders = workflow(
  'sync-orders',
  async ({ step, schedule, workflowId }) => {
    const since = await step.run('read-cursor', async () => {
      const { items } = await engine.getRuns({
        workflowId,
        statuses: [WorkflowStatus.COMPLETED],
        limit: 1,
      })
      return (items[0]?.completedAt ?? new Date(0)).toISOString()
    })

    const orders = await step.run('fetch-orders', async () => {
      const response = await fetch(`https://shop.example.com/orders?updated_since=${since}`)
      return (await response.json()) as { id: string }[]
    })

    return { firedAt: schedule?.timestamp.toISOString(), since, synced: orders.length }
  },
  { schedule: '5m', singleton: true },
)

Don't use getWorkflowLastRun for this. It returns the most recently created run of any status, which inside a scheduled run is the current run.

Overlap. With singleton: true, a scheduled fire is skipped while a previous run is pending or running. Without it, scheduled runs can overlap.