ObzenFlow is a durable execution runtime in Rust, built around durable per-stage journals, typed events, and a DSL for composing auditable flows.
Most systems in this space put infrastructure at the center, and you pay for it in operational overhead. You end up standing up brokers, coordinators, or clusters before you ship a single business event. ObzenFlow takes a different path. It compiles to a single binary, journals every event for replay and audit out of the box, and enforces typed correctness contracts between stages at compile time. Every event is also a wide event, carrying full causality, history, agency, intent, and narrative context automatically, so you spend your time on business logic instead of plumbing. You get production-grade durable execution without the platform tax.
The first example you will work through is char_transform. It takes three short sentences as input, visits each character, capitalizes lowercase letters, turns digits into words, accumulates the transformed stream into one final text value, and prints that result at the sink.
The example is deliberately small, but deceptively powerful. In real systems, this same pattern shows up in problems such as:
- Event validation and normalization before persistence
- Text cleanup, classification prep, or PII redaction
- ETL and CSV cleanup jobs that still need replayable history
- Feature generation and aggregation steps ahead of downstream services or models
With those same building blocks, you can eventually move on to building event-sourced service endpoints, AI inference pipelines with Ollama, external delivery to queues and databases with correctness guarantees, and larger flows that remain replayable, auditable, and bulletproof.
Let’s begin.
Step 0. Designing a character transformer
You just pitched your startup idea to a venture capital firm and it was accepted. “Enterprise Character Transformer Inc.” is born! Now the hard part begins. How do you build this on time, on budget, and in the most correct way possible?
Before writing a single line of code, think about what you actually need. A server. Metrics. Middleware. Simple development and deployments. Correctness and replayability, because you intend to remain SOC II compliant from day one. That’s a lot of plumbing before you even get to the interesting part.
ObzenFlow handles all of this out of the box, so instead of spending your first week wiring up infrastructure, you can spend it thinking about what your system actually does.
And that’s exactly where this tutorial starts. Rather than jumping straight into technical details, we’ll think about the flow first. How do you process text, chunk it into characters, transform those characters, and reassemble them into sentences and paragraphs? If you’re new to event-driven systems design and you’re wondering where to even begin, you’re in the right place.
Designing with values, not places
In classical place-oriented programming (PLOP), you might start by sketching out methods and database schemas. But ObzenFlow is event-driven. That means we need to rethink our approach to design.
Most systems are designed around mutable places. A row, an object, a field holding whatever value was written last. ObzenFlow is designed around immutable facts. A fact is a value that records something that happened at a point in time. Facts are values, not places. When you append facts instead of overwriting places, state becomes a projection of history instead of a single opaque snapshot.
In ObzenFlow, designing a system starts with facts and then progresses naturally into cause and effect. You are not thinking first about schemas. You are thinking about interesting things that happen over time. Once those facts are clear, the flow becomes easier to see. An event happens, the system reacts to that event, and a new event happens. The current state of any identifiable entity along the way is simply the accumulation of prior events: cheque deposited, cash withdrawn, payment posted.
An event is simply an immutable thing that happened and cannot be refuted. We call these facts. Facts are values, not places. Facts don’t change. And the reason we call a fact an event is that a fact by definition incorporates time. From the Latin factum, something done.
What would it mean to build information systems that are about information? It would mean that the systems would be fundamentally about facts. It would be about maintaining facts, manipulating facts, and presenting facts to users to give them leverage so they can make decisions.
In ObzenFlow, facts are the substrate of the entire framework. Every fact is durably journaled, and state changes are treated as first-class values. That history of facts is what makes replay, audit, and correctness possible.
Character transformation as a linear event stream
For a character transformer, three facts emerge:
- A character arrives from the input text.
- A text chunk is produced after transformation rules are applied.
- A final transformed text is assembled from all the chunks.
That stream of events is enough to sketch out our first three Rust types. These represent the initial events of the system.
struct CharInput {
character: char,
}
struct TextChunk {
text: String,
}
struct TransformedText {
text: String,
sentence_count: usize,
character_count: usize,
last_was_sentence_end: bool,
}With the domain events defined, you can wire them into a flow. ObzenFlow lets you express the entire topology without writing a single handler. The placeholder!() macro stands in for logic you have not written yet, and the skeleton still compiles. It’s rare that your design will be correct after the first domain sketch, so this “type-flow skeleton” design approach lets you iterate quickly.
// FlowApplication provides the operational shell: metrics, web
// endpoints, lifecycle management. You write the flow, it runs it.
FlowApplication::run(flow! {
name: "char_transform_skeleton",
// Journals are event-sourced. Actual computation becomes the
// observability substrate, making audits and replay simple.
journals: disk_journals(PathBuf::from("target/char-transform-skeleton-logs")),
// Pluggable middleware comes standard. Circuit breakers, rate
// limiters, and other operational policy slot in here.
middleware: [],
// Stages cover every processing need: sources, transforms,
// stateful accumulators, stream-table joins, sinks, and more.
stages: {
characters = source!(CharInput => placeholder!());
transform_text = transform!(CharInput -> TextChunk => placeholder!());
collect_text = stateful!(TextChunk -> TransformedText => placeholder!());
output = sink!(TransformedText => placeholder!());
},
// Topology wires stages together by type and gives you a
// visual map of your flow in the (upcoming) ObzenFlow Studio.
topology: {
characters |> transform_text;
transform_text |> collect_text;
collect_text |> output;
}
})
.await?;You have not written a single line of business logic, and you already have a compilable pipeline with four stages wired in sequence. That pipeline is far more than a sketch. It comes with:
- Journal configuration, so every event will be durably recorded the moment you fill in the handlers.
- A middleware slot ready for rate limiters, circuit breakers, or whatever operational policy you need.
- Typed contracts between every stage, enforced by the Rust compiler. If you wire
CharInputinto a stage that expectsTextChunk, it will not compile.
This skeleton runs. The placeholders produce no output, but the framework still initializes the journals, validates the topology, and drains the pipeline cleanly. You can even start the server and view the topology in the UI before a single line of handler code exists.
That is the design-first workflow ObzenFlow enables. The framework directly connects event-driven design techniques to production-grade code, not a proof of concept, and helps you get there within a day rather than a month. Name your domain events, sketch the topology, and fill in behaviour when you are ready.
In the next steps, you will clone the repo, replace those placeholders with real logic, and watch the pipeline transform actual input.
ObzenFlow embraces wide events, which moves semantic meaning into the event itself. When naming events, keep in mind that your event names do not need to capture all semantic meaning, which is an anti-pattern that can lead to bloated naming schemes. Wide events ensure that robust information is carried in the envelope, including who created it, why, what came before it, what it was intended to accomplish.
Under the hood, every event your handler touches is carried inside a ChainEvent, the framework’s wide event type. A ChainEvent structurally encodes causality, history, agency, intent, and narrative. It knows where it came from, which stage produced it, what correlation chain it belongs to, whether it is a replay, and what observability context was active at creation time. All of that travels with the event automatically, without you writing a line of plumbing code.
Your domain type, CharInput, sits inside that wide envelope. The name describes what the business cares about. The envelope carries everything else. Name your events the way your business talks about them. The framework carries the rest. To see the chain event envelope for yourself, the code is in chain_event/model.rs.
Step 1. Clone ObzenFlow and set up your environment
The first step is cloning the repo and making sure your environment is set up correctly. Clone the main framework repo, ObzenFlow/obzenflow, which contains all examples referenced in this tutorial.
git clone https://github.com/ObzenFlow/obzenflow
cd obzenflowIf you are integrating ObzenFlow into an existing Rust project instead of starting from the examples, the crate is also published on crates.io. The examples catalog is still a useful companion because it shows the full flow model in runnable form.
- Install Rust with
rustup. The examples on this page target Rust1.93.0. - Most examples run cleanly with
cargo run -p obzenflow --example <name>. --features obzenflow_infra/warp-serverenables HTTP endpoints and/metrics.--features http-pullenables HTTP pull sources.--features "http-pull ai"enables the AI digest example.- Use
RUST_LOG=infoinstead when you want to see stage startup, middleware activity, and more of the framework’s runtime behavior. - If you are new to Rust, use the commands on this page exactly as written and do not worry about understanding every macro or trait on the first pass.
- Read each example in this order: the
flow!block, the topology, then the handler code behind each stage. - Treat the first run as a product tour. Learn the ObzenFlow model first, then come back for Rust details.
Step 2. Exploring the basics of ObzenFlow
This first example will give you a strong sense of the developer experience in ObzenFlow and how the framework feels in practice. Pay attention to the getting started experience itself: cloning the repo and running a single command with no other ceremony is part of the ethos behind the framework. All batteries are included and you could even enable a Prometheus metrics endpoint for this example with one command-line switch.
You can work through this section on its own, or open examples/char_transform.rs and compare each part of the walkthrough to the actual example. Both approaches will help you get up to speed with the framework, and using them side by side is usually the fastest path.
If you would rather run the example first, use:
cargo run -p obzenflow --example char_transformHere is what to expect when you run char_transform:
Input:
hello 2024 world!
42 is the answer.
rust 1 python 0.
Output:
HELLO twozerotwofour WORLD!
fourtwo IS THE ANSWER.
RUST one PYTHON zero.Let’s begin. You will start with the overall shape of a typed pipeline, then work through how flows, stages, events, and handlers fit together, and finish by connecting all of that to topology, execution, and replay.
Learning sources
The source is what determines whether your pipeline is a batch job or a long-running service. ObzenFlow supports both modes with the same stage model, the same journaling, and the same typed contracts. You do not need a separate batch system and a separate streaming system.
Finite sources produce a bounded input and then signal completion. A collection of records, a CSV file, a fixed set of API responses. The pipeline processes everything and terminates cleanly. This is the mode that early frameworks like Apache Spark were built around, and it remains the right choice for ETL jobs, data migrations, report generation, and anything where the input has a known end.
// From a collection
characters = source!(CharInput => sources::finite(char_inputs));
// From a producer function that returns None when exhausted
orders = source!(OrderEvent => sources::finite_from_fn(|index| {
if index >= 100 { None } else { Some(OrderEvent { /* ... */ }) }
}));
// From a CSV file on disk
customers = source!(Customer => CsvSource::typed_from_file::<Customer>(&path)?);Continuous sources run indefinitely. They accept HTTP requests, poll external APIs, or consume from any upstream feed that never ends. The pipeline stays alive and processes events as they arrive. This is the mode you use for live services, real-time dashboards, and event-driven architectures.
// Push-based HTTP ingestion (service-like)
accounts = source!(AccountOpened =>
HttpSource::with_telemetry(accounts_rx, telemetry)
.typed::<AccountOpened>()
);
// HTTP polling with decode logic
stories = source!(Story => HttpPullSource::new(decoder, config));
// A producer that never stops
ticks = source!(Tick => sources::infinite(|batch| {
vec![Tick { ts: now() }]
}));In both modes, the stages downstream of the source are identical. The same transforms, stateful accumulators, joins, and sinks work regardless of whether the source is finite or continuous. The same journals record every event. The same contracts verify delivery at every edge. The only difference is whether the source eventually signals EOF.
char_transform uses a finite source. It reads a fixed set of input sentences, processes them through a chain of stages, and terminates when the input is exhausted. That makes it the simplest starting point, because the pipeline runs, completes, and leaves behind a journal you can inspect at your leisure.
Understanding FlowApplication and defining a flow
Every ObzenFlow application starts with two layers: FlowApplication, which gives you the runtime shell that executes the flow, and flow! (pronounced “flow bang”), which gives you the declarative shape of the flow itself.
FlowApplication handles the runtime lifecycle. flow! is where you define the flow itself: its name, journals, middleware, stages, and topology.
FlowApplication::run(flow! {
name: "char_transform",
// Durable event journals. Disk is the default.
journals: disk_journals(PathBuf::from("target/char-transform-logs")),
// Flow-level operational policy: rate limiters, circuit breakers, etc.
middleware: [],
stages: {
// declare your stages here
},
topology: {
// connect the stages here
}
})
.await?;That block is worth reading in order:
FlowApplication::run(...)starts the runtime shell that executes the flow. It handles lifecycle, metrics, server endpoints, and graceful shutdown.flow! { ... }defines the flow itself. Everything inside this block describes what the flow is, not how it runs.namegives the flow a stable identity.journalschooses where the durable history will live on disk.middlewaresets shared operational policy for the whole flow.stagesdefines the pieces.topologydefines how data moves between those pieces.
The char_transform example leaves middleware empty, but the slot is there for when you need it. Other examples in the repo use it to set flow-level rate limits, circuit breakers, and other operational policy that applies to every stage in the pipeline.
// From the payment_gateway_resilience example
middleware: [
RateLimiterBuilder::new(1.0).build()
]Middleware can be declared once for the whole flow or attached directly to an individual stage when one part of the pipeline needs its own policy. The syntax is the same either way.
Disk journals are the default and the right choice for most workloads. They give you replay and audit for free, and the data survives process restarts. For performance-sensitive or ephemeral workloads where durability is not required, ObzenFlow also provides memory journals. If neither fits, you can implement your own journal backend.
Understanding stages, events, and handlers
Stages are the named pieces of your flow. Every stage you declare inside the stages block becomes a unit of processing with its own journal, its own lifecycle, and its own typed contract with the stages around it. Events are the facts that move between stages. Handlers are the code you write to tell a stage what to do with those facts.
ObzenFlow keeps the stage menu deliberately small:
sourceproduces events into the flowtransformchanges one event into the nextstatefulaccumulates facts across many eventsjoinenriches a stream with reference datasinksends results out of the flow boundary
char_transform uses four of those five. You already saw the shape of these stages in step zero, where every handler was placeholder!(). The pattern is the same here. Every stage follows this shape:
stage_name = transform!(In -> Out => transforms::map(|input: In| {
Out { ... }
}))The type contract sits on the left of =>, the handler sits on the right. A source has one type (its output). Every other stage has an input type and an output type separated by ->. The handler is a function that takes In and returns Out. You are stitching together typed tapes with functions, and the compiler enforces that the types match when you wire the topology.
Now it is time to fill in those placeholders. Step 3 walks through each stage type in the char_transform example one at a time.
Step 3. Writing handlers for each stage
In step zero you sketched the topology with placeholder!() handlers. In step two you learned the shape: In -> Out => handler. Now you will see what real handlers look like for each stage type in the char_transform example.
char_transform only uses a small slice of ObzenFlow’s stage system. The framework ships a broader catalogue of source helpers, transforms, accumulators, emitters, joins, and sinks.
Source
You already covered sources in depth in step two under Learning sources. The char_transform source uses sources::finite to emit one CharInput event per character from a collection.
characters = source!(CharInput => sources::finite(char_inputs));The char_inputs collection is built before the flow starts by breaking the input sentences into individual characters. When the collection is exhausted, the source signals EOF and the pipeline begins to drain.
Transform
A transform is stateless. It looks at one event at a time and produces the next event. It does not remember anything between invocations.
transform_text = transform!(
CharInput -> TextChunk => transforms::map(|input: CharInput| TextChunk {
// `text` is the String field on the `TextChunk` output type.
text: transform_char(input.character),
})
);transforms::map is a 1-to-1 helper. It receives a CharInput, applies transform_char to capitalize letters and convert digits to words, and returns a TextChunk. The framework journals both the input and the output, so you can trace every transformed character back to the original.
Stateful
A stateful stage remembers what came before. This is how you get totals, summaries, counters, rollups, and any result that depends on history rather than a single input.
collect_text = stateful!(
TextChunk -> TransformedText => typed_stateful::reduce(
TransformedText::default(),
// `reduce` folds each chunk into one running accumulator.
// We update `acc` in place rather than allocating a new value each step.
|acc, chunk: &TextChunk| {
acc.text.push_str(&chunk.text);
acc.character_count += chunk.text.chars().count();
let ends_sentence = matches!(chunk.text.as_str(), "." | "!" | "?");
if ends_sentence && !acc.last_was_sentence_end {
acc.sentence_count += 1;
}
acc.last_was_sentence_end = ends_sentence;
}
).emit_on_eof()
);typed_stateful::reduce takes an initial accumulator (TransformedText::default()) and a closure that folds each incoming TextChunk into it. If you know functional programming, this is a left fold: start with one accumulator, then keep accumulating each chunk into that running value. The closure variable is named acc for accumulator. It appends the text, updates the character count, and tracks sentence boundaries.
One of the most obvious things you may have noticed is that there is no return clause in the reducer closure. That is intentional. In a stateful reduction like this, you fold each incoming event into the accumulator and return nothing from the closure itself. The output event is produced later by the emission strategy.
The .emit_on_eof() at the end is that emission strategy. In this example, it tells the framework to take the accumulator, turn it into one output event, and write it to the journal when the upstream source is exhausted. Accumulation and emission are separated on purpose. reduce describes how to build the accumulator. emit_on_eof() describes when that accumulated value becomes output.
Sink
A sink receives the final output and delivers it outside the flow boundary. It is where events leave the system.
output = sink!(TransformedText => sinks::console(format_output));sinks::console takes a formatter function and prints the result. In a real application, the sink might write to a database, call an API, publish to a queue, or write a file. The framework journals the delivery fact regardless of destination.
Step 4. Topology, execution, and replay
You have defined event types, wired stages with handlers, and understand the contract between them. The last piece is the topology, which tells the framework how data flows between stages, and then running the example to see what the framework produces.
Topology
ObzenFlow wears one of its biggest influences openly. Some of the most composable primitives in computing history are Unix pipes and files on disk. ObzenFlow pays respect to that lineage and adds syntactic and conceptual sugar on top.
|> means events flow forward into the next stage. <| means results feed back in the other direction. The pipe is still the organizing idea. What ObzenFlow adds is explicit direction, typed events, and durable journals behind every stage boundary.
topology: {
characters |> transform_text;
transform_text |> collect_text;
collect_text |> output;
}That is the whole pipeline. Characters come in, get transformed, get accumulated, and leave through the sink. The topology reads like an architecture diagram instead of hiding the pipeline shape inside framework glue.
The topology syntax was directly inspired by Unix pipes. Pipes as processing and files as the data substrate are about as perfect as computing gets, all these years later. ObzenFlow takes that same idea and adds typed events, durable journals, and stage contracts on top. If you are comfortable reading cat input | transform | sort | sink, you already understand how to read an ObzenFlow topology.
Running the example
You have seen the types, the stages, the handlers, and the topology. Now run it. The framework defaults to warn log level, so the output is clean.
cargo run -p obzenflow --example char_transformHere is what you should see:
Input
hello 2024 world!
42 is the answer.
rust 1 python 0.
Output
HELLO twozerotwofour WORLD!
fourtwo IS THE ANSWER.
RUST one PYTHON zero.
3 sentences, 72 characters
FlowApplication complete!
To replay, add: --replay-from target/char-transform-logs/flows/<run-id>
(Source config env vars are ignored during replay)The run archive
After you run the example, look under target/char-transform-logs/flows/<run-id>/. ObzenFlow writes a run archive for every execution. Each stage gets its own journal file, and the run_manifest.json records the flow structure, stage metadata, and journal file mappings for that run.
This is the first concrete reason this is more than a for loop over a collection. The business logic stays typed end to end, and the framework durably records every event that passed through the pipeline.
Replay
That run archive is what replay reads. You can point the framework at a previous run and re-execute the flow against the same historical inputs.
cargo run -p obzenflow --example char_transform -- \
--replay-from target/char-transform-logs/flows/<run-id>The runtime swaps the live source for the archived source journal, attaches replay metadata to each event, and reruns the flow against the same inputs. That is how ObzenFlow turns a past execution into something you can inspect, test against, and compare after changing the code. The same archive shape that this tiny example produces is what the larger replay-enabled applications use.
Next steps
At this point, you should have one complete ObzenFlow example in your hands. You have seen the outer shape of a flow, the stage vocabulary, the topology model, the journal-backed run archive, and the line between framework code and application code.
When you want to go further, use the docs section as a tutorials surface for now:
- Model bank transactions as a flow explains HTTP ingress, joins, projections, and why journal-backed accounting fits ObzenFlow so naturally.
- Run live AI inference from a real endpoint explains token estimation, chunking, accumulation, and why model-backed workflows still need deterministic scaffolding around inference.
The philosophy is worth reading once the tutorials feel natural, if you want the architectural rationale behind why ObzenFlow is shaped this way.