Run Live AI Inference from a Real Endpoint

A journaled AI tutorial with token budgeting, chunking, accumulation, and replay.

Use the Hacker News digest example to see how ObzenFlow wraps live HTTP data and model inference in deterministic flow structure, durable journals, and replayable evidence.

Info · Start with the earlier tutorials first

This tutorial assumes you have already worked through Getting Started and Model Bank Transactions as a Flow. Getting Started explains FlowApplication, flow!, stage types, topology, journals, and replay. The bank tutorial adds HTTP ingress, joins, and projections. This page carries that same flow model into a model-backed workload without losing determinism, inspectability, or replay.

This example pulls live data from the Hacker News Firebase API and runs AI inference through a local Ollama model to produce a journaled, replayable digest. This is the kind of program you would reach for to build an AI summarization service, a content triage pipeline, or a retrieval-and-enrichment flow, any model-backed workload that still needs deterministic scaffolding around the model call.

It’s worth pointing out that this example demonstrates how ObzenFlow handles some of the trickiest parts of running live AI inference for you. Reliability comes from built-in circuit breakers that contain a failing or overloaded model instead of letting it stall the whole flow, and token-budgeted chunking splits a large input into pieces that fit inside the model’s context window before anything is sent to it. Out of the box, it connects to Ollama for the inference step itself.

Let’s get started by briefly outlining the program. We’ll grab a number of posts from Hacker News, then run them through local inference to summarize trending topics. This type of summarization is a bread and butter example for local inference.

Step 0. Designing the digest flow

This example started as a little tool I built for myself to test my ideas behind live durable inference, and it was one of the motivating designs in my head before ObzenFlow existed. It grabs trending Hacker News stories, runs all of them through inference, and produces a summary.

This is exactly the kind of workload you might use to summarize support tickets, incidents, or orders, taking a finite or infinite stream of data and working out how to smooth it and chunk it. For instance, you could follow the same kind of approach to emit summaries of support tickets. You could even take it to the next level by leveraging joins to enrich support tickets with detailed customer information and historical support requests. Really, the sky is the limit. This is the basic shape of such a workflow.

Let’s jump right into domain modeling, which we’ve already covered in the previous two tutorials. (If this is too abrupt, please go and re-review those before continuing.)

The facts in this flow

For a Hacker News digest, the interesting things that happen over time are short to list.

  1. A story is fetched from the Hacker News API.
  2. A story is formatted into a stable, display-ready record.
  3. The formatted stories are accumulated into one batch.
  4. A chunk of that batch is summarized by the model, once per token-budgeted group.
  5. The chunk summaries are reduced into one final digest.

This might be a strange way to think about ingesting data, but in terms of auditability, it’s actually critical.

One of the challenges teams will face over time is that inference is very good at summaries, but it’s also prone to hallucinating. Journaling how input data is fetched and formatted, accumulated, chunked, and submitted for inference is the difference between being able to pass an audit and being able to respond intelligently during an incident or a production investigation.

Those facts become the event types in the flow. Here they are in trimmed form, drawn from examples/hn_ai_digest_demo/domain.rs and flow.rs.

// Raw item from the Hacker News Firebase API. Almost every field is
// optional because deleted and malformed items do show up in practice.
struct HnStory {
    id: HnStoryId,
    title: Option<String>,
    url: Option<String>,
    by: Option<String>,
    score: Option<u32>,
    descendants: Option<u32>,
}

// A cleaned, display-ready story.
struct FormattedStory {
    id: HnStoryId,
    title: String,
    url: String,
    author: String,
    points: u32,
    comments: u32,
}

// The accumulated batch the model will work from.
struct HnTopStories {
    stories: Vec<FormattedStory>,
}

// One chunk's partial summary.
struct HnDigestGroupSummary {
    output_markdown: String,
}

// The final evidence record. It carries the inputs, the exact prompts,
// the partial summaries, and the final output together. Shown in full in Step 3.
struct HnDigestSummary { /* ... */ }

Sketching the flow

