horizon-sdk 16.0.0

Canonical Rust data access layer for the Horizon platform
use chrono::Utc;
use tracing::debug;

use crate::types::error::Result;
use crate::types::model::Annotation;

use super::PostgresRepository;

impl PostgresRepository {
    /// Ingest a single annotation (label) for a platform.
    ///
    /// User attribution columns (`created_by_user_id`, `modified_by_user_id`)
    /// are populated by the `trigger__update_user_id` trigger from the active
    /// Postgres session, so they are not exposed in the SDK model.
    pub async fn insert_annotation(&self, annotation: &Annotation) -> Result<Annotation> {
        let data_stream_id = annotation.data_stream_id.ok_or_else(|| {
            sqlx::Error::Protocol("insert_annotation requires a data_stream_id".into())
        })?;
        let now = Utc::now();
        let org = self.organization_id.or(annotation.organization_id);
        let result = sqlx::query_as!(
            Annotation,
            r#"
            INSERT INTO horizon_public.annotation
                (id, platform_id, data_stream_id, time_coordinates, value_coordinates,
                 ontology_class_id, notes, confidence, organization_id,
                 duration_seconds, bearing_time_record_specification_id,
                 spectrogram_specification_id, parent_annotation_id,
                 feed_context, modified_datetime)
            VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8, $9,
                    $10, $11, $12, $13, $14, $15)
            RETURNING
                id, platform_id,
                data_stream_id,
                time_coordinates,
                value_coordinates as "value_coordinates!: Vec<f64>",
                created_datetime, modified_datetime, ontology_class_id, notes,
                confidence, organization_id, duration_seconds,
                bearing_time_record_specification_id,
                spectrogram_specification_id, parent_annotation_id,
                feed_context
            "#,
            annotation.id,
            annotation.platform_id,
            data_stream_id,
            &annotation.time_coordinates,
            &annotation.value_coordinates,
            annotation.ontology_class_id,
            annotation.notes,
            annotation.confidence,
            org,
            annotation.duration_seconds,
            annotation.bearing_time_record_specification_id,
            annotation.spectrogram_specification_id,
            annotation.parent_annotation_id,
            annotation.feed_context,
            now,
        )
        .fetch_one(&self.pool)
        .await?;
        debug!(
            annotation_id = ?result.id,
            platform_id = ?result.platform_id,
            data_stream_id = ?result.data_stream_id,
            "insert_annotation"
        );
        Ok(result)
    }

    /// Ingest a batch of annotations inside a single transaction.
    ///
    /// Returns the inserted rows in input order. Each row is inserted via
    /// its own statement against the shared transaction; if any insert
    /// fails the transaction is rolled back and no annotations from the
    /// batch are committed. Matches the atomicity guarantee of Python's
    /// `insert_batch` in `horizon-data-core`.
    #[allow(
        clippy::cognitive_complexity,
        reason = "long sqlx::query_as! parameter list inflates cognitive complexity"
    )]
    pub async fn insert_annotation_batch(
        &self,
        annotations: &[Annotation],
    ) -> Result<Vec<Annotation>> {
        if annotations.is_empty() {
            return Ok(Vec::new());
        }
        let now = Utc::now();
        let mut tx = self.pool.begin().await?;
        let mut inserted = Vec::with_capacity(annotations.len());
        for annotation in annotations {
            let data_stream_id = annotation.data_stream_id.ok_or_else(|| {
                sqlx::Error::Protocol("insert_annotation_batch requires a data_stream_id".into())
            })?;
            let org = self.organization_id.or(annotation.organization_id);
            let result = sqlx::query_as!(
                Annotation,
                r#"
                INSERT INTO horizon_public.annotation
                    (id, platform_id, data_stream_id, time_coordinates, value_coordinates,
                     ontology_class_id, notes, confidence, organization_id,
                     duration_seconds, bearing_time_record_specification_id,
                     spectrogram_specification_id, parent_annotation_id,
                     feed_context, modified_datetime)
                VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8, $9,
                        $10, $11, $12, $13, $14, $15)
                RETURNING
                    id, platform_id,
                    data_stream_id,
                    time_coordinates,
                    value_coordinates as "value_coordinates!: Vec<f64>",
                    created_datetime, modified_datetime, ontology_class_id, notes,
                    confidence, organization_id, duration_seconds,
                    bearing_time_record_specification_id,
                    spectrogram_specification_id, parent_annotation_id,
                    feed_context
                "#,
                annotation.id,
                annotation.platform_id,
                data_stream_id,
                &annotation.time_coordinates,
                &annotation.value_coordinates,
                annotation.ontology_class_id,
                annotation.notes,
                annotation.confidence,
                org,
                annotation.duration_seconds,
                annotation.bearing_time_record_specification_id,
                annotation.spectrogram_specification_id,
                annotation.parent_annotation_id,
                annotation.feed_context,
                now,
            )
            .fetch_one(&mut *tx)
            .await?;
            debug!(
                annotation_id = ?result.id,
                platform_id = ?result.platform_id,
                data_stream_id = ?result.data_stream_id,
                "insert_annotation_batch row"
            );
            inserted.push(result);
        }
        tx.commit().await?;
        Ok(inserted)
    }
}