Model Bank Transactions as a Flow

A journal-backed checking-account tutorial with HTTP ingress, joins, and projections.

Use the piggy bank example to see how ObzenFlow fits into a real write-side service: HTTP commands in, journal-backed facts through the flow, and a projection out the other side.

Info · Complete the getting started tutorial first

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:

  1. Understand the domain.
  2. Name the commands and events.
  3. Turn those facts into Rust types.
  4. Sketch the flow topology.
  5. 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 deeper pass on this modelling workflow, see Domain Modelling in Practice. The series covers why modelling before code matters, how Event Storming surfaces events, commands, policies, and aggregates, and how those outputs turn into Domain-Driven Design structure.

For a checking-account service, the list of facts is short enough to write from memory. Three are visible at the service boundary.

  1. An account is opened with an account id, an owner, and an initial balance.
  2. A ledger entry is submitted as a credit or a debit for an account.
  3. A checkbook snapshot is observed by anyone reading the account state.

A fourth fact lives entirely inside the flow.

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

Tip · Commands are requests, events are facts

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.

  1. The first is the source declaration, which is a streaming source. Both sources use async_infinite_source!, the live counterpart to the source! 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.

  2. The second is the way accounts participates in the topology. The transaction stream flows forward through tx |> posted. The accounts source 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 obzenflow

Run the piggy bank example:

RUST_LOG=info cargo run -p obzenflow --example http_ingestion_piggy_bank_demo --features obzenflow_infra/warp-server

The 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: 2

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

Tip · Modeling error flows is critical

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.

  1. It owns the HTTP surface, which is the route, the deserialization, the readiness wiring, and the ingress telemetry.
  2. 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.

Tip · Why two ingress endpoints and not one?

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.

  1. The catalog accounts: AccountOpened line declares the reference side of the join. The accounts source declared in Step 2 is held inside this join as a continuously updated catalog of AccountOpened events, keyed by account id. That is why accounts does not appear in the topology block. Catalog sources are not forward edges in the pipeline. They are side inputs to the join stage.
  2. The LedgerEntry -> PostedEntry line declares the forward contract. The join receives LedgerEntry events on its forward edge and emits PostedEntry events downstream. The tx |> posted line in the topology block is what connects the transaction stream to that forward edge.
  3. The joins::inner_live(...) helper supplies three things. The first closure extracts the join key from the catalog side, which is account.account_id. The second closure extracts the join key from the stream side, which is entry.account_id. The third closure is the merge function. It runs once per matched pair and produces the PostedEntry that 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.

Tip · Validation by joining, not by branching

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.

  1. 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. Reduce and GroupBy are 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.
  2. The second reason is that the opening balance is not Default::default(). It arrives inside the very first PostedEntry for that account, carried in the initial_balance_cents field, 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 be Default, so it cannot express the idea that the seed for an account comes from the first event for that account. (ReduceTyped does take an explicit initial value, but it is wrong-shaped for a different reason, covered next.)
  3. The third reason is the shape of emission. ReduceTyped emits the state as a single typed domain event, which is exactly how Getting Started’s char_transform publishes TransformedText. The trouble for a checkbook is that ReduceTyped has one global accumulator, so if you put the per-account HashMap inside it, every emission would dump the entire bank state as one event. The natural shape for per-account aggregation is GroupByTyped, but GroupByTyped::emit wraps 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 with EmitAlways, that would re-emit every account on every PostedEntry, and downstream readers would have to peel the wrapper to find a snapshot. A custom StatefulHandler lets you emit a raw CheckbookSnapshot for only the affected account on each event. The last_snapshot slot inside Checkbook is what records which slice changed, and emit produces 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.

Tip · The same projection feeds many destinations

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/health

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

Warning · Protecting the operational surfaces
The default configuration leaves /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.

  1. 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 503 with a Retry-After header 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.
  2. 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.

Tip · Why this sink is safe to resume

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.