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