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::builder::ExternalPortId;
43#[cfg(feature = "tokio")]
44use crate::compile::ir::DebugInstantiate;
45use crate::compile::ir::{
46    ClusterMembersState, HydroIrOpMetadata, HydroNode, HydroRoot, HydroSource,
47};
48use crate::forward_handle::ForwardRef;
49#[cfg(stageleft_runtime)]
50use crate::forward_handle::{CycleCollection, ForwardHandle};
51use crate::live_collections::boundedness::{Bounded, Unbounded};
52use crate::live_collections::keyed_stream::KeyedStream;
53use crate::live_collections::singleton::Singleton;
54#[cfg(feature = "sim")]
55#[cfg(stageleft_runtime)]
56use crate::live_collections::stream::networking::serialize_bincode;
57use crate::live_collections::stream::{ExactlyOnce, NoOrder, Stream, TotalOrder};
58#[cfg(feature = "tokio")]
59use crate::live_collections::stream::{Ordering, Retries};
60#[cfg(stageleft_runtime)]
61use crate::location::dynamic::DynLocation;
62use crate::location::dynamic::{ClusterConsistency, LocationId};
63#[cfg(feature = "tokio")]
64use crate::location::external_process::{
65    ExternalBincodeBidi, ExternalBincodeSink, ExternalBytesPort, Many, NotMany,
66};
67use crate::nondet::NonDet;
68#[cfg(feature = "tokio")]
69use crate::properties::manual_proof;
70#[cfg(feature = "sim")]
71use crate::sim::SimSender;
72use crate::staging_util::get_this_crate;
73
74pub mod dynamic;
75
76pub mod external_process;
77pub use external_process::External;
78
79pub mod process;
80pub use process::Process;
81
82pub mod cluster;
83pub use cluster::Cluster;
84
85pub mod member_id;
86pub use member_id::{MemberId, TaglessMemberId};
87
88pub mod tick;
89pub use tick::{Atomic, Tick};
90
91/// An event indicating a change in membership status of a location in a group
92/// (e.g. a node in a [`Cluster`] or an external client connection).
93#[derive(PartialEq, Eq, Clone, Debug, Hash, Serialize, Deserialize)]
94pub enum MembershipEvent {
95    /// The member has joined the group and is now active.
96    Joined,
97    /// The member has left the group and is no longer active.
98    Left,
99}
100
101/// A hint for configuring the network transport used by an external connection.
102///
103/// This controls how the underlying TCP listener is set up when binding
104/// external client connections via methods like [`Location::bind_single_client`]
105/// or [`Location::bidi_external_many_bytes`].
106#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
107pub enum NetworkHint {
108    /// Automatically select the network configuration (e.g. an ephemeral port).
109    Auto,
110    /// Use a TCP port, optionally specifying a fixed port number.
111    ///
112    /// If `None`, an available port will be chosen automatically.
113    /// If `Some(port)`, the given port number will be used.
114    TcpPort(Option<u16>),
115}
116
117#[track_caller]
118pub(crate) fn check_matching_location<'a, L: Location<'a>>(l1: &L, l2: &L) {
119    assert_eq!(Location::id(l1), Location::id(l2), "locations do not match");
120}
121
122#[stageleft::export(LocationKey)]
123new_key_type! {
124    /// A unique identifier for a clock tick.
125    pub struct LocationKey;
126}
127
128impl std::fmt::Display for LocationKey {
129    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
130        write!(f, "loc{:?}", self.data()) // `"loc1v1"``
131    }
132}
133
134/// This is used for the ECS membership stream.
135/// TODO(mingwei): Make this more robust?
136impl std::str::FromStr for LocationKey {
137    type Err = Option<ParseIntError>;
138
139    fn from_str(s: &str) -> Result<Self, Self::Err> {
140        let nvn = s.strip_prefix("loc").ok_or(None)?;
141        let (idx, ver) = nvn.split_once('v').ok_or(None)?;
142        let idx: u64 = idx.parse()?;
143        let ver: u64 = ver.parse()?;
144        Ok(slotmap::KeyData::from_ffi((ver << 32) | idx).into())
145    }
146}
147
148impl LocationKey {
149    /// TODO(minwgei): Remove this and avoid magic key for simulator external.
150    /// The first location key, used by the simulator as the default external location.
151    pub const FIRST: Self = Self(slotmap::KeyData::from_ffi(0x0000000100000001)); // `1v1`
152
153    /// A key for testing with index 1.
154    #[cfg(test)]
155    pub const TEST_KEY_1: Self = Self(slotmap::KeyData::from_ffi(0x000000FF00000001)); // `1v255`
156
157    /// A key for testing with index 2.
158    #[cfg(test)]
159    pub const TEST_KEY_2: Self = Self(slotmap::KeyData::from_ffi(0x000000FF00000002)); // `2v255`
160}
161
162/// This is used within `q!` code in docker and ECS.
163impl<Ctx> FreeVariableWithContextWithProps<Ctx, ()> for LocationKey {
164    type O = LocationKey;
165
166    fn to_tokens(self, _ctx: &Ctx) -> (QuoteTokens, ())
167    where
168        Self: Sized,
169    {
170        let root = get_this_crate();
171        let n = Key::data(&self).as_ffi();
172        (
173            QuoteTokens {
174                prelude: None,
175                expr: Some(quote! {
176                    #root::location::LocationKey::from(#root::runtime_support::slotmap::KeyData::from_ffi(#n))
177                }),
178            },
179            (),
180        )
181    }
182}
183
184/// A simple enum for the type of a root location.
185#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Serialize)]
186pub enum LocationType {
187    /// A process (single node).
188    Process,
189    /// A cluster (multiple nodes).
190    Cluster,
191    /// An external client.
192    External,
193}
194
195/// A top-level location (i.e. a [`Process`] or [`Cluster`]) that is outside a tick / atomic region.
196pub trait TopLevel<'a>: Location<'a> {}
197
198#[cfg(feature = "sim")]
199#[cfg(stageleft_runtime)]
200fn register_serialized_external_input<'a, At, L, T>(
201    at: &At,
202    from: &External<'_, L>,
203    deserialize_fn: syn::Expr,
204) -> (
205    ExternalPortId,
206    Stream<T, At::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
207)
208where
209    At: TopLevel<'a> + Sized,
210{
211    let (port_id, stream, sink) = at.register_serialized_single_client::<_, T, ()>(
212        from,
213        serialize_bincode::<()>(false),
214        deserialize_fn,
215    );
216    sink.complete(stream.location().source_iter(q!([])));
217
218    (port_id, stream)
219}
220
221/// A location where data can be materialized and computation can be executed.
222///
223/// Hydro is a **global**, **distributed** programming model. This means that the data
224/// and computation in a Hydro program can be spread across multiple machines, data
225/// centers, and even continents. To achieve this, Hydro uses the concept of
226/// **locations** to keep track of _where_ data is located and computation is executed.
227///
228/// Each live collection type (in [`crate::live_collections`]) has a type parameter `L`
229/// which will always be a type that implements the [`Location`] trait (e.g. [`Process`]
230/// and [`Cluster`]). To create distributed programs, Hydro provides a variety of APIs
231/// to allow live collections to be _moved_ between locations via network send/receive.
232///
233/// See [the Hydro docs](https://hydro.run/docs/hydro/reference/locations/) for more information.
234#[expect(
235    private_bounds,
236    reason = "only internal Hydro code can define location types"
237)]
238pub trait Location<'a>: DynLocation {
239    /// The root location type for this location.
240    ///
241    /// For top-level locations like [`Process`] and [`Cluster`], this is `Self`.
242    /// For nested locations like [`Tick`], this is the root location that contains it.
243    type Root: Location<'a>;
244
245    /// The scope of simulator hook handles (see [`crate::sim_hooks`]) bindable to
246    /// unsafe operators at this location: [`OnProcess<P>`] for a [`Process<P>`](Process)
247    /// root, [`OnCluster<C>`] for a [`Cluster<C>`](Cluster) root. Operators name this in
248    /// their `NonDet` hook payload, so a handle can only bind to a matching location
249    /// kind.
250    ///
251    /// [`OnProcess<P>`]: crate::sim_hooks::OnProcess
252    /// [`OnCluster<C>`]: crate::sim_hooks::OnCluster
253    type SimHookScope: crate::sim_hooks::BindableHookScope;
254
255    /// Location type with consistency guarantees dropped for the live collection on it.
256    type DropConsistency: Location<'a, DropConsistency = Self::DropConsistency, SimHookScope = Self::SimHookScope>;
257
258    /// Returns the root location for this location.
259    ///
260    /// For top-level locations like [`Process`] and [`Cluster`], this returns `self`.
261    /// For nested locations like [`Tick`], this returns the root location that contains it.
262    fn root(&self) -> Self::Root;
263
264    /// This location but with consistency guarantees dropped for the live collection
265    fn drop_consistency(&self) -> Self::DropConsistency;
266    /// Gets the runtime enum variant for the current consistency level, if this is a cluster.
267    fn consistency() -> Option<ClusterConsistency>;
268
269    /// Updates the consistency guarantees to match that of the given location.
270    fn with_consistency_of<L2: Location<'a, DropConsistency = Self::DropConsistency>>(&self) -> L2 {
271        L2::from_drop_consistency(self.drop_consistency())
272    }
273
274    #[doc(hidden)]
275    fn from_drop_consistency(l2: Self::DropConsistency) -> Self;
276
277    /// Attempts to create a new [`Tick`] clock domain at this location.
278    ///
279    /// Returns `Some(Tick)` if this is a top-level location (like [`Process`] or [`Cluster`]),
280    /// or `None` if this location is already inside a tick (nested ticks are not supported).
281    ///
282    /// Prefer using [`Location::tick`] when you know the location is top-level.
283    fn try_tick(&self) -> Option<Tick<Self>> {
284        if Self::is_top_level() {
285            let id = if let LocationId::Atomic { .. } = self.id() {
286                None
287            } else {
288                Some(self.flow_state().borrow_mut().next_clock_id())
289            };
290            Some(Tick {
291                id,
292                l: self.clone(),
293            })
294        } else {
295            None
296        }
297    }
298
299    /// Returns the unique identifier for this location.
300    fn id(&self) -> LocationId {
301        DynLocation::dyn_id(self)
302    }
303
304    /// Creates a new [`Tick`] clock domain at this location.
305    ///
306    /// A tick represents a logical clock that can be used to batch streaming data
307    /// into discrete time steps. This is useful for implementing iterative algorithms
308    /// or for synchronizing data across multiple streams.
309    ///
310    /// # Example
311    /// ```rust
312    /// # #[cfg(feature = "deploy")] {
313    /// # use hydro_lang::prelude::*;
314    /// # use futures::StreamExt;
315    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
316    /// let tick = process.tick();
317    /// let inside_tick = process
318    ///     .source_iter(q!(vec![1, 2, 3, 4]))
319    ///     .batch(&tick, nondet!(/** test */));
320    /// inside_tick.all_ticks()
321    /// # }, |mut stream| async move {
322    /// // 1, 2, 3, 4
323    /// # for w in vec![1, 2, 3, 4] {
324    /// #     assert_eq!(stream.next().await.unwrap(), w);
325    /// # }
326    /// # }));
327    /// # }
328    /// ```
329    fn tick(&self) -> Tick<Self> {
330        self.try_tick().expect("cannot create nested ticks")
331    }
332
333    /// Creates an unbounded stream that continuously emits unit values `()`.
334    ///
335    /// This is useful for driving computations that need to run continuously,
336    /// such as polling or heartbeat mechanisms.
337    ///
338    /// # Example
339    /// ```rust
340    /// # #[cfg(feature = "deploy")] {
341    /// # use hydro_lang::prelude::*;
342    /// # use futures::StreamExt;
343    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
344    /// let tick = process.tick();
345    /// process.spin()
346    ///     .batch(&tick, nondet!(/** test */))
347    ///     .map(q!(|_| 42))
348    ///     .all_ticks()
349    /// # }, |mut stream| async move {
350    /// // 42, 42, 42, ...
351    /// # assert_eq!(stream.next().await.unwrap(), 42);
352    /// # assert_eq!(stream.next().await.unwrap(), 42);
353    /// # assert_eq!(stream.next().await.unwrap(), 42);
354    /// # }));
355    /// # }
356    /// ```
357    fn spin(&self) -> Stream<(), Self, Unbounded, TotalOrder, ExactlyOnce>
358    where
359        Self: TopLevel<'a> + Sized,
360    {
361        Stream::new(
362            self.clone(),
363            HydroNode::Source {
364                source: HydroSource::Spin(),
365                metadata: self.new_node_metadata(Stream::<
366                    (),
367                    Self,
368                    Unbounded,
369                    TotalOrder,
370                    ExactlyOnce,
371                >::collection_kind()),
372            },
373        )
374    }
375
376    /// Creates a stream from an async [`FuturesStream`].
377    ///
378    /// This is useful for integrating with external async data sources,
379    /// such as network connections or file readers.
380    ///
381    /// # Example
382    /// ```rust
383    /// # #[cfg(feature = "deploy")] {
384    /// # use hydro_lang::prelude::*;
385    /// # use futures::StreamExt;
386    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
387    /// process.source_stream(q!(futures::stream::iter(vec![1, 2, 3])))
388    /// # }, |mut stream| async move {
389    /// // 1, 2, 3
390    /// # for w in vec![1, 2, 3] {
391    /// #     assert_eq!(stream.next().await.unwrap(), w);
392    /// # }
393    /// # }));
394    /// # }
395    /// ```
396    fn source_stream<T, E>(
397        &self,
398        e: impl QuotedWithContext<'a, E, Self>,
399    ) -> Stream<T, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>
400    where
401        E: FuturesStream<Item = T> + Unpin,
402        Self: TopLevel<'a> + Sized,
403    {
404        let e = e.splice_untyped_ctx(self);
405
406        let target_location = self.drop_consistency();
407        Stream::new(
408            target_location.clone(),
409            HydroNode::Source {
410                source: HydroSource::Stream(e.into()),
411                metadata: target_location.new_node_metadata(Stream::<
412                    T,
413                    Self::DropConsistency,
414                    Unbounded,
415                    TotalOrder,
416                    ExactlyOnce,
417                >::collection_kind()),
418            },
419        )
420    }
421
422    /// Creates a bounded stream from an iterator.
423    ///
424    /// The iterator is evaluated once at runtime, and all elements are emitted
425    /// in order. This is useful for creating streams from static data or
426    /// for testing.
427    ///
428    /// # Example
429    /// ```rust
430    /// # #[cfg(feature = "deploy")] {
431    /// # use hydro_lang::prelude::*;
432    /// # use futures::StreamExt;
433    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
434    /// process.source_iter(q!(vec![1, 2, 3, 4]))
435    /// # }, |mut stream| async move {
436    /// // 1, 2, 3, 4
437    /// # for w in vec![1, 2, 3, 4] {
438    /// #     assert_eq!(stream.next().await.unwrap(), w);
439    /// # }
440    /// # }));
441    /// # }
442    /// ```
443    fn source_iter<T, E>(
444        &self,
445        e: impl QuotedWithContext<'a, E, Self>,
446    ) -> Stream<T, Self::DropConsistency, Bounded, TotalOrder, ExactlyOnce>
447    where
448        E: IntoIterator<Item = T>,
449        Self: Sized,
450    {
451        let e = e.splice_typed_ctx(self);
452
453        let target_location = self.drop_consistency();
454        Stream::new(
455            target_location.clone(),
456            HydroNode::Source {
457                source: HydroSource::Iter(e.into()),
458                metadata: target_location.new_node_metadata(Stream::<
459                    T,
460                    Self::DropConsistency,
461                    Bounded,
462                    TotalOrder,
463                    ExactlyOnce,
464                >::collection_kind()),
465            },
466        )
467    }
468
469    #[deprecated(note = "use .source_cluster_membership_stream(...) instead")]
470    /// Creates a stream of membership events for a cluster.
471    ///
472    /// This stream emits [`MembershipEvent::Joined`] when a cluster member joins
473    /// and [`MembershipEvent::Left`] when a cluster member leaves. The stream is
474    /// keyed by the [`MemberId`] of the cluster member.
475    ///
476    /// This is useful for implementing protocols that need to track cluster membership,
477    /// such as broadcasting to all members or detecting failures.
478    ///
479    /// # Non-Determinism
480    /// This stream is non-deterministic because the timing of membership events, for example
481    /// if a node leaves, the membership event may not be received if the node left before the
482    /// stream was created.
483    ///
484    /// # Example
485    /// ```rust
486    /// # #[cfg(feature = "deploy")] {
487    /// # use hydro_lang::prelude::*;
488    /// # use futures::StreamExt;
489    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
490    /// let p1 = flow.process::<()>();
491    /// let workers: Cluster<()> = flow.cluster::<()>();
492    /// # // do nothing on each worker
493    /// # workers.source_iter(q!(vec![])).for_each(q!(|_: ()| {}));
494    /// let cluster_members = p1.source_cluster_members(&workers, nondet!(/** late joiners may miss events */));
495    /// # cluster_members.entries().send(&p2, TCP.fail_stop().bincode())
496    /// // if there are 4 members in the cluster, we would see a join event for each
497    /// // { MemberId::<Worker>(0): [MembershipEvent::Join], MemberId::<Worker>(2): [MembershipEvent::Join], ... }
498    /// # }, |mut stream| async move {
499    /// # let mut results = Vec::new();
500    /// # for w in 0..4 {
501    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
502    /// # }
503    /// # results.sort();
504    /// # assert_eq!(results, vec!["(MemberId::<()>(0), Joined)", "(MemberId::<()>(1), Joined)", "(MemberId::<()>(2), Joined)", "(MemberId::<()>(3), Joined)"]);
505    /// # }));
506    /// # }
507    /// ```
508    fn source_cluster_members<C: 'a>(
509        &self,
510        cluster: &Cluster<'a, C>,
511        nondet_start: NonDet,
512    ) -> KeyedStream<MemberId<C>, MembershipEvent, Self::DropConsistency, Unbounded>
513    where
514        Self: TopLevel<'a> + Sized,
515    {
516        self.source_cluster_membership_stream(cluster, nondet_start)
517    }
518
519    /// Creates a stream of membership events for a cluster.
520    ///
521    /// This stream emits [`MembershipEvent::Joined`] when a cluster member joins
522    /// and [`MembershipEvent::Left`] when a cluster member leaves. The stream is
523    /// keyed by the [`MemberId`] of the cluster member.
524    ///
525    /// This is useful for implementing protocols that need to track cluster membership,
526    /// such as broadcasting to all members or detecting failures.
527    ///
528    /// # Non-Determinism
529    /// This stream is non-deterministic because the timing of membership events, for example
530    /// if a node leaves, the membership event may not be received if the node left before the
531    /// stream was created.
532    ///
533    /// # Example
534    /// ```rust
535    /// # #[cfg(feature = "deploy")] {
536    /// # use hydro_lang::prelude::*;
537    /// # use futures::StreamExt;
538    /// # tokio_test::block_on(hydro_lang::test_util::multi_location_test(|flow, p2| {
539    /// let p1 = flow.process::<()>();
540    /// let workers: Cluster<()> = flow.cluster::<()>();
541    /// # // do nothing on each worker
542    /// # workers.source_iter(q!(vec![])).for_each(q!(|_: ()| {}));
543    /// let cluster_members = p1.source_cluster_membership_stream(&workers, nondet!(/** late joiners may miss events */));
544    /// # cluster_members.entries().send(&p2, TCP.fail_stop().bincode())
545    /// // if there are 4 members in the cluster, we would see a join event for each
546    /// // { MemberId::<Worker>(0): [MembershipEvent::Join], MemberId::<Worker>(2): [MembershipEvent::Join], ... }
547    /// # }, |mut stream| async move {
548    /// # let mut results = Vec::new();
549    /// # for w in 0..4 {
550    /// #     results.push(format!("{:?}", stream.next().await.unwrap()));
551    /// # }
552    /// # results.sort();
553    /// # assert_eq!(results, vec!["(MemberId::<()>(0), Joined)", "(MemberId::<()>(1), Joined)", "(MemberId::<()>(2), Joined)", "(MemberId::<()>(3), Joined)"]);
554    /// # }));
555    /// # }
556    /// ```
557    fn source_cluster_membership_stream<C: 'a>(
558        &self,
559        cluster: &Cluster<'a, C>,
560        _nondet_start: NonDet,
561    ) -> KeyedStream<MemberId<C>, MembershipEvent, Self::DropConsistency, Unbounded>
562    where
563        Self: TopLevel<'a> + Sized,
564    {
565        let target_consistency = self.drop_consistency();
566        Stream::new(
567            target_consistency.clone(),
568            HydroNode::Source {
569                source: HydroSource::ClusterMembers(cluster.id(), ClusterMembersState::Uninit),
570                metadata: target_consistency.new_node_metadata(Stream::<
571                    (TaglessMemberId, MembershipEvent),
572                    Self,
573                    Unbounded,
574                    TotalOrder,
575                    ExactlyOnce,
576                >::collection_kind(
577                )),
578            },
579        )
580        .map(q!(|(k, v)| (MemberId::from_tagless(k), v)))
581        .into_keyed()
582    }
583
584    /// Creates a one-way connection from an external process to receive raw bytes.
585    ///
586    /// Returns a port handle for the external process to connect to, and a stream
587    /// of received byte buffers.
588    ///
589    /// For bidirectional communication or typed data, see [`Location::bind_single_client`]
590    /// or [`Location::source_external_bincode`].
591    #[cfg(feature = "tokio")]
592    fn source_external_bytes<L>(
593        &self,
594        from: &External<'_, L>,
595    ) -> (
596        ExternalBytesPort,
597        Stream<BytesMut, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
598    )
599    where
600        Self: TopLevel<'a> + Sized,
601    {
602        let (port, stream, sink) =
603            self.bind_single_client::<_, Bytes, LengthDelimitedCodec>(from, NetworkHint::Auto);
604
605        sink.complete(stream.location().source_iter(q!([])));
606
607        (port, stream)
608    }
609
610    /// Creates a one-way connection from an external process to receive bincode-serialized data.
611    ///
612    /// Returns a sink handle for the external process to send data to, and a stream
613    /// of received values.
614    ///
615    /// For bidirectional communication, see [`Location::bind_single_client_bincode`].
616    #[cfg(feature = "tokio")]
617    fn source_external_bincode<L, T, O: Ordering, R: Retries>(
618        &self,
619        from: &External<'_, L>,
620    ) -> (
621        ExternalBincodeSink<T, NotMany, O, R>,
622        Stream<T, Self::DropConsistency, Unbounded, O, R>,
623    )
624    where
625        Self: TopLevel<'a> + Sized,
626        T: Serialize + DeserializeOwned,
627    {
628        let (port, stream, sink) = self.bind_single_client_bincode::<_, T, ()>(from);
629        sink.complete(stream.location().source_iter(q!([])));
630
631        (
632            ExternalBincodeSink {
633                process_key: from.key,
634                port_id: port.port_id,
635                _phantom: PhantomData,
636            },
637            stream.weaken_ordering().weaken_retries(),
638        )
639    }
640
641    /// Sets up a bincode-encoded simulated input port on this location for testing.
642    ///
643    /// Returns a handle to send messages to the location as well as a stream
644    /// of received messages. Use [`Location::sim_input_with`] to select another codec.
645    /// This is only available when the `sim` feature is enabled.
646    #[cfg(feature = "sim")]
647    fn sim_input<T, O: Ordering, R: Retries>(
648        &self,
649    ) -> (
650        SimSender<T, O, R>,
651        Stream<T, Self::DropConsistency, Unbounded, O, R>,
652    )
653    where
654        Self: TopLevel<'a> + Sized,
655        T: Serialize + DeserializeOwned,
656    {
657        self.sim_input_with::<crate::sim::codec::BincodeCodec, T, O, R>()
658    }
659
660    /// Sets up a simulated input port using the codec `C`.
661    ///
662    /// Returns a handle to send messages to the location as well as a stream
663    /// of received messages. The codec is a type parameter; the message, ordering and
664    /// retries types are usually inferred: `location.sim_input_with::<MyCodec, _, _, _>()`.
665    /// Custom codecs implement [`SimCodec`](crate::sim::codec::SimCodec), which documents
666    /// where they must be defined. This is only available when the `sim` feature is enabled.
667    #[cfg(feature = "sim")]
668    fn sim_input_with<C: crate::sim::codec::SimCodec<T>, T, O: Ordering, R: Retries>(
669        &self,
670    ) -> (
671        SimSender<T, O, R>,
672        Stream<T, Self::DropConsistency, Unbounded, O, R>,
673    )
674    where
675        Self: TopLevel<'a> + Sized,
676    {
677        let external_location: External<'a, ()> = External {
678            key: LocationKey::FIRST,
679            flow_state: self.flow_state().clone(),
680            _phantom: PhantomData,
681        };
682
683        let (external_port_id, stream) = register_serialized_external_input(
684            self,
685            &external_location,
686            crate::sim::codec::staged_deserialize::<T, C>(),
687        );
688
689        (
690            SimSender(external_port_id, PhantomData, C::encode),
691            stream.weaken_ordering().weaken_retries(),
692        )
693    }
694
695    /// Creates an external input stream for embedded deployment mode.
696    ///
697    /// The `name` parameter specifies the name of the generated function parameter
698    /// that will supply data to this stream at runtime. The generated function will
699    /// accept an `impl Stream<Item = T> + Unpin` argument with this name.
700    fn embedded_input<T>(
701        &self,
702        name: impl Into<String>,
703    ) -> Stream<T, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>
704    where
705        Self: TopLevel<'a> + Sized,
706    {
707        let ident = syn::Ident::new(&name.into(), Span::call_site());
708
709        let target_location = self.drop_consistency();
710        Stream::new(
711            target_location.clone(),
712            HydroNode::Source {
713                source: HydroSource::Embedded(ident),
714                metadata: target_location.new_node_metadata(Stream::<
715                    T,
716                    Self,
717                    Unbounded,
718                    TotalOrder,
719                    ExactlyOnce,
720                >::collection_kind()),
721            },
722        )
723    }
724
725    /// Creates an embedded singleton input for embedded deployment mode.
726    ///
727    /// The `name` parameter specifies the name of the generated function parameter
728    /// that will supply data to this singleton at runtime. The generated function will
729    /// accept a plain `T` parameter with this name.
730    fn embedded_singleton_input<T>(
731        &self,
732        name: impl Into<String>,
733    ) -> Singleton<T, Self::DropConsistency, Bounded>
734    where
735        Self: TopLevel<'a> + Sized,
736    {
737        let ident = syn::Ident::new(&name.into(), Span::call_site());
738
739        let target_location = self.drop_consistency();
740        Singleton::new(
741            target_location.clone(),
742            HydroNode::Source {
743                source: HydroSource::EmbeddedSingleton(ident),
744                metadata: target_location
745                    .new_node_metadata(Singleton::<T, Self, Bounded>::collection_kind()),
746            },
747        )
748    }
749
750    /// Establishes a server on this location to receive a bidirectional connection from a single
751    /// client, identified by the given `External` handle. Returns a port handle for the external
752    /// process to connect to, a stream of incoming messages, and a handle to send outgoing
753    /// messages.
754    ///
755    /// # Example
756    /// ```rust
757    /// # #[cfg(feature = "deploy")] {
758    /// # use hydro_lang::prelude::*;
759    /// # use hydro_deploy::Deployment;
760    /// # use futures::{SinkExt, StreamExt};
761    /// # tokio_test::block_on(async {
762    /// # use bytes::Bytes;
763    /// # use hydro_lang::location::NetworkHint;
764    /// # use tokio_util::codec::LengthDelimitedCodec;
765    /// # let mut flow = FlowBuilder::new();
766    /// let node = flow.process::<()>();
767    /// let external = flow.external::<()>();
768    /// let (port, incoming, outgoing) =
769    ///     node.bind_single_client::<_, Bytes, LengthDelimitedCodec>(&external, NetworkHint::Auto);
770    /// outgoing.complete(incoming.map(q!(|data /* : Bytes */| {
771    ///     let mut resp: Vec<u8> = data.into();
772    ///     resp.push(42);
773    ///     resp.into() // : Bytes
774    /// })));
775    ///
776    /// # let mut deployment = Deployment::new();
777    /// let nodes = flow // ... with_process and with_external
778    /// #     .with_process(&node, deployment.Localhost())
779    /// #     .with_external(&external, deployment.Localhost())
780    /// #     .deploy(&mut deployment);
781    ///
782    /// deployment.deploy().await.unwrap();
783    /// deployment.start().await.unwrap();
784    ///
785    /// let (mut external_out, mut external_in) = nodes.connect(port).await;
786    /// external_in.send(vec![1, 2, 3].into()).await.unwrap();
787    /// assert_eq!(
788    ///     external_out.next().await.unwrap().unwrap(),
789    ///     vec![1, 2, 3, 42]
790    /// );
791    /// # });
792    /// # }
793    /// ```
794    #[cfg(feature = "tokio")]
795    #[expect(clippy::type_complexity, reason = "stream markers")]
796    fn bind_single_client<L, T, Codec: Encoder<T> + Decoder>(
797        &self,
798        from: &External<'_, L>,
799        port_hint: NetworkHint,
800    ) -> (
801        ExternalBytesPort<NotMany>,
802        Stream<<Codec as Decoder>::Item, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
803        ForwardHandle<'a, Stream<T, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>>,
804    )
805    where
806        Self: TopLevel<'a> + Sized,
807    {
808        let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
809        let target_consistency = self.drop_consistency();
810
811        let (fwd_ref, to_sink) = target_consistency.forward_ref::<Stream<
812            T,
813            Self::DropConsistency,
814            Unbounded,
815            TotalOrder,
816            ExactlyOnce,
817        >>();
818        let mut flow_state_borrow = self.flow_state().borrow_mut();
819
820        flow_state_borrow.push_root(HydroRoot::SendExternal {
821            to_external_key: from.key,
822            to_port_id: next_external_port_id,
823            to_many: false,
824            unpaired: false,
825            serialize_fn: None,
826            instantiate_fn: DebugInstantiate::Building,
827            input: Box::new(to_sink.ir_node.replace(HydroNode::Placeholder)),
828            op_metadata: HydroIrOpMetadata::new(),
829        });
830        drop(flow_state_borrow);
831
832        let raw_stream: Stream<
833            Result<<Codec as Decoder>::Item, <Codec as Decoder>::Error>,
834            Self::DropConsistency,
835            Unbounded,
836            TotalOrder,
837            ExactlyOnce,
838        > = Stream::new(
839            target_consistency.clone(),
840            HydroNode::ExternalInput {
841                from_external_key: from.key,
842                from_port_id: next_external_port_id,
843                from_many: false,
844                codec_type: quote_type::<Codec>().into(),
845                port_hint,
846                instantiate_fn: DebugInstantiate::Building,
847                deserialize_fn: None,
848                metadata: target_consistency.new_node_metadata(Stream::<
849                    Result<<Codec as Decoder>::Item, <Codec as Decoder>::Error>,
850                    Self::DropConsistency,
851                    Unbounded,
852                    TotalOrder,
853                    ExactlyOnce,
854                >::collection_kind(
855                )),
856            },
857        );
858
859        (
860            ExternalBytesPort {
861                process_key: from.key,
862                port_id: next_external_port_id,
863                _phantom: PhantomData,
864            },
865            raw_stream.flatten_ordered(),
866            fwd_ref,
867        )
868    }
869
870    // TODO: Replace this staged-expression helper with codec-parameterized sink and bidi handles.
871    #[doc(hidden)]
872    #[cfg(feature = "tokio")]
873    #[expect(clippy::type_complexity, reason = "stream markers")]
874    fn register_serialized_single_client<L, InT, OutT>(
875        &self,
876        from: &External<'_, L>,
877        serialize_fn: syn::Expr,
878        deserialize_fn: syn::Expr,
879    ) -> (
880        ExternalPortId,
881        Stream<InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
882        ForwardHandle<'a, Stream<OutT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>>,
883    )
884    where
885        Self: TopLevel<'a> + Sized,
886    {
887        let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
888
889        let target_consistency = self.drop_consistency();
890        let (fwd_ref, to_sink) = target_consistency.forward_ref::<Stream<
891            OutT,
892            Self::DropConsistency,
893            Unbounded,
894            TotalOrder,
895            ExactlyOnce,
896        >>();
897        let mut flow_state_borrow = self.flow_state().borrow_mut();
898
899        flow_state_borrow.push_root(HydroRoot::SendExternal {
900            to_external_key: from.key,
901            to_port_id: next_external_port_id,
902            to_many: false,
903            unpaired: false,
904            serialize_fn: Some(serialize_fn.into()),
905            instantiate_fn: DebugInstantiate::Building,
906            input: Box::new(to_sink.ir_node.replace(HydroNode::Placeholder)),
907            op_metadata: HydroIrOpMetadata::new(),
908        });
909        drop(flow_state_borrow);
910
911        let raw_stream: Stream<InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce> =
912            Stream::new(
913                target_consistency.clone(),
914                HydroNode::ExternalInput {
915                    from_external_key: from.key,
916                    from_port_id: next_external_port_id,
917                    from_many: false,
918                    codec_type: quote_type::<LengthDelimitedCodec>().into(),
919                    port_hint: NetworkHint::Auto,
920                    instantiate_fn: DebugInstantiate::Building,
921                    deserialize_fn: Some(deserialize_fn.into()),
922                    metadata: target_consistency.new_node_metadata(Stream::<
923                        InT,
924                        Self::DropConsistency,
925                        Unbounded,
926                        TotalOrder,
927                        ExactlyOnce,
928                    >::collection_kind(
929                    )),
930                },
931            );
932
933        (next_external_port_id, raw_stream, fwd_ref)
934    }
935
936    /// Establishes a bidirectional connection from a single external client using bincode serialization.
937    ///
938    /// Returns a port handle for the external process to connect to, a stream of incoming messages,
939    /// and a handle to send outgoing messages. This is a convenience wrapper around
940    /// [`Location::bind_single_client`] that uses bincode for serialization.
941    ///
942    /// # Type Parameters
943    /// - `InT`: The type of incoming messages (must implement [`DeserializeOwned`])
944    /// - `OutT`: The type of outgoing messages (must implement [`Serialize`])
945    #[cfg(feature = "tokio")]
946    #[expect(clippy::type_complexity, reason = "stream markers")]
947    fn bind_single_client_bincode<L, InT: DeserializeOwned, OutT: Serialize>(
948        &self,
949        from: &External<'_, L>,
950    ) -> (
951        ExternalBincodeBidi<InT, OutT, NotMany>,
952        Stream<InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
953        ForwardHandle<'a, Stream<OutT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>>,
954    )
955    where
956        Self: TopLevel<'a> + Sized,
957    {
958        let root = get_this_crate();
959
960        let out_t_type = quote_type::<OutT>();
961        let ser_fn: syn::Expr = syn::parse_quote! {
962            #root::runtime_support::stageleft::runtime_support::fn1_type_hint::<#out_t_type, _>(
963                |b| #root::runtime_support::bincode::serialize(&b).unwrap().into()
964            )
965        };
966
967        let in_t_type = quote_type::<InT>();
968        let deser_fn: syn::Expr = syn::parse_quote! {
969            |res| {
970                let b = res.unwrap();
971                #root::runtime_support::bincode::deserialize::<#in_t_type>(&b).unwrap()
972            }
973        };
974
975        let (port_id, raw_stream, fwd_ref) =
976            self.register_serialized_single_client::<_, InT, OutT>(from, ser_fn, deser_fn);
977
978        (
979            ExternalBincodeBidi {
980                process_key: from.key,
981                port_id,
982                _phantom: PhantomData,
983            },
984            raw_stream,
985            fwd_ref,
986        )
987    }
988
989    /// Establishes a server on this location to receive bidirectional connections from multiple
990    /// external clients using raw bytes.
991    ///
992    /// Unlike [`Location::bind_single_client`], this method supports multiple concurrent client
993    /// connections. Each client is assigned a unique `u64` identifier.
994    ///
995    /// Returns:
996    /// - A port handle for external processes to connect to
997    /// - A keyed stream of incoming messages, keyed by client ID
998    /// - A keyed stream of membership events (client joins/leaves), keyed by client ID
999    /// - A handle to send outgoing messages, keyed by client ID
1000    #[cfg(feature = "tokio")]
1001    #[expect(clippy::type_complexity, reason = "stream markers")]
1002    fn bidi_external_many_bytes<L, T, Codec: Encoder<T> + Decoder>(
1003        &self,
1004        from: &External<'_, L>,
1005        port_hint: NetworkHint,
1006    ) -> (
1007        ExternalBytesPort<Many>,
1008        KeyedStream<
1009            u64,
1010            <Codec as Decoder>::Item,
1011            Self::DropConsistency,
1012            Unbounded,
1013            TotalOrder,
1014            ExactlyOnce,
1015        >,
1016        KeyedStream<
1017            u64,
1018            MembershipEvent,
1019            Self::DropConsistency,
1020            Unbounded,
1021            TotalOrder,
1022            ExactlyOnce,
1023        >,
1024        ForwardHandle<
1025            'a,
1026            KeyedStream<u64, T, Self::DropConsistency, Unbounded, NoOrder, ExactlyOnce>,
1027        >,
1028    )
1029    where
1030        Self: TopLevel<'a> + Sized,
1031    {
1032        let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
1033
1034        let target_consistency = self.drop_consistency();
1035        let (fwd_ref, to_sink) = target_consistency.forward_ref::<KeyedStream<
1036            u64,
1037            T,
1038            Self::DropConsistency,
1039            Unbounded,
1040            NoOrder,
1041            ExactlyOnce,
1042        >>();
1043        let to_sink_input = Box::new(to_sink.entries().ir_node.replace(HydroNode::Placeholder));
1044        let mut flow_state_borrow = self.flow_state().borrow_mut();
1045
1046        flow_state_borrow.push_root(HydroRoot::SendExternal {
1047            to_external_key: from.key,
1048            to_port_id: next_external_port_id,
1049            to_many: true,
1050            unpaired: false,
1051            serialize_fn: None,
1052            instantiate_fn: DebugInstantiate::Building,
1053            input: to_sink_input,
1054            op_metadata: HydroIrOpMetadata::new(),
1055        });
1056        drop(flow_state_borrow);
1057
1058        let raw_stream: Stream<
1059            Result<(u64, <Codec as Decoder>::Item), <Codec as Decoder>::Error>,
1060            Self::DropConsistency,
1061            Unbounded,
1062            TotalOrder,
1063            ExactlyOnce,
1064        > = Stream::new(
1065            target_consistency.clone(),
1066            HydroNode::ExternalInput {
1067                from_external_key: from.key,
1068                from_port_id: next_external_port_id,
1069                from_many: true,
1070                codec_type: quote_type::<Codec>().into(),
1071                port_hint,
1072                instantiate_fn: DebugInstantiate::Building,
1073                deserialize_fn: None,
1074                metadata: target_consistency.new_node_metadata(Stream::<
1075                    Result<(u64, <Codec as Decoder>::Item), <Codec as Decoder>::Error>,
1076                    Self::DropConsistency,
1077                    Unbounded,
1078                    TotalOrder,
1079                    ExactlyOnce,
1080                >::collection_kind(
1081                )),
1082            },
1083        );
1084
1085        let membership_stream_ident = syn::Ident::new(
1086            &format!(
1087                "__hydro_deploy_many_{}_{}_membership",
1088                from.key, next_external_port_id
1089            ),
1090            Span::call_site(),
1091        );
1092        let membership_stream_expr: syn::Expr = parse_quote!(#membership_stream_ident);
1093        let raw_membership_stream: KeyedStream<
1094            u64,
1095            bool,
1096            Self::DropConsistency,
1097            Unbounded,
1098            TotalOrder,
1099            ExactlyOnce,
1100        > = KeyedStream::new(
1101            target_consistency.clone(),
1102            HydroNode::Source {
1103                source: HydroSource::Stream(membership_stream_expr.into()),
1104                metadata: target_consistency.new_node_metadata(KeyedStream::<
1105                    u64,
1106                    bool,
1107                    Self::DropConsistency,
1108                    Unbounded,
1109                    TotalOrder,
1110                    ExactlyOnce,
1111                >::collection_kind(
1112                )),
1113            },
1114        );
1115
1116        (
1117            ExternalBytesPort {
1118                process_key: from.key,
1119                port_id: next_external_port_id,
1120                _phantom: PhantomData,
1121            },
1122            raw_stream
1123                .flatten_ordered() // TODO(shadaj): this silently drops framing errors, decide on right defaults
1124                .into_keyed(),
1125            raw_membership_stream.map(q!(|join| {
1126                if join {
1127                    MembershipEvent::Joined
1128                } else {
1129                    MembershipEvent::Left
1130                }
1131            })),
1132            fwd_ref,
1133        )
1134    }
1135
1136    /// Establishes a server on this location to receive bidirectional connections from multiple
1137    /// external clients using bincode serialization.
1138    ///
1139    /// Unlike [`Location::bind_single_client_bincode`], this method supports multiple concurrent
1140    /// client connections. Each client is assigned a unique `u64` identifier.
1141    ///
1142    /// Returns:
1143    /// - A port handle for external processes to connect to
1144    /// - A keyed stream of incoming messages, keyed by client ID
1145    /// - A keyed stream of membership events (client joins/leaves), keyed by client ID
1146    /// - A handle to send outgoing messages, keyed by client ID
1147    ///
1148    /// # Type Parameters
1149    /// - `InT`: The type of incoming messages (must implement [`DeserializeOwned`])
1150    /// - `OutT`: The type of outgoing messages (must implement [`Serialize`])
1151    #[cfg(feature = "tokio")]
1152    #[expect(clippy::type_complexity, reason = "stream markers")]
1153    fn bidi_external_many_bincode<L, InT: DeserializeOwned, OutT: Serialize>(
1154        &self,
1155        from: &External<'_, L>,
1156    ) -> (
1157        ExternalBincodeBidi<InT, OutT, Many>,
1158        KeyedStream<u64, InT, Self::DropConsistency, Unbounded, TotalOrder, ExactlyOnce>,
1159        KeyedStream<
1160            u64,
1161            MembershipEvent,
1162            Self::DropConsistency,
1163            Unbounded,
1164            TotalOrder,
1165            ExactlyOnce,
1166        >,
1167        ForwardHandle<
1168            'a,
1169            KeyedStream<u64, OutT, Self::DropConsistency, Unbounded, NoOrder, ExactlyOnce>,
1170        >,
1171    )
1172    where
1173        Self: TopLevel<'a> + Sized,
1174    {
1175        let next_external_port_id = from.flow_state.borrow_mut().next_external_port();
1176
1177        let target_consistency = self.drop_consistency();
1178        let (fwd_ref, to_sink) = target_consistency.forward_ref::<KeyedStream<
1179            u64,
1180            OutT,
1181            Self::DropConsistency,
1182            Unbounded,
1183            NoOrder,
1184            ExactlyOnce,
1185        >>();
1186        let to_sink_input = Box::new(to_sink.entries().ir_node.replace(HydroNode::Placeholder));
1187        let mut flow_state_borrow = self.flow_state().borrow_mut();
1188
1189        let root = get_this_crate();
1190
1191        let out_t_type = quote_type::<OutT>();
1192        let ser_fn: syn::Expr = syn::parse_quote! {
1193            #root::runtime_support::stageleft::runtime_support::fn1_type_hint::<(u64, #out_t_type), _>(
1194                |(id, b)| (id, #root::runtime_support::bincode::serialize(&b).unwrap().into())
1195            )
1196        };
1197
1198        flow_state_borrow.push_root(HydroRoot::SendExternal {
1199            to_external_key: from.key,
1200            to_port_id: next_external_port_id,
1201            to_many: true,
1202            unpaired: false,
1203            serialize_fn: Some(ser_fn.into()),
1204            instantiate_fn: DebugInstantiate::Building,
1205            input: to_sink_input,
1206            op_metadata: HydroIrOpMetadata::new(),
1207        });
1208        drop(flow_state_borrow);
1209
1210        let in_t_type = quote_type::<InT>();
1211
1212        let deser_fn: syn::Expr = syn::parse_quote! {
1213            |res| {
1214                let (id, b) = res.unwrap();
1215                (id, #root::runtime_support::bincode::deserialize::<#in_t_type>(&b).unwrap())
1216            }
1217        };
1218
1219        let raw_stream: KeyedStream<
1220            u64,
1221            InT,
1222            Self::DropConsistency,
1223            Unbounded,
1224            TotalOrder,
1225            ExactlyOnce,
1226        > = KeyedStream::new(
1227            target_consistency.clone(),
1228            HydroNode::ExternalInput {
1229                from_external_key: from.key,
1230                from_port_id: next_external_port_id,
1231                from_many: true,
1232                codec_type: quote_type::<LengthDelimitedCodec>().into(),
1233                port_hint: NetworkHint::Auto,
1234                instantiate_fn: DebugInstantiate::Building,
1235                deserialize_fn: Some(deser_fn.into()),
1236                metadata: target_consistency.new_node_metadata(KeyedStream::<
1237                    u64,
1238                    InT,
1239                    Self::DropConsistency,
1240                    Unbounded,
1241                    TotalOrder,
1242                    ExactlyOnce,
1243                >::collection_kind(
1244                )),
1245            },
1246        );
1247
1248        let membership_stream_ident = syn::Ident::new(
1249            &format!(
1250                "__hydro_deploy_many_{}_{}_membership",
1251                from.key, next_external_port_id
1252            ),
1253            Span::call_site(),
1254        );
1255        let membership_stream_expr: syn::Expr = parse_quote!(#membership_stream_ident);
1256        let raw_membership_stream: KeyedStream<
1257            u64,
1258            bool,
1259            Self::DropConsistency,
1260            Unbounded,
1261            TotalOrder,
1262            ExactlyOnce,
1263        > = KeyedStream::new(
1264            target_consistency.clone(),
1265            HydroNode::Source {
1266                source: HydroSource::Stream(membership_stream_expr.into()),
1267                metadata: target_consistency.new_node_metadata(KeyedStream::<
1268                    u64,
1269                    bool,
1270                    Self::DropConsistency,
1271                    Unbounded,
1272                    TotalOrder,
1273                    ExactlyOnce,
1274                >::collection_kind(
1275                )),
1276            },
1277        );
1278
1279        (
1280            ExternalBincodeBidi {
1281                process_key: from.key,
1282                port_id: next_external_port_id,
1283                _phantom: PhantomData,
1284            },
1285            raw_stream,
1286            raw_membership_stream.map(q!(|join| {
1287                if join {
1288                    MembershipEvent::Joined
1289                } else {
1290                    MembershipEvent::Left
1291                }
1292            })),
1293            fwd_ref,
1294        )
1295    }
1296
1297    /// Bridges user-owned async code to the dataflow as a **bidirectional sidecar**.
1298    ///
1299    /// The closure is called once at startup and must return a
1300    /// `(Stream<InT>, Sink<OutT>)` pair. The framework reads from the stream
1301    /// (items flowing *into* the dataflow) and writes to the sink (items flowing
1302    /// *out* to the sidecar). The user controls buffering, backpressure, and
1303    /// internal lifecycle — Hydro only sees the stream/sink interface.
1304    ///
1305    /// This will hopefully make it easy to integrate hydro with existing frameworks,
1306    /// for example grpc code generated service endpoints.
1307    ///
1308    /// # Returns
1309    /// - A `Stream<InT>` carrying items from the sidecar into the dataflow.
1310    /// - A [`ForwardHandle`] expecting a `Stream<OutT>` that the user completes
1311    ///   with items destined for the sidecar.
1312    ///
1313    /// # Example
1314    ///
1315    /// ```rust
1316    /// # #[cfg(feature = "deploy")] {
1317    /// # use hydro_lang::prelude::*;
1318    /// # use futures::StreamExt;
1319    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1320    /// // Sidecar that echoes whatever it receives back into the dataflow.
1321    /// let (inbound, response_handle) = process.sidecar_bidi::<String, String, _>(q!(|| {
1322    ///     let (to_df_tx, to_df_rx) = tokio::sync::mpsc::channel::<String>(16);
1323    ///     let (from_df_tx, mut from_df_rx) = tokio::sync::mpsc::channel::<String>(16);
1324    ///
1325    ///     // Spawn the sidecar: echoes items from the dataflow back into it.
1326    ///     tokio::spawn(async move {
1327    ///         while let Some(msg) = from_df_rx.recv().await {
1328    ///             to_df_tx.send(msg).await.ok();
1329    ///         }
1330    ///     });
1331    ///
1332    ///     // Return the framework-facing ends (concrete types, no boxing needed).
1333    ///     let stream = tokio_stream::wrappers::ReceiverStream::new(to_df_rx);
1334    ///     let sink = tokio_util::sync::PollSender::new(from_df_tx);
1335    ///     (stream, sink)
1336    /// }));
1337    ///
1338    /// // Send "hello" into the sidecar via the response channel.
1339    /// let input = process.source_stream(q!(futures::stream::iter(vec!["hello".to_string()])));
1340    /// response_handle.complete(input);
1341    ///
1342    /// // The sidecar echoes it back — assert we get "hello" out.
1343    /// inbound
1344    /// # }, |mut stream| async move {
1345    /// #     assert_eq!(stream.next().await.unwrap(), "hello");
1346    /// # }));
1347    /// # }
1348    /// ```
1349    fn sidecar_bidi<InT: 'static, OutT: 'static, F>(
1350        &self,
1351        sidecar: impl QuotedWithContext<'a, F, Self>,
1352    ) -> (
1353        Stream<InT, Self, Unbounded, TotalOrder, ExactlyOnce>,
1354        ForwardHandle<'a, Stream<OutT, Self, Unbounded, NoOrder, ExactlyOnce>>,
1355    )
1356    where
1357        Self: Sized + TopLevel<'a>,
1358    {
1359        let location_key = Location::id(self).key();
1360
1361        let sidecar_id = self.flow_state().borrow_mut().next_sidecar_id();
1362        let (stream_ident, sink_ident) = sidecar_id.idents();
1363
1364        let sidecar_closure: syn::Expr = sidecar.splice_untyped_ctx(self);
1365        self.flow_state()
1366            .borrow_mut()
1367            .sidecars
1368            .push(crate::compile::builder::Sidecar::Bidi {
1369                location_key,
1370                sidecar_id,
1371                sidecar_closure: Box::new(sidecar_closure),
1372            });
1373
1374        // Inbound stream: reads from the stream returned by the sidecar closure
1375        let source_expr: syn::Expr = parse_quote! {
1376            #stream_ident
1377        };
1378        let inbound: Stream<InT, Self, Unbounded, TotalOrder, ExactlyOnce> = Stream::new(
1379            self.clone(),
1380            HydroNode::Source {
1381                source: HydroSource::Stream(source_expr.into()),
1382                metadata: self.new_node_metadata(Stream::<
1383                    InT,
1384                    Self,
1385                    Unbounded,  // TODO: maybe bounded sidecars are interesting..?
1386                    TotalOrder, // TODO: NoOrder..?
1387                    ExactlyOnce,
1388                >::collection_kind()),
1389            },
1390        );
1391
1392        // Outbound: forward_ref cycle feeding the sink returned by the sidecar closure
1393        let (fwd_ref, to_sink): (
1394            ForwardHandle<'a, Stream<OutT, Self, Unbounded, NoOrder, ExactlyOnce>>,
1395            Stream<OutT, Self, Unbounded, NoOrder, ExactlyOnce>,
1396        ) = self.forward_ref();
1397
1398        let sink_expr: syn::Expr = parse_quote! {
1399            #sink_ident
1400        };
1401
1402        let sink_input_ir = to_sink.ir_node.replace(HydroNode::Placeholder);
1403        self.flow_state()
1404            .borrow_mut()
1405            .try_push_root(HydroRoot::DestSink {
1406                sink: sink_expr.into(),
1407                input: Box::new(sink_input_ir),
1408                op_metadata: HydroIrOpMetadata::new(),
1409            });
1410
1411        (inbound, fwd_ref)
1412    }
1413
1414    /// Constructs a [`Singleton`] materialized at this location with the given static value.
1415    ///
1416    /// See also: [`Tick::singleton`], for creating a singleton _within_ a tick, which requires
1417    /// `T: Clone`.
1418    ///
1419    /// # Example
1420    /// ```rust
1421    /// # #[cfg(feature = "deploy")] {
1422    /// # use hydro_lang::prelude::*;
1423    /// # use futures::StreamExt;
1424    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1425    /// let singleton = process.singleton(q!(5));
1426    /// # singleton.into_stream()
1427    /// # }, |mut stream| async move {
1428    /// // 5
1429    /// # assert_eq!(stream.next().await.unwrap(), 5);
1430    /// # }));
1431    /// # }
1432    /// ```
1433    fn singleton<T>(
1434        &self,
1435        e: impl QuotedWithContext<'a, T, Self>,
1436    ) -> Singleton<T, Self::DropConsistency, Bounded>
1437    where
1438        Self: Sized,
1439    {
1440        let e = e.splice_untyped_ctx(self);
1441
1442        let target_location = self.drop_consistency();
1443        Singleton::new(
1444            target_location.clone(),
1445            HydroNode::SingletonSource {
1446                value: e.into(),
1447                first_tick_only: false,
1448                metadata: target_location.new_node_metadata(Singleton::<
1449                    T,
1450                    Self::DropConsistency,
1451                    Bounded,
1452                >::collection_kind()),
1453            },
1454        )
1455    }
1456
1457    /// Constructs a [`Singleton`] by resolving an async [`Future`] to completion.
1458    ///
1459    /// This is a convenience method equivalent to
1460    /// `self.singleton(future_expr).resolve_future_blocking()`, which is a common
1461    /// pattern when initializing a singleton from an async computation.
1462    ///
1463    /// # Example
1464    /// ```rust
1465    /// # #[cfg(feature = "deploy")] {
1466    /// # use hydro_lang::prelude::*;
1467    /// # use futures::StreamExt;
1468    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1469    /// let singleton = process.singleton_future(q!(async { 42 }));
1470    /// singleton.into_stream()
1471    /// # }, |mut stream| async move {
1472    /// // 42
1473    /// # assert_eq!(stream.next().await.unwrap(), 42);
1474    /// # }));
1475    /// # }
1476    /// ```
1477    ///
1478    /// [`Future`]: std::future::Future
1479    fn singleton_future<F>(
1480        &self,
1481        e: impl QuotedWithContext<'a, F, Self>,
1482    ) -> Singleton<F::Output, Self::DropConsistency, Bounded>
1483    where
1484        F: Future,
1485        Self: Sized,
1486    {
1487        self.singleton(e).resolve_future_blocking()
1488    }
1489
1490    /// Generates a stream that emits `()` at a fixed interval.
1491    ///
1492    /// The first tick completes immediately. Missed ticks will be scheduled
1493    /// as soon as possible.
1494    ///
1495    /// Because this only emits `()`, the non-determinism of *when* events fire
1496    /// is captured by the `AtLeastOnce` retry semantics downstream, so no
1497    /// [`NonDet`] guard is required.
1498    #[cfg(feature = "tokio")]
1499    fn source_interval(
1500        &self,
1501        interval: impl QuotedWithContext<'a, Duration, Self> + Copy + 'a,
1502    ) -> Stream<(), Self, Unbounded, TotalOrder, ExactlyOnce>
1503    where
1504        Self: TopLevel<'a> + Sized,
1505    {
1506        self.source_stream(q!(tokio_stream::StreamExt::map(
1507            tokio_stream::wrappers::IntervalStream::new(tokio::time::interval(interval)),
1508            |_| ()
1509        )))
1510        .assert_has_consistency_of_trusted(
1511            manual_proof!(/** interval does not reveal timestamps */),
1512        )
1513    }
1514
1515    /// Generates a stream that emits `()` at a fixed interval, after an
1516    /// initial delay.
1517    ///
1518    /// Because this only emits `()`, the non-determinism of *when* events fire
1519    /// is captured by the `AtLeastOnce` retry semantics downstream, so no
1520    /// [`NonDet`] guard is required.
1521    #[cfg(feature = "tokio")]
1522    fn source_interval_delayed(
1523        &self,
1524        delay: impl QuotedWithContext<'a, Duration, Self> + Copy + 'a,
1525        interval: impl QuotedWithContext<'a, Duration, Self> + Copy + 'a,
1526    ) -> Stream<(), Self, Unbounded, TotalOrder, ExactlyOnce>
1527    where
1528        Self: TopLevel<'a> + Sized,
1529    {
1530        self.source_stream(q!(tokio_stream::StreamExt::map(
1531            tokio_stream::wrappers::IntervalStream::new(tokio::time::interval_at(
1532                tokio::time::Instant::now() + delay,
1533                interval,
1534            )),
1535            |_| ()
1536        )))
1537        .assert_has_consistency_of_trusted(
1538            manual_proof!(/** interval does not reveal timestamps */),
1539        )
1540    }
1541
1542    /// Creates a forward reference, allowing a stream to be used before its source is defined.
1543    ///
1544    /// Returns a `(handle, placeholder)` pair. Use the placeholder in the dataflow graph,
1545    /// then call `handle.complete(actual_stream)` to wire in the real source.
1546    ///
1547    /// This is useful for mutually-dependent dataflows or when the definition order
1548    /// doesn't match the data flow direction. For feedback loops, prefer [`Tick::cycle`]
1549    /// instead, which automatically defers values by one tick.
1550    ///
1551    /// # Panics
1552    /// Panics if the forward reference creates a synchronous cycle (i.e., the completed
1553    /// stream transitively depends on the placeholder without a `defer_tick` or network
1554    /// hop in between).
1555    ///
1556    /// # Example
1557    /// ```rust
1558    /// # #[cfg(feature = "deploy")] {
1559    /// # use hydro_lang::prelude::*;
1560    /// # use hydro_lang::live_collections::stream::NoOrder;
1561    /// # use futures::StreamExt;
1562    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1563    /// // Create a forward reference to define a stream that will be completed later
1564    /// let (complete, forward_stream) = process.forward_ref::<Stream<i32, _, _, NoOrder>>();
1565    ///
1566    /// // Use the forward reference as input to another computation
1567    /// let output: Stream<_, _, _, NoOrder> = forward_stream.map(q!(|x| x * 2));
1568    ///
1569    /// // Complete the forward reference with the actual source
1570    /// let source: Stream<_, _, Unbounded> = process.source_iter(q!([1, 2, 3])).into();
1571    /// complete.complete(source);
1572    /// output
1573    /// # }, |mut stream| async move {
1574    /// // 2, 4, 6
1575    /// # assert_eq!(stream.next().await.unwrap(), 2);
1576    /// # assert_eq!(stream.next().await.unwrap(), 4);
1577    /// # assert_eq!(stream.next().await.unwrap(), 6);
1578    /// # }));
1579    /// # }
1580    /// ```
1581    fn forward_ref<S>(&self) -> (ForwardHandle<'a, S>, S)
1582    where
1583        S: CycleCollection<'a, ForwardRef, Location = Self>,
1584    {
1585        let cycle_id = self.flow_state().borrow_mut().next_cycle_id();
1586        (
1587            ForwardHandle::new(cycle_id, Location::id(self)),
1588            S::create_source(cycle_id, self.clone()),
1589        )
1590    }
1591}
1592
1593#[cfg(feature = "deploy")]
1594#[cfg(test)]
1595mod tests {
1596    use std::collections::HashSet;
1597
1598    use futures::{SinkExt, StreamExt};
1599    use hydro_deploy::Deployment;
1600    use stageleft::q;
1601    use tokio_util::codec::LengthDelimitedCodec;
1602
1603    use crate::compile::builder::FlowBuilder;
1604    use crate::live_collections::stream::{ExactlyOnce, TotalOrder};
1605    use crate::location::{Location, NetworkHint};
1606    use crate::nondet::nondet;
1607
1608    #[tokio::test]
1609    async fn top_level_singleton_replay_cardinality() {
1610        let mut deployment = Deployment::new();
1611
1612        let mut flow = FlowBuilder::new();
1613        let node = flow.process::<()>();
1614        let external = flow.external::<()>();
1615
1616        let (in_port, input) =
1617            node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
1618        let singleton = node.singleton(q!(123));
1619        let tick = node.tick();
1620        let out = input
1621            .batch(&tick, nondet!(/** test */))
1622            .cross_singleton(singleton.clone().snapshot(&tick, nondet!(/** test */)))
1623            .cross_singleton(
1624                singleton
1625                    .snapshot(&tick, nondet!(/** test */))
1626                    .into_stream()
1627                    .count(),
1628            )
1629            .all_ticks()
1630            .send_bincode_external(&external);
1631
1632        let nodes = flow
1633            .with_process(&node, deployment.Localhost())
1634            .with_external(&external, deployment.Localhost())
1635            .deploy(&mut deployment);
1636
1637        deployment.deploy().await.unwrap();
1638
1639        let mut external_in = nodes.connect(in_port).await;
1640        let mut external_out = nodes.connect(out).await;
1641
1642        deployment.start().await.unwrap();
1643
1644        external_in.send(1).await.unwrap();
1645        assert_eq!(external_out.next().await.unwrap(), ((1, 123), 1));
1646
1647        external_in.send(2).await.unwrap();
1648        assert_eq!(external_out.next().await.unwrap(), ((2, 123), 1));
1649    }
1650
1651    #[tokio::test]
1652    async fn tick_singleton_replay_cardinality() {
1653        let mut deployment = Deployment::new();
1654
1655        let mut flow = FlowBuilder::new();
1656        let node = flow.process::<()>();
1657        let external = flow.external::<()>();
1658
1659        let (in_port, input) =
1660            node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
1661        let tick = node.tick();
1662        let singleton = tick.singleton(q!(123));
1663        let out = input
1664            .batch(&tick, nondet!(/** test */))
1665            .cross_singleton(singleton.clone())
1666            .cross_singleton(singleton.into_stream().count())
1667            .all_ticks()
1668            .send_bincode_external(&external);
1669
1670        let nodes = flow
1671            .with_process(&node, deployment.Localhost())
1672            .with_external(&external, deployment.Localhost())
1673            .deploy(&mut deployment);
1674
1675        deployment.deploy().await.unwrap();
1676
1677        let mut external_in = nodes.connect(in_port).await;
1678        let mut external_out = nodes.connect(out).await;
1679
1680        deployment.start().await.unwrap();
1681
1682        external_in.send(1).await.unwrap();
1683        assert_eq!(external_out.next().await.unwrap(), ((1, 123), 1));
1684
1685        external_in.send(2).await.unwrap();
1686        assert_eq!(external_out.next().await.unwrap(), ((2, 123), 1));
1687    }
1688
1689    #[tokio::test]
1690    async fn external_bytes() {
1691        let mut deployment = Deployment::new();
1692
1693        let mut flow = FlowBuilder::new();
1694        let first_node = flow.process::<()>();
1695        let external = flow.external::<()>();
1696
1697        let (in_port, input) = first_node.source_external_bytes(&external);
1698        let out = input.send_bincode_external(&external);
1699
1700        let nodes = flow
1701            .with_process(&first_node, deployment.Localhost())
1702            .with_external(&external, deployment.Localhost())
1703            .deploy(&mut deployment);
1704
1705        deployment.deploy().await.unwrap();
1706
1707        let mut external_in = nodes.connect(in_port).await.1;
1708        let mut external_out = nodes.connect(out).await;
1709
1710        deployment.start().await.unwrap();
1711
1712        external_in.send(vec![1, 2, 3].into()).await.unwrap();
1713
1714        assert_eq!(external_out.next().await.unwrap(), vec![1, 2, 3]);
1715    }
1716
1717    #[tokio::test]
1718    async fn multi_external_source() {
1719        let mut deployment = Deployment::new();
1720
1721        let mut flow = FlowBuilder::new();
1722        let first_node = flow.process::<()>();
1723        let external = flow.external::<()>();
1724
1725        let (in_port, input, _membership, complete_sink) =
1726            first_node.bidi_external_many_bincode(&external);
1727        let out = input.entries().send_bincode_external(&external);
1728        complete_sink.complete(
1729            first_node
1730                .source_iter::<(u64, ()), _>(q!([]))
1731                .into_keyed()
1732                .weaken_ordering(),
1733        );
1734
1735        let nodes = flow
1736            .with_process(&first_node, deployment.Localhost())
1737            .with_external(&external, deployment.Localhost())
1738            .deploy(&mut deployment);
1739
1740        deployment.deploy().await.unwrap();
1741
1742        let (_, mut external_in_1) = nodes.connect_bincode(in_port.clone()).await;
1743        let (_, mut external_in_2) = nodes.connect_bincode(in_port).await;
1744        let external_out = nodes.connect(out).await;
1745
1746        deployment.start().await.unwrap();
1747
1748        external_in_1.send(123).await.unwrap();
1749        external_in_2.send(456).await.unwrap();
1750
1751        assert_eq!(
1752            external_out.take(2).collect::<HashSet<_>>().await,
1753            vec![(0, 123), (1, 456)].into_iter().collect()
1754        );
1755    }
1756
1757    #[tokio::test]
1758    async fn second_connection_only_multi_source() {
1759        let mut deployment = Deployment::new();
1760
1761        let mut flow = FlowBuilder::new();
1762        let first_node = flow.process::<()>();
1763        let external = flow.external::<()>();
1764
1765        let (in_port, input, _membership, complete_sink) =
1766            first_node.bidi_external_many_bincode(&external);
1767        let out = input.entries().send_bincode_external(&external);
1768        complete_sink.complete(
1769            first_node
1770                .source_iter::<(u64, ()), _>(q!([]))
1771                .into_keyed()
1772                .weaken_ordering(),
1773        );
1774
1775        let nodes = flow
1776            .with_process(&first_node, deployment.Localhost())
1777            .with_external(&external, deployment.Localhost())
1778            .deploy(&mut deployment);
1779
1780        deployment.deploy().await.unwrap();
1781
1782        // intentionally skipped to test stream waking logic
1783        let (_, mut _external_in_1) = nodes.connect_bincode(in_port.clone()).await;
1784        let (_, mut external_in_2) = nodes.connect_bincode(in_port).await;
1785        let mut external_out = nodes.connect(out).await;
1786
1787        deployment.start().await.unwrap();
1788
1789        external_in_2.send(456).await.unwrap();
1790
1791        assert_eq!(external_out.next().await.unwrap(), (1, 456));
1792    }
1793
1794    #[tokio::test]
1795    async fn multi_external_bytes() {
1796        let mut deployment = Deployment::new();
1797
1798        let mut flow = FlowBuilder::new();
1799        let first_node = flow.process::<()>();
1800        let external = flow.external::<()>();
1801
1802        let (in_port, input, _membership, complete_sink) = first_node
1803            .bidi_external_many_bytes::<_, _, LengthDelimitedCodec>(&external, NetworkHint::Auto);
1804        let out = input.entries().send_bincode_external(&external);
1805        complete_sink.complete(
1806            first_node
1807                .source_iter(q!([]))
1808                .into_keyed()
1809                .weaken_ordering(),
1810        );
1811
1812        let nodes = flow
1813            .with_process(&first_node, deployment.Localhost())
1814            .with_external(&external, deployment.Localhost())
1815            .deploy(&mut deployment);
1816
1817        deployment.deploy().await.unwrap();
1818
1819        let mut external_in_1 = nodes.connect(in_port.clone()).await.1;
1820        let mut external_in_2 = nodes.connect(in_port).await.1;
1821        let external_out = nodes.connect(out).await;
1822
1823        deployment.start().await.unwrap();
1824
1825        external_in_1.send(vec![1, 2, 3].into()).await.unwrap();
1826        external_in_2.send(vec![4, 5].into()).await.unwrap();
1827
1828        assert_eq!(
1829            external_out.take(2).collect::<HashSet<_>>().await,
1830            vec![
1831                (0, (&[1u8, 2, 3] as &[u8]).into()),
1832                (1, (&[4u8, 5] as &[u8]).into())
1833            ]
1834            .into_iter()
1835            .collect()
1836        );
1837    }
1838
1839    #[tokio::test]
1840    async fn single_client_external_bytes() {
1841        let mut deployment = Deployment::new();
1842        let mut flow = FlowBuilder::new();
1843        let first_node = flow.process::<()>();
1844        let external = flow.external::<()>();
1845        let (port, input, complete_sink) = first_node
1846            .bind_single_client::<_, _, LengthDelimitedCodec>(&external, NetworkHint::Auto);
1847        complete_sink.complete(input.map(q!(|data| {
1848            let mut resp: Vec<u8> = data.into();
1849            resp.push(42);
1850            resp.into() // : Bytes
1851        })));
1852
1853        let nodes = flow
1854            .with_process(&first_node, deployment.Localhost())
1855            .with_external(&external, deployment.Localhost())
1856            .deploy(&mut deployment);
1857
1858        deployment.deploy().await.unwrap();
1859        deployment.start().await.unwrap();
1860
1861        let (mut external_out, mut external_in) = nodes.connect(port).await;
1862
1863        external_in.send(vec![1, 2, 3].into()).await.unwrap();
1864        assert_eq!(
1865            external_out.next().await.unwrap().unwrap(),
1866            vec![1, 2, 3, 42]
1867        );
1868    }
1869
1870    #[tokio::test]
1871    async fn echo_external_bytes() {
1872        let mut deployment = Deployment::new();
1873
1874        let mut flow = FlowBuilder::new();
1875        let first_node = flow.process::<()>();
1876        let external = flow.external::<()>();
1877
1878        let (port, input, _membership, complete_sink) = first_node
1879            .bidi_external_many_bytes::<_, _, LengthDelimitedCodec>(&external, NetworkHint::Auto);
1880        complete_sink
1881            .complete(input.map(q!(|bytes| { bytes.into_iter().map(|x| x + 1).collect() })));
1882
1883        let nodes = flow
1884            .with_process(&first_node, deployment.Localhost())
1885            .with_external(&external, deployment.Localhost())
1886            .deploy(&mut deployment);
1887
1888        deployment.deploy().await.unwrap();
1889
1890        let (mut external_out_1, mut external_in_1) = nodes.connect(port.clone()).await;
1891        let (mut external_out_2, mut external_in_2) = nodes.connect(port).await;
1892
1893        deployment.start().await.unwrap();
1894
1895        external_in_1.send(vec![1, 2, 3].into()).await.unwrap();
1896        external_in_2.send(vec![4, 5].into()).await.unwrap();
1897
1898        assert_eq!(external_out_1.next().await.unwrap().unwrap(), vec![2, 3, 4]);
1899        assert_eq!(external_out_2.next().await.unwrap().unwrap(), vec![5, 6]);
1900    }
1901
1902    #[tokio::test]
1903    async fn echo_external_bincode() {
1904        let mut deployment = Deployment::new();
1905
1906        let mut flow = FlowBuilder::new();
1907        let first_node = flow.process::<()>();
1908        let external = flow.external::<()>();
1909
1910        let (port, input, _membership, complete_sink) =
1911            first_node.bidi_external_many_bincode(&external);
1912        complete_sink.complete(input.map(q!(|text: String| { text.to_uppercase() })));
1913
1914        let nodes = flow
1915            .with_process(&first_node, deployment.Localhost())
1916            .with_external(&external, deployment.Localhost())
1917            .deploy(&mut deployment);
1918
1919        deployment.deploy().await.unwrap();
1920
1921        let (mut external_out_1, mut external_in_1) = nodes.connect_bincode(port.clone()).await;
1922        let (mut external_out_2, mut external_in_2) = nodes.connect_bincode(port).await;
1923
1924        deployment.start().await.unwrap();
1925
1926        external_in_1.send("hi".to_owned()).await.unwrap();
1927        external_in_2.send("hello".to_owned()).await.unwrap();
1928
1929        assert_eq!(external_out_1.next().await.unwrap(), "HI");
1930        assert_eq!(external_out_2.next().await.unwrap(), "HELLO");
1931    }
1932
1933    #[tokio::test]
1934    async fn closure_location_name() {
1935        let mut deployment = Deployment::new();
1936        let mut flow = FlowBuilder::new();
1937
1938        enum ClosureProcess {}
1939
1940        let node = flow.process::<ClosureProcess>();
1941        let external = flow.external::<()>();
1942
1943        let (in_port, input) =
1944            node.source_external_bincode::<_, i32, TotalOrder, ExactlyOnce>(&external);
1945        let out = input.send_bincode_external(&external);
1946
1947        let nodes = flow
1948            .with_process(&node, deployment.Localhost())
1949            .with_external(&external, deployment.Localhost())
1950            .deploy(&mut deployment);
1951
1952        deployment.deploy().await.unwrap();
1953
1954        let mut external_in = nodes.connect(in_port).await;
1955        let mut external_out = nodes.connect(out).await;
1956
1957        deployment.start().await.unwrap();
1958
1959        external_in.send(42).await.unwrap();
1960        assert_eq!(external_out.next().await.unwrap(), 42);
1961    }
1962}