pub mod native;
pub mod spark;
use crate::error::Result;
use crate::ir::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,
pub output_modes: Vec<String>,
}
#[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(),
});
}
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())
}