Context and state¶
Everything a handler can reach besides its payload arrives through two objects with different lifetimes:
| Level | Type | Lives for | Holds |
|---|---|---|---|
| Application | the state type S |
the whole service | shared resources: pools, clients, configuration |
| Delivery | Context<'_, C, S> |
one message | the channel name, a headers working copy, the broker's typed per-delivery context C (read by key), and the typed shared state S |
The state is produced once, at startup, and is a single typed value of your own choosing. A
Context is built fresh for every delivery and threaded as &mut through the middleware chain into
the handler, so middleware and the handler observe (and can enrich) the same per-message view.
Application level: typed state¶
The shared application state is one typed value S (a struct you define, or () when the service
needs none). It is produced by an on_startup hook - the value the hook returns becomes the state,
fixing the app's state type:
// `on_startup` fixes the app's state type to `AppConfig`; `.layer` then grows the global stack.
#[ruststream::app]
fn app() -> RustStream<Stack<RequestId, Identity>, AppConfig> {
RustStream::new(AppInfo::new("context", "0.1.0"))
.on_startup(async move |()| {
Ok::<_, std::convert::Infallible>(AppConfig {
reject_zero_ids: true,
})
})
.layer(RequestId)
.with_broker(MemoryBroker::new(), |b| b.include(handle))
}
The state type is checked at compile time: a #[subscriber] handler that reads state names it as
the third Context generic (Context<'_, C, S>), and the runtime only lets that handler mount on
an app whose state type matches. A handler that names no state type is generic over it, so it mounts
on any app. publish(..) handlers follow the same rule with one twist: one that ignores the state
omits the Context parameter entirely and still mounts on a stateful app, but one that declares a
Context without naming a state type pins the state to (), so name the app's state type
explicitly to mount such a handler on a stateful app.
Handlers borrow the state with ctx.state(), which returns &S, the typed state itself - no
lookup, no Option, no downcast. The state is shared behind an Arc once the service runs, so
handlers get cheap shared references, not copies; interior mutability (an AtomicU64, a
mutex-guarded map) is the tool when a shared value must change at runtime. For data scoped to one
message rather than the whole service, use the per-delivery context
instead. See Lifespan for the startup-hook contract.
/// Shared configuration: produced once at startup, read by every handler as the typed app state.
#[derive(Debug)]
struct AppConfig {
reject_zero_ids: bool,
}
Injecting dependencies: extractor parameters¶
Reaching for a dependency through ctx.state().field always works, but a handler can also take it
as a parameter. Any handler parameter after the message (and the optional &mut Context) whose type
implements FromContext is an extractor: the runtime resolves it from the delivery before the
body runs, and a failed extraction settles the message by the rejection's HandlerResult without
running the body.
To inject a piece of the state, derive FromRef on the state and take State<T> in the handler -
no extractor impl by hand. State<T> resolves for any field type (T: FromRef<S>), including types
from other crates (a broker publisher, a client pool) that a bare per-field impl could not cover
under the orphan rule:
#[derive(FromRef)]
struct AppState {
create_order: CreateOrder,
}
The handler takes State<FieldType>, with no ctx.state() reach-through:
#[subscriber("orders")]
async fn handle(order: &Order, State(create_order): State<CreateOrder>) -> HandlerResult {
create_order.execute(order);
HandlerResult::Ack
}
A field that should not be injectable, or whose type another field already claims, opts out with
#[from_ref(skip)]; two fields may not share a type, since injection by type would be ambiguous. For
a custom extractor that does more than read the state - an auth guard that rejects, a request-scoped
resolver - implement FromContext directly: it borrows the &mut Context, so it can read headers,
broker fields, or a scratch value a middleware left, and return a Rejection to settle the delivery.
Delivery level: Context¶
A #[subscriber] handler opts in by declaring a second parameter after the payload; omit it when
the handler needs nothing but the message. The macro resolves the type itself, so Context needs
no import when it appears only in handler signatures:
#[subscriber("orders")]
async fn handle(order: &Order, ctx: &mut Context<'_, (), AppConfig>) -> HandlerResult {
// 1. The channel the message arrived on.
println!("received on {}", ctx.name());
// 2. The headers working copy - including what middleware added on the way in.
if let Some(id) = ctx.headers().get("x-request-id") {
println!("request {}", String::from_utf8_lossy(id));
}
// 3. The typed app-level shared state, borrowed through state().
let config = ctx.state();
if config.reject_zero_ids && order.id == 0 {
return HandlerResult::drop();
}
// 4. A post-settle hook: fires after the broker has acked this message, off the delivery
// path, so slow follow-up work never gates the ack or the next delivery. At-most-once:
// a lost hook does not redeliver.
let id = order.id;
ctx.after_ack(async move {
println!("order {id} acked; sending the confirmation");
});
HandlerResult::Ack
}
What the context exposes:
| Method | Returns | Purpose |
|---|---|---|
name() |
&str |
the channel / subject the message arrived on |
headers() |
&Headers |
the working copy of the message headers |
headers_mut() |
&mut Headers |
the same copy, for middleware to enrich |
state() |
&S |
the typed shared application state, borrowed directly |
context(KEY) |
KEY::Value |
a broker field read by compile-time key |
set(KEY, v) |
() |
write a per-delivery scratch value (middleware) |
after(outcome).then(fut) |
() |
a post-settle hook gated on the settlement outcome |
after_ack(fut) / after_settle(fut) |
() |
post-settle hook sugar (after an ack / after any settlement) |
Closure handlers (the manual typed(codec, |msg, ctx| ...) form) always take the context as their
second argument.
Per-delivery context¶
Beside the shared application state, the context carries the broker's typed per-delivery context,
read by compile-time key with no hashing, boxing, or downcasting. A key is a zero-sized selector the
broker exports; ctx.context(KEY) resolves it to a direct field read off the context, so a handler
reads native delivery metadata - a stream id, an offset, a delivery handle - without the broker
serializing it into the byte-only headers. A key implements Field only for the context types that
carry its field, so an inapplicable key is a compile error rather than a runtime miss.
use ruststream::Field;
// A broker crate ships its per-delivery context and the keys that read its fields; an application
// reads a field by key from a handler taking `&mut Context<'_, Delivery>`.
struct Delivery {
offset: u64,
}
#[derive(Clone, Copy)]
struct Offset;
impl Field<Delivery> for Offset {
type Value<'a> = u64;
fn get(self, d: &Delivery) -> u64 {
d.offset
}
}
The context type is built from the message by BuildContext, which the runtime calls once per
delivery; a broker with no per-delivery fields uses (), the default (so a #[subscriber] handler
that names no context type - and takes no Ctx extractor - sees
Context<'_>). Middleware can also carry a typed scratch value to a
downstream handler: a writable key (FieldMut) lets a layer ctx.set(KEY, value) and the handler
ctx.context(KEY) it back - a correlation id, an authenticated user a layer resolved - without
serializing it into the headers. The context is built fresh per delivery, so one delivery's values
never leak into the next.
Context fields as parameters¶
A field can also arrive as a handler argument, the way State<T> injects a state component: the
Ctx<K> extractor binds the value the key K reads. The key implements ContextField - a
Field-style trait that additionally names the context type it reads from and yields an owned
value - so the handler needs no &mut Context parameter at all: the #[subscriber] macro projects
the subscription's context type from the first Ctx key in the signature.
// The broker's per-delivery context, built once per delivery; here it carries the payload
// size, standing in for an offset, a partition, or a delivery tag a real broker exposes.
struct DeliveryMeta {
payload_len: usize,
}
impl BuildContext<MemoryMessage> for DeliveryMeta {
fn build(msg: &MemoryMessage) -> Self {
Self {
payload_len: msg.payload().len(),
}
}
}
// The key: `ContextField` names the context it reads and yields an owned value, which is what
// lets it work as an extractor. Broker crates ship these next to their `Field` keys.
#[derive(Clone, Copy, Default)]
struct PayloadLen;
impl ContextField for PayloadLen {
type Context = DeliveryMeta;
type Value = usize;
fn read(self, src: &DeliveryMeta) -> usize {
src.payload_len
}
}
// The field arrives as an argument; no `&mut Context` parameter, no context type named - the
// macro projects it from the key.
#[subscriber("orders")]
async fn audit(order: &Order, Ctx(len): Ctx<PayloadLen>) -> HandlerResult {
println!("order {} arrived as {len} bytes", order.id);
HandlerResult::Ack
}
Three things to know:
- Values are owned (
ContextField::Valueis'static): extractor values bind before the handler body runs, so borrowing from the context is not an option. Keys yielding borrowed values (a name as&str) stay readable throughctx.context(KEY)with a declared ctx parameter. - With a
&mut Context<'_, C>parameter also present, everyCtxkey must read that sameC; the compiler enforces it through the extractor bounds. - The projection is syntactic: the macro recognizes the literal
Ctx<K>shape (any path ending inCtxwith one type argument). A type alias hides it, and the context type falls back to().
The headers working copy¶
ctx.headers() is not the broker message itself: each delivery clones the incoming headers into a
working copy that lives in the context. That makes it a scratchpad for the dispatch chain -
middleware earlier in the chain can stamp values onto it with headers_mut(), and the handler
reads the enriched result:
/// A layer that stamps a request id onto the context headers before the handler runs.
#[derive(Clone)]
struct RequestId;
struct WithRequestId<H>(H);
impl<H> Layer<H> for RequestId {
type Handler = WithRequestId<H>;
fn layer(&self, inner: H) -> WithRequestId<H> {
WithRequestId(inner)
}
}
static NEXT_REQUEST: AtomicU64 = AtomicU64::new(1);
// The layer is state-agnostic: it threads the context `C` and state `S` through unchanged, so it
// wraps a handler whatever typed state the app declares.
impl<M: Send + Sync, C: Send, S: Send + Sync, H: Handler<M, C, S>> Handler<M, C, S>
for WithRequestId<H>
{
async fn handle(&self, msg: &M, ctx: &mut Context<'_, C, S>) -> Settle {
if ctx.headers().get("x-request-id").is_none() {
let id = format!("req-{}", NEXT_REQUEST.fetch_add(1, Ordering::Relaxed));
ctx.headers_mut().insert("x-request-id", id.into_bytes());
}
self.0.handle(msg, ctx).await
}
}
Mounted globally, the layer runs before every handler, so handle above always finds
x-request-id:
// `on_startup` fixes the app's state type to `AppConfig`; `.layer` then grows the global stack.
#[ruststream::app]
fn app() -> RustStream<Stack<RequestId, Identity>, AppConfig> {
RustStream::new(AppInfo::new("context", "0.1.0"))
.on_startup(async move |()| {
Ok::<_, std::convert::Infallible>(AppConfig {
reject_zero_ids: true,
})
})
.layer(RequestId)
.with_broker(MemoryBroker::new(), |b| b.include(handle))
}
Two boundaries to keep in mind:
- Mutations stay within the delivery: the broker message and other subscribers' deliveries are untouched.
- Outgoing messages do not inherit the copy. Replies and manual publishes start from fresh
headers; attach outgoing metadata in the publish pipeline
(a
PublishTransformorPublishLayer) instead.
Publishing from a handler¶
To publish from inside a handler (beyond the publish(..) reply form), do not put the publisher
in the state: take it as a handler parameter with Out - the pattern Out(out): Out<P> binds
out to a live publisher inside the body. The source is attached where the handler is included,
and the runtime pairs it after the broker connects, so the handler never sees a "not connected"
publisher and the state stays free of connection-bound values. The full pattern and its snippet
live in Publishing from inside a handler.
Post-settle hooks¶
Sometimes a handler needs a side effect to fire after the message has been settled - a non-critical notification, slow follow-up work, a cache warm-up - without it gating the ack decision or affecting redelivery. Register one on the context:
#[subscriber("orders")]
async fn handle(order: &Order, ctx: &mut Context<'_, (), AppConfig>) -> HandlerResult {
// 1. The channel the message arrived on.
println!("received on {}", ctx.name());
// 2. The headers working copy - including what middleware added on the way in.
if let Some(id) = ctx.headers().get("x-request-id") {
println!("request {}", String::from_utf8_lossy(id));
}
// 3. The typed app-level shared state, borrowed through state().
let config = ctx.state();
if config.reject_zero_ids && order.id == 0 {
return HandlerResult::drop();
}
// 4. A post-settle hook: fires after the broker has acked this message, off the delivery
// path, so slow follow-up work never gates the ack or the next delivery. At-most-once:
// a lost hook does not redeliver.
let id = order.id;
ctx.after_ack(async move {
println!("order {id} acked; sending the confirmation");
});
HandlerResult::Ack
}
The handler above ends with ctx.after_ack(..): the continuation runs only once the broker has
acked the message, off the delivery path, so it never delays the ack or the next delivery.
Three forms, all additive:
ctx.after(outcome).then(fut)- runs only if the message settles byoutcome, matched by kind. The four kinds are distinct:Ack,drop()(nack, no requeue),retry()(nack, requeue), andretry_after()(matched regardless of the delay). Drop and retry are separate mechanics, so a hook gated ondrop()does not fire on aretry()settlement, and vice versa.ctx.after_ack(fut)- sugar forctx.after(HandlerResult::Ack).then(fut).ctx.after_settle(fut)- runs after the message settles, whatever the outcome.
A handler can also attach a continuation through its return value: any outcome converts into a
Settle with .and_after(fut), which is how a batch handler gets per-element continuations. See
Post-settle continuations for that form; the semantics
below apply to both.
Multiple registrations accumulate and every matching one runs, on a tracked task set off the
delivery path. The semantics are at-most-once: the message is already settled before any hook
runs, so a hook that panics, or that is lost when the process crashes, never causes a redelivery.
Do not put work whose loss must redeliver the message in a hook; settle by outcome and let the
broker retry instead. A graceful shutdown drains in-flight hooks (bounded by shutdown_timeout);
an aborted shutdown may drop them.
On the batch path a Context is one per batch, so a hook runs after the whole batch has settled.
Because a batch has per-element outcomes, the outcome gate is ill-defined there: only
after_settle hooks fire (the gated after(..) / after_ack forms are ignored on a batch).
Context in middleware¶
Every middleware form receives the same &mut Context the handler will see, which is what makes
the enrichment pattern work:
- A static layer's
Handler::handle(&self, msg, ctx)- as in the example above. - A dynamic
DynMiddleware::handle(&self, input, ctx, next)- inspect or enrich, thennext.run(input, ctx).
The middleware forms themselves are covered in Middleware. The full program for
this page is
examples/context.rs.