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.

An upstream supervised state machine writing to a journal that fans out to two independently supervised downstream state machines, illustrated beside the corresponding ObzenFlow topology declaration.

Building with ObzenFlow

Anatomy of an ObzenFlow Application

Start with the outer shape: the application that runs the flow, the declaration that defines it, and the typed graph the runtime materializes.

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.

stateful_counter.rs
#[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.

stage_vocabulary.rs
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 Transform stage.
inference!
Generates the ChatCompletion effect protocol around one scalar model call and materializes it as one Transform stage.
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.rs — three authoring shapes in one flow
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

Ship an ObzenFlow application as one binary and run it like any other service. Familiar health, metrics, and control endpoints work with the tools your operators already use. Every run leaves a durable journal that you can investigate offline, replay without repeating committed effects, or resume safely after a failure. None of this requires a separate orchestration platform.

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.
main.rs
Step 1 · Define and host the flow
let flow = build_flow();

FlowApplication::run(flow).await?;
terminal
Step 2 · Compile and launch the HTTP service
cargo run -p obzenflow \
  --example http_ingestion_piggy_bank_demo \
  --features obzenflow_infra/warp-server
built-in operator surface
Step 3 · Use the endpoints FlowApplication provides
# 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_key

Ships 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

/health reports process liveness. /ready reports whether the flow is ready to accept work.

Topology

/api/topology exposes 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

/metrics exports measurements while /api/flow/events streams lifecycle events over SSE.

Control

/api/flow/control starts 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.

payment_gateway_resilience — real run transcript
$ 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 differences

Next 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.