Skip to main content

hydro_lang/deploy/
deploy_graph.rs

1//! Deployment backend for Hydro that uses [`hydro_deploy`] to provision and launch services.
2
3use std::cell::RefCell;
4use std::collections::HashMap;
5use std::future::Future;
6use std::io::Error;
7use std::pin::Pin;
8use std::rc::Rc;
9use std::sync::Arc;
10
11use bytes::{Bytes, BytesMut};
12use dfir_lang::graph::{AsCodeOptions, DfirGraph};
13use futures::{Sink, SinkExt, Stream, StreamExt};
14use hydro_deploy::custom_service::CustomClientPort;
15use hydro_deploy::rust_crate::RustCrateService;
16use hydro_deploy::rust_crate::ports::{DemuxSink, RustCrateSink, RustCrateSource, TaggedSource};
17use hydro_deploy::rust_crate::tracing_options::TracingOptions;
18use hydro_deploy::{CustomService, Deployment, Host, RustCrate};
19use hydro_deploy_integration::{ConnectedSink, ConnectedSource};
20use nameof::name_of;
21use proc_macro2::Span;
22use serde::Serialize;
23use serde::de::DeserializeOwned;
24use slotmap::SparseSecondaryMap;
25use stageleft::{QuotedWithContext, RuntimeData};
26use syn::parse_quote;
27
28use super::deploy_runtime::*;
29use crate::compile::builder::ExternalPortId;
30use crate::compile::deploy_provider::{
31    ClusterSpec, Deploy, ExternalSpec, IntoProcessSpec, Node, ProcessSpec, RegisterPort,
32};
33use crate::compile::trybuild::generate::{
34    HYDRO_RUNTIME_FEATURES, LinkingMode, create_graph_trybuild,
35};
36use crate::location::dynamic::LocationId;
37use crate::location::member_id::TaglessMemberId;
38use crate::location::{LocationKey, MembershipEvent, NetworkHint};
39use crate::staging_util::get_this_crate;
40
41/// Deployment backend that uses [`hydro_deploy`] for provisioning and launching.
42///
43/// Automatically used when you call [`crate::compile::builder::FlowBuilder::deploy`] and pass in
44/// an `&mut` reference to [`hydro_deploy::Deployment`] as the deployment context.
45pub enum HydroDeploy {}
46
47impl<'a> Deploy<'a> for HydroDeploy {
48    /// Map from Cluster location ID to member IDs.
49    type Meta = SparseSecondaryMap<LocationKey, Vec<TaglessMemberId>>;
50    type InstantiateEnv = Deployment;
51
52    type Process = DeployNode;
53    type Cluster = DeployCluster;
54    type External = DeployExternal;
55
56    fn o2o_sink_source(
57        _env: &mut Self::InstantiateEnv,
58        _p1: &Self::Process,
59        p1_port: &<Self::Process as Node>::Port,
60        _p2: &Self::Process,
61        p2_port: &<Self::Process as Node>::Port,
62        _name: Option<&str>,
63        networking_info: &crate::networking::NetworkingInfo,
64        _external_types: Option<(&syn::Type, &syn::Type)>,
65    ) -> (syn::Expr, syn::Expr) {
66        match networking_info {
67            crate::networking::NetworkingInfo::Tcp {
68                fault: crate::networking::TcpFault::FailStop,
69            } => {}
70            _ => panic!("Unsupported networking info: {:?}", networking_info),
71        }
72        let p1_port = p1_port.as_str();
73        let p2_port = p2_port.as_str();
74        deploy_o2o(
75            RuntimeData::new("__hydro_lang_trybuild_cli"),
76            p1_port,
77            p2_port,
78        )
79    }
80
81    fn o2o_connect(
82        p1: &Self::Process,
83        p1_port: &<Self::Process as Node>::Port,
84        p2: &Self::Process,
85        p2_port: &<Self::Process as Node>::Port,
86    ) -> Box<dyn FnOnce()> {
87        let p1 = p1.clone();
88        let p1_port = p1_port.clone();
89        let p2 = p2.clone();
90        let p2_port = p2_port.clone();
91
92        Box::new(move || {
93            let self_underlying_borrow = p1.underlying.borrow();
94            let self_underlying = self_underlying_borrow.as_ref().unwrap();
95            let source_port = self_underlying.get_port(p1_port.clone());
96
97            let other_underlying_borrow = p2.underlying.borrow();
98            let other_underlying = other_underlying_borrow.as_ref().unwrap();
99            let recipient_port = other_underlying.get_port(p2_port.clone());
100
101            source_port.send_to(&recipient_port)
102        })
103    }
104
105    fn o2m_sink_source(
106        _env: &mut Self::InstantiateEnv,
107        _p1: &Self::Process,
108        p1_port: &<Self::Process as Node>::Port,
109        _c2: &Self::Cluster,
110        c2_port: &<Self::Cluster as Node>::Port,
111        _name: Option<&str>,
112        networking_info: &crate::networking::NetworkingInfo,
113        _external_types: Option<(&syn::Type, &syn::Type)>,
114    ) -> (syn::Expr, syn::Expr) {
115        match networking_info {
116            crate::networking::NetworkingInfo::Tcp {
117                fault: crate::networking::TcpFault::FailStop,
118            } => {}
119            _ => panic!("Unsupported networking info: {:?}", networking_info),
120        }
121        let p1_port = p1_port.as_str();
122        let c2_port = c2_port.as_str();
123        deploy_o2m(
124            RuntimeData::new("__hydro_lang_trybuild_cli"),
125            p1_port,
126            c2_port,
127        )
128    }
129
130    fn o2m_connect(
131        p1: &Self::Process,
132        p1_port: &<Self::Process as Node>::Port,
133        c2: &Self::Cluster,
134        c2_port: &<Self::Cluster as Node>::Port,
135    ) -> Box<dyn FnOnce()> {
136        let p1 = p1.clone();
137        let p1_port = p1_port.clone();
138        let c2 = c2.clone();
139        let c2_port = c2_port.clone();
140
141        Box::new(move || {
142            let self_underlying_borrow = p1.underlying.borrow();
143            let self_underlying = self_underlying_borrow.as_ref().unwrap();
144            let source_port = self_underlying.get_port(p1_port.clone());
145
146            let recipient_port = DemuxSink {
147                demux: c2
148                    .members
149                    .borrow()
150                    .iter()
151                    .enumerate()
152                    .map(|(id, c)| {
153                        (
154                            id as u32,
155                            Arc::new(c.underlying.get_port(c2_port.clone()))
156                                as Arc<dyn RustCrateSink + 'static>,
157                        )
158                    })
159                    .collect(),
160            };
161
162            source_port.send_to(&recipient_port)
163        })
164    }
165
166    fn m2o_sink_source(
167        _env: &mut Self::InstantiateEnv,
168        _c1: &Self::Cluster,
169        c1_port: &<Self::Cluster as Node>::Port,
170        _p2: &Self::Process,
171        p2_port: &<Self::Process as Node>::Port,
172        _name: Option<&str>,
173        networking_info: &crate::networking::NetworkingInfo,
174        _external_types: Option<(&syn::Type, &syn::Type)>,
175    ) -> (syn::Expr, syn::Expr) {
176        match networking_info {
177            crate::networking::NetworkingInfo::Tcp {
178                fault: crate::networking::TcpFault::FailStop,
179            } => {}
180            _ => panic!("Unsupported networking info: {:?}", networking_info),
181        }
182        let c1_port = c1_port.as_str();
183        let p2_port = p2_port.as_str();
184        deploy_m2o(
185            RuntimeData::new("__hydro_lang_trybuild_cli"),
186            c1_port,
187            p2_port,
188        )
189    }
190
191    fn m2o_connect(
192        c1: &Self::Cluster,
193        c1_port: &<Self::Cluster as Node>::Port,
194        p2: &Self::Process,
195        p2_port: &<Self::Process as Node>::Port,
196    ) -> Box<dyn FnOnce()> {
197        let c1 = c1.clone();
198        let c1_port = c1_port.clone();
199        let p2 = p2.clone();
200        let p2_port = p2_port.clone();
201
202        Box::new(move || {
203            let other_underlying_borrow = p2.underlying.borrow();
204            let other_underlying = other_underlying_borrow.as_ref().unwrap();
205            let recipient_port = other_underlying.get_port(p2_port.clone()).merge();
206
207            for (i, node) in c1.members.borrow().iter().enumerate() {
208                let source_port = node.underlying.get_port(c1_port.clone());
209
210                TaggedSource {
211                    source: Arc::new(source_port),
212                    tag: i as u32,
213                }
214                .send_to(&recipient_port);
215            }
216        })
217    }
218
219    fn m2m_sink_source(
220        _env: &mut Self::InstantiateEnv,
221        _c1: &Self::Cluster,
222        c1_port: &<Self::Cluster as Node>::Port,
223        _c2: &Self::Cluster,
224        c2_port: &<Self::Cluster as Node>::Port,
225        _name: Option<&str>,
226        networking_info: &crate::networking::NetworkingInfo,
227        _external_types: Option<(&syn::Type, &syn::Type)>,
228    ) -> (syn::Expr, syn::Expr) {
229        match networking_info {
230            crate::networking::NetworkingInfo::Tcp {
231                fault: crate::networking::TcpFault::FailStop,
232            } => {}
233            _ => panic!("Unsupported networking info: {:?}", networking_info),
234        }
235        let c1_port = c1_port.as_str();
236        let c2_port = c2_port.as_str();
237        deploy_m2m(
238            RuntimeData::new("__hydro_lang_trybuild_cli"),
239            c1_port,
240            c2_port,
241        )
242    }
243
244    fn m2m_connect(
245        c1: &Self::Cluster,
246        c1_port: &<Self::Cluster as Node>::Port,
247        c2: &Self::Cluster,
248        c2_port: &<Self::Cluster as Node>::Port,
249    ) -> Box<dyn FnOnce()> {
250        let c1 = c1.clone();
251        let c1_port = c1_port.clone();
252        let c2 = c2.clone();
253        let c2_port = c2_port.clone();
254
255        Box::new(move || {
256            for (i, sender) in c1.members.borrow().iter().enumerate() {
257                let source_port = sender.underlying.get_port(c1_port.clone());
258
259                let recipient_port = DemuxSink {
260                    demux: c2
261                        .members
262                        .borrow()
263                        .iter()
264                        .enumerate()
265                        .map(|(id, c)| {
266                            (
267                                id as u32,
268                                Arc::new(c.underlying.get_port(c2_port.clone()).merge())
269                                    as Arc<dyn RustCrateSink + 'static>,
270                            )
271                        })
272                        .collect(),
273                };
274
275                TaggedSource {
276                    source: Arc::new(source_port),
277                    tag: i as u32,
278                }
279                .send_to(&recipient_port);
280            }
281        })
282    }
283
284    fn e2o_many_source(
285        extra_stmts: &mut Vec<syn::Stmt>,
286        _p2: &Self::Process,
287        p2_port: &<Self::Process as Node>::Port,
288        codec_type: &syn::Type,
289        shared_handle: String,
290    ) -> syn::Expr {
291        let connect_ident = syn::Ident::new(
292            &format!("__hydro_deploy_many_{}_connect", shared_handle),
293            Span::call_site(),
294        );
295        let source_ident = syn::Ident::new(
296            &format!("__hydro_deploy_many_{}_source", shared_handle),
297            Span::call_site(),
298        );
299        let sink_ident = syn::Ident::new(
300            &format!("__hydro_deploy_many_{}_sink", shared_handle),
301            Span::call_site(),
302        );
303        let membership_ident = syn::Ident::new(
304            &format!("__hydro_deploy_many_{}_membership", shared_handle),
305            Span::call_site(),
306        );
307
308        let root = get_this_crate();
309
310        extra_stmts.push(syn::parse_quote! {
311            let #connect_ident = __hydro_lang_trybuild_cli
312                .port(#p2_port)
313                .connect::<#root::runtime_support::hydro_deploy_integration::multi_connection::ConnectedMultiConnection<_, _, #codec_type>>();
314        });
315
316        extra_stmts.push(syn::parse_quote! {
317            let #source_ident = #connect_ident.source;
318        });
319
320        extra_stmts.push(syn::parse_quote! {
321            let #sink_ident = #connect_ident.sink;
322        });
323
324        extra_stmts.push(syn::parse_quote! {
325            let #membership_ident = #connect_ident.membership;
326        });
327
328        parse_quote!(#source_ident)
329    }
330
331    fn e2o_many_sink(shared_handle: String) -> syn::Expr {
332        let sink_ident = syn::Ident::new(
333            &format!("__hydro_deploy_many_{}_sink", shared_handle),
334            Span::call_site(),
335        );
336        parse_quote!(#sink_ident)
337    }
338
339    fn e2o_source(
340        extra_stmts: &mut Vec<syn::Stmt>,
341        _p1: &Self::External,
342        _p1_port: &<Self::External as Node>::Port,
343        _p2: &Self::Process,
344        p2_port: &<Self::Process as Node>::Port,
345        codec_type: &syn::Type,
346        shared_handle: String,
347    ) -> syn::Expr {
348        let connect_ident = syn::Ident::new(
349            &format!("__hydro_deploy_{}_connect", shared_handle),
350            Span::call_site(),
351        );
352        let source_ident = syn::Ident::new(
353            &format!("__hydro_deploy_{}_source", shared_handle),
354            Span::call_site(),
355        );
356        let sink_ident = syn::Ident::new(
357            &format!("__hydro_deploy_{}_sink", shared_handle),
358            Span::call_site(),
359        );
360
361        let root = get_this_crate();
362
363        extra_stmts.push(syn::parse_quote! {
364            let #connect_ident = __hydro_lang_trybuild_cli
365                .port(#p2_port)
366                .connect::<#root::runtime_support::hydro_deploy_integration::single_connection::ConnectedSingleConnection<_, _, #codec_type>>();
367        });
368
369        extra_stmts.push(syn::parse_quote! {
370            let #source_ident = #connect_ident.source;
371        });
372
373        extra_stmts.push(syn::parse_quote! {
374            let #sink_ident = #connect_ident.sink;
375        });
376
377        parse_quote!(#source_ident)
378    }
379
380    fn e2o_connect(
381        p1: &Self::External,
382        p1_port: &<Self::External as Node>::Port,
383        p2: &Self::Process,
384        p2_port: &<Self::Process as Node>::Port,
385        _many: bool,
386        server_hint: NetworkHint,
387    ) -> Box<dyn FnOnce()> {
388        let p1 = p1.clone();
389        let p1_port = p1_port.clone();
390        let p2 = p2.clone();
391        let p2_port = p2_port.clone();
392
393        Box::new(move || {
394            let self_underlying_borrow = p1.underlying.borrow();
395            let self_underlying = self_underlying_borrow.as_ref().unwrap();
396            let source_port = self_underlying.declare_many_client();
397
398            let other_underlying_borrow = p2.underlying.borrow();
399            let other_underlying = other_underlying_borrow.as_ref().unwrap();
400            let recipient_port = other_underlying.get_port_with_hint(
401                p2_port.clone(),
402                match server_hint {
403                    NetworkHint::Auto => hydro_deploy::PortNetworkHint::Auto,
404                    NetworkHint::TcpPort(p) => hydro_deploy::PortNetworkHint::TcpPort(p),
405                },
406            );
407
408            source_port.send_to(&recipient_port);
409
410            p1.client_ports
411                .borrow_mut()
412                .insert(p1_port.clone(), source_port);
413        })
414    }
415
416    fn o2e_sink(
417        _p1: &Self::Process,
418        _p1_port: &<Self::Process as Node>::Port,
419        _p2: &Self::External,
420        _p2_port: &<Self::External as Node>::Port,
421        shared_handle: String,
422    ) -> syn::Expr {
423        let sink_ident = syn::Ident::new(
424            &format!("__hydro_deploy_{}_sink", shared_handle),
425            Span::call_site(),
426        );
427        parse_quote!(#sink_ident)
428    }
429
430    fn cluster_ids(
431        of_cluster: LocationKey,
432    ) -> impl QuotedWithContext<'a, &'a [TaglessMemberId], ()> + Clone + 'a {
433        cluster_members(RuntimeData::new("__hydro_lang_trybuild_cli"), of_cluster)
434    }
435
436    fn cluster_self_id() -> impl QuotedWithContext<'a, TaglessMemberId, ()> + Clone + 'a {
437        cluster_self_id(RuntimeData::new("__hydro_lang_trybuild_cli"))
438    }
439
440    fn cluster_membership_stream(
441        _env: &mut Self::InstantiateEnv,
442        _at_location: &LocationId,
443        location_id: &LocationId,
444    ) -> impl QuotedWithContext<'a, Box<dyn Stream<Item = (TaglessMemberId, MembershipEvent)> + Unpin>, ()>
445    {
446        cluster_membership_stream(location_id)
447    }
448}
449
450#[expect(missing_docs, reason = "TODO")]
451pub trait DeployCrateWrapper {
452    fn underlying(&self) -> Arc<RustCrateService>;
453
454    fn stdout(&self) -> tokio::sync::mpsc::UnboundedReceiver<String> {
455        self.underlying().stdout()
456    }
457
458    fn stderr(&self) -> tokio::sync::mpsc::UnboundedReceiver<String> {
459        self.underlying().stderr()
460    }
461
462    fn stdout_filter(
463        &self,
464        prefix: impl Into<String>,
465    ) -> tokio::sync::mpsc::UnboundedReceiver<String> {
466        self.underlying().stdout_filter(prefix.into())
467    }
468
469    fn stderr_filter(
470        &self,
471        prefix: impl Into<String>,
472    ) -> tokio::sync::mpsc::UnboundedReceiver<String> {
473        self.underlying().stderr_filter(prefix.into())
474    }
475}
476
477#[expect(missing_docs, reason = "TODO")]
478#[derive(Clone)]
479pub struct TrybuildHost {
480    host: Arc<dyn Host>,
481    display_name: Option<String>,
482    rustflags: Option<String>,
483    profile: Option<String>,
484    additional_hydro_features: Vec<String>,
485    features: Vec<String>,
486    tracing: Option<TracingOptions>,
487    build_envs: Vec<(String, String)>,
488    env: HashMap<String, String>,
489    pin_to_core: Option<usize>,
490    name_hint: Option<String>,
491    cluster_idx: Option<usize>,
492}
493
494impl From<Arc<dyn Host>> for TrybuildHost {
495    fn from(host: Arc<dyn Host>) -> Self {
496        Self {
497            host,
498            display_name: None,
499            rustflags: None,
500            profile: None,
501            additional_hydro_features: vec![],
502            features: vec![],
503            tracing: None,
504            build_envs: vec![],
505            env: HashMap::new(),
506            pin_to_core: None,
507            name_hint: None,
508            cluster_idx: None,
509        }
510    }
511}
512
513impl<H: Host + 'static> From<Arc<H>> for TrybuildHost {
514    fn from(host: Arc<H>) -> Self {
515        Self {
516            host,
517            display_name: None,
518            rustflags: None,
519            profile: None,
520            additional_hydro_features: vec![],
521            features: vec![],
522            tracing: None,
523            build_envs: vec![],
524            env: HashMap::new(),
525            pin_to_core: None,
526            name_hint: None,
527            cluster_idx: None,
528        }
529    }
530}
531
532#[expect(missing_docs, reason = "TODO")]
533impl TrybuildHost {
534    pub fn new(host: Arc<dyn Host>) -> Self {
535        Self {
536            host,
537            display_name: None,
538            rustflags: None,
539            profile: None,
540            additional_hydro_features: vec![],
541            features: vec![],
542            tracing: None,
543            build_envs: vec![],
544            env: HashMap::new(),
545            pin_to_core: None,
546            name_hint: None,
547            cluster_idx: None,
548        }
549    }
550
551    pub fn display_name(self, display_name: impl Into<String>) -> Self {
552        assert!(
553            self.display_name.is_none(),
554            "{} already set",
555            name_of!(display_name in Self)
556        );
557
558        Self {
559            display_name: Some(display_name.into()),
560            ..self
561        }
562    }
563
564    pub fn rustflags(self, rustflags: impl Into<String>) -> Self {
565        assert!(
566            self.rustflags.is_none(),
567            "{} already set",
568            name_of!(rustflags in Self)
569        );
570
571        Self {
572            rustflags: Some(rustflags.into()),
573            ..self
574        }
575    }
576
577    pub fn profile(self, profile: impl Into<String>) -> Self {
578        assert!(
579            self.profile.is_none(),
580            "{} already set",
581            name_of!(profile in Self)
582        );
583
584        Self {
585            profile: Some(profile.into()),
586            ..self
587        }
588    }
589
590    pub fn additional_hydro_features(
591        mut self,
592        additional_hydro_features: impl IntoIterator<Item = impl Into<String>>,
593    ) -> Self {
594        self.additional_hydro_features
595            .extend(additional_hydro_features.into_iter().map(Into::into));
596        self
597    }
598
599    pub fn additional_hydro_feature(mut self, feature: impl Into<String>) -> Self {
600        self.additional_hydro_features.push(feature.into());
601        self
602    }
603
604    pub fn features(mut self, features: impl IntoIterator<Item = impl Into<String>>) -> Self {
605        self.features.extend(features.into_iter().map(Into::into));
606        self
607    }
608
609    pub fn feature(mut self, feature: impl Into<String>) -> Self {
610        self.features.push(feature.into());
611        self
612    }
613
614    #[must_use]
615    pub fn tracing(self, tracing: TracingOptions) -> Self {
616        assert!(
617            self.tracing.is_none(),
618            "{} already set",
619            name_of!(tracing in Self)
620        );
621
622        Self {
623            tracing: Some(tracing),
624            ..self
625        }
626    }
627
628    pub fn build_env(self, key: impl Into<String>, value: impl Into<String>) -> Self {
629        Self {
630            build_envs: self
631                .build_envs
632                .into_iter()
633                .chain(std::iter::once((key.into(), value.into())))
634                .collect(),
635            ..self
636        }
637    }
638
639    pub fn env(self, key: impl Into<String>, value: impl Into<String>) -> Self {
640        let mut env = self.env;
641        env.insert(key.into(), value.into());
642        Self { env, ..self }
643    }
644
645    #[must_use]
646    pub fn pin_to_core(self, core: usize) -> Self {
647        Self {
648            pin_to_core: Some(core),
649            ..self
650        }
651    }
652}
653
654impl IntoProcessSpec<'_, HydroDeploy> for Arc<dyn Host> {
655    type ProcessSpec = TrybuildHost;
656    fn into_process_spec(self) -> TrybuildHost {
657        TrybuildHost {
658            host: self,
659            display_name: None,
660            rustflags: None,
661            profile: None,
662            additional_hydro_features: vec![],
663            features: vec![],
664            tracing: None,
665            build_envs: vec![],
666            env: HashMap::new(),
667            pin_to_core: None,
668            name_hint: None,
669            cluster_idx: None,
670        }
671    }
672}
673
674impl<H: Host + 'static> IntoProcessSpec<'_, HydroDeploy> for Arc<H> {
675    type ProcessSpec = TrybuildHost;
676    fn into_process_spec(self) -> TrybuildHost {
677        TrybuildHost {
678            host: self,
679            display_name: None,
680            rustflags: None,
681            profile: None,
682            additional_hydro_features: vec![],
683            features: vec![],
684            tracing: None,
685            build_envs: vec![],
686            env: HashMap::new(),
687            pin_to_core: None,
688            name_hint: None,
689            cluster_idx: None,
690        }
691    }
692}
693
694#[expect(missing_docs, reason = "TODO")]
695#[derive(Clone)]
696pub struct DeployExternal {
697    next_port: Rc<RefCell<usize>>,
698    host: Arc<dyn Host>,
699    underlying: Rc<RefCell<Option<Arc<CustomService>>>>,
700    client_ports: Rc<RefCell<HashMap<String, CustomClientPort>>>,
701    allocated_ports: Rc<RefCell<HashMap<ExternalPortId, String>>>,
702}
703
704impl DeployExternal {
705    pub(crate) fn raw_port(&self, external_port_id: ExternalPortId) -> CustomClientPort {
706        self.client_ports
707            .borrow()
708            .get(
709                self.allocated_ports
710                    .borrow()
711                    .get(&external_port_id)
712                    .unwrap(),
713            )
714            .unwrap()
715            .clone()
716    }
717}
718
719impl<'a> RegisterPort<'a, HydroDeploy> for DeployExternal {
720    fn register(&self, external_port_id: ExternalPortId, port: Self::Port) {
721        assert!(
722            self.allocated_ports
723                .borrow_mut()
724                .insert(external_port_id, port.clone())
725                .is_none_or(|old| old == port)
726        );
727    }
728
729    fn as_bytes_bidi(
730        &self,
731        external_port_id: ExternalPortId,
732    ) -> impl Future<
733        Output = (
734            Pin<Box<dyn Stream<Item = Result<BytesMut, Error>>>>,
735            Pin<Box<dyn Sink<Bytes, Error = Error>>>,
736        ),
737    > + 'a {
738        let port = self.raw_port(external_port_id);
739
740        async move {
741            let (source, sink) = port.connect().await.into_source_sink();
742            (
743                Box::pin(source) as Pin<Box<dyn Stream<Item = Result<BytesMut, Error>>>>,
744                Box::pin(sink) as Pin<Box<dyn Sink<Bytes, Error = Error>>>,
745            )
746        }
747    }
748
749    fn as_bincode_bidi<InT, OutT>(
750        &self,
751        external_port_id: ExternalPortId,
752    ) -> impl Future<
753        Output = (
754            Pin<Box<dyn Stream<Item = OutT>>>,
755            Pin<Box<dyn Sink<InT, Error = Error>>>,
756        ),
757    > + 'a
758    where
759        InT: Serialize + 'static,
760        OutT: DeserializeOwned + 'static,
761    {
762        let port = self.raw_port(external_port_id);
763        async move {
764            let (source, sink) = port.connect().await.into_source_sink();
765            (
766                Box::pin(source.map(|item| bincode::deserialize(&item.unwrap()).unwrap()))
767                    as Pin<Box<dyn Stream<Item = OutT>>>,
768                Box::pin(
769                    sink.with(|item| async move { Ok(bincode::serialize(&item).unwrap().into()) }),
770                ) as Pin<Box<dyn Sink<InT, Error = Error>>>,
771            )
772        }
773    }
774
775    fn as_bincode_sink<T: Serialize + 'static>(
776        &self,
777        external_port_id: ExternalPortId,
778    ) -> impl Future<Output = Pin<Box<dyn Sink<T, Error = Error>>>> + 'a {
779        let port = self.raw_port(external_port_id);
780        async move {
781            let sink = port.connect().await.into_sink();
782            Box::pin(sink.with(|item| async move { Ok(bincode::serialize(&item).unwrap().into()) }))
783                as Pin<Box<dyn Sink<T, Error = Error>>>
784        }
785    }
786
787    fn as_bincode_source<T: DeserializeOwned + 'static>(
788        &self,
789        external_port_id: ExternalPortId,
790    ) -> impl Future<Output = Pin<Box<dyn Stream<Item = T>>>> + 'a {
791        let port = self.raw_port(external_port_id);
792        async move {
793            let source = port.connect().await.into_source();
794            Box::pin(source.map(|item| bincode::deserialize(&item.unwrap()).unwrap()))
795                as Pin<Box<dyn Stream<Item = T>>>
796        }
797    }
798}
799
800impl Node for DeployExternal {
801    type Port = String;
802    /// Map from Cluster location ID to member IDs.
803    type Meta = SparseSecondaryMap<LocationKey, Vec<TaglessMemberId>>;
804    type InstantiateEnv = Deployment;
805
806    fn next_port(&self) -> Self::Port {
807        let next_port = *self.next_port.borrow();
808        *self.next_port.borrow_mut() += 1;
809
810        format!("port_{}", next_port)
811    }
812
813    fn instantiate(
814        &self,
815        env: &mut Self::InstantiateEnv,
816        _meta: &mut Self::Meta,
817        _graph: DfirGraph,
818        extra_stmts: &[syn::Stmt],
819        sidecars: &[syn::Expr],
820        _as_code_options: &AsCodeOptions,
821    ) {
822        assert!(extra_stmts.is_empty());
823        assert!(sidecars.is_empty());
824        let service = env.CustomService(self.host.clone(), vec![]);
825        *self.underlying.borrow_mut() = Some(service);
826    }
827
828    fn update_meta(&self, _meta: &Self::Meta) {}
829}
830
831impl ExternalSpec<'_, HydroDeploy> for Arc<dyn Host> {
832    fn build(self, _key: LocationKey, _name_hint: &str) -> DeployExternal {
833        DeployExternal {
834            next_port: Rc::new(RefCell::new(0)),
835            host: self,
836            underlying: Rc::new(RefCell::new(None)),
837            allocated_ports: Rc::new(RefCell::new(HashMap::new())),
838            client_ports: Rc::new(RefCell::new(HashMap::new())),
839        }
840    }
841}
842
843impl<H: Host + 'static> ExternalSpec<'_, HydroDeploy> for Arc<H> {
844    fn build(self, _key: LocationKey, _name_hint: &str) -> DeployExternal {
845        DeployExternal {
846            next_port: Rc::new(RefCell::new(0)),
847            host: self,
848            underlying: Rc::new(RefCell::new(None)),
849            allocated_ports: Rc::new(RefCell::new(HashMap::new())),
850            client_ports: Rc::new(RefCell::new(HashMap::new())),
851        }
852    }
853}
854
855pub(crate) enum CrateOrTrybuild {
856    Crate(RustCrate, Arc<dyn Host>),
857    Trybuild(TrybuildHost),
858}
859
860#[expect(missing_docs, reason = "TODO")]
861#[derive(Clone)]
862pub struct DeployNode {
863    next_port: Rc<RefCell<usize>>,
864    service_spec: Rc<RefCell<Option<CrateOrTrybuild>>>,
865    underlying: Rc<RefCell<Option<Arc<RustCrateService>>>>,
866}
867
868impl DeployCrateWrapper for DeployNode {
869    fn underlying(&self) -> Arc<RustCrateService> {
870        Arc::clone(self.underlying.borrow().as_ref().unwrap())
871    }
872}
873
874impl Node for DeployNode {
875    type Port = String;
876    /// Map from Cluster location ID to member IDs.
877    type Meta = SparseSecondaryMap<LocationKey, Vec<TaglessMemberId>>;
878    type InstantiateEnv = Deployment;
879
880    fn next_port(&self) -> String {
881        let next_port = *self.next_port.borrow();
882        *self.next_port.borrow_mut() += 1;
883
884        format!("port_{}", next_port)
885    }
886
887    fn update_meta(&self, meta: &Self::Meta) {
888        let underlying_node = self.underlying.borrow();
889        underlying_node.as_ref().unwrap().update_meta(HydroMeta {
890            clusters: meta.clone(),
891            cluster_id: None,
892        });
893    }
894
895    fn instantiate(
896        &self,
897        env: &mut Self::InstantiateEnv,
898        _meta: &mut Self::Meta,
899        graph: DfirGraph,
900        extra_stmts: &[syn::Stmt],
901        sidecars: &[syn::Expr],
902        as_code_options: &AsCodeOptions,
903    ) {
904        let (service, host) = match self.service_spec.borrow_mut().take().unwrap() {
905            CrateOrTrybuild::Crate(c, host) => (c, host),
906            CrateOrTrybuild::Trybuild(trybuild) => {
907                // Determine linking mode based on host target type
908                let linking_mode = if !cfg!(target_os = "windows")
909                    && trybuild.host.target_type() == hydro_deploy::HostTargetType::Local
910                    && trybuild.rustflags.is_none()
911                {
912                    // When compiling for local, prefer dynamic linking to reduce binary size
913                    // Windows is currently not supported due to https://github.com/bevyengine/bevy/pull/2016
914                    LinkingMode::Dynamic
915                } else {
916                    LinkingMode::Static
917                };
918                let (bin_name, config) = create_graph_trybuild(
919                    graph,
920                    extra_stmts,
921                    sidecars,
922                    as_code_options,
923                    trybuild.name_hint.as_deref(),
924                    crate::compile::trybuild::generate::DeployMode::HydroDeploy,
925                    linking_mode,
926                );
927                let host = trybuild.host.clone();
928                (
929                    create_trybuild_service(
930                        trybuild,
931                        &config.project_dir,
932                        &config.target_dir,
933                        config.features.as_deref(),
934                        &bin_name,
935                        &config.linking_mode,
936                    ),
937                    host,
938                )
939            }
940        };
941
942        *self.underlying.borrow_mut() = Some(env.add_service(service, host));
943    }
944}
945
946#[expect(missing_docs, reason = "TODO")]
947#[derive(Clone)]
948pub struct DeployClusterNode {
949    underlying: Arc<RustCrateService>,
950}
951
952impl DeployCrateWrapper for DeployClusterNode {
953    fn underlying(&self) -> Arc<RustCrateService> {
954        self.underlying.clone()
955    }
956}
957#[expect(missing_docs, reason = "TODO")]
958#[derive(Clone)]
959pub struct DeployCluster {
960    key: LocationKey,
961    next_port: Rc<RefCell<usize>>,
962    cluster_spec: Rc<RefCell<Option<Vec<CrateOrTrybuild>>>>,
963    members: Rc<RefCell<Vec<DeployClusterNode>>>,
964    name_hint: Option<String>,
965}
966
967impl DeployCluster {
968    #[expect(missing_docs, reason = "TODO")]
969    pub fn members(&self) -> Vec<DeployClusterNode> {
970        self.members.borrow().clone()
971    }
972}
973
974impl Node for DeployCluster {
975    type Port = String;
976    /// Map from Cluster location ID to member IDs.
977    type Meta = SparseSecondaryMap<LocationKey, Vec<TaglessMemberId>>;
978    type InstantiateEnv = Deployment;
979
980    fn next_port(&self) -> String {
981        let next_port = *self.next_port.borrow();
982        *self.next_port.borrow_mut() += 1;
983
984        format!("port_{}", next_port)
985    }
986
987    fn instantiate(
988        &self,
989        env: &mut Self::InstantiateEnv,
990        meta: &mut Self::Meta,
991        graph: DfirGraph,
992        extra_stmts: &[syn::Stmt],
993        sidecars: &[syn::Expr],
994        as_code_options: &AsCodeOptions,
995    ) {
996        let has_trybuild = self
997            .cluster_spec
998            .borrow()
999            .as_ref()
1000            .unwrap()
1001            .iter()
1002            .any(|spec| matches!(spec, CrateOrTrybuild::Trybuild { .. }));
1003
1004        // For clusters, use static linking if ANY host is non-local (conservative approach)
1005        let linking_mode = if !cfg!(target_os = "windows")
1006            && self
1007                .cluster_spec
1008                .borrow()
1009                .as_ref()
1010                .unwrap()
1011                .iter()
1012                .all(|spec| match spec {
1013                    CrateOrTrybuild::Crate(_, _) => true, // crates handle their own linking
1014                    CrateOrTrybuild::Trybuild(t) => {
1015                        t.host.target_type() == hydro_deploy::HostTargetType::Local
1016                            && t.rustflags.is_none()
1017                    }
1018                }) {
1019            // See comment above for Windows exception
1020            LinkingMode::Dynamic
1021        } else {
1022            LinkingMode::Static
1023        };
1024
1025        let maybe_trybuild = if has_trybuild {
1026            Some(create_graph_trybuild(
1027                graph,
1028                extra_stmts,
1029                sidecars,
1030                as_code_options,
1031                self.name_hint.as_deref(),
1032                crate::compile::trybuild::generate::DeployMode::HydroDeploy,
1033                linking_mode,
1034            ))
1035        } else {
1036            None
1037        };
1038
1039        let cluster_nodes = self
1040            .cluster_spec
1041            .borrow_mut()
1042            .take()
1043            .unwrap()
1044            .into_iter()
1045            .map(|spec| {
1046                let (service, host) = match spec {
1047                    CrateOrTrybuild::Crate(c, host) => (c, host),
1048                    CrateOrTrybuild::Trybuild(trybuild) => {
1049                        let (bin_name, config) = maybe_trybuild.as_ref().unwrap();
1050                        let host = trybuild.host.clone();
1051                        (
1052                            create_trybuild_service(
1053                                trybuild,
1054                                &config.project_dir,
1055                                &config.target_dir,
1056                                config.features.as_deref(),
1057                                bin_name,
1058                                &config.linking_mode,
1059                            ),
1060                            host,
1061                        )
1062                    }
1063                };
1064
1065                env.add_service(service, host)
1066            })
1067            .collect::<Vec<_>>();
1068        meta.insert(
1069            self.key,
1070            (0..(cluster_nodes.len() as u32))
1071                .map(TaglessMemberId::from_raw_id)
1072                .collect(),
1073        );
1074        *self.members.borrow_mut() = cluster_nodes
1075            .into_iter()
1076            .map(|n| DeployClusterNode { underlying: n })
1077            .collect();
1078    }
1079
1080    fn update_meta(&self, meta: &Self::Meta) {
1081        for (cluster_id, node) in self.members.borrow().iter().enumerate() {
1082            node.underlying.update_meta(HydroMeta {
1083                clusters: meta.clone(),
1084                cluster_id: Some(TaglessMemberId::from_raw_id(cluster_id as u32)),
1085            });
1086        }
1087    }
1088}
1089
1090#[expect(missing_docs, reason = "TODO")]
1091#[derive(Clone)]
1092pub struct DeployProcessSpec(RustCrate, Arc<dyn Host>);
1093
1094impl DeployProcessSpec {
1095    #[expect(missing_docs, reason = "TODO")]
1096    pub fn new(t: RustCrate, host: Arc<dyn Host>) -> Self {
1097        Self(t, host)
1098    }
1099}
1100
1101impl ProcessSpec<'_, HydroDeploy> for DeployProcessSpec {
1102    fn build(self, _key: LocationKey, _name_hint: &str) -> DeployNode {
1103        DeployNode {
1104            next_port: Rc::new(RefCell::new(0)),
1105            service_spec: Rc::new(RefCell::new(Some(CrateOrTrybuild::Crate(self.0, self.1)))),
1106            underlying: Rc::new(RefCell::new(None)),
1107        }
1108    }
1109}
1110
1111impl ProcessSpec<'_, HydroDeploy> for TrybuildHost {
1112    fn build(mut self, key: LocationKey, name_hint: &str) -> DeployNode {
1113        self.name_hint = Some(format!("{} (process {})", name_hint, key));
1114        DeployNode {
1115            next_port: Rc::new(RefCell::new(0)),
1116            service_spec: Rc::new(RefCell::new(Some(CrateOrTrybuild::Trybuild(self)))),
1117            underlying: Rc::new(RefCell::new(None)),
1118        }
1119    }
1120}
1121
1122#[expect(missing_docs, reason = "TODO")]
1123#[derive(Clone)]
1124pub struct DeployClusterSpec(Vec<(RustCrate, Arc<dyn Host>)>);
1125
1126impl DeployClusterSpec {
1127    #[expect(missing_docs, reason = "TODO")]
1128    pub fn new(crates: Vec<(RustCrate, Arc<dyn Host>)>) -> Self {
1129        Self(crates)
1130    }
1131}
1132
1133impl ClusterSpec<'_, HydroDeploy> for DeployClusterSpec {
1134    fn build(self, key: LocationKey, _name_hint: &str) -> DeployCluster {
1135        DeployCluster {
1136            key,
1137            next_port: Rc::new(RefCell::new(0)),
1138            cluster_spec: Rc::new(RefCell::new(Some(
1139                self.0
1140                    .into_iter()
1141                    .map(|(c, h)| CrateOrTrybuild::Crate(c, h))
1142                    .collect(),
1143            ))),
1144            members: Rc::new(RefCell::new(vec![])),
1145            name_hint: None,
1146        }
1147    }
1148}
1149
1150impl<T: Into<TrybuildHost>, I: IntoIterator<Item = T>> ClusterSpec<'_, HydroDeploy> for I {
1151    fn build(self, key: LocationKey, name_hint: &str) -> DeployCluster {
1152        let name_hint = format!("{} (cluster {})", name_hint, key);
1153        DeployCluster {
1154            key,
1155            next_port: Rc::new(RefCell::new(0)),
1156            cluster_spec: Rc::new(RefCell::new(Some(
1157                self.into_iter()
1158                    .enumerate()
1159                    .map(|(idx, b)| {
1160                        let mut b = b.into();
1161                        b.name_hint = Some(name_hint.clone());
1162                        b.cluster_idx = Some(idx);
1163                        CrateOrTrybuild::Trybuild(b)
1164                    })
1165                    .collect(),
1166            ))),
1167            members: Rc::new(RefCell::new(vec![])),
1168            name_hint: Some(name_hint),
1169        }
1170    }
1171}
1172
1173fn create_trybuild_service(
1174    trybuild: TrybuildHost,
1175    dir: &std::path::Path,
1176    target_dir: &std::path::PathBuf,
1177    features: Option<&[String]>,
1178    bin_name: &str,
1179    linking_mode: &LinkingMode,
1180) -> RustCrate {
1181    // For dynamic linking, use the dylib-examples crate; for static, use the base crate
1182    let crate_dir = match linking_mode {
1183        LinkingMode::Dynamic => dir.join("dylib-examples"),
1184        LinkingMode::Static => dir.to_path_buf(),
1185    };
1186
1187    let mut ret = RustCrate::new(&crate_dir, dir)
1188        .target_dir(target_dir)
1189        .example(bin_name)
1190        .no_default_features();
1191
1192    ret = ret.set_is_dylib(matches!(linking_mode, LinkingMode::Dynamic));
1193
1194    if let Some(display_name) = trybuild.display_name {
1195        ret = ret.display_name(display_name);
1196    } else if let Some(name_hint) = trybuild.name_hint {
1197        if let Some(cluster_idx) = trybuild.cluster_idx {
1198            ret = ret.display_name(format!("{} / {}", name_hint, cluster_idx));
1199        } else {
1200            ret = ret.display_name(name_hint);
1201        }
1202    }
1203
1204    if let Some(rustflags) = trybuild.rustflags {
1205        ret = ret.rustflags(rustflags);
1206    }
1207
1208    if let Some(profile) = trybuild.profile {
1209        ret = ret.profile(profile);
1210    }
1211
1212    if let Some(tracing) = trybuild.tracing {
1213        ret = ret.tracing(tracing);
1214    }
1215
1216    if let Some(core) = trybuild.pin_to_core {
1217        ret = ret.pin_to_core(core);
1218    }
1219
1220    ret = ret.features(
1221        vec!["hydro___feature_deploy_integration".to_owned()]
1222            .into_iter()
1223            .chain(
1224                trybuild
1225                    .additional_hydro_features
1226                    .into_iter()
1227                    .map(|runtime_feature| {
1228                        assert!(
1229                            HYDRO_RUNTIME_FEATURES.iter().any(|f| f == &runtime_feature),
1230                            "{runtime_feature} is not a valid Hydro runtime feature"
1231                        );
1232                        format!("hydro___feature_{runtime_feature}")
1233                    }),
1234            )
1235            .chain(trybuild.features),
1236    );
1237
1238    for (key, value) in trybuild.build_envs {
1239        ret = ret.build_env(key, value);
1240    }
1241
1242    for (key, value) in trybuild.env {
1243        ret = ret.env(key, value);
1244    }
1245
1246    ret = ret.build_env("STAGELEFT_TRYBUILD_BUILD_STAGED", "1");
1247    ret = ret.config("build.incremental = false");
1248
1249    if let Some(features) = features {
1250        ret = ret.features(features);
1251    }
1252
1253    ret
1254}