Repository files navigation

@falcondev-oss/workflow

Durable, type-safe queue workers on Redis. Workflows are plain async functions whose steps are memoized in Redis, so a retried job replays completed steps instead of re-running them.

Installation

npm install @falcondev-oss/workflow

Requires Redis (any version with Lua scripting) and Node 24.

Usage

Workflows live in a WorkflowNamespace, which owns the Redis connection, the cross-workflow concurrency cap, and the shared option defaults.

import{createRedis,WorkflowNamespace}from'@falcondev-oss/workflow'import{z}from'zod'constnamespace=newWorkflowNamespace({id: 'my-app',redis: awaitcreateRedis({url: process.env.REDIS_URL}),logger: console,})constworkflow=namespace.createWorkflow({id: 'example-workflow',schema: z.object({timezone: z.string().default('UTC'),name: z.string(),}),asyncrun({ input, step }){awaitstep.do('send welcome',()=>{console.log(`Welcome, ${input.name}! Timezone: ${input.timezone}`)})awaitstep.wait('wait a lil',60_000)constisEngaged=awaitstep.do('check engagement',()=>Math.random()>0.5)if(!isEngaged)return{engagementLevel: 'low'}awaitstep.do('send tips',()=>{console.log(`Here are some tips to get started, ${input.name}!`)})return{engagementLevel: 'high'}},})// Start a worker for this processawaitworkflow.work()// Enqueue a runconstjob=awaitworkflow.run({name: 'John Doe',timezone: 'America/New_York'})// Wait for completion (works from a pure producer too — no worker needed)constresult=awaitjob.wait()console.log(result.engagementLevel)

Watching jobs

Declare a Standard Schema for progress, then emit its input type from any step. Watchers receive the validated output type on the same stream as lifecycle and terminal events.

constexportPdf=namespace.createWorkflow({id: 'export-pdf',schema: z.object({reportId: z.string()}),progressSchema: z.object({label: z.string(),done: z.number(),total: z.number(),}),asyncrun({ input, step }){constrows=awaitstep.do('fetch rows',async({step: nestedStep})=>{awaitnestedStep.progress({label: 'Fetching rows',done: 0,total: 0})returnloadRows(input.reportId)})awaitstep.progress({label: 'Rendering',done: 0,total: rows.length})returnrenderPdf(rows)},})const{ job, events }=awaitexportPdf.runAndWatch({reportId: '1'})forawait(consteventofevents){if(event.type==='progress')console.log(event.data.label)if(event.type==='completed')console.log(event.output)}

Events published before watch() attaches are lost. Lifecycle and progress events are not persisted, and the watcher reads no progress snapshot. Use runAndWatch() when the first event matters because it subscribes before enqueueing. To attach from another request or process, build a handle from the known id:

constjob=awaitexportPdf.getJob(jobId)constevents=awaitjob.watch({signal: request.signal})

watch() subscribes before it resolves and buffers until iteration starts. Breaking out of the loop unsubscribes. A retry emits another started event, while failed is terminal. A completed job still yields its stored terminal result until the configured result TTL expires.

The library does not derive a percentage, ETA, or step count because workflows have no declared step list. If a workflow knows a total, include it in its progress payload as above. Without a progressSchema, step.progress() is a type error and the progress event arm is absent.

Steps

  • step.do(name, fn) — run once, memoize the result. Replayed from Redis on a retry.
  • step.wait(name, ms) — durable sleep; remaining time is computed from the persisted start.
  • step.waitUntil(name, date) — the same, to an absolute time.

Steps nest: the callback receives its own step scoped under the parent's name.

Scheduling

awaitworkflow.run(input)// nowawaitworkflow.runIn(input,60_000)// in 60sawaitworkflow.runAt(input,newDate('2030-01-01'))// at a time// Cron, keyed by (workflow, scheduleId) — upserting the same id replaces in placeawaitworkflow.upsertSchedule('nightly',{pattern: '0 3 * * *',
input,tz: 'Europe/Berlin',})awaitworkflow.getSchedules()awaitworkflow.removeSchedule('nightly')

Ordering, priority and concurrency

Jobs sharing a groupId run one at a time, in enqueue order. Everything else runs in parallel up to the concurrency caps.

