horizon-sdk 6.7.2

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

use crate::types::error::{HorizonError, PostgresError, Result};
use crate::types::model::BeamgramSpecification;
use sqlx::{Postgres, QueryBuilder};

use super::PostgresRepository;

impl PostgresRepository {
    /// Batch-insert beamgram specs (multi-row INSERT) and return their
    /// server-generated ids.
    pub async fn bulk_insert_beamgram_specification(
        &self,
        specs: &[BeamgramSpecification],
    ) -> Result<Vec<Uuid>> {
        const COLUMN_COUNT: usize = 14;
        #[allow(
            clippy::decimal_literal_representation,
            reason = "Integer representation is more readable for this applications"
        )]
        const MAX_PARAMETER_COUNT: usize = 65_535;
        #[allow(
            clippy::integer_division,
            clippy::integer_division_remainder_used,
            reason = "Can only perform whole writes at once"
        )]
        let max_writes = MAX_PARAMETER_COUNT / COLUMN_COUNT;

        let now = Utc::now();

        let mut id_vector = Vec::with_capacity(specs.len());
        for chunk in specs.chunks(max_writes) {
            let mut query_builder: QueryBuilder<Postgres> = QueryBuilder::new(
                "INSERT INTO horizon_public.beamgram_specification \
             (id, center_bearing, center_bin_width, min_frequency, max_frequency, \
              fft_sample_count, organization_id, update_rate, normalizer, name, \
              elevation_increment, lower_elevation, upper_elevation, modified_datetime) ",
            );
            query_builder.push_values(chunk, |mut b, spec| {
                b.push("COALESCE(")
                    .push_bind_unseparated(spec.id)
                    .push_unseparated(", gen_random_uuid())")
                    .push_bind(spec.center_bearing)
                    .push_bind(spec.center_bin_width)
                    .push_bind(spec.min_frequency)
                    .push_bind(spec.max_frequency)
                    .push_bind(spec.fft_sample_count)
                    .push_bind(self.organization_id.or(spec.organization_id))
                    .push_bind(spec.update_rate_ms)
                    .push_unseparated("::bigint * interval '1 millisecond'")
                    .push_bind(spec.normalizer.clone())
                    .push_bind(spec.name.clone())
                    .push_bind(spec.elevation_increment)
                    .push_bind(spec.lower_elevation)
                    .push_bind(spec.upper_elevation)
                    .push_bind(now);
            });
            query_builder.push(" RETURNING id");
            let id_batch: Vec<Uuid> = query_builder
                .build_query_scalar()
                .fetch_all(&self.pool)
                .await?;
            id_vector.extend(id_batch);
        }
        Ok(id_vector)
    }

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

    /// Create a new beamgram specification.
    pub async fn insert_beamgram_specification(
        &self,
        spec: &BeamgramSpecification,
    ) -> Result<BeamgramSpecification> {
        let now = Utc::now();
        let org = self.organization_id.or(spec.organization_id);
        Ok(sqlx::query_as!(
            BeamgramSpecification,
            r#"
            INSERT INTO horizon_public.beamgram_specification
                (id, center_bearing, center_bin_width, elevation_increment,
                 lower_elevation, upper_elevation, max_frequency,
                 min_frequency, name, fft_sample_count, organization_id,
                 update_rate, normalizer, modified_datetime)
            VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8, $9, $10, $11,
                    $12::bigint * interval '1 millisecond', $13, $14)
            RETURNING id, created_datetime, modified_datetime, center_bearing,
                      center_bin_width, elevation_increment, lower_elevation,
                      upper_elevation, max_frequency, min_frequency,
                      name, fft_sample_count, organization_id,
                      (EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
                      normalizer
            "#,
            spec.id,
            spec.center_bearing,
            spec.center_bin_width,
            spec.elevation_increment,
            spec.lower_elevation,
            spec.upper_elevation,
            spec.max_frequency,
            spec.min_frequency,
            spec.name,
            spec.fft_sample_count,
            org,
            spec.update_rate_ms,
            spec.normalizer,
            now,
        )
        .fetch_one(&self.pool)
        .await?)
    }

    /// List all beamgram specification records.
    pub async fn list_beamgram_specifications(&self) -> Result<Vec<BeamgramSpecification>> {
        Ok(sqlx::query_as!(
            BeamgramSpecification,
            r#"
            SELECT id, created_datetime, modified_datetime, center_bearing,
                   center_bin_width, elevation_increment, lower_elevation,
                   upper_elevation, max_frequency, min_frequency,
                   name, fft_sample_count, organization_id,
                   (EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
                   normalizer
            FROM horizon_public.beamgram_specification
            "#,
        )
        .fetch_all(&self.pool)
        .await?)
    }

    /// Read a beamgram specification by ID.
    pub async fn read_beamgram_specification(
        &self,
        id: Uuid,
    ) -> Result<Option<BeamgramSpecification>> {
        Ok(sqlx::query_as!(
            BeamgramSpecification,
            r#"
            SELECT id, created_datetime, modified_datetime, center_bearing,
                   center_bin_width, elevation_increment, lower_elevation,
                   upper_elevation, max_frequency, min_frequency,
                   name, fft_sample_count, organization_id,
                   (EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
                   normalizer
            FROM horizon_public.beamgram_specification WHERE id = $1
            "#,
            id,
        )
        .fetch_optional(&self.pool)
        .await?)
    }

    /// Update an existing beamgram specification by ID.
    ///
    /// Returns `HorizonError::Postgres(PostgresError::NotFound)` if no row
    /// with the given ID exists. Matches Python `BaseRepository.update()`
    /// semantics. `update_rate_ms` is converted to a Postgres interval on
    /// write and back to milliseconds on read.
    pub async fn update_beamgram_specification(
        &self,
        spec: &BeamgramSpecification,
    ) -> Result<BeamgramSpecification> {
        let now = Utc::now();
        let org = self.organization_id.or(spec.organization_id);
        let result = sqlx::query_as!(
            BeamgramSpecification,
            r#"
            UPDATE horizon_public.beamgram_specification
            SET
                center_bearing = $2,
                center_bin_width = $3,
                elevation_increment = $4,
                lower_elevation = $5,
                upper_elevation = $6,
                max_frequency = $7,
                min_frequency = $8,
                name = $9,
                fft_sample_count = $10,
                organization_id = $11,
                update_rate = $12::bigint * interval '1 millisecond',
                normalizer = $13,
                modified_datetime = $14
            WHERE id = $1
            RETURNING id, created_datetime, modified_datetime, center_bearing,
                      center_bin_width, elevation_increment, lower_elevation,
                      upper_elevation, max_frequency, min_frequency,
                      name, fft_sample_count, organization_id,
                      (EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
                      normalizer
            "#,
            spec.id,
            spec.center_bearing,
            spec.center_bin_width,
            spec.elevation_increment,
            spec.lower_elevation,
            spec.upper_elevation,
            spec.max_frequency,
            spec.min_frequency,
            spec.name,
            spec.fft_sample_count,
            org,
            spec.update_rate_ms,
            spec.normalizer,
            now,
        )
        .fetch_optional(&self.pool)
        .await?;

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

    /// Insert a beamgram specification or update it if a row with the same
    /// ID already exists. `update_rate_ms` is converted to a Postgres
    /// interval on write and back to milliseconds on read.
    pub async fn upsert_beamgram_specification(
        &self,
        spec: &BeamgramSpecification,
    ) -> Result<BeamgramSpecification> {
        let now = Utc::now();
        let org = self.organization_id.or(spec.organization_id);
        Ok(sqlx::query_as!(
            BeamgramSpecification,
            r#"
            INSERT INTO horizon_public.beamgram_specification
                (id, center_bearing, center_bin_width, elevation_increment,
                 lower_elevation, upper_elevation, max_frequency,
                 min_frequency, name, fft_sample_count, organization_id,
                 update_rate, normalizer, modified_datetime)
            VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8, $9, $10, $11,
                    $12::bigint * interval '1 millisecond', $13, $14)
            ON CONFLICT (id) DO UPDATE SET
                center_bearing = EXCLUDED.center_bearing,
                center_bin_width = EXCLUDED.center_bin_width,
                elevation_increment = EXCLUDED.elevation_increment,
                lower_elevation = EXCLUDED.lower_elevation,
                upper_elevation = EXCLUDED.upper_elevation,
                max_frequency = EXCLUDED.max_frequency,
                min_frequency = EXCLUDED.min_frequency,
                name = EXCLUDED.name, fft_sample_count = EXCLUDED.fft_sample_count,
                organization_id = EXCLUDED.organization_id,
                update_rate = EXCLUDED.update_rate,
                normalizer = EXCLUDED.normalizer,
                modified_datetime = EXCLUDED.modified_datetime
            RETURNING id, created_datetime, modified_datetime, center_bearing,
                      center_bin_width, elevation_increment, lower_elevation,
                      upper_elevation, max_frequency, min_frequency,
                      name, fft_sample_count, organization_id,
                      (EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
                      normalizer
            "#,
            spec.id,
            spec.center_bearing,
            spec.center_bin_width,
            spec.elevation_increment,
            spec.lower_elevation,
            spec.upper_elevation,
            spec.max_frequency,
            spec.min_frequency,
            spec.name,
            spec.fft_sample_count,
            org,
            spec.update_rate_ms,
            spec.normalizer,
            now,
        )
        .fetch_one(&self.pool)
        .await?)
    }
}