use std::fmt::{Display, Formatter};
use std::str::FromStr;
use crate::err::{Context, TranslationError};
use crate::tree::ast::identifier::{CompoundIdentifier, Identifier, SimpleIdentifier};
use derive_more::derive::{From, TryUnwrap};
use serde::{Deserialize, Serialize};
use std::sync::Arc;
fn spanless_err(message: &str) -> TranslationError {
TranslationError::new(Context::new(0..=0, message))
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct UnqualifiedDatasetIdentifier {
pub namespace: Vec<SimpleIdentifier>,
pub table: SimpleIdentifier,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct QualifiedDatasetIdentifier {
pub space: SimpleIdentifier,
pub namespace: Vec<SimpleIdentifier>,
pub table: SimpleIdentifier,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, From)]
pub enum DatasetIdentifier {
Unqualified(UnqualifiedDatasetIdentifier),
Qualified(QualifiedDatasetIdentifier),
}
#[derive(Debug, Clone, PartialEq, Eq, From, TryUnwrap, Hash)]
pub enum ParsedDatasetIdentifier {
Valid(DatasetIdentifier),
Error(Arc<TranslationError>),
}
impl ParsedDatasetIdentifier {
pub fn valid_ref(&self) -> Result<&DatasetIdentifier, Arc<TranslationError>> {
match self {
ParsedDatasetIdentifier::Valid(id) => Ok(id),
ParsedDatasetIdentifier::Error(err) => Err(err.clone()),
}
}
pub fn valid(self) -> Result<DatasetIdentifier, Arc<TranslationError>> {
match self {
ParsedDatasetIdentifier::Valid(id) => Ok(id),
ParsedDatasetIdentifier::Error(err) => Err(err),
}
}
}
impl UnqualifiedDatasetIdentifier {
pub fn new(namespace: Vec<SimpleIdentifier>, table: SimpleIdentifier) -> Self {
Self { namespace, table }
}
pub fn qualify(self, default_space: &SimpleIdentifier) -> QualifiedDatasetIdentifier {
QualifiedDatasetIdentifier {
space: default_space.clone(),
namespace: self.namespace,
table: self.table,
}
}
pub fn is_cte_candidate(&self, cte_names: &std::collections::HashSet<Identifier>) -> bool {
self.namespace.is_empty() && cte_names.contains(&self.table.clone().into())
}
}
impl QualifiedDatasetIdentifier {
pub fn new(
space: SimpleIdentifier,
namespace: Vec<SimpleIdentifier>,
table: SimpleIdentifier,
) -> Self {
Self {
space,
namespace,
table,
}
}
pub fn from_canonical_str(s: &str) -> Result<Self, TranslationError> {
let id = DatasetIdentifier::from_ref_str(s)?;
match id {
DatasetIdentifier::Qualified(q) => Ok(q),
DatasetIdentifier::Unqualified(_) => Err(spanless_err(
"catalog dataset reference must include a space prefix (space:path)",
)),
}
}
pub fn from_namespace_parts(
ns_parts: &[SimpleIdentifier],
table: SimpleIdentifier,
) -> Result<Self, TranslationError> {
let space = ns_parts
.first()
.ok_or_else(|| spanless_err("catalog namespace must have at least one segment"))?;
Ok(Self {
space: space.clone(),
namespace: ns_parts[1..].to_vec(),
table,
})
}
pub fn to_namespace_table(&self) -> (Vec<String>, String) {
let mut namespace = vec![self.space.to_string()];
namespace.extend(self.namespace.iter().map(|segment| segment.to_string()));
(namespace, self.table.to_string())
}
pub fn qualified_path_segments(&self) -> Vec<SimpleIdentifier> {
let mut parts = vec![self.space.clone()];
parts.extend(self.namespace.iter().cloned());
parts.push(self.table.clone());
parts
}
pub fn to_execution_identifier(&self) -> Identifier {
let (namespace_parts, table) = self.to_namespace_table();
let schema = namespace_parts.join(".");
CompoundIdentifier::from_two(&schema, &table).into()
}
pub fn from_execution_identifier(id: &Identifier) -> Result<Self, TranslationError> {
let segments = id.segments();
if segments.len() != 2 {
return Err(spanless_err(
"execution identifier must have exactly two segments (schema.table)",
));
}
let schema = segments[0].as_str();
let table = segments[1].clone();
let mut dot_parts = schema.split('.');
let space_str = dot_parts
.next()
.ok_or_else(|| spanless_err("empty schema in execution identifier"))?;
let space = SimpleIdentifier::parse(space_str)?;
let namespace = dot_parts
.map(SimpleIdentifier::parse)
.collect::<Result<Vec<_>, _>>()?;
Ok(Self::new(space, namespace, table))
}
}
impl DatasetIdentifier {
pub fn from_ref_str(s: &str) -> Result<Self, TranslationError> {
let s = s.trim();
if s.is_empty() {
return Err(spanless_err("empty dataset reference"));
}
if let Some((space_str, rest)) = s.split_once(':') {
if space_str.is_empty() {
return Err(spanless_err("empty space in dataset reference"));
}
if rest.is_empty() {
return Err(spanless_err("empty path after ':' in dataset reference"));
}
if rest.contains(':') {
return Err(spanless_err(
"dataset reference may contain at most one ':' (between space and path)",
));
}
let space = SimpleIdentifier::parse(space_str)?;
let (namespace, table) = parse_dot_path(rest)?;
return Ok(Self::Qualified(QualifiedDatasetIdentifier {
space,
namespace,
table,
}));
}
let (namespace, table) = parse_dot_path(s)?;
Ok(Self::Unqualified(UnqualifiedDatasetIdentifier {
namespace,
table,
}))
}
pub fn qualify(
self,
default_space: Option<&SimpleIdentifier>,
) -> Result<QualifiedDatasetIdentifier, TranslationError> {
match self {
Self::Qualified(q) => Ok(q),
Self::Unqualified(u) => {
let space = default_space.ok_or_else(|| {
spanless_err("no default space for unqualified dataset reference")
})?;
Ok(u.qualify(space))
}
}
}
pub fn maybe_qualify(
&self,
default_space: Option<&SimpleIdentifier>,
) -> Option<QualifiedDatasetIdentifier> {
match self {
Self::Qualified(q) => Some(q.clone()),
Self::Unqualified(u) => default_space.map(|space| u.clone().qualify(space)),
}
}
pub fn namespace(&self) -> &[SimpleIdentifier] {
match self {
Self::Unqualified(u) => &u.namespace,
Self::Qualified(q) => &q.namespace,
}
}
pub fn table(&self) -> &SimpleIdentifier {
match self {
Self::Unqualified(u) => &u.table,
Self::Qualified(q) => &q.table,
}
}
pub fn is_cte_candidate(&self, cte_names: &std::collections::HashSet<Identifier>) -> bool {
match self {
Self::Unqualified(u) => u.is_cte_candidate(cte_names),
Self::Qualified(_) => false,
}
}
}
fn parse_dot_path(
rest: &str,
) -> Result<(Vec<SimpleIdentifier>, SimpleIdentifier), TranslationError> {
let segments: Result<Vec<_>, _> = rest
.split('.')
.map(|part| {
if part.is_empty() {
Err(spanless_err("empty segment in dataset reference"))
} else {
SimpleIdentifier::parse(part)
}
})
.collect();
let segments = segments?;
let table = segments
.last()
.ok_or_else(|| spanless_err("dataset reference must include a table name"))?
.clone();
let namespace = if segments.len() <= 1 {
Vec::new()
} else {
segments[..segments.len() - 1].to_vec()
};
Ok((namespace, table))
}
impl Display for UnqualifiedDatasetIdentifier {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
let segments = self.namespace.iter().chain(std::iter::once(&self.table));
for (i, segment) in segments.enumerate() {
if i > 0 {
write!(f, ".")?;
}
write!(f, "{segment}")?;
}
Ok(())
}
}
impl Display for QualifiedDatasetIdentifier {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
write!(f, "{}:", self.space)?;
let segments = self.namespace.iter().chain(std::iter::once(&self.table));
for (i, segment) in segments.enumerate() {
if i > 0 {
write!(f, ".")?;
}
write!(f, "{segment}")?;
}
Ok(())
}
}
impl Display for DatasetIdentifier {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self {
Self::Unqualified(u) => write!(f, "{u}"),
Self::Qualified(q) => write!(f, "{q}"),
}
}
}
impl Display for ParsedDatasetIdentifier {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self.valid_ref() {
Ok(id) => write!(f, "{}", id),
Err(_) => write!(f, "<error>"),
}
}
}
impl FromStr for DatasetIdentifier {
type Err = TranslationError;
fn from_str(s: &str) -> Result<Self, Self::Err> {
Self::from_ref_str(s)
}
}
impl FromStr for QualifiedDatasetIdentifier {
type Err = TranslationError;
fn from_str(s: &str) -> Result<Self, Self::Err> {
Self::from_canonical_str(s)
}
}
impl Serialize for QualifiedDatasetIdentifier {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
self.to_string().serialize(serializer)
}
}
impl<'de> Deserialize<'de> for QualifiedDatasetIdentifier {
fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
let s = String::deserialize(deserializer)?;
Self::from_canonical_str(&s).map_err(serde::de::Error::custom)
}
}
#[cfg(test)]
mod tests {
use super::*;
use pretty_assertions::assert_eq;
#[test]
fn parse_qualified_and_unqualified() {
let q = DatasetIdentifier::from_ref_str("acme:signals.events").unwrap();
assert!(matches!(q, DatasetIdentifier::Qualified(_)));
if let DatasetIdentifier::Qualified(q) = q {
assert_eq!(q.space.as_str(), "acme");
assert_eq!(q.namespace.len(), 1);
assert_eq!(q.table.as_str(), "events");
}
let u = DatasetIdentifier::from_ref_str("events").unwrap();
assert!(matches!(u, DatasetIdentifier::Unqualified(_)));
if let DatasetIdentifier::Unqualified(u) = u {
assert!(u.namespace.is_empty());
assert_eq!(u.table.as_str(), "events");
}
}
#[test]
fn qualify_applies_default_space() {
let u = DatasetIdentifier::from_ref_str("signals.events").unwrap();
let default = SimpleIdentifier::new("acme");
let q = u.qualify(Some(&default)).unwrap();
assert_eq!(q.to_string(), "acme:signals.events");
}
#[test]
fn rejects_multiple_colons() {
assert!(DatasetIdentifier::from_ref_str("a:b:c").is_err());
}
#[test]
fn to_namespace_table_space_table_and_nested() {
let q = QualifiedDatasetIdentifier::from_canonical_str("acme:events").unwrap();
assert_eq!(
q.to_namespace_table(),
(vec!["acme".to_string()], "events".to_string())
);
let q = QualifiedDatasetIdentifier::from_canonical_str("acme:signals.events").unwrap();
assert_eq!(
q.to_namespace_table(),
(
vec!["acme".to_string(), "signals".to_string()],
"events".to_string()
)
);
let q =
QualifiedDatasetIdentifier::from_canonical_str("acme:connectors.s3.events").unwrap();
assert_eq!(
q.to_namespace_table(),
(
vec![
"acme".to_string(),
"connectors".to_string(),
"s3".to_string()
],
"events".to_string()
)
);
}
#[test]
fn to_execution_identifier_uses_flat_schema_name() {
let q = QualifiedDatasetIdentifier::from_canonical_str("acme:signals.events").unwrap();
let id = q.to_execution_identifier();
assert_eq!(id.segments().len(), 2);
assert_eq!(id.segments()[0].as_str(), "acme.signals");
assert_eq!(id.segments()[1].as_str(), "events");
}
#[test]
fn from_execution_identifier_round_trips_canonical() {
let cases = [
"acme:events",
"acme:signals.events",
"acme:connectors.s3.events",
];
for canonical in cases {
let original = QualifiedDatasetIdentifier::from_canonical_str(canonical).unwrap();
let execution_id = original.to_execution_identifier();
let round_tripped =
QualifiedDatasetIdentifier::from_execution_identifier(&execution_id).unwrap();
assert_eq!(round_tripped, original, "round-trip failed for {canonical}");
}
}
#[test]
fn qualified_path_segments_includes_space_and_namespace() {
let q = QualifiedDatasetIdentifier::from_canonical_str("acme:signals.events").unwrap();
assert_eq!(
q.qualified_path_segments()
.iter()
.map(|s| s.as_str())
.collect::<Vec<_>>(),
vec!["acme", "signals", "events"]
);
}
#[test]
fn from_namespace_parts_round_trip() {
let table = SimpleIdentifier::new("events");
let id = QualifiedDatasetIdentifier::from_namespace_parts(
&[
SimpleIdentifier::new("acme"),
SimpleIdentifier::new("signals"),
],
table,
)
.unwrap();
assert_eq!(id.to_string(), "acme:signals.events");
}
#[test]
fn from_namespace_parts_deep_path() {
let table = SimpleIdentifier::new("events");
let id = QualifiedDatasetIdentifier::from_namespace_parts(
&[
SimpleIdentifier::new("acme"),
SimpleIdentifier::new("connectors"),
SimpleIdentifier::new("s3"),
],
table,
)
.unwrap();
assert_eq!(id.to_string(), "acme:connectors.s3.events");
assert_eq!(
id.to_namespace_table(),
(
vec![
"acme".to_string(),
"connectors".to_string(),
"s3".to_string()
],
"events".to_string()
)
);
}
#[test]
fn maybe_qualify_returns_none_without_default_space() {
let u = DatasetIdentifier::from_ref_str("events").unwrap();
assert!(u.maybe_qualify(None).is_none());
}
#[test]
fn maybe_qualify_qualifies_unqualified_with_default_space() {
let u = DatasetIdentifier::from_ref_str("signals.events").unwrap();
let default = SimpleIdentifier::new("acme");
assert_eq!(
u.maybe_qualify(Some(&default)).unwrap().to_string(),
"acme:signals.events"
);
}
#[test]
fn maybe_qualify_passes_through_qualified() {
let q = DatasetIdentifier::from_ref_str("acme:events").unwrap();
assert_eq!(q.maybe_qualify(None).unwrap().to_string(), "acme:events");
}
}