horizon-sdk 10.2.0

Canonical Rust data access layer for the Horizon platform
Documentation
//! Iceberg catalog discovery for field IDs and partition specifications.
//!
//! Connects to an Iceberg REST catalog at SDK initialization time to
//! discover schema metadata. Field IDs and partition specs are immutable
//! once assigned in Iceberg, so they are cached for the SDK lifetime.

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

use iceberg::spec::{NestedFieldRef, Transform, Type};
use iceberg::{Catalog as _, CatalogBuilder as _, TableIdent};
use iceberg_catalog_rest::{RestCatalog, RestCatalogBuilder};
use iceberg_storage_opendal::OpenDalStorageFactory;
use tracing::{debug, info};

use super::partition::{PartitionFieldInfo, PartitionSpec};
use crate::types::arrow::FieldIdMap;
use crate::types::config::IcebergConfig;
use crate::types::constants::{DATA_ROW_EXPECTED_COLUMNS, METADATA_ROW_EXPECTED_COLUMNS};
use crate::types::error::{IcebergError, Result};

/// Look up the expected column list for a discovered table name.
///
/// Returns `None` for custom table names so callers can skip the preflight
/// check without forcing every user-defined table to match a Horizon row
/// shape. The `metadata_row` suffix is checked first because `"metadata_row"`
/// also ends with `"data_row"`.
fn expected_columns_for_table(table_name: &str) -> Option<&'static [&'static str]> {
    if table_name.ends_with("metadata_row") {
        Some(METADATA_ROW_EXPECTED_COLUMNS)
    } else if table_name.ends_with("data_row") {
        Some(DATA_ROW_EXPECTED_COLUMNS)
    } else {
        None
    }
}

/// Discover field ID mappings from an Iceberg REST catalog.
///
/// Connects to the catalog, loads the table schema, and extracts all field IDs
/// including nested types (lists, structs, maps) using dot-notation paths.
///
/// # Errors
///
/// Returns `IcebergError::CatalogConnection` if the catalog cannot be
/// reached, `IcebergError::InvalidTableName` if the identifier is
/// malformed, or `IcebergError::TableLoad` if the table does not exist.
pub async fn discover_field_ids(config: &IcebergConfig, table_name: &str) -> Result<FieldIdMap> {
    info!(
        catalog = config.catalog_uri.as_str(),
        table = table_name,
        "discovering field IDs"
    );

    let catalog = build_catalog(config).await?;
    let table_ident = parse_table_ident(table_name)?;
    let table =
        catalog
            .load_table(&table_ident)
            .await
            .map_err(|source| IcebergError::TableLoad {
                table_name: table_name.to_owned(),
                source: Box::new(source),
            })?;

    let schema = table.metadata().current_schema();
    let mut mapping = HashMap::new();

    for field in schema.as_struct().fields() {
        collect_field_ids(field, &mut mapping, "");
    }

    debug!(
        table = table_name,
        fields = mapping.len(),
        "discovered field IDs"
    );

    let field_id_map = FieldIdMap::new(table_name.to_owned(), mapping);

    if let Some(expected) = expected_columns_for_table(table_name) {
        field_id_map.validate_columns(expected)?;
    }

    Ok(field_id_map)
}

/// Discover the partition specification from an Iceberg REST catalog.
///
/// # Errors
///
/// Returns `IcebergError::CatalogConnection` if the catalog cannot be
/// reached, `IcebergError::InvalidTableName` if the identifier is
/// malformed, or `IcebergError::TableLoad` if the table does not exist.
pub async fn discover_partition_spec(
    config: &IcebergConfig,
    table_name: &str,
) -> Result<PartitionSpec> {
    info!(
        catalog = config.catalog_uri.as_str(),
        table = table_name,
        "discovering partition spec"
    );

    let catalog = build_catalog(config).await?;
    let table_ident = parse_table_ident(table_name)?;
    let table =
        catalog
            .load_table(&table_ident)
            .await
            .map_err(|source| IcebergError::TableLoad {
                table_name: table_name.to_owned(),
                source: Box::new(source),
            })?;

    let schema = table.metadata().current_schema();
    let spec = table.metadata().default_partition_spec();

    let mut fields = Vec::new();
    for partition_field in spec.fields() {
        let source_field = schema
            .field_by_id(partition_field.source_id)
            .ok_or_else(|| IcebergError::PartitionFieldNotFound {
                field_id: partition_field.source_id,
                table_name: table_name.to_owned(),
            })?;

        let transform_name = transform_to_string(partition_field.transform);

        fields.push(PartitionFieldInfo::new(
            &source_field.name,
            &partition_field.name,
            &transform_name,
        ));
    }

    debug!(
        table = table_name,
        partition_fields = fields.len(),
        "discovered partition spec"
    );

    Ok(PartitionSpec::new(table_name, fields))
}

