Skip to main content

hydro_lang/location/
mod.rs

1//! Type definitions for distributed locations, which specify where pieces of a Hydro
2//! program will be executed.
3//!
4//! Hydro is a **global**, **distributed** programming model. This means that the data
5//! and computation in a Hydro program can be spread across multiple machines, data
6//! centers, and even continents. To achieve this, Hydro uses the concept of
7//! **locations** to keep track of _where_ data is located and computation is executed.
8//!
9//! Each live collection type (in [`crate::live_collections`]) has a type parameter `L`
10//! which will always be a type that implements the [`Location`] trait (e.g. [`Process`]
11//! and [`Cluster`]). To create distributed programs, Hydro provides a variety of APIs
12//! to allow live collections to be _moved_ between locations via network send/receive.
13//!
14//! See [the Hydro docs](https://hydro.run/docs/hydro/reference/locations/) for more information.
15
16use 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/// An event indicating a change in membership status of a location in a group
87/// (e.g. a node in a [`Cluster`] or an external client connection).
88#[derive(PartialEq, Eq, Clone, Debug, Hash, Serialize, Deserialize)]
89pub enum MembershipEvent {
90    /// The member has joined the group and is now active.
91    Joined,
92    /// The member has left the group and is no longer active.
93    Left,
94}
95
96/// A hint for configuring the network transport used by an external connection.
97///
98/// This controls how the underlying TCP listener is set up when binding
99/// external client connections via methods like [`Location::bind_single_client`]
100/// or [`Location::bidi_external_many_bytes`].
101#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
102pub enum NetworkHint {
103    /// Automatically select the network configuration (e.g. an ephemeral port).
104    Auto,
105    /// Use a TCP port, optionally specifying a fixed port number.
106    ///
107    /// If `None`, an available port will be chosen automatically.
108    /// If `Some(port)`, the given port number will be used.
109    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    /// A unique identifier for a clock tick.
120    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()) // `"loc1v1"``
126    }
127}
128
129/// This is used for the ECS membership stream.
130/// TODO(mingwei): Make this more robust?
131impl 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    /// TODO(minwgei): Remove this and avoid magic key for simulator external.
145    /// The first location key, used by the simulator as the default external location.
146    pub const FIRST: Self = Self(slotmap::KeyData::from_ffi(0x0000000100000001)); // `1v1`
147
148    /// A key for testing with index 1.
149    #[cfg(test)]
150    pub const TEST_KEY_1: Self = Self(slotmap::KeyData::from_ffi(0x000000FF00000001)); // `1v255`
151
152    /// A key for testing with index 2.
153    #[cfg(test)]
154    pub const TEST_KEY_2: Self = Self(slotmap::KeyData::from_ffi(0x000000FF00000002)); // `2v255`
155}
156
157/// This is used within `q!` code in docker and ECS.
158impl<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/// A simple enum for the type of a root location.
180#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Serialize)]
181pub enum LocationType {
182    /// A process (single node).
183    Process,
184    /// A cluster (multiple nodes).
185    Cluster,
186    /// An external client.
187    External,
188}
189
190/// A top-level location (i.e. a [`Process`] or [`Cluster`]) that is outside a tick / atomic region.
191pub trait TopLevel<'a>: Location<'a> {}
192
193/// A location where data can be materialized and computation can be executed.
194///
195/// Hydro is a **global**, **distributed** programming model. This means that the data
196/// and computation in a Hydro program can be spread across multiple machines, data
197/// centers, and even continents. To achieve this, Hydro uses the concept of
198/// **locations** to keep track of _where_ data is located and computation is executed.
199///
200/// Each live collection type (in [`crate::live_collections`]) has a type parameter `L`
201/// which will always be a type that implements the [`Location`] trait (e.g. [`Process`]
202/// and [`Cluster`]). To create distributed programs, Hydro provides a variety of APIs
203/// to allow live collections to be _moved_ between locations via network send/receive.
204///
205/// See [the Hydro docs](https://hydro.run/docs/hydro/reference/locations/) for more information.
206#[expect(
207    private_bounds,
208    reason = "only internal Hydro code can define location types"
209)]
210pub trait Location<'a>: DynLocation {
211    /// The root location type for this location.
212    ///
213    /// For top-level locations like [`Process`] and [`Cluster`], this is `Self`.
214    /// For nested locations like [`Tick`], this is the root location that contains it.
215    type Root: Location<'a>;
216
217    /// Location type with consistency guarantees dropped for the live collection on it.
218    type DropConsistency: Location<'a, DropConsistency = Self::DropConsistency>;
219
220    /// Returns the root location for this location.
221    ///
222    /// For top-level locations like [`Process`] and [`Cluster`], this returns `self`.
223    /// For nested locations like [`Tick`], this returns the root location that contains it.
224    fn root(&self) -> Self::Root;
225
226    /// This location but with consistency guarantees dropped for the live collection
227    fn drop_consistency(&self) -> Self::DropConsistency;
228    /// Gets the runtime enum variant for the current consistency level, if this is a cluster.
229    fn consistency() -> Option<ClusterConsistency>;
230
231    /// Updates the consistency guarantees to match that of the given location.
232    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    /// Attempts to create a new [`Tick`] clock domain at this location.
240    ///
241    /// Returns `Some(Tick)` if this is a top-level location (like [`Process`] or [`Cluster`]),
242    /// or `None` if this location is already inside a tick (nested ticks are not supported).
243    ///
244    /// Prefer using [`Location::tick`] when you know the location is top-level.
245    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    /// Returns the unique identifier for this location.
262    fn id(&self) -> LocationId {
263        DynLocation::dyn_id(self)
264    }
265
266    /// Creates a new [`Tick`] clock domain at this location.
267    ///
268    /// A tick represents a logical clock that can be used to batch streaming data
269    /// into discrete time steps. This is useful for implementing iterative algorithms
270    /// or for synchronizing data across multiple streams.
271    ///
272    /// # Example
273    /// ```rust
274    /// # #[cfg(feature = "deploy")] {
275    /// # use hydro_lang::prelude::*;
276    /// # use futures::StreamExt;
277    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
278    /// let tick = process.tick();
279    /// let inside_tick = process
280    ///     .source_iter(q!(vec![1, 2, 3, 4]))
281    ///     .batch(&tick, nondet!(/** test */));
282    /// inside_tick.all_ticks()
283    /// # }, |mut stream| async move {
284    /// // 1, 2, 3, 4
285    /// # for w in vec![1, 2, 3, 4] {
286    /// #     assert_eq!(stream.next().await.unwrap(), w);
287    /// # }
288    /// # }));
289    /// # }
290    /// ```
291    fn tick(&self) -> Tick<Self> {
292        self.try_tick().expect("cannot create nested ticks")
293    }
294
295    /// Creates an unbounded stream that continuously emits unit values `()`.
296    ///
297    /// This is useful for driving computations that need to run continuously,
298    /// such as polling or heartbeat mechanisms.
299    ///
300    /// # Example
301    /// ```rust
302    /// # #[cfg(feature = "deploy")] {
303    /// # use hydro_lang::prelude::*;
304    /// # use futures::StreamExt;
305    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
306    /// let tick = process.tick();
307    /// process.spin()
308    ///     .batch(&tick, nondet!(/** test */))
309    ///     .map(q!(|_| 42))
310    ///     .all_ticks()
311    /// # }, |mut stream| async move {
312    /// // 42, 42, 42, ...
313    /// # assert_eq!(stream.next().await.unwrap(), 42);
314    /// # assert_eq!(stream.next().await.unwrap(), 42);
315    /// # assert_eq!(stream.next().await.unwrap(), 42);
316    /// # }));
317    /// # }
318    /// ```
319    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    /// Creates a stream from an async [`FuturesStream`].
339    ///
340    /// This is useful for integrating with external async data sources,
341    /// such as network connections or file readers.
342    ///
343    /// # Example
344    /// ```rust
345    /// # #[cfg(feature = "deploy")] {
346    /// # use hydro_lang::prelude::*;
347    /// # use futures::StreamExt;
348    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
349    /// process.source_stream(q!(futures::stream::iter(vec![1, 2, 3])))
350    /// # }, |mut stream| async move {
351    /// // 1, 2, 3
352    /// # for w in vec![1, 2, 3] {
353    /// #     assert_eq!(stream.next().await.unwrap(), w);
354    /// # }
355    /// # }));
356    /// # }
357    /// ```
358    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    /// Creates a bounded stream from an iterator.
385    ///
386    /// The iterator is evaluated once at runtime, and all elements are emitted
387    /// in order. This is useful for creating streams from static data or
388    /// for testing.
389    ///
390    /// # Example
391    /// ```rust
392    /// # #[cfg(feature = "deploy")] {
393    /// # use hydro_lang::prelude::*;
394    /// # use futures::StreamExt;
395    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
396    /// process.source_iter(q!(vec![1, 2, 3, 4]))
397    /// # }, |mut stream| async move {
398    /// // 1, 2, 3, 4
399    /// # for w in vec![1, 2, 3, 4] {
400    /// #     assert_eq!(stream.next().await.unwrap(), w);
401    /// # }
402    /// # }));
403    /// # }
404    /// ```
405    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    /// Creates a stream of membership events for a cluster.
433    ///
434    /// This stream emits [`MembershipEvent::Joined`] when a cluster member joins
435    /// and [`MembershipEvent::Left`] when a cluster member leaves. The stream is
436    /// keyed by the [`MemberId`] of the cluster member.
437    ///
438    /// This is useful for implementing protocols that need to track cluster membership,
439    /// such as broadcasting to all members or detecting failures.
440    ///
441    /// # Non-Determinism
442    /// This stream is non-deterministic because the timing of membership events, for example
443    /// if a node leaves, the membership event may not be received if the node left before the
444    /// stream was created.
445    ///
446    /// # Example
447    /// ```rust
448    /// # #[cfg(feature = "deploy")] {
449    /// # use hydro_lang::prelude::*;
450    /// # use futures::StreamExt;
451    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
452    /// let p1 = flow.process::<()>();
453    /// let workers: Cluster<()> = flow.cluster::<()>();
454    /// # // do nothing on each worker
455    /// # workers.source_iter(q!(vec![])).for_each(q!(|_: ()| {}));
456    /// let cluster_members = p1.source_cluster_members(&workers, nondet!(/** late joiners may miss events */));
457    /// # cluster_members.entries().send(&p2, TCP.fail_stop().bincode())
458    /// // if there are 4 members in the cluster, we would see a join event for each
459    /// // { MemberId::<Worker>(0): [MembershipEvent::Join], MemberId::<Worker>(2): [MembershipEvent::Join], ... }
460    /// # }, |mut stream| async move {
461    /// # let mut results = Vec::new();
462    /// # for w in 0..4 {
463    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
464    /// # }
465    /// # results.sort();
466    /// # assert_eq!(results, vec!["(MemberId::<()>(0), Joined)", "(MemberId::<()>(1), Joined)", "(MemberId::<()>(2), Joined)", "(MemberId::<()>(3), Joined)"]);
467    /// # }));
468    /// # }
469    /// ```
470    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    /// Creates a stream of membership events for a cluster.
482    ///
483    /// This stream emits [`MembershipEvent::Joined`] when a cluster member joins
484    /// and [`MembershipEvent::Left`] when a cluster member leaves. The stream is
485    /// keyed by the [`MemberId`] of the cluster member.
486    ///
487    /// This is useful for implementing protocols that need to track cluster membership,
488    /// such as broadcasting to all members or detecting failures.
489    ///
490    /// # Non-Determinism
491    /// This stream is non-deterministic because the timing of membership events, for example
492    /// if a node leaves, the membership event may not be received if the node left before the
493    /// stream was created.
494    ///
495    /// # Example
496    /// ```rust
497    /// # #[cfg(feature = "deploy")] {
498    /// # use hydro_lang::prelude::*;
499    /// # use futures::StreamExt;
500    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
501    /// let p1 = flow.process::<()>();
502    /// let workers: Cluster<()> = flow.cluster::<()>();
503    /// # // do nothing on each worker
504    /// # workers.source_iter(q!(vec![])).for_each(q!(|_: ()| {}));
505    /// let cluster_members = p1.source_cluster_membership_stream(&workers, nondet!(/** late joiners may miss events */));
506    /// # cluster_members.entries().send(&p2, TCP.fail_stop().bincode())
507    /// // if there are 4 members in the cluster, we would see a join event for each
508    /// // { MemberId::<Worker>(0): [MembershipEvent::Join], MemberId::<Worker>(2): [MembershipEvent::Join], ... }
509    /// # }, |mut stream| async move {
510    /// # let mut results = Vec::new();
511    /// # for w in 0..4 {
512    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
513    /// # }
514    /// # results.sort();
515    /// # assert_eq!(results, vec!["(MemberId::<()>(0), Joined)", "(MemberId::<()>(1), Joined)", "(MemberId::<()>(2), Joined)", "(MemberId::<()>(3), Joined)"]);
516    /// # }));
517    /// # }
518    /// ```
519    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    /// Creates a one-way connection from an external process to receive raw bytes.
547    ///
548    /// Returns a port handle for the external process to connect to, and a stream
549    /// of received byte buffers.
550    ///
551    /// For bidirectional communication or typed data, see [`Location::bind_single_client`]
552    /// or [`Location::source_external_bincode`].
553    #[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    /// Creates a one-way connection from an external process to receive bincode-serialized data.
573    ///
574    /// Returns a sink handle for the external process to send data to, and a stream
575    /// of received values.
576    ///
577    /// For bidirectional communication, see [`Location::bind_single_client_bincode`].
578    #[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    /// Sets up a simulated input port on this location for testing.
604    ///
605    /// Returns a handle to send messages to the location as well as a stream
606    /// of received messages. This is only available when the `sim` feature is enabled.
607    #[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    /// Creates an external input stream for embedded deployment mode.
630    ///
631    /// The `name` parameter specifies the name of the generated function parameter
632    /// that will supply data to this stream at runtime. The generated function will
633    /// accept an `impl Stream<Item = T> + Unpin` argument with this name.
634    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    /// Creates an embedded singleton input for embedded deployment mode.
660    ///
661    /// The `name` parameter specifies the name of the generated function parameter
662    /// that will supply data to this singleton at runtime. The generated function will
663    /// accept a plain `T` parameter with this name.
664    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    /// Establishes a server on this location to receive a bidirectional connection from a single
685    /// client, identified by the given `External` handle. Returns a port handle for the external
686    /// process to connect to, a stream of incoming messages, and a handle to send outgoing
687    /// messages.
688    ///
689    /// # Example
690    /// ```rust
691    /// # #[cfg(feature = "deploy")] {
692    /// # use hydro_lang::prelude::*;
693    /// # use hydro_deploy::Deployment;
694    /// # use futures::{SinkExt, StreamExt};
695    /// # tokio_test::block_on(async {
696    /// # use bytes::Bytes;
697    /// # use hydro_lang::location::NetworkHint;
698    /// # use tokio_util::codec::LengthDelimitedCodec;
699    /// # let mut flow = FlowBuilder::new();
700    /// let node = flow.process::<()>();
701    /// let external = flow.external::<()>();
702    /// let (port, incoming, outgoing) =
703    ///     node.bind_single_client::<_, Bytes, LengthDelimitedCodec>(&external, NetworkHint::Auto);
704    /// outgoing.complete(incoming.map(q!(|data /* : Bytes */| {
705    ///     let mut resp: Vec<u8> = data.into();
706    ///     resp.push(42);
707    ///     resp.into() // : Bytes
708    /// })));
709    ///
710    /// # let mut deployment = Deployment::new();
711    /// let nodes = flow // ... with_process and with_external
712    /// #     .with_process(&node, deployment.Localhost())
713    /// #     .with_external(&external, deployment.Localhost())
714    /// #     .deploy(&mut deployment);
715    ///
716    /// deployment.deploy().await.unwrap();
717    /// deployment.start().await.unwrap();
718    ///
719    /// let (mut external_out, mut external_in) = nodes.connect(port).await;
720    /// external_in.send(vec![1, 2, 3].into()).await.unwrap();
721    /// assert_eq!(
722    ///     external_out.next().await.unwrap().unwrap(),
723    ///     vec![1, 2, 3, 42]
724    /// );
725    /// # });
726    /// # }
727    /// ```
728    #[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    /// Establishes a bidirectional connection from a single external client using bincode serialization.
805    ///
806    /// Returns a port handle for the external process to connect to, a stream of incoming messages,
807    /// and a handle to send outgoing messages. This is a convenience wrapper around
808    /// [`Location::bind_single_client`] that uses bincode for serialization.
809    ///
810    /// # Type Parameters
811    /// - `InT`: The type of incoming messages (must implement [`DeserializeOwned`])
812    /// - `OutT`: The type of outgoing messages (must implement [`Serialize`])
813    #[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    /// Establishes a server on this location to receive bidirectional connections from multiple
902    /// external clients using raw bytes.
903    ///
904    /// Unlike [`Location::bind_single_client`], this method supports multiple concurrent client
905    /// connections. Each client is assigned a unique `u64` identifier.
906    ///
907    /// Returns:
908    /// - A port handle for external processes to connect to
909    /// - A keyed stream of incoming messages, keyed by client ID
910    /// - A keyed stream of membership events (client joins/leaves), keyed by client ID
911    /// - A handle to send outgoing messages, keyed by client ID
912    #[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() // TODO(shadaj): this silently drops framing errors, decide on right defaults
1036                .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    /// Establishes a server on this location to receive bidirectional connections from multiple
1049    /// external clients using bincode serialization.
1050    ///
1051    /// Unlike [`Location::bind_single_client_bincode`], this method supports multiple concurrent
1052    /// client connections. Each client is assigned a unique `u64` identifier.
1053    ///
1054    /// Returns:
1055    /// - A port handle for external processes to connect to
1056    /// - A keyed stream of incoming messages, keyed by client ID
1057    /// - A keyed stream of membership events (client joins/leaves), keyed by client ID
1058    /// - A handle to send outgoing messages, keyed by client ID
1059    ///
1060    /// # Type Parameters
1061    /// - `InT`: The type of incoming messages (must implement [`DeserializeOwned`])
1062    /// - `OutT`: The type of outgoing messages (must implement [`Serialize`])
1063    #[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    /// Bridges user-owned async code to the dataflow as a **bidirectional sidecar**.
1210    ///
1211    /// The closure is called once at startup and must return a
1212    /// `(Stream<InT>, Sink<OutT>)` pair. The framework reads from the stream
1213    /// (items flowing *into* the dataflow) and writes to the sink (items flowing
1214    /// *out* to the sidecar). The user controls buffering, backpressure, and
1215    /// internal lifecycle — Hydro only sees the stream/sink interface.
1216    ///
1217    /// This will hopefully make it easy to integrate hydro with existing frameworks,
1218    /// for example grpc code generated service endpoints.
1219    ///
1220    /// # Returns
1221    /// - A `Stream<InT>` carrying items from the sidecar into the dataflow.
1222    /// - A [`ForwardHandle`] expecting a `Stream<OutT>` that the user completes
1223    ///   with items destined for the sidecar.
1224    ///
1225    /// # Example
1226    ///
1227    /// ```rust
1228    /// # #[cfg(feature = "deploy")] {
1229    /// # use hydro_lang::prelude::*;
1230    /// # use futures::StreamExt;
1231    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1232    /// // Sidecar that echoes whatever it receives back into the dataflow.
1233    /// let (inbound, response_handle) = process.sidecar_bidi::<String, String, _>(q!(|| {
1234    ///     let (to_df_tx, to_df_rx) = tokio::sync::mpsc::channel::<String>(16);
1235    ///     let (from_df_tx, mut from_df_rx) = tokio::sync::mpsc::channel::<String>(16);
1236    ///
1237    ///     // Spawn the sidecar: echoes items from the dataflow back into it.
1238    ///     tokio::spawn(async move {
1239    ///         while let Some(msg) = from_df_rx.recv().await {
1240    ///             to_df_tx.send(msg).await.ok();
1241    ///         }
1242    ///     });
1243    ///
1244    ///     // Return the framework-facing ends (concrete types, no boxing needed).
1245    ///     let stream = tokio_stream::wrappers::ReceiverStream::new(to_df_rx);
1246    ///     let sink = tokio_util::sync::PollSender::new(from_df_tx);
1247    ///     (stream, sink)
1248    /// }));
1249    ///
1250    /// // Send "hello" into the sidecar via the response channel.
1251    /// let input = process.source_stream(q!(futures::stream::iter(vec!["hello".to_string()])));
1252    /// response_handle.complete(input);
1253    ///
1254    /// // The sidecar echoes it back — assert we get "hello" out.
1255    /// inbound
1256    /// # }, |mut stream| async move {
1257    /// #     assert_eq!(stream.next().await.unwrap(), "hello");
1258    /// # }));
1259    /// # }
1260    /// ```
1261    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        // Inbound stream: reads from the stream returned by the sidecar closure
1287        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,  // TODO: maybe bounded sidecars are interesting..?
1298                    TotalOrder, // TODO: NoOrder..?
1299                    ExactlyOnce,
1300                >::collection_kind()),
1301            },
1302        );
1303
1304        // Outbound: forward_ref cycle feeding the sink returned by the sidecar closure
1305        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    /// Constructs a [`Singleton`] materialized at this location with the given static value.
1327    ///
1328    /// See also: [`Tick::singleton`], for creating a singleton _within_ a tick, which requires
1329    /// `T: Clone`.
1330    ///
1331    /// # Example
1332    /// ```rust
1333    /// # #[cfg(feature = "deploy")] {
1334    /// # use hydro_lang::prelude::*;
1335    /// # use futures::StreamExt;
1336    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1337    /// let singleton = process.singleton(q!(5));
1338    /// # singleton.into_stream()
1339    /// # }, |mut stream| async move {
1340    /// // 5
1341    /// # assert_eq!(stream.next().await.unwrap(), 5);
1342    /// # }));
1343    /// # }
1344    /// ```
1345    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    /// Constructs a [`Singleton`] by resolving an async [`Future`] to completion.
1370    ///
1371    /// This is a convenience method equivalent to
1372    /// `self.singleton(future_expr).resolve_future_blocking()`, which is a common
1373    /// pattern when initializing a singleton from an async computation.
1374    ///
1375    /// # Example
1376    /// ```rust
1377    /// # #[cfg(feature = "deploy")] {
1378    /// # use hydro_lang::prelude::*;
1379    /// # use futures::StreamExt;
1380    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1381    /// let singleton = process.singleton_future(q!(async { 42 }));
1382    /// singleton.into_stream()
1383    /// # }, |mut stream| async move {
1384    /// // 42
1385    /// # assert_eq!(stream.next().await.unwrap(), 42);
1386    /// # }));
1387    /// # }
1388    /// ```
1389    ///
1390    /// [`Future`]: std::future::Future
1391    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    /// Generates a stream that emits `()` at a fixed interval.
1403    ///
1404    /// The first tick completes immediately. Missed ticks will be scheduled
1405    /// as soon as possible.
1406    ///
1407    /// Because this only emits `()`, the non-determinism of *when* events fire
1408    /// is captured by the `AtLeastOnce` retry semantics downstream, so no
1409    /// [`NonDet`] guard is required.
1410    #[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!(/** interval does not reveal timestamps */),
1424        )
1425    }
1426
1427    /// Generates a stream that emits `()` at a fixed interval, after an
1428    /// initial delay.
1429    ///
1430    /// Because this only emits `()`, the non-determinism of *when* events fire
1431    /// is captured by the `AtLeastOnce` retry semantics downstream, so no
1432    /// [`NonDet`] guard is required.
1433    #[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!(/** interval does not reveal timestamps */),
1451        )
1452    }
1453
1454    /// Creates a forward reference, allowing a stream to be used before its source is defined.
1455    ///
1456    /// Returns a `(handle, placeholder)` pair. Use the placeholder in the dataflow graph,
1457    /// then call `handle.complete(actual_stream)` to wire in the real source.
1458    ///
1459    /// This is useful for mutually-dependent dataflows or when the definition order
1460    /// doesn't match the data flow direction. For feedback loops, prefer [`Tick::cycle`]
1461    /// instead, which automatically defers values by one tick.
1462    ///
1463    /// # Panics
1464    /// Panics if the forward reference creates a synchronous cycle (i.e., the completed
1465    /// stream transitively depends on the placeholder without a `defer_tick` or network
1466    /// hop in between).
1467    ///
1468    /// # Example
1469    /// ```rust
1470    /// # #[cfg(feature = "deploy")] {
1471    /// # use hydro_lang::prelude::*;
1472    /// # use hydro_lang::live_collections::stream::NoOrder;
1473    /// # use futures::StreamExt;
1474    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1475    /// // Create a forward reference to define a stream that will be completed later
1476    /// let (complete, forward_stream) = process.forward_ref::<Stream<i32, _, _, NoOrder>>();
1477    ///
1478    /// // Use the forward reference as input to another computation
1479    /// let output: Stream<_, _, _, NoOrder> = forward_stream.map(q!(|x| x * 2));
1480    ///
1481    /// // Complete the forward reference with the actual source
1482    /// let source: Stream<_, _, Unbounded> = process.source_iter(q!([1, 2, 3])).into();
1483    /// complete.complete(source);
1484    /// output
1485    /// # }, |mut stream| async move {
1486    /// // 2, 4, 6
1487    /// # assert_eq!(stream.next().await.unwrap(), 2);
1488    /// # assert_eq!(stream.next().await.unwrap(), 4);
1489    /// # assert_eq!(stream.next().await.unwrap(), 6);
1490    /// # }));
1491    /// # }
1492    /// ```
1493    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!(/** test */))
1534            .cross_singleton(singleton.clone().snapshot(&tick, nondet!(/** test */)))
1535            .cross_singleton(
1536                singleton
1537                    .snapshot(&tick, nondet!(/** test */))
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!(/** test */))
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        // intentionally skipped to test stream waking logic
1695        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() // : Bytes
1763        })));
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}