This tutorial assumes you have already worked through Getting Started and the char_transform walkthrough. Getting Started explains FlowApplication, flow!, stage types, topology, journals, and replay. This page uses those same building blocks in a service-shaped example, then introduces the building blocks you need for live services.
In this tutorial you will build a small checking-account service from end to end. Two HTTP endpoints sit at the front of the service.
- One endpoint accepts the opening of new accounts.
- The other endpoint accepts the credits and debits that move money in and out of those accounts.
- Behind the endpoints, an ObzenFlow pipeline matches each transaction to its account, keeps a running balance per account, and emits a fresh checkbook snapshot every time the state changes.
The example is small enough to read in one sitting, but it still has the shape of a real write-side service. Live commands arrive at the edge. Durable facts move through the flow. A read model appears downstream, and the whole run leaves a replayable journal behind it.
ObzenFlow is a durable execution runtime. It runs your flow and, as it runs, keeps a record solid enough to prove afterward exactly what happened and to replay it step for step. Banking is where that model becomes especially intuitive.
Step 0. Designing the transaction flow
Banking is a natural domain for ObzenFlow because it already thinks in journals.
An account is opened. A credit is posted. A debit is posted.
Your current balance is derived from the entries that came before it, meaning the statement you read at the end of the month is a projection of every posted entry. Deriving state from a journal of facts rather than overwriting a place that holds the latest value is the spirit of every ObzenFlow program.
Designing an ObzenFlow service follows the same flow you saw in the Getting Started tutorial:
- Understand the domain.
- Name the commands and events.
- Turn those facts into Rust types.
- Sketch the flow topology.
- Wire the types together with
placeholder!()handlers.
That sequence lets you model the shape of a real ObzenFlow program in a few dozen lines before writing the business logic. Using placeholder!() syntax means you can compile the skeleton, inspect the topology in ObzenFlow Studio, and eventually replace each placeholder with real logic once the overall design is clear enough. We will use that same process here.
Understanding the domain
Event storming is a useful way to do the first pass because it asks you to name what happened before you name where it is stored. We are not going to run a full event-storming workshop in this tutorial. Credits, debits, and account openings are familiar enough that we can use the technique lightly, but on most projects, you would want to workshop the business domain itself before proceeding.
For a checking-account service, the list of facts is short enough to write from memory. Three are visible at the service boundary.
- An account is opened with an account id, an owner, and an initial balance.
- A ledger entry is submitted as a credit or a debit for an account.
- A checkbook snapshot is observed by anyone reading the account state.
A fourth fact lives entirely inside the flow.
- A posted entry exists when the framework has matched a submitted ledger entry against a known account.
The fourth fact is where the service crosses from request handling into durable accounting. LedgerEntry is the incoming command. PostedEntry is the accepted fact the rest of the flow can trust.
In ObzenFlow, your code acts as commands in event-sourcing terms, and once a command is accepted by the system, it will produce one or more events, which are facts.
Facts cannot be refuted. Every fact is written to a journal in ObzenFlow, and journal entries can be replayed to rebuild the state of a system. Your intent or command can fail, which means that it is not yet an event.
The way to think about it is this: commands can be refuted, but events cannot be refuted; they are facts. The types that you define (for instance, ledger entry, posted entry, etc.) resemble facts, whereas the code that you define in custom handlers represents intent. The fact that we processed a ledger entry from an external system is a fact. A posted entry is a different type of fact.
Now in a real system we would probably have rejected ledger entries which could feed back into a totally different type of reconciliation system. This is why modeling your domain is so critical but for a simple example like this we are assuming somewhat of a simple happy path.
Sketching domain types
With the facts in hand, the Rust types follow naturally. Here are the four types from examples/http_ingestion_piggy_bank_demo/domain.rs:
pub struct AccountOpened {
pub account_id: String,
pub owner: String,
pub initial_balance_cents: i64,
}
pub enum EntryKind { Credit, Debit }
pub struct LedgerEntry {
pub account_id: String,
pub kind: EntryKind,
pub amount_cents: u64,
pub note: Option<String>,
}
pub struct PostedEntry {
pub account_id: String,
pub owner: String,
pub kind: EntryKind,
pub amount_cents: u64,
pub initial_balance_cents: i64,
pub note: Option<String>,
}
pub struct CheckbookSnapshot {
pub account_id: String,
pub owner: String,
pub current_balance_cents: i64,
pub available_balance_cents: i64,
pub total_credits_cents: i64,
pub total_debits_cents: i64,
pub transactions: Vec<CheckbookEntry>,
}Look closely at PostedEntry and you can see the merge written into the type. It carries the owner and initial_balance_cents from AccountOpened along with the kind, amount_cents, and note from LedgerEntry. The shape of the struct mirrors the shape of the join. When two streams meet, the resulting event carries fields from both inputs, and the type system records the merge as part of the contract.
Sketching the flow
With the events named, you can sketch the topology before writing any handler logic. This is the same placeholder-first workflow you saw in Getting Started and in examples/char_transform_skeleton.rs. The framework wires stages together by type, and the compiler checks every connection even while the handlers are still placeholders.
flow! {
name: "http_ingestion_piggy_bank_demo",
journals: disk_journals(PathBuf::from("target/http-ingestion-piggy-bank-demo-logs")),
middleware: [],
stages: {
accounts = async_infinite_source!(AccountOpened => placeholder!());
tx = async_infinite_source!(LedgerEntry => placeholder!());
posted = join!(
catalog accounts: AccountOpened,
LedgerEntry -> PostedEntry => placeholder!()
);
checkbook = stateful!(PostedEntry -> CheckbookSnapshot => placeholder!());
printer = sink!(CheckbookSnapshot => placeholder!());
},
topology: {
tx |> posted;
posted |> checkbook;
checkbook |> printer;
}
}placeholder!() is the design-time handler for a typed stage, not a replacement for the domain types themselves. We won’t cover placeholders in depth as we covered them in the Getting Started tutorial, but I’d like to add that they shine during domain modelling as they allow shaping a topology and verifying the shape in ObzenFlow Studio before committing to implementation detaills. Use placeholders liberally during design!
Now two parts of that sketch are worth a closer look, because both are new compared to char_transform.
The first is the source declaration, which is a streaming source. Both sources use
async_infinite_source!, the live counterpart to thesource!macro from Getting Started. Getting Started’s input was a finite (batch) collection that processes input with a fixed size. However, real inputs like the transaction endpoint in this example are HTTP endpoints that never signal end of stream, they are infinite, so the completed pipeline stays alive for as long as the service is running and drains events as they arrive.The second is the way
accountsparticipates in the topology. The transaction stream flows forward throughtx |> posted. Theaccountssource enters the graph differently. It is named on the catalog line inside the join itself, and the framework holds the account stream there as a continuously updated side input. That is how stream-table joins are expressed in ObzenFlow, and the next several steps return to this idea from different angles.
We will cover both of these concepts in detail later in the tutorial but for now understand that you’re looking at slightly new terminology that you don’t yet need to fully understand. We already have enough implemented that we can visualize the flow in ObzenFlow Studio.
Visualizing the sketch in ObzenFlow Studio
With this very small amount of code, We can already visualize the skeleton. this is what the example already looks like in ObzenFlow Studio.
That’s enough to design a durable execution flow. Everything else we discuss fills in the details: implementation details, operational aspects like middleware, and debugging aspects like this.
To get a flavor of how an ObzenFlow program runs, we’ll run through the completed example. First we’ll start with the conclusion in step one. Steps 2, 3, and 4 will walk through new concepts that we didn’t cover in the getting started tutorial. Step 5 explains the operational surface area, and step 6 closes the loop on durable execution by stopping the service and resuming it from its own journals.
Step 1. Run the service end to end
The fastest way to build a mental model of the service is to run it before reading the code. If you have already cloned the framework repo from Getting Started, you can stay in that working directory. Otherwise, clone it now.
git clone https://github.com/ObzenFlow/obzenflow
cd obzenflowRun the piggy bank example:
RUST_LOG=info cargo run -p obzenflow --example http_ingestion_piggy_bank_demo --features obzenflow_infra/warp-serverThe example loads examples/http_ingestion_piggy_bank_demo/obzenflow.toml at startup. That file is two lines long and only enables the local HTTP server. By default, the service listens on 127.0.0.1:9090. Leave that terminal running.
Open a second terminal, then we’ll open accounts for Alice and Bob via curl. Essentially, this example demonstrates a stream-to-table join. First, we build a table of accounts. Second, we accept transactions against those accounts, then each transaction is matched to the corresponding entry in the populated accounts table.
curl -XPOST http://127.0.0.1:9090/api/bank/accounts/events \
-H 'content-type: application/json' \
-d '{"event_type":"bank.account_opened","data":{"account_id":"acct-1","owner":"Alice","initial_balance_cents":1000}}'
curl -XPOST http://127.0.0.1:9090/api/bank/accounts/events \
-H 'content-type: application/json' \
-d '{"event_type":"bank.account_opened","data":{"account_id":"acct-2","owner":"Bob","initial_balance_cents":0}}'Each successful POST returns an ingestion acknowledgement.
{"accepted":1,"rejected":0,"errors":[]}Now post credits and debits for both accounts.
curl -XPOST http://127.0.0.1:9090/api/bank/tx/events \
-H 'content-type: application/json' \
-d '{"event_type":"bank.ledger_entry","data":{"account_id":"acct-1","kind":"Credit","amount_cents":250,"note":"paycheck"}}'
curl -XPOST http://127.0.0.1:9090/api/bank/tx/events \
-H 'content-type: application/json' \
-d '{"event_type":"bank.ledger_entry","data":{"account_id":"acct-1","kind":"Debit","amount_cents":99,"note":"coffee"}}'
curl -XPOST http://127.0.0.1:9090/api/bank/tx/events \
-H 'content-type: application/json' \
-d '{"event_type":"bank.ledger_entry","data":{"account_id":"acct-2","kind":"Credit","amount_cents":5000,"note":"gift"}}'
curl -XPOST http://127.0.0.1:9090/api/bank/tx/events \
-H 'content-type: application/json' \
-d '{"event_type":"bank.ledger_entry","data":{"account_id":"acct-2","kind":"Debit","amount_cents":1250,"note":"new shoes"}}'The service terminal prints a fresh checkbook snapshot after every posted transaction. After the Alice transactions, you should see an account block with a current balance of $11.51. After the Bob transactions, you should see a block with a current balance of $37.50.
Account: acct-1 (Alice)
Current: $11.51 | Available: $11.51
Credits: $2.50 | Debits: $0.99 | Tx: 2
Account: acct-2 (Bob)
Current: $37.50 | Available: $37.50
Credits: $50.00 | Debits: $12.50 | Tx: 2We call the output you’re seeing above a projection. Essentially, the stream table join matches each transaction against its account, and the stateful stage “folds” the posted entries into an account state. The account state includes a current balance. The sink in this example renders the latest projection after each transaction to the console.
Try posting a transaction for an account that does not exist:
curl -XPOST http://127.0.0.1:9090/api/bank/tx/events \
-H 'content-type: application/json' \
-d '{"event_type":"bank.ledger_entry","data":{"account_id":"acct-unknown","kind":"Credit","amount_cents":100}}'The endpoint still returns {"accepted":1,"rejected":0,"errors":[]} because the JSON parsed and the event was ingested as a LedgerEntry.
In a real-world system you may want to signal back rejection above, but it really depends on the architecture of your system.
Many banks’ personal checking and savings account systems would not signal back rejection at the point of ingestion, because reconciliation would then be handled downstream. These are the nuances that event storming and domain modeling should be helping you to work through.
The rest of the tutorial walks through the code that produced everything you just saw, then circles back to the operational surfaces and replay archive that the framework left behind.
Step 2. Treating an HTTP endpoint as a source
One of ObzenFlow’s core features is the ability to define RESTful endpoints and treat them as streaming surface areas, so that every invocation against an endpoint streams data from ObzenFlow’s perspective. By the time the flow sees the RESTful source, it looks like any other typed source (such as a file, a queue, or anything else). ObzenFlow’s sources support both streaming and batch data, and ObzenFlow turns REST endpoints into unbounded streaming data.
The example splits that work in two. The runner owns the HTTP surface and its lifecycle. The flow only ever sees the typed source the surface produces. Together, those two pieces turn a route into a streaming source.
Config and hosting
Open examples/http_ingestion_piggy_bank_demo/runner.rs and look at the body of run_example:
let accounts_ingress = http_ingress::<AccountOpened>(ingress_config(ACCOUNTS_BASE_PATH));
let accounts_source = accounts_ingress.source();
let tx_ingress = http_ingress::<LedgerEntry>(ingress_config(TX_BASE_PATH));
let tx_source = tx_ingress.source();
FlowApplication::builder()
.with_config_file(CONFIG_FILE)
.with_log_level(LogLevel::Info)
.with_http_ingress(accounts_ingress)
.with_http_ingress(tx_ingress)
.run_blocking(flow::build_flow(accounts_source, tx_source))?;Each call to http_ingress::<T>(...) returns an HttpIngress<T> bundle. That bundle owns two distinct things at once.
- It owns the HTTP surface, which is the route, the deserialization, the readiness wiring, and the ingress telemetry.
- It also owns a typed
HttpSourceTyped<T>, which is the source the flow reads from.
The runner pulls the typed source out first with .source(), and then moves the full ingress bundle into the FlowApplication builder. The result is that flow::build_flow(...) only sees typed sources. It never sees HTTP routes, port numbers, or readiness primitives. The HTTP surface lives entirely in runner.rs, and the flow lives entirely in flow.rs.
That split matters. It keeps flow! source-oriented and lets the same flow definition be hosted under any input boundary that can produce a typed source. In a test, you can drive the same build_flow(...) from in-memory sources without the HTTP server. In production, the HTTP shell adds metrics, an /api/topology endpoint, and the rest of the operational surfaces around the same pipeline.
The example deliberately exposes two HTTP ingress paths, one for accounts and one for transactions, even though both are bank events. Each ingress carries one typed payload, which keeps the JSON contract explicit and the metrics per-endpoint clean. The flow is also free to treat each typed source differently. In this example, accounts feed a join catalog and transactions feed the forward path.
Infinite sources
In Getting Started, the source was a finite in-memory collection. The pipeline ran, drained, and exited. Service-shaped flows need a different kind of source. They need a source that stays alive, accepts events as they arrive, and never signals end of stream on its own.
The piggy bank example uses async_infinite_source! for both inputs:
accounts = async_infinite_source!(AccountOpened => accounts_source);
tx = async_infinite_source!(LedgerEntry => tx_source);The async_infinite_source! macro wraps an asynchronous producer that runs until the application shuts down. That is the right shape for an HTTP feed, an upstream queue consumer, or any other source that does not end on its own.
Downstream of a source, transforms, accumulators, joins, and so forth don’t really care whether a source is finite, like a flat file or a CSV file, or infinite, like a streaming RESTful endpoint. All they see are typed events. Essentially, when a payload is posted against a RESTful endpoint, that post is accepted and written to ObzenFlow’s journal like any other event in the flow.
The only real practical difference between a finite and an infinite source is that a finite source will eventually end naturally when ObzenFlow processes the last entry, whereas an infinite source is never-ending and only stops listening and processing when the program ends.
In this example, the two sources play different roles in the flow. The accounts source feeds the reference catalog inside the join. The tx source is the live transaction stream. Both arrive over HTTP, but the flow treats them differently because of how they are wired downstream, which is the subject of the next step.
Step 3. Validating transactions with stream-table joins
The posted stage is where the most important new concept in this tutorial lives. It is the first time you have seen a join! in this tutorial series, and it is the stage that turns a raw HTTP request into a trusted, enriched fact.
Here is the join from examples/http_ingestion_piggy_bank_demo/flow.rs:
posted = join!(
catalog accounts: AccountOpened,
LedgerEntry -> PostedEntry => joins::inner_live(
|account| account.account_id.clone(),
|entry| entry.account_id.clone(),
|account, entry| PostedEntry {
account_id: entry.account_id,
owner: account.owner,
kind: entry.kind,
amount_cents: entry.amount_cents,
initial_balance_cents: account.initial_balance_cents,
note: entry.note,
},
)
);Read the join in three layers.
- The
catalog accounts: AccountOpenedline declares the reference side of the join. Theaccountssource declared in Step 2 is held inside this join as a continuously updated catalog ofAccountOpenedevents, keyed by account id. That is whyaccountsdoes not appear in the topology block. Catalog sources are not forward edges in the pipeline. They are side inputs to the join stage. - The
LedgerEntry -> PostedEntryline declares the forward contract. The join receivesLedgerEntryevents on its forward edge and emitsPostedEntryevents downstream. Thetx |> postedline in the topology block is what connects the transaction stream to that forward edge. - The
joins::inner_live(...)helper supplies three things. The first closure extracts the join key from the catalog side, which isaccount.account_id. The second closure extracts the join key from the stream side, which isentry.account_id. The third closure is the merge function. It runs once per matched pair and produces thePostedEntrythat downstream stages will see.
The “live” half of inner_live is what makes this a stream-table join rather than a static lookup. The catalog updates continuously as new AccountOpened events arrive. A transaction submitted before its account exists finds no match and produces no PostedEntry. Once the account is opened, subsequent transactions for that account will match and post.
The “inner” half is exactly what it sounds like. Only matched pairs flow downstream. Unmatched transactions produce no posted event. That is a design choice. In a more elaborate service you might keep unmatched entries as their own RejectedEntry event and surface them on a separate sink, but the canonical example keeps the join inner so the rest of the pipeline can assume every event it sees is already validated.
The most ObzenFlow-shaped move in this example is that validation doesn’t live in an if statement. It lives in the join.
A LedgerEntry is a request from the edge, and a PostedEntry is journaled only if the join was successful. The structural fact that downstream stages only ever see PostedEntry types is what keeps the rest of the pipeline simple. Domain invariants are enforced by the type system and by the topology, not by handler code.
It is worth repeating the stages and topology code from earlier so it really makes sense how joins are wired up. Below you can see how accounts and transactions slot into the join and then how the join is wired into the flow itself.
stages: {
accounts = async_infinite_source!(AccountOpened => placeholder!());
tx = async_infinite_source!(LedgerEntry => placeholder!());
posted = join!(
catalog accounts: AccountOpened,
LedgerEntry -> PostedEntry => placeholder!()
);
checkbook = stateful!(PostedEntry -> CheckbookSnapshot => placeholder!());
printer = sink!(CheckbookSnapshot => placeholder!());
},
topology: {
tx |> posted;
posted |> checkbook;
checkbook |> printer;
}One of the things to pay attention to that might be a little new is that when you are reading an ObzenFlow, it is much more important to think about types.
Above, accounts are a table that is being hydrated through a RESTful endpoint. Then, if you look at the topology definition, transactions flow into the posted stage, which is a join. We wire the account catalog into the catalog position of the join, then that second line of the join specifies the stream-side, which is read as input type -> output type => the handler that you will define later on.
The rest of the example is wired together very similarly to the Getting Started character transformer that we saw in the first tutorial.
Step 4. Maintaining bank balances with stateful projections
The checkbook stage is the projection, which receives every PostedEntry, folds it into per-account state, and emits a CheckbookSnapshot for downstream readers.
checkbook = stateful!(
PostedEntry -> CheckbookSnapshot =>
Checkbook::new().with_emission(EmitAlways)
);Two details are worth comparing expanding on.
Choosing an emission strategy: EmitAlways instead of OnEOF
In char_transform, the stateful stage was a built-in Reduce accumulator, which exposes .emit_on_eof() as a one-line shortcut for .with_emission(OnEOF::new()). The reducer accumulated text chunks and emitted exactly one TransformedText event when the upstream source signalled end of stream. That is the right strategy for a finite batch.
A live service does not have an end of stream, and it has read-side consumers that want to see the latest state after every input. The piggy bank therefore picks a different strategy, EmitAlways. Every PostedEntry produces a fresh CheckbookSnapshot immediately, and a downstream sink, database table, dashboard, or another flow can subscribe to that stream and stay current with no extra plumbing.
Because Checkbook is a hand-written StatefulHandler rather than one of the built-in accumulators, it reaches for the general form .with_emission(EmitAlways) instead of a shortcut method.
- The shortcut methods (
.emit_on_eof(),.emit_always(),.emit_every_n(n),.emit_within(duration)) live on the built-in accumulators. - The
.with_emission(...)form is the underlying builder that every stateful handler accepts.
Both forms produce the same wired stage. The split between accumulation and emission is identical to Getting Started. The handler describes how to update the state, and the emission strategy describes when that state becomes an output event.
Why not a built-in accumulator?
The framework ships with Reduce, GroupBy, Conflate, and TopN. Each covers a common shape. Reduce folds every event into one global accumulator. GroupBy keys events by a payload field and keeps one accumulator per key. Conflate keeps the latest value, and TopN keeps the highest-ranked entries. Any of those would let you write a stateful stage as a one-liner, the way Getting Started did with char_transform. A checkbook does not quite fit any of them, and the reasons are worth naming so you know when to reach past the built-ins in your own services.
- The first reason is that each account has a multi-field ledger. The current balance, the running totals for credits and debits, and the list of past transactions all move on every event, and the update branches on
EntryKind. A credit adds to one running total and to the balance. A debit adds to a different running total and subtracts from the balance.ReduceandGroupByare designed around folding into a single accumulator value per scope, so encoding that branch inside a one-liner closure quickly becomes harder to read than a named handler. - The second reason is that the opening balance is not
Default::default(). It arrives inside the very firstPostedEntryfor that account, carried in theinitial_balance_centsfield, and the per-account ledger needs to be seeded from that value when the account is first seen. The keyed built-in,GroupByTyped, requires its per-group state to beDefault, so it cannot express the idea that the seed for an account comes from the first event for that account. (ReduceTypeddoes take an explicitinitialvalue, but it is wrong-shaped for a different reason, covered next.) - The third reason is the shape of emission.
ReduceTypedemits the state as a single typed domain event, which is exactly how Getting Started’schar_transformpublishesTransformedText. The trouble for a checkbook is thatReduceTypedhas one global accumulator, so if you put the per-accountHashMapinside it, every emission would dump the entire bank state as one event. The natural shape for per-account aggregation isGroupByTyped, butGroupByTyped::emitwraps each per-key state in a{key, result}JSON envelope tagged with the state’s event type, and it walks every known key on every emission tick. Paired withEmitAlways, that would re-emit every account on everyPostedEntry, and downstream readers would have to peel the wrapper to find a snapshot. A customStatefulHandlerlets you emit a rawCheckbookSnapshotfor only the affected account on each event. Thelast_snapshotslot insideCheckbookis what records which slice changed, andemitproduces exactly one snapshot event for the account that just moved.
Any one of those constraints could be worked around with the built-ins. All three together push the example over the line where a hand-written StatefulHandler is clearer than a stacked closure.
From a reduce closure to a StatefulHandler trait
Getting Started used typed_stateful::reduce with a closure. The piggy bank uses a hand-written type that implements the StatefulHandler trait. Both forms are valid. The custom handler comes into play when the state and emission logic need names of their own.
The full handler lives in examples/http_ingestion_piggy_bank_demo/handlers.rs. The accumulation logic itself is straightforward arithmetic on cents:
match entry.kind {
EntryKind::Credit => {
ledger.total_credits_cents = ledger
.total_credits_cents
.saturating_add(entry.amount_cents as i64);
ledger.current_balance_cents = ledger
.current_balance_cents
.saturating_add(entry.amount_cents as i64);
}
EntryKind::Debit => {
ledger.total_debits_cents = ledger
.total_debits_cents
.saturating_add(entry.amount_cents as i64);
ledger.current_balance_cents = ledger
.current_balance_cents
.saturating_sub(entry.amount_cents as i64);
}
}The rule of thumb is simple. Reach for typed_stateful::reduce and a closure when the accumulator is a single value or a small struct. Reach for the StatefulHandler trait when you need a typed state, multiple state fields, or non-trivial emission logic.
The sink in this example renders the snapshot to the console so you can see the shape directly. In a larger service, the same CheckbookSnapshot stream could write rows into a database, push to a queue, update a search index, or feed another flow. The projection stage is where read-side state is materialized once and reused everywhere.
Step 5. Operational surfaces, evidence, and replay
In this final section, we’ll explain a little more of the operational surface area of an ObzenFlow program, which includes all sorts of topics from journal manifests to middleware and metrics, and just simply how a program runs end-to-end, and what is included out of the box.
While the service is running, you can inspect the operational surfaces the framework gives you for free. They come from FlowApplication once the HTTP server is enabled, which provides all of the operational machinery that most programs look for, like metrics and middleware.
Below is a collection of endpoints that the example provides to you for free, excluding the endpoints that you define yourself. You can see that we expose details about metrics, the topology, and health checks. This comes wired right out of the box.
curl http://127.0.0.1:9090/metrics
curl http://127.0.0.1:9090/api/topology
curl http://127.0.0.1:9090/api/bank/accounts/health
curl http://127.0.0.1:9090/api/bank/tx/healthWhen you stop the example with Ctrl-C, ObzenFlow leaves a run archive under:
target/http-ingestion-piggy-bank-demo-logs/flows/<run-id>/Every stage has its own journal file, including the two HTTP sources, the join, the stateful stage, and the sink. The run_manifest.json records the flow structure, stage metadata, and journal mappings for the run. That archive is the durable evidence trail for every account fact, ledger entry, posted entry, and emitted checkbook snapshot.
Replay reads the archive as input. The framework swaps the live HTTP sources for the archived source journals and reruns the flow against the same historical inputs.
RUST_LOG=info cargo run -p obzenflow --example http_ingestion_piggy_bank_demo --features obzenflow_infra/warp-server -- \
--replay-from target/http-ingestion-piggy-bank-demo-logs/flows/<run-id>The point worth carrying forward from this tutorial is that the same journal-and-replay model from char_transform applies cleanly to a live HTTP-backed service. The input boundary moved out to HTTP, but every stage still produces a durable journal, and the same archive shape that the character transformer produced is what the piggy bank produces.
/metrics, /api/topology, and the /api/flow/* control plane unauthenticated for local development. To exercise the recommended bearer-token setup, use the obzenflow.auth.toml config bundled with the example.Step 6. Stop the service, then resume it
Replay reruns a finished archive and drains. A live service needs one more verb. When you pressed Ctrl-C earlier, the run was recorded as cancelled, and every fact it had processed was already durable in its journals. --resume-from continues that same run rather than starting a new one, and a cancelled archive is a valid resume input with no extra flags.
Point it at the same run directory you replayed in Step 5:
RUST_LOG=info cargo run -p obzenflow --example http_ingestion_piggy_bank_demo --features obzenflow_infra/warp-server -- \
--resume-from target/http-ingestion-piggy-bank-demo-logs/flows/<run-id>Resume runs in two phases.
- First it catches up. The runtime re-reads the recorded journals through the same reconstruction engine that replay uses, re-folding every posted entry back into the checkbook state. The recorded snapshots render again, each marked as a replayed delivery. While catch-up runs, the two HTTP endpoints return
503with aRetry-Afterheader instead of admitting fresh input, so a new command can never interleave with history. With an archive this small, that window lasts only a moment. - When catch-up reaches the end of the recorded history, each source crosses to live and the endpoints open. From that point the run accepts new input, and the flow behaves as though it had never stopped.
Post one more transaction for Alice.
curl -XPOST http://127.0.0.1:9090/api/bank/tx/events \
-H 'content-type: application/json' \
-d '{"event_type":"bank.ledger_entry","data":{"account_id":"acct-1","kind":"Credit","amount_cents":100,"note":"back online"}}'The snapshot continues from the recorded balance. Alice’s account was at $11.51 when the service stopped, so this credit moves it to $12.51, and her transaction table now shows three entries. The state came from re-folding the recorded journals during catch-up, and the new entry extends that same history.
The proof lives in the resumed run’s archive rather than in the console output. The resumed run writes its own run_manifest.json under a new run directory, and that manifest carries a resume block recording the run’s lineage.
"resume": {
"resumed_from": "target/http-ingestion-piggy-bank-demo-logs/flows/<run-id>",
"resume_generation": 1
}The resumed journals extend the recorded history under the same run identity. That closure is the point of the design. The resumed run can itself be replayed for verification, and if it is interrupted in turn, it can be resumed again.
During catch-up the console sink re-renders the recorded snapshots. That is safe here because a console projection is idempotent, and the framework’s console sink declares itself so. A sink that wrote to a non-idempotent external destination, such as a POST to a payment processor, would be refused at resume until it declared its delivery safety or you explicitly accepted duplicate delivery. The runtime fails loud rather than silently delivering twice.
Next steps
You have now seen ObzenFlow handle write-side HTTP traffic, join live events against reference data, produce a running projection from journaled facts, and resume an interrupted service from its own journals. The same pattern applies to webhook receivers, command-ingestion services, operational read models, and any workflow where facts accumulate over time.
When you are ready to keep going, Run Live AI Inference from a Real Endpoint takes the same flow model into a model-backed workload. That tutorial adds live HTTP pull, token budgeting, ai_map_reduce!, and prompt evidence on top of the flow model you have already seen.
The philosophy is worth reading once the tutorials feel natural, if you want the architectural rationale behind why ObzenFlow is shaped this way.