horizon-sdk 16.0.0

Canonical Rust data access layer for the Horizon platform
//! Iceberg repository with async write path for high-throughput data ingestion.
//!
//! Writes are asynchronous and eventually consistent: data is queued to Kafka
//! and consumed by `horizon-iceberg-sink`, which handles batching, sorting,
//! and atomic commits to Iceberg.
//!
//! **Consistency model:**
//! - Writes return after Kafka acknowledgement (`acks=all`)
//! - Data becomes visible in Iceberg asynchronously (30-60s low throughput, <1s high throughput)
//! - At-least-once semantics (duplicates possible on crash/retry)

use std::collections::BTreeMap;
use std::fmt;
use std::time::Duration;

use futures::future::try_join_all;
use tracing::info;

use super::partition::{PartitionSpec, PartitionValue, ValidatedPartition};
use super::serialization::{concat_batches, serialize_to_arrow_ipc};
use crate::kafka::producer::KafkaProducer;
use crate::types::arrow::{ArrowSerializable, FieldIdMap};
use crate::types::error::Result;

/// Maximum number of partition-group messages produced concurrently in
/// [`IcebergRepository::insert_batch`]. A bounded window lets librdkafka batch
/// and pipeline the sends — rather than paying one serial round-trip per
/// message — while capping in-flight memory (each task holds its serialized
/// Arrow IPC payload until the broker acknowledges it).
const PRODUCE_MAX_IN_FLIGHT: usize = 256;

/// Repository for Iceberg tables with eventual consistency guarantees.
///
/// Writes go through: SDK → Arrow IPC → Kafka → `horizon-iceberg-sink` → Iceberg.
///
/// The repository is configured for a single target table at construction
/// time (field ID map + partition spec + Kafka topic). Individual write
/// methods accept any `T: ArrowSerializable`, so callers must pass a model
/// whose Arrow schema matches the configured table.
pub struct IcebergRepository {
    /// Iceberg column-to-field-ID mapping for Arrow schema annotation.
    field_id_map: FieldIdMap,
    /// Partition specification for computing Kafka partition keys.
    partition_spec: PartitionSpec,
    /// Kafka producer for the target topic.
    producer: KafkaProducer,
}

impl IcebergRepository {
    /// Access the underlying field ID mapping.
    #[must_use]
    pub const fn field_id_map(&self) -> &FieldIdMap {
        &self.field_id_map
    }

    /// Flush all pending messages, waiting up to `timeout`.
    ///
    /// On success, returns the number of messages still in flight after
    /// the timeout; 0 means all messages produced before this call are
    /// durable in Kafka (but NOT yet visible in Iceberg).
    ///
    /// # Errors
    ///
    /// Returns `KafkaError::Flush` if the underlying producer flush
    /// fails; the in-flight count is not reported in that case.
    pub fn flush(&self, timeout: Duration) -> Result<i32> {
        self.producer.flush(timeout)
    }