With the facts named, you can wire the entire flow before writing a single handler. This is the same placeholder-first workflow you saw in Getting Started and in examples/char_transform_skeleton.rs. Every handler is placeholder!(), but the event types and the topology are real, so the skeleton compiles and you can inspect the shape in ObzenFlow Studio before any business logic exists.

flow! {
    name: "hn_ai_digest_demo",
    journals: disk_journals(PathBuf::from("target/hn-ai-digest-logs")),
    middleware: [],

    stages: {
        hn_stories     = async_source!(HnStory => placeholder!());
        formatter      = transform!(HnStory -> FormattedStory => placeholder!());
        batch          = stateful!(FormattedStory -> HnTopStories => placeholder!());
        digest         = transform!(HnTopStories -> HnDigestSummary => placeholder!());
        digest_summary = sink!(HnDigestSummary => placeholder!());
    },

    topology: {
        hn_stories |> formatter;
        formatter  |> batch;
        batch      |> digest;
        digest     |> digest_summary;
    }
}

The topology is worth pausing on. Stories come in, get formatted, get accumulated into one batch, get turned into a digest, and leave through the sink.

Walking the skeleton from left to right, here is the thinking behind each placeholder.

  • Ingestion (hn_stories). The first question is where the data comes from and whether it ever ends. Here, it is the Hacker News API. At design time you do not care whether the bytes arrive from the live endpoint or a local mock, only that the stage produces HnStory events.

  • Formatting (formatter). Raw APIs are messy. Fields are optional, some entries are deleted, and titles or URLs can be missing. Rather than let that mess leak into the rest of the flow, you decide to normalize each item into a clean FormattedStory right away.

  • Batching (batch). Now a design decision that shapes everything after it. Do you want one summary per story, or one digest over all of them? You want the digest, so the per-story stream has to become a single collection before the model sees it. That is a stateful accumulation that folds many FormattedStory events into one HnTopStories value, emitted once.

  • Digest (digest). This is the heart of the program, and the striking thing at design time is how little you have to decide. A batch of stories goes in, one HnDigestSummary comes out. You sketch it as a single placeholder transform and move on. The genuinely hard questions, how to fit a large batch through a model with a fixed context window, are implementation details of this one stage.

  • Output (digest_summary). Finally, where does the result leave the system? In this example the sink prints to the console, but the same HnDigestSummary could go to email, Slack, a database, or another flow. The sink is the flow boundary, and choosing it is the last design decision rather than the first.

Tip · Think in types first before implementation

Because the skeleton compiles, you do not have to hold all of this in your head at once. You can wire these five placeholders, open the flow in ObzenFlow Studio, and look at the real shape before committing to a single line of behaviour. If a stage sits in the wrong place or a type does not line up, you find out now, while the flow is still cheap to change.

Next let’s run the example end-to-end so you can get a feel for what it does.

Step 1. Run the digest end to end

The fastest way to build a mental model of this flow is to run it before reading the code, exactly as you did with the bank service. The default run needs no network connection for its data, because the example ships with a built-in mock Hacker News server. It does still need a local model for inference, so set up Ollama first.

Set up Ollama

brew install ollama
ollama serve
ollama pull llama3.1:8b

You can use any model you like, but a small one such as llama3.1:8b keeps the demo fast. If you would rather use a hosted provider instead of Ollama, set HN_AI_PROVIDER=openai and OPENAI_API_KEY=... before running. When you use a hosted provider, your prompts and the story text are sent to that provider, so be deliberate about it.

Run the mock example

If you have already cloned the framework repo from the earlier tutorials, stay in that working directory. Otherwise clone it now.

git clone https://github.com/ObzenFlow/obzenflow
cd obzenflow

Run the example with its default settings. The data comes from the built-in mock server, and inference runs against your local Ollama model.

cargo run -p obzenflow --example hn_ai_digest_demo --features "http-pull ai"

The mock server only removes the network dependency on the data side. The model call is still real, which is the point.

The first thing you see is the startup banner that with_presentation adds, which records the configuration for the run.

