Skip to content

Repository files navigation

tasking — reliable async tasks, DAG workflows & iterative computations for Go

tasking is an embeddable Go library for running asynchronous background tasks, multi-step DAG workflows, and cyclic bulk-synchronous computations reliably — with retries, timeouts, cancellation, and a subscribable notification stream — surviving process restarts and worker crashes.

It is a library embedded by other applications: it orchestrates execution, but it does not own multi-tenancy or access control — those belong to the embedding application. Each component's DESIGN.md records the rationale for that boundary.

Use cases

  • Background jobs — submit a unit of work and have it run reliably: immediate or scheduled one-shot, with per-attempt retry (exponential backoff), a per-attempt timeout, and cancellation.
  • Multi-step pipelines — express work as a DAG of steps with dependency edges; the engine fans out steps that can run in parallel, dispatches each step once its parents finish, enforces a workflow-wide deadline, and lets you revive a failed workflow or cancel one (e.g. render → thumbnail → publish).
  • Cyclic, iterative computations — express work as a set of communicating vertices that run in lockstep supersteps, exchanging messages and iterating until they converge. A run has no success terminal: it rests when nothing is scheduled and wakes when external input arrives, so it also expresses long-lived interactive work (an LLM agent loop being the case it was built for).
  • Reacting to state changes — subscribe to task, workflow and BSP lifecycle events over Redis pub/sub.

Intentionally out of scope (belongs to the embedding application): multi-tenancy isolation, and access control. Periodic/recurring tasks are also not yet supported (only immediate and scheduled one-shot). See each component's DESIGN.md.

How the pieces fit

flowchart TD
    subgraph app["your application"]
      WF["workflow engine<br>(DAG orchestration + state)"]
      BSP["BSP engine<br>(supersteps + barrier)"]
      TASK["task engine<br>(execution: retry / timeout / cancel)"]
      NOTIFY["notify<br>(Redis pub/sub over the audit log)"]
    end
    DB[("db<br>durable state +<br>SystemEventAudit log")]

    WF -- "every step runs as a task<br>(__EXECUTE_WORKFLOW_STEP__)" --> TASK
    BSP -- "every vertex superstep runs as a task<br>(__EXECUTE_BSP_VERTEX__)" --> TASK
    WF --> DB
    BSP --> DB
    TASK --> DB
    DB -- "un-broadcast audit rows" --> NOTIFY
    NOTIFY -- "step results (creator channel)" --> WF
    NOTIFY -- "vertex results (creator channel)" --> BSP
Loading

Two engines are layered on the task engine, and they are siblings rather than alternatives. The workflow engine owns DAG orchestration and state — acyclic, each step run once, complete when its sinks do. The BSP engine owns iteration — cycles are the point, a vertex runs many times, and the run rests rather than completing. The task engine owns the execution underneath both: per-attempt retry and timeout are its job. Every workflow step and every vertex superstep, whatever its type, is executed as an ordinary task under a reserved name (__EXECUTE_WORKFLOW_STEP__ and __EXECUTE_BSP_VERTEX__), so both engines inherit the task engine's reliability for free. All three write durable, creator-tagged SystemEventAudit rows in the same transaction as each state change; notify turns those rows into a best-effort Redis pub/sub stream — and both schedulers are themselves notify subscribers, which is how their results reach them.

Shared reliability model. Across every component, the database is the source of truth. Every IPC message is a best-effort poke that lets a component act sooner than its periodic maintenance sweep would; lose a poke and work is delayed, never lost. IPC rides crash-safe reliable Redis queues, handlers are idempotent, and unprocessable messages are quarantined rather than crash-looped.

Components

Component Purpose Docs
Task Engine (task) Reliable async execution of a single unit of work README · DESIGN
Workflow Engine (workflow) DAG orchestration of multi-step workflows over the task engine README · DESIGN
BSP Engine (bsp) Cyclic, iterative computation over communicating vertices, in supersteps README · DESIGN
Notifications (notify) Best-effort Redis pub/sub stream over the durable audit log README · DESIGN

