horizon-sdk 10.2.0

Canonical Rust data access layer for the Horizon platform
Documentation
//! Partition specification for computing Kafka partition keys.
//!
//! Partition keys align Kafka topic partitions with Iceberg table partitions
//! for optimal file layout in the sink.

use chrono::{DateTime, Utc};
use uuid::Uuid;

use crate::types::error::{IcebergError, Result};

/// Information about a single partition field.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PartitionFieldInfo {
    pub partition_name: String,
    pub source_field_name: String,
    pub transform: String,
}

impl PartitionFieldInfo {
    #[must_use]
    pub fn new(source_field_name: &str, partition_name: &str, transform: &str) -> Self {
        Self {
            partition_name: partition_name.to_owned(),
            source_field_name: source_field_name.to_owned(),
            transform: transform.to_owned(),
        }
    }
}

/// A partition value in its canonical string representation.
///
/// Transforms are applied at construction time via typed constructors.
/// Use `Option<PartitionValue>` for nullable partition fields.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PartitionValue(String);

impl PartitionValue {
    /// The canonical string representation used in partition keys.
    #[must_use]
    pub fn as_str(&self) -> &str {
        &self.0
    }

    /// Create from a datetime with the specified Iceberg partition transform.
    ///
    /// Supported transforms: `identity`, `day`, `hour`, `month`, `year`, `void`.
    pub fn from_datetime(dt: DateTime<Utc>, transform: &str) -> Result<Self> {
        let formatted = match transform {
            "identity" => dt.to_rfc3339(),
            "day" => dt.format("%Y-%m-%d").to_string(),
            "hour" => dt.format("%Y-%m-%d-%H").to_string(),
            "month" => dt.format("%Y-%m").to_string(),
            "year" => dt.format("%Y").to_string(),
            "void" => String::new(),
            unsupported => {
                return Err(IcebergError::UnsupportedPartitionTransform {
                    transform: unsupported.to_owned(),
                }
                .into());
            }
        };
        Ok(Self(formatted))
    }
}

impl From<&str> for PartitionValue {
    fn from(s: &str) -> Self {
        Self(s.to_owned())
    }
}

impl From<String> for PartitionValue {
    fn from(s: String) -> Self {
        Self(s)
    }
}

impl From<Uuid> for PartitionValue {
    fn from(id: Uuid) -> Self {
        Self(id.to_string())
    }
}

/// Immutable partition specification discovered from the Iceberg catalog.
///
/// Contains partition fields and their transforms. Used to generate Kafka
/// partition keys that align with Iceberg table partitioning.
///
/// Keys are pipe-separated UTF-8 strings, e.g.:
/// `b"uuid-str|2025-01-15|audio|track-uuid"`.
#[derive(Debug, Clone)]
pub struct PartitionSpec {
    /// Discovered partition fields and their transforms.
    fields: Vec<PartitionFieldInfo>,
    /// The Iceberg table this spec was discovered from.
    table_name: String,
}

impl PartitionSpec {
    /// Compute a Kafka partition key from pre-transformed partition values.
    ///
    /// Values must be provided in the same order as `self.fields`.
    /// Use `None` for nullable fields or void transforms.
    ///
    /// # Errors
    ///
    /// Returns `IcebergError::PartitionValueCountMismatch` if the number of
    /// values does not match the number of partition fields.
    pub fn compute_key(&self, values: &[Option<PartitionValue>]) -> Result<Vec<u8>> {
        if values.len() != self.fields.len() {
            return Err(IcebergError::PartitionValueCountMismatch {
                expected: self.fields.len(),
                actual: values.len(),
            }
            .into());
        }

        let parts: Vec<&str> = values
            .iter()
            .map(|v| v.as_ref().map_or("", PartitionValue::as_str))
            .collect();

        Ok(parts.join("|").into_bytes())
    }

    #[must_use]
    pub fn fields(&self) -> &[PartitionFieldInfo] {
        &self.fields
    }

    #[must_use]
    pub fn new(table_name: &str, fields: Vec<PartitionFieldInfo>) -> Self {
        Self {
            fields,
            table_name: table_name.to_owned(),
        }
    }

    #[must_use]
    pub fn table_name(&self) -> &str {
        &self.table_name
    }

    /// Validate that every partition field's source column is in
    /// `supported_columns`, yielding a [`ValidatedPartition`].
    ///
    /// # Errors
    ///
    /// Returns `IcebergError::PartitionColumnUnmapped` for the first partition
    /// field whose source column is absent from `supported_columns`.
    pub(crate) fn validate_source_columns(
        self,
        supported_columns: &[&str],
    ) -> Result<ValidatedPartition> {
        for field in &self.fields {
            if !supported_columns.contains(&field.source_field_name.as_str()) {
                return Err(IcebergError::PartitionColumnUnmapped {
                    column_name: field.source_field_name.clone(),
                    table_name: self.table_name.clone(),
                }
                .into());
            }
        }
        Ok(ValidatedPartition(self))
    }
}

/// A [`PartitionSpec`] proven to reference only a table's supported columns.
///
/// Constructed only by [`PartitionSpec::validate_source_columns`], so requiring
/// one (as `IcebergRepository::new` does) makes an unvalidated spec impossible
/// to pass in.
#[derive(Debug, Clone)]
pub struct ValidatedPartition(PartitionSpec);

impl ValidatedPartition {
    /// Unwrap to the underlying validated specification.
    pub(crate) fn into_inner(self) -> PartitionSpec {
        self.0
    }
}

