A Go task queue + DAG workflow engine built on NATS JetStream.
Function-first ergonomics (Register(reg, MyFunc), Enqueue(c, MyFunc, args...), Await[T](ctx, fut)) on top of a persistent queue with retries, dead-lettering, and optional DAG orchestration - all driven by a single NATS dependency that can run as an embedded in-process server (including 3-node HA cluster) or against an external JetStream deployment.
- NATS-native. No Redis, no Postgres. If you already run NATS you already have ebind's dependencies.
- Single-binary HA. The
embedpackage boots a 3-node JetStream cluster inside your process. One binary per machine, cluster.routes wired automatically. - Function-first API. Pass your function reference, not a string name and a JSON schema. Reflection introspects the signature; runtime arg-type validation happens before publish.
- Typed responses.
Await[Profile](ctx, fut)returns a typed value, notinterface{}. - Durable DAG workflows. Declare dependencies between steps; state lives in a NATS KV bucket so workflows survive producer restarts. Mandatory/optional steps, per-step retry policies, dynamic step addition from within handlers, step placement across workers, pause/resume/cancel, step breakpoints, and human-in-the-loop signals.
go get github.com/f1bonacc1/ebindRequires Go 1.26+. The embed package bundles its own NATS server; running against an external deployment requires NATS Server 2.9+ with JetStream enabled.
package main
import (
"context"
"fmt"
"log"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/f1bonacc1/ebind/client"
"github.com/f1bonacc1/ebind/embed"
"github.com/f1bonacc1/ebind/stream"
"github.com/f1bonacc1/ebind/task"
"github.com/f1bonacc1/ebind/worker"
)
// Any Go function with (context.Context, ...args) (T, error) or (context.Context, ...args) error.
func SendEmail(ctx context.Context, to, subject, body string) (string, error) {
// ... actually send ...
return "msg-id-42", nil
}
func main() {
ctx := context.Background()
// 1. Start an embedded NATS JetStream (dev). In prod, point at an external cluster
// or use embed.StartCluster(embed.ClusterConfig{Size: 3, ...}) for in-process HA.
node, _ := embed.StartNode(embed.NodeConfig{Port: -1, StoreDir: "/tmp/ebind-demo"})
defer node.Shutdown()
nc, _ := nats.Connect(node.ClientURL())
// 2. Create the ebind streams (TASKS, RESP, DLQ).
js, _ := jetstream.New(nc)
_ = stream.EnsureStreams(ctx, js, stream.Config{Replicas: 1})
// 3. Register handlers + start worker.
reg := task.NewRegistry()
task.MustRegister(reg, SendEmail)
w, _ := worker.New(nc, reg, worker.Options{Concurrency: 16})
go w.Run(ctx)
// 4. Enqueue from anywhere in the cluster - producer and worker can be different binaries.
c, _ := client.New(ctx, nc, client.Options{})
defer c.Close()
fut, _ := client.Enqueue(c, SendEmail, "alice@example.com", "hello", "world")
msgID, err := client.Await[string](ctx, fut)
if err != nil {
log.Fatal(err)
}
fmt.Println("sent:", msgID)
}Don't need the result? client.EnqueueAsync(c, SendEmail, ...) is fire-and-forget: it publishes the task and returns just the task ID, with no response consumer involved.
Run the bundled end-to-end demo:
make demoimport (
"github.com/f1bonacc1/ebind/task"
"github.com/f1bonacc1/ebind/workflow"
)
// Handlers are plain functions.
func FetchUser(ctx context.Context, id string) (User, error) { /* ... */ }
func Enrich(ctx context.Context, id string) (Enriched, error) { /* ... */ }
func Combine(ctx context.Context, u User, e Enriched) (Profile, error)
// In main(): after EnsureStreams + worker started, wire the workflow layer.
wf, _ := workflow.NewFromNATS(ctx, nc, 1 /* replicas */)
task.MustRegister(reg, FetchUser)
task.MustRegister(reg, Enrich)
task.MustRegister(reg, Combine)
// Attach the hook + middleware so the worker talks to the workflow layer.
w, _ := worker.New(nc, reg, worker.Options{
StepHook: wf.Hook(),
Middleware: []worker.Middleware{wf.ContextMiddleware()},
})
go w.Run(ctx)
go wf.RunScheduler(ctx)
// Build + submit a DAG.
dag := workflow.New()
a := dag.Step("fetch", FetchUser, userID)
b := dag.StepOpts("enrich", Enrich, []workflow.StepOption{workflow.Optional()}, userID)
c := dag.Step("combine", Combine, a.Ref(), b.RefOrDefault(Enriched{}))
_ = dag.Submit(ctx, wf)
profile, err := workflow.Await[Profile](ctx, wf, dag.ID(), c)Key behaviors:
a.Ref()-combineruns only iffetchsucceeds; cascade-skips otherwise.b.RefOrDefault(v)-combineruns withvsubstituted ifenrichfails or is skipped.workflow.After(steps...)/workflow.AfterAny(steps...)- temporal-only ordering when the handler doesn't consume upstream output:Aftercascade-skips on upstream failure,AfterAnyruns regardless (best-effort/cleanup steps).workflow.Optional()-enrich's failure does not fail the DAG.workflow.WithRetry(policy)- per-DAG default retry;workflow.WithStepRetry(policy)overrides per-step.workflow.WithLabels("billing", "nightly")- immutable topic tags for querying workflow history (see below).- From inside a handler,
workflow.FromContext(ctx).Step(...)adds more steps dynamically. workflow.Pause/workflow.Resume- graceful whole-DAG pause: in-flight steps drain, pending steps are fenced.workflow.Cancel- durable cancel (alsoebctl dag cancel): pending steps flip to canceled immediately; running steps may finish but schedule no follow-on work.workflow.OnTarget("gpu")and friends - pin steps to specific workers (see below).workflow.SignalRef("approval")/workflow.WaitForSignal(...)- gate a step on an external event;workflow.Signaldelivers it with a payload (see below).workflow.BreakBefore("label")/workflow.BreakAfter("label")- per-step breakpoints (see below).
A step can wait for a named external event - a human approval, a webhook, another DAG - and consume the event's payload as a typed handler argument:
type Approval struct{ Approver string }
// Deploy waits for "approval" (implied by the SignalRef arg) and receives its payload.
func Deploy(ctx context.Context, ap Approval, artifact string) (string, error) { /* ... */ }
dag := workflow.New()
build := dag.Step("build", Build)
deploy := dag.Step("deploy", Deploy, workflow.SignalRef("approval"), build.Ref())
_ = dag.Submit(ctx, wf)
// ...build runs; deploy waits. Any process on the cluster can approve:
delivered, _ := workflow.Signal(ctx, wf, dag.ID(), "approval", Approval{Approver: "eugene"})workflow.SignalRef(name)as an argument = wait for the signal AND receive its payload.workflow.WaitForSignal(names...)as aStepOption= wait only (ALL names must arrive, like deps).- One-shot, buffered, first-wins. A signal is delivered at most once per (DAG, name): the first payload is immutable, duplicates are idempotent no-ops (
delivered=false). A signal sent before any step waits on it - including steps added dynamically later - is kept for the DAG's lifetime. - Waiting is durable (NATS KV) and composes with everything else: independent branches keep running, pause/resume and breakpoints are separate gates, and a waiting DAG survives restarts of every process involved.
- Deliver and inspect from the CLI:
ebctl dag signal ls <dag-id> # NAME | DELIVERED | AT | WAITING | PAYLOAD
ebctl dag signal <dag-id> approval '{"approver":"e"}' # deliver with a JSON payload
ebctl dag ls --sig-waiting # only DAGs awaiting a signal (⚑N marker)
ebctl dag watch # live feed includes ⚑ signal eventsSee examples/17-workflow-signals for a runnable approval-flow walkthrough.
Workers advertise logical targets ("claims"); steps can bind to a claim or to wherever another step ran. Useful when a step needs a GPU box, a host with a warm local cache, or the exact worker holding an open resource:
// Worker side: claim logical targets. Claims can change at runtime
// (worker.ClaimsFunc); every worker also always claims its own WorkerID.
w, _ := worker.New(nc, reg, worker.Options{
WorkerID: "gpu-box-1",
Claims: worker.StaticClaims{"gpu"},
})
// DAG side: bind steps to targets.
train := dag.StepOpts("train", Train, []workflow.StepOption{workflow.OnTarget("gpu")}, data.Ref())
report := dag.StepOpts("report", Report, []workflow.StepOption{workflow.ColocateWith(train)}, train.Ref())OnTarget(target)- run on any worker currently claimingtarget(a logical name or a concreteWorkerID).ColocateWith(step)- run on the concrete worker that executedstep(adds an implicit dependency on it).FollowTargetOf(step)- reusestep's logical target expression, re-resolved at execution time - lands wherever that claim lives now.ColocateHere()- for dynamic steps added viaFromContext: pin the child step to the worker executing the current handler.
See examples/12-workflow-placement for all four modes across two workers with moving claims.
Stop a DAG line at a specific step to inspect intermediate state, then continue - debugger semantics for workflows:
// Breakpoint labels are part of the DAG's structure...
dag := workflow.New()
parse := dag.Step("parse", Parse, download.Ref())
upload := dag.StepOpts("upload", Upload,
[]workflow.StepOption{workflow.BreakBefore("BeforeUpload")}, parse.Ref())
// ...but which ones are ARMED is decided when you run it.
_ = dag.Submit(ctx, wf, workflow.WithActiveBreakpoints("BeforeUpload"))
// The line stops with `upload` pending (never dispatched); parallel branches keep running.
// Inspect whatever you need, then continue:
n, _ := workflow.ResumeBreakpoint(ctx, wf, dag.ID(), "BeforeUpload")BreakBefore(labels...)stops before the step executes (it stayspending);BreakAfter(labels...)lets the step complete - result persisted - but holds its direct dependents.- Breakpoints are inactive by default.
WithActiveBreakpoints(labels...)is aSubmitoption, so a statically defined DAG decides at run time which breakpoints to arm; a breakpoint with several labels is armed (and resumable) by any one of them. ResumeBreakpointis "continue", not "disable": the label stays armed, so a later step (including dynamically added ones) carrying it stops again.- Blocked state lives in NATS KV - it survives restarts, and any process (or
ebctl) can list and resume:
ebctl dag bp ls <dag-id> # STEP | POS | LABELS | ARMED | STATE | SINCE | HOLDING
ebctl dag bp resume <dag-id> <label> # release the blocked line(s); label stays armed
ebctl dag watch # live feed includes bp_hit / bp_resumed eventsSee examples/15-workflow-breakpoints for a runnable two-stop walkthrough.
DAG state + step results live in NATS KV. Workers keep running independently of whoever called Await. If the waiting process dies, the DAG continues; results land in KV and stay there. A different process (same NATS cluster) can resume the wait with only the DAG and step IDs:
// Instance A - submitter. Persist these two strings somewhere (DB, Redis, file).
dagID := dag.ID()
stepID := c.ID() // the *Step you'd pass to Await
_ = dag.Submit(ctx, wf)
// ... instance A may exit now ...
// Instance B - resumer. No *Step handle needed.
wfB, _ := workflow.NewFromNATS(ctx, nc, 1)
result, err := workflow.AwaitByID[Profile](ctx, wfB, dagID, stepID)AwaitByID uses NATS KV IncludeHistory() under the hood, so late subscribers still receive results that were written before they started watching. See examples/11-workflow-resume for a runnable two-invocation demo.
Attach immutable string tags at creation to group workflows by topic, then retrieve only the relevant ones from the history:
// Labels are fixed at New and persisted at Submit - immutable for the DAG's lifetime.
dag := workflow.New(workflow.WithLabels("billing", "nightly"))
// ... add steps ...
_ = dag.Submit(ctx, wf)
// Query the history by label (AND semantics: a DAG must carry all given labels),
// newest-first. No labels returns every DAG.
billing, _ := workflow.ListDAGsByLabels(ctx, wf, "billing") // every "billing" workflow
both, _ := workflow.ListDAGsByLabels(ctx, wf, "billing", "nightly") // only those carrying both- Labels are plain string tags - a workflow carries a set of them; a label is shared across many workflows.
ListDAGsByLabelsis a client-side filter over the full DAG history (the same scanebctl dag lsperforms), so it needs no secondary index and survives in NATS KV for the DAG's lifetime.ebctlexposes the same query and surfaces labels in listings:
ebctl dag ls --label billing # only workflows tagged "billing"
ebctl dag ls --label billing --label nightly # AND: carry both labels
ebctl dag ls # LABELS column shown for every DAG
ebctl dag get <dag-id> # includes a `labels:` lineSee examples/16-workflow-labels for a runnable submit-then-query walkthrough.
When a task or DAG step fails terminally (retries exhausted or a non-retryable error), ebind records the failure in three places:
- The DAG step record (NATS KV, lives for the DAG's lifetime): the step status flips to
failedand the record stores botherror_kind(a category — e.g.handler,deadline,panic) anderror_msg(the handler's error text). This is written by the same hook that fails the step, so it's available as soon as the step showsfailed— no DLQ round-trip needed. - The DLQ (
EBIND_DLQstream, default 7-day retention): the fulldlq.Entrywith the completeTaskError, attempt count, and timestamp. - The response stream (
EBIND_RESP, short retention): theTaskErrorsurfaces throughAwait/Future.Getas a Go error.
Inspect either with the ebctl operator CLI:
# DAG step — durable error kind + message straight from the step record
ebctl dag get <dag-id> # every step: status + error kind
ebctl dag step get <dag-id> <step-id> # one step: error_kind + error_msg
# DLQ — terminally failed tasks with the full error message
ebctl dlq ls # table: SEQ, FN, ERROR, ATTEMPT, AGE
ebctl dlq show <seq> # full entry incl. error_message + payloadThe error text persisted into the step record is bounded by Workflow.MaxStepErrorBytes (default 4096; set negative to keep only the error kind — useful when error text may contain sensitive data). The full, untruncated message is always available in the DLQ entry. Note that only terminal failures are recorded; an attempt that will still be retried writes nothing to the step record — use the worker.Log middleware to capture per-attempt errors.
flowchart LR
Producer["<b>Producer</b><br/>client.Enqueue"]
Tasks["<b>NATS JetStream</b><br/>TASKS stream"]
Worker["<b>Worker(s)</b><br/>• reflect.Call<br/>• middleware chain<br/>• StepHook"]
Future["Future.Get() /<br/>Await[T]"]
Scheduler["<b>Scheduler</b><br/>(every worker, leader-gated)<br/>• state in NATS KV<br/>• resync on acquire"]
Producer -- "publish TASKS.<name>" --> Tasks
Tasks -- "pull" --> Worker
Worker -- "RESP.<client_id>.>" --> Future
Worker -- "DAG events" --> Scheduler
task.Registry- name → reflect.Value map;Register(fn)introspects signature.client.Client- one response-consumer per client; routes responses to typed Futures.worker.Worker- pull consumer + middleware chain (Recover,Log, user) + per-task retry policy.embed.StartCluster(3)- in-process 3-node JetStream cluster with loopback routes.workflow- DAG builder + persistent state (KV bucketebind-dags) + event-driven scheduler with leader-elector-gated sweep for stranded recovery.
See CLAUDE.md for the full architectural walk-through.
| Concern | Handled by |
|---|---|
| At-least-once delivery | JetStream AckExplicitPolicy + MaxDeliver |
| Exactly-one enqueue on retries | JetStream Nats-Msg-Id dedupe with 5-min window |
| State consistency | NATS KV Update(key, val, expectedRev) CAS |
| Handler panics | worker.Recover middleware → TaskError{Kind: "panic"} |
| Retry control | task.RetryPolicy on envelope (per-task), worker.Options.DefaultRetryPolicy fallback, optional MaxDeliver hard cap |
| Non-retryable errors | RetryPolicy.NonRetryableErrorKinds OR TaskError{Retryable: false} |
| Dead-lettering | dlq.Publish auto-called on final failure → EBIND_DLQ stream |
| Failure visibility | Failed step record carries error_kind + error_msg (durable in KV, via ebctl dag step get); full TaskError in EBIND_DLQ; cap with Workflow.MaxStepErrorBytes |
| Graceful shutdown | worker.Run(ctx) drains in-flight on ctx-cancel (configurable grace) |
| HA | 3-node embedded cluster with Replicas: 3 streams & KV |
| Stranded DAG recovery | Scheduler sweep on LeaderElector false→true edge |
task/ envelope + registry + RetryPolicy
client/ Enqueue + Future + Await[T]
worker/ consume loop + middleware + StepHook
stream/ JetStream stream setup
dlq/ dead-letter publishing
embed/ in-process NATS server (single + cluster)
workflow/ DAG builder + scheduler + KV-backed state
internal/testutil/ harness for integration tests
cmd/demo/ single-process end-to-end demo
cmd/ebctl/ operator CLI: inspect DAGs/steps/results, DLQ, streams
examples/ holds 17 self-contained runnable programs, each starting its own embedded NATS - from basic enqueue/await through retries, middleware, and cluster HA to every workflow feature (fan-out, optional steps, dynamic steps, temporal deps, placement, cancel, pause/resume, breakpoints, label queries, signals). Start with the index.
make help # list all targets
make build # compile everything
make test # run all tests with -race
make test-count # 3× runs, catch flakes
make lint # golangci-lint
make cover # HTML coverage report
make demo # run cmd/demo end-to-end- v1 (task queue): done - retries, DLQ, middleware, embedded HA cluster.
- v2 (DAG workflows): done - optional steps, retry policies, dynamic DAGs.
- v2.1 (stranded recovery): done - leader-acquisition sweep.
- v2.2 (signals): done - human-in-the-loop waits with payload delivery, cross-DAG capable.
- v2+ (future): phantom-Running detection, saga/compensation, durable timers.
MIT