The Grammar of Durable Flows

Dive deeper into ObzenFlow’s vocabulary.

At the heart of ObzenFlow are stages, specialized handlers, and composites: a small vocabulary for transforming events, joining streams with reference data, accumulating state, and making external effects explicit. You’ll learn how to compose them into sophisticated durable programs that bring the directness of classical application development to real-time stream processing.

Runtime types

Six runtime stage types execute every flow.

These are the framework’s concrete execution categories. Each materialized stage has one type, one supervisor, and its own durable runtime boundary.

source! ingest batch and streaming data in one flow

ObzenFlow gives you four source shapes. Pick one by how its data arrives (bounded or always-on, synchronous or asynchronous) and bring the handler that speaks to your system.

Boundedness

finite
Bounded batch. Exits when drained.
infinite
Runs for the life of the binary.

Execution model

sync
Handler returns synchronously.
async
Handler awaits I/O per item.

Together, these choices provide flexibility and efficiency. Mix finite batch sources with always-on live sources in the same flow, then choose synchronous or asynchronous execution to match how each source produces data.

flow.rs — four source shapes, one topology
stages: {
    // finite + sync
    csv_history = source!(
        OrderEvent => csv_orders("orders.csv")
    );

    // finite + async
    s3_backfill = async_source!(
        OrderEvent => s3_orders(bucket)
    );

    // infinite + sync
    local_live = infinite_source!(
        OrderEvent => local_orders()
    );

    // infinite + async
    kafka_live = async_infinite_source!(
        OrderEvent => kafka_orders(topic)
    );

    audit = sink!(
        OrderEvent => sinks::console(render)
    );
},

topology: {
    csv_history |> audit;
    s3_backfill |> audit;
    local_live  |> audit;
    kafka_live  |> audit;
}

transform! the swiss army knife for stateless data transformation

transform! is ObzenFlow’s bread-and-butter stage. Use it whenever the next facts can be decided from one input fact, without remembering earlier events or touching the outside world. Normalization, validation, classification, projection, filtering, and the facts that drive typed routing all begin here.

The built-in helpers cover the transformations most pipelines repeat.

map
One fact becomes one fact.
filter
The same fact passes or is dropped.
filter_map
Zero or one fact of a new type.
try_map
One fact or a handler error.
custom
Exhaustive, multi-type business outcomes.
flow.rs — one source, three business facts
stages: {
    web_orders_source = source!(
        WebOrderSubmitted => web_orders_feed
    );

    triage_web_orders = transform!(
        WebOrderSubmitted -> {
            RegularOrderAccepted,
            PriorityOrderAccepted,
            OrderHeldForDiscontinuedItem,
        } => TriageWebOrders
    );

    regular_orders_sink = sink!(
        RegularOrderAccepted => deliver_regular_orders
    );
    priority_orders_sink = sink!(
        PriorityOrderAccepted => deliver_priority_orders
    );
    discontinued_orders_sink = sink!(
        OrderHeldForDiscontinuedItem => deliver_discontinued_orders
    );
},

topology: {
    web_orders_source |> triage_web_orders;
    triage_web_orders |> regular_orders_sink;
    triage_web_orders |> priority_orders_sink;
    triage_web_orders |> discontinued_orders_sink;
}
triage.rs — one order, one authored fact
#[derive(Debug, Clone, StageOutputFacts)]
pub enum WebOrderTriageOutcome {
    Regular(RegularOrderAccepted),
    Priority(PriorityOrderAccepted),
    DiscontinuedItem(OrderHeldForDiscontinuedItem),
}

#[derive(Debug, Clone)]
pub struct TriageWebOrders;

impl TypedTransformHandler for TriageWebOrders {
    type Input = WebOrderSubmitted;
    type Output = WebOrderTriageOutcome;

