Skip to main content

agent-workflow-sdk

An async-first, LangGraph-style SDK for building agent workflows from composable pieces. You describe a workflow as a graph of nodes over a typed, reducer-based state, and run it against a pluggable backend — a coding agent (Kiro, Codex, or Claude Code) or a plain LLM (Ollama, any OpenAI-compatible endpoint, or Anthropic).

The import root is agentflow. The distribution name is agent-workflow-sdk.

Why

  • Build any workflow from primitives: nodes, edges, conditional routing, and a shared state. Not a fixed loop, not a linear chain.
  • Swap the backend without touching workflow code. Agents and LLMs share one event vocabulary.
  • Async everywhere: the engine, backends, checkpointing, and streaming.
  • Durable by default: every super-step is checkpointed, runs resume after a crash, and human-in-the-loop interrupts suspend and resume a run.

Install

The core is dependency-free. Backends that need extra libraries are optional extras.

pip install -e .            # core only
pip install -e '.[ollama]'  # + httpx, for the Ollama LLM backend
pip install -e '.[dev]'     # + pytest, pytest-asyncio

Python 3.11+ is required. The Kiro backend needs kiro-cli on PATH.

Quickstart

A graph that loops until a counter reaches a target, then stops.

import asyncio
from typing import Annotated
from agentflow import Graph, START, END, State, add, append

class CountState(State):
    n: Annotated[int, add]        # updates are summed
    log: Annotated[list, append]  # updates are concatenated

async def tick(state, ctx):
    return {"n": 1, "log": f"tick {state.get('n', 0) + 1}"}

def route(state):
    return "again" if state["n"] < 5 else "done"

g = Graph(CountState)
g.add_node("tick", tick)
g.add_edge(START, "tick")
g.add_conditional_edges("tick", route, {"again": "tick", "done": END})
app = g.compile()

print(asyncio.run(app.invoke({"n": 0, "log": []})))

See examples/hello_graph.py and examples/optimize_loop.py for runnable versions.

Core concepts

State and reducers

State is a TypedDict subclass of State. Each field is a channel; annotate it with a reducer that folds each node's update into the current value. A field without a reducer uses last (last-value-wins). Built-in reducers: last, append, add, merge, union. A node returns a partial update; it never mutates the state it was given.

Nodes and edges

A node is an async callable (state, ctx) -> update. Edges are either static (add_edge) or conditional (add_conditional_edges with a router(state) -> key). START and END are sentinels. A conditional router may return a list of keys to fan out to several nodes at once.

Execution

A run proceeds in super-steps. All nodes on the current frontier run concurrently against the same immutable state, their updates are folded through the reducers in a deterministic order, and the next frontier is computed from the edges. A step_limit guards against runaway loops.

Super-steps are all-or-nothing: if any node in a step raises, that step's updates are discarded before they touch the state and no checkpoint is written for it, so the run stays at the last committed checkpoint. Input to a run is validated at start — an undeclared channel key is rejected with a clear error rather than sitting unused in the state.

CompiledGraph gives you:

  • await app.invoke(input, thread=..., timeout=...) — run to completion, return state. timeout bounds the whole run; on expiry it raises RunTimeout and the last checkpoint is preserved for resume.
  • app.stream(input, thread=...) — async-iterate StreamEvents as they occur.
  • await app.resume(thread, value=...) — continue a suspended run.
  • await app.get_state(thread) / app.history(thread) — inspect checkpoints.

For non-async callers there are blocking wrappers — app.invoke_sync(...), app.resume_sync(...), app.stream_sync(...) — which run the coroutine via asyncio.run and refuse to run inside an existing event loop.

Subgraphs

A compiled graph composes as a node in another graph:

parent.add_node("sub", compiled_subgraph)              # channels shared by name
parent.add_subgraph("sub", compiled_subgraph,          # or map explicitly
                    input_map={"query": "q"}, output_map={"answer": "result"})

A subgraph runs on an isolated checkpoint sub-thread and merges its result back through the parent's reducers. Use output_map to route a result into a distinct parent channel and avoid double-counting accumulator channels.

Observability

