Skip to main content

hydro_lang/live_collections/keyed_stream/
networking.rs

1//! Networking APIs for [`KeyedStream`].
2
3use serde::Serialize;
4use serde::de::DeserializeOwned;
5use stageleft::{q, quote_type};
6
7use super::KeyedStream;
8use crate::compile::ir::{DebugInstantiate, HydroNode, NetworkRecv, NetworkSend};
9use crate::live_collections::boundedness::{Boundedness, Unbounded};
10use crate::live_collections::stream::{MinOrder, Ordering, Retries, Stream};
11use crate::location::cluster::{Consistency, NoConsistency};
12#[cfg(stageleft_runtime)]
13use crate::location::dynamic::DynLocation;
14use crate::location::{Cluster, MemberId, Process};
15use crate::networking::{NetworkFor, TCP};
16
17impl<'a, T, L, L2, B: Boundedness, O: Ordering, R: Retries>
18    KeyedStream<MemberId<L2>, T, Process<'a, L>, B, O, R>
19{
20    #[deprecated = "use KeyedStream::demux(..., TCP.fail_stop().bincode()) instead"]
21    /// Sends each group of this stream to a specific member of a cluster, with the [`MemberId`] key
22    /// identifying the recipient for each group and using [`bincode`] to serialize/deserialize messages.
23    ///
24    /// Each key must be a `MemberId<L2>` and each value must be a `T` where the key specifies
25    /// which cluster member should receive the data. Unlike [`Stream::broadcast_bincode`], this
26    /// API allows precise targeting of specific cluster members rather than broadcasting to
27    /// all members.
28    ///
29    /// # Example
30    /// ```rust
31    /// # #[cfg(feature = "deploy")] {
32    /// # use hydro_lang::prelude::*;
33    /// # use futures::StreamExt;
34    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
35    /// let p1 = flow.process::<()>();
36    /// let workers: Cluster<()> = flow.cluster::<()>();
37    /// let numbers: Stream<_, Process<_>, _> = p1.source_iter(q!(vec![0, 1, 2, 3]));
38    /// let on_worker: Stream<_, Cluster<_>, _> = numbers
39    ///     .map(q!(|x| (hydro_lang::location::MemberId::from_raw_id(x), x)))
40    ///     .into_keyed()
41    ///     .demux_bincode(&workers);
42    /// # on_worker.send_bincode(&p2).entries()
43    /// // if there are 4 members in the cluster, each receives one element
44    /// // - MemberId::<()>(0): [0]
45    /// // - MemberId::<()>(1): [1]
46    /// // - MemberId::<()>(2): [2]
47    /// // - MemberId::<()>(3): [3]
48    /// # }, |mut stream| async move {
49    /// # let mut results = Vec::new();
50    /// # for w in 0..4 {
51    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
52    /// # }
53    /// # results.sort();
54    /// # assert_eq!(results, vec!["(MemberId::<()>(0), 0)", "(MemberId::<()>(1), 1)", "(MemberId::<()>(2), 2)", "(MemberId::<()>(3), 3)"]);
55    /// # }));
56    /// # }
57    /// ```
58    pub fn demux_bincode(
59        self,
60        other: &Cluster<'a, L2>,
61    ) -> Stream<T, Cluster<'a, L2>, Unbounded, O, R>
62    where
63        T: Serialize + DeserializeOwned,
64    {
65        self.demux(other, TCP.fail_stop().bincode())
66    }
67
68    /// Sends each group of this stream to a specific member of a cluster, with the [`MemberId`] key
69    /// identifying the recipient for each group and using the configuration in `via` to set up the
70    /// message transport.
71    ///
72    /// Each key must be a `MemberId<L2>` and each value must be a `T` where the key specifies
73    /// which cluster member should receive the data. Unlike [`Stream::broadcast`], this
74    /// API allows precise targeting of specific cluster members rather than broadcasting to
75    /// all members.
76    ///
77    /// # Example
78    /// ```rust
79    /// # #[cfg(feature = "deploy")] {
80    /// # use hydro_lang::prelude::*;
81    /// # use futures::StreamExt;
82    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
83    /// let p1 = flow.process::<()>();
84    /// let workers: Cluster<()> = flow.cluster::<()>();
85    /// let numbers: Stream<_, Process<_>, _> = p1.source_iter(q!(vec![0, 1, 2, 3]));
86    /// let on_worker: Stream<_, Cluster<_>, _> = numbers
87    ///     .map(q!(|x| (hydro_lang::location::MemberId::from_raw_id(x), x)))
88    ///     .into_keyed()
89    ///     .demux(&workers, TCP.fail_stop().bincode());
90    /// # on_worker.send(&p2, TCP.fail_stop().bincode()).entries()
91    /// // if there are 4 members in the cluster, each receives one element
92    /// // - MemberId::<()>(0): [0]
93    /// // - MemberId::<()>(1): [1]
94    /// // - MemberId::<()>(2): [2]
95    /// // - MemberId::<()>(3): [3]
96    /// # }, |mut stream| async move {
97    /// # let mut results = Vec::new();
98    /// # for w in 0..4 {
99    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
100    /// # }
101    /// # results.sort();
102    /// # assert_eq!(results, vec!["(MemberId::<()>(0), 0)", "(MemberId::<()>(1), 1)", "(MemberId::<()>(2), 2)", "(MemberId::<()>(3), 3)"]);
103    /// # }));
104    /// # }
105    /// ```
106    pub fn demux<N: NetworkFor<T>>(
107        self,
108        to: &Cluster<'a, L2>,
109        via: N,
110    ) -> Stream<
111        T,
112        // NoConsistency because there each replica member may receive different streams
113        Cluster<'a, L2, NoConsistency>,
114        Unbounded,
115        <O as MinOrder<N::OrderingGuarantee>>::Min,
116        R,
117    >
118    where
119        O: MinOrder<N::OrderingGuarantee>,
120    {
121        let name = via.name();
122        assert!(
123            !to.multiversioned() || name.is_some(),
124            "Cannot send to a multiversioned location without a channel name. Please provide a name for the network."
125        );
126
127        let (serialize, deserialize) = if N::is_embedded() {
128            (
129                NetworkSend::Embedded {
130                    tag: Some(quote_type::<L2>().into()),
131                    element_type: quote_type::<T>().into(),
132                },
133                NetworkRecv::Embedded {
134                    tag: None,
135                    element_type: quote_type::<T>().into(),
136                },
137            )
138        } else {
139            (
140                NetworkSend::Custom {
141                    serialize_fn: Some(N::serialize_thunk(true).into()),
142                },
143                NetworkRecv::Custom {
144                    deserialize_fn: Some(N::deserialize_thunk(None).into()),
145                },
146            )
147        };
148
149        Stream::new(
150            to.clone(),
151            HydroNode::Network {
152                name: name.map(ToOwned::to_owned),
153                networking_info: N::networking_info(),
154                serialize,
155                deserialize,
156                instantiate_fn: DebugInstantiate::Building,
157                input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
158                metadata: to.new_node_metadata(Stream::<
159                    T,
160                    Cluster<'a, L2>,
161                    Unbounded,
162                    <O as MinOrder<N::OrderingGuarantee>>::Min,
163                    R,
164                >::collection_kind()),
165            },
166        )
167    }
168}
169
170impl<'a, K, T, L, L2, B: Boundedness, O: Ordering, R: Retries>
171    KeyedStream<(MemberId<L2>, K), T, Process<'a, L>, B, O, R>
172{
173    #[deprecated = "use KeyedStream::demux(..., TCP.fail_stop().bincode()) instead"]
174    /// Sends each group of this stream to a specific member of a cluster. The input stream has a
175    /// compound key where the first element is the recipient's [`MemberId`] and the second element
176    /// is a key that will be sent along with the value, using [`bincode`] to serialize/deserialize
177    /// messages.
178    ///
179    /// # Example
180    /// ```rust
181    /// # #[cfg(feature = "deploy")] {
182    /// # use hydro_lang::prelude::*;
183    /// # use futures::StreamExt;
184    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
185    /// let p1 = flow.process::<()>();
186    /// let workers: Cluster<()> = flow.cluster::<()>();
187    /// let to_send: KeyedStream<_, _, Process<_>, _> = p1
188    ///     .source_iter(q!(vec![0, 1, 2, 3]))
189    ///     .map(q!(|x| ((hydro_lang::location::MemberId::from_raw_id(x), x), x + 123)))
190    ///     .into_keyed();
191    /// let on_worker: KeyedStream<_, _, Cluster<_>, _> = to_send.demux_bincode(&workers);
192    /// # on_worker.entries().send_bincode(&p2).entries()
193    /// // if there are 4 members in the cluster, each receives one element
194    /// // - MemberId::<()>(0): { 0: [123] }
195    /// // - MemberId::<()>(1): { 1: [124] }
196    /// // - ...
197    /// # }, |mut stream| async move {
198    /// # let mut results = Vec::new();
199    /// # for w in 0..4 {
200    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
201    /// # }
202    /// # results.sort();
203    /// # assert_eq!(results, vec!["(MemberId::<()>(0), (0, 123))", "(MemberId::<()>(1), (1, 124))", "(MemberId::<()>(2), (2, 125))", "(MemberId::<()>(3), (3, 126))"]);
204    /// # }));
205    /// # }
206    /// ```
207    pub fn demux_bincode(
208        self,
209        other: &Cluster<'a, L2>,
210    ) -> KeyedStream<K, T, Cluster<'a, L2>, Unbounded, O, R>
211    where
212        K: Serialize + DeserializeOwned,
213        T: Serialize + DeserializeOwned,
214    {
215        self.demux(other, TCP.fail_stop().bincode())
216    }
217
218    /// Sends each group of this stream to a specific member of a cluster. The input stream has a
219    /// compound key where the first element is the recipient's [`MemberId`] and the second element
220    /// is a key that will be sent along with the value, using the configuration in `via` to set up
221    /// the message transport.
222    ///
223    /// # Example
224    /// ```rust
225    /// # #[cfg(feature = "deploy")] {
226    /// # use hydro_lang::prelude::*;
227    /// # use futures::StreamExt;
228    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
229    /// let p1 = flow.process::<()>();
230    /// let workers: Cluster<()> = flow.cluster::<()>();
231    /// let to_send: KeyedStream<_, _, Process<_>, _> = p1
232    ///     .source_iter(q!(vec![0, 1, 2, 3]))
233    ///     .map(q!(|x| ((hydro_lang::location::MemberId::from_raw_id(x), x), x + 123)))
234    ///     .into_keyed();
235    /// let on_worker: KeyedStream<_, _, Cluster<_>, _> = to_send.demux(&workers, TCP.fail_stop().bincode());
236    /// # on_worker.entries().send(&p2, TCP.fail_stop().bincode()).entries()
237    /// // if there are 4 members in the cluster, each receives one element
238    /// // - MemberId::<()>(0): { 0: [123] }
239    /// // - MemberId::<()>(1): { 1: [124] }
240    /// // - ...
241    /// # }, |mut stream| async move {
242    /// # let mut results = Vec::new();
243    /// # for w in 0..4 {
244    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
245    /// # }
246    /// # results.sort();
247    /// # assert_eq!(results, vec!["(MemberId::<()>(0), (0, 123))", "(MemberId::<()>(1), (1, 124))", "(MemberId::<()>(2), (2, 125))", "(MemberId::<()>(3), (3, 126))"]);
248    /// # }));
249    /// # }
250    /// ```
251    pub fn demux<N: NetworkFor<(K, T)>>(
252        self,
253        to: &Cluster<'a, L2>,
254        via: N,
255    ) -> KeyedStream<
256        K,
257        T,
258        Cluster<'a, L2, NoConsistency>,
259        Unbounded,
260        <O as MinOrder<N::OrderingGuarantee>>::Min,
261        R,
262    >
263    where
264        O: MinOrder<N::OrderingGuarantee>,
265    {
266        let name = via.name();
267        assert!(
268            !to.multiversioned() || name.is_some(),
269            "Cannot send to a multiversioned location without a channel name. Please provide a name for the network."
270        );
271
272        let (serialize, deserialize) = if N::is_embedded() {
273            (
274                NetworkSend::Embedded {
275                    tag: Some(quote_type::<L2>().into()),
276                    element_type: quote_type::<(K, T)>().into(),
277                },
278                NetworkRecv::Embedded {
279                    tag: None,
280                    element_type: quote_type::<(K, T)>().into(),
281                },
282            )
283        } else {
284            (
285                NetworkSend::Custom {
286                    serialize_fn: Some(N::serialize_thunk(true).into()),
287                },
288                NetworkRecv::Custom {
289                    deserialize_fn: Some(N::deserialize_thunk(None).into()),
290                },
291            )
292        };
293
294        KeyedStream::new(
295            to.clone(),
296            HydroNode::Network {
297                name: name.map(ToOwned::to_owned),
298                networking_info: N::networking_info(),
299                serialize,
300                deserialize,
301                instantiate_fn: DebugInstantiate::Building,
302                input: Box::new(
303                    self.entries()
304                        .map(q!(|((id, k), v)| (id, (k, v))))
305                        .ir_node
306                        .replace(HydroNode::Placeholder),
307                ),
308                metadata: to.new_node_metadata(KeyedStream::<
309                    K,
310                    T,
311                    Cluster<'a, L2>,
312                    Unbounded,
313                    <O as MinOrder<N::OrderingGuarantee>>::Min,
314                    R,
315                >::collection_kind()),
316            },
317        )
318    }
319}
320
321impl<'a, T, L, L2, B: Boundedness, C: Consistency, O: Ordering, R: Retries>
322    KeyedStream<MemberId<L2>, T, Cluster<'a, L, C>, B, O, R>
323{
324    #[deprecated = "use KeyedStream::demux(..., TCP.fail_stop().bincode()) instead"]
325    /// Sends each group of this stream at each source member to a specific member of a destination
326    /// cluster, with the [`MemberId`] key identifying the recipient for each group and using
327    /// [`bincode`] to serialize/deserialize messages.
328    ///
329    /// Each key must be a `MemberId<L2>` and each value must be a `T` where the key specifies
330    /// which cluster member should receive the data. Unlike [`Stream::broadcast_bincode`], this
331    /// API allows precise targeting of specific cluster members rather than broadcasting to all
332    /// members.
333    ///
334    /// Each cluster member sends its local stream elements, and they are collected at each
335    /// destination member as a [`KeyedStream`] where keys identify the source cluster member.
336    ///
337    /// # Example
338    /// ```rust
339    /// # #[cfg(feature = "deploy")] {
340    /// # use hydro_lang::prelude::*;
341    /// # use futures::StreamExt;
342    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
343    /// # type Source = ();
344    /// # type Destination = ();
345    /// let source: Cluster<Source> = flow.cluster::<Source>();
346    /// let to_send: KeyedStream<_, _, Cluster<_>, _> = source
347    ///     .source_iter(q!(vec![0, 1, 2, 3]))
348    ///     .map(q!(|x| (hydro_lang::location::MemberId::from_raw_id(x), x)))
349    ///     .into_keyed();
350    /// let destination: Cluster<Destination> = flow.cluster::<Destination>();
351    /// let all_received = to_send.demux_bincode(&destination); // KeyedStream<MemberId<Source>, i32, ...>
352    /// # all_received.entries().send_bincode(&p2).entries()
353    /// # }, |mut stream| async move {
354    /// // if there are 4 members in the destination cluster, each receives one message from each source member
355    /// // - Destination(0): { Source(0): [0], Source(1): [0], ... }
356    /// // - Destination(1): { Source(0): [1], Source(1): [1], ... }
357    /// // - ...
358    /// # let mut results = Vec::new();
359    /// # for w in 0..16 {
360    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
361    /// # }
362    /// # results.sort();
363    /// # assert_eq!(results, vec![
364    /// #   "(MemberId::<()>(0), (MemberId::<()>(0), 0))", "(MemberId::<()>(0), (MemberId::<()>(1), 0))", "(MemberId::<()>(0), (MemberId::<()>(2), 0))", "(MemberId::<()>(0), (MemberId::<()>(3), 0))",
365    /// #   "(MemberId::<()>(1), (MemberId::<()>(0), 1))", "(MemberId::<()>(1), (MemberId::<()>(1), 1))", "(MemberId::<()>(1), (MemberId::<()>(2), 1))", "(MemberId::<()>(1), (MemberId::<()>(3), 1))",
366    /// #   "(MemberId::<()>(2), (MemberId::<()>(0), 2))", "(MemberId::<()>(2), (MemberId::<()>(1), 2))", "(MemberId::<()>(2), (MemberId::<()>(2), 2))", "(MemberId::<()>(2), (MemberId::<()>(3), 2))",
367    /// #   "(MemberId::<()>(3), (MemberId::<()>(0), 3))", "(MemberId::<()>(3), (MemberId::<()>(1), 3))", "(MemberId::<()>(3), (MemberId::<()>(2), 3))", "(MemberId::<()>(3), (MemberId::<()>(3), 3))"
368    /// # ]);
369    /// # }));
370    /// # }
371    /// ```
372    pub fn demux_bincode(
373        self,
374        other: &Cluster<'a, L2>,
375    ) -> KeyedStream<MemberId<L>, T, Cluster<'a, L2>, Unbounded, O, R>
376    where
377        T: Serialize + DeserializeOwned,
378    {
379        self.demux(other, TCP.fail_stop().bincode())
380    }
381
382    /// Sends each group of this stream at each source member to a specific member of a destination
383    /// cluster, with the [`MemberId`] key identifying the recipient for each group and using the
384    /// configuration in `via` to set up the message transport.
385    ///
386    /// Each key must be a `MemberId<L2>` and each value must be a `T` where the key specifies
387    /// which cluster member should receive the data. Unlike [`Stream::broadcast`], this
388    /// API allows precise targeting of specific cluster members rather than broadcasting to all
389    /// members.
390    ///
391    /// Each cluster member sends its local stream elements, and they are collected at each
392    /// destination member as a [`KeyedStream`] where keys identify the source cluster member.
393    ///
394    /// # Example
395    /// ```rust
396    /// # #[cfg(feature = "deploy")] {
397    /// # use hydro_lang::prelude::*;
398    /// # use futures::StreamExt;
399    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
400    /// # type Source = ();
401    /// # type Destination = ();
402    /// let source: Cluster<Source> = flow.cluster::<Source>();
403    /// let to_send: KeyedStream<_, _, Cluster<_>, _> = source
404    ///     .source_iter(q!(vec![0, 1, 2, 3]))
405    ///     .map(q!(|x| (hydro_lang::location::MemberId::from_raw_id(x), x)))
406    ///     .into_keyed();
407    /// let destination: Cluster<Destination> = flow.cluster::<Destination>();
408    /// let all_received = to_send.demux(&destination, TCP.fail_stop().bincode()); // KeyedStream<MemberId<Source>, i32, ...>
409    /// # all_received.entries().send(&p2, TCP.fail_stop().bincode()).entries()
410    /// # }, |mut stream| async move {
411    /// // if there are 4 members in the destination cluster, each receives one message from each source member
412    /// // - Destination(0): { Source(0): [0], Source(1): [0], ... }
413    /// // - Destination(1): { Source(0): [1], Source(1): [1], ... }
414    /// // - ...
415    /// # let mut results = Vec::new();
416    /// # for w in 0..16 {
417    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
418    /// # }
419    /// # results.sort();
420    /// # assert_eq!(results, vec![
421    /// #   "(MemberId::<()>(0), (MemberId::<()>(0), 0))", "(MemberId::<()>(0), (MemberId::<()>(1), 0))", "(MemberId::<()>(0), (MemberId::<()>(2), 0))", "(MemberId::<()>(0), (MemberId::<()>(3), 0))",
422    /// #   "(MemberId::<()>(1), (MemberId::<()>(0), 1))", "(MemberId::<()>(1), (MemberId::<()>(1), 1))", "(MemberId::<()>(1), (MemberId::<()>(2), 1))", "(MemberId::<()>(1), (MemberId::<()>(3), 1))",
423    /// #   "(MemberId::<()>(2), (MemberId::<()>(0), 2))", "(MemberId::<()>(2), (MemberId::<()>(1), 2))", "(MemberId::<()>(2), (MemberId::<()>(2), 2))", "(MemberId::<()>(2), (MemberId::<()>(3), 2))",
424    /// #   "(MemberId::<()>(3), (MemberId::<()>(0), 3))", "(MemberId::<()>(3), (MemberId::<()>(1), 3))", "(MemberId::<()>(3), (MemberId::<()>(2), 3))", "(MemberId::<()>(3), (MemberId::<()>(3), 3))"
425    /// # ]);
426    /// # }));
427    /// # }
428    /// ```
429    pub fn demux<N: NetworkFor<T>>(
430        self,
431        to: &Cluster<'a, L2>,
432        via: N,
433    ) -> KeyedStream<
434        MemberId<L>,
435        T,
436        Cluster<'a, L2, NoConsistency>,
437        Unbounded,
438        <O as MinOrder<N::OrderingGuarantee>>::Min,
439        R,
440    >
441    where
442        O: MinOrder<N::OrderingGuarantee>,
443    {
444        let name = via.name();
445        assert!(
446            !to.multiversioned() || name.is_some(),
447            "Cannot send to a multiversioned location without a channel name. Please provide a name for the network."
448        );
449
450        let (serialize, deserialize) = if N::is_embedded() {
451            (
452                NetworkSend::Embedded {
453                    tag: Some(quote_type::<L2>().into()),
454                    element_type: quote_type::<T>().into(),
455                },
456                NetworkRecv::Embedded {
457                    tag: Some(quote_type::<L>().into()),
458                    element_type: quote_type::<T>().into(),
459                },
460            )
461        } else {
462            (
463                NetworkSend::Custom {
464                    serialize_fn: Some(N::serialize_thunk(true).into()),
465                },
466                NetworkRecv::Custom {
467                    deserialize_fn: Some(N::deserialize_thunk(Some(&quote_type::<L>())).into()),
468                },
469            )
470        };
471
472        KeyedStream::new(
473            to.clone(),
474            HydroNode::Network {
475                name: name.map(ToOwned::to_owned),
476                networking_info: N::networking_info(),
477                serialize,
478                deserialize,
479                instantiate_fn: DebugInstantiate::Building,
480                input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
481                metadata: to.new_node_metadata(KeyedStream::<
482                    MemberId<L>,
483                    T,
484                    Cluster<'a, L2>,
485                    Unbounded,
486                    <O as MinOrder<N::OrderingGuarantee>>::Min,
487                    R,
488                >::collection_kind()),
489            },
490        )
491    }
492}
493
494impl<'a, K, V, L, B: Boundedness, C: Consistency, O: Ordering, R: Retries>
495    KeyedStream<K, V, Cluster<'a, L, C>, B, O, R>
496{
497    #[deprecated = "use KeyedStream::send(..., TCP.fail_stop().bincode()) instead"]
498    /// "Moves" elements of this keyed stream from a cluster to a process by sending them over the
499    /// network, using [`bincode`] to serialize/deserialize messages. The resulting [`KeyedStream`]
500    /// has a compound key where the first element is the sender's [`MemberId`] and the second
501    /// element is the original key.
502    ///
503    /// # Example
504    /// ```rust
505    /// # #[cfg(feature = "deploy")] {
506    /// # use hydro_lang::prelude::*;
507    /// # use futures::StreamExt;
508    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
509    /// # type Source = ();
510    /// # type Destination = ();
511    /// let source: Cluster<Source> = flow.cluster::<Source>();
512    /// let to_send: KeyedStream<_, _, Cluster<_>, _> = source
513    ///     .source_iter(q!(vec![0, 1, 2, 3]))
514    ///     .map(q!(|x| (x, x + 123)))
515    ///     .into_keyed();
516    /// let destination_process = flow.process::<Destination>();
517    /// let all_received = to_send.send_bincode(&destination_process); // KeyedStream<(MemberId<Source>, i32), i32, ...>
518    /// # all_received.entries().send_bincode(&p2)
519    /// # }, |mut stream| async move {
520    /// // if there are 4 members in the source cluster, the destination process receives four messages from each source member
521    /// // {
522    /// //     (MemberId<Source>(0), 0): [123], (MemberId<Source>(1), 0): [123], ...,
523    /// //     (MemberId<Source>(0), 1): [124], (MemberId<Source>(1), 1): [124], ...,
524    /// //     ...
525    /// // }
526    /// # let mut results = Vec::new();
527    /// # for w in 0..16 {
528    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
529    /// # }
530    /// # results.sort();
531    /// # assert_eq!(results, vec![
532    /// #   "((MemberId::<()>(0), 0), 123)",
533    /// #   "((MemberId::<()>(0), 1), 124)",
534    /// #   "((MemberId::<()>(0), 2), 125)",
535    /// #   "((MemberId::<()>(0), 3), 126)",
536    /// #   "((MemberId::<()>(1), 0), 123)",
537    /// #   "((MemberId::<()>(1), 1), 124)",
538    /// #   "((MemberId::<()>(1), 2), 125)",
539    /// #   "((MemberId::<()>(1), 3), 126)",
540    /// #   "((MemberId::<()>(2), 0), 123)",
541    /// #   "((MemberId::<()>(2), 1), 124)",
542    /// #   "((MemberId::<()>(2), 2), 125)",
543    /// #   "((MemberId::<()>(2), 3), 126)",
544    /// #   "((MemberId::<()>(3), 0), 123)",
545    /// #   "((MemberId::<()>(3), 1), 124)",
546    /// #   "((MemberId::<()>(3), 2), 125)",
547    /// #   "((MemberId::<()>(3), 3), 126)",
548    /// # ]);
549    /// # }));
550    /// # }
551    /// ```
552    pub fn send_bincode<L2>(
553        self,
554        other: &Process<'a, L2>,
555    ) -> KeyedStream<(MemberId<L>, K), V, Process<'a, L2>, Unbounded, O, R>
556    where
557        K: Serialize + DeserializeOwned,
558        V: Serialize + DeserializeOwned,
559    {
560        self.send(other, TCP.fail_stop().bincode())
561    }
562
563    /// "Moves" elements of this keyed stream from a cluster to a process by sending them over the
564    /// network, using the configuration in `via` to set up the message transport. The resulting
565    /// [`KeyedStream`] has a compound key where the first element is the sender's [`MemberId`] and
566    /// the second element is the original key.
567    ///
568    /// # Example
569    /// ```rust
570    /// # #[cfg(feature = "deploy")] {
571    /// # use hydro_lang::prelude::*;
572    /// # use futures::StreamExt;
573    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
574    /// # type Source = ();
575    /// # type Destination = ();
576    /// let source: Cluster<Source> = flow.cluster::<Source>();
577    /// let to_send: KeyedStream<_, _, Cluster<_>, _> = source
578    ///     .source_iter(q!(vec![0, 1, 2, 3]))
579    ///     .map(q!(|x| (x, x + 123)))
580    ///     .into_keyed();
581    /// let destination_process = flow.process::<Destination>();
582    /// let all_received = to_send.send(&destination_process, TCP.fail_stop().bincode()); // KeyedStream<(MemberId<Source>, i32), i32, ...>
583    /// # all_received.entries().send(&p2, TCP.fail_stop().bincode())
584    /// # }, |mut stream| async move {
585    /// // if there are 4 members in the source cluster, the destination process receives four messages from each source member
586    /// // {
587    /// //     (MemberId<Source>(0), 0): [123], (MemberId<Source>(1), 0): [123], ...,
588    /// //     (MemberId<Source>(0), 1): [124], (MemberId<Source>(1), 1): [124], ...,
589    /// //     ...
590    /// // }
591    /// # let mut results = Vec::new();
592    /// # for w in 0..16 {
593    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
594    /// # }
595    /// # results.sort();
596    /// # assert_eq!(results, vec![
597    /// #   "((MemberId::<()>(0), 0), 123)",
598    /// #   "((MemberId::<()>(0), 1), 124)",
599    /// #   "((MemberId::<()>(0), 2), 125)",
600    /// #   "((MemberId::<()>(0), 3), 126)",
601    /// #   "((MemberId::<()>(1), 0), 123)",
602    /// #   "((MemberId::<()>(1), 1), 124)",
603    /// #   "((MemberId::<()>(1), 2), 125)",
604    /// #   "((MemberId::<()>(1), 3), 126)",
605    /// #   "((MemberId::<()>(2), 0), 123)",
606    /// #   "((MemberId::<()>(2), 1), 124)",
607    /// #   "((MemberId::<()>(2), 2), 125)",
608    /// #   "((MemberId::<()>(2), 3), 126)",
609    /// #   "((MemberId::<()>(3), 0), 123)",
610    /// #   "((MemberId::<()>(3), 1), 124)",
611    /// #   "((MemberId::<()>(3), 2), 125)",
612    /// #   "((MemberId::<()>(3), 3), 126)",
613    /// # ]);
614    /// # }));
615    /// # }
616    /// ```
617    pub fn send<L2, N: NetworkFor<(K, V)>>(
618        self,
619        to: &Process<'a, L2>,
620        via: N,
621    ) -> KeyedStream<
622        (MemberId<L>, K),
623        V,
624        Process<'a, L2>,
625        Unbounded,
626        <O as MinOrder<N::OrderingGuarantee>>::Min,
627        R,
628    >
629    where
630        O: MinOrder<N::OrderingGuarantee>,
631    {
632        let name = via.name();
633        assert!(
634            !to.multiversioned() || name.is_some(),
635            "Cannot send to a multiversioned location without a channel name. Please provide a name for the network."
636        );
637
638        let (serialize, deserialize) = if N::is_embedded() {
639            (
640                NetworkSend::Embedded {
641                    tag: None,
642                    element_type: quote_type::<(K, V)>().into(),
643                },
644                NetworkRecv::Embedded {
645                    tag: Some(quote_type::<L>().into()),
646                    element_type: quote_type::<(K, V)>().into(),
647                },
648            )
649        } else {
650            (
651                NetworkSend::Custom {
652                    serialize_fn: Some(N::serialize_thunk(false).into()),
653                },
654                NetworkRecv::Custom {
655                    deserialize_fn: Some(N::deserialize_thunk(Some(&quote_type::<L>())).into()),
656                },
657            )
658        };
659
660        let raw_stream: Stream<
661            (MemberId<L>, (K, V)),
662            Process<'a, L2>,
663            Unbounded,
664            <O as MinOrder<N::OrderingGuarantee>>::Min,
665            R,
666        > = Stream::new(
667            to.clone(),
668            HydroNode::Network {
669                name: name.map(ToOwned::to_owned),
670                networking_info: N::networking_info(),
671                serialize,
672                deserialize,
673                instantiate_fn: DebugInstantiate::Building,
674                input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
675                metadata: to.new_node_metadata(Stream::<
676                    (MemberId<L>, (K, V)),
677                    Cluster<'a, L2>,
678                    Unbounded,
679                    <O as MinOrder<N::OrderingGuarantee>>::Min,
680                    R,
681                >::collection_kind()),
682            },
683        );
684
685        raw_stream
686            .map(q!(|(sender, (k, v))| ((sender, k), v)))
687            .into_keyed()
688    }
689}