#[cfg(test)]
#[allow(clippy::unwrap_used, reason = "test assertions use unwrap for clarity")]
mod tests {
    use chrono::TimeZone as _;

    use super::*;

    #[test]
    fn compute_key_day_transform() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![PartitionFieldInfo::new("datetime", "datetime_day", "day")],
        );
        let dt = Utc.with_ymd_and_hms(2025, 1, 15, 14, 30, 0).unwrap();
        let value = PartitionValue::from_datetime(dt, "day").unwrap();
        let key = spec.compute_key(&[Some(value)]).unwrap();
        assert_eq!(key, b"2025-01-15");
    }

    #[test]
    fn compute_key_hour_transform() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![PartitionFieldInfo::new("datetime", "datetime_hour", "hour")],
        );
        let dt = Utc.with_ymd_and_hms(2025, 1, 15, 14, 30, 0).unwrap();
        let value = PartitionValue::from_datetime(dt, "hour").unwrap();
        let key = spec.compute_key(&[Some(value)]).unwrap();
        assert_eq!(key, b"2025-01-15-14");
    }

    #[test]
    fn compute_key_identity_string() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![PartitionFieldInfo::new("id", "id", "identity")],
        );
        let key = spec
            .compute_key(&[Some(PartitionValue::from("abc123"))])
            .unwrap();
        assert_eq!(key, b"abc123");
    }

    #[test]
    fn compute_key_identity_uuid() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![PartitionFieldInfo::new(
                "stream_id",
                "stream_id",
                "identity",
            )],
        );
        let id = Uuid::nil();
        let key = spec.compute_key(&[Some(PartitionValue::from(id))]).unwrap();
        assert_eq!(key, b"00000000-0000-0000-0000-000000000000");
    }

    #[test]
    fn compute_key_month_transform() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![PartitionFieldInfo::new(
                "datetime",
                "datetime_month",
                "month",
            )],
        );
        let dt = Utc.with_ymd_and_hms(2025, 3, 15, 14, 30, 0).unwrap();
        let value = PartitionValue::from_datetime(dt, "month").unwrap();
        let key = spec.compute_key(&[Some(value)]).unwrap();
        assert_eq!(key, b"2025-03");
    }

    #[test]
    fn compute_key_multiple_fields() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![
                PartitionFieldInfo::new("stream_id", "stream_id", "identity"),
                PartitionFieldInfo::new("datetime", "datetime_day", "day"),
                PartitionFieldInfo::new("data_type", "data_type", "identity"),
            ],
        );
        let stream_id = Uuid::nil();
        let dt = Utc.with_ymd_and_hms(2025, 1, 15, 12, 0, 0).unwrap();
        let key = spec
            .compute_key(&[
                Some(PartitionValue::from(stream_id)),
                Some(PartitionValue::from_datetime(dt, "day").unwrap()),
                Some(PartitionValue::from("audio")),
            ])
            .unwrap();
        let expected = format!("{stream_id}|2025-01-15|audio");
        assert_eq!(key, expected.as_bytes());
    }

    #[test]
    fn compute_key_none_value() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![PartitionFieldInfo::new("optional", "optional", "identity")],
        );
        let key = spec.compute_key(&[None]).unwrap();
        assert_eq!(key, b"");
    }

    #[test]
    fn compute_key_void_transform() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![PartitionFieldInfo::new("value", "value_void", "void")],
        );
        let dt = Utc.with_ymd_and_hms(2025, 1, 15, 14, 30, 0).unwrap();
        let value = PartitionValue::from_datetime(dt, "void").unwrap();
        let key = spec.compute_key(&[Some(value)]).unwrap();
        assert_eq!(key, b"");
    }

    #[test]
    fn compute_key_wrong_value_count_fails() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![
                PartitionFieldInfo::new("a", "a", "identity"),
                PartitionFieldInfo::new("b", "b", "identity"),
            ],
        );
        let result = spec.compute_key(&[Some(PartitionValue::from("only-one"))]);
        assert!(result.is_err());
        let err = result.unwrap_err().to_string();
        assert!(err.contains("expected 2 partition values but got 1"));
    }

    #[test]
    fn compute_key_year_transform() {
        let spec = PartitionSpec::new(
            "test.table",
            vec![PartitionFieldInfo::new("datetime", "datetime_year", "year")],
        );
        let dt = Utc.with_ymd_and_hms(2025, 3, 15, 14, 30, 0).unwrap();
        let value = PartitionValue::from_datetime(dt, "year").unwrap();
        let key = spec.compute_key(&[Some(value)]).unwrap();
        assert_eq!(key, b"2025");
    }

    #[test]
    fn from_datetime_unsupported_transform_fails() {
        let dt = Utc.with_ymd_and_hms(2025, 1, 15, 14, 30, 0).unwrap();
        let result = PartitionValue::from_datetime(dt, "bucket");
        assert!(result.is_err());
        let err = result.unwrap_err().to_string();
        assert!(err.contains("unsupported datetime partition transform: bucket"));
    }

    #[test]
    fn partition_field_info_equality() {
        let a = PartitionFieldInfo::new("col", "col", "identity");
        let b = PartitionFieldInfo::new("col", "col", "identity");
        let c = PartitionFieldInfo::new("col", "col_day", "day");
        assert_eq!(a, b);
        assert_ne!(a, c);
    }
}