    fn process(&self, order: WebOrderSubmitted)
        -> Result<Self::Output, HandlerError>
    {
        if order.contains_discontinued_item() {
            Ok(WebOrderTriageOutcome::DiscontinuedItem(
                OrderHeldForDiscontinuedItem::from(order),
            ))
        } else if order.is_priority_customer() {
            Ok(WebOrderTriageOutcome::Priority(
                PriorityOrderAccepted::from(order),
            ))
        } else {
            Ok(WebOrderTriageOutcome::Regular(
                RegularOrderAccepted::from(order),
            ))
        }
    }
}

join! enrich events with reference data

join! gives a forward-moving event typed context from a second stream. One side maintains a keyed catalog of reference facts; the other carries events through the topology. The handler identifies their shared key and decides what fact a match authors. Readers familiar with Martin Kleppmann’s Designing Data-Intensive Applications will recognize this as the established stream-table join pattern documented there.

At its core, a join has a catalog, a stream, and a key shared by both. Its matching policy decides what to emit when that key matches—or when it does not.

inner
Authors one output for a match and nothing when the key is absent.
left
Always calls the merge function and supplies the reference as an Option.
strict
Treats a missing reference as a flow-integrity failure.

Each matching policy comes in two reference lifecycles.

finite
Hydrates the complete reference catalog to EOF before processing the stream.
live
Accepts reference updates while stream events continue to arrive.
flow.rs — one catalog, one forward stream
stages: {
    opened_accounts_source = async_infinite_source!(
        AccountOpened => opened_accounts_feed
    );
    ledger_entries_source = async_infinite_source!(
        LedgerEntry => ledger_entries_feed
    );

    post_ledger_entries = join!(
        catalog opened_accounts_source: AccountOpened,
        LedgerEntry -> PostedEntry => post_ledger_entry
    );

    posted_entries_sink = sink!(
        PostedEntry => deliver_posted_entries
    );
},

topology: {
    ledger_entries_source |> post_ledger_entries;
    post_ledger_entries |> posted_entries_sink;
}
post_entry.rs — two keys, one authored fact
let post_ledger_entry = joins::inner_live(
    |account: &AccountOpened| account.account_id.clone(),
    |entry: &LedgerEntry| 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,
    },
);

stateful! accumulates across events and controls when to emit

A stateful stage is a fold. Events arrive, state accumulates, and the stage eventually emits that state as an output event. You make two separate choices: how to accumulate and when to emit. Start with the accumulation strategy.

reduce
Accumulates into a running value.
group_by
Keeps per-key state across a keyed dimension.
top_n_by
Maintains the top N items by an extracted score.

Emission is the second choice. A bounded backfill wants one summary when the input ends, a live dashboard wants a snapshot every few seconds, and an audit stream wants every update.

emit_on_eof
Waits until the input is exhausted.
emit_every_n
Emits after each batch of events.
emit_within
Emits on a time window.
emit_always
Emits after every update.

Strategy and emission compose, so the same accumulator can produce end-of-input summaries, per-N snapshots, or tumbling windows with only a line change.

stateful_choices.rs
// reduce: accumulate into a running value
metrics = stateful!(
    UserEvent -> MetricsState => typed_stateful::reduce(
        MetricsState::default(),
        |state, event: &UserEvent| event.update_metrics(state),
    )
    .emit_on_eof()
);

// group_by: per-key state across a keyed dimension
session_tracker = stateful!(typed_stateful::group_by(
    |event: &UserEvent| event.user_id.clone(),
    |session: &mut SessionData, event: &UserEvent| {
        event.update_session(session);
    },
)
.emit_within(Duration::from_secs(3)));

// top_n_by: keep the top N by an extracted score
top_products = stateful!(typed_stateful::top_n_by(
    5,
    |order: &OrderEvent| order.product_id.clone(),
    |order: &OrderEvent| order.total_value,
)
.emit_every_n(5));

sink! delivers committed facts out of the flow