namespace.createWorkflow({id: 'per-user',schema: z.object({userId: z.string()}),getGroupId: (input)=>input.userId,// serialize per userqueueOptions: {concurrency: 10,groupConcurrency: 1},workerOptions: {concurrency: 4,maxAttempts: 3},jobOptions: {priority: 1},// 0…2^21-1, higher runs firstrun: async()=>{},})

Options

WorkflowNamespace options are shared defaults — each is shallow-merged under the matching per-workflow override.

OptionDefaultDescription
idNamespace id; scopes the cross-workflow concurrency cap
redisnew clientShared connection, owned by the namespace
prefixwfGlobal key prefix
concurrencyunlimitedCeiling across all workflows in the namespace
loggernoneInherited by every workflow, queue and worker
autoClosetrueClose (drain workers, disconnect) on SIGINT/SIGTERM
queueOptionsDefaults for every workflow's queue
workerOptionsDefaults for every workflow's worker
jobOptionsDefaults for every enqueued job

Shutdown is handled at the namespace: workers drain in-flight jobs, then the connections close. Set autoClose: false and call await namespace.close() yourself to own it.

Failures and retries

A throwing handler is retried up to maxAttempts with an exponential backoff (expBackoff()), then dead-lettered. Throw NonRecoverableError to dead-letter immediately, skipping the remaining budget — the library does this itself when a job's stored payload no longer validates against the workflow schema (a job enqueued before a schema change), since a retry would only re-read the same payload.

Metrics

workflow.getMetrics() returns point-in-time { active, waiting, delayed } depths. Pass an OpenTelemetry meter to export them as gauges:

constnamespace=newWorkflowNamespace({id: 'my-app',workerOptions: {metrics: { meter,prefix: 'my_app'}},})

Spans are emitted for producers, workers and each step via the global OpenTelemetry tracer.

Inspiration

About

Durable, type-safe queue workers on Redis with memoized steps, cron schedules, group ordering, OpenTelemetry tracing

Topics

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Repository files navigation

@falcondev-oss/workflow

Durable, type-safe queue workers on Redis. Workflows are plain async functions whose steps are memoized in Redis, so a retried job replays completed steps instead of re-running them.

Installation

npm install @falcondev-oss/workflow

Requires Redis (any version with Lua scripting) and Node 24.

Usage

Workflows live in a WorkflowNamespace, which owns the Redis connection, the cross-workflow concurrency cap, and the shared option defaults.

import{createRedis,WorkflowNamespace}from'@falcondev-oss/workflow'import{z}from'zod'constnamespace=newWorkflowNamespace({id: 'my-app',redis: awaitcreateRedis({url: process.env.REDIS_URL}),logger: console,})constworkflow=namespace.createWorkflow({id: 'example-workflow',schema: z.object({timezone: z.string().default('UTC'),name: z.string(),}),asyncrun({ input, step }){awaitstep.do('send welcome',()=>{console.log(`Welcome, ${input.name}! Timezone: ${input.timezone}`)})awaitstep.wait('wait a lil',60_000)constisEngaged=awaitstep.do('check engagement',()=>Math.random()>0.5)if(!isEngaged)return{engagementLevel: 'low'}awaitstep.do('send tips',()=>{console.log(`Here are some tips to get started, ${input.name}!`)})return{engagementLevel: 'high'}},})// Start a worker for this processawaitworkflow.work()// Enqueue a runconstjob=awaitworkflow.run({name: 'John Doe',timezone: 'America/New_York'})// Wait for completion (works from a pure producer too — no worker needed)constresult=awaitjob.wait()console.log(result.engagementLevel)

Watching jobs

Declare a Standard Schema for progress, then emit its input type from any step. Watchers receive the validated output type on the same stream as lifecycle and terminal events.

constexportPdf=namespace.createWorkflow({id: 'export-pdf',schema: z.object({reportId: z.string()}),progressSchema: z.object({label: z.string(),done: z.number(),total: z.number(),}),asyncrun({ input, step }){constrows=awaitstep.do('fetch rows',async({step: nestedStep})=>{awaitnestedStep.progress({label: 'Fetching rows',done: 0,total: 0})returnloadRows(input.reportId)})awaitstep.progress({label: 'Rendering',done: 0,total: rows.length})returnrenderPdf(rows)},})const{ job, events }=awaitexportPdf.runAndWatch({reportId: '1'})forawait(consteventofevents){if(event.type==='progress')console.log(event.data.label)if(event.type==='completed')console.log(event.output)}

Events published before watch() attaches are lost. Lifecycle and progress events are not persisted, and the watcher reads no progress snapshot. Use runAndWatch() when the first event matters because it subscribes before enqueueing. To attach from another request or process, build a handle from the known id:

constjob=awaitexportPdf.getJob(jobId)constevents=awaitjob.watch({signal: request.signal})

watch() subscribes before it resolves and buffers until iteration starts. Breaking out of the loop unsubscribes. A retry emits another started event, while failed is terminal. A completed job still yields its stored terminal result until the configured result TTL expires.

The library does not derive a percentage, ETA, or step count because workflows have no declared step list. If a workflow knows a total, include it in its progress payload as above. Without a progressSchema, step.progress() is a type error and the progress event arm is absent.

Steps

  • step.do(name, fn) — run once, memoize the result. Replayed from Redis on a retry.
  • step.wait(name, ms) — durable sleep; remaining time is computed from the persisted start.
  • step.waitUntil(name, date) — the same, to an absolute time.

Steps nest: the callback receives its own step scoped under the parent's name.

Scheduling

awaitworkflow.run(input)// nowawaitworkflow.runIn(input,60_000)// in 60sawaitworkflow.runAt(input,newDate('2030-01-01'))// at a time// Cron, keyed by (workflow, scheduleId) — upserting the same id replaces in placeawaitworkflow.upsertSchedule('nightly',{pattern: '0 3 * * *',
input,tz: 'Europe/Berlin',})awaitworkflow.getSchedules()awaitworkflow.removeSchedule('nightly')

Ordering, priority and concurrency

Jobs sharing a groupId run one at a time, in enqueue order. Everything else runs in parallel up to the concurrency caps.