Pass a Hooks implementation to compile(hooks=...) for lifecycle callbacks (on_run_start, on_node_start/end/error, on_step_end, on_run_end). A hook raising never breaks a run. RunMetrics is a ready-made Hooks that records per-node call counts, errors, and durations:

from agentflow import RunMetrics
m = RunMetrics()
await g.compile(hooks=m).invoke({...})
print(m.summary())   # {"steps": 3, "completed": True, "nodes": {...}}

For durable telemetry, JsonlTelemetry is a Hooks that writes one JSON line per event (run_start, node_start/end/error, event, step, run_end) — backend events surfaced via ctx.emit are logged too. Compose several listeners with MultiHooks; a failing listener never breaks the run or the others:

from agentflow import MultiHooks, RunMetrics, JsonlTelemetry
metrics = RunMetrics()
telemetry = JsonlTelemetry(".runs")           # dir -> one file per thread
app = g.compile(hooks=MultiHooks(metrics, telemetry))

See examples/telemetry_demo.py for a runnable version.

For distributed tracing, agentflow.otel.OtelHooks (install the otel extra) is a Hooks that emits an OpenTelemetry span per run and per node:

from agentflow.otel import OtelHooks
app = g.compile(hooks=OtelHooks())   # uses the global tracer/provider

For operators who scrape Prometheus instead of running an OTel collector, agentflow.prometheus.PrometheusHooks (install the prometheus extra) is a Hooks that records run/node/step counters, a node-duration histogram, and backend-event counts. It uses a private registry by default; call exposition() to render the text format for a /metrics endpoint (you own the HTTP layer).

from agentflow.prometheus import PrometheusHooks

metrics = PrometheusHooks()
app = g.compile(hooks=metrics)
# ... inside your /metrics handler:
body, content_type = metrics.exposition()

Compose several listeners with MultiHooks(RunMetrics(), OtelHooks(), PrometheusHooks()).

Backends

