use serde::{Deserialize, Serialize};
use super::expectation::Expectation;
use super::streaming::{OutputMode, SourceSpec, Trigger, Watermark, WindowSpec};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ScdType {
Type1,
Type2,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CdcSpec {
pub keys: Vec<String>,
pub sequence_by: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub apply_as_deletes: Option<String>,
pub scd_type: ScdType,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FlowKind {
Batch,
Streaming {
source: SourceSpec,
#[serde(default, skip_serializing_if = "Option::is_none")]
watermark: Option<Watermark>,
#[serde(default)]
trigger: Trigger,
#[serde(default)]
output_mode: OutputMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
window: Option<WindowSpec>,
},
Cdc {
cdc: CdcSpec,
},
}
impl FlowKind {
pub fn is_unbounded(&self) -> bool {
matches!(self, FlowKind::Streaming { .. } | FlowKind::Cdc { .. })
}
pub fn label(&self) -> &'static str {
match self {
FlowKind::Batch => "batch",
FlowKind::Streaming { .. } => "streaming",
FlowKind::Cdc { .. } => "cdc",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Flow {
pub name: String,
pub target: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub reads: Vec<String>,
pub kind: FlowKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub query: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub projection: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub filter: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub expectations: Vec<Expectation>,
}
impl Flow {
pub fn batch(
name: impl Into<String>,
target: impl Into<String>,
reads: impl IntoIterator<Item = impl Into<String>>,
) -> Self {
Flow {
name: name.into(),
target: target.into(),
reads: reads.into_iter().map(Into::into).collect(),
kind: FlowKind::Batch,
query: None,
projection: Vec::new(),
filter: None,
expectations: Vec::new(),
}
}
pub fn streaming(
name: impl Into<String>,
target: impl Into<String>,
source: SourceSpec,
) -> Self {
Flow {
name: name.into(),
target: target.into(),
reads: Vec::new(),
kind: FlowKind::Streaming {
source,
watermark: None,
trigger: Trigger::default(),
output_mode: OutputMode::default(),
window: None,
},
query: None,
projection: Vec::new(),
filter: None,
expectations: Vec::new(),
}
}
pub fn with_query(mut self, sql: impl Into<String>) -> Self {
self.query = Some(sql.into());
self
}
pub fn with_projection<I, S>(mut self, cols: I) -> Self
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
self.projection = cols.into_iter().map(Into::into).collect();
self
}
pub fn with_filter(mut self, predicate: impl Into<String>) -> Self {
self.filter = Some(predicate.into());
self
}
pub fn has_pushdown(&self) -> bool {
!self.projection.is_empty() || self.filter.is_some()
}
pub fn synth_query(&self) -> Option<String> {
if self.query.is_some() || !self.has_pushdown() {
return None;
}
let [input] = self.reads.as_slice() else {
return None;
};
let cols = if self.projection.is_empty() {
"*".to_string()
} else {
self.projection.join(", ")
};
let mut sql = format!("SELECT {cols} FROM {input}");
if let Some(pred) = &self.filter {
sql.push_str(" WHERE ");
sql.push_str(pred);
}
Some(sql)
}
pub fn expect(mut self, e: Expectation) -> Self {
self.expectations.push(e);
self
}
pub fn with_watermark(mut self, wm: Watermark) -> Self {
if let FlowKind::Streaming { watermark, .. } = &mut self.kind {
*watermark = Some(wm);
}
self
}
pub fn with_window(mut self, win: WindowSpec) -> Self {
if let FlowKind::Streaming { window, .. } = &mut self.kind {
*window = Some(win);
}
self
}
pub fn to_sdp(&self) -> knut_pipelines::Flow {
let mut f =
knut_pipelines::Flow::new(self.name.clone(), self.target.clone(), self.reads.clone());
if let Some(q) = self.query.clone().or_else(|| self.synth_query()) {
f = f.with_query(q);
}
if let FlowKind::Streaming {
trigger: Trigger::AvailableNow,
..
} = &self.kind
{
f.once = true;
}
f
}
}