namespace.createWorkflow({id: 'per-user',schema: z.object({userId: z.string()}),getGroupId: (input)=>input.userId,// serialize per userqueueOptions: {concurrency: 10,groupConcurrency: 1},workerOptions: {concurrency: 4,maxAttempts: 3},jobOptions: {priority: 1},// 0…2^21-1, higher runs firstrun: async()=>{},})

Options

WorkflowNamespace options are shared defaults — each is shallow-merged under the matching per-workflow override.

OptionDefaultDescription
idNamespace id; scopes the cross-workflow concurrency cap
redisnew clientShared connection, owned by the namespace
prefixwfGlobal key prefix
concurrencyunlimitedCeiling across all workflows in the namespace
loggernoneInherited by every workflow, queue and worker
autoClosetrueClose (drain workers, disconnect) on SIGINT/SIGTERM
queueOptionsDefaults for every workflow's queue
workerOptionsDefaults for every workflow's worker
jobOptionsDefaults for every enqueued job

Shutdown is handled at the namespace: workers drain in-flight jobs, then the connections close. Set autoClose: false and call await namespace.close() yourself to own it.

Failures and retries

A throwing handler is retried up to maxAttempts with an exponential backoff (expBackoff()), then dead-lettered. Throw NonRecoverableError to dead-letter immediately, skipping the remaining budget — the library does this itself when a job's stored payload no longer validates against the workflow schema (a job enqueued before a schema change), since a retry would only re-read the same payload.

Metrics

workflow.getMetrics() returns point-in-time { active, waiting, delayed } depths. Pass an OpenTelemetry meter to export them as gauges:

constnamespace=newWorkflowNamespace({id: 'my-app',workerOptions: {metrics: { meter,prefix: 'my_app'}},})

Spans are emitted for producers, workers and each step via the global OpenTelemetry tracer.

Inspiration

About

Durable, type-safe queue workers on Redis with memoized steps, cron schedules, group ordering, OpenTelemetry tracing

Topics

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

@falcondev-oss/workflow

Durable, type-safe queue workers on Redis. Workflows are plain async functions whose steps are memoized in Redis, so a retried job replays completed steps instead of re-running them.

Installation

npm install @falcondev-oss/workflow

Requires Redis (any version with Lua scripting) and Node 24.

Usage

Workflows live in a WorkflowNamespace, which owns the Redis connection, the cross-workflow concurrency cap, and the shared option defaults.

import{createRedis,WorkflowNamespace}from'@falcondev-oss/workflow'import{z}from'zod'constnamespace=newWorkflowNamespace({id: 'my-app',redis: awaitcreateRedis({url: process.env.REDIS_URL}),logger: console,})constworkflow=namespace.createWorkflow({id: 'example-workflow',schema: z.object({timezone: z.string().default('UTC'),name: z.string(),}),asyncrun({ input, step }){awaitstep.do('send welcome',()=>{console.log(`Welcome, ${input.name}! Timezone: ${input.timezone}`)})awaitstep.wait('wait a lil',60_000)constisEngaged=awaitstep.do('check engagement',()=>Math.random()>0.5)if(!isEngaged)return{engagementLevel: 'low'}awaitstep.do('send tips',()=>{console.log(`Here are some tips to get started, ${input.name}!`)})return{engagementLevel: 'high'}},})// Start a worker for this processawaitworkflow.work()// Enqueue a runconstjob=awaitworkflow.run({name: 'John Doe',timezone: 'America/New_York'})// Wait for completion (works from a pure producer too — no worker needed)constresult=awaitjob.wait()console.log(result.engagementLevel)

Watching jobs

Declare a Standard Schema for progress, then emit its input type from any step. Watchers receive the validated output type on the same stream as lifecycle and terminal events.

constexportPdf=namespace.createWorkflow({id: 'export-pdf',schema: z.object({reportId: z.string()}),progressSchema: z.object({label: z.string(),done: z.number(),total: z.number(),}),asyncrun({ input, step }){constrows=awaitstep.do('fetch rows',async({step: nestedStep})=>{awaitnestedStep.progress({label: 'Fetching rows',done: 0,total: 0})returnloadRows(input.reportId)})awaitstep.progress({label: 'Rendering',done: 0,total: rows.length})returnrenderPdf(rows)},})const{ job, events }=awaitexportPdf.runAndWatch({reportId: '1'})forawait(consteventofevents){if(event.type==='progress')console.log(event.data.label)if(event.type==='completed')console.log(event.output)}

Events published before watch() attaches are lost. Lifecycle and progress events are not persisted, and the watcher reads no progress snapshot. Use runAndWatch() when the first event matters because it subscribes before enqueueing. To attach from another request or process, build a handle from the known id:

constjob=awaitexportPdf.getJob(jobId)constevents=awaitjob.watch({signal: request.signal})

watch() subscribes before it resolves and buffers until iteration starts. Breaking out of the loop unsubscribes. A retry emits another started event, while failed is terminal. A completed job still yields its stored terminal result until the configured result TTL expires.

The library does not derive a percentage, ETA, or step count because workflows have no declared step list. If a workflow knows a total, include it in its progress payload as above. Without a progressSchema, step.progress() is a type error and the progress event arm is absent.

Steps

  • step.do(name, fn) — run once, memoize the result. Replayed from Redis on a retry.
  • step.wait(name, ms) — durable sleep; remaining time is computed from the persisted start.
  • step.waitUntil(name, date) — the same, to an absolute time.

Steps nest: the callback receives its own step scoped under the parent's name.

Scheduling

awaitworkflow.run(input)// nowawaitworkflow.runIn(input,60_000)// in 60sawaitworkflow.runAt(input,newDate('2030-01-01'))// at a time// Cron, keyed by (workflow, scheduleId) — upserting the same id replaces in placeawaitworkflow.upsertSchedule('nightly',{pattern: '0 3 * * *',
input,tz: 'Europe/Berlin',})awaitworkflow.getSchedules()awaitworkflow.removeSchedule('nightly')

Ordering, priority and concurrency

Jobs sharing a groupId run one at a time, in enqueue order. Everything else runs in parallel up to the concurrency caps.

