Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions docs/content/api/conductor.md
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ const task = conductor.createTask(
removeOnComplete?: { days: number } | false; // Retention policy
removeOnFail?: { days: number } | false; // Retention policy
batch?: { size: number; timeoutMs: number }; // Batch processing config
deadLetter?: { queue: string; task?: Task }; // Final-failure destination
}
```

Expand Down
24 changes: 24 additions & 0 deletions docs/content/task-execution/dead-letter-queue.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
# Dead-letter queues

Configure a destination for executions that fail on their final attempt:

```ts
const failedPayment = conductor.createTask(
{ name: "failed-payment", queue: "payments-dlq" },
{ invocable: true },
async (event) => {},
);

const chargeCard = conductor.createTask(
{
name: "charge-card",
deadLetter: { queue: "payments-dlq", task: failedPayment },
},
{ invocable: true },
async (event) => {},
);
```

The destination is a new execution with the original payload. Its execution row contains machine-readable source execution ID, source queue and task, final error, attempt count, and failure timestamp. The source remains a normal failed execution unless its retention policy removes it.

Retries and cancellation do not deliver to a dead-letter queue. Delivery is transactional, claim-fenced, and idempotent. A destination may have its own retry, retention, and concurrency settings. Chains are supported, but a task cannot target itself directly. The destination task must accept the source payload.
1 change: 1 addition & 0 deletions docs/zensical.toml
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@ nav = [
{ "Cancellation" = "task-execution/cancellation.md" },
{ "Priority" = "task-execution/priority.md" },
{ "Concurrency" = "task-execution/concurrency.md" },
{ "Dead-letter queues" = "task-execution/dead-letter-queue.md" },
{ "Deduplication" = "task-execution/deduplication.md" },
{ "Rate Limiting" = "task-execution/rate-limiting.md" },
{ "Batching" = "task-execution/batching.md" },
Expand Down
42 changes: 38 additions & 4 deletions migrations/0000000001_setup.sql
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,14 @@ create table pgconductor._private_executions (
waiting_step_key text,
parent_execution_id uuid,
singleton_on timestamptz,

-- Dead-letter metadata is denormalized so retained source rows are optional.
dead_letter_source_execution_id uuid,
dead_letter_source_queue text,
dead_letter_source_task_key text,
dead_letter_error text,
dead_letter_attempts integer,
dead_letter_failed_at timestamptz,
primary key (id, queue),
unique (task_key, dedupe_key, queue)
) partition by list (queue);
Expand Down Expand Up @@ -147,10 +155,19 @@ create table pgconductor._private_tasks (
-- NULL means no limit (unlimited concurrency)
concurrency_limit integer,
group_concurrency_limit integer,

-- Destination copied onto each source task registration.
dead_letter_queue text,
dead_letter_task_key text,

constraint positive_concurrency_limits check (
(concurrency_limit is null or concurrency_limit > 0) and
(group_concurrency_limit is null or group_concurrency_limit > 0)
),
constraint dead_letter_not_self check (
dead_letter_queue is null or dead_letter_queue <> queue or
(dead_letter_task_key is not null and dead_letter_task_key <> key)
),

primary key (queue, key)
);
Expand All @@ -168,6 +185,10 @@ create table pgconductor._private_steps (

create index idx_steps_execution_id on pgconductor._private_steps (execution_id);

create unique index idx_executions_dead_letter_delivery
on pgconductor._private_executions (dead_letter_source_execution_id, queue, task_key)
where dead_letter_source_execution_id is not null;

-- Trigger function to manage executions partitions per queue
-- Automatically creates partition when queue is inserted
create or replace function pgconductor._private_manage_queue_partition()
Expand Down Expand Up @@ -329,7 +350,9 @@ create type pgconductor.task_spec as (
window_start timetz,
window_end timetz,
concurrency_limit integer,
group_concurrency_limit integer
group_concurrency_limit integer,
dead_letter_queue text,
dead_letter_task_key text
);

create type pgconductor._private_event_operation as enum (
Expand Down Expand Up @@ -367,8 +390,15 @@ begin
values (p_queue_name)
on conflict (name) do nothing;

-- Dead-letter destinations may not have a worker yet; create their partitions.
insert into pgconductor._private_queues (name)
select distinct spec.dead_letter_queue
from unnest(p_task_specs) as spec
where spec.dead_letter_queue is not null
on conflict (name) do nothing;

-- step 2: register/update tasks
insert into pgconductor._private_tasks (key, queue, max_attempts, remove_on_complete_days, remove_on_fail_days, window_start, window_end, concurrency_limit, group_concurrency_limit)
insert into pgconductor._private_tasks (key, queue, max_attempts, remove_on_complete_days, remove_on_fail_days, window_start, window_end, concurrency_limit, group_concurrency_limit, dead_letter_queue, dead_letter_task_key)
select
spec.key,
coalesce(spec.queue, 'default'),
Expand All @@ -378,7 +408,9 @@ begin
spec.window_start,
spec.window_end,
spec.concurrency_limit,
spec.group_concurrency_limit
spec.group_concurrency_limit,
spec.dead_letter_queue,
spec.dead_letter_task_key
from unnest(p_task_specs) as spec
on conflict (queue, key)
do update set
Expand All @@ -389,7 +421,9 @@ begin
window_start = excluded.window_start,
window_end = excluded.window_end,
concurrency_limit = excluded.concurrency_limit,
group_concurrency_limit = excluded.group_concurrency_limit;
group_concurrency_limit = excluded.group_concurrency_limit,
dead_letter_queue = excluded.dead_letter_queue,
dead_letter_task_key = excluded.dead_letter_task_key;

-- step 3: insert scheduled cron executions
insert into pgconductor._private_executions (task_key, queue, payload, run_at, dedupe_key, cron_expression, "group")
Expand Down
9 changes: 7 additions & 2 deletions packages/pgconductor-js/src/conductor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import {
type ValidateTasksQueue,
type BatchConfig,
type ExecuteFunction,
type ValidateDeadLetterConfiguration,
} from "./task";
import type { TaskContext, BatchTaskContext } from "./task-context";
import {
Expand Down Expand Up @@ -162,7 +163,7 @@ export class Conductor<
},
const TTriggers extends object | readonly object[],
>(
definition: TDef,
definition: TDef & ValidateDeadLetterConfiguration<TDef, ResolvedPayload<Tasks, TDef>>,
triggers: TTriggers & ValidateTriggers<Tasks, TDef["name"], TTriggers, ResolvedQueue<TDef>>,
fn: TDef extends { readonly batch: BatchConfig }
? ResolvedReturns<Tasks, TDef> extends void
Expand Down Expand Up @@ -198,7 +199,11 @@ export class Conductor<
TaskContext<Tasks, Events> & ExtraContext,
TaskEventFromTriggers<TTriggers, ResolvedPayload<Tasks, TDef>, Events, Database>
>(
definition as TaskConfiguration<TDef["name"], ResolvedQueue<TDef>>,
definition as TaskConfiguration<
TDef["name"],
ResolvedQueue<TDef>,
ResolvedPayload<Tasks, TDef>
>,
triggers as NonEmptyArray<Trigger> | Trigger,
fn as ExecuteFunction<
TaskEventFromTriggers<TTriggers, ResolvedPayload<Tasks, TDef>, Events, Database>,
Expand Down
17 changes: 17 additions & 0 deletions packages/pgconductor-js/src/database-client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,19 @@ export interface TaskSpec {
window?: [string, string] | null;
concurrency?: number | null;
groupConcurrency?: number | null;
deadLetterQueue?: string | null;
deadLetterTaskKey?: string | null;
}

export type DeadLetterMetadata = {
sourceExecutionId: string;
sourceQueue: string;
sourceTaskKey: string;
error: string | null;
attempts: number;
failedAt: Date;
};

export interface Execution {
id: string;
task_key: string;
Expand All @@ -63,6 +74,12 @@ export interface Execution {
dedupe_key?: string | null;
cron_expression?: string | null;
group?: string | null;
dead_letter_source_execution_id?: string | null;
dead_letter_source_queue?: string | null;
dead_letter_source_task_key?: string | null;
dead_letter_error?: string | null;
dead_letter_attempts?: number | null;
dead_letter_failed_at?: Date | null;
}

// todo: move all of this to query-builder too or create new types.ts file
Expand Down
42 changes: 38 additions & 4 deletions packages/pgconductor-js/src/generated/sql.ts
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,14 @@ create table pgconductor._private_executions (
waiting_step_key text,
parent_execution_id uuid,
singleton_on timestamptz,

-- Dead-letter metadata is denormalized so retained source rows are optional.
dead_letter_source_execution_id uuid,
dead_letter_source_queue text,
dead_letter_source_task_key text,
dead_letter_error text,
dead_letter_attempts integer,
dead_letter_failed_at timestamptz,
primary key (id, queue),
unique (task_key, dedupe_key, queue)
) partition by list (queue);
Expand Down Expand Up @@ -163,10 +171,19 @@ create table pgconductor._private_tasks (
-- NULL means no limit (unlimited concurrency)
concurrency_limit integer,
group_concurrency_limit integer,

-- Destination copied onto each source task registration.
dead_letter_queue text,
dead_letter_task_key text,

constraint positive_concurrency_limits check (
(concurrency_limit is null or concurrency_limit > 0) and
(group_concurrency_limit is null or group_concurrency_limit > 0)
),
constraint dead_letter_not_self check (
dead_letter_queue is null or dead_letter_queue <> queue or
(dead_letter_task_key is not null and dead_letter_task_key <> key)
),

primary key (queue, key)
);
Expand All @@ -184,6 +201,10 @@ create table pgconductor._private_steps (

create index idx_steps_execution_id on pgconductor._private_steps (execution_id);

create unique index idx_executions_dead_letter_delivery
on pgconductor._private_executions (dead_letter_source_execution_id, queue, task_key)
where dead_letter_source_execution_id is not null;

-- Trigger function to manage executions partitions per queue
-- Automatically creates partition when queue is inserted
create or replace function pgconductor._private_manage_queue_partition()
Expand Down Expand Up @@ -345,7 +366,9 @@ create type pgconductor.task_spec as (
window_start timetz,
window_end timetz,
concurrency_limit integer,
group_concurrency_limit integer
group_concurrency_limit integer,
dead_letter_queue text,
dead_letter_task_key text
);

create type pgconductor._private_event_operation as enum (
Expand Down Expand Up @@ -383,8 +406,15 @@ begin
values (p_queue_name)
on conflict (name) do nothing;

-- Dead-letter destinations may not have a worker yet; create their partitions.
insert into pgconductor._private_queues (name)
select distinct spec.dead_letter_queue
from unnest(p_task_specs) as spec
where spec.dead_letter_queue is not null
on conflict (name) do nothing;

-- step 2: register/update tasks
insert into pgconductor._private_tasks (key, queue, max_attempts, remove_on_complete_days, remove_on_fail_days, window_start, window_end, concurrency_limit, group_concurrency_limit)
insert into pgconductor._private_tasks (key, queue, max_attempts, remove_on_complete_days, remove_on_fail_days, window_start, window_end, concurrency_limit, group_concurrency_limit, dead_letter_queue, dead_letter_task_key)
select
spec.key,
coalesce(spec.queue, 'default'),
Expand All @@ -394,7 +424,9 @@ begin
spec.window_start,
spec.window_end,
spec.concurrency_limit,
spec.group_concurrency_limit
spec.group_concurrency_limit,
spec.dead_letter_queue,
spec.dead_letter_task_key
from unnest(p_task_specs) as spec
on conflict (queue, key)
do update set
Expand All @@ -405,7 +437,9 @@ begin
window_start = excluded.window_start,
window_end = excluded.window_end,
concurrency_limit = excluded.concurrency_limit,
group_concurrency_limit = excluded.group_concurrency_limit;
group_concurrency_limit = excluded.group_concurrency_limit,
dead_letter_queue = excluded.dead_letter_queue,
dead_letter_task_key = excluded.dead_letter_task_key;

-- step 3: insert scheduled cron executions
insert into pgconductor._private_executions (task_key, queue, payload, run_at, dedupe_key, cron_expression, "group")
Expand Down
Loading
Loading