use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::HashMap;
#[derive(Debug, Clone, Deserialize, Serialize, Default)]
pub struct PipelineConfig {
#[serde(skip_serializing_if = "Option::is_none")]
pub filter: Option<FilterConfig>,
#[serde(skip_serializing_if = "Option::is_none")]
pub transform: Option<TransformConfig>,
#[serde(skip_serializing_if = "Option::is_none")]
pub mapping: Option<MappingConfig>,
}
#[derive(Debug, Clone, Deserialize, Serialize, Default)]
pub struct FilterConfig {
#[serde(default)]
pub tables: TableFilter,
#[serde(skip_serializing_if = "Option::is_none")]
pub event_types: Option<Vec<String>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub conditions: Option<Vec<Condition>>,
}
#[derive(Debug, Clone, Deserialize, Serialize, Default)]
pub struct TableFilter {
#[serde(skip_serializing_if = "Option::is_none")]
pub whitelist: Option<Vec<String>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub blacklist: Option<Vec<String>>,
}
#[derive(Debug, Clone, Deserialize, Serialize, Default)]
pub struct TransformConfig {
#[serde(default)]
pub fields: HashMap<String, HashMap<String, FieldTransform>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub global_transforms: Option<Vec<Transform>>,
}
#[derive(Debug, Clone, Deserialize, Serialize, Default)]
pub struct MappingConfig {
#[serde(default)]
pub tables: HashMap<String, TableMapping>,
#[serde(default = "default_unmapped_strategy")]
pub unmapped_fields_strategy: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub unmapped_fields_prefix: Option<String>,
}
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct TableMapping {
#[serde(skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub fields: Option<HashMap<String, FieldMapping>>,
}
#[derive(Debug, Clone, Deserialize, Serialize)]
#[serde(untagged)]
pub enum FieldMapping {
Simple(String),
Complex {
#[serde(skip_serializing_if = "Option::is_none")]
to: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
default: Option<Value>,
#[serde(skip_serializing_if = "Option::is_none")]
sources: Option<Vec<String>>,
#[serde(skip_serializing_if = "Option::is_none")]
separator: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
path: Option<String>,
},
}
#[derive(Debug, Clone, Deserialize, Serialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum FieldTransform {
Rename {
from: String,
to: String,
},
Convert {
field: String,
to_type: String,
},
Extract {
from: String,
path: String,
to: String,
},
Compute {
expression: String,
to: String,
},
}
#[derive(Debug, Clone, Deserialize, Serialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum Transform {
AddField { name: String, value: Value },
RemoveField { name: String },
AddTimestamp { field: String },
Lowercase { fields: Vec<String> },
Uppercase { fields: Vec<String> },
}
#[derive(Debug, Clone, Deserialize, Serialize)]
#[serde(tag = "op", rename_all = "snake_case")]
pub enum Condition {
Equals { field: String, value: Value },
NotEquals { field: String, value: Value },
Contains { field: String, value: Value },
GreaterThan { field: String, value: Value },
LessThan { field: String, value: Value },
In { field: String, values: Vec<Value> },
IsNull { field: String },
IsNotNull { field: String },
And { conditions: Vec<Condition> },
Or { conditions: Vec<Condition> },
}
fn default_unmapped_strategy() -> String {
"include".to_string()
}