horizon-sdk 6.9.0

Canonical Rust data access layer for the Horizon platform
Documentation
//! SDK configuration primitives.
//!
//! Configuration is layered: struct defaults -> env vars -> programmatic overrides.
//! The builder in `lib.rs` consumes these to construct subsystems.

use std::fmt;

use serde::Deserialize;

use super::constants;

/// Coverage `BentoML` service configuration.
#[derive(Debug, Clone, Deserialize)]
pub struct CoverageConfig {
    pub base_url: String,
    #[serde(default = "default_coverage_timeout_secs")]
    pub timeout_secs: u64,
}

/// Iceberg catalog configuration for schema discovery.
#[derive(Debug, Clone, Deserialize)]
pub struct IcebergConfig {
    #[serde(default = "default_catalog_uri")]
    pub catalog_uri: String,
    pub catalog_warehouse: Option<String>,
    /// `OAuth2` client credential (`client_id:client_secret`) for the REST catalog.
    #[serde(default)]
    pub credential: Option<String>,
    #[serde(default = "default_data_row_table")]
    pub data_row_table: String,
    #[serde(default = "default_metadata_row_table")]
    pub metadata_row_table: String,
    /// `OAuth2` token endpoint for the REST catalog.
    #[serde(default)]
    pub oauth2_server_uri: Option<String>,
    /// `OAuth2` scope requested when fetching a catalog token.
    #[serde(default)]
    pub scope: Option<String>,
}

/// Kafka producer configuration (SASL/SSL + topic routing).
#[derive(Clone, Deserialize)]
pub struct KafkaConfig {
    #[serde(default = "default_bootstrap_servers")]
    pub bootstrap_servers: String,
    #[serde(default = "default_compression_type")]
    pub compression_type: String,
    #[serde(default = "default_data_row_topic")]
    pub data_row_topic: String,
    #[serde(default = "default_metadata_row_topic")]
    pub metadata_row_topic: String,
    #[serde(default = "default_sasl_mechanism")]
    pub sasl_mechanism: String,
    pub sasl_password: Option<String>,
    pub sasl_username: Option<String>,
    #[serde(default = "default_security_protocol")]
    pub security_protocol: String,
}

impl fmt::Debug for KafkaConfig {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        f.debug_struct("KafkaConfig")
            .field("bootstrap_servers", &self.bootstrap_servers)
            .field("compression_type", &self.compression_type)
            .field("data_row_topic", &self.data_row_topic)
            .field("metadata_row_topic", &self.metadata_row_topic)
            .field("sasl_mechanism", &self.sasl_mechanism)
            .field(
                "sasl_password",
                &self.sasl_password.as_ref().map(|_| "[REDACTED]"),
            )
            .field(
                "sasl_username",
                &self.sasl_username.as_ref().map(|_| "[REDACTED]"),
            )
            .field("security_protocol", &self.security_protocol)
            .finish()
    }
}

/// `PostgreSQL` connection configuration.
#[derive(Clone, Deserialize)]
pub struct PostgresConfig {
    pub database: String,
    pub host: String,
    pub password: String,
    pub port: u16,
    #[serde(rename = "organization_id")]
    pub possible_organization_id: Option<uuid::Uuid>,
    pub role: Option<String>,
    #[serde(default = "default_sslmode")]
    pub sslmode: String,
    pub user: String,
}

impl fmt::Debug for PostgresConfig {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        f.debug_struct("PostgresConfig")
            .field("database", &self.database)
            .field("host", &self.host)
            .field("password", &"[REDACTED]")
            .field("port", &self.port)
            .field("possible_organization_id", &self.possible_organization_id)
            .field("role", &self.role)
            .field("sslmode", &self.sslmode)
            .field("user", &"[REDACTED]")
            .finish()
    }
}

/// Top-level SDK configuration.
#[derive(Debug, Clone, Deserialize)]
pub struct SdkConfig {
    pub coverage: Option<CoverageConfig>,
    pub iceberg: Option<IcebergConfig>,
    pub kafka: Option<KafkaConfig>,
    pub postgres: PostgresConfig,
}

/// Serde default for `KafkaConfig::bootstrap_servers`.
fn default_bootstrap_servers() -> String {
    constants::DEFAULT_BOOTSTRAP_SERVERS.to_owned()
}

/// Serde default for `IcebergConfig::catalog_uri`.
fn default_catalog_uri() -> String {
    constants::DEFAULT_CATALOG_URI.to_owned()
}

/// Serde default for `KafkaConfig::compression_type`.
fn default_compression_type() -> String {
    constants::DEFAULT_COMPRESSION_TYPE.to_owned()
}

/// Serde default for `CoverageConfig::timeout_secs`.
const fn default_coverage_timeout_secs() -> u64 {
    constants::DEFAULT_COVERAGE_TIMEOUT_SECS
}

/// Serde default for `IcebergConfig::data_row_table`.
fn default_data_row_table() -> String {
    constants::DEFAULT_DATA_ROW_TABLE.to_owned()
}

/// Serde default for `KafkaConfig::data_row_topic`.
fn default_data_row_topic() -> String {
    constants::DEFAULT_DATA_ROW_TOPIC.to_owned()
}

/// Serde default for `IcebergConfig::metadata_row_table`.
fn default_metadata_row_table() -> String {
    constants::DEFAULT_METADATA_ROW_TABLE.to_owned()
}

/// Serde default for `KafkaConfig::metadata_row_topic`.
fn default_metadata_row_topic() -> String {
    constants::DEFAULT_METADATA_ROW_TOPIC.to_owned()
}

/// Serde default for `KafkaConfig::sasl_mechanism`.
fn default_sasl_mechanism() -> String {
    constants::DEFAULT_SASL_MECHANISM.to_owned()
}

/// Serde default for `KafkaConfig::security_protocol`.
fn default_security_protocol() -> String {
    constants::DEFAULT_SECURITY_PROTOCOL.to_owned()
}

/// Serde default for `PostgresConfig::sslmode`.
fn default_sslmode() -> String {
    constants::DEFAULT_SSLMODE.to_owned()
}