Documentation

The full Packtrail reference.

Everything needed to declare flows, write workers, drive executions and operate a deployment. Packtrail is an event-sourced workflow engine on NATS JetStream: every execution is an ordered log of decisions, and everything else — state, indexes, jobs, timers — is derived from it. The Go API reference is on pkg.go.dev.

Getting started

Installation

Packtrail is a Go module. It needs Go 1.26+ and a NATS Server 2.12+ with JetStream enabled: durable timers and cron are JetStream message schedules.

shell
$ go get github.com/henomis/packtrail
$ go install github.com/henomis/packtrail/cmd/packtrail@latest      # CLI
$ go install github.com/henomis/packtrail/cmd/packtrail-ui@latest   # dashboard

# a NATS server with JetStream for local development
$ docker run --rm -p 4222:4222 nats:2.14.2 -js

Packages: packtrail (engine and client), worker (the Go worker SDK), flow (definitions: parse, build, validate) and event (event types, for readers of History and Watch). Packtrail never owns the *nats.Conn you pass in: your code connects and closes it.


Getting started

Quick start

An engine, one worker and a client in a single program. In production they are usually separate processes — they share nothing but NATS.

main.go
const hello = `
name: hello
nodes:
  - {id: greet, type: task, kind: greeter, next: shout}
  - {id: shout, type: task, kind: shouter}
`

nc, _ := nats.Connect(nats.DefaultURL)
defer nc.Close()

eng, err := packtrail.New(nc, packtrail.WithFlowYAML([]byte(hello)))
if err != nil { log.Fatal(err) }
if err = eng.Init(ctx); err != nil { log.Fatal(err) } // provision, register flows
go eng.Run(ctx)

greeter, _ := worker.New(nc, "greeter",
    func(ctx context.Context, j *worker.Job) (*worker.Result, error) {
        var in struct{ Name string }
        if err := j.Input(&in); err != nil {
            return nil, worker.Permanent(err)
        }
        return &worker.Result{Output: map[string]any{"text": "hello " + in.Name}}, nil
    })
go greeter.Run(ctx)

shouter, _ := worker.New(nc, "shouter",
    func(ctx context.Context, j *worker.Job) (*worker.Result, error) {
        var prev struct{ Text string }
        _ = j.Result("greet", &prev) // a previous node's output
        return &worker.Result{Output: map[string]any{"text": strings.ToUpper(prev.Text)}}, nil
    })
go shouter.Run(ctx)

c := eng.Client()
id, _ := c.Start(ctx, "hello", map[string]any{"name": "world"})
st, _ := c.Wait(ctx, id)
fmt.Println(st.Status, string(st.Results["shout"])) // completed {"text":"HELLO WORLD"}
New performs no I/O. It validates the configuration — flows, schedules, namespace, tuning — and returns. Init provisions the namespace (streams, buckets, object stores), checks the server and registers the flows; Run calls it if needed. Call Init explicitly when you want Client() before Run, or provisioning errors to fail fast.

The repository has runnable examples: human in the loop, parallel research with reducers, map-reduce, retries and timeouts, time travel, an agent-style tool loop, subflows, cron and message triggers. make examples runs them all.


Getting started

Core concepts

ConceptDescription
FlowAn immutable, versioned definition (<name>.<hash>). Registering a changed definition creates a new version; executions keep the version they started with.
ExecutionAn event log on <ns>.ev.<partition>.<exec>. Status running, waiting, completed, failed or cancelled.
EngineProcesses commands per partition: load state → decide → append the decision with an expected-sequence check → ack. One engine is active per partition; others take over on failure.
DispatcherTurns stored events into effects: jobs, timer schedules, child starts, parent notifications, index updates, archival.
WorkerServes one kind: consumes jobs, heartbeats long ones, answers with complete, fail or interrupt.
ChannelsTyped state written by tasks as deltas and folded by reducers.
ContextWhat workers and expressions see: input, channels, results, last_node, visits, signals, branches, counters, errors, and item/index/resume when set.

Workflow as data

You declare the graph; packtrail interprets it. That puts it closer to other graph-based engines than to the workflow-as-code model, whose workflows are code replayed deterministically. What it shares with other durable engines is the durability model: event history, durable timers, workers in any language. Your worker code has no determinism rules.

Delivery semantics

Delivery to workers is at-least-once: after a crash a task can run more than once. Make side effects idempotent, or enable the result cache. The engine itself is exactly-once per decision: every result carries the task's generation and attempt, and duplicate or stale commands are absorbed by the fold without producing events.

Work that is no longer wanted is stopped: when an execution ends, a task is cancelled (a join settled, a map aborted) or an attempt times out, the running job's context is cancelled with cause worker.ErrCancelled. The handler should return promptly; its result is dropped.


Flow definitions

Anatomy of a flow

