reifydb_sub_flow/
builder.rs1use std::{collections::HashMap, path::PathBuf, sync::Arc};
5
6#[cfg(all(reifydb_target = "host", not(reifydb_dst)))]
7use reifydb_core::interface::flow::to_bitmask;
8use reifydb_core::{event::operator::OperatorColumn, interface::catalog::flow::OperatorId};
9use reifydb_flow::operator::BoxedHostOperator;
10#[cfg(all(reifydb_target = "host", not(reifydb_dst)))]
11use reifydb_sdk::flow::operator::GuestOperator;
12#[cfg(all(reifydb_target = "host", not(reifydb_dst)))]
13use reifydb_sdk::flow::operator::{OperatorMetadata, column::operator::OperatorColumn as SdkOperatorColumn};
14use reifydb_value::{Result, config::Config};
15
16#[cfg(all(reifydb_target = "host", not(reifydb_dst)))]
17use crate::operator::mount::mount;
18
19pub(crate) type OperatorFactory = Arc<dyn Fn(OperatorId, &Config) -> Result<BoxedHostOperator> + Send + Sync>;
20
21#[derive(Clone)]
22pub struct CustomOperatorEntry {
23 pub factory: OperatorFactory,
24 pub abi: Option<u32>,
25 pub version: String,
26 pub description: String,
27 pub capabilities: u32,
28 pub input: Vec<OperatorColumn>,
29 pub output: Vec<OperatorColumn>,
30}
31
32#[derive(Clone, Default)]
33pub struct CustomOperators {
34 inner: Arc<HashMap<String, CustomOperatorEntry>>,
35}
36
37impl CustomOperators {
38 pub(crate) fn new(map: HashMap<String, CustomOperatorEntry>) -> Self {
39 Self {
40 inner: Arc::new(map),
41 }
42 }
43
44 pub(crate) fn get(&self, name: &str) -> Option<&OperatorFactory> {
45 self.inner.get(name).map(|entry| &entry.factory)
46 }
47
48 pub(crate) fn iter(&self) -> impl Iterator<Item = (&String, &CustomOperatorEntry)> {
49 self.inner.iter()
50 }
51}
52
53#[cfg(all(reifydb_target = "host", not(reifydb_dst)))]
54fn describe_columns(columns: &[SdkOperatorColumn]) -> Vec<OperatorColumn> {
55 columns.iter()
56 .map(|column| OperatorColumn {
57 name: column.name.to_string(),
58 field_type: column.type_constraint.clone(),
59 description: column.description.to_string(),
60 })
61 .collect()
62}
63
64pub struct FlowConfigurator {
65 operators_dir: Option<PathBuf>,
66 custom_operators: HashMap<String, CustomOperatorEntry>,
67}
68
69impl Default for FlowConfigurator {
70 fn default() -> Self {
71 Self::new()
72 }
73}
74
75impl FlowConfigurator {
76 pub fn new() -> Self {
77 Self {
78 operators_dir: None,
79 custom_operators: HashMap::new(),
80 }
81 }
82
83 pub fn operators_dir(mut self, path: PathBuf) -> Self {
84 self.operators_dir = Some(path);
85 self
86 }
87
88 #[cfg(all(reifydb_target = "host", not(reifydb_dst)))]
89 pub fn register_operator<O>(mut self) -> Self
90 where
91 O: GuestOperator + OperatorMetadata + 'static,
92 {
93 self.custom_operators.insert(
94 O::NAME.to_string(),
95 CustomOperatorEntry {
96 factory: Arc::new(|operator, config| {
97 let logic = O::create(operator, config)?;
98 Ok(mount(logic, operator, O::CAPABILITIES))
99 }),
100 abi: None,
101 version: <O as OperatorMetadata>::VERSION.to_string(),
102 description: <O as OperatorMetadata>::DESCRIPTION.to_string(),
103 capabilities: to_bitmask(<O as OperatorMetadata>::CAPABILITIES),
104 input: describe_columns(<O as OperatorMetadata>::INPUT_COLUMNS),
105 output: describe_columns(<O as OperatorMetadata>::OUTPUT_COLUMNS),
106 },
107 );
108 self
109 }
110
111 pub(crate) fn configure(self) -> FlowConfig {
112 FlowConfig {
113 operators_dir: self.operators_dir,
114 custom_operators: CustomOperators::new(self.custom_operators),
115 }
116 }
117}
118
119pub struct FlowConfig {
120 pub operators_dir: Option<PathBuf>,
121
122 pub custom_operators: CustomOperators,
123}