A sink is the final delivery boundary. It consumes facts already committed by an upstream stage and hands them to a destination such as a file, database projection, console, or downstream service. Every sink declares what should happen if recovery delivers the same fact again.

idempotent
The destination absorbs re-delivery. Deterministic files, keyed upserts, and replaceable projections fit here.
non_idempotent
Re-delivery would duplicate work at the destination, so Resume refuses unless the operator explicitly opts into duplicate sink delivery.

A non-idempotent external write the domain cares about does not belong in a sink. It belongs mid-flow as a declared effect, where its outcome becomes a typed fact. What reaches a sink should usually be a projection that can safely be delivered again.

flow.rs
// The destination absorbs re-delivery: resume proceeds.
summaries = sink!(
    OrderSummary => order_summary_upsert(pool)
        .idempotent()
);

// Re-delivery would duplicate at the destination:
// resume refuses unless the run opts in with
// allow_duplicate_sink_delivery.
notify = sink!(
    OrderNotification => notification_post(client)
        .non_idempotent()
);

Specializations

Specialized forms still create exactly one runtime stage.

These declarations change what a handler may do, not how many execution boundaries appear in the graph.

effectful_transform! makes per-event side effects explicit

Re-running deterministic logic is safe. Re-charging a credit card is not. An effectful transform makes the external call a typed part of the stage signature. The stage declares its effects, and the handler hands each call to the runtime through fx.perform rather than performing it inline.

uses
The stage’s declared effect types. An undeclared effect fails at cargo check, and the runtime refuses it again before any I/O happens.
with
The dependency’s resilience unit—breaker, retry, and per-attempt rate limit—attached to the effect it protects.
fx.perform
Hands the live call to the runtime, which journals its outcome as a named fact.
execute
The effect’s own method, where the live call actually runs. The handler never calls it directly.
idempotency_key
Non-idempotent effects must carry a stable key before the call leaves the process; the receiving system enforces deduplication.

Once an outcome is committed, replay reads it from the journal instead of firing the call again. The adapter behind an effect is supplied at the flow level, so tests and hosts can swap implementations without touching the stage.

payment_gateway_resilience/flow.rs
Step 1 · Declare the effect on the stage
// flow.rs — the stage names the effect boundary.
authorize_payment = effectful_transform!(
    ValidatedOrder -> {
        PaymentAuthorized,
        PaymentDeclined,
        CancelledOrder,
        PaymentAuthorizationUnavailable
    }
    uses AuthorizePayment
        with gateway_resilience
    => gateway_transform,
    observers: [/* observation only */]
);
payment_gateway_resilience/gateway.rs
Step 2 · Perform it where the handler touches the world
// gateway.rs — the handler returns the call as a value.
let outcome = fx
    .perform(AuthorizePayment { order: order.clone() })
    .await;

// PaymentAuthorized or PaymentDeclined is already
// journaled; outcome is matched only for consequences.
stripe_gateway.rs
Step 3 · Describe exactly how the world is touched
// stripe_gateway.rs — execute is where the flow touches the world.
#[async_trait]
impl Effect for AuthorizePayment {
    const SAFETY: EffectSafety =
        EffectSafety::NonIdempotentRequiresKey;

    type Outcome = AuthorizePaymentOutcome;

    async fn execute(&self, ctx: &mut EffectContext)
        -> Result<Self::Outcome, EffectError>
    {
        // One real POST to Stripe. The effect's key rides
        // the Idempotency-Key header, so a crash-and-retry
        // can never become a second authorization.
        let intent = self.stripe
            .create_payment_intent(CreatePaymentIntent {
                amount: self.order.amount_cents,
                currency: Currency::Usd,
                payment_method: self.order.payment_method.clone(),
                capture_method: CaptureMethod::Manual,
                confirm: true,
            })
            .idempotency_key(self.idempotency_key())
            .await?;

        match intent.status {
            PaymentIntentStatus::RequiresCapture =>
                Ok(Authorized(/* PaymentAuthorized fact */)),
            _ => Ok(Declined(/* PaymentDeclined fact */)),
        }
    }

