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