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("e_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("e_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}