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