    /// Group models by their computed partition key, preserving input order
    /// within each group.
    fn group_by_partition<'models, T, F>(
        &self,
        models: &'models [T],
        partition_values_fn: F,
    ) -> Result<BTreeMap<Vec<u8>, Vec<&'models T>>>
    where
        F: Fn(&PartitionSpec, &T) -> Result<Vec<Option<PartitionValue>>>,
    {
        let mut groups: BTreeMap<Vec<u8>, Vec<&T>> = BTreeMap::new();
        for model in models {
            let values = partition_values_fn(&self.partition_spec, model)?;
            let key = self.partition_spec.compute_key(&values)?;
            groups.entry(key).or_default().push(model);
        }
        Ok(groups)
    }

    /// Insert a single record to Iceberg asynchronously.
    ///
    /// The record is serialized to Arrow IPC and sent to Kafka. This method
    /// returns after Kafka acknowledgement — the data is NOT yet visible in Iceberg.
    ///
    /// # Errors
    ///
    /// Returns `IcebergError` on serialization or partition-key computation
    /// failure, or `KafkaError::Delivery` if the message cannot be delivered.
    pub async fn insert<T, F>(&self, model: &T, partition_values_fn: F) -> Result<()>
    where
        T: ArrowSerializable + Send + Sync,
        F: Fn(&PartitionSpec, &T) -> Result<Vec<Option<PartitionValue>>>,
    {
        let batch = model.to_record_batch(&self.field_id_map)?;
        let payload = serialize_to_arrow_ipc(&[batch])?;
        let values = partition_values_fn(&self.partition_spec, model)?;
        let key = self.partition_spec.compute_key(&values)?;

        self.producer.produce(&key, &payload).await
    }

    /// Insert a batch of records to Iceberg asynchronously.
    ///
    /// Groups records by partition key, then serializes and produces up to
    /// [`PRODUCE_MAX_IN_FLIGHT`] groups per window so librdkafka can pipeline
    /// the sends. `acks=all` still holds: this resolves only once every message
    /// is acknowledged by the broker.
    ///
    /// # Errors
    ///
    /// Returns `IcebergError` on serialization or partition-key computation
    /// failure, or `KafkaError::Delivery` if a message cannot be delivered.
    pub async fn insert_batch<T, F>(&self, models: &[T], partition_values_fn: F) -> Result<()>
    where
        T: ArrowSerializable + Send + Sync,
        F: Fn(&PartitionSpec, &T) -> Result<Vec<Option<PartitionValue>>>,
    {
        if models.is_empty() {
            return Ok(());
        }

        let group_vec: Vec<(Vec<u8>, Vec<&T>)> = self
            .group_by_partition(models, partition_values_fn)?
            .into_iter()
            .collect();

        // Produce in windows of PRODUCE_MAX_IN_FLIGHT groups. Each window is
        // serialized to owned bytes up front, then its messages are produced
        // concurrently so librdkafka can batch and pipeline the sends. Serializing
        // before the produce futures is what keeps this future `Send`: a `&T` borrow
        // captured by a future (or threaded through a stream-combinator closure)
        // trips a higher-ranked `Send` inference limitation ("FnOnce is not general
        // enough") that would make this method unusable behind a `Send` async trait.
        for window in group_vec.chunks(PRODUCE_MAX_IN_FLIGHT) {
            let mut payloads: Vec<(Vec<u8>, Vec<u8>)> = Vec::with_capacity(window.len());
            for (key, group_model_vec) in window {
                payloads.push((key.clone(), self.serialize_group(group_model_vec)?));
            }
            try_join_all(
                payloads
                    .iter()
                    .map(|(key, payload)| self.producer.produce(key, payload)),
            )
            .await?;
        }

        Ok(())
    }

    /// Create a new Iceberg repository.
    ///
    /// The `producer` must be configured for the correct topic
    /// (e.g. `horizon.data_row` or `horizon.metadata_row`).
    pub fn new(
        producer: KafkaProducer,
        field_id_map: FieldIdMap,
        partition: ValidatedPartition,
    ) -> Self {
        info!(
            topic = producer.topic(),
            table = field_id_map.table_name(),
            "initialized IcebergRepository"
        );
        Self {
            field_id_map,
            partition_spec: partition.into_inner(),
            producer,
        }
    }

    /// Access the partition specification.
    #[must_use]
    pub const fn partition_spec(&self) -> &PartitionSpec {
        &self.partition_spec
    }
    /// Serialize one partition group into a single Arrow IPC payload.
    fn serialize_group<T: ArrowSerializable>(&self, group_models: &[&T]) -> Result<Vec<u8>> {
        let batches = group_models
            .iter()
            .map(|m| m.to_record_batch(&self.field_id_map))
            .collect::<Result<Vec<_>>>()?;
        let combined = concat_batches(&batches)?;
        serialize_to_arrow_ipc(&[combined])
    }
}

impl fmt::Debug for IcebergRepository {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        f.debug_struct("IcebergRepository")
            .field("topic", &self.producer.topic())
            .field("table", &self.field_id_map.table_name())
            .finish_non_exhaustive()
    }
}