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::ir::DebugInstantiate;
43use crate::compile::ir::{
44 ClusterMembersState, HydroIrOpMetadata, HydroNode, HydroRoot, HydroSource,
45};
46use crate::forward_handle::ForwardRef;
47#[cfg(stageleft_runtime)]
48use crate::forward_handle::{CycleCollection, ForwardHandle};
49use crate::live_collections::boundedness::{Bounded, Unbounded};
50use crate::live_collections::keyed_stream::KeyedStream;
51use crate::live_collections::singleton::Singleton;
52use crate::live_collections::stream::{ExactlyOnce, NoOrder, Stream, TotalOrder};
53#[cfg(feature = "tokio")]
54use crate::live_collections::stream::{Ordering, Retries};
55#[cfg(stageleft_runtime)]
56use crate::location::dynamic::DynLocation;
57use crate::location::dynamic::{ClusterConsistency, LocationId};
58#[cfg(feature = "tokio")]
59use crate::location::external_process::{
60 ExternalBincodeBidi, ExternalBincodeSink, ExternalBytesPort, Many, NotMany,
61};
62use crate::nondet::NonDet;
63#[cfg(feature = "tokio")]
64use crate::properties::manual_proof;
65#[cfg(feature = "sim")]
66use crate::sim::SimSender;
67use crate::staging_util::get_this_crate;
68
69pub mod dynamic;
70
71pub mod external_process;
72pub use external_process::External;
73
74pub mod process;
75pub use process::Process;
76
77pub mod cluster;
78pub use cluster::Cluster;
79
80pub mod member_id;
81pub use member_id::{MemberId, TaglessMemberId};
82
83pub mod tick;
84pub use tick::{Atomic, Tick};
85
86#[derive(PartialEq, Eq, Clone, Debug, Hash, Serialize, Deserialize)]
89pub enum MembershipEvent {
90 Joined,
92 Left,
94}
95
96#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
102pub enum NetworkHint {
103 Auto,
105 TcpPort(Option<u16>),
110}
111
112#[track_caller]
113pub(crate) fn check_matching_location<'a, L: Location<'a>>(l1: &L, l2: &L) {
114 assert_eq!(Location::id(l1), Location::id(l2), "locations do not match");
115}
116
117#[stageleft::export(LocationKey)]
118new_key_type! {
119 pub struct LocationKey;
121}
122
123impl std::fmt::Display for LocationKey {
124 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
125 write!(f, "loc{:?}", self.data()) }
127}
128
129impl std::str::FromStr for LocationKey {
132 type Err = Option<ParseIntError>;
133
134 fn from_str(s: &str) -> Result<Self, Self::Err> {
135 let nvn = s.strip_prefix("loc").ok_or(None)?;
136 let (idx, ver) = nvn.split_once("v").ok_or(None)?;
137 let idx: u64 = idx.parse()?;
138 let ver: u64 = ver.parse()?;
139 Ok(slotmap::KeyData::from_ffi((ver << 32) | idx).into())
140 }
141}
142
143impl LocationKey {
144 pub const FIRST: Self = Self(slotmap::KeyData::from_ffi(0x0000000100000001)); #[cfg(test)]
150 pub const TEST_KEY_1: Self = Self(slotmap::KeyData::from_ffi(0x000000FF00000001)); #[cfg(test)]
154 pub const TEST_KEY_2: Self = Self(slotmap::KeyData::from_ffi(0x000000FF00000002)); }
156
157impl<Ctx> FreeVariableWithContextWithProps<Ctx, ()> for LocationKey {
159 type O = LocationKey;
160
161 fn to_tokens(self, _ctx: &Ctx) -> (QuoteTokens, ())
162 where
163 Self: Sized,
164 {
165 let root = get_this_crate();
166 let n = Key::data(&self).as_ffi();
167 (
168 QuoteTokens {
169 prelude: None,
170 expr: Some(quote! {
171 #root::location::LocationKey::from(#root::runtime_support::slotmap::KeyData::from_ffi(#n))
172 }),
173 },
174 (),
175 )
176 }
177}
178
179#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Serialize)]
181pub enum LocationType {
182 Process,
184 Cluster,
186 External,
188}
189
190pub trait TopLevel<'a>: Location<'a> {}
192
193#[expect(
207 private_bounds,
208 reason = "only internal Hydro code can define location types"
209)]
210pub trait Location<'a>: DynLocation {
211 type Root: Location<'a>;
216
217 type DropConsistency: Location<'a, DropConsistency = Self::DropConsistency>;
219
220 fn root(&self) -> Self::Root;
225
226 fn drop_consistency(&self) -> Self::DropConsistency;
228 fn consistency() -> Option<ClusterConsistency>;
230
231 fn with_consistency_of<L2: Location<'a, DropConsistency = Self::DropConsistency>>(&self) -> L2 {
233 L2::from_drop_consistency(self.drop_consistency())
234 }
235
236 #[doc(hidden)]
237 fn from_drop_consistency(l2: Self::DropConsistency) -> Self;
238
239 fn try_tick(&self) -> Option<Tick<Self>> {
246 if Self::is_top_level() {
247 let id = if let LocationId::Atomic { .. } = self.id() {
248 None
249 } else {
250 Some(self.flow_state().borrow_mut().next_clock_id())
251 };
252 Some(Tick {
253 id,
254 l: self.clone(),
255 })
256 } else {
257 None
258 }
259 }
260
261 fn id(&self) -> LocationId {
263 DynLocation::dyn_id(self)
264 }
265
266 fn tick(&self) -> Tick<Self> {
292 self.try_tick().expect("cannot create nested ticks")
293 }
294
295 fn spin(&self) -> Stream<(), Self, Unbounded, TotalOrder, ExactlyOnce>
320 where
321 Self: TopLevel<'a> + Sized,
322 {
323 Stream::new(
324 self.clone(),
325 HydroNode::Source {
326 source: HydroSource::Spin(),
327 metadata: self.new_node_metadata(Stream::<
328 (),
329 Self,
330 Unbounded,
331 TotalOrder,
332 ExactlyOnce,
333 >::collection_kind()),
334 },
335 )
336 }
337
338 fn source_stream<T, E>(
359 &self,
360 e: impl QuotedWithContext<'a, E, Self>,
361 ) -> Stream<T, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>
362 where
363 E: FuturesStream<Item = T> + Unpin,
364 Self: TopLevel<'a> + Sized,
365 {
366 let e = e.splice_untyped_ctx(self);
367
368 let target_location = self.drop_consistency();
369 Stream::new(
370 target_location.clone(),
371 HydroNode::Source {
372 source: HydroSource::Stream(e.into()),
373 metadata: target_location.new_node_metadata(Stream::<
374 T,
375 Self::DropConsistency,
376 Unbounded,
377 TotalOrder,
378 ExactlyOnce,
379 >::collection_kind()),
380 },
381 )
382 }
383
384 fn source_iter<T, E>(
406 &self,
407 e: impl QuotedWithContext<'a, E, Self>,
408 ) -> Stream<T, Self::DropConsistency, Bounded, TotalOrder, ExactlyOnce>
409 where
410 E: IntoIterator<Item = T>,
411 Self: Sized,
412 {
413 let e = e.splice_typed_ctx(self);
414
415 let target_location = self.drop_consistency();
416 Stream::new(
417 target_location.clone(),
418 HydroNode::Source {
419 source: HydroSource::Iter(e.into()),
420 metadata: target_location.new_node_metadata(Stream::<
421 T,
422 Self::DropConsistency,
423 Bounded,
424 TotalOrder,
425 ExactlyOnce,
426 >::collection_kind()),
427 },
428 )
429 }
430
431 #[deprecated(note = "use .source_cluster_membership_stream(...) instead")]
432 fn source_cluster_members<C: 'a>(
471 &self,
472 cluster: &Cluster<'a, C>,
473 nondet_start: NonDet,
474 ) -> KeyedStream<MemberId<C>, MembershipEvent, Self::DropConsistency, Unbounded>
475 where
476 Self: TopLevel<'a> + Sized,
477 {
478 self.source_cluster_membership_stream(cluster, nondet_start)
479 }
480
481 fn source_cluster_membership_stream<C: 'a>(
520 &self,
521 cluster: &Cluster<'a, C>,
522 _nondet_start: NonDet,
523 ) -> KeyedStream<MemberId<C>, MembershipEvent, Self::DropConsistency, Unbounded>
524 where
525 Self: TopLevel<'a> + Sized,
526 {
527 let target_consistency = self.drop_consistency();
528 Stream::new(
529 target_consistency.clone(),
530 HydroNode::Source {
531 source: HydroSource::ClusterMembers(cluster.id(), ClusterMembersState::Uninit),
532 metadata: target_consistency.new_node_metadata(Stream::<
533 (TaglessMemberId, MembershipEvent),
534 Self,
535 Unbounded,
536 TotalOrder,
537 ExactlyOnce,
538 >::collection_kind(
539 )),
540 },
541 )
542 .map(q!(|(k, v)| (MemberId::from_tagless(k), v)))
543 .into_keyed()
544 }
545
546 #[cfg(feature = "tokio")]
554 fn source_external_bytes<L>(
555 &self,
556 from: &External<'_, L>,
557 ) -> (
558 ExternalBytesPort,
559 Stream<BytesMut, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
560 )
561 where
562 Self: TopLevel<'a> + Sized,
563 {
564 let (port, stream, sink) =
565 self.bind_single_client::<_, Bytes, LengthDelimitedCodec>(from, NetworkHint::Auto);
566
567 sink.complete(stream.location().source_iter(q!([])));
568
569 (port, stream)
570 }
571
572 #[cfg(feature = "tokio")]
579 fn source_external_bincode<L, T, O: Ordering, R: Retries>(
580 &self,
581 from: &External<'_, L>,
582 ) -> (
583 ExternalBincodeSink<T, NotMany, O, R>,
584 Stream<T, Self::DropConsistency, Unbounded, O, R>,
585 )
586 where
587 Self: TopLevel<'a> + Sized,
588 T: Serialize + DeserializeOwned,
589 {
590 let (port, stream, sink) = self.bind_single_client_bincode::<_, T, ()>(from);
591 sink.complete(stream.location().source_iter(q!([])));
592
593 (
594 ExternalBincodeSink {
595 process_key: from.key,
596 port_id: port.port_id,
597 _phantom: PhantomData,
598 },
599 stream.weaken_ordering().weaken_retries(),
600 )
601 }
602
603 #[cfg(feature = "sim")]
608 fn sim_input<T, O: Ordering, R: Retries>(
609 &self,
610 ) -> (
611 SimSender<T, O, R>,
612 Stream<T, Self::DropConsistency, Unbounded, O, R>,
613 )
614 where
615 Self: TopLevel<'a> + Sized,
616 T: Serialize + DeserializeOwned,
617 {
618 let external_location: External<'a, ()> = External {
619 key: LocationKey::FIRST,
620 flow_state: self.flow_state().clone(),
621 _phantom: PhantomData,
622 };
623
624 let (external, stream) = self.source_external_bincode(&external_location);
625
626 (SimSender(external.port_id, PhantomData), stream)
627 }
628
629 fn embedded_input<T>(
635 &self,
636 name: impl Into<String>,
637 ) -> Stream<T, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>
638 where
639 Self: TopLevel<'a> + Sized,
640 {
641 let ident = syn::Ident::new(&name.into(), Span::call_site());
642
643 let target_location = self.drop_consistency();
644 Stream::new(
645 target_location.clone(),
646 HydroNode::Source {
647 source: HydroSource::Embedded(ident),
648 metadata: target_location.new_node_metadata(Stream::<
649 T,
650 Self,
651 Unbounded,
652 TotalOrder,
653 ExactlyOnce,
654 >::collection_kind()),
655 },
656 )
657 }
658
659 fn embedded_singleton_input<T>(
665 &self,
666 name: impl Into<String>,
667 ) -> Singleton<T, Self::DropConsistency, Bounded>
668 where
669 Self: TopLevel<'a> + Sized,
670 {
671 let ident = syn::Ident::new(&name.into(), Span::call_site());
672
673 let target_location = self.drop_consistency();
674 Singleton::new(
675 target_location.clone(),
676 HydroNode::Source {
677 source: HydroSource::EmbeddedSingleton(ident),
678 metadata: target_location
679 .new_node_metadata(Singleton::<T, Self, Bounded>::collection_kind()),
680 },
681 )
682 }
683
684 #[cfg(feature = "tokio")]
729 #[expect(clippy::type_complexity, reason = "stream markers")]
730 fn bind_single_client<L, T, Codec: Encoder<T> + Decoder>(
731 &self,
732 from: &External<'_, L>,
733 port_hint: NetworkHint,
734 ) -> (
735 ExternalBytesPort<NotMany>,
736 Stream<<Codec as Decoder>::Item, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
737 ForwardHandle<'a, Stream<T, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>>,
738 )
739 where
740 Self: TopLevel<'a> + Sized,
741 {
742 let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
743 let target_consistency = self.drop_consistency();
744
745 let (fwd_ref, to_sink) = target_consistency.forward_ref::<Stream<
746 T,
747 Self::DropConsistency,
748 Unbounded,
749 TotalOrder,
750 ExactlyOnce,
751 >>();
752 let mut flow_state_borrow = self.flow_state().borrow_mut();
753
754 flow_state_borrow.push_root(HydroRoot::SendExternal {
755 to_external_key: from.key,
756 to_port_id: next_external_port_id,
757 to_many: false,
758 unpaired: false,
759 serialize_fn: None,
760 instantiate_fn: DebugInstantiate::Building,
761 input: Box::new(to_sink.ir_node.replace(HydroNode::Placeholder)),
762 op_metadata: HydroIrOpMetadata::new(),
763 });
764 drop(flow_state_borrow);
765
766 let raw_stream: Stream<
767 Result<<Codec as Decoder>::Item, <Codec as Decoder>::Error>,
768 Self::DropConsistency,
769 Unbounded,
770 TotalOrder,
771 ExactlyOnce,
772 > = Stream::new(
773 target_consistency.clone(),
774 HydroNode::ExternalInput {
775 from_external_key: from.key,
776 from_port_id: next_external_port_id,
777 from_many: false,
778 codec_type: quote_type::<Codec>().into(),
779 port_hint,
780 instantiate_fn: DebugInstantiate::Building,
781 deserialize_fn: None,
782 metadata: target_consistency.new_node_metadata(Stream::<
783 Result<<Codec as Decoder>::Item, <Codec as Decoder>::Error>,
784 Self::DropConsistency,
785 Unbounded,
786 TotalOrder,
787 ExactlyOnce,
788 >::collection_kind(
789 )),
790 },
791 );
792
793 (
794 ExternalBytesPort {
795 process_key: from.key,
796 port_id: next_external_port_id,
797 _phantom: PhantomData,
798 },
799 raw_stream.flatten_ordered(),
800 fwd_ref,
801 )
802 }
803
804 #[cfg(feature = "tokio")]
814 #[expect(clippy::type_complexity, reason = "stream markers")]
815 fn bind_single_client_bincode<L, InT: DeserializeOwned, OutT: Serialize>(
816 &self,
817 from: &External<'_, L>,
818 ) -> (
819 ExternalBincodeBidi<InT, OutT, NotMany>,
820 Stream<InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
821 ForwardHandle<'a, Stream<OutT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>>,
822 )
823 where
824 Self: TopLevel<'a> + Sized,
825 {
826 let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
827
828 let target_consistency = self.drop_consistency();
829 let (fwd_ref, to_sink) = target_consistency.forward_ref::<Stream<
830 OutT,
831 Self::DropConsistency,
832 Unbounded,
833 TotalOrder,
834 ExactlyOnce,
835 >>();
836 let mut flow_state_borrow = self.flow_state().borrow_mut();
837
838 let root = get_this_crate();
839
840 let out_t_type = quote_type::<OutT>();
841 let ser_fn: syn::Expr = syn::parse_quote! {
842 #root::runtime_support::stageleft::runtime_support::fn1_type_hint::<#out_t_type, _>(
843 |b| #root::runtime_support::bincode::serialize(&b).unwrap().into()
844 )
845 };
846
847 flow_state_borrow.push_root(HydroRoot::SendExternal {
848 to_external_key: from.key,
849 to_port_id: next_external_port_id,
850 to_many: false,
851 unpaired: false,
852 serialize_fn: Some(ser_fn.into()),
853 instantiate_fn: DebugInstantiate::Building,
854 input: Box::new(to_sink.ir_node.replace(HydroNode::Placeholder)),
855 op_metadata: HydroIrOpMetadata::new(),
856 });
857 drop(flow_state_borrow);
858
859 let in_t_type = quote_type::<InT>();
860
861 let deser_fn: syn::Expr = syn::parse_quote! {
862 |res| {
863 let b = res.unwrap();
864 #root::runtime_support::bincode::deserialize::<#in_t_type>(&b).unwrap()
865 }
866 };
867
868 let raw_stream: Stream<InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce> =
869 Stream::new(
870 target_consistency.clone(),
871 HydroNode::ExternalInput {
872 from_external_key: from.key,
873 from_port_id: next_external_port_id,
874 from_many: false,
875 codec_type: quote_type::<LengthDelimitedCodec>().into(),
876 port_hint: NetworkHint::Auto,
877 instantiate_fn: DebugInstantiate::Building,
878 deserialize_fn: Some(deser_fn.into()),
879 metadata: target_consistency.new_node_metadata(Stream::<
880 InT,
881 Self::DropConsistency,
882 Unbounded,
883 TotalOrder,
884 ExactlyOnce,
885 >::collection_kind(
886 )),
887 },
888 );
889
890 (
891 ExternalBincodeBidi {
892 process_key: from.key,
893 port_id: next_external_port_id,
894 _phantom: PhantomData,
895 },
896 raw_stream,
897 fwd_ref,
898 )
899 }
900
901 #[cfg(feature = "tokio")]
913 #[expect(clippy::type_complexity, reason = "stream markers")]
914 fn bidi_external_many_bytes<L, T, Codec: Encoder<T> + Decoder>(
915 &self,
916 from: &External<'_, L>,
917 port_hint: NetworkHint,
918 ) -> (
919 ExternalBytesPort<Many>,
920 KeyedStream<
921 u64,
922 <Codec as Decoder>::Item,
923 Self::DropConsistency,
924 Unbounded,
925 TotalOrder,
926 ExactlyOnce,
927 >,
928 KeyedStream<
929 u64,
930 MembershipEvent,
931 Self::DropConsistency,
932 Unbounded,
933 TotalOrder,
934 ExactlyOnce,
935 >,
936 ForwardHandle<
937 'a,
938 KeyedStream<u64, T, Self::DropConsistency, Unbounded, NoOrder, ExactlyOnce>,
939 >,
940 )
941 where
942 Self: TopLevel<'a> + Sized,
943 {
944 let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
945
946 let target_consistency = self.drop_consistency();
947 let (fwd_ref, to_sink) = target_consistency.forward_ref::<KeyedStream<
948 u64,
949 T,
950 Self::DropConsistency,
951 Unbounded,
952 NoOrder,
953 ExactlyOnce,
954 >>();
955 let to_sink_input = Box::new(to_sink.entries().ir_node.replace(HydroNode::Placeholder));
956 let mut flow_state_borrow = self.flow_state().borrow_mut();
957
958 flow_state_borrow.push_root(HydroRoot::SendExternal {
959 to_external_key: from.key,
960 to_port_id: next_external_port_id,
961 to_many: true,
962 unpaired: false,
963 serialize_fn: None,
964 instantiate_fn: DebugInstantiate::Building,
965 input: to_sink_input,
966 op_metadata: HydroIrOpMetadata::new(),
967 });
968 drop(flow_state_borrow);
969
970 let raw_stream: Stream<
971 Result<(u64, <Codec as Decoder>::Item), <Codec as Decoder>::Error>,
972 Self::DropConsistency,
973 Unbounded,
974 TotalOrder,
975 ExactlyOnce,
976 > = Stream::new(
977 target_consistency.clone(),
978 HydroNode::ExternalInput {
979 from_external_key: from.key,
980 from_port_id: next_external_port_id,
981 from_many: true,
982 codec_type: quote_type::<Codec>().into(),
983 port_hint,
984 instantiate_fn: DebugInstantiate::Building,
985 deserialize_fn: None,
986 metadata: target_consistency.new_node_metadata(Stream::<
987 Result<(u64, <Codec as Decoder>::Item), <Codec as Decoder>::Error>,
988 Self::DropConsistency,
989 Unbounded,
990 TotalOrder,
991 ExactlyOnce,
992 >::collection_kind(
993 )),
994 },
995 );
996
997 let membership_stream_ident = syn::Ident::new(
998 &format!(
999 "__hydro_deploy_many_{}_{}_membership",
1000 from.key, next_external_port_id
1001 ),
1002 Span::call_site(),
1003 );
1004 let membership_stream_expr: syn::Expr = parse_quote!(#membership_stream_ident);
1005 let raw_membership_stream: KeyedStream<
1006 u64,
1007 bool,
1008 Self::DropConsistency,
1009 Unbounded,
1010 TotalOrder,
1011 ExactlyOnce,
1012 > = KeyedStream::new(
1013 target_consistency.clone(),
1014 HydroNode::Source {
1015 source: HydroSource::Stream(membership_stream_expr.into()),
1016 metadata: target_consistency.new_node_metadata(KeyedStream::<
1017 u64,
1018 bool,
1019 Self::DropConsistency,
1020 Unbounded,
1021 TotalOrder,
1022 ExactlyOnce,
1023 >::collection_kind(
1024 )),
1025 },
1026 );
1027
1028 (
1029 ExternalBytesPort {
1030 process_key: from.key,
1031 port_id: next_external_port_id,
1032 _phantom: PhantomData,
1033 },
1034 raw_stream
1035 .flatten_ordered() .into_keyed(),
1037 raw_membership_stream.map(q!(|join| {
1038 if join {
1039 MembershipEvent::Joined
1040 } else {
1041 MembershipEvent::Left
1042 }
1043 })),
1044 fwd_ref,
1045 )
1046 }
1047
1048 #[cfg(feature = "tokio")]
1064 #[expect(clippy::type_complexity, reason = "stream markers")]
1065 fn bidi_external_many_bincode<L, InT: DeserializeOwned, OutT: Serialize>(
1066 &self,
1067 from: &External<'_, L>,
1068 ) -> (
1069 ExternalBincodeBidi<InT, OutT, Many>,
1070 KeyedStream<u64, InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
1071 KeyedStream<
1072 u64,
1073 MembershipEvent,
1074 Self::DropConsistency,
1075 Unbounded,
1076 TotalOrder,
1077 ExactlyOnce,
1078 >,
1079 ForwardHandle<
1080 'a,
1081 KeyedStream<u64, OutT, Self::DropConsistency, Unbounded, NoOrder, ExactlyOnce>,
1082 >,
1083 )
1084 where
1085 Self: TopLevel<'a> + Sized,
1086 {
1087 let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
1088
1089 let target_consistency = self.drop_consistency();
1090 let (fwd_ref, to_sink) = target_consistency.forward_ref::<KeyedStream<
1091 u64,
1092 OutT,
1093 Self::DropConsistency,
1094 Unbounded,
1095 NoOrder,
1096 ExactlyOnce,
1097 >>();
1098 let to_sink_input = Box::new(to_sink.entries().ir_node.replace(HydroNode::Placeholder));
1099 let mut flow_state_borrow = self.flow_state().borrow_mut();
1100
1101 let root = get_this_crate();
1102
1103 let out_t_type = quote_type::<OutT>();
1104 let ser_fn: syn::Expr = syn::parse_quote! {
1105 #root::runtime_support::stageleft::runtime_support::fn1_type_hint::<(u64, #out_t_type), _>(
1106 |(id, b)| (id, #root::runtime_support::bincode::serialize(&b).unwrap().into())
1107 )
1108 };
1109
1110 flow_state_borrow.push_root(HydroRoot::SendExternal {
1111 to_external_key: from.key,
1112 to_port_id: next_external_port_id,
1113 to_many: true,
1114 unpaired: false,
1115 serialize_fn: Some(ser_fn.into()),
1116 instantiate_fn: DebugInstantiate::Building,
1117 input: to_sink_input,
1118 op_metadata: HydroIrOpMetadata::new(),
1119 });
1120 drop(flow_state_borrow);
1121
1122 let in_t_type = quote_type::<InT>();
1123
1124 let deser_fn: syn::Expr = syn::parse_quote! {
1125 |res| {
1126 let (id, b) = res.unwrap();
1127 (id, #root::runtime_support::bincode::deserialize::<#in_t_type>(&b).unwrap())
1128 }
1129 };
1130
1131 let raw_stream: KeyedStream<
1132 u64,
1133 InT,
1134 Self::DropConsistency,
1135 Unbounded,
1136 TotalOrder,
1137 ExactlyOnce,
1138 > = KeyedStream::new(
1139 target_consistency.clone(),
1140 HydroNode::ExternalInput {
1141 from_external_key: from.key,
1142 from_port_id: next_external_port_id,
1143 from_many: true,
1144 codec_type: quote_type::<LengthDelimitedCodec>().into(),
1145 port_hint: NetworkHint::Auto,
1146 instantiate_fn: DebugInstantiate::Building,
1147 deserialize_fn: Some(deser_fn.into()),
1148 metadata: target_consistency.new_node_metadata(KeyedStream::<
1149 u64,
1150 InT,
1151 Self::DropConsistency,
1152 Unbounded,
1153 TotalOrder,
1154 ExactlyOnce,
1155 >::collection_kind(
1156 )),
1157 },
1158 );
1159
1160 let membership_stream_ident = syn::Ident::new(
1161 &format!(
1162 "__hydro_deploy_many_{}_{}_membership",
1163 from.key, next_external_port_id
1164 ),
1165 Span::call_site(),
1166 );
1167 let membership_stream_expr: syn::Expr = parse_quote!(#membership_stream_ident);
1168 let raw_membership_stream: KeyedStream<
1169 u64,
1170 bool,
1171 Self::DropConsistency,
1172 Unbounded,
1173 TotalOrder,
1174 ExactlyOnce,
1175 > = KeyedStream::new(
1176 target_consistency.clone(),
1177 HydroNode::Source {
1178 source: HydroSource::Stream(membership_stream_expr.into()),
1179 metadata: target_consistency.new_node_metadata(KeyedStream::<
1180 u64,
1181 bool,
1182 Self::DropConsistency,
1183 Unbounded,
1184 TotalOrder,
1185 ExactlyOnce,
1186 >::collection_kind(
1187 )),
1188 },
1189 );
1190
1191 (
1192 ExternalBincodeBidi {
1193 process_key: from.key,
1194 port_id: next_external_port_id,
1195 _phantom: PhantomData,
1196 },
1197 raw_stream,
1198 raw_membership_stream.map(q!(|join| {
1199 if join {
1200 MembershipEvent::Joined
1201 } else {
1202 MembershipEvent::Left
1203 }
1204 })),
1205 fwd_ref,
1206 )
1207 }
1208
1209 fn sidecar_bidi<InT: 'static, OutT: 'static, F>(
1262 &self,
1263 sidecar: impl QuotedWithContext<'a, F, Self>,
1264 ) -> (
1265 Stream<InT, Self, Unbounded, TotalOrder, ExactlyOnce>,
1266 ForwardHandle<'a, Stream<OutT, Self, Unbounded, NoOrder, ExactlyOnce>>,
1267 )
1268 where
1269 Self: Sized + TopLevel<'a>,
1270 {
1271 let location_key = Location::id(self).key();
1272
1273 let sidecar_id = self.flow_state().borrow_mut().next_sidecar_id();
1274 let (stream_ident, sink_ident) = sidecar_id.idents();
1275
1276 let sidecar_closure: syn::Expr = sidecar.splice_untyped_ctx(self);
1277 self.flow_state()
1278 .borrow_mut()
1279 .sidecars
1280 .push(crate::compile::builder::Sidecar::Bidi {
1281 location_key,
1282 sidecar_id,
1283 sidecar_closure: Box::new(sidecar_closure),
1284 });
1285
1286 let source_expr: syn::Expr = parse_quote! {
1288 #stream_ident
1289 };
1290 let inbound: Stream<InT, Self, Unbounded, TotalOrder, ExactlyOnce> = Stream::new(
1291 self.clone(),
1292 HydroNode::Source {
1293 source: HydroSource::Stream(source_expr.into()),
1294 metadata: self.new_node_metadata(Stream::<
1295 InT,
1296 Self,
1297 Unbounded, TotalOrder, ExactlyOnce,
1300 >::collection_kind()),
1301 },
1302 );
1303
1304 let (fwd_ref, to_sink): (
1306 ForwardHandle<'a, Stream<OutT, Self, Unbounded, NoOrder, ExactlyOnce>>,
1307 Stream<OutT, Self, Unbounded, NoOrder, ExactlyOnce>,
1308 ) = self.forward_ref();
1309
1310 let sink_expr: syn::Expr = parse_quote! {
1311 #sink_ident
1312 };
1313
1314 let sink_input_ir = to_sink.ir_node.replace(HydroNode::Placeholder);
1315 self.flow_state()
1316 .borrow_mut()
1317 .try_push_root(HydroRoot::DestSink {
1318 sink: sink_expr.into(),
1319 input: Box::new(sink_input_ir),
1320 op_metadata: HydroIrOpMetadata::new(),
1321 });
1322
1323 (inbound, fwd_ref)
1324 }
1325
1326 fn singleton<T>(
1346 &self,
1347 e: impl QuotedWithContext<'a, T, Self>,
1348 ) -> Singleton<T, Self::DropConsistency, Bounded>
1349 where
1350 Self: Sized,
1351 {
1352 let e = e.splice_untyped_ctx(self);
1353
1354 let target_location = self.drop_consistency();
1355 Singleton::new(
1356 target_location.clone(),
1357 HydroNode::SingletonSource {
1358 value: e.into(),
1359 first_tick_only: false,
1360 metadata: target_location.new_node_metadata(Singleton::<
1361 T,
1362 Self::DropConsistency,
1363 Bounded,
1364 >::collection_kind()),
1365 },
1366 )
1367 }
1368
1369 fn singleton_future<F>(
1392 &self,
1393 e: impl QuotedWithContext<'a, F, Self>,
1394 ) -> Singleton<F::Output, Self::DropConsistency, Bounded>
1395 where
1396 F: Future,
1397 Self: Sized,
1398 {
1399 self.singleton(e).resolve_future_blocking()
1400 }
1401
1402 #[cfg(feature = "tokio")]
1411 fn source_interval(
1412 &self,
1413 interval: impl QuotedWithContext<'a, Duration, Self> + Copy + 'a,
1414 ) -> Stream<(), Self, Unbounded, TotalOrder, ExactlyOnce>
1415 where
1416 Self: TopLevel<'a> + Sized,
1417 {
1418 self.source_stream(q!(tokio_stream::StreamExt::map(
1419 tokio_stream::wrappers::IntervalStream::new(tokio::time::interval(interval)),
1420 |_| ()
1421 )))
1422 .assert_has_consistency_of_trusted(
1423 manual_proof!(),
1424 )
1425 }
1426
1427 #[cfg(feature = "tokio")]
1434 fn source_interval_delayed(
1435 &self,
1436 delay: impl QuotedWithContext<'a, Duration, Self> + Copy + 'a,
1437 interval: impl QuotedWithContext<'a, Duration, Self> + Copy + 'a,
1438 ) -> Stream<(), Self, Unbounded, TotalOrder, ExactlyOnce>
1439 where
1440 Self: TopLevel<'a> + Sized,
1441 {
1442 self.source_stream(q!(tokio_stream::StreamExt::map(
1443 tokio_stream::wrappers::IntervalStream::new(tokio::time::interval_at(
1444 tokio::time::Instant::now() + delay,
1445 interval,
1446 )),
1447 |_| ()
1448 )))
1449 .assert_has_consistency_of_trusted(
1450 manual_proof!(),
1451 )
1452 }
1453
1454 fn forward_ref<S>(&self) -> (ForwardHandle<'a, S>, S)
1494 where
1495 S: CycleCollection<'a, ForwardRef, Location = Self>,
1496 {
1497 let cycle_id = self.flow_state().borrow_mut().next_cycle_id();
1498 (
1499 ForwardHandle::new(cycle_id, Location::id(self)),
1500 S::create_source(cycle_id, self.clone()),
1501 )
1502 }
1503}
1504
1505#[cfg(feature = "deploy")]
1506#[cfg(test)]
1507mod tests {
1508 use std::collections::HashSet;
1509
1510 use futures::{SinkExt, StreamExt};
1511 use hydro_deploy::Deployment;
1512 use stageleft::q;
1513 use tokio_util::codec::LengthDelimitedCodec;
1514
1515 use crate::compile::builder::FlowBuilder;
1516 use crate::live_collections::stream::{ExactlyOnce, TotalOrder};
1517 use crate::location::{Location, NetworkHint};
1518 use crate::nondet::nondet;
1519
1520 #[tokio::test]
1521 async fn top_level_singleton_replay_cardinality() {
1522 let mut deployment = Deployment::new();
1523
1524 let mut flow = FlowBuilder::new();
1525 let node = flow.process::<()>();
1526 let external = flow.external::<()>();
1527
1528 let (in_port, input) =
1529 node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
1530 let singleton = node.singleton(q!(123));
1531 let tick = node.tick();
1532 let out = input
1533 .batch(&tick, nondet!())
1534 .cross_singleton(singleton.clone().snapshot(&tick, nondet!()))
1535 .cross_singleton(
1536 singleton
1537 .snapshot(&tick, nondet!())
1538 .into_stream()
1539 .count(),
1540 )
1541 .all_ticks()
1542 .send_bincode_external(&external);
1543
1544 let nodes = flow
1545 .with_process(&node, deployment.Localhost())
1546 .with_external(&external, deployment.Localhost())
1547 .deploy(&mut deployment);
1548
1549 deployment.deploy().await.unwrap();
1550
1551 let mut external_in = nodes.connect(in_port).await;
1552 let mut external_out = nodes.connect(out).await;
1553
1554 deployment.start().await.unwrap();
1555
1556 external_in.send(1).await.unwrap();
1557 assert_eq!(external_out.next().await.unwrap(), ((1, 123), 1));
1558
1559 external_in.send(2).await.unwrap();
1560 assert_eq!(external_out.next().await.unwrap(), ((2, 123), 1));
1561 }
1562
1563 #[tokio::test]
1564 async fn tick_singleton_replay_cardinality() {
1565 let mut deployment = Deployment::new();
1566
1567 let mut flow = FlowBuilder::new();
1568 let node = flow.process::<()>();
1569 let external = flow.external::<()>();
1570
1571 let (in_port, input) =
1572 node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
1573 let tick = node.tick();
1574 let singleton = tick.singleton(q!(123));
1575 let out = input
1576 .batch(&tick, nondet!())
1577 .cross_singleton(singleton.clone())
1578 .cross_singleton(singleton.into_stream().count())
1579 .all_ticks()
1580 .send_bincode_external(&external);
1581
1582 let nodes = flow
1583 .with_process(&node, deployment.Localhost())
1584 .with_external(&external, deployment.Localhost())
1585 .deploy(&mut deployment);
1586
1587 deployment.deploy().await.unwrap();
1588
1589 let mut external_in = nodes.connect(in_port).await;
1590 let mut external_out = nodes.connect(out).await;
1591
1592 deployment.start().await.unwrap();
1593
1594 external_in.send(1).await.unwrap();
1595 assert_eq!(external_out.next().await.unwrap(), ((1, 123), 1));
1596
1597 external_in.send(2).await.unwrap();
1598 assert_eq!(external_out.next().await.unwrap(), ((2, 123), 1));
1599 }
1600
1601 #[tokio::test]
1602 async fn external_bytes() {
1603 let mut deployment = Deployment::new();
1604
1605 let mut flow = FlowBuilder::new();
1606 let first_node = flow.process::<()>();
1607 let external = flow.external::<()>();
1608
1609 let (in_port, input) = first_node.source_external_bytes(&external);
1610 let out = input.send_bincode_external(&external);
1611
1612 let nodes = flow
1613 .with_process(&first_node, deployment.Localhost())
1614 .with_external(&external, deployment.Localhost())
1615 .deploy(&mut deployment);
1616
1617 deployment.deploy().await.unwrap();
1618
1619 let mut external_in = nodes.connect(in_port).await.1;
1620 let mut external_out = nodes.connect(out).await;
1621
1622 deployment.start().await.unwrap();
1623
1624 external_in.send(vec![1, 2, 3].into()).await.unwrap();
1625
1626 assert_eq!(external_out.next().await.unwrap(), vec![1, 2, 3]);
1627 }
1628
1629 #[tokio::test]
1630 async fn multi_external_source() {
1631 let mut deployment = Deployment::new();
1632
1633 let mut flow = FlowBuilder::new();
1634 let first_node = flow.process::<()>();
1635 let external = flow.external::<()>();
1636
1637 let (in_port, input, _membership, complete_sink) =
1638 first_node.bidi_external_many_bincode(&external);
1639 let out = input.entries().send_bincode_external(&external);
1640 complete_sink.complete(
1641 first_node
1642 .source_iter::<(u64, ()), _>(q!([]))
1643 .into_keyed()
1644 .weaken_ordering(),
1645 );
1646
1647 let nodes = flow
1648 .with_process(&first_node, deployment.Localhost())
1649 .with_external(&external, deployment.Localhost())
1650 .deploy(&mut deployment);
1651
1652 deployment.deploy().await.unwrap();
1653
1654 let (_, mut external_in_1) = nodes.connect_bincode(in_port.clone()).await;
1655 let (_, mut external_in_2) = nodes.connect_bincode(in_port).await;
1656 let external_out = nodes.connect(out).await;
1657
1658 deployment.start().await.unwrap();
1659
1660 external_in_1.send(123).await.unwrap();
1661 external_in_2.send(456).await.unwrap();
1662
1663 assert_eq!(
1664 external_out.take(2).collect::<HashSet<_>>().await,
1665 vec![(0, 123), (1, 456)].into_iter().collect()
1666 );
1667 }
1668
1669 #[tokio::test]
1670 async fn second_connection_only_multi_source() {
1671 let mut deployment = Deployment::new();
1672
1673 let mut flow = FlowBuilder::new();
1674 let first_node = flow.process::<()>();
1675 let external = flow.external::<()>();
1676
1677 let (in_port, input, _membership, complete_sink) =
1678 first_node.bidi_external_many_bincode(&external);
1679 let out = input.entries().send_bincode_external(&external);
1680 complete_sink.complete(
1681 first_node
1682 .source_iter::<(u64, ()), _>(q!([]))
1683 .into_keyed()
1684 .weaken_ordering(),
1685 );
1686
1687 let nodes = flow
1688 .with_process(&first_node, deployment.Localhost())
1689 .with_external(&external, deployment.Localhost())
1690 .deploy(&mut deployment);
1691
1692 deployment.deploy().await.unwrap();
1693
1694 let (_, mut _external_in_1) = nodes.connect_bincode(in_port.clone()).await;
1696 let (_, mut external_in_2) = nodes.connect_bincode(in_port).await;
1697 let mut external_out = nodes.connect(out).await;
1698
1699 deployment.start().await.unwrap();
1700
1701 external_in_2.send(456).await.unwrap();
1702
1703 assert_eq!(external_out.next().await.unwrap(), (1, 456));
1704 }
1705
1706 #[tokio::test]
1707 async fn multi_external_bytes() {
1708 let mut deployment = Deployment::new();
1709
1710 let mut flow = FlowBuilder::new();
1711 let first_node = flow.process::<()>();
1712 let external = flow.external::<()>();
1713
1714 let (in_port, input, _membership, complete_sink) = first_node
1715 .bidi_external_many_bytes::<_, _, LengthDelimitedCodec>(&external, NetworkHint::Auto);
1716 let out = input.entries().send_bincode_external(&external);
1717 complete_sink.complete(
1718 first_node
1719 .source_iter(q!([]))
1720 .into_keyed()
1721 .weaken_ordering(),
1722 );
1723
1724 let nodes = flow
1725 .with_process(&first_node, deployment.Localhost())
1726 .with_external(&external, deployment.Localhost())
1727 .deploy(&mut deployment);
1728
1729 deployment.deploy().await.unwrap();
1730
1731 let mut external_in_1 = nodes.connect(in_port.clone()).await.1;
1732 let mut external_in_2 = nodes.connect(in_port).await.1;
1733 let external_out = nodes.connect(out).await;
1734
1735 deployment.start().await.unwrap();
1736
1737 external_in_1.send(vec![1, 2, 3].into()).await.unwrap();
1738 external_in_2.send(vec![4, 5].into()).await.unwrap();
1739
1740 assert_eq!(
1741 external_out.take(2).collect::<HashSet<_>>().await,
1742 vec![
1743 (0, (&[1u8, 2, 3] as &[u8]).into()),
1744 (1, (&[4u8, 5] as &[u8]).into())
1745 ]
1746 .into_iter()
1747 .collect()
1748 );
1749 }
1750
1751 #[tokio::test]
1752 async fn single_client_external_bytes() {
1753 let mut deployment = Deployment::new();
1754 let mut flow = FlowBuilder::new();
1755 let first_node = flow.process::<()>();
1756 let external = flow.external::<()>();
1757 let (port, input, complete_sink) = first_node
1758 .bind_single_client::<_, _, LengthDelimitedCodec>(&external, NetworkHint::Auto);
1759 complete_sink.complete(input.map(q!(|data| {
1760 let mut resp: Vec<u8> = data.into();
1761 resp.push(42);
1762 resp.into() })));
1764
1765 let nodes = flow
1766 .with_process(&first_node, deployment.Localhost())
1767 .with_external(&external, deployment.Localhost())
1768 .deploy(&mut deployment);
1769
1770 deployment.deploy().await.unwrap();
1771 deployment.start().await.unwrap();
1772
1773 let (mut external_out, mut external_in) = nodes.connect(port).await;
1774
1775 external_in.send(vec![1, 2, 3].into()).await.unwrap();
1776 assert_eq!(
1777 external_out.next().await.unwrap().unwrap(),
1778 vec![1, 2, 3, 42]
1779 );
1780 }
1781
1782 #[tokio::test]
1783 async fn echo_external_bytes() {
1784 let mut deployment = Deployment::new();
1785
1786 let mut flow = FlowBuilder::new();
1787 let first_node = flow.process::<()>();
1788 let external = flow.external::<()>();
1789
1790 let (port, input, _membership, complete_sink) = first_node
1791 .bidi_external_many_bytes::<_, _, LengthDelimitedCodec>(&external, NetworkHint::Auto);
1792 complete_sink
1793 .complete(input.map(q!(|bytes| { bytes.into_iter().map(|x| x + 1).collect() })));
1794
1795 let nodes = flow
1796 .with_process(&first_node, deployment.Localhost())
1797 .with_external(&external, deployment.Localhost())
1798 .deploy(&mut deployment);
1799
1800 deployment.deploy().await.unwrap();
1801
1802 let (mut external_out_1, mut external_in_1) = nodes.connect(port.clone()).await;
1803 let (mut external_out_2, mut external_in_2) = nodes.connect(port).await;
1804
1805 deployment.start().await.unwrap();
1806
1807 external_in_1.send(vec![1, 2, 3].into()).await.unwrap();
1808 external_in_2.send(vec![4, 5].into()).await.unwrap();
1809
1810 assert_eq!(external_out_1.next().await.unwrap().unwrap(), vec![2, 3, 4]);
1811 assert_eq!(external_out_2.next().await.unwrap().unwrap(), vec![5, 6]);
1812 }
1813
1814 #[tokio::test]
1815 async fn echo_external_bincode() {
1816 let mut deployment = Deployment::new();
1817
1818 let mut flow = FlowBuilder::new();
1819 let first_node = flow.process::<()>();
1820 let external = flow.external::<()>();
1821
1822 let (port, input, _membership, complete_sink) =
1823 first_node.bidi_external_many_bincode(&external);
1824 complete_sink.complete(input.map(q!(|text: String| { text.to_uppercase() })));
1825
1826 let nodes = flow
1827 .with_process(&first_node, deployment.Localhost())
1828 .with_external(&external, deployment.Localhost())
1829 .deploy(&mut deployment);
1830
1831 deployment.deploy().await.unwrap();
1832
1833 let (mut external_out_1, mut external_in_1) = nodes.connect_bincode(port.clone()).await;
1834 let (mut external_out_2, mut external_in_2) = nodes.connect_bincode(port).await;
1835
1836 deployment.start().await.unwrap();
1837
1838 external_in_1.send("hi".to_owned()).await.unwrap();
1839 external_in_2.send("hello".to_owned()).await.unwrap();
1840
1841 assert_eq!(external_out_1.next().await.unwrap(), "HI");
1842 assert_eq!(external_out_2.next().await.unwrap(), "HELLO");
1843 }
1844
1845 #[tokio::test]
1846 async fn closure_location_name() {
1847 let mut deployment = Deployment::new();
1848 let mut flow = FlowBuilder::new();
1849
1850 enum ClosureProcess {}
1851
1852 let node = flow.process::<ClosureProcess>();
1853 let external = flow.external::<()>();
1854
1855 let (in_port, input) =
1856 node.source_external_bincode::<_, i32, TotalOrder, ExactlyOnce>(&external);
1857 let out = input.send_bincode_external(&external);
1858
1859 let nodes = flow
1860 .with_process(&node, deployment.Localhost())
1861 .with_external(&external, deployment.Localhost())
1862 .deploy(&mut deployment);
1863
1864 deployment.deploy().await.unwrap();
1865
1866 let mut external_in = nodes.connect(in_port).await;
1867 let mut external_out = nodes.connect(out).await;
1868
1869 deployment.start().await.unwrap();
1870
1871 external_in.send(42).await.unwrap();
1872 assert_eq!(external_out.next().await.unwrap(), 42);
1873 }
1874}