How ObzenFlow Works
Build durable software from a small, composable vocabulary.
At the heart of ObzenFlow are stages, specialized handlers, and composites: a small vocabulary for transforming events, joining streams with reference data, accumulating state, and making external effects explicit. You’ll learn how to compose them into sophisticated durable programs that bring the directness of classical application development to real-time stream processing.

Building with ObzenFlow
See the whole ObzenFlow program in one declaration.
The application hosts one declaration of journals, middleware, typed stages, and topology.
Operating ObzenFlow
Run every flow like an operable service.
You define the flow. FlowApplication handles lifecycle, metrics, server endpoints, and graceful shutdown.
Building with ObzenFlow
Anatomy of an ObzenFlow Application
The shape of an ObzenFlow program
Read the program from the outside in. FlowApplication hosts one flow! declaration, and that declaration names four things: journals, middleware, stages, and topology.
- FlowApplication
- Host boundary that runs the flow. Its lifecycle and endpoints live on Operating ObzenFlow.
- flow!
- Macro that declares everything inside.
- journals
- Where each stage writes its append-only journal.
- middleware
- Flow-wide runtime protections. Stages inherit unless they override.
- stages
- The typed units of work.
- topology
- Graph of stage edges. Forward edges use
|>. Backflow edges use<|.
The macro expands at compile time into fully type-checked async Rust. If a stage emits one type but the next stage expects another, the compiler catches it before the binary exists. When the binary starts, FlowApplication materializes the declaration into the supervised runtime stages described next.
#[tokio::main]
async fn main() -> Result<()> {
// The runner: a plain async main that starts the runtime.
FlowApplication::run(flow! {
// The flow: types and shape, not plumbing.
name: "stateful_counter",
journals: disk_journals(PathBuf::from("target/counter-logs")),
middleware: [rate_limit(5.0)],
stages: {
numbers = source!(
NumberEvent => number_source(10)
);
evens = transform!(
NumberEvent -> NumberEvent => even_filter()
);
counter = stateful!(
NumberEvent -> AggregationResult => CounterHandler::new()
);
output = sink!(
AggregationResult => sinks::json_pretty()
);
},
topology: {
numbers |> evens;
evens |> counter;
counter |> output;
}
})
.await?;
Ok(())
}Six runtime stage types execute the graph
StageType is the runtime’s concrete classification. A declaration may offer richer authoring syntax, but every ordinary stage ultimately materializes as exactly one of these six types.
- FiniteSource
- Produces a bounded input and eventually reaches end-of-file.
- InfiniteSource
- Admits an always-on input for the life of the process.
- Transform
- Processes one input at a time without framework-managed persistent state.
- Stateful
- Folds events into state that survives across inputs.
- Join
- Maintains keyed reference facts beside a forward-moving stream.
- Sink
- Terminates a branch by delivering committed facts downstream.
The DSL presents the two source types as one source family because boundedness and sync-versus-async production are decisions at ingress. It also offers specialized declarations that reuse Transform; those are called out separately rather than pretending the runtime gained another stage type.
stages: {
// StageType::FiniteSource
products = source!(Product => product_catalog);
// StageType::InfiniteSource
orders = async_infinite_source!(Order => live_orders);
// StageType::Transform
validated = transform!(Order -> ValidatedOrder => validate());
// StageType::Join
enriched = join!(
catalog products: Product,
ValidatedOrder -> EnrichedOrder => enrich
);
// StageType::Stateful
totals = stateful!(
EnrichedOrder -> OrderTotals => typed_stateful::reduce(
OrderTotals::default(),
accumulate,
)
.emit_on_eof()
);
// StageType::Sink
output = sink!(OrderTotals => write_totals);
}
topology: {
orders |> validated |> enriched |> totals |> output;
}Specializations reuse stages while composites expand into a graph
The six runtime types are the concrete execution vocabulary. ObzenFlow also provides higher-level declarations that either strengthen one of those types or compose several of them into one reusable capability.
In the example, an investigation request first loads evidence from an observability API. The resulting IncidentEvidenceLoaded fact fans out: one model call assesses severity, while AI MapReduce works through the full incident timeline to build a report. IncidentEvidenceNotFound follows its own typed route instead of becoming a generic error DTO.
- effectful_transform!
- Adds declared effects and typed outcome facts to one
Transformstage. - inference!
- Generates the
ChatCompletioneffect protocol around one scalar model call and materializes it as oneTransformstage. - ai_map_reduce!
- Expands one declaration into chunk, map, collect, and finalize stages connected by typed internal feeds.
A specialization still produces one descriptor, one supervisor, and one journal-owning execution boundary. A composite produces a named subgraph whose member stages retain their own supervisors and journals. Neither creates another hidden runtime type.
flow! {
name: "incident_intelligence",
journals: disk_journals(journal_base),
stages: {
investigations_source = source!(
IncidentInvestigationRequested => investigation_feed
);
load_incident_evidence = effectful_transform!(
IncidentInvestigationRequested -> {
IncidentEvidenceLoaded,
IncidentEvidenceNotFound,
}
uses LoadIncidentEvidence
via observability_api
with telemetry_resilience()
=> evidence_loader
);
assess_incident_severity = inference!(
IncidentEvidenceLoaded -> IncidentSeverityAssessed
uses at_least_once(ChatCompletion)
via chat
with ai_resilience()
=> severity_assessor
);
build_incident_report = ai_map_reduce!(
IncidentEvidenceLoaded -> IncidentReport => {
map: [IncidentTimelineEntry] -> TimelineGroupSummary
uses at_least_once(ChatCompletion)
via chat
with ai_resilience()
=> timeline_summarizer,
reduce: (
IncidentEvidenceLoaded,
[TimelineGroupSummary]
) -> IncidentReport
uses at_least_once(ChatCompletion)
via chat
with ai_resilience()
=> report_writer,
},
chunking: by_budget {
items: |incident: &IncidentEvidenceLoaded| {
incident.timeline.clone()
},
render: |entry: &IncidentTimelineEntry, _ctx| {
entry.render_for_model()
},
budget: obzenflow_core::ai::TokenCount::new(4_000),
max_items: Some(50),
oversize: error,
}
);
severity_sink = sink!(
IncidentSeverityAssessed => store_severity
);
report_sink = sink!(IncidentReport => store_report);
missing_evidence_sink = sink!(
IncidentEvidenceNotFound => record_missing_evidence
);
},
topology: {
investigations_source |> load_incident_evidence;
load_incident_evidence |> assess_incident_severity;
load_incident_evidence |> build_incident_report;
load_incident_evidence |> missing_evidence_sink;
assess_incident_severity |> severity_sink;
build_incident_report |> report_sink;
}
}Operating ObzenFlow
How to Run an ObzenFlow Application
The operational substrate is built in.
You define the business facts, stages, and topology. FlowApplication turns that flow into a supervised service: it manages startup and shutdown, exposes a standard operator surface, and records a durable journal for every stage. The whole application ships as one Rust binary.
- Run anywhere
- Deploy the binary as an ordinary process or container. It works with existing schedulers and process signals, with no ObzenFlow control plane to install.
- Observe and control
- Check health and readiness, scrape metrics, inspect topology and resolved configuration, follow live events, and control the flow through the built-in operator surface.
- Recover safely
- Inspect durable stage history, verify a replay against its recorded run, or resume from the durable frontier without repeating committed effects.
let flow = build_flow();
FlowApplication::run(flow).await?;cargo run -p obzenflow \
--example http_ingestion_piggy_bank_demo \
--features obzenflow_infra/warp-server# Availability and telemetry
GET /health
GET /ready
GET /metrics
# Runtime model and lifecycle
GET /api/topology
GET /api/flow/events
POST /api/flow/control
# Resolved configuration
GET /api/config
GET /api/config/overlay
GET /api/config/effective
GET /api/config/schema
GET /api/config/diff
GET /api/config/flows/:flow_id
GET /api/config/flows/:flow_id/stages/:stage_keyShips with sensible integrations out of the box.
ObzenFlow keeps infrastructure optional. Compile the capabilities your application needs into the same binary, without adopting a separate platform or control plane.
- HTTP ingress and operations
Serve the built-in operator API, typed ingress, and application-defined HTTP endpoints from the same binary, with managed authentication, CORS, request limits, timeouts, readiness, and graceful shutdown.
features = ["warp"]- HTTP pull
Add Reqwest-backed pull and polling sources for remote feeds.
features = ["http-pull"]- AI
Add Rig-backed model providers and tiktoken budgeting for replayable inference.
features = ["ai"]- PostgreSQL
Add typed PostgreSQL delivery for repeat-safe projections and downstream integration.
features = ["postgres"]- Metrics and diagnostics
Export Prometheus-compatible runtime metrics and optionally inspect asynchronous work through Tokio Console.
features = ["console"]- Studio discovery
Let running applications register and renew their presence with ObzenFlow Studio.
features = ["studio"]
Durable execution with the operator surface your SRE team needs from day one.
Compile with the obzenflow_infra/warp-server feature, then enable the server through configuration or --server. One managed HTTP surface then answers the first questions from orchestrators, operators, development tools, and ObzenFlow Studio.
- Availability
/healthreports process liveness./readyreports whether the flow is ready to accept work.- Topology
/api/topologyexposes the materialized graph, stage types, typed edges, contracts, and composite membership.- Configuration
/api/config/*exposes the resolved, redacted configuration and its effective flow and stage views.- Telemetry
/metricsexports measurements while/api/flow/eventsstreams lifecycle events over SSE.- Control
/api/flow/controlstarts a manual-mode flow or requests cancellation and graceful stop.- Hosted ingress
Optional typed ingress surfaces accept individual events and batches while following the flow's readiness and shutdown lifecycle.
Run live, replay with proof, or resume
The same application exposes four operator actions over the history it already records.
- Live
- Accepts and processes new work.
- Replay
- Reconstructs an archived run without polling its live sources or repeating committed effects.
- Verify
- Compares a completed replay with the original archive and reports whether they agree.
- Resume
- Reconstructs the durable prefix, then reopens live work beyond its recorded frontier.
The runtime itself resolves three execution modes: Live, Replay, and Resume. --verify is the certification step applied to Replay, not a fourth execution mode.
$ cargo run -p obzenflow \
--example payment_gateway_resilience
payment_gateway_resilience_demo completed.
Journal: …/flow_01M1CMS6JG1J7944YDWYFJ8S35
To verify with a bounded replay, add:
--replay-from <journal>
To continue this run live from where it left off, add:
--resume-from <journal>
$ cargo run -p obzenflow \
--example payment_gateway_resilience -- \
--replay-from <journal> --verify
Archived outcomes were reconstructed:
source config ignored
effects suppressed
recorded facts reused
output matched the original run, 0 differencesNext Up
Visualize your running workflow on one canvas.
ObzenFlow Studio turns topology, metrics, live events, middleware state, and contract evidence into one visual operating surface.