pub mod native;
pub mod spark;
#[cfg(feature = "native")]
pub(crate) mod checkpoint;
use crate::error::Result;
use crate::ir::{Dataset, OutputType, Pipeline};
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Capabilities {
pub name: String,
pub batch: bool,
pub streaming: bool,
pub event_time: bool,
pub cdc: bool,
pub expectations: bool,
#[serde(default)]
pub incremental_state: bool,
#[serde(default)]
pub partition_transforms: bool,
#[serde(default)]
pub graph_sinks: bool,
pub output_modes: Vec<String>,
}
pub(crate) fn is_partition_transform(col: &str) -> bool {
col.contains('(')
}
pub(crate) fn is_graph_sink(ds: &Dataset) -> bool {
if ds.output_type != OutputType::Sink {
return false;
}
match ds.format.as_deref() {
Some("falkordb-node") => ds.properties.contains_key("knut.node.label"),
Some("falkordb-edge") => ds.properties.contains_key("knut.edge.rel"),
_ => false,
}
}
pub(crate) fn graph_sink_target(ds: &Dataset) -> String {
if let Some(label) = ds.properties.get("knut.node.label") {
format!("node label `{label}`")
} else if let Some(rel) = ds.properties.get("knut.edge.rel") {
format!("edge relationship `{rel}`")
} else {
"graph element".to_string()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SdpFlowDrop {
pub flow: String,
pub dropped: Vec<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct SdpLoweringReport {
pub backend: String,
pub spark_4_2_plus: bool,
pub flow_drops: Vec<SdpFlowDrop>,
}
impl SdpLoweringReport {
pub fn is_lossless(&self) -> bool {
self.flow_drops.is_empty()
}
pub fn flow(&self, flow: &str) -> Option<&SdpFlowDrop> {
self.flow_drops.iter().find(|d| d.flow == flow)
}
pub fn all_dropped(&self) -> Vec<String> {
self.flow_drops
.iter()
.flat_map(|d| d.dropped.iter().cloned())
.collect()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RunEvent {
pub timestamp: Option<String>,
pub element: Option<String>,
pub message: String,
pub phase: Option<RunPhase>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RunPhase {
Queued,
Planning,
Running,
Completed,
Failed,
}
pub trait RunHandle {
fn poll_events(&mut self) -> Result<Vec<RunEvent>>;
fn cancel(&mut self) -> Result<()>;
}
pub trait ExecBackend {
type Run: RunHandle;
fn capabilities(&self) -> Capabilities;
fn check(&self, pipeline: &Pipeline) -> Result<()> {
let caps = self.capabilities();
if pipeline.is_streaming() && !caps.streaming {
return Err(crate::error::ThundError::Unsupported {
backend: leak(caps.name.clone()),
what: "streaming flows".into(),
});
}
if !caps.partition_transforms {
for d in &pipeline.datasets {
if let Some(t) = d.partition_cols.iter().find(|c| is_partition_transform(c)) {
return Err(crate::error::ThundError::Unsupported {
backend: leak(caps.name.clone()),
what: format!(
"partition transform `{t}` on dataset `{}` \
(this backend takes identity partition columns only — \
materialize the value as a column and partition by that)",
d.name
),
});
}
}
}
if !caps.graph_sinks {
let graph_sinks: Vec<&Dataset> = pipeline
.datasets
.iter()
.filter(|d| is_graph_sink(d))
.collect();
if let Some(d) = graph_sinks.first() {
return Err(crate::error::ThundError::Unsupported {
backend: leak(caps.name.clone()),
what: format!(
"graph-sink pipelines — SparkBackend cannot execute the {n} graph-sink \
dataset(s) in this pipeline (a rel2graph Cypher graph load). \
First offender: dataset `{name}` (format `{fmt}`, {target}). \
WHY: SDP over Spark Connect has no shape to carry a graph write — an SDP \
flow lowers to a SQL relation and an SDP output to a table, neither of \
which can hold the idempotent Cypher `UNWIND … MERGE` a graph load is; \
`DefineOutput`-ing the unresolvable `{fmt}` format would only fail \
obscurely server-side after a graph was already created. \
WHAT TO DO: run this pipeline on the NativeBackend (build \
`--features native,rel2graph`), which executes the load through \
`knut_bifrost::sink::FalkorDbSink`; OR render the PySpark projection \
(`GraphPlan::to_pyspark`) and run that graph load as a Spark job.",
n = graph_sinks.len(),
name = d.name,
fmt = d.format.as_deref().unwrap_or(""),
target = graph_sink_target(d),
),
});
}
}
for f in &pipeline.flows {
if matches!(f.kind, crate::ir::FlowKind::Cdc { .. }) && !caps.cdc {
return Err(crate::error::ThundError::Unsupported {
backend: leak(caps.name.clone()),
what: format!("CDC flow `{}`", f.name),
});
}
if !f.expectations.is_empty() && !caps.expectations {
return Err(crate::error::ThundError::Unsupported {
backend: leak(caps.name.clone()),
what: format!("expectations on flow `{}`", f.name),
});
}
}
Ok(())
}
fn run(&self, pipeline: &Pipeline) -> Result<Self::Run>;
}
fn leak(s: String) -> &'static str {
Box::leak(s.into_boxed_str())
}