1mod github_context;
7mod spec;
8
9pub use github_context::GhOutput;
10pub use github_context::GhToRust;
11pub use github_context::RustToGh;
12
13use self::steps::ado::AdoRuntimeVar;
14use self::steps::ado::AdoStepServices;
15use self::steps::github::GhStepBuilder;
16use self::steps::rust::RustRuntimeServices;
17use self::user_facing::ClaimedGhParam;
18use self::user_facing::GhPermission;
19use self::user_facing::GhPermissionValue;
20use crate::node::github_context::GhContextVarReader;
21use github_context::state::Root;
22use serde::Deserialize;
23use serde::Serialize;
24use serde::de::DeserializeOwned;
25use std::cell::RefCell;
26use std::collections::BTreeMap;
27use std::path::PathBuf;
28use std::rc::Rc;
29use user_facing::GhParam;
30
31pub mod user_facing {
34 pub use super::ClaimVar;
35 pub use super::ClaimedReadVar;
36 pub use super::ClaimedWriteVar;
37 pub use super::ConfigField;
38 pub use super::ConfigMerge;
39 pub use super::ConfigVar;
40 pub use super::FlowArch;
41 pub use super::FlowBackend;
42 pub use super::FlowNode;
43 pub use super::FlowNodeWithConfig;
44 pub use super::FlowPlatform;
45 pub use super::FlowPlatformKind;
46 pub use super::GhUserSecretVar;
47 pub use super::ImportCtx;
48 pub use super::IntoConfig;
49 pub use super::IntoRequest;
50 pub use super::NodeCtx;
51 pub use super::ReadVar;
52 pub use super::SideEffect;
53 pub use super::SimpleFlowNode;
54 pub use super::StepCtx;
55 pub use super::VarClaimed;
56 pub use super::VarEqBacking;
57 pub use super::VarNotClaimed;
58 pub use super::WriteVar;
59 pub use super::steps::ado::AdoResourcesRepositoryId;
60 pub use super::steps::ado::AdoRuntimeVar;
61 pub use super::steps::ado::AdoStepServices;
62 pub use super::steps::github::ClaimedGhParam;
63 pub use super::steps::github::GhParam;
64 pub use super::steps::github::GhPermission;
65 pub use super::steps::github::GhPermissionValue;
66 pub use super::steps::rust::RustRuntimeServices;
67 pub use crate::flowey_config;
68 pub use crate::flowey_request;
69 pub use crate::new_flow_node;
70 pub use crate::new_flow_node_with_config;
71 pub use crate::new_simple_flow_node;
72 pub use crate::node::FlowPlatformLinuxDistro;
73 pub use crate::pipeline::Artifact;
74 pub use crate::pipeline::ArtifactType;
75
76 pub fn same_across_all_reqs<T: PartialEq>(
108 req_name: &str,
109 var: &mut Option<T>,
110 new: T,
111 ) -> anyhow::Result<()> {
112 match (var.as_ref(), new) {
113 (None, v) => *var = Some(v),
114 (Some(old), new) => {
115 if *old != new {
116 anyhow::bail!("`{}` must be consistent across requests", req_name);
117 }
118 }
119 }
120
121 Ok(())
122 }
123
124 pub fn same_across_all_reqs_backing_var<V: VarEqBacking>(
128 req_name: &str,
129 var: &mut Option<V>,
130 new: V,
131 ) -> anyhow::Result<()> {
132 match (var.as_ref(), new) {
133 (None, v) => *var = Some(v),
134 (Some(old), new) => {
135 if !old.eq(&new) {
136 anyhow::bail!("`{}` must be consistent across requests", req_name);
137 }
138 }
139 }
140
141 Ok(())
142 }
143
144 #[macro_export]
148 macro_rules! match_arch {
149 ($host_arch:expr, $match_arch:pat, $expr:expr) => {
150 if matches!($host_arch, $match_arch) {
151 $expr
152 } else {
153 anyhow::bail!("Linux distro not supported on host arch {}", $host_arch);
154 }
155 };
156 }
157
158 #[macro_export]
160 macro_rules! claim_vars {
161 ($ctx:ident, ($($var:ident),* $(,)?)) => {
162 $(let $var = $var.claim($ctx);)*
163 };
164 }
165
166 #[macro_export]
168 macro_rules! read_vars {
169 ($rt:ident, ($($var:ident),* $(,)?)) => {
170 $(let $var = $rt.read($var);)*
171 };
172 }
173}
174
175pub trait VarEqBacking {
199 fn eq(&self, other: &Self) -> bool;
201}
202
203impl<T> VarEqBacking for WriteVar<T>
204where
205 T: Serialize + DeserializeOwned,
206{
207 fn eq(&self, other: &Self) -> bool {
208 self.backing_var == other.backing_var
209 }
210}
211
212impl<T> VarEqBacking for ReadVar<T>
213where
214 T: Serialize + DeserializeOwned + PartialEq + Eq + Clone,
215{
216 fn eq(&self, other: &Self) -> bool {
217 self.backing_var == other.backing_var
218 }
219}
220
221impl<T, U> VarEqBacking for (T, U)
223where
224 T: VarEqBacking,
225 U: VarEqBacking,
226{
227 fn eq(&self, other: &Self) -> bool {
228 (self.0.eq(&other.0)) && (self.1.eq(&other.1))
229 }
230}
231
232#[derive(Serialize, Deserialize)]
250#[serde(bound(serialize = "T: Serialize", deserialize = "T: DeserializeOwned"))]
251pub struct ConfigVar<T>(pub ReadVar<T>);
252
253impl<T: Serialize + DeserializeOwned> Clone for ConfigVar<T> {
254 fn clone(&self) -> Self {
255 ConfigVar(self.0.clone())
256 }
257}
258
259impl<T> std::fmt::Debug for ConfigVar<T> {
260 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
261 f.debug_tuple("ConfigVar").finish()
262 }
263}
264
265impl<T: Serialize + DeserializeOwned + PartialEq + Eq + Clone> PartialEq for ConfigVar<T> {
266 fn eq(&self, other: &Self) -> bool {
267 VarEqBacking::eq(&self.0, &other.0)
268 }
269}
270
271impl<T: Serialize + DeserializeOwned + PartialEq + Eq + Clone> ClaimVar for ConfigVar<T> {
272 type Claimed = ClaimedReadVar<T>;
273
274 fn claim(self, ctx: &mut StepCtx<'_>) -> ClaimedReadVar<T> {
275 self.0.claim(ctx)
276 }
277}
278
279impl<T: Serialize + DeserializeOwned + PartialEq + Eq + Clone> From<ReadVar<T>> for ConfigVar<T> {
280 fn from(v: ReadVar<T>) -> Self {
281 ConfigVar(v)
282 }
283}
284
285pub type SideEffect = ();
292
293#[derive(Clone, Debug, Serialize, Deserialize)]
296pub enum VarNotClaimed {}
297
298#[derive(Clone, Debug, Serialize, Deserialize)]
301pub enum VarClaimed {}
302
303#[derive(Debug, Serialize, Deserialize)]
323pub struct WriteVar<T: Serialize + DeserializeOwned, C = VarNotClaimed> {
324 backing_var: String,
325 is_side_effect: bool,
328
329 #[serde(skip)]
330 _kind: core::marker::PhantomData<(T, C)>,
331}
332
333pub type ClaimedWriteVar<T> = WriteVar<T, VarClaimed>;
336
337impl<T: Serialize + DeserializeOwned> WriteVar<T, VarNotClaimed> {
338 fn into_claimed(self) -> WriteVar<T, VarClaimed> {
340 let Self {
341 backing_var,
342 is_side_effect,
343 _kind,
344 } = self;
345
346 WriteVar {
347 backing_var,
348 is_side_effect,
349 _kind: std::marker::PhantomData,
350 }
351 }
352
353 #[track_caller]
355 pub fn write_static(self, ctx: &mut NodeCtx<'_>, val: T)
356 where
357 T: 'static,
358 {
359 let val = ReadVar::from_static(val);
360 val.write_into(ctx, self);
361 }
362
363 pub(crate) fn into_json(self) -> WriteVar<serde_json::Value> {
364 WriteVar {
365 backing_var: self.backing_var,
366 is_side_effect: self.is_side_effect,
367 _kind: std::marker::PhantomData,
368 }
369 }
370}
371
372impl WriteVar<SideEffect, VarNotClaimed> {
373 pub fn discard_result<T: Serialize + DeserializeOwned>(self) -> WriteVar<T> {
378 WriteVar {
379 backing_var: self.backing_var,
380 is_side_effect: true,
381 _kind: std::marker::PhantomData,
382 }
383 }
384}
385
386pub trait ClaimVar {
394 type Claimed;
396 fn claim(self, ctx: &mut StepCtx<'_>) -> Self::Claimed;
398}
399
400pub trait ReadVarValue {
406 type Value;
408 fn read_value(self, rt: &mut RustRuntimeServices<'_>) -> Self::Value;
410}
411
412impl<T: Serialize + DeserializeOwned> ClaimVar for ReadVar<T> {
413 type Claimed = ClaimedReadVar<T>;
414
415 fn claim(self, ctx: &mut StepCtx<'_>) -> ClaimedReadVar<T> {
416 if let ReadVarBacking::RuntimeVar {
417 var,
418 is_side_effect: _,
419 } = &self.backing_var
420 {
421 ctx.backend.borrow_mut().on_claimed_runtime_var(var, true);
422 }
423 self.into_claimed()
424 }
425}
426
427impl<T: Serialize + DeserializeOwned> ClaimVar for WriteVar<T> {
428 type Claimed = ClaimedWriteVar<T>;
429
430 fn claim(self, ctx: &mut StepCtx<'_>) -> ClaimedWriteVar<T> {
431 ctx.backend
432 .borrow_mut()
433 .on_claimed_runtime_var(&self.backing_var, false);
434 self.into_claimed()
435 }
436}
437
438impl<T: Serialize + DeserializeOwned> ReadVarValue for ClaimedReadVar<T> {
439 type Value = T;
440
441 fn read_value(self, rt: &mut RustRuntimeServices<'_>) -> Self::Value {
442 match self.backing_var {
443 ReadVarBacking::RuntimeVar {
444 var,
445 is_side_effect,
446 } => {
447 let data = rt.get_var(&var, is_side_effect);
449 if is_side_effect {
450 serde_json::from_slice(b"null").expect("should be deserializing into ()")
454 } else {
455 serde_json::from_slice(&data).expect("improve this error path")
457 }
458 }
459 ReadVarBacking::Inline(val) => val,
460 }
461 }
462}
463
464impl<T: ClaimVar> ClaimVar for Vec<T> {
465 type Claimed = Vec<T::Claimed>;
466
467 fn claim(self, ctx: &mut StepCtx<'_>) -> Vec<T::Claimed> {
468 self.into_iter().map(|v| v.claim(ctx)).collect()
469 }
470}
471
472impl<T: ReadVarValue> ReadVarValue for Vec<T> {
473 type Value = Vec<T::Value>;
474
475 fn read_value(self, rt: &mut RustRuntimeServices<'_>) -> Self::Value {
476 self.into_iter().map(|v| v.read_value(rt)).collect()
477 }
478}
479
480impl<T: ClaimVar> ClaimVar for Option<T> {
481 type Claimed = Option<T::Claimed>;
482
483 fn claim(self, ctx: &mut StepCtx<'_>) -> Option<T::Claimed> {
484 self.map(|x| x.claim(ctx))
485 }
486}
487
488impl<T: ReadVarValue> ReadVarValue for Option<T> {
489 type Value = Option<T::Value>;
490
491 fn read_value(self, rt: &mut RustRuntimeServices<'_>) -> Self::Value {
492 self.map(|x| x.read_value(rt))
493 }
494}
495
496impl<U: Ord, T: ClaimVar> ClaimVar for BTreeMap<U, T> {
497 type Claimed = BTreeMap<U, T::Claimed>;
498
499 fn claim(self, ctx: &mut StepCtx<'_>) -> BTreeMap<U, T::Claimed> {
500 self.into_iter().map(|(k, v)| (k, v.claim(ctx))).collect()
501 }
502}
503
504impl<U: Ord, T: ReadVarValue> ReadVarValue for BTreeMap<U, T> {
505 type Value = BTreeMap<U, T::Value>;
506
507 fn read_value(self, rt: &mut RustRuntimeServices<'_>) -> Self::Value {
508 self.into_iter()
509 .map(|(k, v)| (k, v.read_value(rt)))
510 .collect()
511 }
512}
513
514macro_rules! impl_tuple_claim {
515 ($($T:tt)*) => {
516 impl<$($T,)*> $crate::node::ClaimVar for ($($T,)*)
517 where
518 $($T: $crate::node::ClaimVar,)*
519 {
520 type Claimed = ($($T::Claimed,)*);
521
522 #[expect(non_snake_case)]
523 fn claim(self, ctx: &mut $crate::node::StepCtx<'_>) -> Self::Claimed {
524 let ($($T,)*) = self;
525 ($($T.claim(ctx),)*)
526 }
527 }
528
529 impl<$($T,)*> $crate::node::ReadVarValue for ($($T,)*)
530 where
531 $($T: $crate::node::ReadVarValue,)*
532 {
533 type Value = ($($T::Value,)*);
534
535 #[expect(non_snake_case)]
536 fn read_value(self, rt: &mut $crate::node::RustRuntimeServices<'_>) -> Self::Value {
537 let ($($T,)*) = self;
538 ($($T.read_value(rt),)*)
539 }
540 }
541 };
542}
543
544impl_tuple_claim!(A B C D E F G H I J);
545impl_tuple_claim!(A B C D E F G H I);
546impl_tuple_claim!(A B C D E F G H);
547impl_tuple_claim!(A B C D E F G);
548impl_tuple_claim!(A B C D E F);
549impl_tuple_claim!(A B C D E);
550impl_tuple_claim!(A B C D);
551impl_tuple_claim!(A B C);
552impl_tuple_claim!(A B);
553impl_tuple_claim!(A);
554
555impl ClaimVar for () {
556 type Claimed = ();
557
558 fn claim(self, _ctx: &mut StepCtx<'_>) -> Self::Claimed {}
559}
560
561impl ReadVarValue for () {
562 type Value = ();
563
564 fn read_value(self, _rt: &mut RustRuntimeServices<'_>) -> Self::Value {}
565}
566
567#[derive(Serialize, Deserialize, Clone)]
572pub struct GhUserSecretVar(pub(crate) String);
573
574#[derive(Debug, Serialize, Deserialize)]
593pub struct ReadVar<T, C = VarNotClaimed> {
594 backing_var: ReadVarBacking<T>,
595 #[serde(skip)]
596 _kind: std::marker::PhantomData<C>,
597}
598
599pub type ClaimedReadVar<T> = ReadVar<T, VarClaimed>;
602
603impl<T: Serialize + DeserializeOwned, C> Clone for ReadVar<T, C> {
605 fn clone(&self) -> Self {
606 ReadVar {
607 backing_var: self.backing_var.clone(),
608 _kind: std::marker::PhantomData,
609 }
610 }
611}
612
613#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
614enum ReadVarBacking<T> {
615 RuntimeVar {
616 var: String,
617 is_side_effect: bool,
624 },
625 Inline(T),
626}
627
628impl<T: Serialize + DeserializeOwned> Clone for ReadVarBacking<T> {
630 fn clone(&self) -> Self {
631 match self {
632 Self::RuntimeVar {
633 var,
634 is_side_effect,
635 } => Self::RuntimeVar {
636 var: var.clone(),
637 is_side_effect: *is_side_effect,
638 },
639 Self::Inline(v) => {
640 Self::Inline(serde_json::from_value(serde_json::to_value(v).unwrap()).unwrap())
641 }
642 }
643 }
644}
645
646impl<T: Serialize + DeserializeOwned> ReadVar<T> {
647 fn into_claimed(self) -> ReadVar<T, VarClaimed> {
649 let Self { backing_var, _kind } = self;
650
651 ReadVar {
652 backing_var,
653 _kind: std::marker::PhantomData,
654 }
655 }
656
657 #[must_use]
666 pub fn into_side_effect(self) -> ReadVar<SideEffect> {
667 ReadVar {
668 backing_var: match self.backing_var {
669 ReadVarBacking::RuntimeVar {
670 var,
671 is_side_effect: _,
672 } => ReadVarBacking::RuntimeVar {
673 var,
674 is_side_effect: true,
675 },
676 ReadVarBacking::Inline(_) => ReadVarBacking::Inline(()),
677 },
678 _kind: std::marker::PhantomData,
679 }
680 }
681
682 #[track_caller]
685 #[must_use]
686 pub fn map<F, U>(&self, ctx: &mut NodeCtx<'_>, f: F) -> ReadVar<U>
687 where
688 T: 'static,
689 U: Serialize + DeserializeOwned + 'static,
690 F: FnOnce(T) -> U + 'static,
691 {
692 let (read_from, write_into) = ctx.new_var();
693 self.write_into_with(ctx, write_into, f);
694 read_from
695 }
696
697 #[track_caller]
700 pub fn write_into_with<F, U>(&self, ctx: &mut NodeCtx<'_>, write_into: WriteVar<U>, f: F)
701 where
702 T: 'static,
703 U: Serialize + DeserializeOwned + 'static,
704 F: FnOnce(T) -> U + 'static,
705 {
706 let this = self.clone();
707 ctx.emit_minor_rust_step("🌼 write_into Var", move |ctx| {
708 let this = this.claim(ctx);
709 let write_into = write_into.claim(ctx);
710 move |rt| {
711 let this = rt.read(this);
712 rt.write(write_into, &f(this));
713 }
714 });
715 }
716
717 #[track_caller]
719 pub fn write_into(&self, ctx: &mut NodeCtx<'_>, write_into: WriteVar<T>)
720 where
721 T: 'static,
722 {
723 self.write_into_with(ctx, write_into, |x| x);
724 }
725
726 #[track_caller]
729 #[must_use]
730 pub fn zip<U>(&self, ctx: &mut NodeCtx<'_>, other: ReadVar<U>) -> ReadVar<(T, U)>
731 where
732 T: 'static,
733 U: Serialize + DeserializeOwned + 'static,
734 {
735 let (read_from, write_into) = ctx.new_var();
736 let this = self.clone();
737 ctx.emit_minor_rust_step("🌼 Zip Vars", move |ctx| {
738 let this = this.claim(ctx);
739 let other = other.claim(ctx);
740 let write_into = write_into.claim(ctx);
741 move |rt| {
742 let this = rt.read(this);
743 let other = rt.read(other);
744 rt.write(write_into, &(this, other));
745 }
746 });
747 read_from
748 }
749
750 #[track_caller]
755 #[must_use]
756 pub fn from_static(val: T) -> ReadVar<T>
757 where
758 T: 'static,
759 {
760 ReadVar {
761 backing_var: ReadVarBacking::Inline(val),
762 _kind: std::marker::PhantomData,
763 }
764 }
765
766 pub fn get_static(&self) -> Option<T> {
775 match self.clone().backing_var {
776 ReadVarBacking::Inline(v) => Some(v),
777 _ => None,
778 }
779 }
780
781 #[track_caller]
783 #[must_use]
784 pub fn transpose_vec(ctx: &mut NodeCtx<'_>, vec: Vec<ReadVar<T>>) -> ReadVar<Vec<T>>
785 where
786 T: 'static,
787 {
788 let (read_from, write_into) = ctx.new_var();
789 ctx.emit_minor_rust_step("🌼 Transpose Vec<ReadVar<T>>", move |ctx| {
790 let vec = vec.claim(ctx);
791 let write_into = write_into.claim(ctx);
792 move |rt| {
793 let mut v = Vec::new();
794 for var in vec {
795 v.push(rt.read(var));
796 }
797 rt.write(write_into, &v);
798 }
799 });
800 read_from
801 }
802
803 #[must_use]
819 pub fn depending_on<U>(&self, ctx: &mut NodeCtx<'_>, other: &ReadVar<U>) -> Self
820 where
821 T: 'static,
822 U: Serialize + DeserializeOwned + 'static,
823 {
824 ctx.emit_minor_rust_stepv("🌼 Add dependency", |ctx| {
827 let this = self.clone().claim(ctx);
828 other.clone().claim(ctx);
829 move |rt| rt.read(this)
830 })
831 }
832
833 pub fn claim_unused(self, ctx: &mut NodeCtx<'_>) {
836 match self.backing_var {
837 ReadVarBacking::RuntimeVar {
838 var,
839 is_side_effect: _,
840 } => ctx.backend.borrow_mut().on_unused_read_var(&var),
841 ReadVarBacking::Inline(_) => {}
842 }
843 }
844
845 pub(crate) fn into_json(self) -> ReadVar<serde_json::Value> {
846 match self.backing_var {
847 ReadVarBacking::RuntimeVar {
848 var,
849 is_side_effect,
850 } => ReadVar {
851 backing_var: ReadVarBacking::RuntimeVar {
852 var,
853 is_side_effect,
854 },
855 _kind: std::marker::PhantomData,
856 },
857 ReadVarBacking::Inline(v) => ReadVar {
858 backing_var: ReadVarBacking::Inline(serde_json::to_value(v).unwrap()),
859 _kind: std::marker::PhantomData,
860 },
861 }
862 }
863}
864
865#[must_use]
871pub fn thin_air_read_runtime_var<T>(backing_var: String) -> ReadVar<T>
872where
873 T: Serialize + DeserializeOwned,
874{
875 ReadVar {
876 backing_var: ReadVarBacking::RuntimeVar {
877 var: backing_var,
878 is_side_effect: false,
879 },
880 _kind: std::marker::PhantomData,
881 }
882}
883
884#[must_use]
890pub fn thin_air_write_runtime_var<T>(backing_var: String) -> WriteVar<T>
891where
892 T: Serialize + DeserializeOwned,
893{
894 WriteVar {
895 backing_var,
896 is_side_effect: false,
897 _kind: std::marker::PhantomData,
898 }
899}
900
901pub fn read_var_internals<T: Serialize + DeserializeOwned, C>(
907 var: &ReadVar<T, C>,
908) -> (Option<String>, bool) {
909 match var.backing_var {
910 ReadVarBacking::RuntimeVar {
911 var: ref s,
912 is_side_effect,
913 } => (Some(s.clone()), is_side_effect),
914 ReadVarBacking::Inline(_) => (None, false),
915 }
916}
917
918pub trait ImportCtxBackend {
919 fn on_possible_dep(&mut self, node_handle: NodeHandle);
920}
921
922pub struct ImportCtx<'a> {
924 backend: &'a mut dyn ImportCtxBackend,
925}
926
927impl ImportCtx<'_> {
928 pub fn import<N: FlowNodeBase + 'static>(&mut self) {
930 self.backend.on_possible_dep(NodeHandle::from_type::<N>())
931 }
932}
933
934pub fn new_import_ctx(backend: &mut dyn ImportCtxBackend) -> ImportCtx<'_> {
935 ImportCtx { backend }
936}
937
938#[derive(Debug)]
939pub enum CtxAnchor {
940 PostJob,
941}
942
943pub trait NodeCtxBackend {
944 fn current_node(&self) -> NodeHandle;
946
947 fn on_new_var(&mut self) -> String;
952
953 fn on_claimed_runtime_var(&mut self, var: &str, is_read: bool);
955
956 fn on_unused_read_var(&mut self, var: &str);
958
959 fn on_request(&mut self, node_handle: NodeHandle, req: anyhow::Result<Box<[u8]>>);
967
968 fn on_config(&mut self, node_handle: NodeHandle, config: anyhow::Result<Box<[u8]>>);
972
973 fn on_emit_rust_step(
974 &mut self,
975 label: &str,
976 can_merge: bool,
977 code: Box<dyn for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) -> anyhow::Result<()>>,
978 );
979
980 fn on_emit_ado_step(
981 &mut self,
982 label: &str,
983 yaml_snippet: Box<dyn for<'a> FnOnce(&'a mut AdoStepServices<'_>) -> String>,
984 inline_script: Option<
985 Box<dyn for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) -> anyhow::Result<()>>,
986 >,
987 condvar: Option<String>,
988 );
989
990 fn on_emit_gh_step(
991 &mut self,
992 label: &str,
993 uses: &str,
994 with: BTreeMap<String, ClaimedGhParam>,
995 condvar: Option<String>,
996 outputs: BTreeMap<String, Vec<GhOutput>>,
997 permissions: BTreeMap<GhPermission, GhPermissionValue>,
998 gh_to_rust: Vec<GhToRust>,
999 rust_to_gh: Vec<RustToGh>,
1000 );
1001
1002 fn on_emit_side_effect_step(&mut self);
1003
1004 fn backend(&mut self) -> FlowBackend;
1005 fn platform(&mut self) -> FlowPlatform;
1006 fn arch(&mut self) -> FlowArch;
1007
1008 fn persistent_dir_path_var(&mut self) -> Option<String>;
1012}
1013
1014pub fn new_node_ctx(backend: &mut dyn NodeCtxBackend) -> NodeCtx<'_> {
1015 NodeCtx {
1016 backend: Rc::new(RefCell::new(backend)),
1017 }
1018}
1019
1020#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
1022pub enum FlowBackend {
1023 Local,
1025 Ado,
1027 Github,
1029}
1030
1031#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
1033pub enum FlowPlatformKind {
1034 Windows,
1035 Unix,
1036}
1037
1038#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
1040pub enum FlowPlatformLinuxDistro {
1041 Fedora,
1043 Ubuntu,
1045 AzureLinux,
1047 Arch,
1049 Nix,
1051 Unknown,
1053}
1054
1055#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
1057#[non_exhaustive]
1058pub enum FlowPlatform {
1059 Windows,
1061 Linux(FlowPlatformLinuxDistro),
1063 MacOs,
1065}
1066
1067impl FlowPlatform {
1068 pub fn kind(&self) -> FlowPlatformKind {
1069 match self {
1070 Self::Windows => FlowPlatformKind::Windows,
1071 Self::Linux(_) | Self::MacOs => FlowPlatformKind::Unix,
1072 }
1073 }
1074
1075 fn as_str(&self) -> &'static str {
1076 match self {
1077 Self::Windows => "windows",
1078 Self::Linux(_) => "linux",
1079 Self::MacOs => "macos",
1080 }
1081 }
1082
1083 pub fn exe_suffix(&self) -> &'static str {
1085 if self == &Self::Windows { ".exe" } else { "" }
1086 }
1087
1088 pub fn binary(&self, name: &str) -> String {
1090 format!("{}{}", name, self.exe_suffix())
1091 }
1092}
1093
1094impl std::fmt::Display for FlowPlatform {
1095 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1096 f.pad(self.as_str())
1097 }
1098}
1099
1100#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
1102#[non_exhaustive]
1103pub enum FlowArch {
1104 X86_64,
1105 Aarch64,
1106}
1107
1108impl FlowArch {
1109 fn as_str(&self) -> &'static str {
1110 match self {
1111 Self::X86_64 => "x86_64",
1112 Self::Aarch64 => "aarch64",
1113 }
1114 }
1115}
1116
1117impl std::fmt::Display for FlowArch {
1118 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1119 f.pad(self.as_str())
1120 }
1121}
1122
1123pub struct StepCtx<'a> {
1125 backend: Rc<RefCell<&'a mut dyn NodeCtxBackend>>,
1126}
1127
1128impl StepCtx<'_> {
1129 pub fn backend(&self) -> FlowBackend {
1132 self.backend.borrow_mut().backend()
1133 }
1134
1135 pub fn platform(&self) -> FlowPlatform {
1138 self.backend.borrow_mut().platform()
1139 }
1140}
1141
1142const NO_ADO_INLINE_SCRIPT: Option<
1143 for<'a> fn(&'a mut RustRuntimeServices<'_>) -> anyhow::Result<()>,
1144> = None;
1145
1146pub struct NodeCtx<'a> {
1148 backend: Rc<RefCell<&'a mut dyn NodeCtxBackend>>,
1149}
1150
1151impl<'ctx> NodeCtx<'ctx> {
1152 pub fn emit_rust_step<F, G>(&mut self, label: impl AsRef<str>, code: F) -> ReadVar<SideEffect>
1158 where
1159 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> G,
1160 G: for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) -> anyhow::Result<()> + 'static,
1161 {
1162 self.emit_rust_step_inner(label.as_ref(), false, code)
1163 }
1164
1165 pub fn emit_minor_rust_step<F, G>(
1171 &mut self,
1172 label: impl AsRef<str>,
1173 code: F,
1174 ) -> ReadVar<SideEffect>
1175 where
1176 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> G,
1177 G: for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) + 'static,
1178 {
1179 self.emit_rust_step_inner(label.as_ref(), true, |ctx| {
1180 let f = code(ctx);
1181 |rt| {
1182 f(rt);
1183 Ok(())
1184 }
1185 })
1186 }
1187
1188 #[must_use]
1209 #[track_caller]
1210 pub fn emit_rust_stepv<T, F, G>(&mut self, label: impl AsRef<str>, code: F) -> ReadVar<T>
1211 where
1212 T: Serialize + DeserializeOwned + 'static,
1213 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> G,
1214 G: for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) -> anyhow::Result<T> + 'static,
1215 {
1216 self.emit_rust_stepv_inner(label.as_ref(), false, code)
1217 }
1218
1219 #[must_use]
1243 #[track_caller]
1244 pub fn emit_minor_rust_stepv<T, F, G>(&mut self, label: impl AsRef<str>, code: F) -> ReadVar<T>
1245 where
1246 T: Serialize + DeserializeOwned + 'static,
1247 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> G,
1248 G: for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) -> T + 'static,
1249 {
1250 self.emit_rust_stepv_inner(label.as_ref(), true, |ctx| {
1251 let f = code(ctx);
1252 |rt| Ok(f(rt))
1253 })
1254 }
1255
1256 fn emit_rust_step_inner<F, G>(
1257 &mut self,
1258 label: &str,
1259 can_merge: bool,
1260 code: F,
1261 ) -> ReadVar<SideEffect>
1262 where
1263 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> G,
1264 G: for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) -> anyhow::Result<()> + 'static,
1265 {
1266 let (read, write) = self.new_prefixed_var("auto_se");
1267
1268 let ctx = &mut StepCtx {
1269 backend: self.backend.clone(),
1270 };
1271 write.claim(ctx);
1272
1273 let code = code(ctx);
1274 self.backend
1275 .borrow_mut()
1276 .on_emit_rust_step(label.as_ref(), can_merge, Box::new(code));
1277 read
1278 }
1279
1280 #[must_use]
1281 #[track_caller]
1282 fn emit_rust_stepv_inner<T, F, G>(
1283 &mut self,
1284 label: impl AsRef<str>,
1285 can_merge: bool,
1286 code: F,
1287 ) -> ReadVar<T>
1288 where
1289 T: Serialize + DeserializeOwned + 'static,
1290 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> G,
1291 G: for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) -> anyhow::Result<T> + 'static,
1292 {
1293 let (read, write) = self.new_var();
1294
1295 let ctx = &mut StepCtx {
1296 backend: self.backend.clone(),
1297 };
1298 let write = write.claim(ctx);
1299
1300 let code = code(ctx);
1301 self.backend.borrow_mut().on_emit_rust_step(
1302 label.as_ref(),
1303 can_merge,
1304 Box::new(|rt| {
1305 let val = code(rt)?;
1306 rt.write(write, &val);
1307 Ok(())
1308 }),
1309 );
1310 read
1311 }
1312
1313 #[track_caller]
1315 #[must_use]
1316 pub fn get_ado_variable(&mut self, ado_var: AdoRuntimeVar) -> ReadVar<String> {
1317 let (var, write_var) = self.new_var();
1318 self.emit_ado_step(format!("🌼 read {}", ado_var.as_raw_var_name()), |ctx| {
1319 let write_var = write_var.claim(ctx);
1320 |rt| {
1321 rt.set_var(write_var, ado_var);
1322 "".into()
1323 }
1324 });
1325 var
1326 }
1327
1328 pub fn emit_ado_step<F, G>(&mut self, display_name: impl AsRef<str>, yaml_snippet: F)
1330 where
1331 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> G,
1332 G: for<'a> FnOnce(&'a mut AdoStepServices<'_>) -> String + 'static,
1333 {
1334 self.emit_ado_step_inner(display_name, None, |ctx| {
1335 (yaml_snippet(ctx), NO_ADO_INLINE_SCRIPT)
1336 })
1337 }
1338
1339 pub fn emit_ado_step_with_condition<F, G>(
1342 &mut self,
1343 display_name: impl AsRef<str>,
1344 cond: ReadVar<bool>,
1345 yaml_snippet: F,
1346 ) where
1347 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> G,
1348 G: for<'a> FnOnce(&'a mut AdoStepServices<'_>) -> String + 'static,
1349 {
1350 self.emit_ado_step_inner(display_name, Some(cond), |ctx| {
1351 (yaml_snippet(ctx), NO_ADO_INLINE_SCRIPT)
1352 })
1353 }
1354
1355 pub fn emit_ado_step_with_condition_optional<F, G>(
1358 &mut self,
1359 display_name: impl AsRef<str>,
1360 cond: Option<ReadVar<bool>>,
1361 yaml_snippet: F,
1362 ) where
1363 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> G,
1364 G: for<'a> FnOnce(&'a mut AdoStepServices<'_>) -> String + 'static,
1365 {
1366 self.emit_ado_step_inner(display_name, cond, |ctx| {
1367 (yaml_snippet(ctx), NO_ADO_INLINE_SCRIPT)
1368 })
1369 }
1370
1371 pub fn emit_ado_step_with_inline_script<F, G, H>(
1400 &mut self,
1401 display_name: impl AsRef<str>,
1402 yaml_snippet: F,
1403 ) where
1404 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> (G, H),
1405 G: for<'a> FnOnce(&'a mut AdoStepServices<'_>) -> String + 'static,
1406 H: for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) -> anyhow::Result<()> + 'static,
1407 {
1408 self.emit_ado_step_inner(display_name, None, |ctx| {
1409 let (f, g) = yaml_snippet(ctx);
1410 (f, Some(g))
1411 })
1412 }
1413
1414 fn emit_ado_step_inner<F, G, H>(
1415 &mut self,
1416 display_name: impl AsRef<str>,
1417 cond: Option<ReadVar<bool>>,
1418 yaml_snippet: F,
1419 ) where
1420 F: for<'a> FnOnce(&'a mut StepCtx<'_>) -> (G, Option<H>),
1421 G: for<'a> FnOnce(&'a mut AdoStepServices<'_>) -> String + 'static,
1422 H: for<'a> FnOnce(&'a mut RustRuntimeServices<'_>) -> anyhow::Result<()> + 'static,
1423 {
1424 let condvar = match cond.map(|c| c.backing_var) {
1425 Some(ReadVarBacking::Inline(cond)) => {
1427 if !cond {
1428 return;
1429 } else {
1430 None
1431 }
1432 }
1433 Some(ReadVarBacking::RuntimeVar {
1434 var,
1435 is_side_effect,
1436 }) => {
1437 assert!(!is_side_effect);
1438 self.backend.borrow_mut().on_claimed_runtime_var(&var, true);
1439 Some(var)
1440 }
1441 None => None,
1442 };
1443
1444 let (yaml_snippet, inline_script) = yaml_snippet(&mut StepCtx {
1445 backend: self.backend.clone(),
1446 });
1447 self.backend.borrow_mut().on_emit_ado_step(
1448 display_name.as_ref(),
1449 Box::new(yaml_snippet),
1450 if let Some(inline_script) = inline_script {
1451 Some(Box::new(inline_script))
1452 } else {
1453 None
1454 },
1455 condvar,
1456 );
1457 }
1458
1459 #[track_caller]
1461 #[must_use]
1462 pub fn get_gh_context_var(&mut self) -> GhContextVarReader<'ctx, Root> {
1463 GhContextVarReader {
1464 ctx: NodeCtx {
1465 backend: self.backend.clone(),
1466 },
1467 _state: std::marker::PhantomData,
1468 }
1469 }
1470
1471 pub fn emit_gh_step(
1473 &mut self,
1474 display_name: impl AsRef<str>,
1475 uses: impl AsRef<str>,
1476 ) -> GhStepBuilder {
1477 GhStepBuilder::new(display_name, uses)
1478 }
1479
1480 fn emit_gh_step_inner(
1481 &mut self,
1482 display_name: impl AsRef<str>,
1483 cond: Option<ReadVar<bool>>,
1484 uses: impl AsRef<str>,
1485 with: Option<BTreeMap<String, GhParam>>,
1486 outputs: BTreeMap<String, Vec<WriteVar<String>>>,
1487 run_after: Vec<ReadVar<SideEffect>>,
1488 permissions: BTreeMap<GhPermission, GhPermissionValue>,
1489 ) {
1490 let condvar = match cond.map(|c| c.backing_var) {
1491 Some(ReadVarBacking::Inline(cond)) => {
1493 if !cond {
1494 return;
1495 } else {
1496 None
1497 }
1498 }
1499 Some(ReadVarBacking::RuntimeVar {
1500 var,
1501 is_side_effect,
1502 }) => {
1503 assert!(!is_side_effect);
1504 self.backend.borrow_mut().on_claimed_runtime_var(&var, true);
1505 Some(var)
1506 }
1507 None => None,
1508 };
1509
1510 let with = with
1511 .unwrap_or_default()
1512 .into_iter()
1513 .map(|(k, v)| {
1514 (
1515 k.clone(),
1516 v.claim(&mut StepCtx {
1517 backend: self.backend.clone(),
1518 }),
1519 )
1520 })
1521 .collect();
1522
1523 for var in run_after {
1524 var.claim(&mut StepCtx {
1525 backend: self.backend.clone(),
1526 });
1527 }
1528
1529 let outputvars = outputs
1530 .into_iter()
1531 .map(|(name, vars)| {
1532 (
1533 name,
1534 vars.into_iter()
1535 .map(|var| {
1536 let var = var.claim(&mut StepCtx {
1537 backend: self.backend.clone(),
1538 });
1539 GhOutput {
1540 backing_var: var.backing_var,
1541 is_secret: false,
1542 is_object: false,
1543 }
1544 })
1545 .collect(),
1546 )
1547 })
1548 .collect();
1549
1550 self.backend.borrow_mut().on_emit_gh_step(
1551 display_name.as_ref(),
1552 uses.as_ref(),
1553 with,
1554 condvar,
1555 outputvars,
1556 permissions,
1557 Vec::new(),
1558 Vec::new(),
1559 );
1560 }
1561
1562 pub fn emit_side_effect_step(
1570 &mut self,
1571 use_side_effects: impl IntoIterator<Item = ReadVar<SideEffect>>,
1572 resolve_side_effects: impl IntoIterator<Item = WriteVar<SideEffect>>,
1573 ) {
1574 let mut backend = self.backend.borrow_mut();
1575 for var in use_side_effects.into_iter() {
1576 if let ReadVarBacking::RuntimeVar {
1577 var,
1578 is_side_effect: _,
1579 } = &var.backing_var
1580 {
1581 backend.on_claimed_runtime_var(var, true);
1582 }
1583 }
1584
1585 for var in resolve_side_effects.into_iter() {
1586 backend.on_claimed_runtime_var(&var.backing_var, false);
1587 }
1588
1589 backend.on_emit_side_effect_step();
1590 }
1591
1592 pub fn backend(&self) -> FlowBackend {
1595 self.backend.borrow_mut().backend()
1596 }
1597
1598 pub fn platform(&self) -> FlowPlatform {
1601 self.backend.borrow_mut().platform()
1602 }
1603
1604 pub fn arch(&self) -> FlowArch {
1606 self.backend.borrow_mut().arch()
1607 }
1608
1609 pub fn req<R>(&mut self, req: R)
1611 where
1612 R: IntoRequest + 'static,
1613 {
1614 let mut backend = self.backend.borrow_mut();
1615 backend.on_request(
1616 NodeHandle::from_type::<R::Node>(),
1617 serde_json::to_vec(&req.into_request())
1618 .map(Into::into)
1619 .map_err(Into::into),
1620 );
1621 }
1622
1623 pub fn config<C>(&mut self, config: C)
1628 where
1629 C: IntoConfig + 'static,
1630 {
1631 let mut backend = self.backend.borrow_mut();
1632 backend.on_config(
1633 NodeHandle::from_type::<C::Node>(),
1634 serde_json::to_vec(&config)
1635 .map(Into::into)
1636 .map_err(Into::into),
1637 );
1638 }
1639
1640 #[track_caller]
1643 #[must_use]
1644 pub fn reqv<T, R>(&mut self, f: impl FnOnce(WriteVar<T>) -> R) -> ReadVar<T>
1645 where
1646 T: Serialize + DeserializeOwned,
1647 R: IntoRequest + 'static,
1648 {
1649 let (read, write) = self.new_var();
1650 self.req::<R>(f(write));
1651 read
1652 }
1653
1654 pub fn requests<N>(&mut self, reqs: impl IntoIterator<Item = N::Request>)
1656 where
1657 N: FlowNodeBase + 'static,
1658 {
1659 let mut backend = self.backend.borrow_mut();
1660 for req in reqs.into_iter() {
1661 backend.on_request(
1662 NodeHandle::from_type::<N>(),
1663 serde_json::to_vec(&req).map(Into::into).map_err(Into::into),
1664 );
1665 }
1666 }
1667
1668 #[track_caller]
1671 #[must_use]
1672 pub fn new_var<T>(&self) -> (ReadVar<T>, WriteVar<T>)
1673 where
1674 T: Serialize + DeserializeOwned,
1675 {
1676 self.new_prefixed_var("")
1677 }
1678
1679 #[track_caller]
1680 #[must_use]
1681 fn new_prefixed_var<T>(&self, prefix: &'static str) -> (ReadVar<T>, WriteVar<T>)
1682 where
1683 T: Serialize + DeserializeOwned,
1684 {
1685 let caller = std::panic::Location::caller().file().replace('\\', "/");
1692
1693 let caller = caller
1709 .split_once("flowey/")
1710 .expect("due to a known limitation with flowey, all flowey code must have an ancestor dir called 'flowey/' somewhere in its full path")
1711 .1;
1712
1713 let colon = if prefix.is_empty() { "" } else { ":" };
1714 let ordinal = self.backend.borrow_mut().on_new_var();
1715 let backing_var = format!("{prefix}{colon}{ordinal}:{caller}");
1716
1717 (
1718 ReadVar {
1719 backing_var: ReadVarBacking::RuntimeVar {
1720 var: backing_var.clone(),
1721 is_side_effect: false,
1722 },
1723 _kind: std::marker::PhantomData,
1724 },
1725 WriteVar {
1726 backing_var,
1727 is_side_effect: false,
1728 _kind: std::marker::PhantomData,
1729 },
1730 )
1731 }
1732
1733 #[track_caller]
1744 #[must_use]
1745 pub fn new_post_job_side_effect(&self) -> (ReadVar<SideEffect>, WriteVar<SideEffect>) {
1746 self.new_prefixed_var("post_job")
1747 }
1748
1749 #[track_caller]
1762 #[must_use]
1763 pub fn persistent_dir(&mut self) -> Option<ReadVar<PathBuf>> {
1764 let path: ReadVar<PathBuf> = ReadVar {
1765 backing_var: ReadVarBacking::RuntimeVar {
1766 var: self.backend.borrow_mut().persistent_dir_path_var()?,
1767 is_side_effect: false,
1768 },
1769 _kind: std::marker::PhantomData,
1770 };
1771
1772 let folder_name = self
1773 .backend
1774 .borrow_mut()
1775 .current_node()
1776 .modpath()
1777 .replace("::", "__");
1778
1779 Some(
1780 self.emit_rust_stepv("🌼 Create persistent store dir", |ctx| {
1781 let path = path.claim(ctx);
1782 |rt| {
1783 let dir = rt.read(path).join(folder_name);
1784 fs_err::create_dir_all(&dir)?;
1785 Ok(dir)
1786 }
1787 }),
1788 )
1789 }
1790
1791 pub fn supports_persistent_dir(&mut self) -> bool {
1793 self.backend
1794 .borrow_mut()
1795 .persistent_dir_path_var()
1796 .is_some()
1797 }
1798}
1799
1800pub trait RuntimeVarDb {
1803 fn get_var(&mut self, var_name: &str) -> (Vec<u8>, bool) {
1804 self.try_get_var(var_name)
1805 .unwrap_or_else(|| panic!("db is missing var {}", var_name))
1806 }
1807
1808 fn try_get_var(&mut self, var_name: &str) -> Option<(Vec<u8>, bool)>;
1809 fn set_var(&mut self, var_name: &str, is_secret: bool, value: Vec<u8>);
1810}
1811
1812impl RuntimeVarDb for Box<dyn RuntimeVarDb> {
1813 fn try_get_var(&mut self, var_name: &str) -> Option<(Vec<u8>, bool)> {
1814 (**self).try_get_var(var_name)
1815 }
1816
1817 fn set_var(&mut self, var_name: &str, is_secret: bool, value: Vec<u8>) {
1818 (**self).set_var(var_name, is_secret, value)
1819 }
1820}
1821
1822pub mod steps {
1823 pub mod ado {
1824 use crate::node::ClaimedReadVar;
1825 use crate::node::ClaimedWriteVar;
1826 use crate::node::ReadVarBacking;
1827 use serde::Deserialize;
1828 use serde::Serialize;
1829 use std::borrow::Cow;
1830
1831 #[derive(Debug, Clone, Serialize, Deserialize)]
1837 pub struct AdoResourcesRepositoryId {
1838 pub(crate) repo_id: String,
1839 }
1840
1841 impl AdoResourcesRepositoryId {
1842 pub fn new_self() -> Self {
1848 Self {
1849 repo_id: "self".into(),
1850 }
1851 }
1852
1853 pub fn dangerous_get_raw_id(&self) -> &str {
1859 &self.repo_id
1860 }
1861
1862 pub fn dangerous_new(repo_id: &str) -> Self {
1868 Self {
1869 repo_id: repo_id.into(),
1870 }
1871 }
1872 }
1873
1874 #[derive(Clone, Debug, Serialize, Deserialize)]
1879 pub struct AdoRuntimeVar {
1880 is_secret: bool,
1881 ado_var: Cow<'static, str>,
1882 }
1883
1884 impl AdoRuntimeVar {
1885 pub const BUILD_SOURCE_BRANCH: AdoRuntimeVar = AdoRuntimeVar::new("build.SourceBranch");
1891
1892 pub const BUILD_BUILD_NUMBER: AdoRuntimeVar = AdoRuntimeVar::new("build.BuildNumber");
1894
1895 pub const SYSTEM_ACCESS_TOKEN: AdoRuntimeVar =
1897 AdoRuntimeVar::new_secret("System.AccessToken");
1898
1899 pub const SYSTEM_JOB_ATTEMPT: AdoRuntimeVar =
1901 AdoRuntimeVar::new_secret("System.JobAttempt");
1902
1903 pub const PIPELINE_WORKSPACE: AdoRuntimeVar = AdoRuntimeVar::new("Pipeline.Workspace");
1905 }
1906
1907 impl AdoRuntimeVar {
1908 const fn new(s: &'static str) -> Self {
1909 Self {
1910 is_secret: false,
1911 ado_var: Cow::Borrowed(s),
1912 }
1913 }
1914
1915 const fn new_secret(s: &'static str) -> Self {
1916 Self {
1917 is_secret: true,
1918 ado_var: Cow::Borrowed(s),
1919 }
1920 }
1921
1922 pub fn is_secret(&self) -> bool {
1924 self.is_secret
1925 }
1926
1927 pub fn as_raw_var_name(&self) -> String {
1929 self.ado_var.as_ref().into()
1930 }
1931
1932 pub fn dangerous_from_global(ado_var_name: impl AsRef<str>, is_secret: bool) -> Self {
1940 Self {
1941 is_secret,
1942 ado_var: ado_var_name.as_ref().to_owned().into(),
1943 }
1944 }
1945 }
1946
1947 pub fn new_ado_step_services(
1948 fresh_ado_var: &mut dyn FnMut() -> String,
1949 ) -> AdoStepServices<'_> {
1950 AdoStepServices {
1951 fresh_ado_var,
1952 ado_to_rust: Vec::new(),
1953 rust_to_ado: Vec::new(),
1954 }
1955 }
1956
1957 pub struct CompletedAdoStepServices {
1958 pub ado_to_rust: Vec<(String, String, bool)>,
1959 pub rust_to_ado: Vec<(String, String)>,
1960 }
1961
1962 impl CompletedAdoStepServices {
1963 pub fn from_ado_step_services(access: AdoStepServices<'_>) -> Self {
1964 let AdoStepServices {
1965 fresh_ado_var: _,
1966 ado_to_rust,
1967 rust_to_ado,
1968 } = access;
1969
1970 Self {
1971 ado_to_rust,
1972 rust_to_ado,
1973 }
1974 }
1975 }
1976
1977 pub struct AdoStepServices<'a> {
1978 fresh_ado_var: &'a mut dyn FnMut() -> String,
1979 ado_to_rust: Vec<(String, String, bool)>,
1980 rust_to_ado: Vec<(String, String)>,
1981 }
1982
1983 impl AdoStepServices<'_> {
1984 pub fn resolve_repository_id(&self, repo_id: AdoResourcesRepositoryId) -> String {
1987 repo_id.repo_id
1988 }
1989
1990 pub fn set_var(&mut self, var: ClaimedWriteVar<String>, from_ado_var: AdoRuntimeVar) {
1996 self.ado_to_rust.push((
1997 from_ado_var.ado_var.into(),
1998 var.backing_var,
1999 from_ado_var.is_secret,
2000 ))
2001 }
2002
2003 pub fn get_var(&mut self, var: ClaimedReadVar<String>) -> AdoRuntimeVar {
2005 let backing_var = if let ReadVarBacking::RuntimeVar {
2006 var,
2007 is_side_effect,
2008 } = &var.backing_var
2009 {
2010 assert!(!is_side_effect);
2011 var
2012 } else {
2013 todo!("support inline ado read vars")
2014 };
2015
2016 let new_ado_var_name = (self.fresh_ado_var)();
2017
2018 self.rust_to_ado
2019 .push((backing_var.clone(), new_ado_var_name.clone()));
2020 AdoRuntimeVar::dangerous_from_global(new_ado_var_name, false)
2021 }
2022 }
2023 }
2024
2025 pub mod github {
2026 use crate::node::ClaimVar;
2027 use crate::node::NodeCtx;
2028 use crate::node::ReadVar;
2029 use crate::node::ReadVarBacking;
2030 use crate::node::SideEffect;
2031 use crate::node::StepCtx;
2032 use crate::node::VarClaimed;
2033 use crate::node::VarNotClaimed;
2034 use crate::node::WriteVar;
2035 use std::collections::BTreeMap;
2036
2037 pub struct GhStepBuilder {
2038 display_name: String,
2039 cond: Option<ReadVar<bool>>,
2040 uses: String,
2041 with: Option<BTreeMap<String, GhParam>>,
2042 outputs: BTreeMap<String, Vec<WriteVar<String>>>,
2043 run_after: Vec<ReadVar<SideEffect>>,
2044 permissions: BTreeMap<GhPermission, GhPermissionValue>,
2045 }
2046
2047 impl GhStepBuilder {
2048 pub fn new(display_name: impl AsRef<str>, uses: impl AsRef<str>) -> Self {
2063 Self {
2064 display_name: display_name.as_ref().into(),
2065 cond: None,
2066 uses: uses.as_ref().into(),
2067 with: None,
2068 outputs: BTreeMap::new(),
2069 run_after: Vec::new(),
2070 permissions: BTreeMap::new(),
2071 }
2072 }
2073
2074 pub fn condition(mut self, cond: ReadVar<bool>) -> Self {
2081 self.cond = Some(cond);
2082 self
2083 }
2084
2085 pub fn with(mut self, k: impl AsRef<str>, v: impl Into<GhParam>) -> Self {
2111 self.with.get_or_insert_with(BTreeMap::new);
2112 if let Some(with) = &mut self.with {
2113 with.insert(k.as_ref().to_string(), v.into());
2114 }
2115 self
2116 }
2117
2118 pub fn output(mut self, k: impl AsRef<str>, v: WriteVar<String>) -> Self {
2127 self.outputs
2128 .entry(k.as_ref().to_string())
2129 .or_default()
2130 .push(v);
2131 self
2132 }
2133
2134 pub fn run_after(mut self, side_effect: ReadVar<SideEffect>) -> Self {
2136 self.run_after.push(side_effect);
2137 self
2138 }
2139
2140 pub fn requires_permission(
2145 mut self,
2146 perm: GhPermission,
2147 value: GhPermissionValue,
2148 ) -> Self {
2149 self.permissions.insert(perm, value);
2150 self
2151 }
2152
2153 #[track_caller]
2155 pub fn finish(self, ctx: &mut NodeCtx<'_>) -> ReadVar<SideEffect> {
2156 let (side_effect, claim_side_effect) = ctx.new_prefixed_var("auto_se");
2157 ctx.backend
2158 .borrow_mut()
2159 .on_claimed_runtime_var(&claim_side_effect.backing_var, false);
2160
2161 ctx.emit_gh_step_inner(
2162 self.display_name,
2163 self.cond,
2164 self.uses,
2165 self.with,
2166 self.outputs,
2167 self.run_after,
2168 self.permissions,
2169 );
2170
2171 side_effect
2172 }
2173 }
2174
2175 #[derive(Clone, Debug)]
2176 pub enum GhParam<C = VarNotClaimed> {
2177 Static(String),
2178 FloweyVar(ReadVar<String, C>),
2179 }
2180
2181 impl From<String> for GhParam {
2182 fn from(param: String) -> GhParam {
2183 GhParam::Static(param)
2184 }
2185 }
2186
2187 impl From<&str> for GhParam {
2188 fn from(param: &str) -> GhParam {
2189 GhParam::Static(param.to_string())
2190 }
2191 }
2192
2193 impl From<ReadVar<String>> for GhParam {
2194 fn from(param: ReadVar<String>) -> GhParam {
2195 GhParam::FloweyVar(param)
2196 }
2197 }
2198
2199 pub type ClaimedGhParam = GhParam<VarClaimed>;
2200
2201 impl ClaimVar for GhParam {
2202 type Claimed = ClaimedGhParam;
2203
2204 fn claim(self, ctx: &mut StepCtx<'_>) -> ClaimedGhParam {
2205 match self {
2206 GhParam::Static(s) => ClaimedGhParam::Static(s),
2207 GhParam::FloweyVar(var) => match &var.backing_var {
2208 ReadVarBacking::RuntimeVar { is_side_effect, .. } => {
2209 assert!(!is_side_effect);
2210 ClaimedGhParam::FloweyVar(var.claim(ctx))
2211 }
2212 ReadVarBacking::Inline(var) => ClaimedGhParam::Static(var.clone()),
2213 },
2214 }
2215 }
2216 }
2217
2218 #[derive(Debug, Clone, PartialEq, Eq, PartialOrd)]
2223 pub enum GhPermissionValue {
2224 None = 0,
2225 Read = 1,
2226 Write = 2,
2227 }
2228
2229 #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
2235 pub enum GhPermission {
2236 Actions,
2237 Attestations,
2238 Checks,
2239 Contents,
2240 Deployments,
2241 Discussions,
2242 IdToken,
2243 Issues,
2244 Packages,
2245 Pages,
2246 PullRequests,
2247 RepositoryProjects,
2248 SecurityEvents,
2249 Statuses,
2250 }
2251 }
2252
2253 pub mod rust {
2254 use crate::node::ClaimedWriteVar;
2255 use crate::node::FlowArch;
2256 use crate::node::FlowBackend;
2257 use crate::node::FlowPlatform;
2258 use crate::node::ReadVarValue;
2259 use crate::node::RuntimeVarDb;
2260 use crate::shell::FloweyShell;
2261 use serde::Serialize;
2262 use serde::de::DeserializeOwned;
2263
2264 pub fn new_rust_runtime_services(
2265 runtime_var_db: &mut dyn RuntimeVarDb,
2266 backend: FlowBackend,
2267 platform: FlowPlatform,
2268 arch: FlowArch,
2269 ) -> anyhow::Result<RustRuntimeServices<'_>> {
2270 Ok(RustRuntimeServices {
2271 runtime_var_db,
2272 backend,
2273 platform,
2274 arch,
2275 has_read_secret: false,
2276 sh: FloweyShell::new()?,
2277 })
2278 }
2279
2280 pub struct RustRuntimeServices<'a> {
2281 runtime_var_db: &'a mut dyn RuntimeVarDb,
2282 backend: FlowBackend,
2283 platform: FlowPlatform,
2284 arch: FlowArch,
2285 has_read_secret: bool,
2286 pub sh: FloweyShell,
2292 }
2293
2294 impl RustRuntimeServices<'_> {
2295 pub fn backend(&self) -> FlowBackend {
2298 self.backend
2299 }
2300
2301 pub fn platform(&self) -> FlowPlatform {
2304 self.platform
2305 }
2306
2307 pub fn arch(&self) -> FlowArch {
2309 self.arch
2310 }
2311
2312 pub fn write<T>(&mut self, var: ClaimedWriteVar<T>, val: &T)
2320 where
2321 T: Serialize + DeserializeOwned,
2322 {
2323 self.write_maybe_secret(var, val, self.has_read_secret)
2324 }
2325
2326 pub fn write_secret<T>(&mut self, var: ClaimedWriteVar<T>, val: &T)
2332 where
2333 T: Serialize + DeserializeOwned,
2334 {
2335 self.write_maybe_secret(var, val, true)
2336 }
2337
2338 pub fn write_not_secret<T>(&mut self, var: ClaimedWriteVar<T>, val: &T)
2345 where
2346 T: Serialize + DeserializeOwned,
2347 {
2348 self.write_maybe_secret(var, val, false)
2349 }
2350
2351 fn write_maybe_secret<T>(&mut self, var: ClaimedWriteVar<T>, val: &T, is_secret: bool)
2352 where
2353 T: Serialize + DeserializeOwned,
2354 {
2355 let val = if var.is_side_effect {
2356 b"null".to_vec()
2357 } else {
2358 serde_json::to_vec(val).expect("improve this error path")
2359 };
2360 self.runtime_var_db
2361 .set_var(&var.backing_var, is_secret, val);
2362 }
2363
2364 pub fn write_all<T>(
2365 &mut self,
2366 vars: impl IntoIterator<Item = ClaimedWriteVar<T>>,
2367 val: &T,
2368 ) where
2369 T: Serialize + DeserializeOwned,
2370 {
2371 for var in vars {
2372 self.write(var, val)
2373 }
2374 }
2375
2376 pub fn read<T: ReadVarValue>(&mut self, var: T) -> T::Value {
2377 var.read_value(self)
2378 }
2379
2380 pub(crate) fn get_var(&mut self, var: &str, is_side_effect: bool) -> Vec<u8> {
2381 let (v, is_secret) = self.runtime_var_db.get_var(var);
2382 self.has_read_secret |= is_secret && !is_side_effect;
2383 v
2384 }
2385
2386 pub fn dangerous_gh_set_global_env_var(
2393 &mut self,
2394 var: String,
2395 gh_env_var: String,
2396 ) -> anyhow::Result<()> {
2397 if !matches!(self.backend, FlowBackend::Github) {
2398 return Err(anyhow::anyhow!(
2399 "dangerous_set_gh_env_var can only be used on GitHub Actions"
2400 ));
2401 }
2402
2403 let gh_env_file_path = std::env::var("GITHUB_ENV")?;
2404 let mut gh_env_file = fs_err::OpenOptions::new()
2405 .append(true)
2406 .open(gh_env_file_path)?;
2407 let gh_env_var_assignment = format!(
2408 r#"{}<<EOF
2409{}
2410EOF
2411"#,
2412 gh_env_var, var
2413 );
2414 std::io::Write::write_all(&mut gh_env_file, gh_env_var_assignment.as_bytes())?;
2415
2416 Ok(())
2417 }
2418 }
2419 }
2420}
2421
2422pub trait FlowNodeBase {
2427 type Request: Serialize + DeserializeOwned;
2428
2429 fn imports(&mut self, ctx: &mut ImportCtx<'_>);
2430 fn emit(
2431 &mut self,
2432 config_bytes: Vec<Box<[u8]>>,
2433 requests: Vec<Self::Request>,
2434 ctx: &mut NodeCtx<'_>,
2435 ) -> anyhow::Result<()>;
2436
2437 fn i_know_what_im_doing_with_this_manual_impl(&mut self);
2443}
2444
2445pub mod erased {
2446 use crate::node::FlowNodeBase;
2447 use crate::node::NodeCtx;
2448 use crate::node::user_facing::*;
2449
2450 pub struct ErasedNode<N: FlowNodeBase>(pub N);
2451
2452 impl<N: FlowNodeBase> ErasedNode<N> {
2453 pub fn from_node(node: N) -> Self {
2454 Self(node)
2455 }
2456 }
2457
2458 impl<N> FlowNodeBase for ErasedNode<N>
2459 where
2460 N: FlowNodeBase,
2461 {
2462 type Request = Box<[u8]>;
2464
2465 fn imports(&mut self, ctx: &mut ImportCtx<'_>) {
2466 self.0.imports(ctx)
2467 }
2468
2469 fn emit(
2470 &mut self,
2471 config_bytes: Vec<Box<[u8]>>,
2472 requests: Vec<Box<[u8]>>,
2473 ctx: &mut NodeCtx<'_>,
2474 ) -> anyhow::Result<()> {
2475 let mut converted_requests = Vec::new();
2476 for req in requests {
2477 converted_requests.push(serde_json::from_slice(&req)?)
2478 }
2479
2480 self.0.emit(config_bytes, converted_requests, ctx)
2481 }
2482
2483 fn i_know_what_im_doing_with_this_manual_impl(&mut self) {}
2484 }
2485}
2486
2487#[derive(Clone, Copy, PartialEq, Eq, Hash)]
2489pub struct NodeHandle(std::any::TypeId);
2490
2491impl Ord for NodeHandle {
2492 fn cmp(&self, other: &Self) -> std::cmp::Ordering {
2493 self.modpath().cmp(other.modpath())
2494 }
2495}
2496
2497impl PartialOrd for NodeHandle {
2498 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
2499 Some(self.cmp(other))
2500 }
2501}
2502
2503impl std::fmt::Debug for NodeHandle {
2504 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2505 std::fmt::Debug::fmt(&self.try_modpath(), f)
2506 }
2507}
2508
2509impl NodeHandle {
2510 pub fn from_type<N: FlowNodeBase + 'static>() -> NodeHandle {
2511 NodeHandle(std::any::TypeId::of::<N>())
2512 }
2513
2514 pub fn from_modpath(modpath: &str) -> NodeHandle {
2515 node_luts::erased_node_by_modpath().get(modpath).unwrap().0
2516 }
2517
2518 pub fn try_from_modpath(modpath: &str) -> Option<NodeHandle> {
2519 node_luts::erased_node_by_modpath()
2520 .get(modpath)
2521 .map(|(s, _)| *s)
2522 }
2523
2524 pub fn new_erased_node(&self) -> Box<dyn FlowNodeBase<Request = Box<[u8]>>> {
2525 let ctor = node_luts::erased_node_by_typeid().get(self).unwrap();
2526 ctor()
2527 }
2528
2529 pub fn modpath(&self) -> &'static str {
2530 node_luts::modpath_by_node_typeid().get(self).unwrap()
2531 }
2532
2533 pub fn try_modpath(&self) -> Option<&'static str> {
2534 node_luts::modpath_by_node_typeid().get(self).cloned()
2535 }
2536
2537 pub fn dummy() -> NodeHandle {
2540 NodeHandle(std::any::TypeId::of::<()>())
2541 }
2542}
2543
2544pub fn list_all_registered_nodes() -> impl Iterator<Item = NodeHandle> {
2545 node_luts::modpath_by_node_typeid().keys().cloned()
2546}
2547
2548mod node_luts {
2563 use super::FlowNodeBase;
2564 use super::NodeHandle;
2565 use std::collections::HashMap;
2566 use std::sync::OnceLock;
2567
2568 pub(super) fn modpath_by_node_typeid() -> &'static HashMap<NodeHandle, &'static str> {
2569 static TYPEID_TO_MODPATH: OnceLock<HashMap<NodeHandle, &'static str>> = OnceLock::new();
2570
2571 TYPEID_TO_MODPATH.get_or_init(|| {
2572 let mut lookup = HashMap::new();
2573 for crate::node::private::FlowNodeMeta {
2574 module_path,
2575 ctor: _,
2576 typeid,
2577 } in crate::node::private::FLOW_NODES
2578 {
2579 let existing = lookup.insert(
2580 NodeHandle(*typeid),
2581 module_path
2582 .strip_suffix("::_only_one_call_to_flowey_node_per_module")
2583 .unwrap(),
2584 );
2585 assert!(existing.is_none())
2588 }
2589
2590 lookup
2591 })
2592 }
2593
2594 pub(super) fn erased_node_by_typeid()
2595 -> &'static HashMap<NodeHandle, fn() -> Box<dyn FlowNodeBase<Request = Box<[u8]>>>> {
2596 static LOOKUP: OnceLock<
2597 HashMap<NodeHandle, fn() -> Box<dyn FlowNodeBase<Request = Box<[u8]>>>>,
2598 > = OnceLock::new();
2599
2600 LOOKUP.get_or_init(|| {
2601 let mut lookup = HashMap::new();
2602 for crate::node::private::FlowNodeMeta {
2603 module_path: _,
2604 ctor,
2605 typeid,
2606 } in crate::node::private::FLOW_NODES
2607 {
2608 let existing = lookup.insert(NodeHandle(*typeid), *ctor);
2609 assert!(existing.is_none())
2612 }
2613
2614 lookup
2615 })
2616 }
2617
2618 pub(super) fn erased_node_by_modpath() -> &'static HashMap<
2619 &'static str,
2620 (
2621 NodeHandle,
2622 fn() -> Box<dyn FlowNodeBase<Request = Box<[u8]>>>,
2623 ),
2624 > {
2625 static MODPATH_LOOKUP: OnceLock<
2626 HashMap<
2627 &'static str,
2628 (
2629 NodeHandle,
2630 fn() -> Box<dyn FlowNodeBase<Request = Box<[u8]>>>,
2631 ),
2632 >,
2633 > = OnceLock::new();
2634
2635 MODPATH_LOOKUP.get_or_init(|| {
2636 let mut lookup = HashMap::new();
2637 for crate::node::private::FlowNodeMeta { module_path, ctor, typeid } in crate::node::private::FLOW_NODES {
2638 let existing = lookup.insert(module_path.strip_suffix("::_only_one_call_to_flowey_node_per_module").unwrap(), (NodeHandle(*typeid), *ctor));
2639 if existing.is_some() {
2640 panic!("conflicting node registrations at {module_path}! please ensure there is a single node per module!")
2641 }
2642 }
2643 lookup
2644 })
2645 }
2646}
2647
2648#[doc(hidden)]
2649pub mod private {
2650 pub use linkme;
2651
2652 pub struct FlowNodeMeta {
2653 pub module_path: &'static str,
2654 pub ctor: fn() -> Box<dyn super::FlowNodeBase<Request = Box<[u8]>>>,
2655 pub typeid: std::any::TypeId,
2656 }
2657
2658 #[linkme::distributed_slice]
2659 pub static FLOW_NODES: [FlowNodeMeta] = [..];
2660
2661 #[expect(unsafe_code)]
2663 #[linkme::distributed_slice(FLOW_NODES)]
2664 static DUMMY_FLOW_NODE: FlowNodeMeta = FlowNodeMeta {
2665 module_path: "<dummy>::_only_one_call_to_flowey_node_per_module",
2666 ctor: || unreachable!(),
2667 typeid: std::any::TypeId::of::<()>(),
2668 };
2669}
2670
2671#[doc(hidden)]
2672#[macro_export]
2673macro_rules! new_flow_node_base {
2674 (struct Node) => {
2675 #[non_exhaustive]
2677 pub struct Node;
2678
2679 mod _only_one_call_to_flowey_node_per_module {
2680 const _: () = {
2681 use $crate::node::private::linkme;
2682
2683 fn new_erased() -> Box<dyn $crate::node::FlowNodeBase<Request = Box<[u8]>>> {
2684 Box::new($crate::node::erased::ErasedNode(super::Node))
2685 }
2686
2687 #[linkme::distributed_slice($crate::node::private::FLOW_NODES)]
2688 #[linkme(crate = linkme)]
2689 static FLOW_NODE: $crate::node::private::FlowNodeMeta =
2690 $crate::node::private::FlowNodeMeta {
2691 module_path: module_path!(),
2692 ctor: new_erased,
2693 typeid: std::any::TypeId::of::<super::Node>(),
2694 };
2695 };
2696 }
2697 };
2698}
2699
2700pub trait FlowNode {
2787 type Request: Serialize + DeserializeOwned;
2791
2792 fn imports(ctx: &mut ImportCtx<'_>);
2804
2805 fn emit(requests: Vec<Self::Request>, ctx: &mut NodeCtx<'_>) -> anyhow::Result<()>;
2808}
2809
2810#[macro_export]
2811macro_rules! new_flow_node {
2812 (struct Node) => {
2813 $crate::new_flow_node_base!(struct Node);
2814
2815 impl $crate::node::FlowNodeBase for Node
2816 where
2817 Node: FlowNode,
2818 {
2819 type Request = <Node as FlowNode>::Request;
2820
2821 fn imports(&mut self, dep: &mut $crate::node::ImportCtx<'_>) {
2822 <Node as FlowNode>::imports(dep)
2823 }
2824
2825 fn emit(
2826 &mut self,
2827 _config_bytes: Vec<Box<[u8]>>,
2828 requests: Vec<Self::Request>,
2829 ctx: &mut $crate::node::NodeCtx<'_>,
2830 ) -> anyhow::Result<()> {
2831 <Node as FlowNode>::emit(requests, ctx)
2832 }
2833
2834 fn i_know_what_im_doing_with_this_manual_impl(&mut self) {}
2835 }
2836 };
2837}
2838
2839pub trait SimpleFlowNode {
2860 type Request: Serialize + DeserializeOwned;
2861
2862 fn imports(ctx: &mut ImportCtx<'_>);
2874
2875 fn process_request(request: Self::Request, ctx: &mut NodeCtx<'_>) -> anyhow::Result<()>;
2877}
2878
2879#[macro_export]
2880macro_rules! new_simple_flow_node {
2881 (struct Node) => {
2882 $crate::new_flow_node_base!(struct Node);
2883
2884 impl $crate::node::FlowNodeBase for Node
2885 where
2886 Node: $crate::node::SimpleFlowNode,
2887 {
2888 type Request = <Node as $crate::node::SimpleFlowNode>::Request;
2889
2890 fn imports(&mut self, dep: &mut $crate::node::ImportCtx<'_>) {
2891 <Node as $crate::node::SimpleFlowNode>::imports(dep)
2892 }
2893
2894 fn emit(
2895 &mut self,
2896 _config_bytes: Vec<Box<[u8]>>,
2897 requests: Vec<Self::Request>,
2898 ctx: &mut $crate::node::NodeCtx<'_>,
2899 ) -> anyhow::Result<()> {
2900 for req in requests {
2901 <Node as $crate::node::SimpleFlowNode>::process_request(req, ctx)?
2902 }
2903
2904 Ok(())
2905 }
2906
2907 fn i_know_what_im_doing_with_this_manual_impl(&mut self) {}
2908 }
2909 };
2910}
2911
2912pub trait FlowNodeWithConfig {
2959 type Request: Serialize + DeserializeOwned;
2961
2962 type Config: ConfigMerge;
2969
2970 fn imports(ctx: &mut ImportCtx<'_>);
2972
2973 fn emit(
2975 config: Self::Config,
2976 requests: Vec<Self::Request>,
2977 ctx: &mut NodeCtx<'_>,
2978 ) -> anyhow::Result<()>;
2979}
2980
2981#[macro_export]
2982macro_rules! new_flow_node_with_config {
2983 (struct Node) => {
2984 $crate::new_flow_node_base!(struct Node);
2985
2986 impl $crate::node::FlowNodeBase for Node
2987 where
2988 Node: $crate::node::FlowNodeWithConfig,
2989 {
2990 type Request = <Node as $crate::node::FlowNodeWithConfig>::Request;
2991
2992 fn imports(&mut self, dep: &mut $crate::node::ImportCtx<'_>) {
2993 <Node as $crate::node::FlowNodeWithConfig>::imports(dep)
2994 }
2995
2996 fn emit(
2997 &mut self,
2998 config_bytes: Vec<Box<[u8]>>,
2999 requests: Vec<Self::Request>,
3000 ctx: &mut $crate::node::NodeCtx<'_>,
3001 ) -> anyhow::Result<()> {
3002 use $crate::node::ConfigMerge;
3003
3004 type C = <Node as $crate::node::FlowNodeWithConfig>::Config;
3005
3006 let mut merged = <C as Default>::default();
3007 for bytes in config_bytes {
3008 let partial: C = serde_json::from_slice(&bytes)?;
3009 merged.merge(partial)?;
3010 }
3011
3012 <Node as $crate::node::FlowNodeWithConfig>::emit(merged, requests, ctx)
3013 }
3014
3015 fn i_know_what_im_doing_with_this_manual_impl(&mut self) {}
3016 }
3017 };
3018}
3019
3020pub trait IntoRequest {
3028 type Node: FlowNodeBase;
3029 fn into_request(self) -> <Self::Node as FlowNodeBase>::Request;
3030
3031 #[doc(hidden)]
3034 #[expect(nonstandard_style)]
3035 fn do_not_manually_impl_this_trait__use_the_flowey_request_macro_instead(&mut self);
3036}
3037
3038pub trait IntoConfig: Serialize {
3044 type Node: FlowNodeBase;
3045
3046 #[doc(hidden)]
3049 #[expect(nonstandard_style)]
3050 fn do_not_manually_impl_this_trait__use_the_flowey_config_macro_instead(&mut self);
3051}
3052
3053pub trait ConfigMerge: Serialize + DeserializeOwned + Default {
3056 fn merge(&mut self, other: Self) -> anyhow::Result<()>;
3059}
3060
3061pub trait ConfigField {
3068 fn merge_field(&mut self, field_name: &str, other: Self) -> anyhow::Result<()>;
3069}
3070
3071impl<T: PartialEq> ConfigField for Option<T> {
3072 fn merge_field(&mut self, field_name: &str, other: Self) -> anyhow::Result<()> {
3073 if let Some(new) = other {
3074 match self {
3075 None => *self = Some(new),
3076 Some(old) if *old == new => {}
3077 Some(_) => {
3078 anyhow::bail!("config field `{field_name}` mismatch");
3079 }
3080 }
3081 }
3082 Ok(())
3083 }
3084}
3085
3086impl<K: Ord + std::fmt::Debug, V: PartialEq> ConfigField for BTreeMap<K, V> {
3087 fn merge_field(&mut self, field_name: &str, other: Self) -> anyhow::Result<()> {
3088 for (k, v) in other {
3089 use std::collections::btree_map::Entry;
3090 match self.entry(k) {
3091 Entry::Vacant(e) => {
3092 e.insert(v);
3093 }
3094 Entry::Occupied(e) if *e.get() == v => {}
3095 Entry::Occupied(e) => {
3096 anyhow::bail!("config field `{field_name}` mismatch for key {:?}", e.key(),);
3097 }
3098 }
3099 }
3100 Ok(())
3101 }
3102}
3103
3104#[doc(hidden)]
3105#[macro_export]
3106macro_rules! __flowey_request_inner {
3107 (@emit_struct [$req:ident]
3111 $(#[$a:meta])*
3112 $variant:ident($($tt:tt)*),
3113 $($rest:tt)*
3114 ) => {
3115 $(#[$a])*
3116 #[derive($crate::reexports::Serialize, $crate::reexports::Deserialize)]
3117 pub struct $variant($($tt)*);
3118
3119 impl IntoRequest for $variant {
3120 type Node = Node;
3121 fn into_request(self) -> $req {
3122 $req::$variant(self)
3123 }
3124 fn do_not_manually_impl_this_trait__use_the_flowey_request_macro_instead(&mut self) {}
3125 }
3126
3127 $crate::__flowey_request_inner!(@emit_struct [$req] $($rest)*);
3128 };
3129 (@emit_struct [$req:ident]
3130 $(#[$a:meta])*
3131 $variant:ident { $($tt:tt)* },
3132 $($rest:tt)*
3133 ) => {
3134 $(#[$a])*
3135 #[derive($crate::reexports::Serialize, $crate::reexports::Deserialize)]
3136 pub struct $variant {
3137 $($tt)*
3138 }
3139
3140 impl IntoRequest for $variant {
3141 type Node = Node;
3142 fn into_request(self) -> $req {
3143 $req::$variant(self)
3144 }
3145 fn do_not_manually_impl_this_trait__use_the_flowey_request_macro_instead(&mut self) {}
3146 }
3147
3148 $crate::__flowey_request_inner!(@emit_struct [$req] $($rest)*);
3149 };
3150 (@emit_struct [$req:ident]
3151 $(#[$a:meta])*
3152 $variant:ident,
3153 $($rest:tt)*
3154 ) => {
3155 $(#[$a])*
3156 #[derive(Serialize, Deserialize)]
3157 pub struct $variant;
3158
3159 impl IntoRequest for $variant {
3160 type Node = Node;
3161 fn into_request(self) -> $req {
3162 $req::$variant(self)
3163 }
3164 fn do_not_manually_impl_this_trait__use_the_flowey_request_macro_instead(&mut self) {}
3165 }
3166
3167 $crate::__flowey_request_inner!(@emit_struct [$req] $($rest)*);
3168 };
3169 (@emit_struct [$req:ident]
3170 ) => {};
3171
3172 (@emit_req_enum [$req:ident($($root_a:meta,)*), $($prev:ident[$($prev_a:meta,)*])*]
3176 $(#[$a:meta])*
3177 $variant:ident($($tt:tt)*),
3178 $($rest:tt)*
3179 ) => {
3180 $crate::__flowey_request_inner!(@emit_req_enum [$req($($root_a,)*), $($prev[$($prev_a,)*])* $variant[$($a,)*]] $($rest)*);
3181 };
3182 (@emit_req_enum [$req:ident($($root_a:meta,)*), $($prev:ident[$($prev_a:meta,)*])*]
3183 $(#[$a:meta])*
3184 $variant:ident { $($tt:tt)* },
3185 $($rest:tt)*
3186 ) => {
3187 $crate::__flowey_request_inner!(@emit_req_enum [$req($($root_a,)*), $($prev[$($prev_a,)*])* $variant[$($a,)*]] $($rest)*);
3188 };
3189 (@emit_req_enum [$req:ident($($root_a:meta,)*), $($prev:ident[$($prev_a:meta,)*])*]
3190 $(#[$a:meta])*
3191 $variant:ident,
3192 $($rest:tt)*
3193 ) => {
3194 $crate::__flowey_request_inner!(@emit_req_enum [$req($($root_a,)*), $($prev[$($prev_a,)*])* $variant[$($a,)*]] $($rest)*);
3195 };
3196 (@emit_req_enum [$req:ident($($root_a:meta,)*), $($prev:ident[$($prev_a:meta,)*])*]
3197 ) => {
3198 #[derive(Serialize, Deserialize)]
3199 pub enum $req {$(
3200 $(#[$prev_a])*
3201 $prev(self::req::$prev),
3202 )*}
3203
3204 impl IntoRequest for $req {
3205 type Node = Node;
3206 fn into_request(self) -> $req {
3207 self
3208 }
3209 fn do_not_manually_impl_this_trait__use_the_flowey_request_macro_instead(&mut self) {}
3210 }
3211 };
3212}
3213
3214#[macro_export]
3262macro_rules! flowey_request {
3263 (
3264 $(#[$root_a:meta])*
3265 pub enum_struct $req:ident {
3266 $($tt:tt)*
3267 }
3268 ) => {
3269 $crate::__flowey_request_inner!(@emit_req_enum [$req($($root_a,)*),] $($tt)*);
3270 pub mod req {
3271 use super::*;
3272 $crate::__flowey_request_inner!(@emit_struct [$req] $($tt)*);
3273 }
3274 };
3275
3276 (
3277 $(#[$a:meta])*
3278 pub enum $req:ident {
3279 $($tt:tt)*
3280 }
3281 ) => {
3282 $(#[$a])*
3283 #[derive($crate::reexports::Serialize, $crate::reexports::Deserialize)]
3284 pub enum $req {
3285 $($tt)*
3286 }
3287
3288 impl $crate::node::IntoRequest for $req {
3289 type Node = Node;
3290 fn into_request(self) -> $req {
3291 self
3292 }
3293 fn do_not_manually_impl_this_trait__use_the_flowey_request_macro_instead(&mut self) {}
3294 }
3295 };
3296
3297 (
3298 $(#[$a:meta])*
3299 pub struct $req:ident {
3300 $($tt:tt)*
3301 }
3302 ) => {
3303 $(#[$a])*
3304 #[derive($crate::reexports::Serialize, $crate::reexports::Deserialize)]
3305 pub struct $req {
3306 $($tt)*
3307 }
3308
3309 impl $crate::node::IntoRequest for $req {
3310 type Node = Node;
3311 fn into_request(self) -> $req {
3312 self
3313 }
3314 fn do_not_manually_impl_this_trait__use_the_flowey_request_macro_instead(&mut self) {}
3315 }
3316 };
3317
3318 (
3319 $(#[$a:meta])*
3320 pub struct $req:ident($($tt:tt)*);
3321 ) => {
3322 $(#[$a])*
3323 #[derive($crate::reexports::Serialize, $crate::reexports::Deserialize)]
3324 pub struct $req($($tt)*);
3325
3326 impl $crate::node::IntoRequest for $req {
3327 type Node = Node;
3328 fn into_request(self) -> $req {
3329 self
3330 }
3331 fn do_not_manually_impl_this_trait__use_the_flowey_request_macro_instead(&mut self) {}
3332 }
3333 };
3334}
3335
3336#[macro_export]
3374macro_rules! flowey_config {
3375 (
3376 $(#[$meta:meta])*
3377 pub struct $Config:ident {
3378 $(
3379 $(#[$field_meta:meta])*
3380 pub $field:ident : $ty:ty
3381 ),* $(,)?
3382 }
3383 ) => {
3384 $(#[$meta])*
3385 #[derive(
3386 $crate::reexports::Serialize,
3387 $crate::reexports::Deserialize,
3388 Default,
3389 )]
3390 pub struct $Config {
3391 $(
3392 $(#[$field_meta])*
3393 pub $field: $ty,
3394 )*
3395 }
3396
3397 impl $crate::node::ConfigMerge for $Config {
3398 fn merge(&mut self, other: Self) -> anyhow::Result<()> {
3399 $(
3400 $crate::node::ConfigField::merge_field(
3401 &mut self.$field,
3402 stringify!($field),
3403 other.$field,
3404 )?;
3405 )*
3406 Ok(())
3407 }
3408 }
3409
3410 impl $crate::node::IntoConfig for $Config {
3411 type Node = Node;
3412
3413 fn do_not_manually_impl_this_trait__use_the_flowey_config_macro_instead(&mut self) {}
3414 }
3415 };
3416}
3417
3418#[macro_export]
3435macro_rules! shell_cmd {
3436 ($rt:expr, $cmd:literal) => {{
3437 let flowey_sh = &$rt.sh;
3438 #[expect(clippy::disallowed_macros)]
3439 flowey_sh.wrap($crate::reexports::xshell::cmd!(flowey_sh.xshell(), $cmd))
3440 }};
3441}