namespace.createWorkflow({id: 'per-user',schema: z.object({userId: z.string()}),getGroupId: (input)=>input.userId,// serialize per userqueueOptions: {concurrency: 10,groupConcurrency: 1},workerOptions: {concurrency: 4,maxAttempts: 3},jobOptions: {priority: 1},// 0…2^21-1, higher runs firstrun: async()=>{},})

Options

WorkflowNamespace options are shared defaults — each is shallow-merged under the matching per-workflow override.

OptionDefaultDescription
idNamespace id; scopes the cross-workflow concurrency cap
redisnew clientShared connection, owned by the namespace
prefixwfGlobal key prefix
concurrencyunlimitedCeiling across all workflows in the namespace
loggernoneInherited by every workflow, queue and worker
autoClosetrueClose (drain workers, disconnect) on SIGINT/SIGTERM
queueOptionsDefaults for every workflow's queue
workerOptionsDefaults for every workflow's worker
jobOptionsDefaults for every enqueued job

Shutdown is handled at the namespace: workers drain in-flight jobs, then the connections close. Set autoClose: false and call await namespace.close() yourself to own it.

Failures and retries

A throwing handler is retried up to maxAttempts with an exponential backoff (expBackoff()), then dead-lettered. Throw NonRecoverableError to dead-letter immediately, skipping the remaining budget — the library does this itself when a job's stored payload no longer validates against the workflow schema (a job enqueued before a schema change), since a retry would only re-read the same payload.

Metrics

workflow.getMetrics() returns point-in-time { active, waiting, delayed } depths. Pass an OpenTelemetry meter to export them as gauges:

constnamespace=newWorkflowNamespace({id: 'my-app',workerOptions: {metrics: { meter,prefix: 'my_app'}},})

Spans are emitted for producers, workers and each step via the global OpenTelemetry tracer.

Inspiration

About

Durable, type-safe queue workers on Redis with memoized steps, cron schedules, group ordering, OpenTelemetry tracing

Topics

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

@falcondev-oss/workflow

Durable, type-safe queue workers on Redis. Workflows are plain async functions whose steps are memoized in Redis, so a retried job replays completed steps instead of re-running them.

Installation

npm install @falcondev-oss/workflow

Requires Redis (any version with Lua scripting) and Node 24.

Usage

Workflows live in a WorkflowNamespace, which owns the Redis connection, the cross-workflow concurrency cap, and the shared option defaults.

import{createRedis,WorkflowNamespace}from'@falcondev-oss/workflow'import{z}from'zod'constnamespace=newWorkflowNamespace({id: 'my-app',redis: awaitcreateRedis({url: process.env.REDIS_URL}),logger: console,})constworkflow=namespace.createWorkflow({id: 'example-workflow',schema: z.object({timezone: z.string().default('UTC'),name: z.string(),}),asyncrun({ input, step }){awaitstep.do('send welcome',()=>{console.log(`Welcome, ${input.name}! Timezone: ${input.timezone}`)})awaitstep.wait('wait a lil',60_000)constisEngaged=awaitstep.do('check engagement',()=>Math.random()>0.5)if(!isEngaged)return{engagementLevel: 'low'}awaitstep.do('send tips',()=>{console.log(`Here are some tips to get started, ${input.name}!`)})return{engagementLevel: 'high'}},})// Start a worker for this processawaitworkflow.work()// Enqueue a runconstjob=awaitworkflow.run({name: 'John Doe',timezone: 'America/New_York'})// Wait for completion (works from a pure producer too — no worker needed)constresult=awaitjob.wait()console.log(result.engagementLevel)

Watching jobs

Declare a Standard Schema for progress, then emit its input type from any step. Watchers receive the validated output type on the same stream as lifecycle and terminal events.

constexportPdf=namespace.createWorkflow({id: 'export-pdf',schema: z.object({reportId: z.string()}),progressSchema: z.object({label: z.string(),done: z.number(),total: z.number(),}),asyncrun({ input, step }){constrows=awaitstep.do('fetch rows',async({step: nestedStep})=>{awaitnestedStep.progress({label: 'Fetching rows',done: 0,total: 0})returnloadRows(input.reportId)})awaitstep.progress({label: 'Rendering',done: 0,total: rows.length})returnrenderPdf(rows)},})const{ job, events }=awaitexportPdf.runAndWatch({reportId: '1'})forawait(consteventofevents){if(event.type==='progress')console.log(event.data.label)if(event.type==='completed')console.log(event.output)}

Events published before watch() attaches are lost. Lifecycle and progress events are not persisted, and the watcher reads no progress snapshot. Use runAndWatch() when the first event matters because it subscribes before enqueueing. To attach from another request or process, build a handle from the known id:

constjob=awaitexportPdf.getJob(jobId)constevents=awaitjob.watch({signal: request.signal})

watch() subscribes before it resolves and buffers until iteration starts. Breaking out of the loop unsubscribes. A retry emits another started event, while failed is terminal. A completed job still yields its stored terminal result until the configured result TTL expires.

The library does not derive a percentage, ETA, or step count because workflows have no declared step list. If a workflow knows a total, include it in its progress payload as above. Without a progressSchema, step.progress() is a type error and the progress event arm is absent.

Steps

  • step.do(name, fn) — run once, memoize the result. Replayed from Redis on a retry.
  • step.wait(name, ms) — durable sleep; remaining time is computed from the persisted start.
  • step.waitUntil(name, date) — the same, to an absolute time.

Steps nest: the callback receives its own step scoped under the parent's name.

Scheduling

awaitworkflow.run(input)// nowawaitworkflow.runIn(input,60_000)// in 60sawaitworkflow.runAt(input,newDate('2030-01-01'))// at a time// Cron, keyed by (workflow, scheduleId) — upserting the same id replaces in placeawaitworkflow.upsertSchedule('nightly',{pattern: '0 3 * * *',
input,tz: 'Europe/Berlin',})awaitworkflow.getSchedules()awaitworkflow.removeSchedule('nightly')

Ordering, priority and concurrency

Jobs sharing a groupId run one at a time, in enqueue order. Everything else runs in parallel up to the concurrency caps.

