Internals: Subscribe & Fanout
Picks up where Internals: The Publish Path leaves
off — a DeliveryEnvelope has just been cloned into a subscriber’s mpsc
channel. This page traces what happens from there through to bytes on the
wire, plus the subscribe handshake that set that channel up in the first
place.
The cast of types
Section titled “The cast of types”| Type | Where | What it is |
|---|---|---|
SubscriptionReceiver |
crates/felix-broker/src/lib.rs |
The broker-core side of a subscriber’s channel; yields DeliveryEnvelopes |
WriterLaneManager |
services/broker/src/transport/quic/handlers/subscribe/lane.rs |
Owns a fixed set of writer lanes and the per-connection writer tasks they feed |
LaneCommand |
same | Register / Delivery / Unregister, sent from a subscription’s feeder to its assigned lane |
ConnectionCommand |
same | Same three variants, one hop further — sent from a lane to the connection that owns the subscriber’s QUIC stream |
run_lane_feeder |
same | One task per subscription; reads DeliveryEnvelopes, encodes (once), dispatches LaneCommands |
run_writer_lane |
same | One task per lane; receives LaneCommands, forwards to the right connection |
run_connection_writer |
same | One task per QUIC connection; owns the actual SendStreams, does the writing |
Three hops, deliberately: SubscriptionReceiver (broker-core, per
subscriber) → lane (shared across many subscribers, bounded parallelism) →
connection writer (one per physical connection, since a connection’s QUIC
streams can’t be written from multiple tasks concurrently without
coordination).
Subscribe handshake
Section titled “Subscribe handshake”File: services/broker/src/transport/quic/handlers/subscribe.rs,
handle_subscribe_message
- Client sends
Message::Subscribeon the control (bi) stream. - Broker calls
Broker::subscribe(tenant, namespace, stream)→StreamState::register_subscriber(), which allocates a slot in aSlab, creates thempsc::channel::<DeliveryEnvelope>(subscriber_queue_capacity), and rebuilds the lock-freesubscribers_snapshot(see Internals: The Publish Path for why that snapshot exists). - Broker opens a new unidirectional stream (
connection.open_uni()) — this one is uni, unlike publish — and writesMessage::EventStreamHello { subscription_id }as the first frame. This is the only frame that ever carries the subscription id; every event frame after it doesn’t need to, because the client already knows which stream is bound to which subscription. This is what makes the shared-frame optimization below possible — see Wire Protocol: Shared Binary EventBatch. - Broker computes
lane_idx = manager.select_lane(subscription_id, connection_id)and sendsLaneCommand::Registerto that lane — registering the subscriber’sSendStreamwith the writer-lane pipeline. If this fails (lane queue full), the broker repliesMessage::Errorand stops here — the client never seesSubscribedfor a registration that didn’t actually take. - Only now does the broker reply
Message::Subscribedon the control stream — confirming registration succeeded, not just that the request was received. - Broker spawns
run_lane_feeder, the task that will pullDeliveryEnvelopes out of this subscriber’sSubscriptionReceiverfor the rest of the subscription’s life. Ifcore_shardsis enabled, this task is spawned on the shard owning the stream (resolved viaresolve_stream_handle+shards.handle_for(handle.id())), not on the default runtime — see Internals: Backpressure & Core Sharding.
Lane assignment: select_lane
Section titled “Lane assignment: select_lane”File: subscribe/lane.rs, WriterLaneManager::select_lane /
lane_for_subscriber / lane_for_connection
Controlled by subscriber_lane_shard:
subscriber_id_hash:hash64(subscriber_id) % lane_count— independent of connection topology.connection_id_hash:hash64(connection_id) % lane_count— useful when many subscribers share one connection and you want them on the same lane.round_robin_pin: assigned once at subscribe time, pinned for the life of the subscription — preserves ordering, can skew under uneven churn.auto(default): connection-aware when a connection id is known (equivalent toconnection_id_hash), else falls back to subscriber id.
subscriber_single_writer_per_conn: true forces every subscriber on a
connection onto the same lane regardless of the shard policy — the
latency-profile default, trading lane parallelism for strict per-connection
ordering.
run_lane_feeder: where encode-once actually happens
Section titled “run_lane_feeder: where encode-once actually happens”File: subscribe/feeder.rs, run_lane_feeder
async fn run_lane_feeder( mut event_rx: SubscriptionReceiver, manager: Arc<WriterLaneManager>, lane_idx: usize, connection_id: Option<u64>, config: EventWriterConfig,) { loop { let envelope = event_rx.recv().await; // blocks until broker core sends one // ... coalesce with more envelopes up to max_events / max_bytes ... let frame = envelope.shared_event_frame()?; // <-- the encode-once call enqueue_lane_frame(&manager, lane_idx, config.subscription_id, frame, ..).await; }}shared_event_frame() (on DeliveryEnvelope, crates/felix-broker/src/lib.rs)
is a lazily-populated cache: the first subscriber’s feeder to call it pays
the real encode cost (felix_wire::binary::encode_shared_event_batch_bytes)
and stores the result in Mutex<Option<Bytes>> inside the envelope; every
other subscriber calling it on the same envelope gets a cheap Bytes::clone
(a refcount bump, not a copy). Since publish fanout hands the same
DeliveryEnvelope (via Arc) to every subscriber
(see Internals: The Publish Path),
one publish batch is encoded once, total, regardless of fanout — not once
per subscriber. This is the change that took fanout cost from O(fanout) to
O(1) per publish.
Coalescing here is governed by EventWriterConfig: max_events,
max_bytes, flush_delay, and single_event_mode (forced when
fanout_batch_size <= 1, i.e. the latency profile — one event per frame,
immediate flush, no batching delay).
run_writer_lane → run_connection_writer
Section titled “run_writer_lane → run_connection_writer”File: subscribe/writer.rs (lane routing in subscribe/lane.rs)
A lane’s job is small: receive LaneCommands and forward them as
ConnectionCommands to whichever connection the subscriber belongs to
(ensure_connection_writer/enqueue_connection, which lazily spawns a
run_connection_writer task per connection the first time it’s needed).
This hop exists because a QUIC connection’s streams can’t be written
concurrently from independent tasks without a single owner coordinating it.
run_connection_writer is where the actual send.write_all() happens, and
it’s the part of this pipeline that changed most this session — worth
understanding in detail if you’re touching write scheduling.
The old design: a round barrier
Section titled “The old design: a round barrier”Originally, each pass through the writer loop built one write per
subscriber with pending data, launched them all concurrently via
FuturesUnordered, then waited for every one of them to complete before
starting the next round. That’s fine when every subscriber’s stream is fast,
but one backpressured or slow QUIC stream would stall the next round for
every other subscriber sharing that connection — a straggler problem.
flowchart LR
subgraph barrier["Old: round barrier"]
direction TB
OR(["round starts"]) o1@--> OA["write to A"]
OR o2@--> OB["write to B"]
OR o3@--> OC["write to C<br/><small>slow / backpressured</small>"]
OA o4@--> OW{{"wait for<br/>all three"}}
OB o5@--> OW
OC o6@--> OW
OW o7@--> ON(["next round<br/><small>A and B sat idle</small>"])
end
subgraph pipelined["Current: continuous pipelining"]
direction TB
NA["A completes"] n1@--> NA2(["A's next write<br/><small>starts immediately</small>"])
NB["B completes"] n2@--> NB2(["B's next write<br/><small>starts immediately</small>"])
NC["C still in flight"] n3@--> NC2(["finishes later,<br/><small>blocks nobody</small>"])
end
o1@{ animation: slow }
o2@{ animation: slow }
o3@{ animation: slow }
o4@{ animation: slow }
o5@{ animation: slow }
o6@{ animation: slow }
o7@{ animation: slow }
n1@{ animation: fast }
n2@{ animation: fast }
n3@{ animation: slow }
classDef step fill:#e8f0fe,stroke:#4a6fa5,color:#1a2b40
classDef gate fill:#fdeaea,stroke:#b04a4a,color:#3d1414
classDef slowc fill:#fdf0e3,stroke:#b07d3a,color:#3d2a12
classDef ok fill:#e9f5ec,stroke:#4a8a5e,color:#16301f
class OA,OB,NA,NB step
class OW gate
class OC,NC,ON,NC2 slowc
class OR,NA2,NB2 ok
C moves at the same speed in both halves — the difference is only whether A and B are made to wait for it.
The current design: continuous pipelining
Section titled “The current design: continuous pipelining”let mut in_flight: HashSet<u64> = HashSet::new();let mut writes = FuturesUnordered::new();loop { // Start a write for every subscriber that has queued data AND isn't // already mid-write. let ready: Vec<u64> = deliveries.iter() .filter_map(|(id, q)| (!q.is_empty() && !in_flight.contains(id)).then_some(*id)) .collect(); for subscriber_id in ready { in_flight.insert(subscriber_id); writes.push(async move { /* coalesce + write */ }); } let Some((subscriber_id, .., write_result)) = writes.next().await else { break; // nothing ready, nothing in flight — this connection is drained }; in_flight.remove(&subscriber_id); // handle result; if Ok, this subscriber becomes eligible again next loop}The difference: as soon as any subscriber’s write completes, the loop
immediately checks whether that subscriber has more queued data and — if
so — starts its next write right away, without waiting for other in-flight
writes to finish. A slow subscriber’s write can still be in flight while
three other subscribers race ahead independently. This matters most when
multiple subscribers share one physical connection (sub_conns small
relative to fanout in the benchmark harness, or subscriber_single_writer_per_conn: true in production) — see Benchmarks for the
measured effect.
Worked example
Section titled “Worked example”Three subscribers (A, B, C) on the same QUIC connection, one publish batch
lands as one DeliveryEnvelope:
- Broker core sends the same envelope (3
Arcclones) to A’s, B’s, and C’sSubscriptionReceivers. - Three independent
run_lane_feedertasks wake up (possibly on different lanes, or the same lane ifsubscriber_single_writer_per_connis set). Say A’s feeder runs first: it callsenvelope.shared_event_frame(), pays the encode cost, getsBytes. B’s and C’s feeders call the same method microseconds later and get the cachedBytesfor free. - Each feeder dispatches a
LaneCommand::Deliverycarrying that (shared)Bytes— cloningBytesis cheap (refcount), so no re-serialization happens even though three separate lane commands now exist. - The lane(s) forward
ConnectionCommand::Deliveryto the one connection writer for this connection. - The connection writer’s loop sees three subscribers ready, starts three concurrent writes. If A’s QUIC stream is flow-control-blocked, B’s and C’s writes still complete and — if they have more queued data — start their next write immediately, not waiting on A.
If you want to change…
Section titled “If you want to change…”| You want to… | Look at |
|---|---|
| Change event batching/coalescing thresholds | EventWriterConfig construction in handle_subscribe_message; the coalescing loop in run_lane_feeder |
| Change lane assignment policy | SubscriberLaneShard in services/broker/src/config.rs; WriterLaneManager::select_lane in subscribe/lane.rs |
| Change subscriber backpressure policy | SubQueuePolicy — two separate checkpoints: subscriber_queue_policy (broker-core, Broker::publish_batch_to_handle) and subscriber_lane_queue_policy (lane ingress, WriterLaneManager::enqueue/enqueue_connection). See Internals: Backpressure |
| Change write scheduling/fairness across subscribers on one connection | run_connection_writer’s in_flight/FuturesUnordered loop, subscribe/writer.rs |
| Change the wire format for event delivery | encode_shared_event_batch_bytes/decode_shared_event_batch, crates/felix-wire/src/lib.rs; update Wire Protocol too |
| Add a new lane→connection routing mode | WriterLaneManager::ensure_connection_writer/enqueue_connection, subscribe/lane.rs |
Next: Internals: Backpressure & Core Sharding
ties the publish-side and subscribe-side admission/queue layers together
into the full picture, and covers the core_shards thread-per-core design.
