1use 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
41pub enum HydroDeploy {}
46
47impl<'a> Deploy<'a> for HydroDeploy {
48 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 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 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 let linking_mode = if !cfg!(target_os = "windows")
909 && trybuild.host.target_type() == hydro_deploy::HostTargetType::Local
910 && trybuild.rustflags.is_none()
911 {
912 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 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 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, CrateOrTrybuild::Trybuild(t) => {
1015 t.host.target_type() == hydro_deploy::HostTargetType::Local
1016 && t.rustflags.is_none()
1017 }
1018 }) {
1019 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 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}