horizon-sdk 16.0.0

Canonical Rust data access layer for the Horizon platform
//! Arrow IPC serialization traits for Iceberg writes.
//!
//! Types that implement `ArrowSerializable` can be converted to Arrow record batches
//! with Iceberg-compatible field ID metadata (`PARQUET:field_id`). This is required
//! for the Kafka-backed write path: SDK → Arrow IPC → Kafka → horizon-iceberg-sink → Iceberg.

use std::collections::HashMap;
use std::sync::Arc;

use arrow_array::RecordBatch;
use arrow_schema::{Field, Schema};

use super::error::{IcebergError, Result};

/// Trait for types that can be serialized to Arrow record batches for Iceberg writes.
pub trait ArrowSerializable {
    /// Build the Arrow schema with Iceberg field ID annotations.
    ///
    /// # Errors
    ///
    /// Returns `IcebergError::ColumnNotFound` if field IDs cannot be resolved.
    fn arrow_schema(field_id_map: &FieldIdMap) -> Result<Arc<Schema>>;

    /// Convert this instance to an Arrow `RecordBatch`.
    ///
    /// # Errors
    ///
    /// Returns `IcebergError::ColumnNotFound` on serialization failure.
    fn to_record_batch(&self, field_id_map: &FieldIdMap) -> Result<RecordBatch>;
}

/// Mapping of column names to Iceberg field IDs.
///
/// Field IDs are discovered from the Iceberg REST catalog at SDK initialization
/// time and are immutable for the lifetime of a table schema. They are used to
/// annotate Arrow schemas with `PARQUET:field_id` metadata so the sink writes
/// Parquet files that Iceberg can read.
///
/// Dot-notation is used for nested fields (e.g. `"vector.element"`).
#[derive(Debug, Clone)]
pub struct FieldIdMap {
    /// Map from column path (e.g. `"datetime"`, `"vector.element"`) to Iceberg field ID.
    column_to_field_id: HashMap<String, i32>,
    /// The Iceberg table name this mapping was discovered from.
    table_name: String,
}

impl FieldIdMap {
    /// Get the field ID for a column path, or return an error if not found.
    ///
    /// # Errors
    ///
    /// Returns `IcebergError::ColumnNotFound` if the column path is not in
    /// the mapping.
    pub fn get_field_id(&self, column_path: &str) -> Result<&i32> {
        self.column_to_field_id.get(column_path).ok_or_else(|| {
            IcebergError::ColumnNotFound {
                column_name: column_path.to_owned(),
                table_name: self.table_name.clone(),
            }
            .into()
        })
    }

    /// Whether the mapping is empty.
    #[must_use]
    pub fn is_empty(&self) -> bool {
        self.column_to_field_id.is_empty()
    }

    /// The number of field ID mappings (including nested paths).
    #[must_use]
    pub fn len(&self) -> usize {
        self.column_to_field_id.len()
    }

    #[must_use]
    pub const fn new(table_name: String, column_to_field_id: HashMap<String, i32>) -> Self {
        Self {
            column_to_field_id,
            table_name,
        }
    }

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

    /// Validate that all required columns exist in the mapping.
    ///
    /// Callers are expected to pass top-level column names.
    ///
    /// # Errors
    ///
    /// Returns `IcebergError::MissingColumns` listing any missing columns.
    pub(crate) fn validate_columns(&self, required_columns: &[&str]) -> Result<()> {
        let missing_columns: Vec<String> = required_columns
            .iter()
            .filter_map(|column| {
                if self.column_to_field_id.contains_key(*column) {
                    None
                } else {
                    Some((*column).to_owned())
                }
            })
            .collect();
        if missing_columns.is_empty() {
            Ok(())
        } else {
            Err(IcebergError::MissingColumns {
                missing_columns,
                table_name: self.table_name.clone(),
            }
            .into())
        }
    }
}

/// Set `PARQUET:field_id` metadata on an Arrow field in place.
pub(crate) fn update_field_id(field: &mut Field, field_id: i32) {
    field
        .metadata_mut()
        .insert("PARQUET:field_id".to_owned(), field_id.to_string());
}

#[cfg(test)]
#[allow(
    clippy::expect_used,
    clippy::panic,
    reason = "tests use expect/panic for brevity"
)]
mod tests {
    use super::*;
    use crate::types::error::HorizonError;

    fn sample_map() -> FieldIdMap {
        let mut mapping = HashMap::new();
        mapping.insert("datetime".to_owned(), 1_i32);
        mapping.insert("vector".to_owned(), 2_i32);
        mapping.insert("vector.element".to_owned(), 3_i32);
        FieldIdMap::new("test.table".to_owned(), mapping)
    }

    #[test]
    fn validate_columns_accepts_matching_top_level_fields() {
        let map = sample_map();
        map.validate_columns(&["datetime", "vector"])
            .expect("matching columns should validate");
    }

    #[test]
    fn validate_columns_accepts_empty_required_list() {
        let map = sample_map();
        map.validate_columns(&[])
            .expect("empty required list should validate");
    }

    #[test]
    fn validate_columns_reports_missing_columns_with_table_name() {
        let map = sample_map();
        let error = map
            .validate_columns(&["datetime", "ghost"])
            .expect_err("missing column should error");
        let HorizonError::Iceberg(IcebergError::MissingColumns {
            missing_columns,
            table_name,
        }) = error
        else {
            panic!("expected IcebergError::MissingColumns variant");
        };
        assert_eq!(missing_columns, vec!["ghost".to_owned()]);
        assert_eq!(table_name, "test.table");
    }

    #[test]
    fn get_field_id_returns_column_not_found_for_missing_column() {
        let map = sample_map();
        let error = map
            .get_field_id("ghost")
            .expect_err("missing column should error");
        let HorizonError::Iceberg(IcebergError::ColumnNotFound {
            column_name,
            table_name,
        }) = error
        else {
            panic!("expected IcebergError::ColumnNotFound variant");
        };
        assert_eq!(column_name, "ghost");
        assert_eq!(table_name, "test.table");
    }
}