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.
Finite and infinite sources, transforms, stateful stages, joins, and sinks are the supervised units the runtime actually materializes.
Specializations
Specialized forms still create exactly one stage.
Declared effects and scalar AI inference strengthen a handler contract while reusing a Transform runtime stage.
Composites
Compose a multi-stage job as one typed capability.
AI MapReduce chunks a large input, maps each group, and reduces the typed partials into one final fact.
Runtime types
Six runtime stage types execute every flow.
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.
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.
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;
}#[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.
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;
}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.
// 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.
// 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.
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.
// 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 */]
);// 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 — 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
ChatCompletioneffect 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.
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.
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.
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.