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