HN AI Digest Demo
=================
Fetch top HN stories, then generate a markdown digest via Rig-backed LLM transforms.

  mode: mock
  base_url: http://127.0.0.1:49667/
  max_stories: 60
  poll_timeout: 120s

  AI:
    provider: ollama
    model: llama3.1:8b
    token_estimator: Heuristic
    token_estimator_backend: tiktoken
    token_estimator_fallback_reason: model_not_supported_by_tokenizer
    token_estimator_fallback_detail: no tiktoken encoding for model 'llama3.1:8b'
    context_window: 128000

  group_budget_tokens: 2500
  group_max_stories: 10
  source_rate_limit: 10 events/sec

Even the choice of token estimator is recorded, along with the reason for it. The flow asked for an exact tokenizer, the tiktoken backend had no encoding for llama3.1:8b, so it fell back to a byte-length heuristic and wrote down exactly why. Nothing about the setup is left implicit.

When the flow finishes, the sink prints the evidence trail and the final digest. The header carries the run metadata and the exact prompts.

HN AI Digest — Summary
======================
mode: mock
base_url: http://127.0.0.1:49667/
ai_provider: ollama
ai_model: llama3.1:8b
token_estimator: Heuristic
stories_fetched: 55
groups: 6 (budget_per_group: 2500)

Chat prompt (system)
--------------------
You write concise, skimmable Hacker News digests from a
list of headlines + URLs. Be neutral, avoid hype, and do
not invent facts beyond what the titles imply.

Chat prompt (user)
------------------
Write a concise Markdown digest of the following Hacker
News chunk summaries.

- Do not invent facts that are not implied by the titles.
- Start the response immediately with "## What's topical
  today" (no intro).
- Include: Thesis, Themes (cite story numbers), Notable
  stories, Watch.
- Avoid generic wrap-ups.

Chunk summaries:
(the six group summaries the map step produced are appended here)

Everything in that header is recorded as part of the run. The provider, the model, the token estimator, the chunking budget, the number of groups, and the exact system and user prompts are all fields on a typed event. You’ll learn a little bit more about this in the final section of the tutorial when we cover operations and dig into topics like the manifest.

After the header, the sink prints the input stories it fed to the model, and then the digest the model produced. The block below comes from one real run against llama3.1:8b, wrapped here so it fits without scrolling. Because the data side is mocked the stories are synthetic, and because the model is non-deterministic your wording will differ.

Input data (stories)
--------------------
1. Rust: tooling update #100 (180 points, 0 comments) by author1
    https://example.com/hn/100
2. AI: benchmark #101 (178 points, 1 comments) by author2
    https://example.com/hn/101
3. Databases: postmortem #102 (176 points, 2 comments) by author3
    https://example.com/hn/102
...

Output (markdown)
-----------------
## What's topical today

### Thesis
The current Hacker News discussions revolve around various
topics in computer science, including programming languages,
AI, databases, web development, operating systems, security,
browsers, Linux, compilers, robotics, and distributed
systems.

### Themes