Underpinning all four: the db package (db.Client and transaction support — persistence and the SystemEventAudit log) and the models package (configuration structs, wire types, and the engine's state enums).

Sitting above them, the root tasking package provides one-call constructors for the wiring below — NewPersistenceClient, NewTaskWorkerReceiver, NewTaskScheduler, NewTaskClient, DefineTaskQueueForWorkflowStep, NewWorkflowScheduler, NewWorkflowClient, DefineTaskQueueForBSPVertex, NewBSPScheduler, NewBSPClient, and NewNotificationProducer. Each fills in the standard Redis IPC and executor factories; drop to the per-package constructors (task.NewReceiver, workflow.NewClient, bsp.NewBSPScheduler, …) when you need to substitute one, as tests do.

Wiring it together

A single embedding application typically stands up every component it needs, using the root package's constructors. The snippet below runs all four — take the workflow half or the BSP half alone if that is all you use. It is otherwise complete; errors are discarded into _ purely to keep the composition readable.

The composition point to notice is DefineTaskQueueForWorkflowStep, and its BSP twin DefineTaskQueueForBSPVertex: each builds that engine's runner and hands it back as an ordinary task queue config, running under a reserved task name (__EXECUTE_WORKFLOW_STEP__, __EXECUTE_BSP_VERTEX__). Those queues join your own in one []models.TaskQueueConfig, which is how both engines inherit the task engine for execution, retry and timeout. Processors are supplied declaratively — each is declared on its task, inside the queue that runs it — so there is no runtime registration call: the receiver builds each queue's task name → processor map from the config and hands it to that queue's executor.

// --- One persistence client, shared by every component in this process ---
dbClient, _ := tasking.NewPersistenceClient(postgres.Open(dsn), logger.Warn)

// --- Workflow step handlers: one per step Type ---
stepHandlers := map[string]models.WorkflowStepProcessor{
    "render-html":     renderHandler{},
    "make-thumbnails": thumbnailHandler{},
    "push-cdn":        cdnHandler{},
}

// --- BSP vertex handlers: one per vertex definition ProcessorType ---
vertexHandlers := map[string]models.BSPVertexComputeProcessor{
    "reason": reasonerCompute{},
    "search": searchToolCompute{},
}

// --- Where workflow plugs into task: this builds the Step Runner and returns an ordinary task
//     queue which runs it under the reserved name __EXECUTE_WORKFLOW_STEP__ ---
workflowQueue, _ := tasking.DefineTaskQueueForWorkflowStep(
    "workflow-queue", 4 /* workers */, 16 /* buffered requests */, dbClient, stepHandlers,
)

// --- The same move for BSP: the Vertex Runner, under __EXECUTE_BSP_VERTEX__ ---
bspQueue, _ := tasking.DefineTaskQueueForBSPVertex(
    "bsp-queue", 4 /* workers */, 16 /* buffered requests */, dbClient, vertexHandlers,
)

// --- Your own queues, plus the two engine ones. Task names must be unique across the whole set ---
taskQueues := []models.TaskQueueConfig{
    {
        Name:           "default-queue",
        Workers:        4,
        BufferRequests: 16,
        SupportedTasks: []models.TaskDefinition{
            {
                TaskName:  "resize-image",
                Retry:     models.RetryParam{InitialDelaySec: 5, MaxRetries: 3},
                Processor: resizeProcessor{},
            },
        },
    },
    workflowQueue,
    bspQueue,
}

// --- Task Receiver: runs the tasks. A nil OnFatal keeps the default (log and terminate) ---
receiver, _ := tasking.NewTaskWorkerReceiver(ctx, models.TaskReceiverConfig{
    Name: "worker-01", SchedulerQueue: "task-scheduler", TaskQueues: taskQueues,
}, dbClient, redisClient, nil)
_ = receiver.Initialize(ctx, nil) // MUST run before Start: reconciles buffered work after a crash
_ = receiver.Start(ctx)
defer receiver.Stop(ctx)

// --- Task Scheduler: the single writer of task state. It gets the SAME taskQueues — that is what
//     tells it which queue to poke for __EXECUTE_WORKFLOW_STEP__, __EXECUTE_BSP_VERTEX__, and for
//     your own task names ---
taskScheduler, _ := tasking.NewTaskScheduler(ctx, models.TaskSchedulerConfig{
    MaintenanceTimerIntSecs: 30, SchedulerQueue: "task-scheduler", TaskQueues: taskQueues,
}, dbClient, redisClient, nil)
_ = taskScheduler.Start(ctx)
defer taskScheduler.Stop(ctx)

// --- notify Producer: broadcasts audit rows. EmitCreator is REQUIRED for workflow AND BSP feedback
//     — both schedulers subscribe on notify:creator:<engine-creator>, and without it every step and
//     every vertex sits in RUNNING until its deadline expires ---
producer, _ := tasking.NewNotificationProducer(ctx, models.NotificationProducerConfig{
    PollIntervalSecs: 5, BatchSize: 100, EmitCreator: true,
}, dbClient, redisClient)
_ = producer.Start(ctx)
defer producer.Stop(ctx)

// --- Workflow Scheduler: single writer of workflow state. It builds its own task client (to
//     dispatch steps) and notify consumer (to receive their results) internally ---
wfScheduler, _ := tasking.NewWorkflowScheduler(
    ctx,
    models.WorkflowSchedulerConfig{
        MaintenanceTimerIntSecs: 30, SchedulerQueue: "workflow-scheduler",
    },
    "task-scheduler", // the task scheduler's queue, where step tasks are submitted
    dbClient, redisClient, nil,
)
_ = wfScheduler.Start(ctx)
defer wfScheduler.Stop(ctx)

// --- Submit workflows. The step Types listed here must match the handler map above ---
wfClient, _ := tasking.NewWorkflowClient(
    ctx, "my-app", "my-app",
    models.WorkflowClientConfig{SchedulerQueue: "workflow-scheduler"},
    []string{"render-html", "make-thumbnails", "push-cdn"},
    dbClient, redisClient,
)
wf, _ := wfClient.DefineAndRunWorkflow(ctx, workflow.DefineWorkflowParams{ /* ... */ }, nil)
_ = wf

// --- BSP Scheduler: single writer of branch and vertex state. Like the workflow one it builds its
//     own task client and notify consumer internally — and here that matters more, since the client
//     must be scoped to the engine's reserved creator or every cancel is silently refused ---
bspScheduler, _ := tasking.NewBSPScheduler(
    ctx,
    models.BSPSchedulerConfig{
        MaintenanceTimerIntSecs: 30, BranchStallLimitSecs: 300, SchedulerQueue: "bsp-scheduler",
    },
    "task-scheduler", // the task scheduler's queue, where vertex tasks are submitted
    dbClient, redisClient, nil,
)
_ = bspScheduler.Start(ctx)
defer bspScheduler.Stop(ctx)

// --- Submit BSP groups. The ProcessorTypes listed here must match the handler map above ---
bspClient, _ := tasking.NewBSPClient(
    ctx, "my-app", "my-app",
    models.BSPClientConfig{SchedulerQueue: "bsp-scheduler"},
    []string{"reason", "search"},
    dbClient, redisClient,
)
group, _, rootBranch, _ := bspClient.DefineAndRunBSPGroup(
    ctx, bsp.DefineBSPGroupParams{ /* ... */ }, nil,
)
_, _ = group, rootBranch

// --- Submit ordinary tasks. Only needed if your app dispatches work of its own; neither engine's
//     scheduler uses this client ---
taskClient, _ := tasking.NewTaskClient(
    ctx, "my-app", "my-app",
    models.TaskClientConfig{SchedulerQueue: "task-scheduler", TaskQueues: taskQueues},
    dbClient, redisClient,
)
_ = taskClient

Routing. Because a workflow step and a vertex superstep are both ordinary tasks, the queues returned by DefineTaskQueueForWorkflowStep and DefineTaskQueueForBSPVertex must reach both the receiver's TaskReceiverConfig and the scheduler's TaskSchedulerConfig — the receiver runs those tasks, the scheduler decides which queue to dispatch them to. Give one to the receiver alone and its tasks are defined but never routed, so workflows stall with nothing reported and BSP branches hold one position until the sweep settles them. The reserved names are the engines' own: never declare __EXECUTE_WORKFLOW_STEP__ or __EXECUTE_BSP_VERTEX__ in a queue of your own. See task/README.md, workflow/README.md and bsp/README.md for the full flow.

Requirements & getting started

  • Module: github.com/alwitt/tasking (Go 1.26+).
  • PostgreSQL-compatible database — durable state and the SystemEventAudit log, via the db package.
  • Redis — IPC message queues (engine coordination) and pub/sub (notifications).

Build and test targets live in the Makefile.

License

Released under the MIT License.

About

Embeddable Go library for reliable background tasks and DAG workflows

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Used by

Contributors

Languages