feldera_adapterlib/
preprocess.rs1use crate::ConnectorMetadata;
13use crate::format::{ParseError, Splitter};
14use feldera_types::preprocess::PreprocessorConfig;
15use std::collections::BTreeMap;
16use std::fmt::{Display, Formatter, Result as FmtResult};
17use std::sync::Arc;
18
19#[derive(Debug, Clone, PartialEq, Eq)]
21pub enum PreprocessorCreateError {
22 ConfigurationError(String),
24 FactoryNotFound(String),
26}
27
28impl Display for PreprocessorCreateError {
29 fn fmt(&self, f: &mut Formatter<'_>) -> FmtResult {
30 match self {
31 PreprocessorCreateError::ConfigurationError(msg) => {
32 write!(f, "Configuration error: {}", msg)
33 }
34 PreprocessorCreateError::FactoryNotFound(msg) => {
35 write!(
36 f,
37 "Could not locate factory generating preprocessor: {}",
38 msg
39 )
40 }
41 }
42 }
43}
44
45impl std::error::Error for PreprocessorCreateError {}
46
47pub trait Preprocessor: Send + Sync {
49 fn process(&mut self, data: &[u8]) -> (Vec<u8>, Vec<ParseError>) {
60 self.process_with_metadata(data, None)
61 }
62
63 fn process_with_metadata(
78 &mut self,
79 data: &[u8],
80 _metadata: Option<&ConnectorMetadata>,
81 ) -> (Vec<u8>, Vec<ParseError>) {
82 self.process(data)
83 }
84
85 fn fork(&self) -> Box<dyn Preprocessor>;
90
91 fn splitter(&self) -> Option<Box<dyn Splitter>>;
95}
96
97pub trait PreprocessorFactory: Send + Sync {
99 fn create(
105 &self,
106 config: &PreprocessorConfig,
107 ) -> Result<Box<dyn Preprocessor>, PreprocessorCreateError>;
108}
109
110#[derive(Default)]
112pub struct PreprocessorRegistry {
113 registered: BTreeMap<&'static str, Arc<dyn PreprocessorFactory>>,
114}
115
116impl PreprocessorRegistry {
117 pub fn new() -> Self {
118 Self {
119 registered: BTreeMap::new(),
120 }
121 }
122
123 pub fn register(&mut self, name: &'static str, factory: Box<dyn PreprocessorFactory>) {
125 self.registered.insert(name, Arc::from(factory));
126 }
127
128 pub fn get(&self, name: &str) -> Option<Arc<dyn PreprocessorFactory>> {
129 self.registered.get(name).cloned()
130 }
131}
132
133#[cfg(test)]
134mod tests {
135 use super::*;
136 use feldera_sqllib::{SqlString, Variant};
137 use serde_json::json;
138
139 struct TopicPrefixPreprocessor;
144
145 impl Preprocessor for TopicPrefixPreprocessor {
146 fn process_with_metadata(
147 &mut self,
148 data: &[u8],
149 metadata: Option<&ConnectorMetadata>,
150 ) -> (Vec<u8>, Vec<ParseError>) {
151 let Some(metadata) = metadata else {
152 return (data.to_vec(), vec![]);
154 };
155 let mut output = Vec::new();
156 if let Some(Variant::String(topic)) = metadata.get_by_name("topic") {
157 output.extend_from_slice(topic.str().as_bytes());
158 output.push(b':');
159 }
160 output.extend_from_slice(data);
161 (output, vec![])
162 }
163
164 fn fork(&self) -> Box<dyn Preprocessor> {
165 Box::new(TopicPrefixPreprocessor)
166 }
167
168 fn splitter(&self) -> Option<Box<dyn Splitter>> {
169 None
170 }
171 }
172
173 struct UppercasePreprocessor;
175
176 impl Preprocessor for UppercasePreprocessor {
177 fn process(&mut self, data: &[u8]) -> (Vec<u8>, Vec<ParseError>) {
178 (data.to_ascii_uppercase(), vec![])
179 }
180
181 fn fork(&self) -> Box<dyn Preprocessor> {
182 Box::new(UppercasePreprocessor)
183 }
184
185 fn splitter(&self) -> Option<Box<dyn Splitter>> {
186 None
187 }
188 }
189
190 struct UppercasePreprocessorFactory;
191
192 impl PreprocessorFactory for UppercasePreprocessorFactory {
193 fn create(
194 &self,
195 _config: &PreprocessorConfig,
196 ) -> Result<Box<dyn Preprocessor>, PreprocessorCreateError> {
197 Ok(Box::new(UppercasePreprocessor))
198 }
199 }
200
201 fn make_config(name: &str) -> PreprocessorConfig {
202 PreprocessorConfig {
203 name: name.to_string(),
204 message_oriented: false,
205 config: json!({}),
206 }
207 }
208
209 fn make_metadata() -> ConnectorMetadata {
210 let mut metadata = ConnectorMetadata::new();
211 metadata.insert("topic", Variant::String(SqlString::from("events")));
212 metadata
213 }
214
215 #[test]
216 fn test_process_metadata_transforms_using_metadata() {
217 let mut preprocessor = TopicPrefixPreprocessor;
218
219 let metadata = make_metadata();
221 let (output, errors) = preprocessor.process_with_metadata(b"payload", Some(&metadata));
222 assert!(errors.is_empty(), "unexpected errors: {errors:?}");
223 assert_eq!(output, b"events:payload");
224
225 let (output, errors) = preprocessor.process_with_metadata(b"payload", None);
227 assert!(errors.is_empty(), "unexpected errors: {errors:?}");
228 assert_eq!(output, b"payload");
229 }
230
231 #[test]
236 fn test_process_defaults_to_process_with_metadata() {
237 let mut preprocessor = TopicPrefixPreprocessor;
238
239 let (output, errors) = preprocessor.process(b"payload");
240 assert!(errors.is_empty(), "unexpected errors: {errors:?}");
241 assert_eq!(output, b"payload");
242 }
243
244 #[test]
245 fn test_registry_register_and_get() {
246 let mut registry = PreprocessorRegistry::new();
247 registry.register("upper", Box::new(UppercasePreprocessorFactory));
248
249 let factory = registry.get("upper").expect("factory must be registered");
250 let mut preprocessor = factory.create(&make_config("upper")).unwrap();
251 let (output, errors) = preprocessor.process(b"xyz");
252 assert!(errors.is_empty(), "unexpected errors: {errors:?}");
253 assert_eq!(output, b"XYZ");
254
255 assert!(registry.get("missing").is_none());
256 }
257}