Replay fanout
broadcast_with_replay extends Stream::broadcast with a sliding replay tail.
#![allow(unused)] fn main() { use id_effect::{Stream, broadcast_with_replay, run_async}; let src = Stream::from_iterable(1..=100); let (mut branches, pump) = run_async( broadcast_with_replay(src, /* hub */ 64, /* replay tail */ 16, /* branches */ 2), (), ) .await?; // Run `pump` concurrently with pulls on each branch (same pattern as `broadcast`). }
hub_capacity— slidingPubSubring size.replay_len— number of recent items retained in the shared replay buffer (also seeded into each branch buffer as the pump runs).Stream::broadcast_replay(cap, branches)— convenience wrapper usingreplay_len == cap.
Return shape matches broadcast: (Vec<Stream<…>>, pump_effect).