use crate::ConnectorMetadata;
use crate::format::{ParseError, Splitter};
use feldera_types::preprocess::PreprocessorConfig;
use std::collections::BTreeMap;
use std::fmt::{Display, Formatter, Result as FmtResult};
use std::sync::Arc;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PreprocessorCreateError {
ConfigurationError(String),
FactoryNotFound(String),
}
impl Display for PreprocessorCreateError {
fn fmt(&self, f: &mut Formatter<'_>) -> FmtResult {
match self {
PreprocessorCreateError::ConfigurationError(msg) => {
write!(f, "Configuration error: {}", msg)
}
PreprocessorCreateError::FactoryNotFound(msg) => {
write!(
f,
"Could not locate factory generating preprocessor: {}",
msg
)
}
}
}
}
impl std::error::Error for PreprocessorCreateError {}
pub trait Preprocessor: Send + Sync {
fn process(&mut self, data: &[u8]) -> (Vec<u8>, Vec<ParseError>) {
self.process_with_metadata(data, None)
}
fn process_with_metadata(
&mut self,
data: &[u8],
_metadata: Option<&ConnectorMetadata>,
) -> (Vec<u8>, Vec<ParseError>) {
self.process(data)
}
fn fork(&self) -> Box<dyn Preprocessor>;
fn splitter(&self) -> Option<Box<dyn Splitter>>;
}
pub trait PreprocessorFactory: Send + Sync {
fn create(
&self,
config: &PreprocessorConfig,
) -> Result<Box<dyn Preprocessor>, PreprocessorCreateError>;
}
#[derive(Default)]
pub struct PreprocessorRegistry {
registered: BTreeMap<&'static str, Arc<dyn PreprocessorFactory>>,
}
impl PreprocessorRegistry {
pub fn new() -> Self {
Self {
registered: BTreeMap::new(),
}
}
pub fn register(&mut self, name: &'static str, factory: Box<dyn PreprocessorFactory>) {
self.registered.insert(name, Arc::from(factory));
}
pub fn get(&self, name: &str) -> Option<Arc<dyn PreprocessorFactory>> {
self.registered.get(name).cloned()
}
}
#[cfg(test)]
mod tests {
use super::*;
use feldera_sqllib::{SqlString, Variant};
use serde_json::json;
struct TopicPrefixPreprocessor;
impl Preprocessor for TopicPrefixPreprocessor {
fn process_with_metadata(
&mut self,
data: &[u8],
metadata: Option<&ConnectorMetadata>,
) -> (Vec<u8>, Vec<ParseError>) {
let Some(metadata) = metadata else {
return (data.to_vec(), vec![]);
};
let mut output = Vec::new();
if let Some(Variant::String(topic)) = metadata.get_by_name("topic") {
output.extend_from_slice(topic.str().as_bytes());
output.push(b':');
}
output.extend_from_slice(data);
(output, vec![])
}
fn fork(&self) -> Box<dyn Preprocessor> {
Box::new(TopicPrefixPreprocessor)
}
fn splitter(&self) -> Option<Box<dyn Splitter>> {
None
}
}
struct UppercasePreprocessor;
impl Preprocessor for UppercasePreprocessor {
fn process(&mut self, data: &[u8]) -> (Vec<u8>, Vec<ParseError>) {
(data.to_ascii_uppercase(), vec![])
}
fn fork(&self) -> Box<dyn Preprocessor> {
Box::new(UppercasePreprocessor)
}
fn splitter(&self) -> Option<Box<dyn Splitter>> {
None
}
}
struct UppercasePreprocessorFactory;
impl PreprocessorFactory for UppercasePreprocessorFactory {
fn create(
&self,
_config: &PreprocessorConfig,
) -> Result<Box<dyn Preprocessor>, PreprocessorCreateError> {
Ok(Box::new(UppercasePreprocessor))
}
}
fn make_config(name: &str) -> PreprocessorConfig {
PreprocessorConfig {
name: name.to_string(),
message_oriented: false,
config: json!({}),
}
}
fn make_metadata() -> ConnectorMetadata {
let mut metadata = ConnectorMetadata::new();
metadata.insert("topic", Variant::String(SqlString::from("events")));
metadata
}
#[test]
fn test_process_metadata_transforms_using_metadata() {
let mut preprocessor = TopicPrefixPreprocessor;
let metadata = make_metadata();
let (output, errors) = preprocessor.process_with_metadata(b"payload", Some(&metadata));
assert!(errors.is_empty(), "unexpected errors: {errors:?}");
assert_eq!(output, b"events:payload");
let (output, errors) = preprocessor.process_with_metadata(b"payload", None);
assert!(errors.is_empty(), "unexpected errors: {errors:?}");
assert_eq!(output, b"payload");
}
#[test]
fn test_process_defaults_to_process_with_metadata() {
let mut preprocessor = TopicPrefixPreprocessor;
let (output, errors) = preprocessor.process(b"payload");
assert!(errors.is_empty(), "unexpected errors: {errors:?}");
assert_eq!(output, b"payload");
}
#[test]
fn test_registry_register_and_get() {
let mut registry = PreprocessorRegistry::new();
registry.register("upper", Box::new(UppercasePreprocessorFactory));
let factory = registry.get("upper").expect("factory must be registered");
let mut preprocessor = factory.create(&make_config("upper")).unwrap();
let (output, errors) = preprocessor.process(b"xyz");
assert!(errors.is_empty(), "unexpected errors: {errors:?}");
assert_eq!(output, b"XYZ");
assert!(registry.get("missing").is_none());
}
}