#![allow(
clippy::multiple_inherent_impl,
reason = "builder methods are grouped by subsystem for readability"
)]
pub mod bench_support;
pub mod coverage;
pub mod iceberg;
pub mod kafka;
pub mod postgres;
pub mod types;
use std::fmt;
use tracing::info;
use crate::coverage::client::CoverageClient;
use crate::iceberg::discovery::discover_table_metadata;
use crate::iceberg::partition::PartitionValue;
use crate::iceberg::repository::IcebergRepository;
use crate::kafka::producer::KafkaProducer;
use crate::postgres::pool::create_pool;
use crate::postgres::repository::PostgresRepository;
use crate::types::config::{CoverageConfig, IcebergConfig, KafkaConfig, PostgresConfig, SdkConfig};
use crate::types::error::{ConfigError, HorizonError, Result};
use crate::types::model::{
AudioSpecification, BeamgramSpecification, BearingTimeRecordSpecification, DataRow, DataStream,
DirectionalSpectrogramSpecification, MetadataRow, Mission, MissionPlatform, Ontology,
OntologyClass, Platform, PlatformAudioSpecification, PlatformBeamgramSpecification,
PlatformBearingTimeRecordSpecification, PlatformInformation, PlatformKind,
PlatformSpectrogramSpecification, SpectrogramSpecification,
};
pub struct HorizonSdk {
possible_coverage_client: Option<CoverageClient>,
possible_data_row_iceberg: Option<IcebergRepository>,
possible_metadata_row_iceberg: Option<IcebergRepository>,
postgres_repo: PostgresRepository,
}
impl HorizonSdk {
#[must_use]
pub fn builder() -> HorizonSdkBuilder {
HorizonSdkBuilder::default()
}
#[must_use]
pub const fn coverage(&self) -> Option<&CoverageClient> {
self.possible_coverage_client.as_ref()
}
pub async fn create_audio_specification_batch(
&self,
specs: &[AudioSpecification],
) -> Result<Vec<AudioSpecification>> {
self.postgres_repo
.insert_audio_specification_batch(specs)
.await
}
pub async fn create_beamgram_specification_batch(
&self,
specs: &[BeamgramSpecification],
) -> Result<Vec<BeamgramSpecification>> {
self.postgres_repo
.insert_beamgram_specification_batch(specs)
.await
}
pub async fn create_bearing_time_record_specification_batch(
&self,
specs: &[BearingTimeRecordSpecification],
) -> Result<Vec<BearingTimeRecordSpecification>> {
self.postgres_repo
.insert_bearing_time_record_specification_batch(specs)
.await
}
pub async fn create_data_row(&self, row: &DataRow) -> Result<()> {
self.postgres_repo.insert_data_row(row).await?;
if let Some(iceberg_repo) = &self.possible_data_row_iceberg {
iceberg_repo
.insert(row, &data_row_partition_values(row))
.await
.map_err(|source| HorizonError::PartialDualWrite {
source: Box::new(source),
})?;
}
Ok(())
}
pub async fn create_data_row_batch(&self, rows: &[DataRow]) -> Result<()> {
self.postgres_repo.insert_data_row_batch(rows).await?;
if let Some(iceberg_repo) = &self.possible_data_row_iceberg {
iceberg_repo
.insert_batch(rows, data_row_partition_values)
.await
.map_err(|source| HorizonError::PartialDualWrite {
source: Box::new(source),
})?;
}
Ok(())
}
pub async fn create_data_row_batch_iceberg_only(&self, rows: &[DataRow]) -> Result<()> {
let iceberg_repo =
self.possible_data_row_iceberg
.as_ref()
.ok_or(ConfigError::MissingSubsystem {
subsystem: "Iceberg",
hint: "call with_iceberg() and with_kafka() on the builder",
})?;
iceberg_repo
.insert_batch(rows, data_row_partition_values)
.await
}
pub async fn create_data_row_iceberg_only(&self, row: &DataRow) -> Result<()> {
let iceberg_repo =
self.possible_data_row_iceberg
.as_ref()
.ok_or(ConfigError::MissingSubsystem {
subsystem: "Iceberg",
hint: "call with_iceberg() and with_kafka() on the builder",
})?;
iceberg_repo
.insert(row, &data_row_partition_values(row))
.await
}
pub async fn create_data_stream_batch(
&self,
data_streams: &[DataStream],
) -> Result<Vec<DataStream>> {
self.postgres_repo
.insert_data_stream_batch(data_streams)
.await
}
pub async fn create_directional_spectrogram_specification_batch(
&self,
specs: &[DirectionalSpectrogramSpecification],
) -> Result<Vec<DirectionalSpectrogramSpecification>> {
self.postgres_repo
.insert_directional_spectrogram_specification_batch(specs)
.await
}
pub async fn create_metadata_row(&self, row: &MetadataRow) -> Result<()> {
self.postgres_repo.insert_metadata_row(row).await?;
if let Some(iceberg_repo) = &self.possible_metadata_row_iceberg {
iceberg_repo
.insert(row, &metadata_row_partition_values(row))
.await
.map_err(|source| HorizonError::PartialDualWrite {
source: Box::new(source),
})?;
}
Ok(())
}
pub async fn create_metadata_row_batch(&self, rows: &[MetadataRow]) -> Result<()> {
self.postgres_repo.insert_metadata_row_batch(rows).await?;
if let Some(iceberg_repo) = &self.possible_metadata_row_iceberg {
iceberg_repo
.insert_batch(rows, metadata_row_partition_values)
.await
.map_err(|source| HorizonError::PartialDualWrite {
source: Box::new(source),
})?;
}
Ok(())
}
pub async fn create_mission_batch(&self, missions: &[Mission]) -> Result<Vec<Mission>> {
self.postgres_repo.insert_mission_batch(missions).await
}
pub async fn create_mission_platform_batch(
&self,
mission_platforms: &[MissionPlatform],
) -> Result<Vec<MissionPlatform>> {
self.postgres_repo
.insert_mission_platform_batch(mission_platforms)
.await
}
pub async fn create_ontology_batch(&self, ontologies: &[Ontology]) -> Result<Vec<Ontology>> {
self.postgres_repo.insert_ontology_batch(ontologies).await
}
pub async fn create_ontology_class_batch(
&self,
ontology_classes: &[OntologyClass],
) -> Result<Vec<OntologyClass>> {
self.postgres_repo
.insert_ontology_class_batch(ontology_classes)
.await
}
pub async fn create_platform_audio_specification_batch(
&self,
links: &[PlatformAudioSpecification],
) -> Result<Vec<PlatformAudioSpecification>> {
self.postgres_repo
.insert_platform_audio_specification_batch(links)
.await
}
pub async fn create_platform_batch(&self, platforms: &[Platform]) -> Result<Vec<Platform>> {
self.postgres_repo.insert_platform_batch(platforms).await
}
pub async fn create_platform_beamgram_specification_batch(
&self,
links: &[PlatformBeamgramSpecification],
) -> Result<Vec<PlatformBeamgramSpecification>> {
self.postgres_repo
.insert_platform_beamgram_specification_batch(links)
.await
}
pub async fn create_platform_bearing_time_record_specification_batch(
&self,
links: &[PlatformBearingTimeRecordSpecification],
) -> Result<Vec<PlatformBearingTimeRecordSpecification>> {
self.postgres_repo
.insert_platform_bearing_time_record_specification_batch(links)
.await
}
pub async fn create_platform_information_batch(
&self,
records: &[PlatformInformation],
) -> Result<Vec<PlatformInformation>> {
self.postgres_repo
.insert_platform_information_batch(records)
.await
}
pub async fn create_platform_kind_batch(
&self,
platform_kinds: &[PlatformKind],
) -> Result<Vec<PlatformKind>> {
self.postgres_repo
.insert_platform_kind_batch(platform_kinds)
.await
}
pub async fn create_platform_spectrogram_specification_batch(
&self,
links: &[PlatformSpectrogramSpecification],
) -> Result<Vec<PlatformSpectrogramSpecification>> {
self.postgres_repo
.insert_platform_spectrogram_specification_batch(links)
.await
}
pub async fn create_spectrogram_specification_batch(
&self,
specs: &[SpectrogramSpecification],
) -> Result<Vec<SpectrogramSpecification>> {
self.postgres_repo
.insert_spectrogram_specification_batch(specs)
.await
}
#[must_use]
pub const fn has_iceberg(&self) -> bool {
self.possible_data_row_iceberg.is_some()
}
#[must_use]
pub const fn postgres(&self) -> &PostgresRepository {
&self.postgres_repo
}
}
impl fmt::Debug for HorizonSdk {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("HorizonSdk")
.field("coverage_enabled", &self.possible_coverage_client.is_some())
.field("iceberg_enabled", &self.has_iceberg())
.field("postgres", &self.postgres_repo)
.finish_non_exhaustive()
}
}
#[derive(Default)]
pub struct HorizonSdkBuilder {
coverage_config: Option<CoverageConfig>,
iceberg_config: Option<IcebergConfig>,
kafka_config: Option<KafkaConfig>,
organization_id: Option<uuid::Uuid>,
postgres_config: Option<PostgresConfig>,
postgres_url: Option<String>,
}
impl HorizonSdkBuilder {
pub async fn build(self) -> Result<HorizonSdk> {
if self.iceberg_config.is_some() && self.kafka_config.is_none() {
return Err(ConfigError::MissingSubsystem {
subsystem: "Kafka",
hint: "required when Iceberg is configured for the write path",
}
.into());
}
let pg_config = self.resolve_postgres_config()?;
let pool = create_pool(&pg_config).await?;
let org_id = self.organization_id.or(pg_config.possible_organization_id);
let postgres_repo = PostgresRepository::new(pool, org_id);
info!("initialized Postgres repository");
let (possible_data_row_iceberg, possible_metadata_row_iceberg) =
if let Some(iceberg_config) = &self.iceberg_config {
let kafka_config = self.kafka_config.as_ref().ok_or({
ConfigError::MissingSubsystem {
subsystem: "Kafka",
hint: "required when Iceberg is configured for the write path",
}
})?;
let (dr_field_ids, dr_partition) =
discover_table_metadata(iceberg_config, &iceberg_config.data_row_table).await?;
let dr_producer = KafkaProducer::new(kafka_config, &kafka_config.data_row_topic)?;
let dr_repo = IcebergRepository::new(dr_producer, dr_field_ids, dr_partition);
let (mr_field_ids, mr_partition) =
discover_table_metadata(iceberg_config, &iceberg_config.metadata_row_table)
.await?;
let mr_producer =
KafkaProducer::new(kafka_config, &kafka_config.metadata_row_topic)?;
let mr_repo = IcebergRepository::new(mr_producer, mr_field_ids, mr_partition);
(Some(dr_repo), Some(mr_repo))
} else {
(None, None)
};
let possible_coverage_client = self
.coverage_config
.as_ref()
.map(CoverageClient::new)
.transpose()?;
info!(
dual_write = possible_data_row_iceberg.is_some(),
coverage = possible_coverage_client.is_some(),
"HorizonSdk initialized"
);
Ok(HorizonSdk {
possible_coverage_client,
possible_data_row_iceberg,
possible_metadata_row_iceberg,
postgres_repo,
})
}
fn resolve_postgres_config(&self) -> Result<PostgresConfig> {
if let Some(config) = &self.postgres_config {
return Ok(config.clone());
}
if let Some(url) = &self.postgres_url {
let parsed = url::Url::parse(url).map_err(|source| ConfigError::InvalidUrl {
field: "Postgres",
url: url.clone(),
source,
})?;
let host = parsed.host_str().unwrap_or("localhost").to_owned();
let port = parsed.port().unwrap_or(5432);
let user = if parsed.username().is_empty() {
"postgres".to_owned()
} else {
parsed.username().to_owned()
};
let password = parsed.password().unwrap_or_default().to_owned();
let database = parsed.path().trim_start_matches('/').to_owned();
let sslmode = parsed
.query_pairs()
.find(|(k, _)| k == "sslmode")
.map_or_else(|| "require".to_owned(), |(_, v)| v.into_owned());
return Ok(PostgresConfig {
database,
host,
password,
port,
possible_organization_id: None,
role: None,
sslmode,
user,
});
}
Err(ConfigError::MissingSubsystem {
subsystem: "Postgres",
hint: "call with_postgres_config() or with_postgres_url()",
}
.into())
}
#[must_use]
pub fn with_config(mut self, config: SdkConfig) -> Self {
self.coverage_config = config.coverage;
self.iceberg_config = config.iceberg;
self.kafka_config = config.kafka;
self.postgres_config = Some(config.postgres);
self
}
#[must_use]
pub fn with_coverage(mut self, config: CoverageConfig) -> Self {
self.coverage_config = Some(config);
self
}
#[must_use]
pub fn with_iceberg(mut self, config: IcebergConfig) -> Self {
self.iceberg_config = Some(config);
self
}
#[must_use]
pub fn with_kafka(mut self, config: KafkaConfig) -> Self {
self.kafka_config = Some(config);
self
}
#[must_use]
pub const fn with_organization_id(mut self, id: uuid::Uuid) -> Self {
self.organization_id = Some(id);
self
}
#[must_use]
pub fn with_postgres_config(mut self, config: PostgresConfig) -> Self {
self.postgres_config = Some(config);
self
}
#[must_use]
pub fn with_postgres_url(mut self, url: &str) -> Self {
self.postgres_url = Some(url.to_owned());
self
}
}
pub fn load_config(config_path: Option<&str>) -> Result<SdkConfig> {
let mut builder = config::Config::builder();
if let Some(path) = config_path {
builder = builder.add_source(config::File::with_name(path));
}
builder = builder.add_source(
config::Environment::with_prefix("HORIZON")
.separator("__")
.try_parsing(true),
);
builder
.build()
.and_then(config::Config::try_deserialize)
.map_err(|source| ConfigError::from(source).into())
}
fn data_row_partition_values(row: &DataRow) -> Vec<Option<PartitionValue>> {
vec![
Some(PartitionValue::from(row.data_stream_id)),
Some(PartitionValue::from(
row.datetime.format("%Y-%m-%d").to_string(),
)),
Some(PartitionValue::from(row.data_type.as_str())),
Some(PartitionValue::from(row.specification_id)),
]
}
fn metadata_row_partition_values(row: &MetadataRow) -> Vec<Option<PartitionValue>> {
vec![
Some(PartitionValue::from(row.data_stream_id)),
Some(PartitionValue::from(
row.datetime.format("%Y-%m-%d").to_string(),
)),
]
}
#[cfg(test)]
#[allow(
clippy::absolute_paths,
clippy::assertions_on_result_states,
clippy::panic_in_result_fn,
clippy::single_call_fn,
clippy::unwrap_used,
reason = "unit tests use assertions, unwrap, and test helpers freely"
)]
mod tests {
use super::*;
fn parse_postgres_url(url: &str) -> Result<PostgresConfig> {
HorizonSdkBuilder::default()
.with_postgres_url(url)
.resolve_postgres_config()
}
#[test]
fn builder_organization_id_propagates() {
let id = uuid::Uuid::new_v4();
let builder = HorizonSdk::builder().with_organization_id(id);
assert_eq!(builder.organization_id, Some(id));
}
#[test]
fn builder_requires_postgres() {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
let result = rt.block_on(HorizonSdk::builder().build());
assert!(result.is_err());
let err = result.unwrap_err().to_string();
assert!(err.contains("Postgres configuration is required"));
}
#[test]
fn builder_with_config_sets_all_subsystems() {
let config = SdkConfig {
coverage: Some(CoverageConfig {
base_url: "http://localhost:3000".to_owned(),
timeout_secs: 30,
}),
iceberg: None,
kafka: Some(KafkaConfig {
bootstrap_servers: "localhost:9092".to_owned(),
compression_type: "lz4".to_owned(),
data_row_topic: "horizon.data_row".to_owned(),
metadata_row_topic: "horizon.metadata_row".to_owned(),
sasl_mechanism: "PLAIN".to_owned(),
sasl_password: None,
sasl_username: None,
security_protocol: "PLAINTEXT".to_owned(),
}),
postgres: PostgresConfig {
database: "horizon".to_owned(),
host: "localhost".to_owned(),
password: "test".to_owned(),
port: 5432,
possible_organization_id: None,
role: None,
sslmode: "disable".to_owned(),
user: "test".to_owned(),
},
};
let builder = HorizonSdk::builder().with_config(config);
assert!(builder.postgres_config.is_some());
assert!(builder.kafka_config.is_some());
assert!(builder.iceberg_config.is_none());
assert!(builder.coverage_config.is_some());
}
#[test]
fn parse_postgres_url_defaults() {
let config = parse_postgres_url("postgresql://localhost/horizon").unwrap();
assert_eq!(config.host, "localhost");
assert_eq!(config.port, 5432);
assert_eq!(config.user, "postgres");
assert_eq!(config.password, "");
assert_eq!(config.database, "horizon");
assert_eq!(config.sslmode, "require");
}
#[test]
fn parse_postgres_url_full() {
let config =
parse_postgres_url("postgresql://user:pass@db.example.com:5433/mydb?sslmode=disable")
.unwrap();
assert_eq!(config.host, "db.example.com");
assert_eq!(config.port, 5433);
assert_eq!(config.user, "user");
assert_eq!(config.password, "pass");
assert_eq!(config.database, "mydb");
assert_eq!(config.sslmode, "disable");
}
#[test]
fn parse_postgres_url_invalid() {
let result = parse_postgres_url("not a url");
assert!(result.is_err());
let err = result.unwrap_err().to_string();
assert!(err.contains("invalid Postgres URL"));
}
}