* Programming Languages (1, 7): Updates to the Rust
  toolchain and other language-related developments. (#1)
* AI and Machine Learning (2): Benchmarking different
  frameworks' performance and in-depth looks at AI
  techniques. (#2, #48)
* Databases (3, 12, 21, 28, 49): Postmortem analyses of
  database failures, new database systems, and migration
  guides for distributed systems. (#3, #12, #21, #28, #49)
* Web Development (5): Launch of a new browser with
  improved features. (#5)
* Operating Systems (6, 10, 51): Guides for migrating to
  Linux, benchmark results for distributed systems, and
  case studies on Linux development. (#6, #10, #51)
* Security (4, 13, 22, 31): Deep dives into security
  vulnerabilities, migration guides from one framework to
  another, and benchmarking security performance.
  (#4, #13, #22, #31)

### Notable stories

* A new version of the Rust toolchain has been released
  with various updates. (#1)
* An AI benchmark has been published to compare different
  frameworks' performance. (#2)
* A postmortem analysis of a database failure has been
  shared. (#3)
* A new browser has been launched with improved features. (#5)
* A guide for migrating to Linux is available online. (#10)

### Watch
Keep an eye on the discussions around Rust's tooling
ecosystem, AI benchmarks, and distributed systems design.

The wording of the output will differ on every run, because model output is not deterministic. That is precisely why the flow records the inputs and the prompts alongside the output. The non-deterministic part is captured next to everything that produced it.

Tip · ObzenFlow's philosophy on AI

The prevailing narrative around AI is that LLMs will take over all of software development and make developers irrelevant. ObzenFlow’s philosophy is different.

LLMs will become a utility, like a faucet or, better yet, like the electrical grid. They will be available everywhere, cheaply, and the interesting question will not be whether to use them but how to use them responsibly in systems that need to explain what they did.

It is nearly inevitable that a large share of software will wrap around AI, and AI will become an implementation detail. That’s our approach here.

Switch to the live Hacker News feed

When you are ready to point the flow at the real API, set HN_LIVE=1. The decoder first fetches v0/topstories.json, then follows it with one v0/item/{id}.json request per story until it reaches HN_MAX_STORIES.

HN_LIVE=1 cargo run -p obzenflow --example hn_ai_digest_demo --features "http-pull ai"

In live mode the header shows mode: live and a real base URL, and stories_fetched reflects whatever the front page held at the time. Everything downstream of the source is identical to the mock run.

Warning · Be responsible with the live Hacker News feed

Larger runs create meaningfully more upstream traffic, because each story is a separate item request. Use the mock server while you iterate on prompts and models, and start small when you switch to live mode.

This example puts an explicit rate limiter on the Hacker News source so live fetches stay polite. The default is 10.0 events per second, and you can lower it from the command line with HN_SOURCE_RATE_LIMIT. That limiter complements the live HttpPullSource, which already treats 429 responses as rate limiting, honours Retry-After when it is present, and applies retry with backoff for transient failures.

This is double entry accounting for AI. The input and the output are captured together, as part of the same journaled flow. If the model changes, you can see what changed. If the input changes, you can see how. If you want to replay this run a year from now, the framework allows you to point ObzenFlow to a previous run and replay from a point in time.

Step 2. Understanding live inference flows

The earlier tutorials used bounded input and a single ingestion path. What sets this example apart is that it pulls from an external API, runs inference through a local model, and produces a structured digest from potentially hundreds of stories.

The flow shape is still the same FlowApplication and flow! pattern you already know. Where things really change is that now we’re composing a more complex flow from the building blocks ObzenFlow provides, including:

  • async sources that pull from a live RESTful endpoint
  • a transform stage to convert raw input into a type that the rest of the flow can understand
  • a batch stage that accumulates formatted stories into a single batch
  • the purpose of the example, which is the AI map reduce stage that handles mapping, reducing, chunking, and middleware

Finally, we’ll wire that stage into a sink for output.

Most frameworks put the LLM call at the center. ObzenFlow follows a different philosophy: live-inference LLM calls are not any more special than other stages, except that live inference tends to produce non-deterministic, unstructured output. It also involves operational challenges such as token estimation, context window size, batching, and chunking.

ObzenFlow is not an AI framework. It’s a durable execution runtime. It happens to coordinate durable execution of inference, providing resilience for the most common use case. If you need something more exotic, you’ll likely have to drop down into custom handler code, which is completely supported but beyond this tutorial’s scope.

Let’s get started, and we’ll cover all the stages involved.

Reading the stages

Open examples/hn_ai_digest_demo/flow.rs to follow along. The stages block has five entries.

stages: {
    // Ingest: pull stories from the HN API, rate-limited to stay polite.
    hn_stories = async_source!(HnStory => HttpPullSource::new(decoder, config), [
        RateLimiterBuilder::new(source_rate_limit).build()
    ]);

    // Format each raw story into a stable, display-ready event.
    formatter = transform!(HnStory -> FormattedStory => formatter);

    // Accumulate every formatted story into one batch, emitted at EOF.
    batch = stateful!(FormattedStory -> HnTopStories => digest_seed);

    // The whole map-reduce: chunk by token budget, summarize each chunk
    // through the LLM, then reduce the chunk summaries into one digest.
    digest = ai_map_reduce!(
        HnTopStories -> HnDigestSummary => {
            map:    [FormattedStory] -> HnDigestGroupSummary => map_llm_handler,
            reduce: (HnTopStories, [HnDigestGroupSummary]) -> HnDigestSummary => digest_llm_handler,
        },
        chunking: by_budget { /* shown in Step 3 */ },
        middleware: {
            map: ai_circuit_breaker(),
            reduce: ai_circuit_breaker(),
        }
    );

    // Print the full evidence trail and the final markdown digest.
    digest_summary = sink!(HnDigestSummary => sinks::console(format_digest_summary_for_console));
}

Each stage maps to one of the facts from Step 0.

StageWhat it does
hn_storiesPulls stories from the Hacker News Firebase API using HttpPullSource, rate-limited so live fetches stay polite. In live mode it walks topstories.json and then fetches each item; in the default mode it reads from the built-in mock server.
formatterTurns each raw HnStory into a typed FormattedStory with a stable text representation. Stories that fail to format are dropped rather than crashing the flow.
batchAccumulates every formatted story into one HnTopStories value and emits it once, when the source signals end of input. This is the same accumulate-then-emit pattern as char_transform.
digestThe ai_map_reduce! stage. It chunks the batch by token budget, summarizes each chunk through the model, and reduces the chunk summaries into the final HnDigestSummary. Steps 3 and 4 open it up.
digest_summaryThe sink prints the final HnDigestSummary, including the evidence trail you saw in Step 1.

The forward path connects them in order, which is the topology you already sketched in Step 0. There are no oversize branches or feedback edges in the topology, because the oversize handling lives inside the digest stage rather than as separate stages.

Tuning the scaling knobs

A couple of environment variables control how much work enters the flow. The example resolves them in DemoConfig::from_env() at startup, then hands the resolved DemoConfig to the flow.

// examples/hn_ai_digest_demo/config.rs, in DemoConfig::from_env()
let max_stories       = env_var_or::<usize>("HN_MAX_STORIES", DEFAULT_HN_MAX_STORIES)?;            // default 60
let source_rate_limit = env_var_or::<f64>("HN_SOURCE_RATE_LIMIT", DEFAULT_HN_SOURCE_RATE_LIMIT)?;  // default 10.0

max_stories is passed to the source decoder (hn_story_decoder(base_url, max_stories) in flow.rs), and source_rate_limit configures the rate limiter on the source stage. HN_MAX_STORIES changes how many items the source pulls into the flow, and HN_SOURCE_RATE_LIMIT changes how aggressively the source is allowed to pull them. You can fetch 10 stories for a quick run or a few hundred for a much larger digest, and you can keep the live source conservative while you do it. The chunking, accumulation, and final synthesis stages absorb the change without any edit to the topology.

# A small live run, then a larger and gentler one.
HN_MAX_STORIES=10 HN_LIVE=1 cargo run -p obzenflow --example hn_ai_digest_demo --features "http-pull ai"
HN_SOURCE_RATE_LIMIT=5.0 HN_MAX_STORIES=200 HN_LIVE=1 cargo run -p obzenflow --example hn_ai_digest_demo --features "http-pull ai"

Whether the source fetched 10 stories or 200, the map-reduce structure stays the same. You’ll learn how the framework takes care of the hardest details for you in the next section.

Step 3. Inside the digest stage

Now it is time to open up the one stage that makes this more than a wrapper around a single LLM call. You already saw this stage in Step 2, but it is worth putting back in front of you, because the rest of this step is a guided tour of its parts.

digest = ai_map_reduce!(
    HnTopStories -> HnDigestSummary => {
        map:    [FormattedStory] -> HnDigestGroupSummary => map_llm_handler,
        reduce: (HnTopStories, [HnDigestGroupSummary]) -> HnDigestSummary => digest_llm_handler,
    },
    chunking: by_budget { /* shown in Step 3 */ },
    middleware: {
        map: ai_circuit_breaker(),
        reduce: ai_circuit_breaker(),
    }
);

The input and output types are the same HnTopStories and HnDigestSummary you wired back in Step 0, and the middleware block is the per-arm circuit breaker we get to in Step 4.

The one thing that stands out is the line that is still a comment, chunking: by_budget { /* shown in Step 3 */ }. That placeholder is exactly what this step fills in.

So the digest stage really comes down to three parts you provide:

  • Chunking (chunking: by_budget { ... }), the strategy that splits the batch into token-budgeted groups. This is the /* shown in Step 3 */ placeholder, and filling it in is most of this step.
  • The map arm (map: [FormattedStory] -> HnDigestGroupSummary), the handler that summarizes one chunk at a time.
  • The reduce arm (reduce: (HnTopStories, [HnDigestGroupSummary]) -> HnDigestSummary), the handler that folds those chunk summaries into the final digest.

We’ll start with chunking.

Token estimation and chunking

That chunking: by_budget { ... } placeholder is the framework’s answer to the one hard limit every model imposes.

Every LLM has a context window, a hard limit on how much text it can process in one call, measured in tokens. A token is roughly three quarters of a word in English, but the exact mapping depends on the model’s tokenizer. If you send more tokens than the window allows, the model either rejects the request or, worse, truncates silently. That silent truncation is the exact failure from Step 0, where sending thousands of stories in one shot gives you a confident summary of a sliver of the input and never tells you the rest was dropped.

Two pieces work together so that never happens:

  • Token estimation measures how many tokens a piece of text will cost before you send it.
  • Chunking splits the input into groups that each fit a token budget.

Together they let the flow handle arbitrarily large input without guessing, truncating, or hoping for the best. The chunking: by_budget { ... } block you just saw in the macro is where the digest stage does exactly that, so let’s fill it in.

chunking: by_budget {
    estimator: estimator.clone(),
    items:  |seed: &HnTopStories| seed.stories.clone(),
    render: |story: &FormattedStory, ctx| render_story_line(ctx.item_ordinal + 1, story),
    budget: budget_per_group,
    max_items: max_stories_per_group,
    oversize: decompose { max_depth: 5, exhaustion: fail },
    snapshot_excluded_items_limit: 25,
}

Read it field by field.

  • estimator is the token estimator. When you enable the top-level ai feature, ObzenFlow uses an exact tokenizer when it knows the model and falls back to a byte-length heuristic otherwise. The example derives it from the active model configuration with ai.estimator(), which is why the evidence header reports a token_estimator of Heuristic for the Ollama default.
  • items extracts the things to chunk out of the accumulated seed, here the Vec<FormattedStory> inside HnTopStories.
  • render turns each item into the stable text line the estimator measures and the model reads. The ctx.item_ordinal gives each story a stable position, and + 1 turns it into the one-based number the model cites in its output.
  • budget is the per-group token budget. The example defaults to 2500 tokens for Ollama and 6000 for hosted providers, both overridable with HN_AI_GROUP_BUDGET_TOKENS.
  • max_items caps how many stories land in a single group, regardless of budget. With short headlines, this cap is usually what decides the group count.
  • oversize decides what happens when a single item is larger than the entire budget. This is covered just below.
  • snapshot_excluded_items_limit bounds how many excluded item positions the evidence snapshot records, so a pathological input cannot bloat the journal.

The planner walks the items in order, estimates each one, and packs them into groups until the next item would exceed the budget or the max_items cap. With the default 60 stories and a cap of 10 per group, you get six groups, which is the groups: 6 you saw in the header.

Tip · Oversize handling is primarily a safety net

For Hacker News, a single item is a title plus a URL, so one story never comes close to a 2500-token budget and the oversize path effectively never fires. It is there as a correctness guarantee for arbitrary input. If you fed this same stage long documents instead of headlines, oversize: decompose { max_depth: 5, exhaustion: fail } would re-render an over-budget item at increasing decomposition depths up to five levels, and if it still could not make progress, it would surface a planning error rather than silently dropping or truncating the item. The alternative exhaustion: exclude would drop the item and record it in the planning stats instead. The flow refuses to quietly lose data either way.

The map arm

The map arm runs once per chunk. It builds a prompt from the rendered story lines and parses the model’s reply into a partial summary.

let map_llm_handler = ai
    .chat()
    .system(system_prompt.clone())
    .temperature(0.2)
    .max_tokens(800)
    .context(DigestMapCtx { interests: interests.clone() })
    .build_map_items_with_chunk_info(digest_map_prompt, digest_map_parse)
    .await?;

digest_map_prompt receives the chunk’s stories together with a ChunkInfo value that carries the chunk index, the chunk count, and the rendered story lines. It uses the Prompt builder to assemble a structured, rule-bearing prompt, and it reads the numbered lines straight from the chunk.

fn digest_map_prompt(ctx: &DigestMapCtx, stories: &[FormattedStory], chunk_info: &ChunkInfo)
    -> Result<UserPrompt, HandlerError>
{
    let mut p = Prompt::new();
    p.text_if(ctx.interests.as_deref(), |i| format!("My interests: {i}"))
        .text("Summarise these Hacker News stories (titles + URLs are provided as input).")
        .rules(/* do not invent facts, cite by number, ... */)
        .labeled("Output format (follow exactly)", /* Themes + Notable stories */)
        .fenced_lines("Input stories (numbered; do not repeat)", chunk_info.iter_rendered());
    Ok(p.finish())
}

Each chunk’s reply parses into a small HnDigestGroupSummary that holds only the markdown the model returned.

struct HnDigestGroupSummary {
    output_markdown: String,
}

The reduce arm

The reduce arm runs once, after all chunks have been summarized. It receives the original batch as a seed and the full list of chunk summaries, builds the final prompt, and packages everything into the evidence record.

let digest_llm_handler = ai
    .chat()
    .system(system_prompt.clone())
    .temperature(0.2)
    .max_tokens(800)
    .context(DigestReduceCtx { /* mode, provider, model, prompts, budget, ... */ })
    .build_reduce_seeded_with_prompt(digest_reduce_prompt, digest_reduce_parse)
    .await?;

The reduce parse assembles the final HnDigestSummary, which is the complete record for the run.

struct HnDigestSummary {
    mode: String,
    base_url: String,
    ai_provider: String,
    ai_model: String,
    token_estimator: EstimateSource,
    stories_fetched: usize,
    budget_per_group: TokenCount,
    groups: usize,
    interests: Option<String>,
    chat_prompt_system: SystemPrompt,
    chat_prompt_user: UserPrompt,
    input: HnTopStories,
    group_summaries: Vec<HnDigestGroupSummary>,
    output_markdown: String,
}

This is the part that is easiest to miss when you read the example quickly. ObzenFlow is not just calling an LLM twice. The map arm produces bounded, explainable partial results, the framework collects them in order, and the reduce arm composes the final answer from those partials rather than from one giant, unstable prompt. The map-reduce shape is what lets the same stage handle 10 stories or 200 without changing.

Step 4. Stage-level middleware

This example puts no middleware at the flow level. The flow-level middleware: [] is empty, and the protections are applied per stage instead, which gives precise control over which stage gets which policy.

A rate limiter on the source

The source stage carries a rate limiter, declared in the brackets after the source expression.

hn_stories = async_source!(HnStory => HttpPullSource::new(decoder, config), [
    RateLimiterBuilder::new(source_rate_limit).build()
]);

The limiter uses a token bucket internally, so the stage processes up to source_rate_limit events per second with tokens refilling at a steady rate. The default is 10 per second, lowerable from the command line with HN_SOURCE_RATE_LIMIT. It works alongside HttpPullSource, which already handles 429 responses, honours Retry-After, and retries transient failures. The explicit limiter means you stay polite before the upstream ever has to push back.

A circuit breaker on each model call

The map and reduce arms of the digest stage each carry a circuit breaker, declared in the middleware block of ai_map_reduce!.

middleware: {
    map: ai_circuit_breaker(),
    reduce: ai_circuit_breaker(),
}

Here is the body of ai_circuit_breaker().

pub fn ai_circuit_breaker() -> Box<dyn MiddlewareFactory> {
    CircuitBreakerBuilder::new(5)              // trip after 5 consecutive failures
        .with_retry_exponential(3)             // retry each failed event up to 3 times
        .with_retry_limits(RetryLimits::default()) // bounded backoff and a total wall-time cap
        .build()
}

The breaker trips open after five consecutive failures on its arm, which prevents the flow from hammering a model that is down or overloaded. Before tripping, it retries each failed event up to three times with exponential backoff under the default retry limits, which bound both a single retry and the total wall-clock time spent retrying. This matters in a pipeline with multiple model calls, because a single outage could otherwise stall the whole flow. With the breaker in place, the failure is contained to the affected arm and the flow can recover.

Applying middleware per stage rather than globally is deliberate. A rate limit makes sense on the source but not on the formatter. A circuit breaker makes sense on the model calls but not on the batch accumulator. ObzenFlow lets you put the protection exactly where it belongs.

Tip · Studying the explicit middleware path

If you want to see flow-level versus stage-level middleware inheritance in isolation, the flow_middleware_config example is built for exactly that.

Step 5. Evidence and replay

The whole reason to wrap a model call in a flow is so the run leaves evidence behind. This example writes journals under target/hn-ai-digest-logs, the same archive shape the earlier tutorials produced.

What gets journaled

The final HnDigestSummary event is the complete record for the run. It carries the provider, the model, the token estimator, the story count, the per-group budget, the group count, the exact system prompt, the exact final user prompt, the input stories, every chunk summary, and the final markdown output. The console sink prints much of that directly, which is what you saw in Step 1, but the durable copy is the journal, not the console.

Because every stage writes its own journal, you can inspect the intermediate chunk summaries, see whether oversize decomposition happened, and compare outputs across models or across runs. The run archive lives under target/hn-ai-digest-logs/flows/<run-id>/, and its run_manifest.json records the flow structure, the stage metadata, and the journal file mappings, the same manifest pattern from the earlier tutorials.

This is also where the single digest stage shows its true shape. The ai_map_reduce! stage you wrote as one line in the topology lowers into four journaled stages, one each for chunking, mapping, collecting, and finalizing, so every step of the map-reduce is independently inspectable on disk.

target/hn-ai-digest-logs/flows/<run-id>/
├── run_manifest.json
├── FiniteSource_hn_stories_stage_*.log
├── Transform_formatter_stage_*.log
├── Stateful_batch_stage_*.log
├── Transform_digest__chunk_stage_*.log
├── Transform_digest__map_stage_*.log
├── Stateful_digest__collect_stage_*.log
├── Transform_digest__finalize_stage_*.log
└── Sink_digest_summary_stage_*.log

A matching per-stage error journal and a system.log are written alongside these. The point is that the convenience of one declarative stage costs you nothing in inspectability. The map step that summarized each chunk and the finalize step that produced the digest each left their own durable record.

You do not have to guess the run id. The footer prints it for you when the flow completes, along with the exact replay command.

hn_ai_digest_demo completed. Journal: target/hn-ai-digest-logs/flows/flow_01KSR664VR8TZXPKDJMW7WNQ7X
To replay, add: --replay-from target/hn-ai-digest-logs/flows/flow_01KSR664VR8TZXPKDJMW7WNQ7X
(Source config env vars are ignored during replay)

Replay

That archive is what replay reads. You point the framework at a previous run and re-execute the flow against the same recorded inputs.

cargo run -p obzenflow --example hn_ai_digest_demo --features "http-pull ai" -- \
    --replay-from target/hn-ai-digest-logs/flows/<run-id>

The runtime swaps the live source for the archived source journal and reruns the flow against the same stories. If you stored that archive for years, you would still have the evidence needed to answer what data went in, how it was chunked, which prompts were used, which model ran, and what answer came out. The model call is the one non-deterministic part of the flow, and it is the part the journal records most carefully.

Next steps

You have now worked through all three tutorials. You have seen ObzenFlow handle a finite batch transform, a live HTTP-backed service with joins and projections, and a model-backed digest with token budgeting, chunking, accumulation, and replay.

If the bank transactions tutorial feels fuzzy after seeing this one, go back to Model Bank Transactions as a Flow and read the projection section again. The two examples reinforce each other, because both build a stable accumulated value and then derive output from it.

From here, the philosophy gives the deeper architectural rationale behind why ObzenFlow is shaped this way.