namespace.createWorkflow({id: 'per-user',schema: z.object({userId: z.string()}),getGroupId: (input)=>input.userId,// serialize per userqueueOptions: {concurrency: 10,groupConcurrency: 1},workerOptions: {concurrency: 4,maxAttempts: 3},jobOptions: {priority: 1},// 0…2^21-1, higher runs firstrun: async()=>{},})

Options

WorkflowNamespace options are shared defaults — each is shallow-merged under the matching per-workflow override.

OptionDefaultDescription
idNamespace id; scopes the cross-workflow concurrency cap
redisnew clientShared connection, owned by the namespace
prefixwfGlobal key prefix
concurrencyunlimitedCeiling across all workflows in the namespace
loggernoneInherited by every workflow, queue and worker
autoClosetrueClose (drain workers, disconnect) on SIGINT/SIGTERM
queueOptionsDefaults for every workflow's queue
workerOptionsDefaults for every workflow's worker
jobOptionsDefaults for every enqueued job

Shutdown is handled at the namespace: workers drain in-flight jobs, then the connections close. Set autoClose: false and call await namespace.close() yourself to own it.

Failures and retries

A throwing handler is retried up to maxAttempts with an exponential backoff (expBackoff()), then dead-lettered. Throw NonRecoverableError to dead-letter immediately, skipping the remaining budget — the library does this itself when a job's stored payload no longer validates against the workflow schema (a job enqueued before a schema change), since a retry would only re-read the same payload.

Metrics

workflow.getMetrics() returns point-in-time { active, waiting, delayed } depths. Pass an OpenTelemetry meter to export them as gauges:

constnamespace=newWorkflowNamespace({id: 'my-app',workerOptions: {metrics: { meter,prefix: 'my_app'}},})

Spans are emitted for producers, workers and each step via the global OpenTelemetry tracer.

Inspiration

About

Durable, type-safe queue workers on Redis with memoized steps, cron schedules, group ordering, OpenTelemetry tracing

Topics

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Repository files navigation

@falcondev-oss/workflow

Durable, type-safe queue workers on Redis. Workflows are plain async functions whose steps are memoized in Redis, so a retried job replays completed steps instead of re-running them.

Installation

npm install @falcondev-oss/workflow

Requires Redis (any version with Lua scripting) and Node 24.

Usage

Workflows live in a WorkflowNamespace, which owns the Redis connection, the cross-workflow concurrency cap, and the shared option defaults.

import{createRedis,WorkflowNamespace}from'@falcondev-oss/workflow'import{z}from'zod'constnamespace=newWorkflowNamespace({id: 'my-app',redis: awaitcreateRedis({url: process.env.REDIS_URL}),logger: console,})constworkflow=namespace.createWorkflow({id: 'example-workflow',schema: z.object({timezone: z.string().default('UTC'),name: z.string(),}),asyncrun({ input, step }){awaitstep.do('send welcome',()=>{console.log(`Welcome, ${input.name}! Timezone: ${input.timezone}`)})awaitstep.wait('wait a lil',60_000)constisEngaged=awaitstep.do('check engagement',()=>Math.random()>0.5)if(!isEngaged)return{engagementLevel: 'low'}awaitstep.do('send tips',()=>{console.log(`Here are some tips to get started, ${input.name}!`)})return{engagementLevel: 'high'}},})// Start a worker for this processawaitworkflow.work()// Enqueue a runconstjob=awaitworkflow.run({name: 'John Doe',timezone: 'America/New_York'})// Wait for completion (works from a pure producer too — no worker needed)constresult=awaitjob.wait()console.log(result.engagementLevel)

Watching jobs

Declare a Standard Schema for progress, then emit its input type from any step. Watchers receive the validated output type on the same stream as lifecycle and terminal events.

constexportPdf=namespace.createWorkflow({id: 'export-pdf',schema: z.object({reportId: z.string()}),progressSchema: z.object({label: z.string(),done: z.number(),total: z.number(),}),asyncrun({ input, step }){constrows=awaitstep.do('fetch rows',async({step: nestedStep})=>{awaitnestedStep.progress({label: 'Fetching rows',done: 0,total: 0})returnloadRows(input.reportId)})awaitstep.progress({label: 'Rendering',done: 0,total: rows.length})returnrenderPdf(rows)},})const{ job, events }=awaitexportPdf.runAndWatch({reportId: '1'})forawait(consteventofevents){if(event.type==='progress')console.log(event.data.label)if(event.type==='completed')console.log(event.output)}

Events published before watch() attaches are lost. Lifecycle and progress events are not persisted, and the watcher reads no progress snapshot. Use runAndWatch() when the first event matters because it subscribes before enqueueing. To attach from another request or process, build a handle from the known id:

constjob=awaitexportPdf.getJob(jobId)constevents=awaitjob.watch({signal: request.signal})

watch() subscribes before it resolves and buffers until iteration starts. Breaking out of the loop unsubscribes. A retry emits another started event, while failed is terminal. A completed job still yields its stored terminal result until the configured result TTL expires.

The library does not derive a percentage, ETA, or step count because workflows have no declared step list. If a workflow knows a total, include it in its progress payload as above. Without a progressSchema, step.progress() is a type error and the progress event arm is absent.

Steps

  • step.do(name, fn) — run once, memoize the result. Replayed from Redis on a retry.
  • step.wait(name, ms) — durable sleep; remaining time is computed from the persisted start.
  • step.waitUntil(name, date) — the same, to an absolute time.

Steps nest: the callback receives its own step scoped under the parent's name.

Scheduling

awaitworkflow.run(input)// nowawaitworkflow.runIn(input,60_000)// in 60sawaitworkflow.runAt(input,newDate('2030-01-01'))// at a time// Cron, keyed by (workflow, scheduleId) — upserting the same id replaces in placeawaitworkflow.upsertSchedule('nightly',{pattern: '0 3 * * *',
input,tz: 'Europe/Berlin',})awaitworkflow.getSchedules()awaitworkflow.removeSchedule('nightly')

Ordering, priority and concurrency

Jobs sharing a groupId run one at a time, in enqueue order. Everything else runs in parallel up to the concurrency caps.

