Skip to main content

hydro_lang/live_collections/stream/
networking.rs

1//! Networking APIs for [`Stream`].
2
3use std::marker::PhantomData;
4
5use serde::Serialize;
6use serde::de::DeserializeOwned;
7use stageleft::{q, quote_type};
8use syn::parse_quote;
9
10use super::{ExactlyOnce, MinOrder, Ordering, Stream, TotalOrder};
11use crate::compile::builder::ExternalPortId;
12use crate::compile::ir::{
13    DebugInstantiate, HydroIrOpMetadata, HydroNode, HydroRoot, NetworkRecv, NetworkSend,
14};
15use crate::live_collections::boundedness::{Boundedness, Unbounded};
16use crate::live_collections::keyed_singleton::{KeyedSingleton, MonotonicKeys};
17use crate::live_collections::keyed_stream::KeyedStream;
18use crate::live_collections::sliced::sliced;
19use crate::live_collections::stream::Retries;
20#[cfg(feature = "sim")]
21use crate::location::LocationKey;
22use crate::location::cluster::{ClusterIds, Consistency, NoConsistency};
23#[cfg(stageleft_runtime)]
24use crate::location::dynamic::DynLocation;
25use crate::location::external_process::ExternalBincodeStream;
26use crate::location::{Cluster, External, Location, MemberId, MembershipEvent, Process};
27use crate::networking::{NetworkFor, TCP};
28use crate::nondet::{NonDet, nondet};
29use crate::properties::manual_proof;
30#[cfg(feature = "sim")]
31use crate::sim::SimReceiver;
32use crate::staging_util::get_this_crate;
33
34// same as the one in `hydro_std`, but internal use only
35fn track_membership<'a, C, L: Location<'a>>(
36    membership: KeyedStream<MemberId<C>, MembershipEvent, L, Unbounded>,
37) -> KeyedSingleton<MemberId<C>, bool, L, MonotonicKeys> {
38    membership.fold(
39        q!(|| false),
40        q!(|present, event| {
41            match event {
42                MembershipEvent::Joined => *present = true,
43                MembershipEvent::Left => *present = false,
44            }
45        }),
46    )
47}
48
49fn serialize_bincode_with_type(is_demux: bool, t_type: &syn::Type) -> syn::Expr {
50    let root = get_this_crate();
51
52    if is_demux {
53        parse_quote! {
54            #root::runtime_support::stageleft::runtime_support::fn1_type_hint::<(#root::__staged::location::MemberId<_>, #t_type), _>(
55                |(id, data)| {
56                    (id.into_tagless(), #root::runtime_support::bincode::serialize(&data).unwrap().into())
57                }
58            )
59        }
60    } else {
61        parse_quote! {
62            #root::runtime_support::stageleft::runtime_support::fn1_type_hint::<#t_type, _>(
63                |data| {
64                    #root::runtime_support::bincode::serialize(&data).unwrap().into()
65                }
66            )
67        }
68    }
69}
70
71pub(crate) fn serialize_bincode<T: Serialize>(is_demux: bool) -> syn::Expr {
72    serialize_bincode_with_type(is_demux, &quote_type::<T>())
73}
74
75fn deserialize_bincode_with_type(tagged: Option<&syn::Type>, t_type: &syn::Type) -> syn::Expr {
76    let root = get_this_crate();
77    if let Some(c_type) = tagged {
78        parse_quote! {
79            |res| {
80                let (id, b) = res.unwrap();
81                (#root::__staged::location::MemberId::<#c_type>::from_tagless(id as #root::__staged::location::TaglessMemberId), #root::runtime_support::bincode::deserialize::<#t_type>(&b).unwrap())
82            }
83        }
84    } else {
85        parse_quote! {
86            |res| {
87                #root::runtime_support::bincode::deserialize::<#t_type>(&res.unwrap()).unwrap()
88            }
89        }
90    }
91}
92
93pub(crate) fn deserialize_bincode<T: DeserializeOwned>(tagged: Option<&syn::Type>) -> syn::Expr {
94    deserialize_bincode_with_type(tagged, &quote_type::<T>())
95}
96
97impl<'a, T, L, B: Boundedness, O: Ordering, R: Retries> Stream<T, Process<'a, L>, B, O, R> {
98    #[deprecated = "use Stream::send(..., TCP.fail_stop().bincode()) instead"]
99    /// "Moves" elements of this stream to a new distributed location by sending them over the network,
100    /// using [`bincode`] to serialize/deserialize messages.
101    ///
102    /// The returned stream captures the elements received at the destination, where values will
103    /// asynchronously arrive over the network. Sending from a [`Process`] to another [`Process`]
104    /// preserves ordering and retries guarantees by using a single TCP channel to send the values. The
105    /// recipient is guaranteed to receive a _prefix_ or the sent messages; if the TCP connection is
106    /// dropped no further messages will be sent.
107    ///
108    /// # Example
109    /// ```rust
110    /// # #[cfg(feature = "deploy")] {
111    /// # use hydro_lang::prelude::*;
112    /// # use futures::StreamExt;
113    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p_out| {
114    /// let p1 = flow.process::<()>();
115    /// let numbers: Stream<_, Process<_>, Bounded> = p1.source_iter(q!(vec![1, 2, 3]));
116    /// let p2 = flow.process::<()>();
117    /// let on_p2: Stream<_, Process<_>, Unbounded> = numbers.send_bincode(&p2);
118    /// // 1, 2, 3
119    /// # on_p2.send_bincode(&p_out)
120    /// # }, |mut stream| async move {
121    /// # for w in 1..=3 {
122    /// #     assert_eq!(stream.next().await, Some(w));
123    /// # }
124    /// # }));
125    /// # }
126    /// ```
127    pub fn send_bincode<L2>(
128        self,
129        other: &Process<'a, L2>,
130    ) -> Stream<T, Process<'a, L2>, Unbounded, O, R>
131    where
132        T: Serialize + DeserializeOwned,
133    {
134        self.send(other, TCP.fail_stop().bincode())
135    }
136
137    /// "Moves" elements of this stream to a new distributed location by sending them over the network,
138    /// using the configuration in `via` to set up the message transport.
139    ///
140    /// The returned stream captures the elements received at the destination, where values will
141    /// asynchronously arrive over the network. Sending from a [`Process`] to another [`Process`]
142    /// preserves ordering and retries guarantees when using a single TCP channel to send the values.
143    /// The recipient is guaranteed to receive a _prefix_ or the sent messages; if the connection is
144    /// dropped no further messages will be sent.
145    ///
146    /// # Example
147    /// ```rust
148    /// # #[cfg(feature = "deploy")] {
149    /// # use hydro_lang::prelude::*;
150    /// # use futures::StreamExt;
151    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p_out| {
152    /// let p1 = flow.process::<()>();
153    /// let numbers: Stream<_, Process<_>, Bounded> = p1.source_iter(q!(vec![1, 2, 3]));
154    /// let p2 = flow.process::<()>();
155    /// let on_p2: Stream<_, Process<_>, Unbounded> = numbers.send(&p2, TCP.fail_stop().bincode());
156    /// // 1, 2, 3
157    /// # on_p2.send(&p_out, TCP.fail_stop().bincode())
158    /// # }, |mut stream| async move {
159    /// # for w in 1..=3 {
160    /// #     assert_eq!(stream.next().await, Some(w));
161    /// # }
162    /// # }));
163    /// # }
164    /// ```
165    pub fn send<L2, N: NetworkFor<T>>(
166        self,
167        to: &Process<'a, L2>,
168        via: N,
169    ) -> Stream<T, Process<'a, L2>, Unbounded, <O as MinOrder<N::OrderingGuarantee>>::Min, R>
170    where
171        O: MinOrder<N::OrderingGuarantee>,
172    {
173        let name = via.name();
174        assert!(
175            !to.multiversioned() || name.is_some(),
176            "Cannot send to a multiversioned location without a channel name. Please provide a name for the network."
177        );
178
179        let (serialize, deserialize) = if N::is_embedded() {
180            (
181                NetworkSend::Embedded {
182                    tag: None,
183                    element_type: quote_type::<T>().into(),
184                },
185                NetworkRecv::Embedded {
186                    tag: None,
187                    element_type: quote_type::<T>().into(),
188                },
189            )
190        } else {
191            (
192                NetworkSend::Custom {
193                    serialize_fn: Some(N::serialize_thunk(false).into()),
194                },
195                NetworkRecv::Custom {
196                    deserialize_fn: Some(N::deserialize_thunk(None).into()),
197                },
198            )
199        };
200
201        Stream::new(
202            to.clone(),
203            HydroNode::Network {
204                name: name.map(ToOwned::to_owned),
205                networking_info: N::networking_info(),
206                serialize,
207                deserialize,
208                instantiate_fn: DebugInstantiate::Building,
209                input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
210                metadata: to.new_node_metadata(Stream::<
211                    T,
212                    Process<'a, L2>,
213                    Unbounded,
214                    <O as MinOrder<N::OrderingGuarantee>>::Min,
215                    R,
216                >::collection_kind()),
217            },
218        )
219    }
220
221    #[deprecated = "use Stream::broadcast(..., TCP.fail_stop().bincode()) instead"]
222    /// Broadcasts elements of this stream to all members of a cluster by sending them over the network,
223    /// using [`bincode`] to serialize/deserialize messages.
224    ///
225    /// Each element in the stream will be sent to **every** member of the cluster based on the latest
226    /// membership information. This is a common pattern in distributed systems for broadcasting data to
227    /// all nodes in a cluster. Unlike [`Stream::demux_bincode`], which requires `(MemberId, T)` tuples to
228    /// target specific members, `broadcast_bincode` takes a stream of **only data elements** and sends
229    /// each element to all cluster members.
230    ///
231    /// # Non-Determinism
232    /// The set of cluster members may asynchronously change over time. Each element is only broadcast
233    /// to the current cluster members _at that point in time_. Depending on when we are notified of
234    /// membership changes, we will broadcast each element to different members.
235    ///
236    /// # Example
237    /// ```rust
238    /// # #[cfg(feature = "deploy")] {
239    /// # use hydro_lang::prelude::*;
240    /// # use futures::StreamExt;
241    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
242    /// let p1 = flow.process::<()>();
243    /// let workers: Cluster<()> = flow.cluster::<()>();
244    /// let numbers: Stream<_, Process<_>, _> = p1.source_iter(q!(vec![123]));
245    /// let on_worker: Stream<_, Cluster<_>, _> = numbers.broadcast_bincode(&workers, nondet!(/** assuming stable membership */));
246    /// # on_worker.send_bincode(&p2).entries()
247    /// // if there are 4 members in the cluster, each receives one element
248    /// // - MemberId::<()>(0): [123]
249    /// // - MemberId::<()>(1): [123]
250    /// // - MemberId::<()>(2): [123]
251    /// // - MemberId::<()>(3): [123]
252    /// # }, |mut stream| async move {
253    /// # let mut results = Vec::new();
254    /// # for w in 0..4 {
255    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
256    /// # }
257    /// # results.sort();
258    /// # assert_eq!(results, vec!["(MemberId::<()>(0), 123)", "(MemberId::<()>(1), 123)", "(MemberId::<()>(2), 123)", "(MemberId::<()>(3), 123)"]);
259    /// # }));
260    /// # }
261    /// ```
262    pub fn broadcast_bincode<L2: 'a>(
263        self,
264        other: &Cluster<'a, L2>,
265        nondet_membership: NonDet,
266    ) -> Stream<T, Cluster<'a, L2>, Unbounded, O, R>
267    where
268        T: Clone + Serialize + DeserializeOwned,
269    {
270        self.broadcast(other, TCP.fail_stop().bincode(), nondet_membership)
271    }
272
273    /// Broadcasts elements of this stream to all members of a cluster by sending them over the network,
274    /// using the configuration in `via` to set up the message transport.
275    ///
276    /// Each element in the stream will be sent to **every** member of the cluster based on the latest
277    /// membership information. This is a common pattern in distributed systems for broadcasting data to
278    /// all nodes in a cluster. Unlike [`Stream::demux`], which requires `(MemberId, T)` tuples to
279    /// target specific members, `broadcast` takes a stream of **only data elements** and sends
280    /// each element to all cluster members.
281    ///
282    /// # Non-Determinism
283    /// The set of cluster members may asynchronously change over time. Each element is only broadcast
284    /// to the current cluster members _at that point in time_. Depending on when we are notified of
285    /// membership changes, we will broadcast each element to different members.
286    ///
287    /// # Example
288    /// ```rust
289    /// # #[cfg(feature = "deploy")] {
290    /// # use hydro_lang::prelude::*;
291    /// # use futures::StreamExt;
292    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
293    /// let p1 = flow.process::<()>();
294    /// let workers: Cluster<()> = flow.cluster::<()>();
295    /// let numbers: Stream<_, Process<_>, _> = p1.source_iter(q!(vec![123]));
296    /// let on_worker: Stream<_, Cluster<_>, _> = numbers.broadcast(&workers, TCP.fail_stop().bincode(), nondet!(/** assuming stable membership */));
297    /// # on_worker.send(&p2, TCP.fail_stop().bincode()).entries()
298    /// // if there are 4 members in the cluster, each receives one element
299    /// // - MemberId::<()>(0): [123]
300    /// // - MemberId::<()>(1): [123]
301    /// // - MemberId::<()>(2): [123]
302    /// // - MemberId::<()>(3): [123]
303    /// # }, |mut stream| async move {
304    /// # let mut results = Vec::new();
305    /// # for w in 0..4 {
306    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
307    /// # }
308    /// # results.sort();
309    /// # assert_eq!(results, vec!["(MemberId::<()>(0), 123)", "(MemberId::<()>(1), 123)", "(MemberId::<()>(2), 123)", "(MemberId::<()>(3), 123)"]);
310    /// # }));
311    /// # }
312    /// ```
313    pub fn broadcast<L2: 'a, N: NetworkFor<T>>(
314        self,
315        to: &Cluster<'a, L2>,
316        via: N,
317        nondet_membership: NonDet,
318    ) -> Stream<T, Cluster<'a, L2>, Unbounded, <O as MinOrder<N::OrderingGuarantee>>::Min, R>
319    where
320        T: Clone,
321        O: MinOrder<N::OrderingGuarantee>,
322    {
323        // TODO(#1875): the membership snapshot below is over a `KeyedSingleton`, and keyed
324        // sim hooks do not exist yet. Once they do, expose a composite hook payload here
325        // (`NonDet<(Option<KeyedSnapshotHook<..>>, Option<BatchHook<T, O, R>>)>`) so tests
326        // can script the membership snapshot and the element batching independently.
327        let ids = track_membership(self.location.source_cluster_membership_stream(
328            to,
329            nondet!(/** dropped prefixes don't affect broadcast */),
330        ));
331        sliced! {
332            let members_snapshot = use::snapshot(ids, nondet!(
333                /// membership timing is captured by the caller's guard
334                nondet_membership
335            ));
336            let elements = use::batch(self, nondet!(
337                /// batching timing is captured by the caller's guard
338                nondet_membership
339            ));
340
341            let current_members = members_snapshot.filter(q!(|b| *b));
342            elements.repeat_with_keys(current_members)
343        }
344        .demux(to, via)
345    }
346
347    /// Broadcasts elements of this stream to all members of a cluster,
348    /// assuming membership is closed (fixed at deploy time).
349    ///
350    /// Unlike [`Stream::broadcast`], this does not require a [`NonDet`] guard.
351    /// The membership set is obtained from deploy metadata via
352    /// [`ClusterIds`], producing a
353    /// `Bounded` stream. The cross-product of data × members is fully
354    /// deterministic.
355    ///
356    /// The consistency guarantee of the output depends on the network's failure policy
357    /// ([`NetworkFor::ConsistencyGuarantee`]). Policies like `fail_stop` and
358    /// `lossy_delayed_forever` guarantee that every live member eventually materializes the same
359    /// elements, so the output is
360    /// [`EventualConsistency`](crate::location::cluster::EventualConsistency). A plain `lossy`
361    /// policy can drop individual messages for some members while delivering them to others, so
362    /// replicas may permanently diverge and the output only has
363    /// [`NoConsistency`].
364    ///
365    /// This is only available in deployment targets with static cluster
366    /// membership (legacy Hydro Deploy and simulation). There are no late
367    /// joiners in that context, so broadcast receivers are guaranteed to
368    /// get data from the start of the stream. On dynamic targets
369    /// (e.g. ECS), use [`Stream::broadcast`] instead.
370    ///
371    /// # Example
372    /// ```rust
373    /// # #[cfg(feature = "deploy")] {
374    /// # use hydro_lang::prelude::*;
375    /// # use futures::StreamExt;
376    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
377    /// let p1 = flow.process::<()>();
378    /// let workers: Cluster<()> = flow.cluster::<()>();
379    /// let numbers: Stream<_, Process<_>, _> = p1.source_iter(q!(vec![123]));
380    /// let on_worker = numbers.broadcast_closed(&workers, TCP.fail_stop().bincode());
381    /// # on_worker.send(&p2, TCP.fail_stop().bincode()).entries()
382    /// // each of the 4 cluster members receives 123
383    /// # }, |mut stream| async move {
384    /// # let mut results = Vec::new();
385    /// # for _ in 0..4 {
386    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
387    /// # }
388    /// # results.sort();
389    /// # assert_eq!(results, vec!["(MemberId::<()>(0), 123)", "(MemberId::<()>(1), 123)", "(MemberId::<()>(2), 123)", "(MemberId::<()>(3), 123)"]);
390    /// # }));
391    /// # }
392    /// ```
393    pub fn broadcast_closed<L2: 'a, N: NetworkFor<T>>(
394        self,
395        to: &Cluster<'a, L2>,
396        via: N,
397    ) -> Stream<
398        T,
399        Cluster<'a, L2, N::ConsistencyGuarantee>,
400        Unbounded,
401        <O as MinOrder<N::OrderingGuarantee>>::Min,
402        R,
403    >
404    where
405        T: Clone,
406        O: MinOrder<N::OrderingGuarantee>,
407    {
408        let cluster_ids = ClusterIds {
409            key: to.key,
410            _phantom: PhantomData,
411        };
412        let member_ids = self.location.source_iter(q!(cluster_ids
413            .iter()
414            .map(|id| MemberId::from_tagless(id.clone()))));
415
416        // Late joiners will receive no data from this broadcast, which is
417        // future-monotone and eventually consistent (a safe under-approximation).
418        self.cross_product(member_ids)
419            .map(q!(|(data, member_id)| (member_id, data)))
420            .into_keyed()
421            .demux(to, via)
422            .assert_has_consistency_of_trusted(manual_proof!(
423                /// With a network whose failure policy delivers the same messages to every live
424                /// member (tracked by `NetworkFor::ConsistencyGuarantee`), a closed broadcast
425                /// will materialize the same elements on each member.
426            ))
427    }
428
429    /// Sends the elements of this stream to an external (non-Hydro) process, using [`bincode`]
430    /// serialization. The external process can receive these elements by establishing a TCP
431    /// connection and decoding using [`tokio_util::codec::LengthDelimitedCodec`].
432    ///
433    /// # Example
434    /// ```rust
435    /// # #[cfg(feature = "deploy")] {
436    /// # use hydro_lang::prelude::*;
437    /// # use futures::StreamExt;
438    /// # tokio_test::block_on(async move {
439    /// let mut flow = FlowBuilder::new();
440    /// let process = flow.process::<()>();
441    /// let numbers: Stream<_, Process<_>, Bounded> = process.source_iter(q!(vec![1, 2, 3]));
442    /// let external = flow.external::<()>();
443    /// let external_handle = numbers.send_bincode_external(&external);
444    ///
445    /// let mut deployment = hydro_deploy::Deployment::new();
446    /// let nodes = flow
447    ///     .with_process(&process, deployment.Localhost())
448    ///     .with_external(&external, deployment.Localhost())
449    ///     .deploy(&mut deployment);
450    ///
451    /// deployment.deploy().await.unwrap();
452    /// // establish the TCP connection
453    /// let mut external_recv_stream = nodes.connect(external_handle).await;
454    /// deployment.start().await.unwrap();
455    ///
456    /// for w in 1..=3 {
457    ///     assert_eq!(external_recv_stream.next().await, Some(w));
458    /// }
459    /// # });
460    /// # }
461    /// ```
462    pub fn send_bincode_external<L2>(
463        self,
464        other: &External<'_, L2>,
465    ) -> ExternalBincodeStream<T, O, R>
466    where
467        T: Serialize + DeserializeOwned,
468    {
469        let external_port_id =
470            self.register_serialized_external_port(other, serialize_bincode::<T>(false));
471
472        ExternalBincodeStream {
473            process_key: other.key,
474            port_id: external_port_id,
475            _phantom: PhantomData,
476        }
477    }
478
479    // TODO: Add a codec-parameterized external stream handle once deployment supports custom codecs.
480    fn register_serialized_external_port<L2>(
481        self,
482        other: &External<'_, L2>,
483        serialize_pipeline: syn::Expr,
484    ) -> ExternalPortId {
485        let mut flow_state_borrow = self.location.flow_state().borrow_mut();
486
487        let external_port_id = flow_state_borrow.next_external_port();
488
489        flow_state_borrow.push_root(HydroRoot::SendExternal {
490            to_external_key: other.key,
491            to_port_id: external_port_id,
492            to_many: false,
493            unpaired: true,
494            serialize_fn: Some(serialize_pipeline.into()),
495            instantiate_fn: DebugInstantiate::Building,
496            input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
497            op_metadata: HydroIrOpMetadata::new(),
498        });
499
500        external_port_id
501    }
502
503    #[cfg(feature = "sim")]
504    /// Sets up a bincode-encoded simulation output port for this stream, allowing test code to
505    /// receive elements sent to this stream during simulation. Use [`Stream::sim_output_with`] to
506    /// select another codec.
507    pub fn sim_output(self) -> SimReceiver<T, O, R>
508    where
509        T: Serialize + DeserializeOwned,
510    {
511        self.sim_output_with::<crate::sim::codec::BincodeCodec>()
512    }
513
514    #[cfg(feature = "sim")]
515    /// Sets up a simulation output port using the codec `C`, allowing test code to receive
516    /// elements sent to this stream during simulation: `stream.sim_output_with::<MyCodec>()`.
517    /// Custom codecs implement [`SimCodec`](crate::sim::codec::SimCodec), which documents
518    /// where they must be defined.
519    pub fn sim_output_with<C>(self) -> SimReceiver<T, O, R>
520    where
521        C: crate::sim::codec::SimCodec<T>,
522    {
523        let external_location: External<'a, ()> = External {
524            key: LocationKey::FIRST,
525            flow_state: self.location.flow_state().clone(),
526            _phantom: PhantomData,
527        };
528
529        let external_port_id = self.register_serialized_external_port(
530            &external_location,
531            crate::sim::codec::staged_serialize::<T, C>(),
532        );
533
534        SimReceiver(external_port_id, PhantomData, C::decode)
535    }
536}
537
538impl<'a, T, L: Location<'a>, B: Boundedness> Stream<T, L, B, TotalOrder, ExactlyOnce> {
539    /// Creates an external output for embedded deployment mode.
540    ///
541    /// The `name` parameter specifies the name of the field in the generated
542    /// `EmbeddedOutputs` struct that will receive elements from this stream.
543    /// The generated function will accept an `EmbeddedOutputs` struct with an
544    /// `impl FnMut(T)` field with this name.
545    pub fn embedded_output(self, name: impl Into<String>) {
546        let ident = syn::Ident::new(&name.into(), proc_macro2::Span::call_site());
547
548        self.location
549            .flow_state()
550            .borrow_mut()
551            .push_root(HydroRoot::EmbeddedOutput {
552                ident,
553                input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
554                op_metadata: HydroIrOpMetadata::new(),
555            });
556    }
557}
558
559impl<'a, T, L, L2, B: Boundedness, O: Ordering, R: Retries>
560    Stream<(MemberId<L2>, T), Process<'a, L>, B, O, R>
561{
562    #[deprecated = "use Stream::demux(..., TCP.fail_stop().bincode()) instead"]
563    /// Sends elements of this stream to specific members of a cluster, identified by a [`MemberId`],
564    /// using [`bincode`] to serialize/deserialize messages.
565    ///
566    /// Each element in the stream must be a tuple `(MemberId<L2>, T)` where the first element
567    /// specifies which cluster member should receive the data. Unlike [`Stream::broadcast_bincode`],
568    /// this API allows precise targeting of specific cluster members rather than broadcasting to
569    /// all members.
570    ///
571    /// # Example
572    /// ```rust
573    /// # #[cfg(feature = "deploy")] {
574    /// # use hydro_lang::prelude::*;
575    /// # use futures::StreamExt;
576    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
577    /// let p1 = flow.process::<()>();
578    /// let workers: Cluster<()> = flow.cluster::<()>();
579    /// let numbers: Stream<_, Process<_>, _> = p1.source_iter(q!(vec![0, 1, 2, 3]));
580    /// let on_worker: Stream<_, Cluster<_>, _> = numbers
581    ///     .map(q!(|x| (hydro_lang::location::MemberId::from_raw_id(x), x)))
582    ///     .demux_bincode(&workers);
583    /// # on_worker.send_bincode(&p2).entries()
584    /// // if there are 4 members in the cluster, each receives one element
585    /// // - MemberId::<()>(0): [0]
586    /// // - MemberId::<()>(1): [1]
587    /// // - MemberId::<()>(2): [2]
588    /// // - MemberId::<()>(3): [3]
589    /// # }, |mut stream| async move {
590    /// # let mut results = Vec::new();
591    /// # for w in 0..4 {
592    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
593    /// # }
594    /// # results.sort();
595    /// # assert_eq!(results, vec!["(MemberId::<()>(0), 0)", "(MemberId::<()>(1), 1)", "(MemberId::<()>(2), 2)", "(MemberId::<()>(3), 3)"]);
596    /// # }));
597    /// # }
598    /// ```
599    pub fn demux_bincode(
600        self,
601        other: &Cluster<'a, L2>,
602    ) -> Stream<T, Cluster<'a, L2>, Unbounded, O, R>
603    where
604        T: Serialize + DeserializeOwned,
605    {
606        self.demux(other, TCP.fail_stop().bincode())
607    }
608
609    /// Sends elements of this stream to specific members of a cluster, identified by a [`MemberId`],
610    /// using the configuration in `via` to set up the message transport.
611    ///
612    /// Each element in the stream must be a tuple `(MemberId<L2>, T)` where the first element
613    /// specifies which cluster member should receive the data. Unlike [`Stream::broadcast`],
614    /// this API allows precise targeting of specific cluster members rather than broadcasting to
615    /// all members.
616    ///
617    /// # Example
618    /// ```rust
619    /// # #[cfg(feature = "deploy")] {
620    /// # use hydro_lang::prelude::*;
621    /// # use futures::StreamExt;
622    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
623    /// let p1 = flow.process::<()>();
624    /// let workers: Cluster<()> = flow.cluster::<()>();
625    /// let numbers: Stream<_, Process<_>, _> = p1.source_iter(q!(vec![0, 1, 2, 3]));
626    /// let on_worker: Stream<_, Cluster<_>, _> = numbers
627    ///     .map(q!(|x| (hydro_lang::location::MemberId::from_raw_id(x), x)))
628    ///     .demux(&workers, TCP.fail_stop().bincode());
629    /// # on_worker.send(&p2, TCP.fail_stop().bincode()).entries()
630    /// // if there are 4 members in the cluster, each receives one element
631    /// // - MemberId::<()>(0): [0]
632    /// // - MemberId::<()>(1): [1]
633    /// // - MemberId::<()>(2): [2]
634    /// // - MemberId::<()>(3): [3]
635    /// # }, |mut stream| async move {
636    /// # let mut results = Vec::new();
637    /// # for w in 0..4 {
638    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
639    /// # }
640    /// # results.sort();
641    /// # assert_eq!(results, vec!["(MemberId::<()>(0), 0)", "(MemberId::<()>(1), 1)", "(MemberId::<()>(2), 2)", "(MemberId::<()>(3), 3)"]);
642    /// # }));
643    /// # }
644    /// ```
645    pub fn demux<N: NetworkFor<T>>(
646        self,
647        to: &Cluster<'a, L2>,
648        via: N,
649    ) -> Stream<
650        T,
651        Cluster<'a, L2, NoConsistency>,
652        Unbounded,
653        <O as MinOrder<N::OrderingGuarantee>>::Min,
654        R,
655    >
656    where
657        O: MinOrder<N::OrderingGuarantee>,
658    {
659        self.into_keyed().demux(to, via)
660    }
661}
662
663impl<'a, T, L, B: Boundedness> Stream<T, Process<'a, L>, B, TotalOrder, ExactlyOnce> {
664    #[deprecated = "use Stream::round_robin(..., TCP.fail_stop().bincode()) instead"]
665    /// Distributes elements of this stream to cluster members in a round-robin fashion, using
666    /// [`bincode`] to serialize/deserialize messages.
667    ///
668    /// This provides load balancing by evenly distributing work across cluster members. The
669    /// distribution is deterministic based on element order - the first element goes to member 0,
670    /// the second to member 1, and so on, wrapping around when reaching the end of the member list.
671    ///
672    /// # Non-Determinism
673    /// The set of cluster members may asynchronously change over time. Each element is distributed
674    /// based on the current cluster membership _at that point in time_. Depending on when cluster
675    /// members join and leave, the round-robin pattern will change. Furthermore, even when the
676    /// membership is stable, the order of members in the round-robin pattern may change across runs.
677    ///
678    /// # Ordering Requirements
679    /// This method is only available on streams with [`TotalOrder`] and [`ExactlyOnce`], since the
680    /// order of messages and retries affects the round-robin pattern.
681    ///
682    /// # Example
683    /// ```rust
684    /// # #[cfg(feature = "deploy")] {
685    /// # use hydro_lang::prelude::*;
686    /// # use hydro_lang::live_collections::stream::{TotalOrder, ExactlyOnce};
687    /// # use futures::StreamExt;
688    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
689    /// let p1 = flow.process::<()>();
690    /// let workers: Cluster<()> = flow.cluster::<()>();
691    /// let numbers: Stream<_, Process<_>, _, TotalOrder, ExactlyOnce> = p1.source_iter(q!(vec![1, 2, 3, 4]));
692    /// let on_worker: Stream<_, Cluster<_>, _> = numbers.round_robin_bincode(&workers, nondet!(/** assuming stable membership */));
693    /// on_worker.send_bincode(&p2)
694    /// # .first().values() // we use first to assert that each member gets one element
695    /// // with 4 cluster members, elements are distributed (with a non-deterministic round-robin order):
696    /// // - MemberId::<()>(?): [1]
697    /// // - MemberId::<()>(?): [2]
698    /// // - MemberId::<()>(?): [3]
699    /// // - MemberId::<()>(?): [4]
700    /// # }, |mut stream| async move {
701    /// # let mut results = Vec::new();
702    /// # for w in 0..4 {
703    /// #     results.push(stream.next().await.unwrap());
704    /// # }
705    /// # results.sort();
706    /// # assert_eq!(results, vec![1, 2, 3, 4]);
707    /// # }));
708    /// # }
709    /// ```
710    pub fn round_robin_bincode<L2: 'a>(
711        self,
712        other: &Cluster<'a, L2>,
713        nondet_membership: NonDet,
714    ) -> Stream<T, Cluster<'a, L2>, Unbounded, TotalOrder, ExactlyOnce>
715    where
716        T: Serialize + DeserializeOwned,
717    {
718        self.round_robin(other, TCP.fail_stop().bincode(), nondet_membership)
719    }
720
721    /// Distributes elements of this stream to cluster members in a round-robin fashion, using
722    /// the configuration in `via` to set up the message transport.
723    ///
724    /// This provides load balancing by evenly distributing work across cluster members. The
725    /// distribution is deterministic based on element order - the first element goes to member 0,
726    /// the second to member 1, and so on, wrapping around when reaching the end of the member list.
727    ///
728    /// # Non-Determinism
729    /// The set of cluster members may asynchronously change over time. Each element is distributed
730    /// based on the current cluster membership _at that point in time_. Depending on when cluster
731    /// members join and leave, the round-robin pattern will change. Furthermore, even when the
732    /// membership is stable, the order of members in the round-robin pattern may change across runs.
733    ///
734    /// # Ordering Requirements
735    /// This method is only available on streams with [`TotalOrder`] and [`ExactlyOnce`], since the
736    /// order of messages and retries affects the round-robin pattern.
737    ///
738    /// # Example
739    /// ```rust
740    /// # #[cfg(feature = "deploy")] {
741    /// # use hydro_lang::prelude::*;
742    /// # use hydro_lang::live_collections::stream::{TotalOrder, ExactlyOnce};
743    /// # use futures::StreamExt;
744    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
745    /// let p1 = flow.process::<()>();
746    /// let workers: Cluster<()> = flow.cluster::<()>();
747    /// let numbers: Stream<_, Process<_>, _, TotalOrder, ExactlyOnce> = p1.source_iter(q!(vec![1, 2, 3, 4]));
748    /// let on_worker: Stream<_, Cluster<_>, _> = numbers.round_robin(&workers, TCP.fail_stop().bincode(), nondet!(/** assuming stable membership */));
749    /// on_worker.send(&p2, TCP.fail_stop().bincode())
750    /// # .first().values() // we use first to assert that each member gets one element
751    /// // with 4 cluster members, elements are distributed (with a non-deterministic round-robin order):
752    /// // - MemberId::<()>(?): [1]
753    /// // - MemberId::<()>(?): [2]
754    /// // - MemberId::<()>(?): [3]
755    /// // - MemberId::<()>(?): [4]
756    /// # }, |mut stream| async move {
757    /// # let mut results = Vec::new();
758    /// # for w in 0..4 {
759    /// #     results.push(stream.next().await.unwrap());
760    /// # }
761    /// # results.sort();
762    /// # assert_eq!(results, vec![1, 2, 3, 4]);
763    /// # }));
764    /// # }
765    /// ```
766    pub fn round_robin<L2: 'a, N: NetworkFor<T>>(
767        self,
768        to: &Cluster<'a, L2>,
769        via: N,
770        nondet_membership: NonDet,
771    ) -> Stream<T, Cluster<'a, L2>, Unbounded, N::OrderingGuarantee, ExactlyOnce> {
772        // TODO(#1875): the membership snapshot below is over a `KeyedSingleton`, and keyed
773        // sim hooks do not exist yet. Once they do, expose a composite hook payload here
774        // (`NonDet<(Option<KeyedSnapshotHook<..>>, Option<BatchHook<T, O, R>>)>`) so tests
775        // can script the membership snapshot and the element batching independently.
776        let ids = track_membership(self.location.source_cluster_membership_stream(
777            to,
778            nondet!(/** dropped prefixes don't affect broadcast */),
779        ));
780        sliced! {
781            let members_snapshot = use::snapshot(ids, nondet!(
782                /// membership timing is captured by the caller's guard
783                nondet_membership
784            ));
785            let elements = use::batch(self.enumerate(), nondet!(
786                /// batching timing is captured by the caller's guard
787                nondet_membership
788            ));
789
790            let current_members = members_snapshot
791                .filter(q!(|b| *b))
792                .keys()
793                .assume_ordering::<TotalOrder>(nondet!(/** membership timing is captured by the caller guard */ nondet_membership))
794                .collect_vec();
795
796            elements
797                .cross_singleton(current_members)
798                .filter_map(q!(|(data, members)| {
799                    if members.is_empty() {
800                        None
801                    } else {
802                        Some((members[data.0 % members.len()].clone(), data.1))
803                    }
804                }))
805        }
806        .demux(to, via)
807    }
808}
809
810impl<'a, T, L, B: Boundedness, C: Consistency>
811    Stream<T, Cluster<'a, L, C>, B, TotalOrder, ExactlyOnce>
812{
813    #[deprecated = "use Stream::round_robin(..., TCP.fail_stop().bincode()) instead"]
814    /// Distributes elements of this stream to cluster members in a round-robin fashion, using
815    /// [`bincode`] to serialize/deserialize messages.
816    ///
817    /// This provides load balancing by evenly distributing work across cluster members. The
818    /// distribution is deterministic based on element order - the first element goes to member 0,
819    /// the second to member 1, and so on, wrapping around when reaching the end of the member list.
820    ///
821    /// # Non-Determinism
822    /// The set of cluster members may asynchronously change over time. Each element is distributed
823    /// based on the current cluster membership _at that point in time_. Depending on when cluster
824    /// members join and leave, the round-robin pattern will change. Furthermore, even when the
825    /// membership is stable, the order of members in the round-robin pattern may change across runs.
826    ///
827    /// # Ordering Requirements
828    /// This method is only available on streams with [`TotalOrder`] and [`ExactlyOnce`], since the
829    /// order of messages and retries affects the round-robin pattern.
830    ///
831    /// # Example
832    /// ```rust
833    /// # #[cfg(feature = "deploy")] {
834    /// # use hydro_lang::prelude::*;
835    /// # use hydro_lang::live_collections::stream::{TotalOrder, ExactlyOnce, NoOrder};
836    /// # use hydro_lang::location::MemberId;
837    /// # use futures::StreamExt;
838    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
839    /// let p1 = flow.process::<()>();
840    /// let workers1: Cluster<()> = flow.cluster::<()>();
841    /// let workers2: Cluster<()> = flow.cluster::<()>();
842    /// let numbers: Stream<_, Process<_>, _, TotalOrder, ExactlyOnce> = p1.source_iter(q!(0..=16));
843    /// let on_worker1: Stream<_, Cluster<_>, _> = numbers.round_robin_bincode(&workers1, nondet!(/** assuming stable membership */));
844    /// let on_worker2: Stream<_, Cluster<_>, _> = on_worker1.round_robin_bincode(&workers2, nondet!(/** assuming stable membership */)).entries().assume_ordering(nondet!(/** assuming stable membership */));
845    /// on_worker2.send_bincode(&p2)
846    /// # .entries()
847    /// # .map(q!(|(w2, (w1, v))| ((w2, w1), v)))
848    /// # }, |mut stream| async move {
849    /// # let mut results = Vec::new();
850    /// # let mut locations = std::collections::HashSet::new();
851    /// # for w in 0..=16 {
852    /// #     let (location, v) = stream.next().await.unwrap();
853    /// #     locations.insert(location);
854    /// #     results.push(v);
855    /// # }
856    /// # results.sort();
857    /// # assert_eq!(results, (0..=16).collect::<Vec<_>>());
858    /// # assert_eq!(locations.len(), 16);
859    /// # }));
860    /// # }
861    /// ```
862    pub fn round_robin_bincode<L2: 'a>(
863        self,
864        other: &Cluster<'a, L2>,
865        nondet_membership: NonDet,
866    ) -> KeyedStream<MemberId<L>, T, Cluster<'a, L2>, Unbounded, TotalOrder, ExactlyOnce>
867    where
868        T: Serialize + DeserializeOwned,
869    {
870        self.round_robin(other, TCP.fail_stop().bincode(), nondet_membership)
871    }
872
873    /// Distributes elements of this stream to cluster members in a round-robin fashion, using
874    /// the configuration in `via` to set up the message transport.
875    ///
876    /// This provides load balancing by evenly distributing work across cluster members. The
877    /// distribution is deterministic based on element order - the first element goes to member 0,
878    /// the second to member 1, and so on, wrapping around when reaching the end of the member list.
879    ///
880    /// # Non-Determinism
881    /// The set of cluster members may asynchronously change over time. Each element is distributed
882    /// based on the current cluster membership _at that point in time_. Depending on when cluster
883    /// members join and leave, the round-robin pattern will change. Furthermore, even when the
884    /// membership is stable, the order of members in the round-robin pattern may change across runs.
885    ///
886    /// # Ordering Requirements
887    /// This method is only available on streams with [`TotalOrder`] and [`ExactlyOnce`], since the
888    /// order of messages and retries affects the round-robin pattern.
889    ///
890    /// # Example
891    /// ```rust
892    /// # #[cfg(feature = "deploy")] {
893    /// # use hydro_lang::prelude::*;
894    /// # use hydro_lang::live_collections::stream::{TotalOrder, ExactlyOnce, NoOrder};
895    /// # use hydro_lang::location::MemberId;
896    /// # use futures::StreamExt;
897    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
898    /// let p1 = flow.process::<()>();
899    /// let workers1: Cluster<()> = flow.cluster::<()>();
900    /// let workers2: Cluster<()> = flow.cluster::<()>();
901    /// let numbers: Stream<_, Process<_>, _, TotalOrder, ExactlyOnce> = p1.source_iter(q!(0..=16));
902    /// let on_worker1: Stream<_, Cluster<_>, _> = numbers.round_robin(&workers1, TCP.fail_stop().bincode(), nondet!(/** assuming stable membership */));
903    /// let on_worker2: Stream<_, Cluster<_>, _> = on_worker1.round_robin(&workers2, TCP.fail_stop().bincode(), nondet!(/** assuming stable membership */)).entries().assume_ordering(nondet!(/** assuming stable membership */));
904    /// on_worker2.send(&p2, TCP.fail_stop().bincode())
905    /// # .entries()
906    /// # .map(q!(|(w2, (w1, v))| ((w2, w1), v)))
907    /// # }, |mut stream| async move {
908    /// # let mut results = Vec::new();
909    /// # let mut locations = std::collections::HashSet::new();
910    /// # for w in 0..=16 {
911    /// #     let (location, v) = stream.next().await.unwrap();
912    /// #     locations.insert(location);
913    /// #     results.push(v);
914    /// # }
915    /// # results.sort();
916    /// # assert_eq!(results, (0..=16).collect::<Vec<_>>());
917    /// # assert_eq!(locations.len(), 16);
918    /// # }));
919    /// # }
920    /// ```
921    pub fn round_robin<L2: 'a, N: NetworkFor<T>>(
922        self,
923        to: &Cluster<'a, L2>,
924        via: N,
925        nondet_membership: NonDet,
926    ) -> KeyedStream<MemberId<L>, T, Cluster<'a, L2>, Unbounded, N::OrderingGuarantee, ExactlyOnce>
927    {
928        // TODO(#1875): the membership snapshot below is over a `KeyedSingleton`, and keyed
929        // sim hooks do not exist yet. Once they do, expose a composite hook payload here
930        // (`NonDet<(Option<KeyedSnapshotHook<..>>, Option<BatchHook<T, O, R>>)>`) so tests
931        // can script the membership snapshot and the element batching independently.
932        let ids = track_membership(self.location.source_cluster_membership_stream(
933            to,
934            nondet!(/** dropped prefixes don't affect broadcast */),
935        ));
936        sliced! {
937            let members_snapshot = use::snapshot(ids, nondet!(
938                /// membership timing is captured by the caller's guard
939                nondet_membership
940            ));
941            let elements = use::batch(self.enumerate(), nondet!(
942                /// batching timing is captured by the caller's guard
943                nondet_membership
944            ));
945
946            let current_members = members_snapshot
947                .filter(q!(|b| *b))
948                .keys()
949                .assume_ordering::<TotalOrder>(nondet!(/** membership timing is captured by the caller guard */ nondet_membership))
950                .collect_vec();
951
952            elements
953                .cross_singleton(current_members)
954                .filter_map(q!(|(data, members)| {
955                    if members.is_empty() {
956                        None
957                    } else {
958                        Some((members[data.0 % members.len()].clone(), data.1))
959                    }
960                }))
961        }
962        .demux(to, via)
963    }
964}
965
966impl<'a, T, L, B: Boundedness, C: Consistency, O: Ordering, R: Retries>
967    Stream<T, Cluster<'a, L, C>, B, O, R>
968{
969    #[deprecated = "use Stream::send(..., TCP.fail_stop().bincode()) instead"]
970    /// "Moves" elements of this stream from a cluster to a process by sending them over the network,
971    /// using [`bincode`] to serialize/deserialize messages.
972    ///
973    /// Each cluster member sends its local stream elements, and they are collected at the destination
974    /// as a [`KeyedStream`] where keys identify the source cluster member.
975    ///
976    /// # Example
977    /// ```rust
978    /// # #[cfg(feature = "deploy")] {
979    /// # use hydro_lang::prelude::*;
980    /// # use futures::StreamExt;
981    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, process| {
982    /// let workers: Cluster<()> = flow.cluster::<()>();
983    /// let numbers: Stream<_, Cluster<_>, _> = workers.source_iter(q!(vec![1]));
984    /// let all_received = numbers.send_bincode(&process); // KeyedStream<MemberId<()>, i32, ...>
985    /// # all_received.entries()
986    /// # }, |mut stream| async move {
987    /// // if there are 4 members in the cluster, we should receive 4 elements
988    /// // { MemberId::<()>(0): [1], MemberId::<()>(1): [1], MemberId::<()>(2): [1], MemberId::<()>(3): [1] }
989    /// # let mut results = Vec::new();
990    /// # for w in 0..4 {
991    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
992    /// # }
993    /// # results.sort();
994    /// # assert_eq!(results, vec!["(MemberId::<()>(0), 1)", "(MemberId::<()>(1), 1)", "(MemberId::<()>(2), 1)", "(MemberId::<()>(3), 1)"]);
995    /// # }));
996    /// # }
997    /// ```
998    ///
999    /// If you don't need to know the source for each element, you can use `.values()`
1000    /// to get just the data:
1001    /// ```rust
1002    /// # #[cfg(feature = "deploy")] {
1003    /// # use hydro_lang::prelude::*;
1004    /// # use hydro_lang::live_collections::stream::NoOrder;
1005    /// # use futures::StreamExt;
1006    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, process| {
1007    /// # let workers: Cluster<()> = flow.cluster::<()>();
1008    /// # let numbers: Stream<_, Cluster<_>, _> = workers.source_iter(q!(vec![1]));
1009    /// let values: Stream<i32, _, _, NoOrder> = numbers.send_bincode(&process).values();
1010    /// # values
1011    /// # }, |mut stream| async move {
1012    /// # let mut results = Vec::new();
1013    /// # for w in 0..4 {
1014    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
1015    /// # }
1016    /// # results.sort();
1017    /// // if there are 4 members in the cluster, we should receive 4 elements
1018    /// // 1, 1, 1, 1
1019    /// # assert_eq!(results, vec!["1", "1", "1", "1"]);
1020    /// # }));
1021    /// # }
1022    /// ```
1023    pub fn send_bincode<L2>(
1024        self,
1025        other: &Process<'a, L2>,
1026    ) -> KeyedStream<MemberId<L>, T, Process<'a, L2>, Unbounded, O, R>
1027    where
1028        T: Serialize + DeserializeOwned,
1029    {
1030        self.send(other, TCP.fail_stop().bincode())
1031    }
1032
1033    /// "Moves" elements of this stream from a cluster to a process by sending them over the network,
1034    /// using the configuration in `via` to set up the message transport.
1035    ///
1036    /// Each cluster member sends its local stream elements, and they are collected at the destination
1037    /// as a [`KeyedStream`] where keys identify the source cluster member.
1038    ///
1039    /// # Example
1040    /// ```rust
1041    /// # #[cfg(feature = "deploy")] {
1042    /// # use hydro_lang::prelude::*;
1043    /// # use futures::StreamExt;
1044    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, process| {
1045    /// let workers: Cluster<()> = flow.cluster::<()>();
1046    /// let numbers: Stream<_, Cluster<_>, _> = workers.source_iter(q!(vec![1]));
1047    /// let all_received = numbers.send(&process, TCP.fail_stop().bincode()); // KeyedStream<MemberId<()>, i32, ...>
1048    /// # all_received.entries()
1049    /// # }, |mut stream| async move {
1050    /// // if there are 4 members in the cluster, we should receive 4 elements
1051    /// // { MemberId::<()>(0): [1], MemberId::<()>(1): [1], MemberId::<()>(2): [1], MemberId::<()>(3): [1] }
1052    /// # let mut results = Vec::new();
1053    /// # for w in 0..4 {
1054    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
1055    /// # }
1056    /// # results.sort();
1057    /// # assert_eq!(results, vec!["(MemberId::<()>(0), 1)", "(MemberId::<()>(1), 1)", "(MemberId::<()>(2), 1)", "(MemberId::<()>(3), 1)"]);
1058    /// # }));
1059    /// # }
1060    /// ```
1061    ///
1062    /// If you don't need to know the source for each element, you can use `.values()`
1063    /// to get just the data:
1064    /// ```rust
1065    /// # #[cfg(feature = "deploy")] {
1066    /// # use hydro_lang::prelude::*;
1067    /// # use hydro_lang::live_collections::stream::NoOrder;
1068    /// # use futures::StreamExt;
1069    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, process| {
1070    /// # let workers: Cluster<()> = flow.cluster::<()>();
1071    /// # let numbers: Stream<_, Cluster<_>, _> = workers.source_iter(q!(vec![1]));
1072    /// let values: Stream<i32, _, _, NoOrder> =
1073    ///     numbers.send(&process, TCP.fail_stop().bincode()).values();
1074    /// # values
1075    /// # }, |mut stream| async move {
1076    /// # let mut results = Vec::new();
1077    /// # for w in 0..4 {
1078    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
1079    /// # }
1080    /// # results.sort();
1081    /// // if there are 4 members in the cluster, we should receive 4 elements
1082    /// // 1, 1, 1, 1
1083    /// # assert_eq!(results, vec!["1", "1", "1", "1"]);
1084    /// # }));
1085    /// # }
1086    /// ```
1087    pub fn send<L2, N: NetworkFor<T>>(
1088        self,
1089        to: &Process<'a, L2>,
1090        via: N,
1091    ) -> KeyedStream<
1092        MemberId<L>,
1093        T,
1094        Process<'a, L2>,
1095        Unbounded,
1096        <O as MinOrder<N::OrderingGuarantee>>::Min,
1097        R,
1098    >
1099    where
1100        O: MinOrder<N::OrderingGuarantee>,
1101    {
1102        let name = via.name();
1103        assert!(
1104            !to.multiversioned() || name.is_some(),
1105            "Cannot send to a multiversioned location without a channel name. Please provide a name for the network."
1106        );
1107
1108        let (serialize, deserialize) = if N::is_embedded() {
1109            (
1110                NetworkSend::Embedded {
1111                    tag: None,
1112                    element_type: quote_type::<T>().into(),
1113                },
1114                NetworkRecv::Embedded {
1115                    tag: Some(quote_type::<L>().into()),
1116                    element_type: quote_type::<T>().into(),
1117                },
1118            )
1119        } else {
1120            (
1121                NetworkSend::Custom {
1122                    serialize_fn: Some(N::serialize_thunk(false).into()),
1123                },
1124                NetworkRecv::Custom {
1125                    deserialize_fn: Some(N::deserialize_thunk(Some(&quote_type::<L>())).into()),
1126                },
1127            )
1128        };
1129
1130        let raw_stream: Stream<
1131            (MemberId<L>, T),
1132            Process<'a, L2>,
1133            Unbounded,
1134            <O as MinOrder<N::OrderingGuarantee>>::Min,
1135            R,
1136        > = Stream::new(
1137            to.clone(),
1138            HydroNode::Network {
1139                name: name.map(ToOwned::to_owned),
1140                networking_info: N::networking_info(),
1141                serialize,
1142                deserialize,
1143                instantiate_fn: DebugInstantiate::Building,
1144                input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1145                metadata: to.new_node_metadata(Stream::<
1146                    (MemberId<L>, T),
1147                    Process<'a, L2>,
1148                    Unbounded,
1149                    <O as MinOrder<N::OrderingGuarantee>>::Min,
1150                    R,
1151                >::collection_kind()),
1152            },
1153        );
1154
1155        raw_stream.into_keyed()
1156    }
1157
1158    #[deprecated = "use Stream::broadcast(..., TCP.fail_stop().bincode()) instead"]
1159    /// Broadcasts elements of this stream at each source member to all members of a destination
1160    /// cluster, using [`bincode`] to serialize/deserialize messages.
1161    ///
1162    /// Each source member sends each of its stream elements to **every** member of the cluster
1163    /// based on its latest membership information. Unlike [`Stream::demux_bincode`], which requires
1164    /// `(MemberId, T)` tuples to target specific members, `broadcast_bincode` takes a stream of
1165    /// **only data elements** and sends each element to all cluster members.
1166    ///
1167    /// # Non-Determinism
1168    /// The set of cluster members may asynchronously change over time. Each element is only broadcast
1169    /// to the current cluster members known _at that point in time_ at the source member. Depending
1170    /// on when each source member is notified of membership changes, it will broadcast each element
1171    /// to different members.
1172    ///
1173    /// # Example
1174    /// ```rust
1175    /// # #[cfg(feature = "deploy")] {
1176    /// # use hydro_lang::prelude::*;
1177    /// # use hydro_lang::location::MemberId;
1178    /// # use futures::StreamExt;
1179    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
1180    /// # type Source = ();
1181    /// # type Destination = ();
1182    /// let source: Cluster<Source> = flow.cluster::<Source>();
1183    /// let numbers: Stream<_, Cluster<Source>, _> = source.source_iter(q!(vec![123]));
1184    /// let destination: Cluster<Destination> = flow.cluster::<Destination>();
1185    /// let on_destination: KeyedStream<MemberId<Source>, _, Cluster<Destination>, _> = numbers.broadcast_bincode(&destination, nondet!(/** assuming stable membership */));
1186    /// # on_destination.entries().send_bincode(&p2).entries()
1187    /// // if there are 4 members in the desination, each receives one element from each source member
1188    /// // - Destination(0): { Source(0): [123], Source(1): [123], ... }
1189    /// // - Destination(1): { Source(0): [123], Source(1): [123], ... }
1190    /// // - ...
1191    /// # }, |mut stream| async move {
1192    /// # let mut results = Vec::new();
1193    /// # for w in 0..16 {
1194    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
1195    /// # }
1196    /// # results.sort();
1197    /// # assert_eq!(results, vec![
1198    /// #   "(MemberId::<()>(0), (MemberId::<()>(0), 123))", "(MemberId::<()>(0), (MemberId::<()>(1), 123))", "(MemberId::<()>(0), (MemberId::<()>(2), 123))", "(MemberId::<()>(0), (MemberId::<()>(3), 123))",
1199    /// #   "(MemberId::<()>(1), (MemberId::<()>(0), 123))", "(MemberId::<()>(1), (MemberId::<()>(1), 123))", "(MemberId::<()>(1), (MemberId::<()>(2), 123))", "(MemberId::<()>(1), (MemberId::<()>(3), 123))",
1200    /// #   "(MemberId::<()>(2), (MemberId::<()>(0), 123))", "(MemberId::<()>(2), (MemberId::<()>(1), 123))", "(MemberId::<()>(2), (MemberId::<()>(2), 123))", "(MemberId::<()>(2), (MemberId::<()>(3), 123))",
1201    /// #   "(MemberId::<()>(3), (MemberId::<()>(0), 123))", "(MemberId::<()>(3), (MemberId::<()>(1), 123))", "(MemberId::<()>(3), (MemberId::<()>(2), 123))", "(MemberId::<()>(3), (MemberId::<()>(3), 123))"
1202    /// # ]);
1203    /// # }));
1204    /// # }
1205    /// ```
1206    pub fn broadcast_bincode<L2: 'a>(
1207        self,
1208        other: &Cluster<'a, L2>,
1209        nondet_membership: NonDet,
1210    ) -> KeyedStream<MemberId<L>, T, Cluster<'a, L2>, Unbounded, O, R>
1211    where
1212        T: Clone + Serialize + DeserializeOwned,
1213    {
1214        self.broadcast(other, TCP.fail_stop().bincode(), nondet_membership)
1215    }
1216
1217    /// Broadcasts elements of this stream at each source member to all members of a destination
1218    /// cluster, using the configuration in `via` to set up the message transport.
1219    ///
1220    /// Each source member sends each of its stream elements to **every** member of the cluster
1221    /// based on its latest membership information. Unlike [`Stream::demux`], which requires
1222    /// `(MemberId, T)` tuples to target specific members, `broadcast` takes a stream of
1223    /// **only data elements** and sends each element to all cluster members.
1224    ///
1225    /// # Non-Determinism
1226    /// The set of cluster members may asynchronously change over time. Each element is only broadcast
1227    /// to the current cluster members known _at that point in time_ at the source member. Depending
1228    /// on when each source member is notified of membership changes, it will broadcast each element
1229    /// to different members.
1230    ///
1231    /// # Example
1232    /// ```rust
1233    /// # #[cfg(feature = "deploy")] {
1234    /// # use hydro_lang::prelude::*;
1235    /// # use hydro_lang::location::MemberId;
1236    /// # use futures::StreamExt;
1237    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
1238    /// # type Source = ();
1239    /// # type Destination = ();
1240    /// let source: Cluster<Source> = flow.cluster::<Source>();
1241    /// let numbers: Stream<_, Cluster<Source>, _> = source.source_iter(q!(vec![123]));
1242    /// let destination: Cluster<Destination> = flow.cluster::<Destination>();
1243    /// let on_destination: KeyedStream<MemberId<Source>, _, Cluster<Destination>, _> = numbers.broadcast(&destination, TCP.fail_stop().bincode(), nondet!(/** assuming stable membership */));
1244    /// # on_destination.entries().send(&p2, TCP.fail_stop().bincode()).entries()
1245    /// // if there are 4 members in the desination, each receives one element from each source member
1246    /// // - Destination(0): { Source(0): [123], Source(1): [123], ... }
1247    /// // - Destination(1): { Source(0): [123], Source(1): [123], ... }
1248    /// // - ...
1249    /// # }, |mut stream| async move {
1250    /// # let mut results = Vec::new();
1251    /// # for w in 0..16 {
1252    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
1253    /// # }
1254    /// # results.sort();
1255    /// # assert_eq!(results, vec![
1256    /// #   "(MemberId::<()>(0), (MemberId::<()>(0), 123))", "(MemberId::<()>(0), (MemberId::<()>(1), 123))", "(MemberId::<()>(0), (MemberId::<()>(2), 123))", "(MemberId::<()>(0), (MemberId::<()>(3), 123))",
1257    /// #   "(MemberId::<()>(1), (MemberId::<()>(0), 123))", "(MemberId::<()>(1), (MemberId::<()>(1), 123))", "(MemberId::<()>(1), (MemberId::<()>(2), 123))", "(MemberId::<()>(1), (MemberId::<()>(3), 123))",
1258    /// #   "(MemberId::<()>(2), (MemberId::<()>(0), 123))", "(MemberId::<()>(2), (MemberId::<()>(1), 123))", "(MemberId::<()>(2), (MemberId::<()>(2), 123))", "(MemberId::<()>(2), (MemberId::<()>(3), 123))",
1259    /// #   "(MemberId::<()>(3), (MemberId::<()>(0), 123))", "(MemberId::<()>(3), (MemberId::<()>(1), 123))", "(MemberId::<()>(3), (MemberId::<()>(2), 123))", "(MemberId::<()>(3), (MemberId::<()>(3), 123))"
1260    /// # ]);
1261    /// # }));
1262    /// # }
1263    /// ```
1264    pub fn broadcast<L2: 'a, N: NetworkFor<T>>(
1265        self,
1266        to: &Cluster<'a, L2>,
1267        via: N,
1268        nondet_membership: NonDet,
1269    ) -> KeyedStream<
1270        MemberId<L>,
1271        T,
1272        Cluster<'a, L2>,
1273        Unbounded,
1274        <O as MinOrder<N::OrderingGuarantee>>::Min,
1275        R,
1276    >
1277    where
1278        T: Clone,
1279        O: MinOrder<N::OrderingGuarantee>,
1280    {
1281        // TODO(#1875): the membership snapshot below is over a `KeyedSingleton`, and keyed
1282        // sim hooks do not exist yet. Once they do, expose a composite hook payload here
1283        // (`NonDet<(Option<KeyedSnapshotHook<..>>, Option<BatchHook<T, O, R>>)>`) so tests
1284        // can script the membership snapshot and the element batching independently.
1285        let ids = track_membership(self.location.source_cluster_membership_stream(
1286            to,
1287            nondet!(/** dropped prefixes don't affect broadcast */),
1288        ));
1289        sliced! {
1290            let members_snapshot = use::snapshot(ids, nondet!(
1291                /// membership timing is captured by the caller's guard
1292                nondet_membership
1293            ));
1294            let elements = use::batch(self, nondet!(
1295                /// batching timing is captured by the caller's guard
1296                nondet_membership
1297            ));
1298
1299            let current_members = members_snapshot.filter(q!(|b| *b));
1300            elements.repeat_with_keys(current_members)
1301        }
1302        .demux(to, via)
1303    }
1304
1305    /// Broadcasts elements of this stream at each source member to all members of a destination
1306    /// cluster, assuming membership is closed (fixed at deploy time).
1307    ///
1308    /// Unlike [`Stream::broadcast`], this does not require a [`NonDet`] guard.
1309    /// The membership set is obtained from deploy metadata via [`ClusterIds`], making the
1310    /// broadcast fully deterministic.
1311    ///
1312    /// The consistency guarantee of the output depends on the network's failure policy
1313    /// ([`NetworkFor::ConsistencyGuarantee`]). Policies like `fail_stop` and
1314    /// `lossy_delayed_forever` guarantee that every live destination member eventually
1315    /// materializes the same elements from each source, so the output is
1316    /// [`EventualConsistency`](crate::location::cluster::EventualConsistency). A plain `lossy`
1317    /// policy can drop individual messages for some
1318    /// members while delivering them to others, so replicas may permanently diverge and the
1319    /// output only has [`NoConsistency`].
1320    ///
1321    /// This is only available in deployment targets with static cluster membership
1322    /// (legacy Hydro Deploy and simulation). On dynamic targets, use [`Stream::broadcast`].
1323    pub fn broadcast_closed<L2: 'a, N: NetworkFor<T>>(
1324        self,
1325        to: &Cluster<'a, L2>,
1326        via: N,
1327    ) -> KeyedStream<
1328        MemberId<L>,
1329        T,
1330        Cluster<'a, L2, N::ConsistencyGuarantee>,
1331        Unbounded,
1332        <O as MinOrder<N::OrderingGuarantee>>::Min,
1333        R,
1334    >
1335    where
1336        T: Clone,
1337        O: MinOrder<N::OrderingGuarantee>,
1338    {
1339        let cluster_ids = ClusterIds {
1340            key: to.key,
1341            _phantom: PhantomData,
1342        };
1343        let member_ids = self
1344            .location
1345            .source_iter(q!(cluster_ids
1346                .iter()
1347                .map(|id| MemberId::from_tagless(id.clone()))))
1348            .assert_has_consistency_of_trusted::<Cluster<'a, L, C>>(manual_proof!(
1349                /// ClusterIds is deploy-time metadata, identical on every cluster member.
1350            ));
1351
1352        self.cross_product(member_ids)
1353            .map(q!(|(data, member_id)| (member_id, data)))
1354            .into_keyed()
1355            .demux(to, via)
1356            .assert_has_consistency_of_trusted(manual_proof!(
1357                /// Closed broadcast with fixed membership: every source member sends to every
1358                /// destination member, and the network's failure policy (tracked by
1359                /// `NetworkFor::ConsistencyGuarantee`) delivers the same messages to every live
1360                /// member, so all destinations materialize the same elements.
1361            ))
1362    }
1363
1364    #[cfg(feature = "sim")]
1365    fn register_serialized_external_port<L2>(
1366        self,
1367        other: &External<'_, L2>,
1368        serialize_pipeline: syn::Expr,
1369    ) -> ExternalPortId {
1370        let mut flow_state_borrow = self.location.flow_state().borrow_mut();
1371
1372        let external_port_id = flow_state_borrow.next_external_port();
1373
1374        flow_state_borrow.push_root(HydroRoot::SendExternal {
1375            to_external_key: other.key,
1376            to_port_id: external_port_id,
1377            to_many: false,
1378            unpaired: true,
1379            serialize_fn: Some(serialize_pipeline.into()),
1380            instantiate_fn: DebugInstantiate::Building,
1381            input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1382            op_metadata: HydroIrOpMetadata::new(),
1383        });
1384
1385        external_port_id
1386    }
1387
1388    #[cfg(feature = "sim")]
1389    /// Sets up a bincode-encoded simulation output port for this cluster stream, allowing test
1390    /// code to receive `(member_id, T)` pairs during simulation. Use
1391    /// [`Stream::sim_cluster_output_with`] to select another codec.
1392    pub fn sim_cluster_output(self) -> crate::sim::SimClusterReceiver<T, O, R>
1393    where
1394        T: Serialize + DeserializeOwned,
1395    {
1396        self.sim_cluster_output_with::<crate::sim::codec::BincodeCodec>()
1397    }
1398
1399    #[cfg(feature = "sim")]
1400    /// Sets up a simulation output port for this cluster stream using the codec `Codec`,
1401    /// allowing test code to receive `(member_id, T)` pairs during simulation:
1402    /// `stream.sim_cluster_output_with::<MyCodec>()`. Custom codecs implement
1403    /// [`SimCodec`](crate::sim::codec::SimCodec), which documents where they must be defined.
1404    pub fn sim_cluster_output_with<Codec>(self) -> crate::sim::SimClusterReceiver<T, O, R>
1405    where
1406        Codec: crate::sim::codec::SimCodec<T>,
1407    {
1408        let external_location: External<'a, ()> = External {
1409            key: LocationKey::FIRST,
1410            flow_state: self.location.flow_state().clone(),
1411            _phantom: PhantomData,
1412        };
1413
1414        let external_port_id = self.register_serialized_external_port(
1415            &external_location,
1416            crate::sim::codec::staged_serialize::<T, Codec>(),
1417        );
1418
1419        crate::sim::SimClusterReceiver(external_port_id, PhantomData, Codec::decode)
1420    }
1421}
1422
1423impl<'a, T, L, L2, B: Boundedness, C: Consistency, O: Ordering, R: Retries>
1424    Stream<(MemberId<L2>, T), Cluster<'a, L, C>, B, O, R>
1425{
1426    #[deprecated = "use Stream::demux(..., TCP.fail_stop().bincode()) instead"]
1427    /// Sends elements of this stream at each source member to specific members of a destination
1428    /// cluster, identified by a [`MemberId`], using [`bincode`] to serialize/deserialize messages.
1429    ///
1430    /// Each element in the stream must be a tuple `(MemberId<L2>, T)` where the first element
1431    /// specifies which cluster member should receive the data. Unlike [`Stream::broadcast_bincode`],
1432    /// this API allows precise targeting of specific cluster members rather than broadcasting to
1433    /// all members.
1434    ///
1435    /// Each cluster member sends its local stream elements, and they are collected at each
1436    /// destination member as a [`KeyedStream`] where keys identify the source cluster member.
1437    ///
1438    /// # Example
1439    /// ```rust
1440    /// # #[cfg(feature = "deploy")] {
1441    /// # use hydro_lang::prelude::*;
1442    /// # use futures::StreamExt;
1443    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
1444    /// # type Source = ();
1445    /// # type Destination = ();
1446    /// let source: Cluster<Source> = flow.cluster::<Source>();
1447    /// let to_send: Stream<_, Cluster<_>, _> = source
1448    ///     .source_iter(q!(vec![0, 1, 2, 3]))
1449    ///     .map(q!(|x| (hydro_lang::location::MemberId::from_raw_id(x), x)));
1450    /// let destination: Cluster<Destination> = flow.cluster::<Destination>();
1451    /// let all_received = to_send.demux_bincode(&destination); // KeyedStream<MemberId<Source>, i32, ...>
1452    /// # all_received.entries().send_bincode(&p2).entries()
1453    /// # }, |mut stream| async move {
1454    /// // if there are 4 members in the destination cluster, each receives one message from each source member
1455    /// // - Destination(0): { Source(0): [0], Source(1): [0], ... }
1456    /// // - Destination(1): { Source(0): [1], Source(1): [1], ... }
1457    /// // - ...
1458    /// # let mut results = Vec::new();
1459    /// # for w in 0..16 {
1460    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
1461    /// # }
1462    /// # results.sort();
1463    /// # assert_eq!(results, vec![
1464    /// #   "(MemberId::<()>(0), (MemberId::<()>(0), 0))", "(MemberId::<()>(0), (MemberId::<()>(1), 0))", "(MemberId::<()>(0), (MemberId::<()>(2), 0))", "(MemberId::<()>(0), (MemberId::<()>(3), 0))",
1465    /// #   "(MemberId::<()>(1), (MemberId::<()>(0), 1))", "(MemberId::<()>(1), (MemberId::<()>(1), 1))", "(MemberId::<()>(1), (MemberId::<()>(2), 1))", "(MemberId::<()>(1), (MemberId::<()>(3), 1))",
1466    /// #   "(MemberId::<()>(2), (MemberId::<()>(0), 2))", "(MemberId::<()>(2), (MemberId::<()>(1), 2))", "(MemberId::<()>(2), (MemberId::<()>(2), 2))", "(MemberId::<()>(2), (MemberId::<()>(3), 2))",
1467    /// #   "(MemberId::<()>(3), (MemberId::<()>(0), 3))", "(MemberId::<()>(3), (MemberId::<()>(1), 3))", "(MemberId::<()>(3), (MemberId::<()>(2), 3))", "(MemberId::<()>(3), (MemberId::<()>(3), 3))"
1468    /// # ]);
1469    /// # }));
1470    /// # }
1471    /// ```
1472    pub fn demux_bincode(
1473        self,
1474        other: &Cluster<'a, L2>,
1475    ) -> KeyedStream<MemberId<L>, T, Cluster<'a, L2>, Unbounded, O, R>
1476    where
1477        T: Serialize + DeserializeOwned,
1478    {
1479        self.demux(other, TCP.fail_stop().bincode())
1480    }
1481
1482    /// Sends elements of this stream at each source member to specific members of a destination
1483    /// cluster, identified by a [`MemberId`], using the configuration in `via` to set up the
1484    /// message transport.
1485    ///
1486    /// Each element in the stream must be a tuple `(MemberId<L2>, T)` where the first element
1487    /// specifies which cluster member should receive the data. Unlike [`Stream::broadcast`],
1488    /// this API allows precise targeting of specific cluster members rather than broadcasting to
1489    /// all members.
1490    ///
1491    /// Each cluster member sends its local stream elements, and they are collected at each
1492    /// destination member as a [`KeyedStream`] where keys identify the source cluster member.
1493    ///
1494    /// # Example
1495    /// ```rust
1496    /// # #[cfg(feature = "deploy")] {
1497    /// # use hydro_lang::prelude::*;
1498    /// # use futures::StreamExt;
1499    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
1500    /// # type Source = ();
1501    /// # type Destination = ();
1502    /// let source: Cluster<Source> = flow.cluster::<Source>();
1503    /// let to_send: Stream<_, Cluster<_>, _> = source
1504    ///     .source_iter(q!(vec![0, 1, 2, 3]))
1505    ///     .map(q!(|x| (hydro_lang::location::MemberId::from_raw_id(x), x)));
1506    /// let destination: Cluster<Destination> = flow.cluster::<Destination>();
1507    /// let all_received = to_send.demux(&destination, TCP.fail_stop().bincode()); // KeyedStream<MemberId<Source>, i32, ...>
1508    /// # all_received.entries().send(&p2, TCP.fail_stop().bincode()).entries()
1509    /// # }, |mut stream| async move {
1510    /// // if there are 4 members in the destination cluster, each receives one message from each source member
1511    /// // - Destination(0): { Source(0): [0], Source(1): [0], ... }
1512    /// // - Destination(1): { Source(0): [1], Source(1): [1], ... }
1513    /// // - ...
1514    /// # let mut results = Vec::new();
1515    /// # for w in 0..16 {
1516    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
1517    /// # }
1518    /// # results.sort();
1519    /// # assert_eq!(results, vec![
1520    /// #   "(MemberId::<()>(0), (MemberId::<()>(0), 0))", "(MemberId::<()>(0), (MemberId::<()>(1), 0))", "(MemberId::<()>(0), (MemberId::<()>(2), 0))", "(MemberId::<()>(0), (MemberId::<()>(3), 0))",
1521    /// #   "(MemberId::<()>(1), (MemberId::<()>(0), 1))", "(MemberId::<()>(1), (MemberId::<()>(1), 1))", "(MemberId::<()>(1), (MemberId::<()>(2), 1))", "(MemberId::<()>(1), (MemberId::<()>(3), 1))",
1522    /// #   "(MemberId::<()>(2), (MemberId::<()>(0), 2))", "(MemberId::<()>(2), (MemberId::<()>(1), 2))", "(MemberId::<()>(2), (MemberId::<()>(2), 2))", "(MemberId::<()>(2), (MemberId::<()>(3), 2))",
1523    /// #   "(MemberId::<()>(3), (MemberId::<()>(0), 3))", "(MemberId::<()>(3), (MemberId::<()>(1), 3))", "(MemberId::<()>(3), (MemberId::<()>(2), 3))", "(MemberId::<()>(3), (MemberId::<()>(3), 3))"
1524    /// # ]);
1525    /// # }));
1526    /// # }
1527    /// ```
1528    pub fn demux<N: NetworkFor<T>>(
1529        self,
1530        to: &Cluster<'a, L2>,
1531        via: N,
1532    ) -> KeyedStream<
1533        MemberId<L>,
1534        T,
1535        Cluster<'a, L2, NoConsistency>,
1536        Unbounded,
1537        <O as MinOrder<N::OrderingGuarantee>>::Min,
1538        R,
1539    >
1540    where
1541        O: MinOrder<N::OrderingGuarantee>,
1542    {
1543        self.into_keyed().demux(to, via)
1544    }
1545}
1546
1547#[cfg(test)]
1548mod tests {
1549    #[cfg(feature = "sim")]
1550    use stageleft::q;
1551
1552    #[cfg(feature = "sim")]
1553    use crate::live_collections::sliced::sliced;
1554    #[cfg(feature = "sim")]
1555    use crate::location::{Location, MemberId};
1556    #[cfg(feature = "sim")]
1557    use crate::networking::TCP;
1558    #[cfg(feature = "sim")]
1559    use crate::nondet::nondet;
1560    #[cfg(feature = "sim")]
1561    use crate::prelude::FlowBuilder;
1562
1563    #[cfg(feature = "sim")]
1564    #[test]
1565    fn sim_send_bincode_o2o() {
1566        use crate::networking::TCP;
1567
1568        let mut flow = FlowBuilder::new();
1569        let node = flow.process::<()>();
1570        let node2 = flow.process::<()>();
1571
1572        let (in_send, input) = node.sim_input();
1573
1574        let out_recv = input
1575            .send(&node2, TCP.fail_stop().bincode())
1576            .batch(&node2.tick(), nondet!(/** test */))
1577            .count()
1578            .all_ticks()
1579            .sim_output();
1580
1581        let instances = flow.sim().exhaustive(async || {
1582            in_send.send(());
1583            in_send.send(());
1584            in_send.send(());
1585
1586            let received = out_recv.collect::<Vec<_>>().await;
1587            assert!(received.into_iter().sum::<usize>() == 3);
1588        });
1589
1590        assert_eq!(instances, 4); // 2^{3 - 1}
1591    }
1592
1593    #[cfg(feature = "sim")]
1594    #[test]
1595    fn sim_send_bincode_m2o() {
1596        let mut flow = FlowBuilder::new();
1597        let cluster = flow.cluster::<()>();
1598        let node = flow.process::<()>();
1599
1600        let input = cluster.source_iter(q!(vec![1]));
1601
1602        let out_recv = input
1603            .send(&node, TCP.fail_stop().bincode())
1604            .entries()
1605            .batch(&node.tick(), nondet!(/** test */))
1606            .all_ticks()
1607            .sim_output();
1608
1609        let instances = flow
1610            .sim()
1611            .with_cluster_size(&cluster, 4)
1612            .exhaustive(async || {
1613                out_recv
1614                    .assert_yields_only_unordered(vec![
1615                        (MemberId::from_raw_id(0), 1),
1616                        (MemberId::from_raw_id(1), 1),
1617                        (MemberId::from_raw_id(2), 1),
1618                        (MemberId::from_raw_id(3), 1),
1619                    ])
1620                    .await
1621            });
1622
1623        assert_eq!(instances, 75); // ∑ (k=1 to 4) S(4,k) × k! = 75
1624    }
1625
1626    #[cfg(feature = "sim")]
1627    #[test]
1628    fn sim_send_bincode_multiple_m2o() {
1629        let mut flow = FlowBuilder::new();
1630        let cluster1 = flow.cluster::<()>();
1631        let cluster2 = flow.cluster::<()>();
1632        let node = flow.process::<()>();
1633
1634        let out_recv_1 = cluster1
1635            .source_iter(q!(vec![1]))
1636            .send(&node, TCP.fail_stop().bincode())
1637            .entries()
1638            .sim_output();
1639
1640        let out_recv_2 = cluster2
1641            .source_iter(q!(vec![2]))
1642            .send(&node, TCP.fail_stop().bincode())
1643            .entries()
1644            .sim_output();
1645
1646        let instances = flow
1647            .sim()
1648            .with_cluster_size(&cluster1, 3)
1649            .with_cluster_size(&cluster2, 4)
1650            .exhaustive(async || {
1651                out_recv_1
1652                    .assert_yields_only_unordered(vec![
1653                        (MemberId::from_raw_id(0), 1),
1654                        (MemberId::from_raw_id(1), 1),
1655                        (MemberId::from_raw_id(2), 1),
1656                    ])
1657                    .await;
1658
1659                out_recv_2
1660                    .assert_yields_only_unordered(vec![
1661                        (MemberId::from_raw_id(0), 2),
1662                        (MemberId::from_raw_id(1), 2),
1663                        (MemberId::from_raw_id(2), 2),
1664                        (MemberId::from_raw_id(3), 2),
1665                    ])
1666                    .await;
1667            });
1668
1669        assert_eq!(instances, 1);
1670    }
1671
1672    #[cfg(feature = "sim")]
1673    #[test]
1674    fn sim_send_bincode_o2m() {
1675        let mut flow = FlowBuilder::new();
1676        let cluster = flow.cluster::<()>();
1677        let node = flow.process::<()>();
1678
1679        let input = node.source_iter(q!(vec![
1680            (MemberId::from_raw_id(0), 123),
1681            (MemberId::from_raw_id(1), 456),
1682        ]));
1683
1684        let out_recv = input
1685            .demux(&cluster, TCP.fail_stop().bincode())
1686            .map(q!(|x| x + 1))
1687            .send(&node, TCP.fail_stop().bincode())
1688            .entries()
1689            .sim_output();
1690
1691        flow.sim()
1692            .with_cluster_size(&cluster, 4)
1693            .exhaustive(async || {
1694                out_recv
1695                    .assert_yields_only_unordered(vec![
1696                        (MemberId::from_raw_id(0), 124),
1697                        (MemberId::from_raw_id(1), 457),
1698                    ])
1699                    .await
1700            });
1701    }
1702
1703    #[cfg(feature = "sim")]
1704    #[test]
1705    fn sim_broadcast_bincode_o2m() {
1706        let mut flow = FlowBuilder::new();
1707        let cluster = flow.cluster::<()>();
1708        let node = flow.process::<()>();
1709
1710        let input = node.source_iter(q!(vec![123, 456]));
1711
1712        let out_recv = input
1713            .broadcast(&cluster, TCP.fail_stop().bincode(), nondet!(/** test */))
1714            .map(q!(|x| x + 1))
1715            .send(&node, TCP.fail_stop().bincode())
1716            .entries()
1717            .sim_output();
1718
1719        let mut c_1_produced = false;
1720        let mut c_2_produced = false;
1721        let mut c_1_saw_457_but_not_124 = false;
1722
1723        flow.sim()
1724            .with_cluster_size(&cluster, 2)
1725            .exhaustive(async || {
1726                let all_out = out_recv.collect_sorted::<Vec<_>>().await;
1727
1728                // check that order is preserved
1729                if all_out.contains(&(MemberId::from_raw_id(0), 124)) {
1730                    assert!(all_out.contains(&(MemberId::from_raw_id(0), 457)));
1731                    c_1_produced = true;
1732                }
1733
1734                if all_out.contains(&(MemberId::from_raw_id(1), 124)) {
1735                    assert!(all_out.contains(&(MemberId::from_raw_id(1), 457)));
1736                    c_2_produced = true;
1737                }
1738
1739                if all_out.contains(&(MemberId::from_raw_id(0), 457))
1740                    && !all_out.contains(&(MemberId::from_raw_id(0), 124))
1741                {
1742                    c_1_saw_457_but_not_124 = true;
1743                }
1744            });
1745
1746        assert!(c_1_produced && c_2_produced); // in at least one execution each, the cluster member received both messages
1747
1748        // in at least one execution, the cluster member received 457 but not 124, this tests
1749        // that the simulator properly explores dynamic membership additions (a member that joins after 123 is broadcast)
1750        assert!(c_1_saw_457_but_not_124);
1751    }
1752
1753    #[cfg(feature = "sim")]
1754    #[test]
1755    fn sim_send_bincode_m2m() {
1756        let mut flow = FlowBuilder::new();
1757        let cluster = flow.cluster::<()>();
1758        let node = flow.process::<()>();
1759
1760        let input = node.source_iter(q!(vec![
1761            (MemberId::from_raw_id(0), 123),
1762            (MemberId::from_raw_id(1), 456),
1763        ]));
1764
1765        let out_recv = input
1766            .demux(&cluster, TCP.fail_stop().bincode())
1767            .map(q!(|x| x + 1))
1768            .flat_map_ordered(q!(|x| vec![
1769                (MemberId::from_raw_id(0), x),
1770                (MemberId::from_raw_id(1), x),
1771            ]))
1772            .demux(&cluster, TCP.fail_stop().bincode())
1773            .entries()
1774            .send(&node, TCP.fail_stop().bincode())
1775            .entries()
1776            .sim_output();
1777
1778        flow.sim()
1779            .with_cluster_size(&cluster, 4)
1780            .exhaustive(async || {
1781                out_recv
1782                    .assert_yields_only_unordered(vec![
1783                        (MemberId::from_raw_id(0), (MemberId::from_raw_id(0), 124)),
1784                        (MemberId::from_raw_id(0), (MemberId::from_raw_id(1), 457)),
1785                        (MemberId::from_raw_id(1), (MemberId::from_raw_id(0), 124)),
1786                        (MemberId::from_raw_id(1), (MemberId::from_raw_id(1), 457)),
1787                    ])
1788                    .await
1789            });
1790    }
1791
1792    #[cfg(feature = "sim")]
1793    #[test]
1794    fn sim_lossy_delayed_forever_o2o() {
1795        use std::collections::HashSet;
1796
1797        use crate::properties::manual_proof;
1798
1799        let mut flow = FlowBuilder::new();
1800        let node = flow.process::<()>();
1801        let node2 = flow.process::<()>();
1802
1803        let received = node
1804            .source_iter(q!(0..3_u32))
1805            .send(&node2, TCP.lossy_delayed_forever().bincode())
1806            .fold(
1807                q!(|| std::collections::HashSet::<u32>::new()),
1808                q!(
1809                    |set, v| {
1810                        set.insert(v);
1811                    },
1812                    commutative = manual_proof!(/** set insert is commutative */)
1813                ),
1814            );
1815
1816        let out_recv = sliced! {
1817            let snapshot = use::snapshot(received, nondet!(/** test */));
1818            snapshot.into_stream()
1819        }
1820        .sim_output();
1821
1822        let mut saw_non_contiguous = false;
1823
1824        flow.sim().test_safety_only().exhaustive(async || {
1825            let snapshots = out_recv.collect::<Vec<HashSet<u32>>>().await;
1826
1827            // Check each individual snapshot for a non-contiguous subset.
1828            for set in &snapshots {
1829                #[expect(clippy::disallowed_methods, reason = "min / max are deterministic")]
1830                if set.len() >= 2 && set.len() < 3 {
1831                    let min = *set.iter().min().unwrap();
1832                    let max = *set.iter().max().unwrap();
1833                    if set.len() < (max - min + 1) as usize {
1834                        saw_non_contiguous = true;
1835                    }
1836                }
1837            }
1838        });
1839
1840        assert!(
1841            saw_non_contiguous,
1842            "Expected at least one execution with a non-contiguous subset of inputs"
1843        );
1844    }
1845
1846    #[cfg(feature = "sim")]
1847    #[test]
1848    fn sim_udp_lossy_delayed_forever_o2o() {
1849        use std::collections::HashSet;
1850
1851        use crate::networking::UDP;
1852        use crate::properties::manual_proof;
1853
1854        let mut flow = FlowBuilder::new();
1855        let node = flow.process::<()>();
1856        let node2 = flow.process::<()>();
1857
1858        let received = node
1859            .source_iter(q!(0..3_u32))
1860            .send(&node2, UDP.lossy_delayed_forever().bincode())
1861            .fold(
1862                q!(|| std::collections::HashSet::<u32>::new()),
1863                q!(
1864                    |set, v| {
1865                        set.insert(v);
1866                    },
1867                    commutative = manual_proof!(/** set insert is commutative */)
1868                ),
1869            );
1870
1871        let out_recv = sliced! {
1872            let snapshot = use::snapshot(received, nondet!(/** test */));
1873            snapshot.into_stream()
1874        }
1875        .sim_output();
1876
1877        let mut saw_non_contiguous = false;
1878
1879        flow.sim().test_safety_only().exhaustive(async || {
1880            let snapshots = out_recv.collect::<Vec<HashSet<u32>>>().await;
1881
1882            // Check each individual snapshot for a non-contiguous subset.
1883            for set in &snapshots {
1884                #[expect(clippy::disallowed_methods, reason = "min / max are deterministic")]
1885                if set.len() >= 2 && set.len() < 3 {
1886                    let min = *set.iter().min().unwrap();
1887                    let max = *set.iter().max().unwrap();
1888                    if set.len() < (max - min + 1) as usize {
1889                        saw_non_contiguous = true;
1890                    }
1891                }
1892            }
1893        });
1894
1895        assert!(
1896            saw_non_contiguous,
1897            "Expected at least one execution with a non-contiguous subset of inputs"
1898        );
1899    }
1900
1901    #[cfg(feature = "sim")]
1902    #[test]
1903    fn sim_broadcast_closed_o2m() {
1904        let mut flow = FlowBuilder::new();
1905        let cluster = flow.cluster::<()>();
1906        let node = flow.process::<()>();
1907
1908        let input = node.source_iter(q!(vec![123, 456]));
1909
1910        let out_recv = input
1911            .broadcast_closed(&cluster, TCP.fail_stop().bincode())
1912            .send(&node, TCP.fail_stop().bincode())
1913            .entries()
1914            .sim_output();
1915
1916        flow.sim()
1917            .with_cluster_size(&cluster, 2)
1918            .exhaustive(async || {
1919                out_recv
1920                    .assert_yields_only_unordered(vec![
1921                        (MemberId::from_raw_id(0), 123),
1922                        (MemberId::from_raw_id(0), 456),
1923                        (MemberId::from_raw_id(1), 123),
1924                        (MemberId::from_raw_id(1), 456),
1925                    ])
1926                    .await
1927            });
1928    }
1929
1930    #[cfg(feature = "sim")]
1931    #[test]
1932    fn sim_broadcast_closed_m2m() {
1933        let mut flow = FlowBuilder::new();
1934        let source = flow.cluster::<()>();
1935        let dest: crate::location::Cluster<'_, ()> = flow.cluster::<()>();
1936        let node = flow.process::<()>();
1937
1938        let input = source.source_iter(q!(vec![123]));
1939
1940        // Broadcast from source cluster to dest cluster, then collect at a process.
1941        let out_recv = input
1942            .broadcast_closed(&dest, TCP.fail_stop().bincode())
1943            .entries()
1944            .send(&node, TCP.fail_stop().bincode())
1945            .entries()
1946            .sim_output();
1947
1948        flow.sim()
1949            .with_cluster_size(&source, 2)
1950            .with_cluster_size(&dest, 2)
1951            .exhaustive(async || {
1952                // Each source member (0, 1) broadcasts 123 to each dest member (0, 1).
1953                // The dest members then send to the process keyed by dest member id.
1954                // Each dest member receives (source_0, 123) and (source_1, 123).
1955                out_recv
1956                    .assert_yields_only_unordered(vec![
1957                        (MemberId::from_raw_id(0), (MemberId::from_raw_id(0), 123)),
1958                        (MemberId::from_raw_id(0), (MemberId::from_raw_id(1), 123)),
1959                        (MemberId::from_raw_id(1), (MemberId::from_raw_id(0), 123)),
1960                        (MemberId::from_raw_id(1), (MemberId::from_raw_id(1), 123)),
1961                    ])
1962                    .await
1963            });
1964    }
1965
1966    /// Compile-time check that the consistency guarantee of `broadcast_closed` output tracks
1967    /// the network's failure policy: `fail_stop` and `lossy_delayed_forever` preserve
1968    /// [`EventualConsistency`], while plain `lossy` only provides [`NoConsistency`].
1969    #[cfg(feature = "sim")]
1970    #[test]
1971    fn broadcast_closed_consistency_tracks_failure_policy() {
1972        use crate::live_collections::keyed_stream::KeyedStream;
1973        use crate::live_collections::stream::Stream;
1974        use crate::location::Cluster;
1975        use crate::location::cluster::{EventualConsistency, NoConsistency};
1976
1977        let mut flow = FlowBuilder::new();
1978        let cluster = flow.cluster::<()>();
1979        let source = flow.cluster::<()>();
1980        let node = flow.process::<()>();
1981
1982        // `fail_stop` models a failed connection as the recipient having failed, preserving
1983        // eventual consistency across live members.
1984        let _: Stream<u32, Cluster<'_, (), EventualConsistency>, _, _, _> = node
1985            .source_iter(q!(vec![1u32]))
1986            .broadcast_closed(&cluster, TCP.fail_stop().bincode());
1987
1988        // `lossy_delayed_forever` models drops as indefinite delays, preserving eventual
1989        // consistency.
1990        let _: Stream<u32, Cluster<'_, (), EventualConsistency>, _, _, _> = node
1991            .source_iter(q!(vec![1u32]))
1992            .broadcast_closed(&cluster, TCP.lossy_delayed_forever().bincode());
1993
1994        // Plain `lossy` can drop messages for some members while delivering them to others,
1995        // so replicas may permanently diverge.
1996        let _: Stream<u32, Cluster<'_, (), NoConsistency>, _, _, _> = node
1997            .source_iter(q!(vec![1u32]))
1998            .broadcast_closed(&cluster, TCP.lossy(nondet!(/** test */)).bincode());
1999
2000        // The same applies to cluster-to-cluster closed broadcasts.
2001        let _: KeyedStream<MemberId<()>, u32, Cluster<'_, (), EventualConsistency>, _, _, _> =
2002            source
2003                .source_iter(q!(vec![1u32]))
2004                .broadcast_closed(&cluster, TCP.fail_stop().bincode());
2005
2006        let _: KeyedStream<MemberId<()>, u32, Cluster<'_, (), NoConsistency>, _, _, _> = source
2007            .source_iter(q!(vec![1u32]))
2008            .broadcast_closed(&cluster, TCP.lossy(nondet!(/** test */)).bincode());
2009
2010        let _ = flow.finalize();
2011    }
2012}