1use std::fmt::Debug;
17use std::future::Future;
18#[cfg(feature = "tokio")]
19use std::marker::PhantomData;
20use std::num::ParseIntError;
21#[cfg(feature = "tokio")]
22use std::time::Duration;
23
24#[cfg(feature = "tokio")]
25use bytes::{Bytes, BytesMut};
26use futures::stream::Stream as FuturesStream;
27use proc_macro2::Span;
28use quote::quote;
29#[cfg(feature = "tokio")]
30use serde::de::DeserializeOwned;
31use serde::{Deserialize, Serialize};
32use slotmap::{Key, new_key_type};
33#[cfg(feature = "tokio")]
34use stageleft::quote_type;
35use stageleft::runtime_support::{FreeVariableWithContextWithProps, QuoteTokens};
36use stageleft::{QuotedWithContext, q};
37use syn::parse_quote;
38#[cfg(feature = "tokio")]
39use tokio_util::codec::{Decoder, Encoder, LengthDelimitedCodec};
40
41#[cfg(feature = "tokio")]
42use crate::compile::builder::ExternalPortId;
43#[cfg(feature = "tokio")]
44use crate::compile::ir::DebugInstantiate;
45use crate::compile::ir::{
46 ClusterMembersState, HydroIrOpMetadata, HydroNode, HydroRoot, HydroSource,
47};
48use crate::forward_handle::ForwardRef;
49#[cfg(stageleft_runtime)]
50use crate::forward_handle::{CycleCollection, ForwardHandle};
51use crate::live_collections::boundedness::{Bounded, Unbounded};
52use crate::live_collections::keyed_stream::KeyedStream;
53use crate::live_collections::singleton::Singleton;
54#[cfg(feature = "sim")]
55#[cfg(stageleft_runtime)]
56use crate::live_collections::stream::networking::serialize_bincode;
57use crate::live_collections::stream::{ExactlyOnce, NoOrder, Stream, TotalOrder};
58#[cfg(feature = "tokio")]
59use crate::live_collections::stream::{Ordering, Retries};
60#[cfg(stageleft_runtime)]
61use crate::location::dynamic::DynLocation;
62use crate::location::dynamic::{ClusterConsistency, LocationId};
63#[cfg(feature = "tokio")]
64use crate::location::external_process::{
65 ExternalBincodeBidi, ExternalBincodeSink, ExternalBytesPort, Many, NotMany,
66};
67use crate::nondet::NonDet;
68#[cfg(feature = "tokio")]
69use crate::properties::manual_proof;
70#[cfg(feature = "sim")]
71use crate::sim::SimSender;
72use crate::staging_util::get_this_crate;
73
74pub mod dynamic;
75
76pub mod external_process;
77pub use external_process::External;
78
79pub mod process;
80pub use process::Process;
81
82pub mod cluster;
83pub use cluster::Cluster;
84
85pub mod member_id;
86pub use member_id::{MemberId, TaglessMemberId};
87
88pub mod tick;
89pub use tick::{Atomic, Tick};
90
91#[derive(PartialEq, Eq, Clone, Debug, Hash, Serialize, Deserialize)]
94pub enum MembershipEvent {
95 Joined,
97 Left,
99}
100
101#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
107pub enum NetworkHint {
108 Auto,
110 TcpPort(Option<u16>),
115}
116
117#[track_caller]
118pub(crate) fn check_matching_location<'a, L: Location<'a>>(l1: &L, l2: &L) {
119 assert_eq!(Location::id(l1), Location::id(l2), "locations do not match");
120}
121
122#[stageleft::export(LocationKey)]
123new_key_type! {
124 pub struct LocationKey;
126}
127
128impl std::fmt::Display for LocationKey {
129 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
130 write!(f, "loc{:?}", self.data()) }
132}
133
134impl std::str::FromStr for LocationKey {
137 type Err = Option<ParseIntError>;
138
139 fn from_str(s: &str) -> Result<Self, Self::Err> {
140 let nvn = s.strip_prefix("loc").ok_or(None)?;
141 let (idx, ver) = nvn.split_once('v').ok_or(None)?;
142 let idx: u64 = idx.parse()?;
143 let ver: u64 = ver.parse()?;
144 Ok(slotmap::KeyData::from_ffi((ver << 32) | idx).into())
145 }
146}
147
148impl LocationKey {
149 pub const FIRST: Self = Self(slotmap::KeyData::from_ffi(0x0000000100000001)); #[cfg(test)]
155 pub const TEST_KEY_1: Self = Self(slotmap::KeyData::from_ffi(0x000000FF00000001)); #[cfg(test)]
159 pub const TEST_KEY_2: Self = Self(slotmap::KeyData::from_ffi(0x000000FF00000002)); }
161
162impl<Ctx> FreeVariableWithContextWithProps<Ctx, ()> for LocationKey {
164 type O = LocationKey;
165
166 fn to_tokens(self, _ctx: &Ctx) -> (QuoteTokens, ())
167 where
168 Self: Sized,
169 {
170 let root = get_this_crate();
171 let n = Key::data(&self).as_ffi();
172 (
173 QuoteTokens {
174 prelude: None,
175 expr: Some(quote! {
176 #root::location::LocationKey::from(#root::runtime_support::slotmap::KeyData::from_ffi(#n))
177 }),
178 },
179 (),
180 )
181 }
182}
183
184#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Serialize)]
186pub enum LocationType {
187 Process,
189 Cluster,
191 External,
193}
194
195pub trait TopLevel<'a>: Location<'a> {}
197
198#[cfg(feature = "sim")]
199#[cfg(stageleft_runtime)]
200fn register_serialized_external_input<'a, At, L, T>(
201 at: &At,
202 from: &External<'_, L>,
203 deserialize_fn: syn::Expr,
204) -> (
205 ExternalPortId,
206 Stream<T, At::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
207)
208where
209 At: TopLevel<'a> + Sized,
210{
211 let (port_id, stream, sink) = at.register_serialized_single_client::<_, T, ()>(
212 from,
213 serialize_bincode::<()>(false),
214 deserialize_fn,
215 );
216 sink.complete(stream.location().source_iter(q!([])));
217
218 (port_id, stream)
219}
220
221#[expect(
235 private_bounds,
236 reason = "only internal Hydro code can define location types"
237)]
238pub trait Location<'a>: DynLocation {
239 type Root: Location<'a>;
244
245 type SimHookScope: crate::sim_hooks::BindableHookScope;
254
255 type DropConsistency: Location<'a, DropConsistency = Self::DropConsistency, SimHookScope = Self::SimHookScope>;
257
258 fn root(&self) -> Self::Root;
263
264 fn drop_consistency(&self) -> Self::DropConsistency;
266 fn consistency() -> Option<ClusterConsistency>;
268
269 fn with_consistency_of<L2: Location<'a, DropConsistency = Self::DropConsistency>>(&self) -> L2 {
271 L2::from_drop_consistency(self.drop_consistency())
272 }
273
274 #[doc(hidden)]
275 fn from_drop_consistency(l2: Self::DropConsistency) -> Self;
276
277 fn try_tick(&self) -> Option<Tick<Self>> {
284 if Self::is_top_level() {
285 let id = if let LocationId::Atomic { .. } = self.id() {
286 None
287 } else {
288 Some(self.flow_state().borrow_mut().next_clock_id())
289 };
290 Some(Tick {
291 id,
292 l: self.clone(),
293 })
294 } else {
295 None
296 }
297 }
298
299 fn id(&self) -> LocationId {
301 DynLocation::dyn_id(self)
302 }
303
304 fn tick(&self) -> Tick<Self> {
330 self.try_tick().expect("cannot create nested ticks")
331 }
332
333 fn spin(&self) -> Stream<(), Self, Unbounded, TotalOrder, ExactlyOnce>
358 where
359 Self: TopLevel<'a> + Sized,
360 {
361 Stream::new(
362 self.clone(),
363 HydroNode::Source {
364 source: HydroSource::Spin(),
365 metadata: self.new_node_metadata(Stream::<
366 (),
367 Self,
368 Unbounded,
369 TotalOrder,
370 ExactlyOnce,
371 >::collection_kind()),
372 },
373 )
374 }
375
376 fn source_stream<T, E>(
397 &self,
398 e: impl QuotedWithContext<'a, E, Self>,
399 ) -> Stream<T, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>
400 where
401 E: FuturesStream<Item = T> + Unpin,
402 Self: TopLevel<'a> + Sized,
403 {
404 let e = e.splice_untyped_ctx(self);
405
406 let target_location = self.drop_consistency();
407 Stream::new(
408 target_location.clone(),
409 HydroNode::Source {
410 source: HydroSource::Stream(e.into()),
411 metadata: target_location.new_node_metadata(Stream::<
412 T,
413 Self::DropConsistency,
414 Unbounded,
415 TotalOrder,
416 ExactlyOnce,
417 >::collection_kind()),
418 },
419 )
420 }
421
422 fn source_iter<T, E>(
444 &self,
445 e: impl QuotedWithContext<'a, E, Self>,
446 ) -> Stream<T, Self::DropConsistency, Bounded, TotalOrder, ExactlyOnce>
447 where
448 E: IntoIterator<Item = T>,
449 Self: Sized,
450 {
451 let e = e.splice_typed_ctx(self);
452
453 let target_location = self.drop_consistency();
454 Stream::new(
455 target_location.clone(),
456 HydroNode::Source {
457 source: HydroSource::Iter(e.into()),
458 metadata: target_location.new_node_metadata(Stream::<
459 T,
460 Self::DropConsistency,
461 Bounded,
462 TotalOrder,
463 ExactlyOnce,
464 >::collection_kind()),
465 },
466 )
467 }
468
469 #[deprecated(note = "use .source_cluster_membership_stream(...) instead")]
470 fn source_cluster_members<C: 'a>(
509 &self,
510 cluster: &Cluster<'a, C>,
511 nondet_start: NonDet,
512 ) -> KeyedStream<MemberId<C>, MembershipEvent, Self::DropConsistency, Unbounded>
513 where
514 Self: TopLevel<'a> + Sized,
515 {
516 self.source_cluster_membership_stream(cluster, nondet_start)
517 }
518
519 fn source_cluster_membership_stream<C: 'a>(
558 &self,
559 cluster: &Cluster<'a, C>,
560 _nondet_start: NonDet,
561 ) -> KeyedStream<MemberId<C>, MembershipEvent, Self::DropConsistency, Unbounded>
562 where
563 Self: TopLevel<'a> + Sized,
564 {
565 let target_consistency = self.drop_consistency();
566 Stream::new(
567 target_consistency.clone(),
568 HydroNode::Source {
569 source: HydroSource::ClusterMembers(cluster.id(), ClusterMembersState::Uninit),
570 metadata: target_consistency.new_node_metadata(Stream::<
571 (TaglessMemberId, MembershipEvent),
572 Self,
573 Unbounded,
574 TotalOrder,
575 ExactlyOnce,
576 >::collection_kind(
577 )),
578 },
579 )
580 .map(q!(|(k, v)| (MemberId::from_tagless(k), v)))
581 .into_keyed()
582 }
583
584 #[cfg(feature = "tokio")]
592 fn source_external_bytes<L>(
593 &self,
594 from: &External<'_, L>,
595 ) -> (
596 ExternalBytesPort,
597 Stream<BytesMut, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
598 )
599 where
600 Self: TopLevel<'a> + Sized,
601 {
602 let (port, stream, sink) =
603 self.bind_single_client::<_, Bytes, LengthDelimitedCodec>(from, NetworkHint::Auto);
604
605 sink.complete(stream.location().source_iter(q!([])));
606
607 (port, stream)
608 }
609
610 #[cfg(feature = "tokio")]
617 fn source_external_bincode<L, T, O: Ordering, R: Retries>(
618 &self,
619 from: &External<'_, L>,
620 ) -> (
621 ExternalBincodeSink<T, NotMany, O, R>,
622 Stream<T, Self::DropConsistency, Unbounded, O, R>,
623 )
624 where
625 Self: TopLevel<'a> + Sized,
626 T: Serialize + DeserializeOwned,
627 {
628 let (port, stream, sink) = self.bind_single_client_bincode::<_, T, ()>(from);
629 sink.complete(stream.location().source_iter(q!([])));
630
631 (
632 ExternalBincodeSink {
633 process_key: from.key,
634 port_id: port.port_id,
635 _phantom: PhantomData,
636 },
637 stream.weaken_ordering().weaken_retries(),
638 )
639 }
640
641 #[cfg(feature = "sim")]
647 fn sim_input<T, O: Ordering, R: Retries>(
648 &self,
649 ) -> (
650 SimSender<T, O, R>,
651 Stream<T, Self::DropConsistency, Unbounded, O, R>,
652 )
653 where
654 Self: TopLevel<'a> + Sized,
655 T: Serialize + DeserializeOwned,
656 {
657 self.sim_input_with::<crate::sim::codec::BincodeCodec, T, O, R>()
658 }
659
660 #[cfg(feature = "sim")]
668 fn sim_input_with<C: crate::sim::codec::SimCodec<T>, T, O: Ordering, R: Retries>(
669 &self,
670 ) -> (
671 SimSender<T, O, R>,
672 Stream<T, Self::DropConsistency, Unbounded, O, R>,
673 )
674 where
675 Self: TopLevel<'a> + Sized,
676 {
677 let external_location: External<'a, ()> = External {
678 key: LocationKey::FIRST,
679 flow_state: self.flow_state().clone(),
680 _phantom: PhantomData,
681 };
682
683 let (external_port_id, stream) = register_serialized_external_input(
684 self,
685 &external_location,
686 crate::sim::codec::staged_deserialize::<T, C>(),
687 );
688
689 (
690 SimSender(external_port_id, PhantomData, C::encode),
691 stream.weaken_ordering().weaken_retries(),
692 )
693 }
694
695 fn embedded_input<T>(
701 &self,
702 name: impl Into<String>,
703 ) -> Stream<T, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>
704 where
705 Self: TopLevel<'a> + Sized,
706 {
707 let ident = syn::Ident::new(&name.into(), Span::call_site());
708
709 let target_location = self.drop_consistency();
710 Stream::new(
711 target_location.clone(),
712 HydroNode::Source {
713 source: HydroSource::Embedded(ident),
714 metadata: target_location.new_node_metadata(Stream::<
715 T,
716 Self,
717 Unbounded,
718 TotalOrder,
719 ExactlyOnce,
720 >::collection_kind()),
721 },
722 )
723 }
724
725 fn embedded_singleton_input<T>(
731 &self,
732 name: impl Into<String>,
733 ) -> Singleton<T, Self::DropConsistency, Bounded>
734 where
735 Self: TopLevel<'a> + Sized,
736 {
737 let ident = syn::Ident::new(&name.into(), Span::call_site());
738
739 let target_location = self.drop_consistency();
740 Singleton::new(
741 target_location.clone(),
742 HydroNode::Source {
743 source: HydroSource::EmbeddedSingleton(ident),
744 metadata: target_location
745 .new_node_metadata(Singleton::<T, Self, Bounded>::collection_kind()),
746 },
747 )
748 }
749
750 #[cfg(feature = "tokio")]
795 #[expect(clippy::type_complexity, reason = "stream markers")]
796 fn bind_single_client<L, T, Codec: Encoder<T> + Decoder>(
797 &self,
798 from: &External<'_, L>,
799 port_hint: NetworkHint,
800 ) -> (
801 ExternalBytesPort<NotMany>,
802 Stream<<Codec as Decoder>::Item, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
803 ForwardHandle<'a, Stream<T, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>>,
804 )
805 where
806 Self: TopLevel<'a> + Sized,
807 {
808 let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
809 let target_consistency = self.drop_consistency();
810
811 let (fwd_ref, to_sink) = target_consistency.forward_ref::<Stream<
812 T,
813 Self::DropConsistency,
814 Unbounded,
815 TotalOrder,
816 ExactlyOnce,
817 >>();
818 let mut flow_state_borrow = self.flow_state().borrow_mut();
819
820 flow_state_borrow.push_root(HydroRoot::SendExternal {
821 to_external_key: from.key,
822 to_port_id: next_external_port_id,
823 to_many: false,
824 unpaired: false,
825 serialize_fn: None,
826 instantiate_fn: DebugInstantiate::Building,
827 input: Box::new(to_sink.ir_node.replace(HydroNode::Placeholder)),
828 op_metadata: HydroIrOpMetadata::new(),
829 });
830 drop(flow_state_borrow);
831
832 let raw_stream: Stream<
833 Result<<Codec as Decoder>::Item, <Codec as Decoder>::Error>,
834 Self::DropConsistency,
835 Unbounded,
836 TotalOrder,
837 ExactlyOnce,
838 > = Stream::new(
839 target_consistency.clone(),
840 HydroNode::ExternalInput {
841 from_external_key: from.key,
842 from_port_id: next_external_port_id,
843 from_many: false,
844 codec_type: quote_type::<Codec>().into(),
845 port_hint,
846 instantiate_fn: DebugInstantiate::Building,
847 deserialize_fn: None,
848 metadata: target_consistency.new_node_metadata(Stream::<
849 Result<<Codec as Decoder>::Item, <Codec as Decoder>::Error>,
850 Self::DropConsistency,
851 Unbounded,
852 TotalOrder,
853 ExactlyOnce,
854 >::collection_kind(
855 )),
856 },
857 );
858
859 (
860 ExternalBytesPort {
861 process_key: from.key,
862 port_id: next_external_port_id,
863 _phantom: PhantomData,
864 },
865 raw_stream.flatten_ordered(),
866 fwd_ref,
867 )
868 }
869
870 #[doc(hidden)]
872 #[cfg(feature = "tokio")]
873 #[expect(clippy::type_complexity, reason = "stream markers")]
874 fn register_serialized_single_client<L, InT, OutT>(
875 &self,
876 from: &External<'_, L>,
877 serialize_fn: syn::Expr,
878 deserialize_fn: syn::Expr,
879 ) -> (
880 ExternalPortId,
881 Stream<InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
882 ForwardHandle<'a, Stream<OutT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>>,
883 )
884 where
885 Self: TopLevel<'a> + Sized,
886 {
887 let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
888
889 let target_consistency = self.drop_consistency();
890 let (fwd_ref, to_sink) = target_consistency.forward_ref::<Stream<
891 OutT,
892 Self::DropConsistency,
893 Unbounded,
894 TotalOrder,
895 ExactlyOnce,
896 >>();
897 let mut flow_state_borrow = self.flow_state().borrow_mut();
898
899 flow_state_borrow.push_root(HydroRoot::SendExternal {
900 to_external_key: from.key,
901 to_port_id: next_external_port_id,
902 to_many: false,
903 unpaired: false,
904 serialize_fn: Some(serialize_fn.into()),
905 instantiate_fn: DebugInstantiate::Building,
906 input: Box::new(to_sink.ir_node.replace(HydroNode::Placeholder)),
907 op_metadata: HydroIrOpMetadata::new(),
908 });
909 drop(flow_state_borrow);
910
911 let raw_stream: Stream<InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce> =
912 Stream::new(
913 target_consistency.clone(),
914 HydroNode::ExternalInput {
915 from_external_key: from.key,
916 from_port_id: next_external_port_id,
917 from_many: false,
918 codec_type: quote_type::<LengthDelimitedCodec>().into(),
919 port_hint: NetworkHint::Auto,
920 instantiate_fn: DebugInstantiate::Building,
921 deserialize_fn: Some(deserialize_fn.into()),
922 metadata: target_consistency.new_node_metadata(Stream::<
923 InT,
924 Self::DropConsistency,
925 Unbounded,
926 TotalOrder,
927 ExactlyOnce,
928 >::collection_kind(
929 )),
930 },
931 );
932
933 (next_external_port_id, raw_stream, fwd_ref)
934 }
935
936 #[cfg(feature = "tokio")]
946 #[expect(clippy::type_complexity, reason = "stream markers")]
947 fn bind_single_client_bincode<L, InT: DeserializeOwned, OutT: Serialize>(
948 &self,
949 from: &External<'_, L>,
950 ) -> (
951 ExternalBincodeBidi<InT, OutT, NotMany>,
952 Stream<InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
953 ForwardHandle<'a, Stream<OutT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>>,
954 )
955 where
956 Self: TopLevel<'a> + Sized,
957 {
958 let root = get_this_crate();
959
960 let out_t_type = quote_type::<OutT>();
961 let ser_fn: syn::Expr = syn::parse_quote! {
962 #root::runtime_support::stageleft::runtime_support::fn1_type_hint::<#out_t_type, _>(
963 |b| #root::runtime_support::bincode::serialize(&b).unwrap().into()
964 )
965 };
966
967 let in_t_type = quote_type::<InT>();
968 let deser_fn: syn::Expr = syn::parse_quote! {
969 |res| {
970 let b = res.unwrap();
971 #root::runtime_support::bincode::deserialize::<#in_t_type>(&b).unwrap()
972 }
973 };
974
975 let (port_id, raw_stream, fwd_ref) =
976 self.register_serialized_single_client::<_, InT, OutT>(from, ser_fn, deser_fn);
977
978 (
979 ExternalBincodeBidi {
980 process_key: from.key,
981 port_id,
982 _phantom: PhantomData,
983 },
984 raw_stream,
985 fwd_ref,
986 )
987 }
988
989 #[cfg(feature = "tokio")]
1001 #[expect(clippy::type_complexity, reason = "stream markers")]
1002 fn bidi_external_many_bytes<L, T, Codec: Encoder<T> + Decoder>(
1003 &self,
1004 from: &External<'_, L>,
1005 port_hint: NetworkHint,
1006 ) -> (
1007 ExternalBytesPort<Many>,
1008 KeyedStream<
1009 u64,
1010 <Codec as Decoder>::Item,
1011 Self::DropConsistency,
1012 Unbounded,
1013 TotalOrder,
1014 ExactlyOnce,
1015 >,
1016 KeyedStream<
1017 u64,
1018 MembershipEvent,
1019 Self::DropConsistency,
1020 Unbounded,
1021 TotalOrder,
1022 ExactlyOnce,
1023 >,
1024 ForwardHandle<
1025 'a,
1026 KeyedStream<u64, T, Self::DropConsistency, Unbounded, NoOrder, ExactlyOnce>,
1027 >,
1028 )
1029 where
1030 Self: TopLevel<'a> + Sized,
1031 {
1032 let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
1033
1034 let target_consistency = self.drop_consistency();
1035 let (fwd_ref, to_sink) = target_consistency.forward_ref::<KeyedStream<
1036 u64,
1037 T,
1038 Self::DropConsistency,
1039 Unbounded,
1040 NoOrder,
1041 ExactlyOnce,
1042 >>();
1043 let to_sink_input = Box::new(to_sink.entries().ir_node.replace(HydroNode::Placeholder));
1044 let mut flow_state_borrow = self.flow_state().borrow_mut();
1045
1046 flow_state_borrow.push_root(HydroRoot::SendExternal {
1047 to_external_key: from.key,
1048 to_port_id: next_external_port_id,
1049 to_many: true,
1050 unpaired: false,
1051 serialize_fn: None,
1052 instantiate_fn: DebugInstantiate::Building,
1053 input: to_sink_input,
1054 op_metadata: HydroIrOpMetadata::new(),
1055 });
1056 drop(flow_state_borrow);
1057
1058 let raw_stream: Stream<
1059 Result<(u64, <Codec as Decoder>::Item), <Codec as Decoder>::Error>,
1060 Self::DropConsistency,
1061 Unbounded,
1062 TotalOrder,
1063 ExactlyOnce,
1064 > = Stream::new(
1065 target_consistency.clone(),
1066 HydroNode::ExternalInput {
1067 from_external_key: from.key,
1068 from_port_id: next_external_port_id,
1069 from_many: true,
1070 codec_type: quote_type::<Codec>().into(),
1071 port_hint,
1072 instantiate_fn: DebugInstantiate::Building,
1073 deserialize_fn: None,
1074 metadata: target_consistency.new_node_metadata(Stream::<
1075 Result<(u64, <Codec as Decoder>::Item), <Codec as Decoder>::Error>,
1076 Self::DropConsistency,
1077 Unbounded,
1078 TotalOrder,
1079 ExactlyOnce,
1080 >::collection_kind(
1081 )),
1082 },
1083 );
1084
1085 let membership_stream_ident = syn::Ident::new(
1086 &format!(
1087 "__hydro_deploy_many_{}_{}_membership",
1088 from.key, next_external_port_id
1089 ),
1090 Span::call_site(),
1091 );
1092 let membership_stream_expr: syn::Expr = parse_quote!(#membership_stream_ident);
1093 let raw_membership_stream: KeyedStream<
1094 u64,
1095 bool,
1096 Self::DropConsistency,
1097 Unbounded,
1098 TotalOrder,
1099 ExactlyOnce,
1100 > = KeyedStream::new(
1101 target_consistency.clone(),
1102 HydroNode::Source {
1103 source: HydroSource::Stream(membership_stream_expr.into()),
1104 metadata: target_consistency.new_node_metadata(KeyedStream::<
1105 u64,
1106 bool,
1107 Self::DropConsistency,
1108 Unbounded,
1109 TotalOrder,
1110 ExactlyOnce,
1111 >::collection_kind(
1112 )),
1113 },
1114 );
1115
1116 (
1117 ExternalBytesPort {
1118 process_key: from.key,
1119 port_id: next_external_port_id,
1120 _phantom: PhantomData,
1121 },
1122 raw_stream
1123 .flatten_ordered() .into_keyed(),
1125 raw_membership_stream.map(q!(|join| {
1126 if join {
1127 MembershipEvent::Joined
1128 } else {
1129 MembershipEvent::Left
1130 }
1131 })),
1132 fwd_ref,
1133 )
1134 }
1135
1136 #[cfg(feature = "tokio")]
1152 #[expect(clippy::type_complexity, reason = "stream markers")]
1153 fn bidi_external_many_bincode<L, InT: DeserializeOwned, OutT: Serialize>(
1154 &self,
1155 from: &External<'_, L>,
1156 ) -> (
1157 ExternalBincodeBidi<InT, OutT, Many>,
1158 KeyedStream<u64, InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
1159 KeyedStream<
1160 u64,
1161 MembershipEvent,
1162 Self::DropConsistency,
1163 Unbounded,
1164 TotalOrder,
1165 ExactlyOnce,
1166 >,
1167 ForwardHandle<
1168 'a,
1169 KeyedStream<u64, OutT, Self::DropConsistency, Unbounded, NoOrder, ExactlyOnce>,
1170 >,
1171 )
1172 where
1173 Self: TopLevel<'a> + Sized,
1174 {
1175 let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
1176
1177 let target_consistency = self.drop_consistency();
1178 let (fwd_ref, to_sink) = target_consistency.forward_ref::<KeyedStream<
1179 u64,
1180 OutT,
1181 Self::DropConsistency,
1182 Unbounded,
1183 NoOrder,
1184 ExactlyOnce,
1185 >>();
1186 let to_sink_input = Box::new(to_sink.entries().ir_node.replace(HydroNode::Placeholder));
1187 let mut flow_state_borrow = self.flow_state().borrow_mut();
1188
1189 let root = get_this_crate();
1190
1191 let out_t_type = quote_type::<OutT>();
1192 let ser_fn: syn::Expr = syn::parse_quote! {
1193 #root::runtime_support::stageleft::runtime_support::fn1_type_hint::<(u64, #out_t_type), _>(
1194 |(id, b)| (id, #root::runtime_support::bincode::serialize(&b).unwrap().into())
1195 )
1196 };
1197
1198 flow_state_borrow.push_root(HydroRoot::SendExternal {
1199 to_external_key: from.key,
1200 to_port_id: next_external_port_id,
1201 to_many: true,
1202 unpaired: false,
1203 serialize_fn: Some(ser_fn.into()),
1204 instantiate_fn: DebugInstantiate::Building,
1205 input: to_sink_input,
1206 op_metadata: HydroIrOpMetadata::new(),
1207 });
1208 drop(flow_state_borrow);
1209
1210 let in_t_type = quote_type::<InT>();
1211
1212 let deser_fn: syn::Expr = syn::parse_quote! {
1213 |res| {
1214 let (id, b) = res.unwrap();
1215 (id, #root::runtime_support::bincode::deserialize::<#in_t_type>(&b).unwrap())
1216 }
1217 };
1218
1219 let raw_stream: KeyedStream<
1220 u64,
1221 InT,
1222 Self::DropConsistency,
1223 Unbounded,
1224 TotalOrder,
1225 ExactlyOnce,
1226 > = KeyedStream::new(
1227 target_consistency.clone(),
1228 HydroNode::ExternalInput {
1229 from_external_key: from.key,
1230 from_port_id: next_external_port_id,
1231 from_many: true,
1232 codec_type: quote_type::<LengthDelimitedCodec>().into(),
1233 port_hint: NetworkHint::Auto,
1234 instantiate_fn: DebugInstantiate::Building,
1235 deserialize_fn: Some(deser_fn.into()),
1236 metadata: target_consistency.new_node_metadata(KeyedStream::<
1237 u64,
1238 InT,
1239 Self::DropConsistency,
1240 Unbounded,
1241 TotalOrder,
1242 ExactlyOnce,
1243 >::collection_kind(
1244 )),
1245 },
1246 );
1247
1248 let membership_stream_ident = syn::Ident::new(
1249 &format!(
1250 "__hydro_deploy_many_{}_{}_membership",
1251 from.key, next_external_port_id
1252 ),
1253 Span::call_site(),
1254 );
1255 let membership_stream_expr: syn::Expr = parse_quote!(#membership_stream_ident);
1256 let raw_membership_stream: KeyedStream<
1257 u64,
1258 bool,
1259 Self::DropConsistency,
1260 Unbounded,
1261 TotalOrder,
1262 ExactlyOnce,
1263 > = KeyedStream::new(
1264 target_consistency.clone(),
1265 HydroNode::Source {
1266 source: HydroSource::Stream(membership_stream_expr.into()),
1267 metadata: target_consistency.new_node_metadata(KeyedStream::<
1268 u64,
1269 bool,
1270 Self::DropConsistency,
1271 Unbounded,
1272 TotalOrder,
1273 ExactlyOnce,
1274 >::collection_kind(
1275 )),
1276 },
1277 );
1278
1279 (
1280 ExternalBincodeBidi {
1281 process_key: from.key,
1282 port_id: next_external_port_id,
1283 _phantom: PhantomData,
1284 },
1285 raw_stream,
1286 raw_membership_stream.map(q!(|join| {
1287 if join {
1288 MembershipEvent::Joined
1289 } else {
1290 MembershipEvent::Left
1291 }
1292 })),
1293 fwd_ref,
1294 )
1295 }
1296
1297 fn sidecar_bidi<InT: 'static, OutT: 'static, F>(
1350 &self,
1351 sidecar: impl QuotedWithContext<'a, F, Self>,
1352 ) -> (
1353 Stream<InT, Self, Unbounded, TotalOrder, ExactlyOnce>,
1354 ForwardHandle<'a, Stream<OutT, Self, Unbounded, NoOrder, ExactlyOnce>>,
1355 )
1356 where
1357 Self: Sized + TopLevel<'a>,
1358 {
1359 let location_key = Location::id(self).key();
1360
1361 let sidecar_id = self.flow_state().borrow_mut().next_sidecar_id();
1362 let (stream_ident, sink_ident) = sidecar_id.idents();
1363
1364 let sidecar_closure: syn::Expr = sidecar.splice_untyped_ctx(self);
1365 self.flow_state()
1366 .borrow_mut()
1367 .sidecars
1368 .push(crate::compile::builder::Sidecar::Bidi {
1369 location_key,
1370 sidecar_id,
1371 sidecar_closure: Box::new(sidecar_closure),
1372 });
1373
1374 let source_expr: syn::Expr = parse_quote! {
1376 #stream_ident
1377 };
1378 let inbound: Stream<InT, Self, Unbounded, TotalOrder, ExactlyOnce> = Stream::new(
1379 self.clone(),
1380 HydroNode::Source {
1381 source: HydroSource::Stream(source_expr.into()),
1382 metadata: self.new_node_metadata(Stream::<
1383 InT,
1384 Self,
1385 Unbounded, TotalOrder, ExactlyOnce,
1388 >::collection_kind()),
1389 },
1390 );
1391
1392 let (fwd_ref, to_sink): (
1394 ForwardHandle<'a, Stream<OutT, Self, Unbounded, NoOrder, ExactlyOnce>>,
1395 Stream<OutT, Self, Unbounded, NoOrder, ExactlyOnce>,
1396 ) = self.forward_ref();
1397
1398 let sink_expr: syn::Expr = parse_quote! {
1399 #sink_ident
1400 };
1401
1402 let sink_input_ir = to_sink.ir_node.replace(HydroNode::Placeholder);
1403 self.flow_state()
1404 .borrow_mut()
1405 .try_push_root(HydroRoot::DestSink {
1406 sink: sink_expr.into(),
1407 input: Box::new(sink_input_ir),
1408 op_metadata: HydroIrOpMetadata::new(),
1409 });
1410
1411 (inbound, fwd_ref)
1412 }
1413
1414 fn singleton<T>(
1434 &self,
1435 e: impl QuotedWithContext<'a, T, Self>,
1436 ) -> Singleton<T, Self::DropConsistency, Bounded>
1437 where
1438 Self: Sized,
1439 {
1440 let e = e.splice_untyped_ctx(self);
1441
1442 let target_location = self.drop_consistency();
1443 Singleton::new(
1444 target_location.clone(),
1445 HydroNode::SingletonSource {
1446 value: e.into(),
1447 first_tick_only: false,
1448 metadata: target_location.new_node_metadata(Singleton::<
1449 T,
1450 Self::DropConsistency,
1451 Bounded,
1452 >::collection_kind()),
1453 },
1454 )
1455 }
1456
1457 fn singleton_future<F>(
1480 &self,
1481 e: impl QuotedWithContext<'a, F, Self>,
1482 ) -> Singleton<F::Output, Self::DropConsistency, Bounded>
1483 where
1484 F: Future,
1485 Self: Sized,
1486 {
1487 self.singleton(e).resolve_future_blocking()
1488 }
1489
1490 #[cfg(feature = "tokio")]
1499 fn source_interval(
1500 &self,
1501 interval: impl QuotedWithContext<'a, Duration, Self> + Copy + 'a,
1502 ) -> Stream<(), Self, Unbounded, TotalOrder, ExactlyOnce>
1503 where
1504 Self: TopLevel<'a> + Sized,
1505 {
1506 self.source_stream(q!(tokio_stream::StreamExt::map(
1507 tokio_stream::wrappers::IntervalStream::new(tokio::time::interval(interval)),
1508 |_| ()
1509 )))
1510 .assert_has_consistency_of_trusted(
1511 manual_proof!(),
1512 )
1513 }
1514
1515 #[cfg(feature = "tokio")]
1522 fn source_interval_delayed(
1523 &self,
1524 delay: impl QuotedWithContext<'a, Duration, Self> + Copy + 'a,
1525 interval: impl QuotedWithContext<'a, Duration, Self> + Copy + 'a,
1526 ) -> Stream<(), Self, Unbounded, TotalOrder, ExactlyOnce>
1527 where
1528 Self: TopLevel<'a> + Sized,
1529 {
1530 self.source_stream(q!(tokio_stream::StreamExt::map(
1531 tokio_stream::wrappers::IntervalStream::new(tokio::time::interval_at(
1532 tokio::time::Instant::now() + delay,
1533 interval,
1534 )),
1535 |_| ()
1536 )))
1537 .assert_has_consistency_of_trusted(
1538 manual_proof!(),
1539 )
1540 }
1541
1542 fn forward_ref<S>(&self) -> (ForwardHandle<'a, S>, S)
1582 where
1583 S: CycleCollection<'a, ForwardRef, Location = Self>,
1584 {
1585 let cycle_id = self.flow_state().borrow_mut().next_cycle_id();
1586 (
1587 ForwardHandle::new(cycle_id, Location::id(self)),
1588 S::create_source(cycle_id, self.clone()),
1589 )
1590 }
1591}
1592
1593#[cfg(feature = "deploy")]
1594#[cfg(test)]
1595mod tests {
1596 use std::collections::HashSet;
1597
1598 use futures::{SinkExt, StreamExt};
1599 use hydro_deploy::Deployment;
1600 use stageleft::q;
1601 use tokio_util::codec::LengthDelimitedCodec;
1602
1603 use crate::compile::builder::FlowBuilder;
1604 use crate::live_collections::stream::{ExactlyOnce, TotalOrder};
1605 use crate::location::{Location, NetworkHint};
1606 use crate::nondet::nondet;
1607
1608 #[tokio::test]
1609 async fn top_level_singleton_replay_cardinality() {
1610 let mut deployment = Deployment::new();
1611
1612 let mut flow = FlowBuilder::new();
1613 let node = flow.process::<()>();
1614 let external = flow.external::<()>();
1615
1616 let (in_port, input) =
1617 node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
1618 let singleton = node.singleton(q!(123));
1619 let tick = node.tick();
1620 let out = input
1621 .batch(&tick, nondet!())
1622 .cross_singleton(singleton.clone().snapshot(&tick, nondet!()))
1623 .cross_singleton(
1624 singleton
1625 .snapshot(&tick, nondet!())
1626 .into_stream()
1627 .count(),
1628 )
1629 .all_ticks()
1630 .send_bincode_external(&external);
1631
1632 let nodes = flow
1633 .with_process(&node, deployment.Localhost())
1634 .with_external(&external, deployment.Localhost())
1635 .deploy(&mut deployment);
1636
1637 deployment.deploy().await.unwrap();
1638
1639 let mut external_in = nodes.connect(in_port).await;
1640 let mut external_out = nodes.connect(out).await;
1641
1642 deployment.start().await.unwrap();
1643
1644 external_in.send(1).await.unwrap();
1645 assert_eq!(external_out.next().await.unwrap(), ((1, 123), 1));
1646
1647 external_in.send(2).await.unwrap();
1648 assert_eq!(external_out.next().await.unwrap(), ((2, 123), 1));
1649 }
1650
1651 #[tokio::test]
1652 async fn tick_singleton_replay_cardinality() {
1653 let mut deployment = Deployment::new();
1654
1655 let mut flow = FlowBuilder::new();
1656 let node = flow.process::<()>();
1657 let external = flow.external::<()>();
1658
1659 let (in_port, input) =
1660 node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
1661 let tick = node.tick();
1662 let singleton = tick.singleton(q!(123));
1663 let out = input
1664 .batch(&tick, nondet!())
1665 .cross_singleton(singleton.clone())
1666 .cross_singleton(singleton.into_stream().count())
1667 .all_ticks()
1668 .send_bincode_external(&external);
1669
1670 let nodes = flow
1671 .with_process(&node, deployment.Localhost())
1672 .with_external(&external, deployment.Localhost())
1673 .deploy(&mut deployment);
1674
1675 deployment.deploy().await.unwrap();
1676
1677 let mut external_in = nodes.connect(in_port).await;
1678 let mut external_out = nodes.connect(out).await;
1679
1680 deployment.start().await.unwrap();
1681
1682 external_in.send(1).await.unwrap();
1683 assert_eq!(external_out.next().await.unwrap(), ((1, 123), 1));
1684
1685 external_in.send(2).await.unwrap();
1686 assert_eq!(external_out.next().await.unwrap(), ((2, 123), 1));
1687 }
1688
1689 #[tokio::test]
1690 async fn external_bytes() {
1691 let mut deployment = Deployment::new();
1692
1693 let mut flow = FlowBuilder::new();
1694 let first_node = flow.process::<()>();
1695 let external = flow.external::<()>();
1696
1697 let (in_port, input) = first_node.source_external_bytes(&external);
1698 let out = input.send_bincode_external(&external);
1699
1700 let nodes = flow
1701 .with_process(&first_node, deployment.Localhost())
1702 .with_external(&external, deployment.Localhost())
1703 .deploy(&mut deployment);
1704
1705 deployment.deploy().await.unwrap();
1706
1707 let mut external_in = nodes.connect(in_port).await.1;
1708 let mut external_out = nodes.connect(out).await;
1709
1710 deployment.start().await.unwrap();
1711
1712 external_in.send(vec![1, 2, 3].into()).await.unwrap();
1713
1714 assert_eq!(external_out.next().await.unwrap(), vec![1, 2, 3]);
1715 }
1716
1717 #[tokio::test]
1718 async fn multi_external_source() {
1719 let mut deployment = Deployment::new();
1720
1721 let mut flow = FlowBuilder::new();
1722 let first_node = flow.process::<()>();
1723 let external = flow.external::<()>();
1724
1725 let (in_port, input, _membership, complete_sink) =
1726 first_node.bidi_external_many_bincode(&external);
1727 let out = input.entries().send_bincode_external(&external);
1728 complete_sink.complete(
1729 first_node
1730 .source_iter::<(u64, ()), _>(q!([]))
1731 .into_keyed()
1732 .weaken_ordering(),
1733 );
1734
1735 let nodes = flow
1736 .with_process(&first_node, deployment.Localhost())
1737 .with_external(&external, deployment.Localhost())
1738 .deploy(&mut deployment);
1739
1740 deployment.deploy().await.unwrap();
1741
1742 let (_, mut external_in_1) = nodes.connect_bincode(in_port.clone()).await;
1743 let (_, mut external_in_2) = nodes.connect_bincode(in_port).await;
1744 let external_out = nodes.connect(out).await;
1745
1746 deployment.start().await.unwrap();
1747
1748 external_in_1.send(123).await.unwrap();
1749 external_in_2.send(456).await.unwrap();
1750
1751 assert_eq!(
1752 external_out.take(2).collect::<HashSet<_>>().await,
1753 vec![(0, 123), (1, 456)].into_iter().collect()
1754 );
1755 }
1756
1757 #[tokio::test]
1758 async fn second_connection_only_multi_source() {
1759 let mut deployment = Deployment::new();
1760
1761 let mut flow = FlowBuilder::new();
1762 let first_node = flow.process::<()>();
1763 let external = flow.external::<()>();
1764
1765 let (in_port, input, _membership, complete_sink) =
1766 first_node.bidi_external_many_bincode(&external);
1767 let out = input.entries().send_bincode_external(&external);
1768 complete_sink.complete(
1769 first_node
1770 .source_iter::<(u64, ()), _>(q!([]))
1771 .into_keyed()
1772 .weaken_ordering(),
1773 );
1774
1775 let nodes = flow
1776 .with_process(&first_node, deployment.Localhost())
1777 .with_external(&external, deployment.Localhost())
1778 .deploy(&mut deployment);
1779
1780 deployment.deploy().await.unwrap();
1781
1782 let (_, mut _external_in_1) = nodes.connect_bincode(in_port.clone()).await;
1784 let (_, mut external_in_2) = nodes.connect_bincode(in_port).await;
1785 let mut external_out = nodes.connect(out).await;
1786
1787 deployment.start().await.unwrap();
1788
1789 external_in_2.send(456).await.unwrap();
1790
1791 assert_eq!(external_out.next().await.unwrap(), (1, 456));
1792 }
1793
1794 #[tokio::test]
1795 async fn multi_external_bytes() {
1796 let mut deployment = Deployment::new();
1797
1798 let mut flow = FlowBuilder::new();
1799 let first_node = flow.process::<()>();
1800 let external = flow.external::<()>();
1801
1802 let (in_port, input, _membership, complete_sink) = first_node
1803 .bidi_external_many_bytes::<_, _, LengthDelimitedCodec>(&external, NetworkHint::Auto);
1804 let out = input.entries().send_bincode_external(&external);
1805 complete_sink.complete(
1806 first_node
1807 .source_iter(q!([]))
1808 .into_keyed()
1809 .weaken_ordering(),
1810 );
1811
1812 let nodes = flow
1813 .with_process(&first_node, deployment.Localhost())
1814 .with_external(&external, deployment.Localhost())
1815 .deploy(&mut deployment);
1816
1817 deployment.deploy().await.unwrap();
1818
1819 let mut external_in_1 = nodes.connect(in_port.clone()).await.1;
1820 let mut external_in_2 = nodes.connect(in_port).await.1;
1821 let external_out = nodes.connect(out).await;
1822
1823 deployment.start().await.unwrap();
1824
1825 external_in_1.send(vec![1, 2, 3].into()).await.unwrap();
1826 external_in_2.send(vec![4, 5].into()).await.unwrap();
1827
1828 assert_eq!(
1829 external_out.take(2).collect::<HashSet<_>>().await,
1830 vec![
1831 (0, (&[1u8, 2, 3] as &[u8]).into()),
1832 (1, (&[4u8, 5] as &[u8]).into())
1833 ]
1834 .into_iter()
1835 .collect()
1836 );
1837 }
1838
1839 #[tokio::test]
1840 async fn single_client_external_bytes() {
1841 let mut deployment = Deployment::new();
1842 let mut flow = FlowBuilder::new();
1843 let first_node = flow.process::<()>();
1844 let external = flow.external::<()>();
1845 let (port, input, complete_sink) = first_node
1846 .bind_single_client::<_, _, LengthDelimitedCodec>(&external, NetworkHint::Auto);
1847 complete_sink.complete(input.map(q!(|data| {
1848 let mut resp: Vec<u8> = data.into();
1849 resp.push(42);
1850 resp.into() })));
1852
1853 let nodes = flow
1854 .with_process(&first_node, deployment.Localhost())
1855 .with_external(&external, deployment.Localhost())
1856 .deploy(&mut deployment);
1857
1858 deployment.deploy().await.unwrap();
1859 deployment.start().await.unwrap();
1860
1861 let (mut external_out, mut external_in) = nodes.connect(port).await;
1862
1863 external_in.send(vec![1, 2, 3].into()).await.unwrap();
1864 assert_eq!(
1865 external_out.next().await.unwrap().unwrap(),
1866 vec![1, 2, 3, 42]
1867 );
1868 }
1869
1870 #[tokio::test]
1871 async fn echo_external_bytes() {
1872 let mut deployment = Deployment::new();
1873
1874 let mut flow = FlowBuilder::new();
1875 let first_node = flow.process::<()>();
1876 let external = flow.external::<()>();
1877
1878 let (port, input, _membership, complete_sink) = first_node
1879 .bidi_external_many_bytes::<_, _, LengthDelimitedCodec>(&external, NetworkHint::Auto);
1880 complete_sink
1881 .complete(input.map(q!(|bytes| { bytes.into_iter().map(|x| x + 1).collect() })));
1882
1883 let nodes = flow
1884 .with_process(&first_node, deployment.Localhost())
1885 .with_external(&external, deployment.Localhost())
1886 .deploy(&mut deployment);
1887
1888 deployment.deploy().await.unwrap();
1889
1890 let (mut external_out_1, mut external_in_1) = nodes.connect(port.clone()).await;
1891 let (mut external_out_2, mut external_in_2) = nodes.connect(port).await;
1892
1893 deployment.start().await.unwrap();
1894
1895 external_in_1.send(vec![1, 2, 3].into()).await.unwrap();
1896 external_in_2.send(vec![4, 5].into()).await.unwrap();
1897
1898 assert_eq!(external_out_1.next().await.unwrap().unwrap(), vec![2, 3, 4]);
1899 assert_eq!(external_out_2.next().await.unwrap().unwrap(), vec![5, 6]);
1900 }
1901
1902 #[tokio::test]
1903 async fn echo_external_bincode() {
1904 let mut deployment = Deployment::new();
1905
1906 let mut flow = FlowBuilder::new();
1907 let first_node = flow.process::<()>();
1908 let external = flow.external::<()>();
1909
1910 let (port, input, _membership, complete_sink) =
1911 first_node.bidi_external_many_bincode(&external);
1912 complete_sink.complete(input.map(q!(|text: String| { text.to_uppercase() })));
1913
1914 let nodes = flow
1915 .with_process(&first_node, deployment.Localhost())
1916 .with_external(&external, deployment.Localhost())
1917 .deploy(&mut deployment);
1918
1919 deployment.deploy().await.unwrap();
1920
1921 let (mut external_out_1, mut external_in_1) = nodes.connect_bincode(port.clone()).await;
1922 let (mut external_out_2, mut external_in_2) = nodes.connect_bincode(port).await;
1923
1924 deployment.start().await.unwrap();
1925
1926 external_in_1.send("hi".to_owned()).await.unwrap();
1927 external_in_2.send("hello".to_owned()).await.unwrap();
1928
1929 assert_eq!(external_out_1.next().await.unwrap(), "HI");
1930 assert_eq!(external_out_2.next().await.unwrap(), "HELLO");
1931 }
1932
1933 #[tokio::test]
1934 async fn closure_location_name() {
1935 let mut deployment = Deployment::new();
1936 let mut flow = FlowBuilder::new();
1937
1938 enum ClosureProcess {}
1939
1940 let node = flow.process::<ClosureProcess>();
1941 let external = flow.external::<()>();
1942
1943 let (in_port, input) =
1944 node.source_external_bincode::<_, i32, TotalOrder, ExactlyOnce>(&external);
1945 let out = input.send_bincode_external(&external);
1946
1947 let nodes = flow
1948 .with_process(&node, deployment.Localhost())
1949 .with_external(&external, deployment.Localhost())
1950 .deploy(&mut deployment);
1951
1952 deployment.deploy().await.unwrap();
1953
1954 let mut external_in = nodes.connect(in_port).await;
1955 let mut external_out = nodes.connect(out).await;
1956
1957 deployment.start().await.unwrap();
1958
1959 external_in.send(42).await.unwrap();
1960 assert_eq!(external_out.next().await.unwrap(), 42);
1961 }
1962}