A flow is YAML (one flow per document) or a flow.Flow struct. Both go through the same validation: unknown node types, fields that do not belong to a node type, dangling routes, unreachable nodes, routes into a fan-out branch, invalid expressions and bad reducer defaults are all rejected before anything runs. packtrail validate flows/*.yaml does it offline.

research.yaml
name: research
description: Ask three sources in parallel, write when two answered
channels:
  notes:   {reducer: append}
  sources: {reducer: merge}
  cost:    {reducer: sum}
budget: {cost: 10}
max_steps: 100
search_attributes: {topic: "input.topic"}
retention: 720h
nodes:
  - {id: gather, type: fanout, branches: [archive, journals, forums], next: enough}
  - {id: archive,  type: task, kind: source, timeout: 5s}
  - {id: journals, type: task, kind: source, timeout: 5s}
  - {id: forums,   type: task, kind: source, timeout: 5s}
  - {id: enough, type: join, policy: "quorum:2", next: write}
  - {id: write, type: task, kind: writer}
FieldDescription
nameFlow name, [A-Za-z0-9_-]{1,128}. Required.
versionDefinition schema version. Omit it; the current schema is assumed.
descriptionFree text.
startEntry node. When omitted, the unique node with no inbound route.
channelsTyped state: {reducer, default} per channel — see channels.
outputAn expression evaluated on completion that yields what the execution returns: an object or null. Default: every channel, or the last node's result when the flow declares none — see channels.
nodesThe graph. Node ids are tokens and unique; next is the static successor (empty on a terminal node).
budgetCaps on generic counters reported by workers (Result.Usage, worker.WithUsage for runs that fail or interrupt) and the built-in steps. Exceeding one fails the execution with budget_exceeded.
max_stepsRecursion limit: total node entries (default 1000). Exceeding it fails with max_steps.
search_attributesExpressions evaluated on the start input and indexed, for List / packtrail list -attr.
retentionArchive a finished execution after this long (default: never). Archived executions stay readable.
triggersStart the flow when a message arrives on a subject — see triggers.

An execution completes when it reaches a terminal node and nothing else is in flight. Durations are Go duration strings (30s, 5m, 72h).


Flow definitions

Channels & reducers

Channels are the shared, typed state of an execution. A worker returns Writes — deltas keyed by channel — and each channel folds them with its reducer. Because writes are folded, parallel branches never overwrite each other. Writing an undeclared channel, or a value that does not fit the reducer, fails the attempt.

ReducerBehaviour
replaceThe default: the last write wins.
appendThe channel is a list; a write appends a value (or each element of a list).
mergeThe channel is an object; a write is merged in key by key.
sumThe channel is a number; a write adds to it.

Channels can also be written from outside the graph: synchronously on a running execution with Client.Update, or on a fork with WithForkWrites. Both go through the reducers.

A completed execution returns its channels as an object keyed by channel name (a flow with no channels returns its last node's result). To keep working state such as an agent's memory out of the output, set the flow's output expression: output: "{answer: channels.answer}" returns only answer, and output: "results[last_node]" returns the last node's result even though channels exist. Quote it, since YAML would read a bare {…} as a map. Channels the output leaves out stay in the state, so State, history and forks still see them. An output that is not an object or null fails the execution with reason expression_error.

last_node is the last task, join, map, subflow or await to settle: an await's result is the signal payload, and a node that failed into on_failure counts too. A choice never sets it. So when a flow ends from a choice the default output is the result of whatever ran just before it: in an evaluator-optimizer loop or behind a guardrail that is the verdict, not the work being judged, and after an approval await it is the signal payload. Name the result you want, as in output: "results.generate".


Flow definitions

Expressions & context

Choice rules, map sources, subflow inputs, concurrency keys, search attributes and the flow's output are expr expressions, compiled when the flow is validated. They see the execution context:

VariableContents
inputThe start payload (a JSON object).
channelsThe reduced state channels.
resultsEach settled node's latest output, by node id. A join's result is its branch outputs; a map's, its outputs in item order; a subflow's, the child's output.
last_nodeThe most recently settled node.
visitsHow many times each node was entered.
signalsConsumed signal payloads, by signal name.
branchesBranch status of the last fan-out, by branch id.
countersGeneric usage counters (and steps).
errorsThe last permanent failure of each node: {error, reason, attempt, index?}.
item / indexThe current element of a map (map instances only).

Expressions are restricted to a bounded, straight-line subset: comparisons, arithmetic, boolean logic, membership (in), regex matches, field and index access and len(). Other function calls, ranges and iteration are rejected at validation, so a flow author cannot write an expression that runs unbounded inside the engine.


Flow definitions · node types

task

A task is a job for a worker of the given kind, published on <ns>.work.<kind> with the execution context. Its output becomes results.<node>.

FieldDescription
kindWorker kind. Required.
timeoutPer-attempt deadline. A late result is ignored; the attempt counts as failed and is retried by the policy.
retry{max_attempts, backoff, delay, max_delay}. max_attempts is the total (1 = no retry, at most 64). backoff: fixed (default), linear or exponential; delay defaults to 1s, max_delay to 5m. Backoff waits are durable timers.
output_schemaA JSON Schema the output must match; a mismatch fails the attempt with invalid_output. External $refs are not loaded.
cache{ttl}: reuse a previous result for the same flow version, node and inputs (input, channels, results, signals, item, resume) within the TTL, instead of dispatching a job. A result is stored once its node completes, so runs with the same inputs that start while the first is still working each miss and run the job too: the cache saves repeated work over time, it does not merge concurrent work.
concurrency{key, max}: at most max jobs of this node run at once per value of the key expression, across all executions (e.g. per customer).
dynamicNodes the worker may route to by returning Result.Next (a dynamic edge). Without a next from the worker, next applies.
nextStatic successor.
on_failureHandler node for a permanent failure — see failure routing.
metaFree-form configuration for the worker (any JSON object), handed over as Job.Meta. See below.

Any node may carry meta; the engine ignores it. It is part of the definition, so it changes the flow hash and is returned by Client.Flow(name, version), and each job of a task or map node carries the meta of the version its execution runs: a redeploy never changes it under a running execution, a fork or a rerun. Keep secrets out of it, since it is stored with the flow and sent with every job.

agent loop with dynamic edges
name: agent
max_steps: 30
budget: {calls: 5}
channels:
  messages: {reducer: append}
nodes:
  - {id: plan, type: task, kind: planner, dynamic: [search, calculate, answer]}
  - {id: search,    type: task, kind: tool, next: plan}
  - {id: calculate, type: task, kind: tool, next: plan}
  - {id: answer,    type: task, kind: tool}
start: plan

Flow definitions · node types

choice

Routes on the first rule whose when expression is true; exactly one rule must be the default. Every evaluation is recorded as a ChoiceEvaluated event, so the route taken is part of the history.

yaml
- id: route
  type: choice
  rules:
    - {when: "signals['manager-decision'].approved", to: pay}
    - {when: "results.check.amount > 1000 && visits.check < 3", to: check}
    - {default: true, to: reject}

A rule's to may be $end instead of a node: the execution completes right there, with the flow's output. $end is not a node, so it is never entered and does not count towards visits or max_steps. In a subflow's child it ends the child; the parent goes on. This is how a loop such as evaluator-optimizer finishes without a no-op task:

yaml
start: generate
output: "results.generate"
nodes:
  - {id: generate, type: task, kind: writer, next: evaluate}
  - {id: evaluate, type: task, kind: judge, next: check}
  - id: check
    type: choice
    rules:
      - {when: "results.evaluate.pass || visits.generate >= 3", to: $end}
      - {default: true, to: generate}

on_error decides what an expression error (typically a missing field) does. By default the rule simply does not match and evaluation moves on; with on_error: fail the execution fails with expression_error.


Flow definitions · node types

fanout & join

A fanout starts its branches in parallel; its next must be the join that closes it. Each branch is a task or a subflow node with no successor of its own — the join owns what happens after the fan. A subflow branch runs any graph (a chain, an await, a nested fan-out) as a child execution and settles with the child. At most 256 branches; use a map for wider or data-driven fan-outs.

join fieldDescription
policyall (default): every branch must succeed. any: the first success settles it. quorum:N: N successes. Once settled, branches still running are cancelled — a subflow branch's child too, unless its on_parent_close is abandon: then it carries on, unlinked from the join.
wait_forThe branches the policy counts (default: all branches of the fanout).
nextSuccessor once the policy is met.
on_failureHandler when the policy can no longer be met; without it the execution fails with join_failed.

After the join, results.<join> holds the branch outputs and branches their status. Branches typically contribute to shared state through reducer channels.


Flow definitions · node types

await

Waits for a named signal. The timeout is mandatory: a durable timer, so an execution parked for days holds no engine resources. On timeout the execution follows on_timeout, or fails with await_timeout when there is none. The consumed payload is available as signals.<name>.

yaml
- id: approval
  type: await
  signal: manager-decision
  timeout: 72h
  on_timeout: escalate
  next: route

A signal sent before the execution reaches the await is buffered, never lost. Signals carry an id (WithSignalID) and delivering the same id twice is a no-op.


Flow definitions · node types

map

Runs one task of kind per element of the list the over expression returns, computed at run time. Each instance sees its element as item and its position as index (j.Item(&v), j.Context.Index in Go). results.<map> is the list of outputs in item order.

With flow instead of kind, each item runs that flow as a child execution (LangGraph's Send to a subgraph). The child's input is the item itself when it is an object, or what the input expression returns — it sees item and index besides the parent's context, e.g. input: "{'name': item, 'slot': index}". on_parent_close applies to every item's child; timeout, retry and output_schema belong on the child flow's nodes.

wordcount.yaml
name: wordcount
channels:
  total:   {reducer: sum}
  longest: {reducer: append}
nodes:
  - {id: count, type: map, kind: counter, over: input.docs, max_parallel: 4}

max_parallel bounds the instances in flight (default 100, at most 256). timeout, retry and output_schema apply to each instance. When an item fails permanently (or its child does not complete) the map is aborted: items in flight are cancelled (children too, unless abandoned) and on_failure (if set) is entered with errors.<map>.index naming the item.


Flow definitions · node types

subflow

Runs another registered flow as a child execution with its own id, log and state. The child's output becomes results.<node> and its usage counters are added to the parent's.

FieldDescription
flowThe child flow name. Required.
inputExpression for the child's input (default: the parent's input).
on_parent_closecancel (default): cancelling or failing the parent cancels a running child. abandon: the child keeps running.
next / on_failureSuccessor, and handler when the child does not complete (otherwise child_failed).
yaml
- {id: identity, type: subflow, flow: verify, input: "input.person", next: welcome}

Flow definitions

Failure routing

task, map, subflow and join nodes may declare on_failure: <node>. When the node fails permanently — retries exhausted, a permanent error, an invalid or unencodable output, a last attempt that timed out, a map item that failed, a child that did not complete, a join whose policy was not met — the execution enters the handler instead of failing. This is how sagas and compensation are expressed in the graph.

yaml
- {id: reserve, type: task, kind: inventory, next: charge}
- {id: charge,  type: task, kind: payments,
   retry: {max_attempts: 3, backoff: exponential}, on_failure: release, next: ship}
- {id: release, type: task, kind: inventory}   # undo the reservation
- {id: ship,    type: task, kind: shipping}

The handler sees last_node = the failed node and errors.<node> = {error, reason, attempt, index?}; a later success of that node clears its error. Execution-level stops — budget, max_steps, cancel — never route. Fan-out branches cannot declare on_failure: their join's policy settles them.


Flow definitions

Schedules & triggers

Cron schedules

A schedule starts a flow on a cron expression: six fields (second minute hour day-of-month month day-of-week) or a predefined form such as @hourly or @every 30s. Each firing starts execution <schedule>-<sequence>. Schedules are JetStream message schedules, so they survive restarts and fire once however many engines run.

go
// installed at Init, validated at New
packtrail.WithSchedule("nightly", "report", "0 0 2 * * *", map[string]any{"scope": "all"})

// or at run time
c.Schedule(ctx, "nightly", "report", "0 0 2 * * *", input,
    packtrail.ScheduleTimeZone("Europe/Rome"), packtrail.ScheduleVersion(hash))
c.Unschedule(ctx, "nightly")

Message triggers

A flow's triggers start it for every message on a subject; the message body is the input. With stream set, the trigger is a durable JetStream consumer of that stream (at-least-once); without it, a core NATS queue subscription (at-most-once). The execution id is derived from the flow and the message's Nats-Msg-Id (packtrail.TriggerExecID(flow, msgID), <flow>-t<digest>), so a duplicate message starts nothing twice and several flows triggered by the same message each run once; without a Nats-Msg-Id it is derived from the stream and sequence (stream) or fresh (core). A stream message that still cannot be started after its redeliveries is dead-lettered (trigger). Create the stream yourself: until it exists the trigger keeps retrying and Engine.Ready() stays open.

yaml
name: order
triggers:
  - {subject: orders.created, stream: ORDERS}
nodes:
  - {id: fulfil, type: task, kind: work}

Running

Workers

A worker serves one kind. Every process of a kind shares one durable consumer; each process fetches no more jobs than its free slots (WithConcurrency, default 4). Long jobs are heartbeated automatically. A panic in the handler is recovered and reported as a retryable failure — it never takes the process down.

worker.go
w, err := worker.New(nc, "source", handle,
    worker.WithNamespace("acme"), worker.WithConcurrency(16))
if err != nil { log.Fatal(err) }
log.Fatal(w.Run(ctx)) // drains in-flight jobs when ctx ends

func handle(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)       // fail now, no retry
    }
    _ = j.Progress(map[string]any{"step": "fetching"}) // optional, never stored

    hits, err := search(ctx, in.Topic)
    if err != nil {
        return nil, err                         // retried by the node's policy
    }
    return &worker.Result{
        Output: map[string]any{"hits": len(hits)},    // results.<node>
        Writes: map[string]any{"notes": hits[0]},     // channel deltas
        Usage:  map[string]float64{"cost": 0.4},    // budget counters
    }, nil
}

The job

MemberDescription
ExecID, Flow, Node, KindWhere the job comes from.
Attempt, GenerationAttempt number (from 1) and the instance generation. A redelivered job keeps both.
DeliveriesHow many times this message was delivered (redelivery after a crash).
ContextThe raw execution view (context variables).
Input, Channel, Result, Signal, ItemDecode a part of the context into a Go value.
Resumed(&v)Whether this run resumes an interrupt, and the resume value.
Interrupted(&v)On a resumed run (every attempt), the payload this node interrupted with.
Progress(v)Publish an intermediate result on <ns>.progress.<exec>.<node>.
Store()The long-term store of the worker's namespace, the same one Client.Store() reads: Put, Get, Delete, Keys. Writes are outside the log, so keep them idempotent under redelivery.
TraceparentW3C trace context propagated from the start command.
Meta, DecodeMeta(&v)The node's meta from the execution's flow version (raw JSON, nil when absent), and a decoder for it.

Outcomes

ReturnEffect
*Result, nilComplete: Output (a JSON object or nil), Writes, Next (a declared dynamic successor), Usage. A result that cannot be encoded as JSON (e.g. NaN in Usage) fails the task permanently.
nil, errFail the attempt; retried by the node's policy.
worker.Permanent(err)Fail without retry; on_failure or execution failure follows.
worker.Interrupt(payload)Pause at this node until Client.Resume — see human in the loop.
worker.WithUsage(err, usage)Any of the errors above, plus the usage the run spent: added to the counters like Result.Usage, for every attempt. It composes with Permanent, Interrupt and %w in any order (the outermost counts). A counter past its budget fails the execution at this node: no retry, no pause. A run that never returns (cancelled, timed out) reports nothing.

Cancellation

When an execution ends, a task is cancelled or an attempt times out while a job runs, its context is cancelled with cause worker.ErrCancelled (check context.Cause(ctx)). Return promptly: the result is dropped. Stop notices are best effort; a periodic check of the execution's log head (worker.WithLiveCheck) covers a lost notice.

Workers in other languages implement the protocol; PT_CONFORMANCE_WORKER="<command>" runs the conformance suite against them.


Running

The client

eng.Client() returns a client bound to the engine's namespace; from another process use packtrail.NewClient(nc, packtrail.WithClientNamespace("acme")). It is safe for concurrent use and performs no I/O until the first call.

MethodDescription
Start(flow, input, …)Start an execution; returns its id. WithExecutionID makes it idempotent on that id, WithVersion pins a flow version, WithTraceparent propagates a trace.
Signal(exec, name, payload)Deliver a signal (buffered if no await is waiting yet); WithSignalID deduplicates.
Update(exec, writes)Write channels synchronously and return the state right after.
Resume(exec, node, value)Re-run an interrupted node with a value.
Cancel(exec, reason)Cancel; a no-op on a finished execution.
Get(exec)Current state (archived executions are read from the archive).
Wait(exec)Block until the execution finishes and return the final state. Only ctx ends it early.
Watch(exec, from)Stream events from a sequence.
WatchTerminal(from)Stream the end of every execution (status and sequence), from now or a sequence.
Progress(exec)Stream workers' progress messages from now on.
List(filter)Indexed executions by status, flow or search attribute, newest first.
History, StateAt, OutputHistorySee time travel.
Fork, RerunSee time travel.
Register, Flows, FlowStore a definition (new version when changed), list versions, read one.
Schedule, Unschedule, SchedulesManage cron schedules.
DeadLetters, RedriveInspect and re-publish dead letters (a trigger one restarts only its own flow).
Quarantined, UnquarantineSee quarantine.
Store()A namespaced key-value store for application data that outlives executions: Put, Get, Delete, Keys.

State carries Status, Input, Channels, Results, LastNode, Visits, Steps, Signals, Counters, Branches, Errors, the work in flight, and on a finished execution Output, Error, Reason and FailedNode.


Running

Human in the loop

Three mechanisms, for three situations:

MechanismUse it when
await + SignalThe graph knows it has to wait for an external event (an approval, a webhook). Durable timeout and route.
Interrupt + ResumeA worker discovers mid-task that it needs input. It returns worker.Interrupt(question); the execution is waiting; Client.Resume(exec, node, answer) runs the node again, where j.Resumed(&answer) is true and j.Interrupted(&question) returns the question, on every attempt.
UpdateSomeone outside the graph corrects the state of a running execution. Validated against the channels, applied once per update id, synchronous.
go
// worker side
var answer struct{ Approved bool }
if ok, _ := j.Resumed(&answer); !ok {
    return nil, worker.Interrupt(map[string]any{"question": "refund 1200€?"})
}

// operator side
_ = c.Resume(ctx, id, "refund", map[string]any{"approved": true})
st, err := c.Update(ctx, id, map[string]any{"notes": "checked by ops"},
    packtrail.WithUpdateID("ticket-4711"))

Running

Time travel & forks

The events of one decision are stored as one message, so a stream sequence identifies a decision and every stored message is a consistent cut. All of these read the log; none of them touch the source execution.

CallDescription
History(exec)Every event in order, including archived segments of long executions.
StateAt(exec, seq)The state right after the decision at seq.
OutputHistory(exec, node)Every output the node produced, oldest first — one per visit.
Fork(exec, seq, …)A new execution from the state at seq, continuing from there: what was in flight is dispatched again. WithForkWrites edits its channels first; WithForkID names it.
Rerun(exec, node, …)Fork at the point the node was last entered, so it (and what follows) runs again — the usual move after fixing a worker bug. The node must still be in flight when the decision that entered it ends (a task, map, subflow, fan-out or waiting await): a node that settled right away (a choice, a join, an await whose signal was already buffered) is refused with ErrInvalidArgument; rerun the node before it.
go
evs, _ := c.History(ctx, id)
past, _ := c.StateAt(ctx, id, 12)
fork, _ := c.Fork(ctx, id, 12, packtrail.WithForkWrites(map[string]any{"notes": "try again"}))
again, _ := c.Rerun(ctx, id, "draft")

Running

Statuses & failures

StatusMeaning
runningWork is in flight or about to be.
waitingOnly awaits or interrupted tasks are open: nothing will happen without a signal, a resume or a timer.
completedReached a terminal node; Output is set.
failedError, Reason and FailedNode say why.
cancelledCancelled by a client, or by its parent.
ReasonCause
errorA task failed permanently or ran out of attempts.
timeoutThe last attempt timed out.
invalid_outputThe output did not match output_schema.
await_timeoutAn await timed out with no on_timeout.
join_failedA join's policy could not be met.
child_failedA subflow's child did not complete.
expression_errorA choice rule failed with on_error: fail, a map's over or a subflow or map input could not be built, or the flow's output failed or was not an object.
budget_exceededA counter passed its budget.
max_stepsThe recursion limit was reached.
delivery_failedA job exhausted its deliveries (worker crashes); it is dead-lettered.
unknown_flowA subflow names a flow that is not registered.
decision_too_largeOne decision would exceed 1 000 events.

Node failures — error, timeout, invalid_output, delivery_failed, join_failed, child_failed — route to on_failure when the failed node declares one. Budget, max_steps, an await timeout without a route, an expression error and an oversized decision always end the execution.


Operations

Engine options

OptionDescription
WithNamespace(ns)Prefix of every NATS resource (default packtrail); independent deployments can share a cluster.
WithFlow, WithFlowYAML, WithFlowsDirFlows registered at Init (a struct, YAML bytes, every *.yaml/*.yml of a directory).
WithSchedule(name, flow, cron, input)A cron schedule installed at Init.
WithPartitions(n)Partition count of a new namespace (default 64, at most 1024). Permanent.
WithOwnedPartitions(ps…)Restrict this engine to some partitions (default: compete for all; one is active per partition).
WithReplicas(n)Replication of every stream, bucket and object store. Use 3 in production.
WithoutDispatcher, WithoutCommandsRun only one role, to scale them independently.
WithHistoryLimit(n)Continue as new past n live events (default 10 000; negative disables).
WithSnapshotEvery(n)Snapshot states every n events (default 100; negative disables).
WithStateCacheSize(n)In-memory state cache entries (default 10 000).
WithDrainTimeout(d)Graceful drain budget on shutdown (default 30s).
WithAckWait(d)How long a message held by a crashed engine waits before another gets it (default 5s).
WithPullExpiry(d)Pull request lifetime (default 5s); shorter hands partitions over faster.
WithReadTimeout(d)One batched read of an event log (default 5s).
WithBlobThreshold(n), WithBlobTimeout(d)Claim-check: bodies above n bytes go to the object store (default: server max payload minus a margin); transfer timeout (default 2m).
WithLogger, WithClockA *slog.Logger; the decision clock (tests).

packtrail.ValidateOptions(opts...) runs every check New does without a connection, so CI can reject a bad configuration. ValidateNamespace, ValidateName and ValidateCron expose the identifier and cron rules.

Engine.Ready() and Worker.Ready() return a channel that closes once Run has provisioned and every consumer of the process (commands, dispatch, cron, triggers; a worker's jobs and stop notices) is pulling: wire it to a readiness probe. It stays open when Run fails; once Run returns, Ready() hands out a new, open channel, so call it at each probe. It reports start-up, not later connectivity.


Operations

Worker options

OptionDescription
WithNamespace(ns)Must match the engine's (default packtrail).
WithConcurrency(n)Parallel jobs in this process (default 4); it never prefetches more.
WithAckWait(d)How long a job may go without a heartbeat before redelivery (default 30s, at least 1s). Kind-wide.
WithMaxDeliver(n)Deliveries before a job is dead-lettered and its node failed (default 5). Kind-wide.
WithMaxAckPending(n)Jobs of this kind in flight across every process (default 10 000). Kind-wide.
WithReconfigure()Apply this worker's kind-wide settings to the shared consumer. Without it the first worker's settings stay.
WithLiveCheck(d)Period of the log-head check that stops jobs whose execution ended.
WithDrainTimeout(d), WithBlobTimeout(d), WithLogger(l)Shutdown drain, claim-check transfer timeout, logger.

Operations

Running in production

Replication

Use WithReplicas(3) on a JetStream cluster: the events stream is the source of truth. Without the option existing resources keep their replication, so an engine started without it never scales a deployment down. Asking for replicas on a standalone server fails at Init.

Partitions are permanent

The partition is part of every event subject, so the count cannot change for a namespace. Pick it up front; to change it, run a new namespace alongside and let the old one drain.

Scaling

Run as many engines as you like: pinned consumers make one active per partition and fail over automatically. WithoutDispatcher / WithoutCommands split the roles. Workers scale per kind, independently.

Quarantine

An execution whose events the dispatcher can never process (deterministic failures: an undecodable event, an unknown flow version, a request the server rejects as bad) is quarantined instead of stalling its partition: packtrail quarantined lists them, and they can still be cancelled. Once the cause is fixed — typically by an upgrade — packtrail unquarantine <exec> replays what was skipped. Transient outages never quarantine: the dispatcher waits and retries forever.

Metrics & alerting

Engine.Metrics(ctx) returns counters (commands, conflicts, events appended, jobs dispatched, cache hits, timers, dead letters, archived), ProjectionLag and DispatchStall: the longest time a partition with pending events has gone without progress. Alert when it exceeds a minute or so.

Long histories

Once an execution's live log holds 10 000 events (WithHistoryLimit), it is continued as new under the same id: the log so far is archived as a segment and replaced by one event carrying the state. Nothing else changes; History, StateAt, Fork and Rerun still see every event.

Retention & the index

A flow's retention archives finished executions to the archive object store; Engine.Archive does it immediately. The visibility index is a projection: Engine.RebuildIndex (packtrail rebuild-index) rebuilds it from the log.

Benchmarks

On a laptop (Ryzen 7 5825U, embedded nats-server, file storage): ≈ 2 100 executions/s of a three-step flow on one engine, ≈ 0.86 ms per step for a single execution, ≈ 1 330 executions/s on a three-node R3 cluster. 10 000 parked executions survive an engine restart and all complete once signalled.


Operations

CLI

packtrail [-server URL] [-ns NAMESPACE] <command>. It is a client of the deployment, plus run to host an engine.

packtrail
# engine
run [-flows DIR] [-partitions N]       # run an engine until interrupted
init [-flows DIR] [-partitions N]      # provision the namespace and register flows
rebuild-index                          # rebuild the visibility index from the log
archive <exec>                         # archive a finished execution now

# flows
validate <file.yaml>...                # validate flow files offline
register <file.yaml>...                # register flow versions
flows                                  # list registered versions

# executions
start <flow> [-input JSON] [-id ID] [-version HASH] [-wait]
signal <exec> <name> [-payload JSON] [-id ID]
update <exec> -writes JSON [-id ID]
resume <exec> <node> [-value JSON]
cancel <exec> [-reason TEXT]
get <exec> [-seq N]                    # state now, or right after event N
history <exec>
watch <exec>                           # stream events until the execution ends
progress <exec>                        # stream workers' progress messages
wait <exec>
fork <exec> <seq> [-id ID] [-writes JSON]
rerun <exec> <node> [-id ID]
list [-status S] [-flow F] [-attr K=V] [-limit N]

# schedules and dead letters
schedule <name> <flow> <cron> [-input JSON] [-tz ZONE]
unschedule <name>
schedules
dlq [-limit N]
redrive <seq>
quarantined
unquarantine <exec>

Operations

packtrail-ui

A debugging dashboard: execution list, the flow graph coloured by node state, the event timeline with the state at any sequence, and actions — signal, resume, cancel, fork, rerun, dead-letter redrive. It needs only NATS (NATS_URL, default nats://localhost:4222).

FlagDescription
-addrListen address (default 127.0.0.1:8088, env PACKTRAIL_UI_ADDR).
-namespaceNamespace selected first (default packtrail, env PACKTRAIL_NAMESPACE).
-namespacesComma-separated allowlist (env PACKTRAIL_UI_NAMESPACES); default: every namespace on the NATS account.
No built-in authentication. The UI can drive executions. It binds to loopback by default; put an authenticating reverse proxy in front before exposing it.

Reference

Protocol

Everything packtrail exchanges is plain NATS + JSON, so clients and workers can be written in any language. Messages carry Pt-Protocol-Version: 1; receivers dead-letter versions they do not know. The Go packages are the reference implementation.

Resources

ResourceKind · subjects / keys
<ns>-eventsStream, one message per decision: <ns>.ev.<partition>.<exec>
<ns>-cmdWork-queue stream with message schedules: <ns>.cmd.<p>.<exec>, <ns>.timer.<exec>.<timer>, <ns>.sched.<name>, <ns>.cron.<name>
<ns>-workWork-queue stream: <ns>.work.<kind>, <ns>.workdelay.<kind>.<id>
<ns>-dlqStream, 30 days: <ns>.dlq.<kind>.<key>
<ns>-snapshots, -flows, -index, -cache, -sem, -storeKV buckets
<ns>-blobs, <ns>-archiveObject stores

The partition of an execution is FNV-1a 32-bit of its id modulo the partition count (stored in the <ns>-cmd stream metadata, key packtrail.partitions). Identifiers that become subject tokens match [A-Za-z0-9_-]{1,128}.

Claim-check. A body larger than the threshold (server max_payload − 16 KiB by default) is stored in <ns>-blobs and the message carries Pt-Blob: <name> with an empty body. Every reader resolves it before decoding.

Commands

Published to <ns>.cmd.<p>.<exec> with Nats-Msg-Id = the command id (deduplication) and Pt-Cmd-Type:

command
{"id": "…", "type": "complete", "exec_id": "e1", "data": { … }, "reply": "…"}
Typedata
startflow, version?, input?, parent? — idempotent per execution id (command id start.<exec>)
forkfrom, seq, writes?
completekey, generation, attempt, output?, writes?, next?, usage?, cache_key?
failkey, generation, attempt, error, retryable, reason?
interruptkey, generation, attempt, payload?
resumenode, key?, value?
signalname, payload? — the command id is the signal id
updatewrites — send with reply; the command id is the update id
cancelreason?

With reply set, the engine answers once the command is decided: {"ok": true, "seq": N} or {"ok": false, "code": "invalid" | "rejected" | "not_found" | "failed", "error": "…"}; the Nats-Msg-Id is then <id>@<reply> so a retry with a new inbox still gets an answer. Commands are idempotent; one that can never apply is dead-lettered.

Jobs and conforming workers

A worker of kind K consumes <ns>.work.K through the shared durable consumer <ns>-work-K (explicit ack). A job body:

job
{"exec_id": "e1", "flow": "f", "flow_hash": "…", "node": "n", "kind": "K",
 "key": "n", "index": 0, "generation": 3, "attempt": 1,
 "context": {"input": {}, "channels": {}, "results": {}, "last_node": "", "visits": {},
             "signals": {}, "branches": {}, "counters": {}, "errors": {},
             "item": …, "index": …, "resume": …, "interrupt": …},
 "reply": "<ns>.cmd.<p>.e1", "cache_key": "…", "concurrency": {"key": "…", "max": 2},
 "traceparent": "…"}

A conforming worker:

  1. Resolves Pt-Blob and checks Pt-Protocol-Version.
  2. If concurrency is set, takes a slot in KV <ns>-sem under the key (JSON {"holders": {"<job id>": "<expiry>"}}, compare-and-swap, a lease of 3 × ack wait renewed while the job runs). When the key is full, one job per key per process waits; others are handed back by republishing a copy to <ns>.workdelay.<kind>.<id> with Nats-Schedule: @at …, Nats-Schedule-Target: <ns>.work.<kind> and Pt-Busy (delays 1, 2, 4, 8 s), then acking the original.
  3. Sends InProgress every ack wait / 3 while it holds the job.
  4. Publishes exactly one command to reply — complete, fail or interrupt — with key, generation, attempt copied from the job and command id w.<exec>.<key>.<generation>.<attempt>, then acks. If publishing fails, it naks.
  5. When deliveries exceed its cap, publishes a non-retryable fail (reason delivery_failed), records a dead letter and terminates the message.

A handler crash is reported as a retryable fail.

Progress and stop notices (optional)

Progress goes to core NATS <ns>.progress.<exec>.<node> as {exec_id, node, key, generation, attempt, seq, time, data}, seq counting from 1 per attempt. It is never stored. Stop notices arrive on <ns>.ctl.stop.<exec> as {"exec_id", "all": true} or {"exec_id", "tasks": [{key, generation, attempt}]}; because they may be lost, a worker should also check the last message of the execution's log while a job runs (a terminal type in Pt-Event-Types, or no message twice in a row, means it ended). A stopped job sends no command and is acked.

Clients, timers, dead letters

  • Read by folding the events (a live log may start with ExecutionContinued, whose state is the base) or from the summary in <ns>-index (x.<exec>; membership keys s.<status>.<exec>, f.<flow>.<exec>, a.<attr>.<sha256(value)[:16]>.<exec>). Flows are in <ns>-flows: <name>.latest → hash, <name>.<hash> → definition.
  • Timers are messages on <ns>.timer.<exec>.<timer> with Nats-Schedule: @at <RFC 3339> targeting the execution's command subject. Cron lives on <ns>.sched.<name> targeting <ns>.cron.<name>; each firing starts <name>-<stream sequence>.
  • Dead letters (cmd, job, trigger, dispatch) have body {kind, key, flow?, reason, deliveries, time, subject, header, body}; redrive re-publishes body to subject without the Nats-* headers, except a trigger dead letter: it becomes the start of its flow with execution id key, so other consumers of the subject do not see the message again.

Reference

Event log

The events of one decision are stored as one message: a JSON array of envelopes, appended with Nats-Expected-Last-Subject-Sequence = the execution's previous decision, so there is a single writer without locks. Header Pt-Event-Types lists the types in order; Pt-Cmd-Id names the command that produced it. A decision holds at most 1 000 events.

envelope
{"type": "NodeCompleted", "v": 1, "time": "…", "i": 8, "cmd_id": "w.e1.t.3.1",
 "trace": "00-…-01", "data": { … }}

i is the 1-based index of the event in its execution (the fold rejects a gap). Readers reject unknown types and versions rather than fold wrongly; types are added, never repurposed.

TypeData
ExecutionStartedexec_id, flow, flow_hash, input, parent?, attrs?
ExecutionForkedexec_id, from, seq, state
ExecutionContinuedstate, segment{object, first_index, last_index, first_seq, last_seq}
ExecutionCompletedoutput
ExecutionFailederror, reason, node?, cancel_children[]
ExecutionCancelledreason, cancel_children[]
NodeEnterednode
NodeScheduledkey, node, kind, index?, generation, attempt, owner?, item?, resume?
NodeCompletedkey, node, generation, attempt, output?, writes?, next?, usage?, cached?, cache_key?
NodeFailedkey, node, generation, attempt, error, reason, will_retry
NodeInterruptedkey, node, generation, payload?
NodeCancelledkey, node, reason
ChoiceEvaluatednode, rule, to, error?
FanoutStartednode, join, branches[]
JoinCompletednode, fanout, succeeded[], failed[], ok, output, cancel_children[]?
AwaitStartednode, signal, timer_id
AwaitTimedOutnode, to?
SignalReceivedname, id, payload?
SignalConsumednode, name, id
ChannelsUpdatedid, writes
MapStartednode, items[], max_parallel
MapCompletednode
MapAbortednode, index, cancel_children[]?
ChildStartednode, child_id, child_n, flow, input, policy, key?, owner?, index?
ChildCompletednode, key?, child_id, status, output?, error?, counters?
TimerScheduledid, n, at, purpose, key?, node?, generation?, attempt?
TimerFiredid

Identity and staleness. Every new visit, resume and fork re-dispatch gets a fresh generation from a per-execution counter; retries keep it and bump attempt. A completion, failure, interrupt or timer that does not match the active instance is stale and produces no events. Timer ids (t<n>) and child ids (<exec>-c<n>) come from the same counter, so re-deriving them is idempotent.


Reference

Errors

The public API returns sentinel errors; test them with errors.Is.

ErrorMeaning
ErrInvalidArgumentWraps every validation error of caller input: ids, names, payloads, flows, options.
ErrNotFoundThe execution does not exist and was never archived.
ErrArchivedThe execution is archived: readable, but it cannot be driven any more.
ErrUnknownFlowThe flow or version is not registered.
ErrTerminalThe execution already finished (e.g. Update on a completed execution).
worker.ErrCancelledCause of a job context cancelled because its work is no longer wanted.

Reference

Testing

Tests run against a real embedded nats-server — no mocks. The acceptance suite injects faults: engine kills, NATS restarts, duplicate and reordered completions.

shell
$ make check      # go test -race + golangci-lint + go vet
$ make examples   # run every example against $NATS_URL
$ PT_CONFORMANCE_WORKER="python worker.py" go test ./internal/conformance/

PT_CONFORMANCE_WORKER runs the protocol conformance suite against a worker written in another language. In your own tests, WithClock injects the decision clock.