namespace.createWorkflow({id: 'per-user',schema: z.object({userId: z.string()}),getGroupId: (input)=>input.userId,// serialize per userqueueOptions: {concurrency: 10,groupConcurrency: 1},workerOptions: {concurrency: 4,maxAttempts: 3},jobOptions: {priority: 1},// 0…2^21-1, higher runs firstrun: async()=>{},})

Options

WorkflowNamespace options are shared defaults — each is shallow-merged under the matching per-workflow override.

OptionDefaultDescription
idNamespace id; scopes the cross-workflow concurrency cap
redisnew clientShared connection, owned by the namespace
prefixwfGlobal key prefix
concurrencyunlimitedCeiling across all workflows in the namespace
loggernoneInherited by every workflow, queue and worker
autoClosetrueClose (drain workers, disconnect) on SIGINT/SIGTERM
queueOptionsDefaults for every workflow's queue
workerOptionsDefaults for every workflow's worker
jobOptionsDefaults for every enqueued job

Shutdown is handled at the namespace: workers drain in-flight jobs, then the connections close. Set autoClose: false and call await namespace.close() yourself to own it.

Failures and retries

A throwing handler is retried up to maxAttempts with an exponential backoff (expBackoff()), then dead-lettered. Throw NonRecoverableError to dead-letter immediately, skipping the remaining budget — the library does this itself when a job's stored payload no longer validates against the workflow schema (a job enqueued before a schema change), since a retry would only re-read the same payload.

Metrics

workflow.getMetrics() returns point-in-time { active, waiting, delayed } depths. Pass an OpenTelemetry meter to export them as gauges:

constnamespace=newWorkflowNamespace({id: 'my-app',workerOptions: {metrics: { meter,prefix: 'my_app'}},})

Spans are emitted for producers, workers and each step via the global OpenTelemetry tracer.

Inspiration

About

Durable, type-safe queue workers on Redis with memoized steps, cron schedules, group ordering, OpenTelemetry tracing

Topics

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

@falcondev-oss/workflow

Durable, type-safe queue workers on Redis. Workflows are plain async functions whose steps are memoized in Redis, so a retried job replays completed steps instead of re-running them.

Installation

npm install @falcondev-oss/workflow

Requires Redis (any version with Lua scripting) and Node 24.

Usage

Workflows live in a WorkflowNamespace, which owns the Redis connection, the cross-workflow concurrency cap, and the shared option defaults.

import{createRedis,WorkflowNamespace}from'@falcondev-oss/workflow'import{z}from'zod'constnamespace=newWorkflowNamespace({id: 'my-app',redis: awaitcreateRedis({url: process.env.REDIS_URL}),logger: console,})constworkflow=namespace.createWorkflow({id: 'example-workflow',schema: z.object({timezone: z.string().default('UTC'),name: z.string(),}),asyncrun({ input, step }){awaitstep.do('send welcome',()=>{console.log(`Welcome, ${input.name}! Timezone: ${input.timezone}`)})awaitstep.wait('wait a lil',60_000)constisEngaged=awaitstep.do('check engagement',()=>Math.random()>0.5)if(!isEngaged)return{engagementLevel: 'low'}awaitstep.do('send tips',()=>{console.log(`Here are some tips to get started, ${input.name}!`)})return{engagementLevel: 'high'}},})// Start a worker for this processawaitworkflow.work()// Enqueue a runconstjob=awaitworkflow.run({name: 'John Doe',timezone: 'America/New_York'})// Wait for completion (works from a pure producer too — no worker needed)constresult=awaitjob.wait()console.log(result.engagementLevel)

Watching jobs

Declare a Standard Schema for progress, then emit its input type from any step. Watchers receive the validated output type on the same stream as lifecycle and terminal events.

constexportPdf=namespace.createWorkflow({id: 'export-pdf',schema: z.object({reportId: z.string()}),progressSchema: z.object({label: z.string(),done: z.number(),total: z.number(),}),asyncrun({ input, step }){constrows=awaitstep.do('fetch rows',async({step: nestedStep})=>{awaitnestedStep.progress({label: 'Fetching rows',done: 0,total: 0})returnloadRows(input.reportId)})awaitstep.progress({label: 'Rendering',done: 0,total: rows.length})returnrenderPdf(rows)},})const{ job, events }=awaitexportPdf.runAndWatch({reportId: '1'})forawait(consteventofevents){if(event.type==='progress')console.log(event.data.label)if(event.type==='completed')console.log(event.output)}

Events published before watch() attaches are lost. Lifecycle and progress events are not persisted, and the watcher reads no progress snapshot. Use runAndWatch() when the first event matters because it subscribes before enqueueing. To attach from another request or process, build a handle from the known id:

constjob=awaitexportPdf.getJob(jobId)constevents=awaitjob.watch({signal: request.signal})

watch() subscribes before it resolves and buffers until iteration starts. Breaking out of the loop unsubscribes. A retry emits another started event, while failed is terminal. A completed job still yields its stored terminal result until the configured result TTL expires.

The library does not derive a percentage, ETA, or step count because workflows have no declared step list. If a workflow knows a total, include it in its progress payload as above. Without a progressSchema, step.progress() is a type error and the progress event arm is absent.

Steps

  • step.do(name, fn) — run once, memoize the result. Replayed from Redis on a retry.
  • step.wait(name, ms) — durable sleep; remaining time is computed from the persisted start.
  • step.waitUntil(name, date) — the same, to an absolute time.

Steps nest: the callback receives its own step scoped under the parent's name.

Scheduling

awaitworkflow.run(input)// nowawaitworkflow.runIn(input,60_000)// in 60sawaitworkflow.runAt(input,newDate('2030-01-01'))// at a time// Cron, keyed by (workflow, scheduleId) — upserting the same id replaces in placeawaitworkflow.upsertSchedule('nightly',{pattern: '0 3 * * *',
input,tz: 'Europe/Berlin',})awaitworkflow.getSchedules()awaitworkflow.removeSchedule('nightly')

Ordering, priority and concurrency

Jobs sharing a groupId run one at a time, in enqueue order. Everything else runs in parallel up to the concurrency caps.