    fn idempotency_key(&self) -> Option<IdempotencyKey> {
        Some(IdempotencyKey(format!(
            "payment-authorize:{}", self.order.order_id
        )))
    }
}

inference! one model call is still one transform stage

inference! is the generated single-stage form for one replay-safe model call. It wraps an InferenceHandler in the ChatCompletion effect protocol, requires resilience on that dependency, and materializes the result as one Transform runtime stage.

uses
Declares the generated ChatCompletion effect and its execution safety.
via
Binds the model provider used for the live call.
with
Attaches the resilience unit that protects the provider boundary.

Use this form when the input is already reduced and one model decision is the whole operation. When the job must chunk input, fan work out, collect partials, and finalize a result, that is no longer one specialized transform. It is a composite.

one_shot_inference_demo/main.rs
fn build_flow_definition(input: ReducedEvidence, journal_path: PathBuf) -> FlowDefinition {
    FlowDefinition::materialize(move |runtime_config| {
        let chat = ChatEffectBinding::from_config(&runtime_config.ai_models())?;
        let evidence = sources::once(input);
        let generate_brief = GenerateBrief;
        let display = sinks::console(|brief: &DecisionBrief| {
            format!("{}\n\n{}", brief.question, brief.recommendation.trim())
        });

        Ok(flow! {
            name: "one_shot_inference_demo",
            journals: disk_journals(journal_path),
            stages: {
                evidence = source!(ReducedEvidence => evidence);
                brief = inference!(
                    ReducedEvidence -> DecisionBrief
                    uses at_least_once(ChatCompletion)
                        via chat
                        with ai_resilience()
                    => generate_brief
                );
                display = sink!(DecisionBrief => display);
            },

            topology: {
                evidence |> brief;
                brief |> display;
            }
        })
    })
}

#[tokio::main]
async fn main() -> Result<()> {
    let input = ReducedEvidence {
        question: "Should the release use one-shot inference or map-reduce?".to_string(),
        evidence: vec![
            "The input is already reduced and bounded.".to_string(),
            "Exactly one model decision is required.".to_string(),
            "No fan-out or fan-in is needed.".to_string(),
        ],
    };

    FlowApplication::builder()
        .run_async(build_flow_definition(
            input,
            PathBuf::from("target/one_shot_inference_demo_journal"),
        ))
        .await?;

    Ok(())
}

Composites

Compose a multi-stage job as one typed capability.

Reach for a composite when the work has a repeatable internal workflow but the surrounding flow should see one clear input and output.

ai_map_reduce! turns large inputs into one final result

Some inputs are too large for one model call. Silently trimming them produces a confident answer from incomplete evidence; hand-wiring the chunking, fan-out, collection, and final synthesis repeats the same difficult machinery in every flow. ai_map_reduce! makes that entire job one typed declaration.

map
Runs once per chunk and authors one typed partial result.
reduce
Combines the original input with the ordered map results to author the final fact.
chunking
Groups items to fit the model’s token and item limits.
hn_ai_digest_demo/flow.rs — one typed digest
digest = ai_map_reduce!(
    HnTopStories -> HnDigestSummary => {
        map: [FormattedStory] -> HnDigestGroupSummary
            uses at_least_once(ChatCompletion)
                via chat
                with ai_resilience()
            => map_role,

        reduce: (HnTopStories, [HnDigestGroupSummary]) -> HnDigestSummary
            uses at_least_once(ChatCompletion)
                via chat
                with ai_resilience()
            => finalise_role,
    },
    chunking: by_budget {
        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,
    }
);

Next Up

Now see the running flow through an operational lens.

ObzenFlow Studio turns topology, middleware state, metrics, and contract evidence into one coherent operational view.