Durable workflow engine · Go · NATS JetStream

Every execution is a log.
Everything else follows.

Packtrail interprets declarative flow graphs — task, choice, fanout/join, await, map and subflow nodes, in YAML or Go structs — on nothing but NATS. Each execution is an ordered event log; state, indexes, jobs and timers are derived from it, so time travel, forks, replays and audit come for free. Workers are plain NATS consumers, in any language.

  • Event-sourced
  • NATS-only, no database
  • Workers in any language
  • Human in the loop
  • Time travel & forks
  • Real-server tests, fault injection
Architectural principles

The log is the truth.

Four commitments shape every line of the engine.

01

Event-sourced

Each execution is an ordered log on its own subject. State, snapshots, the visibility index, jobs and timers are projections of it — and can always be rebuilt from it.

02

Workflow as data

You declare the graph; packtrail interprets it. No determinism rules for your code, no replay of your functions. Executions keep the flow version they started with.

03

Workers in any language

A task is a job on <ns>.work.<kind>. A worker answers with one command over a small, versioned JSON protocol. The Go SDK is the reference.

04

Only NATS

One message per decision with an expected-sequence check for single-writer appends, message schedules for timers and cron, KV and object stores for the rest. Nothing else to run.

Capabilities

A small graph language that does the hard parts.

Tasks, retries & timeouts

Per-attempt timeouts, fixed / linear / exponential backoff on durable timers, output JSON Schema, a result cache and per-key concurrency limits across executions.

Fan-out, join & map

Static branches joined with all, any or quorum:N; dynamic maps over a list computed at run time, with bounded parallelism and ordered results.

Typed state channels

Tasks write deltas to channels folded by reducers — replace, append, merge, sum — so parallel branches never overwrite each other. Choices route on expr expressions.

Human in the loop

await nodes wait for a signal with a mandatory timeout route; a worker can interrupt with a question and be resumed with the answer. Update writes state synchronously.

Subflows, dynamic edges & sagas

Call another flow as a child execution; let a worker pick its successor from a declared list; route permanent failures to a compensation node with on_failure.

Time travel & forks

Read the state right after any event, fork a new execution from the past (optionally editing its state), or rerun a node after a fix. Cron schedules and message triggers start flows.

Architecture

Commands in. Decisions logged. Effects derived.

The engine processes commands per partition: load state, decide, append the decision with optimistic concurrency, ack. The dispatcher turns stored events into effects — jobs, timers, child starts, index updates. Workers answer with commands. Nothing is decided twice.

commands Client · CLI · UI · workers start, signal, update, resume, cancel, fork · complete / fail / interrupt
engine Command processors one active per partition · pure decide/fold · expected-sequence append
Dispatcher → effects
jobs on work.<kind>
durable timers
child executions
visibility index
JetStream · KV · object store · message schedules
<ns>-events stream · one message per decision, the source of truth
<ns>-cmd stream · commands, timers, cron
<ns>-work stream · jobs per worker kind
<ns>-index KV · projection for List / Query
<ns>-snapshots KV · fold snapshots every 100 events
Flow definition

Declare the graph. Packtrail walks it.

YAML or Go structs, validated before anything runs: unreachable nodes, routes into a fan-out branch, bad expressions and unknown reducers are rejected up front.

review.yaml
name: review
channels:
  notes: {reducer: append}
nodes:
  - id: draft
    type: task
    kind: writer
    timeout: 1m
    retry: {max_attempts: 3, backoff: exponential}
    next: check
  - id: check
    type: choice
    rules:
      - {when: "results.draft.score >= 8", to: approve}
      - {default: true, to: draft}
  - id: approve
    type: await
    signal: approval
    timeout: 48h
    next: publish
  - {id: publish, type: task, kind: publisher}
start: draft
taskA job for a worker of kind; retries, timeout, cache, concurrency, output schema, dynamic edges.
choiceFirst matching when rule wins; exactly one default.
fanout · joinParallel task branches, closed by a join with all, any or quorum:N.
awaitWait for a named signal; timeout is mandatory, on_timeout routes.
mapOne task per element of over, at most max_parallel at a time.
subflowRun another flow as a child; its output and counters flow back.
Embedding · packages packtrail & worker

An engine, a worker, a client.

Run them in one binary or in a hundred. The engine provisions its namespace on first run; workers share a durable consumer per kind; the client drives executions.

main.go
eng, _ := packtrail.New(nc, packtrail.WithFlowsDir("flows"))
go eng.Run(ctx) // provisions the namespace, processes commands

w, _ := worker.New(nc, "writer",
    func(ctx context.Context, j *worker.Job) (*worker.Result, error) {
        var in struct{ Topic string }
        if err := j.Input(&in); err != nil {
            return nil, worker.Permanent(err) // no retry
        }
        return &worker.Result{
            Output: map[string]any{"text": "…", "score": 9},
            Writes: map[string]any{"notes": "drafted " + in.Topic},
        }, nil
    })
go w.Run(ctx)

c := eng.Client()
id, _ := c.Start(ctx, "review", map[string]any{"topic": "NATS"})
_ = c.Signal(ctx, id, "approval", map[string]any{"by": "ana"})
st, _ := c.Wait(ctx, id) // channels, results, counters, …
worker.Permanent(err)Fail the attempt without retry; the node's on_failure route (if any) takes over.
worker.Interrupt(payload)Pause the execution with a question; Client.Resume runs the node again with the answer.
Result.WritesDeltas for state channels, folded by their reducers.
Result.NextRoute to a declared dynamic successor.
Result.UsageGeneric counters checked against the flow's budget.
j.Progress(v)Stream intermediate results; never stored, followed with Client.Progress.
worker.ErrCancelledThe job's context is cancelled when its work is no longer wanted.
Time travel · CLI

Rewind, branch, retry.

Because the log is the truth, any past point is a valid state. A stream sequence identifies a decision; StateAt, Fork and Rerun cut between decisions. The packtrail CLI exposes all of it.

shell
$ packtrail run -flows flows/ &
$ packtrail start review -input '{"topic":"NATS"}' -wait
$ packtrail history <exec>          # every event
$ packtrail get <exec> -seq 12      # state right after event 12
$ packtrail fork <exec> 12          # continue from there, new execution
$ packtrail fork <exec> 12 -writes '{"notes":"try again"}'
$ packtrail update <exec> -writes '{"notes":"from ops"}'
$ packtrail rerun <exec> draft      # run a node (and what follows) again
HistoryEvery event, including archived segments of long executions.
StateAt(seq)The folded state right after a decision.
Fork(seq, WithForkWrites)A new execution from that state, optionally with edited channels; the source is untouched.
Rerun(node)Fork at the point the node was last entered — the usual move after fixing a worker.
OutputHistory(node)Every output a node produced across its visits.
Observability · packtrail-ui

See every decision.

packtrail-ui is a debugging dashboard and a plain client of the deployment: it needs only NATS. It serves every namespace on the NATS account, switchable from the header.

What it shows

Execution list

Filter by status, flow or search attribute; quarantined executions are flagged.

Flow graph coloured by node state

Every node and route, including on_failure and dynamic edges.

Timeline & state at any sequence

Click an event to see channels, results and counters exactly as they were.

Actions

Signal, resume, cancel, fork, rerun, and redrive dead letters.

start packtrail-ui
# NATS_URL is honoured (default nats://localhost:4222)
$ packtrail-ui -addr 127.0.0.1:8088 \
    -namespaces orders,billing
No built-in authentication. The UI can drive executions, so it binds to loopback by default. Put an authenticating reverse proxy in front before exposing it.