horizon-sdk 16.0.0

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

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

use super::PostgresRepository;

impl PostgresRepository {
    /// If `created_datetime` is `None`, defaults to now (allows backfilling historical data).
    pub async fn insert_metadata_row(&self, metadata_row: &MetadataRow) -> Result<MetadataRow> {
        let now = Utc::now();
        Ok(sqlx::query_as!(
            MetadataRow,
            r#"
            INSERT INTO horizon_public.metadata_row
                (data_stream_id, datetime, latitude, longitude, altitude,
                 speed, heading, pitch, roll, speed_over_ground,
                 created_datetime, modified_datetime)
            VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)
            RETURNING created_datetime, modified_datetime, data_stream_id,
                      datetime, latitude, longitude, altitude, speed,
                      heading, pitch, roll, speed_over_ground
            "#,
            metadata_row.data_stream_id,
            metadata_row.datetime,
            metadata_row.latitude,
            metadata_row.longitude,
            metadata_row.altitude,
            metadata_row.speed,
            metadata_row.heading,
            metadata_row.pitch,
            metadata_row.roll,
            metadata_row.speed_over_ground,
            metadata_row.created_datetime.unwrap_or(now),
            now,
        )
        .fetch_one(&self.pool)
        .await?)
    }

    /// Insert a batch of metadata rows inside a single transaction.
    ///
    /// Each row is inserted via its own statement against the shared
    /// transaction; if any insert fails the transaction is rolled back and
    /// no rows from the batch are committed. Matches the atomicity
    /// guarantee of Python's `insert_batch` in `horizon-data-core`.
    pub async fn insert_metadata_row_batch(&self, rows: &[MetadataRow]) -> Result<()> {
        if rows.is_empty() {
            return Ok(());
        }
        let now = Utc::now();
        let mut tx = self.pool.begin().await?;
        for row in rows {
            sqlx::query!(
                r#"
                INSERT INTO horizon_public.metadata_row
                    (data_stream_id, datetime, latitude, longitude, altitude,
                     speed, heading, pitch, roll, speed_over_ground,
                     created_datetime, modified_datetime)
                VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)
                "#,
                row.data_stream_id,
                row.datetime,
                row.latitude,
                row.longitude,
                row.altitude,
                row.speed,
                row.heading,
                row.pitch,
                row.roll,
                row.speed_over_ground,
                row.created_datetime.unwrap_or(now),
                now,
            )
            .execute(&mut *tx)
            .await?;
        }
        tx.commit().await?;
        Ok(())
    }

    /// Ordered by datetime. Joins through `data_stream` for RLS scoping.
    pub async fn list_metadata_rows(&self, data_stream_id: Uuid) -> Result<Vec<MetadataRow>> {
        Ok(sqlx::query_as!(
            MetadataRow,
            r#"
            SELECT mr.created_datetime, mr.modified_datetime, mr.data_stream_id,
                   mr.datetime, mr.latitude, mr.longitude, mr.altitude, mr.speed,
                   mr.heading, mr.pitch, mr.roll, mr.speed_over_ground
            FROM horizon_public.metadata_row mr
            JOIN horizon_public.data_stream ds ON mr.data_stream_id = ds.id
            WHERE mr.data_stream_id = $1
            ORDER BY mr.datetime
            "#,
            data_stream_id,
        )
        .fetch_all(&self.pool)
        .await?)
    }
}