Skip to main content

reifydb_sub_flow/
builder.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}