horizon-sdk 16.0.0

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

use crate::types::error::{HorizonError, PostgresError, Result};
use crate::types::filter::SpectrogramSpecificationFilter;
use crate::types::model::SpectrogramSpecification;

use super::PostgresRepository;

impl PostgresRepository {
    /// Delete a spectrogram specification by ID. Returns the number of rows removed.
    pub async fn delete_spectrogram_specification(&self, id: Uuid) -> Result<u64> {
        let result = sqlx::query!(
            "DELETE FROM horizon_public.spectrogram_specification WHERE id = $1",
            id
        )
        .execute(&self.pool)
        .await?;
        Ok(result.rows_affected())
    }

    /// Create a new spectrogram specification.
    pub async fn insert_spectrogram_specification(
        &self,
        spec: &SpectrogramSpecification,
    ) -> Result<SpectrogramSpecification> {
        let now = Utc::now();
        let org = self.organization_id.or(spec.organization_id);
        Ok(sqlx::query_as!(
            SpectrogramSpecification,
            r#"
            INSERT INTO horizon_public.spectrogram_specification
                (id, amplitude_unit_mode, channel, frequency_spacing, name,
                 fft_sample_count, fft_sample_overlap_count,
                 organization_id, channel_role, frequency_bin_count, normalizer,
                 modified_datetime)
            VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)
            RETURNING id, amplitude_unit_mode, created_datetime, modified_datetime, channel,
                      frequency_spacing, name, fft_sample_count, fft_sample_overlap_count,
                      organization_id, channel_role, frequency_bin_count, normalizer
            "#,
            spec.id,
            spec.amplitude_unit_mode,
            spec.channel,
            spec.frequency_spacing,
            spec.name,
            spec.fft_sample_count,
            spec.fft_sample_overlap_count,
            org,
            spec.channel_role,
            spec.frequency_bin_count,
            spec.normalizer,
            now,
        )
        .fetch_one(&self.pool)
        .await?)
    }

    /// Create a batch of spectrogram specifications in a single multi-row INSERT.
    ///
    /// Per-column slices are passed to `PostgreSQL` `UNNEST(...)` so the
    /// entire batch becomes one SQL statement and one network round-trip.
    /// Missing IDs are filled by `gen_random_uuid()` server-side, and
    /// `modified_datetime` is set to the same `now()` value for every
    /// row.
    ///
    /// Returns the inserted rows in input order. Wraps the statement in
    /// a transaction so a per-row failure rolls back the entire batch.
    #[allow(
        clippy::as_conversions,
        reason = "sqlx::query_as! requires `&Vec<T> as &[T]` ascription to bind a nullable Postgres array; conversion direction is unambiguous (same element type)"
    )]
    pub async fn insert_spectrogram_specification_batch(
        &self,
        specs: &[SpectrogramSpecification],
    ) -> Result<Vec<SpectrogramSpecification>> {
        if specs.is_empty() {
            return Ok(Vec::new());
        }
        let now = Utc::now();
        let ids: Vec<Option<Uuid>> = specs.iter().map(|spec| spec.id).collect();
        let amplitude_unit_modes: Vec<Option<String>> = specs
            .iter()
            .map(|spec| spec.amplitude_unit_mode.clone())
            .collect();
        let channels: Vec<Option<i32>> = specs.iter().map(|spec| spec.channel).collect();
        let frequency_spacings: Vec<Option<f64>> =
            specs.iter().map(|spec| spec.frequency_spacing).collect();
        let names: Vec<Option<String>> = specs.iter().map(|spec| spec.name.clone()).collect();
        let fft_sample_counts: Vec<Option<i32>> =
            specs.iter().map(|spec| spec.fft_sample_count).collect();
        let fft_sample_overlap_counts: Vec<Option<i32>> = specs
            .iter()
            .map(|spec| spec.fft_sample_overlap_count)
            .collect();
        let orgs: Vec<Option<Uuid>> = specs
            .iter()
            .map(|spec| self.organization_id.or(spec.organization_id))
            .collect();
        let channel_roles: Vec<Option<String>> =
            specs.iter().map(|spec| spec.channel_role.clone()).collect();
        let frequency_bin_counts: Vec<Option<i32>> =
            specs.iter().map(|spec| spec.frequency_bin_count).collect();
        let normalizers: Vec<Option<String>> =
            specs.iter().map(|spec| spec.normalizer.clone()).collect();
        let mut tx = self.pool.begin().await?;
        let rows = sqlx::query_as!(
            SpectrogramSpecification,
            r#"
            WITH source AS (
                SELECT
                    COALESCE(id, gen_random_uuid()) AS id,
                    amplitude_unit_mode, channel, frequency_spacing, name,
                    fft_sample_count, fft_sample_overlap_count,
                    organization_id, channel_role, frequency_bin_count, normalizer, ord
                FROM UNNEST(
                    $1::uuid[], $2::text[], $3::integer[], $4::float8[], $5::text[],
                    $6::integer[], $7::integer[], $8::uuid[], $9::text[], $10::integer[],
                    $11::text[]
                ) WITH ORDINALITY
                    AS batch(id, amplitude_unit_mode, channel, frequency_spacing, name,
                             fft_sample_count, fft_sample_overlap_count,
                             organization_id, channel_role, frequency_bin_count, normalizer, ord)
            ),
            inserted AS (
                INSERT INTO horizon_public.spectrogram_specification
                    (id, amplitude_unit_mode, channel, frequency_spacing, name,
                     fft_sample_count, fft_sample_overlap_count,
                     organization_id, channel_role, frequency_bin_count, normalizer,
                     modified_datetime)
                SELECT
                    id, amplitude_unit_mode, channel, frequency_spacing, name,
                    fft_sample_count, fft_sample_overlap_count,
                    organization_id, channel_role, frequency_bin_count, normalizer, $12
                FROM source
                RETURNING id, amplitude_unit_mode, created_datetime, modified_datetime, channel,
                          frequency_spacing, name, fft_sample_count, fft_sample_overlap_count,
                          organization_id, channel_role, frequency_bin_count, normalizer
            )
            SELECT inserted.id, inserted.amplitude_unit_mode, inserted.created_datetime,
                   inserted.modified_datetime, inserted.channel, inserted.frequency_spacing,
                   inserted.name, inserted.fft_sample_count, inserted.fft_sample_overlap_count,
                   inserted.organization_id, inserted.channel_role,
                   inserted.frequency_bin_count, inserted.normalizer
            FROM inserted
            JOIN source ON source.id = inserted.id
            ORDER BY source.ord
            "#,
            &ids as &[Option<Uuid>],
            &amplitude_unit_modes as &[Option<String>],
            &channels as &[Option<i32>],
            &frequency_spacings as &[Option<f64>],
            &names as &[Option<String>],
            &fft_sample_counts as &[Option<i32>],
            &fft_sample_overlap_counts as &[Option<i32>],
            &orgs as &[Option<Uuid>],
            &channel_roles as &[Option<String>],
            &frequency_bin_counts as &[Option<i32>],
            &normalizers as &[Option<String>],
            now,
        )
        .fetch_all(&mut *tx)
        .await?;
        tx.commit().await?;
        debug!(
            batch_size = rows.len(),
            "insert_spectrogram_specification_batch"
        );
        Ok(rows)
    }

    /// List spectrogram specification records matching every set filter field.
    pub async fn list_spectrogram_specifications(
        &self,
        filter: &SpectrogramSpecificationFilter,
    ) -> Result<Vec<SpectrogramSpecification>> {
        Ok(sqlx::query_as!(
            SpectrogramSpecification,
            r#"
            SELECT id, amplitude_unit_mode, created_datetime, modified_datetime, channel,
                   frequency_spacing, name, fft_sample_count, fft_sample_overlap_count,
                   organization_id, channel_role, frequency_bin_count, normalizer
            FROM horizon_public.spectrogram_specification
            WHERE ($1::text IS NULL OR amplitude_unit_mode = $1)
              AND ($2::integer IS NULL OR channel = $2)
              AND ($3::text IS NULL OR channel_role = $3)
              AND ($4::integer IS NULL OR fft_sample_count = $4)
              AND ($5::integer IS NULL OR fft_sample_overlap_count = $5)
              AND ($6::float8 IS NULL OR frequency_spacing = $6)
              AND ($7::text IS NULL OR name = $7)
              AND ($8::uuid IS NULL OR organization_id = $8)
            "#,
            filter.amplitude_unit_mode,
            filter.channel,
            filter.channel_role,
            filter.fft_sample_count,
            filter.fft_sample_overlap_count,
            filter.frequency_spacing,
            filter.name,
            filter.organization_id,
        )
        .fetch_all(&self.pool)
        .await?)
    }

    /// Read a spectrogram specification by ID.
    pub async fn read_spectrogram_specification(
        &self,
        id: Uuid,
    ) -> Result<Option<SpectrogramSpecification>> {
        Ok(sqlx::query_as!(
            SpectrogramSpecification,
            r#"
            SELECT id, amplitude_unit_mode, created_datetime, modified_datetime, channel,
                   frequency_spacing, name, fft_sample_count, fft_sample_overlap_count,
                   organization_id, channel_role, frequency_bin_count, normalizer
            FROM horizon_public.spectrogram_specification WHERE id = $1
            "#,
            id,
        )
        .fetch_optional(&self.pool)
        .await?)
    }

    /// Update an existing spectrogram specification by ID.
    ///
    /// Returns `HorizonError::Postgres(PostgresError::NotFound)` if no row
    /// with the given ID exists. Matches Python `BaseRepository.update()`
    /// semantics.
    pub async fn update_spectrogram_specification(
        &self,
        spec: &SpectrogramSpecification,
    ) -> Result<SpectrogramSpecification> {
        let now = Utc::now();
        let org = self.organization_id.or(spec.organization_id);
        let result = sqlx::query_as!(
            SpectrogramSpecification,
            r#"
            UPDATE horizon_public.spectrogram_specification
            SET
                amplitude_unit_mode = $2,
                channel = $3,
                frequency_spacing = $4,
                name = $5,
                fft_sample_count = $6,
                fft_sample_overlap_count = $7,
                organization_id = $8,
                channel_role = $9,
                frequency_bin_count = $10,
                normalizer = $11,
                modified_datetime = $12
            WHERE id = $1
            RETURNING id, amplitude_unit_mode, created_datetime, modified_datetime, channel,
                      frequency_spacing, name, fft_sample_count, fft_sample_overlap_count,
                      organization_id, channel_role, frequency_bin_count, normalizer
            "#,
            spec.id,
            spec.amplitude_unit_mode,
            spec.channel,
            spec.frequency_spacing,
            spec.name,
            spec.fft_sample_count,
            spec.fft_sample_overlap_count,
            org,
            spec.channel_role,
            spec.frequency_bin_count,
            spec.normalizer,
            now,
        )
        .fetch_optional(&self.pool)
        .await?;

        result.ok_or_else(|| {
            HorizonError::Postgres(PostgresError::NotFound {
                entity: "spectrogram_specification".to_owned(),
                id: spec.id.into(),
            })
        })
    }

    /// Insert a spectrogram specification or update it if a row with the
    /// same ID already exists. Sets `modified_datetime` to now and falls
    /// back to the repository's `organization_id` when the model does not
    /// provide one.
    pub async fn upsert_spectrogram_specification(
        &self,
        spec: &SpectrogramSpecification,
    ) -> Result<SpectrogramSpecification> {
        let now = Utc::now();
        let org = self.organization_id.or(spec.organization_id);
        Ok(sqlx::query_as!(
            SpectrogramSpecification,
            r#"
            INSERT INTO horizon_public.spectrogram_specification
                (id, amplitude_unit_mode, channel, frequency_spacing, name,
                 fft_sample_count, fft_sample_overlap_count,
                 organization_id, channel_role, frequency_bin_count, normalizer,
                 modified_datetime)
            VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)
            ON CONFLICT (id) DO UPDATE SET
                amplitude_unit_mode = EXCLUDED.amplitude_unit_mode,
                channel = EXCLUDED.channel, frequency_spacing = EXCLUDED.frequency_spacing,
                name = EXCLUDED.name, fft_sample_count = EXCLUDED.fft_sample_count, fft_sample_overlap_count = EXCLUDED.fft_sample_overlap_count,
                organization_id = EXCLUDED.organization_id,
                channel_role = EXCLUDED.channel_role,
                frequency_bin_count = EXCLUDED.frequency_bin_count,
                normalizer = EXCLUDED.normalizer,
                modified_datetime = EXCLUDED.modified_datetime
            RETURNING id, amplitude_unit_mode, created_datetime, modified_datetime, channel,
                      frequency_spacing, name, fft_sample_count, fft_sample_overlap_count,
                      organization_id, channel_role, frequency_bin_count, normalizer
            "#,
            spec.id,
            spec.amplitude_unit_mode,
            spec.channel,
            spec.frequency_spacing,
            spec.name,
            spec.fft_sample_count,
            spec.fft_sample_overlap_count,
            org,
            spec.channel_role,
            spec.frequency_bin_count,
            spec.normalizer,
            now,
        )
        .fetch_one(&self.pool)
        .await?)
    }
}