namespace.createWorkflow({id: 'per-user',schema: z.object({userId: z.string()}),getGroupId: (input)=>input.userId,// serialize per userqueueOptions: {concurrency: 10,groupConcurrency: 1},workerOptions: {concurrency: 4,maxAttempts: 3},jobOptions: {priority: 1},// 0…2^21-1, higher runs firstrun: async()=>{},})

Options

WorkflowNamespace options are shared defaults — each is shallow-merged under the matching per-workflow override.

OptionDefaultDescription
idNamespace id; scopes the cross-workflow concurrency cap
redisnew clientShared connection, owned by the namespace
prefixwfGlobal key prefix
concurrencyunlimitedCeiling across all workflows in the namespace
loggernoneInherited by every workflow, queue and worker
autoClosetrueClose (drain workers, disconnect) on SIGINT/SIGTERM
queueOptionsDefaults for every workflow's queue
workerOptionsDefaults for every workflow's worker
jobOptionsDefaults for every enqueued job

Shutdown is handled at the namespace: workers drain in-flight jobs, then the connections close. Set autoClose: false and call await namespace.close() yourself to own it.

Failures and retries

A throwing handler is retried up to maxAttempts with an exponential backoff (expBackoff()), then dead-lettered. Throw NonRecoverableError to dead-letter immediately, skipping the remaining budget — the library does this itself when a job's stored payload no longer validates against the workflow schema (a job enqueued before a schema change), since a retry would only re-read the same payload.

Metrics

workflow.getMetrics() returns point-in-time { active, waiting, delayed } depths. Pass an OpenTelemetry meter to export them as gauges:

constnamespace=newWorkflowNamespace({id: 'my-app',workerOptions: {metrics: { meter,prefix: 'my_app'}},})

Spans are emitted for producers, workers and each step via the global OpenTelemetry tracer.

Inspiration

About

Durable, type-safe queue workers on Redis with memoized steps, cron schedules, group ordering, OpenTelemetry tracing

Topics

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

@falcondev-oss/workflow

Durable, type-safe queue workers on Redis. Workflows are plain async functions whose steps are memoized in Redis, so a retried job replays completed steps instead of re-running them.

Installation

npm install @falcondev-oss/workflow

Requires Redis (any version with Lua scripting) and Node 24.

Usage

Workflows live in a WorkflowNamespace, which owns the Redis connection, the cross-workflow concurrency cap, and the shared option defaults.

import{createRedis,WorkflowNamespace}from'@falcondev-oss/workflow'import{z}from'zod'constnamespace=newWorkflowNamespace({id: 'my-app',redis: awaitcreateRedis({url: process.env.REDIS_URL}),logger: console,})constworkflow=namespace.createWorkflow({id: 'example-workflow',schema: z.object({timezone: z.string().default('UTC'),name: z.string(),}),asyncrun({ input, step }){awaitstep.do('send welcome',()=>{console.log(`Welcome, ${input.name}! Timezone: ${input.timezone}`)})awaitstep.wait('wait a lil',60_000)constisEngaged=awaitstep.do('check engagement',()=>Math.random()>0.5)if(!isEngaged)return{engagementLevel: 'low'}awaitstep.do('send tips',()=>{console.log(`Here are some tips to get started, ${input.name}!`)})return{engagementLevel: 'high'}},})// Start a worker for this processawaitworkflow.work()// Enqueue a runconstjob=awaitworkflow.run({name: 'John Doe',timezone: 'America/New_York'})// Wait for completion (works from a pure producer too — no worker needed)constresult=awaitjob.wait()console.log(result.engagementLevel)

Watching jobs

Declare a Standard Schema for progress, then emit its input type from any step. Watchers receive the validated output type on the same stream as lifecycle and terminal events.

constexportPdf=namespace.createWorkflow({id: 'export-pdf',schema: z.object({reportId: z.string()}),progressSchema: z.object({label: z.string(),done: z.number(),total: z.number(),}),asyncrun({ input, step }){constrows=awaitstep.do('fetch rows',async({step: nestedStep})=>{awaitnestedStep.progress({label: 'Fetching rows',done: 0,total: 0})returnloadRows(input.reportId)})awaitstep.progress({label: 'Rendering',done: 0,total: rows.length})returnrenderPdf(rows)},})const{ job, events }=awaitexportPdf.runAndWatch({reportId: '1'})forawait(consteventofevents){if(event.type==='progress')console.log(event.data.label)if(event.type==='completed')console.log(event.output)}

Events published before watch() attaches are lost. Lifecycle and progress events are not persisted, and the watcher reads no progress snapshot. Use runAndWatch() when the first event matters because it subscribes before enqueueing. To attach from another request or process, build a handle from the known id:

constjob=awaitexportPdf.getJob(jobId)constevents=awaitjob.watch({signal: request.signal})

watch() subscribes before it resolves and buffers until iteration starts. Breaking out of the loop unsubscribes. A retry emits another started event, while failed is terminal. A completed job still yields its stored terminal result until the configured result TTL expires.

The library does not derive a percentage, ETA, or step count because workflows have no declared step list. If a workflow knows a total, include it in its progress payload as above. Without a progressSchema, step.progress() is a type error and the progress event arm is absent.

Steps

  • step.do(name, fn) — run once, memoize the result. Replayed from Redis on a retry.
  • step.wait(name, ms) — durable sleep; remaining time is computed from the persisted start.
  • step.waitUntil(name, date) — the same, to an absolute time.

Steps nest: the callback receives its own step scoped under the parent's name.

Scheduling

awaitworkflow.run(input)// nowawaitworkflow.runIn(input,60_000)// in 60sawaitworkflow.runAt(input,newDate('2030-01-01'))// at a time// Cron, keyed by (workflow, scheduleId) — upserting the same id replaces in placeawaitworkflow.upsertSchedule('nightly',{pattern: '0 3 * * *',
input,tz: 'Europe/Berlin',})awaitworkflow.getSchedules()awaitworkflow.removeSchedule('nightly')

Ordering, priority and concurrency

Jobs sharing a groupId run one at a time, in enqueue order. Everything else runs in parallel up to the concurrency caps.

