use serde::{Deserialize, Serialize};
use super::schema::DatasetSchema;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OutputType {
MaterializedView,
Table,
TemporaryView,
Sink,
}
impl OutputType {
pub fn label(self) -> &'static str {
match self {
OutputType::MaterializedView => "materialized_view",
OutputType::Table => "table",
OutputType::TemporaryView => "temporary_view",
OutputType::Sink => "sink",
}
}
pub fn to_sdp(self) -> knut_pipelines::OutputType {
match self {
OutputType::MaterializedView => knut_pipelines::OutputType::MaterializedView,
OutputType::Table => knut_pipelines::OutputType::Table,
OutputType::TemporaryView => knut_pipelines::OutputType::TemporaryView,
OutputType::Sink => knut_pipelines::OutputType::Sink,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Materialization {
Full,
Incremental,
}
impl Default for Materialization {
fn default() -> Self {
Materialization::Full
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct SourceLocation {
pub file_name: Option<String>,
pub line_number: Option<i32>,
}
impl SourceLocation {
pub fn is_empty(&self) -> bool {
self.file_name.is_none() && self.line_number.is_none()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Dataset {
pub name: String,
pub output_type: OutputType,
#[serde(default)]
pub materialization: Materialization,
#[serde(default, skip_serializing_if = "DatasetSchema::is_empty")]
pub schema: DatasetSchema,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub format: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub partition_cols: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub comment: Option<String>,
#[serde(default, skip_serializing_if = "SourceLocation::is_empty")]
pub source_location: SourceLocation,
}
impl Dataset {
pub fn new(name: impl Into<String>, output_type: OutputType) -> Self {
Dataset {
name: name.into(),
output_type,
materialization: Materialization::default(),
schema: DatasetSchema::new(),
format: None,
partition_cols: Vec::new(),
comment: None,
source_location: SourceLocation::default(),
}
}
pub fn with_schema(mut self, schema: DatasetSchema) -> Self {
self.schema = schema;
self
}
pub fn incremental(mut self) -> Self {
self.materialization = Materialization::Incremental;
self
}
pub fn with_format(mut self, fmt: impl Into<String>) -> Self {
self.format = Some(fmt.into());
self
}
pub fn to_sdp(&self) -> knut_pipelines::Dataset {
let mut d = knut_pipelines::Dataset::new(self.name.clone(), self.output_type.to_sdp());
d.comment = self.comment.clone();
d.format = self.format.clone();
d.partition_cols = self.partition_cols.clone();
d.schema = self.schema.to_sdp_schema_string();
d
}
}