Documentation
Multi-agent architectures on NATS, defined in YAML or Go.
Overview
stiggy deploys multi-agent architectures on NATS by gluing two libraries: phero (agents, LLMs, tools, memory, the NATS Agent Protocol) and packtrail (a durable, event-sourced workflow engine on JetStream). It adds crews, supervisors, swarms, human review and replay on a substrate that survives crashes and scales out.
A fleet is defined once, as YAML (the stiggy binary) or in Go (the library), and both build the same model.
- models, tools, knowledge: any LLM provider, tool, memory backend or vector store phero offers, plus MCP servers;
- agents: a model, a role/goal/backstory, tools, memory, knowledge and delegates;
- flows: packtrail flows whose steps run agents, with every packtrail feature: choices, fan-out/join, maps, subflows, waits, retries, budgets, fork, rerun and time travel;
- crews and patterns: task crews and the supervisor, swarm, evaluator-optimizer, debate and plan-execute patterns, expanded into flows.
How it works
- A flow step that runs an agent is a packtrail task whose
metaholds the step (prompt, writes, routes, ...). It is versioned with the flow, so a running execution keeps the configuration it started with. - Each agent is a packtrail worker kind. Every job builds a fresh phero agent, runs it, and returns its output, channel writes and usage.
- Agent-chosen routing,
ask_human, structured output and execution memory live in packtrail state, so they are durable and correct under fork, rerun and time travel. - stiggy never touches NATS directly. See Architecture rules.
Install & build
Requires Go 1.26.4 (the toolchain is fetched automatically) and a NATS server with JetStream enabled.
go install github.com/henomis/stiggy/cmd/stiggy@latest
# from a checkout
make build # bin/stiggy (static, honors GOOS/GOARCH)
make test # all tests with -race
make check # test + examples + lint + vetThe module depends on the released phero/v2 v2.0.0 and packtrail v1.0.0.
Quick start
nats-server -js &Write a fleet:
# fleet.yaml
namespace: newsroom
models:
fast: {provider: ollama, model: gemma4:cloud}
agents:
researcher: {model: fast, role: Researcher, goal: "Find accurate, recent facts."}
writer: {model: fast, role: Writer, goal: "Write short, clear articles."}
crews:
article:
tasks:
- {id: research, agent: researcher, description: "Research {{.input.topic}}."}
- {id: write, agent: writer, description: Write a 200-word article from the research.}Check it, look at what it compiles to, serve it, and start an execution:
stiggy validate fleet.yaml # every problem, with its line
stiggy compile fleet.yaml # the packtrail flows it becomes
stiggy run fleet.yaml & # serve it
stiggy start -ns newsroom article -input '{"topic": "NATS"}' -waitstart prints the execution id; with -wait it then prints the final state. The same fleet in Go is shown under Go library.
Fleet structure
A fleet is one YAML document, or the equivalent Go options. Decoding is strict: an unknown key is an error, reported with its line. stiggy validate reports every problem at once.
version: "1" # optional
namespace: newsroom # packtrail namespace; default "stiggy"
models: {}
embedders: {}
vectorstores: {}
knowledge: {}
tools: {}
agents: {}
remotes: {}
flows: {}
crews: {}
patterns: {}
expose: {}
schedules: {}
Names (keys of every section, node ids) must match [A-Za-z0-9_-]. Agent, owner and exposed names are NATS subject tokens of at most 64 characters.
Environment variables in the file
String values outside flows may use ${VAR}, ${VAR:-default} and $$ (a literal $). A missing variable without a default is an error. Quote a reference inside a {...} mapping (api_key: "${OPENAI_API_KEY}"), where YAML would otherwise read its braces. stiggy validate -no-env skips expansion.
Models
models:
fast:
provider: openai # openai | anthropic | ollama (or a registered one)
model: gpt-4o-mini # provider default when empty; required for ollama
api_key: ${OPENAI_API_KEY}
base_url: https://... # OpenAI-compatible endpoints
temperature: 0.2
max_tokens: 1024
retries: 2 # retry a failed call (exponential back-off)
rate_limit: {requests_per_minute: 60, max_concurrent: 4} # per processollama uses Ollama's OpenAI-compatible endpoint (default http://localhost:11434/v1), for example {provider: ollama, model: gemma4:cloud}. In Go, WithModel(name, llm) registers a ready phero LLM and WithModelSpec the same struct as YAML.
Tools
tools:
shell: {type: bash, allowlist: [ls, cat], timeout: 30s, working_dir: /work}
files: {type: file_read, working_dir: /work} # also file_write, file_edit, glob, grep
skills: {type: skill, root: ./skills}
scratch: {type: kv, bucket: scratchpad} # NATS KV read/write tools
browser: {type: mcp, command: npx, args: [-y, "@playwright/mcp"], allow: [browser_navigate]}
remote: {type: mcp, url: https://mcp.example.com/mcp, deny: [delete_everything]}Built-in types: bash, file_read, file_write, file_edit, glob, grep, skill, kv and mcp. Options are checked strictly per type. An mcp entry is either a stdio command (with args and env; only a minimal environment is passed on) or a streamable-HTTP url, filtered with allow and deny.
Agents
agents:
writer:
model: fast
system: You write for a newsroom. # system prompt
role: Writer # appended to system
goal: Short, accurate articles.
backstory: Ten years at a daily paper.
description: Writes articles. # for agents that delegate to it
tools: [shell, files]
knowledge: [handbook] # one search_<name> tool each
delegates: [researcher, legal] # agents or remotes, as ask_<name> tools
max_iterations: 10 # default 25
stream: true # progress events while it runs
memory: execution # see below
Agents run in-process inside packtrail workers of kind agent-<name>. Each job builds a fresh phero agent.
- Steering tools are added from the step:
route_to_<node>for each route andask_human(see Routing & human in the loop). - Delegates appear as
ask_<name>tools; their usage is charged to the calling step. - Streaming (
stream: true) publishes text and tool calls as progress, readable withstiggy progress.
memory
| Type | Kept in | Notes |
|---|---|---|
none (default) | nowhere | Every step starts fresh |
execution | the execution's state (hidden channel _mem_<agent>) | Correct under fork, rerun and time travel |
simple | the worker process | max_items (default 100) |
jsonfile | <dir>/<session>.json | dir |
nats | NATS KV | bucket |
psql | PostgreSQL | dsn, table, ensure_schema |
rag | a knowledge index | knowledge |
Long-term types (all but none and execution) also take:
session: default the agent's name, shared by every execution;per_execution: true: append the execution id to the session;summarize: {model, threshold, size}.
memory may be a bare type (memory: execution) or a mapping (memory: {type: nats, bucket: crew-memory, per_execution: true}).
Embedders, vector stores & knowledge
embedders:
small: {provider: openai, model: text-embedding-3-small, api_key: "${OPENAI_API_KEY}"} # or ollama
vectorstores:
docs: {type: qdrant, host: localhost, port: 6334, collection: handbook, vector_size: 1536, distance: cosine}
pg: {type: psql, dsn: "${DATABASE_URL}", collection: handbook, vector_size: 1536}
wv: {type: weaviate, host: localhost:8080, class: Handbook}
knowledge:
handbook:
embedder: small
vectorstore: docs
description: The editorial handbook.
top_k: 4
sources:
- {path: docs/*.md, chunk_size: 1000, chunk_overlap: 100} # splitter: markdown | recursiveAn agent lists knowledge: [handbook] and gets a search_handbook tool. The engine role ingests the sources at startup when the store is empty.
Remotes
Agents served by others on the NATS Agent Protocol, found by discovery:
remotes:
legal: {owner: legal-team, name: reviewer, description: Reviews legal risks.}Steps run them with remote: legal (worker kind remote-legal); agents list them in delegates. The protocol is final-answer only, so remote agents do not stream.
Flows & steps
A flow is a packtrail flow: start, channels (with reducers replace, append, merge, sum), output (an expression; default all channels, or the last node's result), budget, max_steps, retention, triggers, and nodes of types task, choice (to: $end ends the flow), fanout, join, await, map and subflow. Refer to packtrail's documentation for everything on those nodes.
A task or map node runs a stiggy step when it sets one of these fields:
| Field | Meaning |
|---|---|
agent | Run a fleet agent |
activity | Run a Go activity (library mode) |
remote | Prompt a remote agent |
prompt | Go template over the job context; default: the input or map item, plus the context outputs |
expected_output | Appended to the prompt |
context | Nodes whose outputs the prompt includes; default the previous node |
writes | {channel: output} or {channel: output.field} |
routes | Nodes the agent may choose next, and human for ask_human (task nodes only, not in fan-outs or maps) |
human_input | A human reviews the answer: {"approve": true} or feedback |
A step runs exactly one of agent, activity and remote. The node's own output_schema, retry, timeout, cache and concurrency still apply. With output_schema the agent answers structured JSON.
A hand-written flow
flows:
greet:
nodes:
- {id: hello, type: task, agent: greeter, prompt: "Greet {{.input.name}}."}
Prompt templates
Templates see .input, .channels, .results, .item, .index, .last_node, .visits, .signals, .errors, .counters and .resume, plus the json function: {{json .channels.notes}}. Missing keys render as zero values.
Without a prompt, the agent receives the Item (in a map) or the Input, followed by the output of each context node, or of the previous node when context is empty. A step with nothing to show gets Begin..
Outputs
An agent step's output is {"text": "<answer>"}, or the parsed object when it has an output_schema. writes copies all of it (output) or a field (output.verdict) into a channel, applying the channel's reducer.
Usage counters and budgets
Budgets in budget (flow, crew or pattern) can use agent_steps, tokens_in, tokens_out and cost_usd. They are charged for every run, including failed and paused ones, and checked on completion, failure and pause. max_steps bounds the number of node visits.
A larger example
agents:
planner: {model: m, knowledge: [handbook], delegates: [critic, legal]}
researcher: {model: m, memory: {type: execution}}
critic: {model: m, stream: true}
flows:
crew:
start: plan
channels:
notes: {reducer: append}
verdict: {}
budget: {tokens_out: 50000, agent_steps: 40}
max_steps: 200
nodes:
- id: plan
type: task
agent: planner
output_schema:
type: object
properties: {topics: {type: array, items: {type: string}}}
required: [topics]
next: research
- id: research
type: map
over: results.plan.topics
max_parallel: 4
agent: researcher
prompt: "Research {{.item}}. Context: {{json .channels.notes}}"
writes: {notes: output}
next: review
- id: review
type: task
agent: critic
context: [plan, research]
expected_output: A verdict and the reason for it.
human_input: true
writes: {verdict: output.verdict}
routes: [plan, human]Routing & human in the loop
Agent-chosen routing
routes: [a, b] gives the agent a route_to_a and a route_to_b tool and makes both nodes dynamic edges of the step. The choice is stored in packtrail state, so it is durable and replayable. Routes work on task nodes only, not inside fan-outs or maps. See examples/router.
Asking a human
Add human to routes to give the agent an ask_human tool, or set human_input: true to have a human review every answer. Either way the step pauses. Its interrupt payload is in the execution state, at State.Tasks[...].Interrupt:
{kind: "question" | "review", agent, node, question | answer}| Kind | Resume with |
|---|---|
question | "text" or {"answer": "..."} |
review | {"approve": true}, or "text" / {"feedback": "..."} to have the agent revise its answer |
stiggy resume -ns newsroom <exec> <node> -value '"answer"'
stiggy resume -ns newsroom <exec> <node> -value '{"approve": true}'In Go, use Client().Resume. Exchanges with the human are kept in the payload (the event log), so forks and reruns see exactly their own history. await nodes with stiggy signal cover waits that are not agent questions. See examples/hitl.
Crews
crews:
launch:
process: sequential # or hierarchical (needs manager)
manager: boss
max_steps: 40
budget: {tokens_out: 50000}
tasks:
- id: research
agent: researcher
description: "Research {{.input.product}}."
expected_output: Three facts.
async: true # runs in parallel with neighbouring async tasks
- id: draft
agent: writer
description: Write the post.
context: [research]
output_schema: {...}
human_input: true
guardrail: {agent: editor, max_retries: 2}
- sequential (default): tasks run in order; consecutive
asynctasks run as a fan-out. - hierarchical: a
manageragent assigns the tasks. - guardrail: a checker agent reviews the answer; on failure the task is redone with the feedback, up to
max_retries.
A crew runs as a flow named after it. stiggy compile shows it. See examples/crew.
Patterns
patterns:
team: {type: supervisor, supervisor: boss, workers: [researcher, writer]}
swarm: {type: swarm, agents: [triage, billing, tech]}
polish: {type: evaluator_optimizer, generator: writer, evaluator: editor, max_rounds: 3}
argue: {type: debate, debaters: [optimist, skeptic], judge: editor, rounds: 2}
project: {type: plan_execute, planner: boss, executor: researcher, synthesizer: writer, max_parallel: 4}| Type | Behavior |
|---|---|
supervisor | The supervisor routes to one worker at a time; every worker reports back until the supervisor answers without routing. Node ids are the agents' names. |
swarm | Every agent may hand off to any other. The first starts; the flow ends with the first agent that answers without handing off. |
evaluator_optimizer | The evaluator judges with structured output {pass, feedback}; the generator regenerates with the feedback until it passes or max_rounds generations ran. |
debate | Debaters argue in parallel for rounds rounds, each seeing the transcript so far; then the judge decides. |
plan_execute | The planner lists steps (structured output {steps}), the executor runs each (max_parallel at a time), and the optional synthesizer combines the results. |
All patterns take prompt, max_steps and budget. Each runs as a flow named after it; stiggy compile shows the graph, and validation errors are mapped back to the pattern.
Expose
expose:
owner: newsroom # default: the namespace
agents: [writer] # each prompt runs a fresh agent
flows: [launch] # each prompt starts an execution, answered with its outputExposed agents and flows are served on the NATS Agent Protocol (phero's nats.Server), so any protocol client can call them. Discover them with stiggy agents. Readiness waits until they are discoverable.
Exposed flows take the prompt as input: a JSON object as is, otherwise {"prompt": text}. A Stiggy-Execution-Id request header makes retries join the same execution.
Schedules
schedules:
morning: {flow: launch, cron: "0 0 7 * * *", input: {product: stiggy}}The cron expression has a leading seconds field. Schedules are registered by the engine role.
Deployment & roles
Every process of a deployment runs the same fleet file and picks its part:
stiggy run -role engine -http :8080 fleet.yaml # one or more engines
stiggy run -role workers -http :8080 fleet.yaml # scale out
stiggy run -role workers -only researcher,writer fleet.yaml # dedicated hosts| Role | Runs |
|---|---|
all (default) | Engine and every worker: one process, typical in development |
engine | The packtrail engine: provisions the namespace, registers flows and schedules, ingests knowledge, drives executions |
workers | The agent, activity and remote workers. Scale it out freely |
-only restricts a process to the named agents, activities and remotes. On shutdown, workers drain before the engine stops.
deploy/ has a Dockerfile (build from the repo root) and a docker-compose setup with NATS, one engine, scalable workers and a CLI profile:
docker build -f deploy/Dockerfile -t stiggy .
OPENAI_API_KEY=... docker compose -f deploy/docker-compose.yaml up --scale workers=3
docker compose -f deploy/docker-compose.yaml run --rm cli \
start -ns newsroom article -input '{"prompt": "NATS"}' -wait
Health endpoints
-http ADDR serves /healthz (200 while the process runs) and /readyz (503 until the app is ready, then 200).
Worker kinds
agent-<name>, act-<name> and remote-<name> belong to stiggy. A plain task node with any other kind is left to external packtrail workers, in any language.
CLI reference
Flags may appear before, between or after positional arguments. Online commands take -server URL (default $NATS_URL, else nats://127.0.0.1:4222) and -creds FILE. Execution commands also take -ns NAMESPACE (default stiggy).
Offline
| Command | Purpose |
|---|---|
validate [-activities a,b] [-no-env] <fleet.yaml>... | Check fleet files; every problem is reported with its line |
compile [-activities a,b] [-no-env] [-format yaml|json] <fleet.yaml> | Print the packtrail flows and the workers the fleet compiles to |
-activities names the Go activities a library program registers, so a fleet written for the library can be checked by the binary.
Online
| Command | Purpose |
|---|---|
run [-role all|engine|workers] [-only a,b] [-http ADDR] <fleet.yaml> | Serve the fleet until interrupted |
start [-ns NS] [-input JSON] [-id ID] [-wait] <flow> | Start an execution and print its id; -id makes the start idempotent; -wait prints the final state |
agents [-owner OWNER] [-agent FRAMEWORK] | List NATS Agent Protocol agents on the bus: exposed stiggy agents and flows, phero agents, any v0.3 agent |
Executions
| Command | Purpose |
|---|---|
get <exec> [-seq N] | State now, or right after event N (time travel) |
history <exec> | The event log, as JSON lines |
progress <exec> | Streamed agent progress, as JSON lines |
signal <exec> <name> [-payload JSON] | Signal an await node |
resume <exec> <node> [-value JSON] | Answer an agent waiting for a human |
update <exec> -writes JSON | Write channels, applied with their reducers (edit state while paused) |
cancel <exec> [-reason TEXT] | Cancel an execution |
fork <exec> <seq> | Branch an execution from a point in its history |
rerun <exec> <node> | Run again from a node |
Executions can also be inspected, signalled, resumed, forked and cancelled with packtrail's own CLI and packtrail-ui on the fleet's namespace.
Go library
github.com/henomis/stiggy builds the same fleet from options. New validates and compiles without I/O; models and tools are built, and NATS used, only by Run.
nc, _ := nats.Connect(nats.DefaultURL)
app, err := stiggy.New(nc,
stiggy.WithNamespace("newsroom"),
stiggy.WithModel("fast", openai.New("ollama",
openai.WithBaseURL(openai.OllamaBaseURL), openai.WithModel("gemma4:cloud"))),
stiggy.WithAgent("researcher", spec.Agent{Model: "fast", Role: "Researcher"}),
stiggy.WithAgent("writer", spec.Agent{Model: "fast", Role: "Writer"}),
stiggy.WithCrew("article", spec.Crew{Tasks: []spec.Task{
{ID: "research", Agent: "researcher", Description: "Research {{.input.topic}}."},
{ID: "write", Agent: "writer", Description: "Write a 200-word article from the research."},
}}),
)
go app.Run(ctx)
<-app.Ready()
id, _ := app.Client().Start(ctx, "article", map[string]any{"topic": "NATS"})
st, _ := app.Client().Wait(ctx, id)
app.Client() is a packtrail client: signal, resume, update, fork, rerun, watch and stream progress from it. A YAML fleet loads with config and passes to WithFleet.
Options
| Option | Purpose |
|---|---|
WithFleet(*spec.Fleet) | Merge a whole fleet (for example one loaded from YAML) |
WithNamespace | Packtrail namespace |
WithModel / WithModelSpec | A phero LLM instance, or a spec resolved by the registry |
WithTool / WithToolSpec | Ready *llm.Tool values, or a spec |
WithAgent | An agent (spec.Agent) |
WithEmbedder, WithVectorStore, WithKnowledge | RAG building blocks |
WithCrew, WithPattern | Crews and patterns, expanded into flows |
WithFlow(*flow.Flow) | A packtrail flow; build agent nodes with AgentTask |
WithActivity(name, handler) | A Go function run as a step (worker kind act-<name>) |
WithRemote, WithExpose, WithSchedule | Remote agents, exposure and cron schedules |
WithRegistry | Register custom LLM providers, tool types, memory backends, embedders, vector stores |
WithRole, WithOnly | RoleAll, RoleEngine, RoleWorkers; restrict the workers |
WithLogger | A *slog.Logger |
WithEngineOptions, WithWorkerOptions | Pass packtrail options through |
Defining a name twice is an error (ErrOption). AgentTask(id, agent, step) and ActivityTask(id, activity, step) build task nodes for flows written in Go. App.Plan() returns the compiled plan.
Testing
stiggytest.Start(t, opts...) boots an embedded NATS server, runs the app and returns a harness with the client. Fake models (Reply, Replies, Echo) script agent answers deterministically, so end-to-end tests need no API key.
Examples
Each example runs against $NATS_URL, with a real model when OPENAI_API_KEY, ANTHROPIC_API_KEY or OLLAMA_MODEL is set, and a scripted one otherwise.
| Example | Shows |
|---|---|
examples/hello | The smallest YAML fleet |
examples/router | Agent-chosen routing, built with the Go API |
examples/mapreduce | Structured output, a parallel map, an append channel |
examples/hitl | ask_human and human review, resumed with Client.Resume |
examples/crew | A crew with a guardrail, also served as a protocol agent |
examples/remote | Another team's phero agent as a step and as a delegate |
go run ./examples/router
OLLAMA_MODEL=gemma4:cloud go run ./examples/mapreduce
make examples # scripted models, embedded server
make examples-live # a real model, default OLLAMA_MODEL=gemma4:cloudArchitecture rules
No direct NATS access
NATS is the only transport, and stiggy's own state lives in NATS (through packtrail). Agent storage you configure may be any backend phero offers. stiggy reaches NATS only through the public Go APIs of phero and packtrail:
- it dials NATS only in
internal/natsconnand passes the connection to phero and packtrail constructors; - it never imports
nats.go/jetstream,nats.go/microornats-server, never publishes or subscribes on a subject, and never reads or writes phero's or packtrail's subjects, buckets or wire formats.
internal/archguard enforces this with a type-based check over the whole module, tests included (make guard).
phero and packtrail are read-only
When a feature needs something they don't expose, stiggy leaves it out or makes the compiler reject it with a clear error.
Layout
| Package | Role |
|---|---|
spec/ | The fleet model and ApplyStep |
config/ | YAML loader: strict decoding, ${ENV}, error lines |
compile/ | Pure Compile(fleet, catalog) reporting every problem |
registry/ | Providers, tool types, memory backends, embedders, vector stores, activities |
agentrun/ | The worker handler that runs a phero agent per job |
patterns/ | Crews and patterns as macros that expand into flows |
stiggytest/ | Embedded NATS, the app and fake LLMs |
cmd/stiggy/ | The CLI |
Licensed under Apache 2.0.