/// Discover both field IDs and partition spec from a single catalog load.
///
/// More efficient than calling `discover_field_ids` and `discover_partition_spec`
/// separately, as it only connects to the catalog and loads the table once.
///
/// # Errors
///
/// Returns `IcebergError::CatalogConnection` if the catalog cannot be
/// reached, `IcebergError::InvalidTableName` if the identifier is
/// malformed, or `IcebergError::TableLoad` if the table does not exist.
pub async fn discover_table_metadata(
    config: &IcebergConfig,
    table_name: &str,
) -> Result<(FieldIdMap, PartitionSpec)> {
    info!(
        catalog = config.catalog_uri.as_str(),
        table = table_name,
        "discovering table metadata"
    );

    let catalog = build_catalog(config).await?;
    let table_ident = parse_table_ident(table_name)?;
    let table =
        catalog
            .load_table(&table_ident)
            .await
            .map_err(|source| IcebergError::TableLoad {
                table_name: table_name.to_owned(),
                source: Box::new(source),
            })?;

    let schema = table.metadata().current_schema();

    let mut mapping = HashMap::new();
    for field in schema.as_struct().fields() {
        collect_field_ids(field, &mut mapping, "");
    }
    let field_id_map = FieldIdMap::new(table_name.to_owned(), mapping);

    if let Some(expected) = expected_columns_for_table(table_name) {
        field_id_map.validate_columns(expected)?;
    }

    let spec = table.metadata().default_partition_spec();
    let mut fields = Vec::new();
    for partition_field in spec.fields() {
        let source_field = schema
            .field_by_id(partition_field.source_id)
            .ok_or_else(|| IcebergError::PartitionFieldNotFound {
                field_id: partition_field.source_id,
                table_name: table_name.to_owned(),
            })?;
        fields.push(PartitionFieldInfo::new(
            &source_field.name,
            &partition_field.name,
            &transform_to_string(partition_field.transform),
        ));
    }
    let partition_spec = PartitionSpec::new(table_name, fields);

    debug!(
        table = table_name,
        field_ids = field_id_map.len(),
        partition_fields = partition_spec.fields().len(),
        "discovered table metadata"
    );

    Ok((field_id_map, partition_spec))
}

/// Build a REST catalog connection from SDK configuration.
async fn build_catalog(config: &IcebergConfig) -> Result<RestCatalog> {
    let mut props: HashMap<String, String> = HashMap::new();
    props.insert("uri".to_owned(), config.catalog_uri.clone());
    if let Some(warehouse) = &config.catalog_warehouse {
        props.insert("warehouse".to_owned(), warehouse.clone());
    }

    // OAuth2 client-credentials for a secured REST catalog. When unset the
    // catalog is reached unauthenticated (open Lakekeeper / dev). Mirrors the
    // sink's ICEBERG__CREDENTIAL / ICEBERG__OAUTH2_SERVER_URI / ICEBERG__SCOPE.
    if let Some(credential) = &config.credential {
        props.insert("credential".to_owned(), credential.clone());
    }
    if let Some(oauth2_server_uri) = &config.oauth2_server_uri {
        props.insert("oauth2-server-uri".to_owned(), oauth2_server_uri.clone());
    }
    if let Some(scope) = &config.scope {
        props.insert("scope".to_owned(), scope.clone());
    }

    // iceberg 0.9 split storage backends into a separate crate; the REST catalog
    // requires an explicit StorageFactory before any FileIO operations succeed.
    // Wire the OpenDAL S3 factory so the catalog can read/write S3-backed
    // warehouses; credentials are vended by the REST server.
    let storage_factory = Arc::new(OpenDalStorageFactory::S3 {
        configured_scheme: "s3".to_owned(),
        customized_credential_load: None,
    });

    RestCatalogBuilder::default()
        .with_storage_factory(storage_factory)
        .load("rest", props)
        .await
        .map_err(|source| {
            IcebergError::CatalogConnection {
                catalog_uri: config.catalog_uri.clone(),
                source: Box::new(source),
            }
            .into()
        })
}

/// Recursively collect field IDs from a schema field, including nested types.
fn collect_field_ids(field: &NestedFieldRef, mapping: &mut HashMap<String, i32>, path: &str) {
    let full_path = if path.is_empty() {
        field.name.clone()
    } else {
        format!("{path}.{}", field.name)
    };
    mapping.insert(full_path.clone(), field.id);
    collect_nested_type_ids(&field.field_type, mapping, &full_path);
}

/// Recursively collect field IDs from nested types (lists, structs, maps).
fn collect_nested_type_ids(field_type: &Type, mapping: &mut HashMap<String, i32>, path: &str) {
    match field_type {
        Type::List(list_type) => {
            let element_path = format!("{path}.element");
            mapping.insert(element_path.clone(), list_type.element_field.id);
            collect_nested_type_ids(&list_type.element_field.field_type, mapping, &element_path);
        }
        Type::Struct(struct_type) => {
            for nested_field in struct_type.fields() {
                collect_field_ids(nested_field, mapping, path);
            }
        }
        Type::Map(map_type) => {
            let key_path = format!("{path}.key");
            let value_path = format!("{path}.value");
            mapping.insert(key_path.clone(), map_type.key_field.id);
            mapping.insert(value_path.clone(), map_type.value_field.id);
            collect_nested_type_ids(&map_type.key_field.field_type, mapping, &key_path);
            collect_nested_type_ids(&map_type.value_field.field_type, mapping, &value_path);
        }
        Type::Primitive(_) => {}
    }
}

/// Parse a dot-separated table name into an Iceberg `TableIdent`.
fn parse_table_ident(table_name: &str) -> Result<TableIdent> {
    TableIdent::from_strs(table_name.split('.')).map_err(|source| {
        IcebergError::InvalidTableName {
            table_name: table_name.to_owned(),
            source: Box::new(source),
        }
        .into()
    })
}

/// Convert an Iceberg `Transform` to its string representation.
fn transform_to_string(transform: Transform) -> String {
    match transform {
        Transform::Identity => "identity".to_owned(),
        Transform::Day => "day".to_owned(),
        Transform::Hour => "hour".to_owned(),
        Transform::Month => "month".to_owned(),
        Transform::Year => "year".to_owned(),
        Transform::Bucket(n) => format!("bucket[{n}]"),
        Transform::Truncate(n) => format!("truncate[{n}]"),
        Transform::Void => "void".to_owned(),
        Transform::Unknown => "unknown".to_owned(),
    }
}