namespace.createWorkflow({id: 'per-user',schema: z.object({userId: z.string()}),getGroupId: (input)=>input.userId,// serialize per userqueueOptions: {concurrency: 10,groupConcurrency: 1},workerOptions: {concurrency: 4,maxAttempts: 3},jobOptions: {priority: 1},// 0…2^21-1, higher runs firstrun: async()=>{},})

Options

WorkflowNamespace options are shared defaults — each is shallow-merged under the matching per-workflow override.

OptionDefaultDescription
idNamespace id; scopes the cross-workflow concurrency cap
redisnew clientShared connection, owned by the namespace
prefixwfGlobal key prefix
concurrencyunlimitedCeiling across all workflows in the namespace
loggernoneInherited by every workflow, queue and worker
autoClosetrueClose (drain workers, disconnect) on SIGINT/SIGTERM
queueOptionsDefaults for every workflow's queue
workerOptionsDefaults for every workflow's worker
jobOptionsDefaults for every enqueued job

Shutdown is handled at the namespace: workers drain in-flight jobs, then the connections close. Set autoClose: false and call await namespace.close() yourself to own it.

Failures and retries

A throwing handler is retried up to maxAttempts with an exponential backoff (expBackoff()), then dead-lettered. Throw NonRecoverableError to dead-letter immediately, skipping the remaining budget — the library does this itself when a job's stored payload no longer validates against the workflow schema (a job enqueued before a schema change), since a retry would only re-read the same payload.

Metrics

workflow.getMetrics() returns point-in-time { active, waiting, delayed } depths. Pass an OpenTelemetry meter to export them as gauges:

constnamespace=newWorkflowNamespace({id: 'my-app',workerOptions: {metrics: { meter,prefix: 'my_app'}},})

Spans are emitted for producers, workers and each step via the global OpenTelemetry tracer.

Inspiration

About

Durable, type-safe queue workers on Redis with memoized steps, cron schedules, group ordering, OpenTelemetry tracing

Topics

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Repository files navigation

@falcondev-oss/workflow

Durable, type-safe queue workers on Redis. Workflows are plain async functions whose steps are memoized in Redis, so a retried job replays completed steps instead of re-running them.

Installation

npm install @falcondev-oss/workflow

Requires Redis (any version with Lua scripting) and Node 24.

Usage

Workflows live in a WorkflowNamespace, which owns the Redis connection, the cross-workflow concurrency cap, and the shared option defaults.

import{createRedis,WorkflowNamespace}from'@falcondev-oss/workflow'import{z}from'zod'constnamespace=newWorkflowNamespace({id: 'my-app',redis: awaitcreateRedis({url: process.env.REDIS_URL}),logger: console,})constworkflow=namespace.createWorkflow({id: 'example-workflow',schema: z.object({timezone: z.string().default('UTC'),name: z.string(),}),asyncrun({ input, step }){awaitstep.do('send welcome',()=>{console.log(`Welcome, ${input.name}! Timezone: ${input.timezone}`)})awaitstep.wait('wait a lil',60_000)constisEngaged=awaitstep.do('check engagement',()=>Math.random()>0.5)if(!isEngaged)return{engagementLevel: 'low'}awaitstep.do('send tips',()=>{console.log(`Here are some tips to get started, ${input.name}!`)})return{engagementLevel: 'high'}},})// Start a worker for this processawaitworkflow.work()// Enqueue a runconstjob=awaitworkflow.run({name: 'John Doe',timezone: 'America/New_York'})// Wait for completion (works from a pure producer too — no worker needed)constresult=awaitjob.wait()console.log(result.engagementLevel)

Watching jobs

Declare a Standard Schema for progress, then emit its input type from any step. Watchers receive the validated output type on the same stream as lifecycle and terminal events.

constexportPdf=namespace.createWorkflow({id: 'export-pdf',schema: z.object({reportId: z.string()}),progressSchema: z.object({label: z.string(),done: z.number(),total: z.number(),}),asyncrun({ input, step }){constrows=awaitstep.do('fetch rows',async({step: nestedStep})=>{awaitnestedStep.progress({label: 'Fetching rows',done: 0,total: 0})returnloadRows(input.reportId)})awaitstep.progress({label: 'Rendering',done: 0,total: rows.length})returnrenderPdf(rows)},})const{ job, events }=awaitexportPdf.runAndWatch({reportId: '1'})forawait(consteventofevents){if(event.type==='progress')console.log(event.data.label)if(event.type==='completed')console.log(event.output)}

Events published before watch() attaches are lost. Lifecycle and progress events are not persisted, and the watcher reads no progress snapshot. Use runAndWatch() when the first event matters because it subscribes before enqueueing. To attach from another request or process, build a handle from the known id:

constjob=awaitexportPdf.getJob(jobId)constevents=awaitjob.watch({signal: request.signal})

watch() subscribes before it resolves and buffers until iteration starts. Breaking out of the loop unsubscribes. A retry emits another started event, while failed is terminal. A completed job still yields its stored terminal result until the configured result TTL expires.

The library does not derive a percentage, ETA, or step count because workflows have no declared step list. If a workflow knows a total, include it in its progress payload as above. Without a progressSchema, step.progress() is a type error and the progress event arm is absent.

Steps

  • step.do(name, fn) — run once, memoize the result. Replayed from Redis on a retry.
  • step.wait(name, ms) — durable sleep; remaining time is computed from the persisted start.
  • step.waitUntil(name, date) — the same, to an absolute time.

Steps nest: the callback receives its own step scoped under the parent's name.

Scheduling

awaitworkflow.run(input)// nowawaitworkflow.runIn(input,60_000)// in 60sawaitworkflow.runAt(input,newDate('2030-01-01'))// at a time// Cron, keyed by (workflow, scheduleId) — upserting the same id replaces in placeawaitworkflow.upsertSchedule('nightly',{pattern: '0 3 * * *',
input,tz: 'Europe/Berlin',})awaitworkflow.getSchedules()awaitworkflow.removeSchedule('nightly')

Ordering, priority and concurrency

Jobs sharing a groupId run one at a time, in enqueue order. Everything else runs in parallel up to the concurrency caps.

namespace.createWorkflow({id: 'per-user',schema: z.object({userId: z.string()}),getGroupId: (input)=>input.userId,// serialize per userqueueOptions: {concurrency: 10,groupConcurrency: 1},workerOptions: {concurrency: 4,maxAttempts: 3},jobOptions: {priority: 1},// 0…2^21-1, higher runs firstrun: async()=>{},})

Options

WorkflowNamespace options are shared defaults — each is shallow-merged under the matching per-workflow override.

OptionDefaultDescription
idNamespace id; scopes the cross-workflow concurrency cap
redisnew clientShared connection, owned by the namespace
prefixwfGlobal key prefix
concurrencyunlimitedCeiling across all workflows in the namespace
loggernoneInherited by every workflow, queue and worker
autoClosetrueClose (drain workers, disconnect) on SIGINT/SIGTERM
queueOptionsDefaults for every workflow's queue
workerOptionsDefaults for every workflow's worker
jobOptionsDefaults for every enqueued job

Shutdown is handled at the namespace: workers drain in-flight jobs, then the connections close. Set autoClose: false and call await namespace.close() yourself to own it.

Failures and retries

A throwing handler is retried up to maxAttempts with an exponential backoff (expBackoff()), then dead-lettered. Throw NonRecoverableError to dead-letter immediately, skipping the remaining budget — the library does this itself when a job's stored payload no longer validates against the workflow schema (a job enqueued before a schema change), since a retry would only re-read the same payload.

Metrics

workflow.getMetrics() returns point-in-time { active, waiting, delayed } depths. Pass an OpenTelemetry meter to export them as gauges:

constnamespace=newWorkflowNamespace({id: 'my-app',workerOptions: {metrics: { meter,prefix: 'my_app'}},})

Spans are emitted for producers, workers and each step via the global OpenTelemetry tracer.

Inspiration

About

Durable, type-safe queue workers on Redis with memoized steps, cron schedules, group ordering, OpenTelemetry tracing

Topics

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Contributors

Languages