Two kinds of backend share one event stream, so a node calls either the same way.

  • Agent backends (AgentBackend) drive a session that runs its own tools and asks permission. Available: KiroBackend (persistent ACP session over kiro-cli), CodexBackend (codex exec --json), ClaudeCodeBackend (claude -p --output-format stream-json).
  • LLM backends (LLMBackend) are stateless: messages in, token stream out. They never run tools themselves — a tool call is a request the graph fulfils. Available: OllamaBackend (local /api/chat), OpenAIBackend (any /v1/chat/completions endpoint — OpenAI, Groq, Together, vLLM, LM Studio, and Ollama's own /v1 shim), and AnthropicBackend (Messages API).
from agentflow.backends.kiro import KiroBackend
from agentflow.backends.codex import CodexBackend
from agentflow.backends.claude_code import ClaudeCodeBackend
from agentflow.backends.ollama import OllamaBackend
from agentflow.backends.openai import OpenAIBackend
from agentflow.backends.anthropic import AnthropicBackend

agent = KiroBackend("vibe", engine="v3")          # persistent session
codex = CodexBackend(sandbox="read-only")          # one-shot per turn
claude = ClaudeCodeBackend(model="sonnet")         # one-shot per turn
llm = OllamaBackend("llama3.2")                     # local HTTP
openai = OpenAIBackend("gpt-4o-mini", api_key="...")  # any OpenAI-compatible API
anthropic = AnthropicBackend("claude-sonnet-4", api_key="...")  # Messages API

await agent.start()
async for event in agent.prompt("summarize the repo"):
    ...  # TextChunk, ToolCall, ToolResult, ..., TurnEnd
await agent.close()

The three agent backends have different lifecycles under one interface. Kiro holds a long-lived JSON-RPC session; Codex and Claude Code run a fresh subprocess per turn and carry a resumable session id between turns. Either way you call start(), prompt(text), close() and consume the same event stream. Backends are constructed by you and passed into your nodes. The core never imports a backend, so importing agentflow pulls in no subprocess or HTTP dependency.

The HTTP LLM backends retry transient failures (connection errors, timeouts, 429/5xx) on the connect/initial-response phase only — never mid-stream, so a partial token stream is never replayed. Retry-After is honored. Tune it with a RetryPolicy:

from agentflow.backends import RetryPolicy
from agentflow.backends.openai import OpenAIBackend

llm = OpenAIBackend("gpt-4o-mini", api_key="...",
                    retry=RetryPolicy(max_retries=4, backoff=0.5))

Run python examples/agents_demo.py to smoke every backend installed on your machine (set KIRO_AGENT to a valid agent id, e.g. vibe, to include Kiro).

Permission policies

A PermissionPolicy is required when constructing an agent backend — there is no default, because auto-approving an agent's tool use is a security choice the caller must make explicitly. Options: AllowAll (trusted local sandbox only), DenyAll, ToolAllowlist({"read_file", ...}, fallback=DenyAll()) for least-privilege scoping, and two ways to ask a human — Interactive (suspends the run via an engine interrupt; resume with a decision) and Callback(async fn) (asks inline within the same turn, no interrupt — the right choice for a persistent-session backend like Kiro).

from agentflow.backends import AllowAll, ToolAllowlist, DenyAll
KiroBackend("vibe", permission=ToolAllowlist({"read_file"}, fallback=DenyAll()))

Capability matrix — how the policy is applied:

Backend How PermissionPolicy applies
KiroBackend Runtime gate — ACP session/request_permission is resolved by the policy per tool call (Interactive/Callback ask a human; ToolAllowlist filters).
CodexBackend Translated to a launch flag — the policy sets Codex's --sandbox when you don't set it yourself: AllowAll → workspace-write, DenyAll/ToolAllowlist → read-only. No per-tool runtime gate.
ClaudeCodeBackend Translated to launch flags — ToolAllowlist → --allowed-tools, DenyAll → --permission-mode plan (read-only), when you don't set those yourself. No per-tool runtime gate.

For the one-shot CLI agents the policy is a best-effort upfront mapping to their own controls, not our runtime interception — an explicit sandbox= / allowed_tools= you pass always wins over the derived value.

Checkpointing and human-in-the-loop

Pass a checkpointer to compile to make runs durable.

from agentflow import FileCheckpointer
app = g.compile(checkpointer=FileCheckpointer(".runs"))

MemoryCheckpointer is for tests; FileCheckpointer writes one atomic JSON file per super-step under .runs/<thread>/ (single-writer / single-process); SqliteCheckpointer stores checkpoints transactionally with atomic per-step revisions; RedisCheckpointer (install the redis extra) persists to a Redis server for distributed / multi-process runners; and PostgresCheckpointer (install the postgres extra) persists to a Postgres table. All implement the same Checkpointer protocol, so they are interchangeable at compile(checkpointer=...).

from agentflow.checkpoint import PostgresCheckpointer
app = g.compile(checkpointer=PostgresCheckpointer("postgresql://localhost/app"))

For production durability prefer PostgresCheckpointer: a committed row survives a crash by design. Redis is only crash-durable when you configure AOF/RDB persistence.

Optimistic concurrency and retention

Each checkpoint carries a revision (the write count for its (thread, step), set when you read it back). To guard a read-modify-write against a concurrent resume, pass the revision you read as if_revision; a stale value raises CheckpointConflict instead of silently overwriting.

cp = await checkpointer.get(thread, step)
# ... resume work off cp ...
await checkpointer.put(new_cp, if_revision=cp.revision)  # raises on conflict

Two retention methods trim old state: delete_thread(thread) drops a whole thread, and prune(thread, before_step=..., older_than=...) deletes older checkpoints and returns how many it removed.

A node calls await ctx.interrupt(payload) to suspend the run for a human. The runtime writes an interrupted checkpoint and stops. Later, await app.resume(thread, value=answer) continues the run, and the same interrupt call returns answer. This works across process restarts, since the frontier is persisted.

Cross-thread memory (Store)

A checkpointer persists one thread's execution state so it can resume. A Store is the other axis: durable key-value data shared across threads such as user profiles, learned facts, or long-term agent memory. Items live under a namespace (a tuple of path segments) and a string key; the value is any JSON-serializable object.

from agentflow import MemoryStore

store = MemoryStore()
await store.put(("users", "u1", "memories"), "favorite_color", "blue")
item = await store.get(("users", "u1", "memories"), "favorite_color")
recent = await store.search(("users", "u1"))  # everything under that prefix

MemoryStore is ephemeral (tests, single process); PostgresStore (install the postgres extra) is durable and shared across processes. Both implement the same Store protocol. put(..., ttl=seconds) expires an item; expired items never surface from get or search.

from agentflow.store import PostgresStore

store = PostgresStore("postgresql://localhost/app")
await store.put(("cache",), "doc-42", payload, ttl=3600)

Control plane (run queue + workers)

invoke/stream/resume run a graph inline in your code. The control plane lets you manage runs from outside that process: enqueue a run, let a pool of workers execute it, and list, cancel, or resume it from anywhere sharing the same queue. Runs survive a process restart and scale across workers.

A run request references a graph by name (graphs are code, not data), so an enqueuer and its workers share a GraphRegistry mapping names to compiled graphs, and the same checkpointer so workers see the run's state.

from agentflow import GraphRegistry, MemoryRunQueue, Worker

registry = GraphRegistry()
registry.register("summarize", lambda: build_graph().compile(checkpointer=cp))

queue = MemoryRunQueue()
run = await queue.enqueue("summarize", {"text": "..."})

# A worker (usually a separate long-running process) drains the queue.
worker = Worker(queue, registry)
await worker.run_once()          # or: await worker.run_forever()

record = await queue.get(run.run_id)   # status: succeeded / interrupted / ...

Run lifecycle: queued -> running -> succeeded | interrupted | failed | cancelled. An interrupted run (a graph that hit ctx.interrupt) is resumed by re-queuing it with the human answer:

await queue.enqueue_resume(run.run_id, value="approved")  # back to queued

Cancellation is cooperative — await queue.request_cancel(run_id) stops a running graph at the next super-step boundary, leaving a consistent checkpoint.

MemoryRunQueue is for a single process/tests. PostgresRunQueue (install the postgres extra) is durable and safe for many concurrent workers: claim hands each run to exactly one worker via FOR UPDATE SKIP LOCKED, and a time-bounded lease means a crashed worker's run is re-claimed once the lease expires.

For monitoring, await queue.stats() returns a QueueStats snapshot (depth by status plus expired_leases — running runs whose lease is past due, i.e. work stranded by a crashed worker awaiting re-claim), and pool.health() returns a PoolHealth snapshot (workers/alive/busy/idle, and healthy when every worker loop is alive). Both are plain values you can expose through your own /healthz and /metrics handlers.

stats = await queue.stats()          # stats.queued, stats.running, stats.expired_leases
if not pool.health().healthy:
    ...                              # a worker loop died — pool is degraded

An HTTP layer (each queue method maps 1:1 to an endpoint) is planned and not part of this release.

Prebuilt patterns

agentflow.prebuilt.iterate_until_converged(work, ...) compiles the classic baseline-then-iterate optimization loop as a graph. You supply an async work(state) -> Candidate; the loop owns the best-so-far, the no-improvement streak, and the stop decision (done, optimal, converged, exhausted).

from agentflow.prebuilt import Candidate, iterate_until_converged

async def work(state):
    best = state.get("best_score", 0.0)
    return Candidate(score=min(best + 25.0, 100.0))

app = iterate_until_converged(work, perfect_score=100.0, patience=4)
result = await app.invoke({})

agentflow.prebuilt.tool_loop(llm, tools) compiles the function-calling agent loop as a two-node graph: an agent node calls the LLMBackend, a tools node runs any tool the model requested and appends the result to the conversation, and a conditional edge repeats until the model answers without a tool call. The LLM never executes a tool itself — the graph does.

from agentflow.prebuilt import Tool, tool_loop
from agentflow.events import Message
from agentflow.backends.ollama import OllamaBackend

async def multiply(a: float, b: float) -> str:
    return str(a * b)

tools = [Tool("multiply", multiply, description="Multiply two numbers",
              schema={"type": "object",
                      "properties": {"a": {"type": "number"}, "b": {"type": "number"}},
                      "required": ["a", "b"]})]

llm = OllamaBackend("gpt-oss")
await llm.start()
app = tool_loop(llm, tools, max_turns=6)
out = await app.invoke({"messages": [Message("user", "What is 23 * 19?")]})
print(out["messages"][-1].content)   # -> the model's final answer
await llm.close()

See examples/tool_loop_ollama.py for a runnable version, and examples/mixed_backends.py for a single workflow that combines an agent backend (Codex/Claude/Kiro) with an LLM tool loop, composing a compiled sub-graph inside a node.

agentflow.prebuilt also ships with_retry(node, retries=..., on=...) and with_timeout(node, seconds=...) — composable wrappers that add retry with exponential backoff and per-node timeouts. Both let an InterruptError through untouched, so human-in-the-loop suspends are never retried or timed out.

Project layout

agentflow/
  state.py            channels, reducers, update merge
  graph.py            Graph builder + validation
  runtime.py          super-step scheduler, Context, interrupts
  compiled.py         CompiledGraph runnable
  events.py           messages, requests, streaming events
  errors.py           exception hierarchy
  observability.py    Hooks + RunMetrics
  telemetry.py        MultiHooks + JsonlTelemetry (durable JSONL)
  backends/           base protocols + kiro/codex/claude_code/ollama/openai/anthropic
  checkpoint/         Checkpointer protocol + memory/file/sqlite/redis/postgres
  store/              Store protocol (cross-thread memory) + memory/postgres
  controlplane/       RunQueue + GraphRegistry + Worker (memory/postgres)
  prebuilt/           iterate_until_converged, tool_loop, with_retry/with_timeout
examples/             runnable examples
tests/                pytest suite (async; offline by default, `-m live` for real backends)
DESIGN.md             architecture and contracts

Development

pip install -e '.[ollama,dev]'
pytest -q                 # offline suite only (hermetic, fast)
pytest --cov              # with coverage (source=agentflow, branch)

Tests are async and run under pytest-asyncio in auto mode, so no per-test decorator is needed. Coverage is opt-in via --cov to keep the default run fast.

Lint, format, and type checks (run in CI):

ruff check agentflow tests examples
ruff format --check agentflow tests examples
mypy agentflow
pip-audit

See CHANGELOG.md for the release history.

The suite is split by a live marker. The default run skips live tests (addopts = -m 'not live') so CI stays hermetic — everything mocks its transport. Live smoke tests spawn the real agent CLIs and hit a real Ollama server; run them explicitly:

pytest -m live            # runs only backends that are installed/reachable
KIRO_AGENT=vibe pytest -m live   # include the Kiro backend

Each live test skips itself when its backend is absent, so pytest -m live never fails on a missing CLI.

License

This project is licensed under the Apache License, Version 2.0 (Apache-2.0). See the LICENSE and NOTICE files for details.

Copyright 2026 Daniel Sorkin

Metadata

Release files for agent-workflow-sdk 0.1.1

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for agent-workflow-sdk 0.1.1
File Size Uploaded
agent_workflow_sdk-0.1.1.tar.gz 142.7 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for agent-workflow-sdk 0.1.1
File Interpreter ABI Platform
agent_workflow_sdk-0.1.1-py3-none-any.whl Python 3 none any Details

Total release size: 271.5 kB

Release files / agent_workflow_sdk-0.1.1.tar.gz

Download URL agent_workflow_sdk-0.1.1.tar.gz
Size 142.7 kB
Tags Source
SHA-256 checksum
How to use checksums
75879c9f132e663d5729f83e78d8f86a46af221eeecf4444671d8c80ab2e709a
BLAKE2b-256 checksum
How to use checksums
95f1dbe16a4ec2825c3644721e48b12a41fbcf490acb0a7d2ef92316228d3e6b
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.15

Release files / agent_workflow_sdk-0.1.1-py3-none-any.whl

Download URL agent_workflow_sdk-0.1.1-py3-none-any.whl
Size 128.8 kB
Tags Python 3
SHA-256 checksum
How to use checksums
d7cf35acd0e433d55f28fc588d7b0c0ebf96e2a3f1385b6da3989cb59f3a8e97
BLAKE2b-256 checksum
How to use checksums
e7269f483607810e6cfd3ebb2976b8387407c6deae7a6af448cfe182348d9de4
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.15

Release history Release notifications | RSS feed

0.2.0

2 release files

This release

0.1.1 This release

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page