horizon-sdk 9.0.0

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

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

use super::PostgresRepository;

impl PostgresRepository {
    /// Create a new platform information record.
    pub async fn insert_platform_information(
        &self,
        pi: &PlatformInformation,
    ) -> Result<PlatformInformation> {
        let now = Utc::now();
        let org = self.organization_id.or(pi.organization_id);
        Ok(sqlx::query_as!(
            PlatformInformation,
            r#"
            INSERT INTO horizon_public.platform_information
                (id, platform_id, organization_id, properties, modified_datetime)
            VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5)
            RETURNING id, created_datetime, modified_datetime,
                      platform_id, organization_id, properties
            "#,
            pi.id,
            pi.platform_id,
            org,
            pi.properties,
            now,
        )
        .fetch_one(&self.pool)
        .await?)
    }

    /// Create a batch of platform information records 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.
    /// `properties` is a `jsonb` column; sqlx encodes a slice of
    /// `serde_json::Value` as a `jsonb[]` parameter element-wise.
    ///
    /// 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_platform_information_batch(
        &self,
        records: &[PlatformInformation],
    ) -> Result<Vec<PlatformInformation>> {
        if records.is_empty() {
            return Ok(Vec::new());
        }
        let now = Utc::now();
        let ids: Vec<Option<Uuid>> = records.iter().map(|record| record.id).collect();
        let platform_ids: Vec<Uuid> = records.iter().map(|record| record.platform_id).collect();
        let orgs: Vec<Option<Uuid>> = records
            .iter()
            .map(|record| self.organization_id.or(record.organization_id))
            .collect();
        let properties: Vec<serde_json::Value> = records
            .iter()
            .map(|record| record.properties.clone())
            .collect();
        let mut tx = self.pool.begin().await?;
        let rows = sqlx::query_as!(
            PlatformInformation,
            r#"
            WITH source AS (
                SELECT
                    COALESCE(id, gen_random_uuid()) AS id,
                    platform_id, organization_id, properties, ord
                FROM UNNEST($1::uuid[], $2::uuid[], $3::uuid[], $4::jsonb[]) WITH ORDINALITY
                    AS batch(id, platform_id, organization_id, properties, ord)
            ),
            inserted AS (
                INSERT INTO horizon_public.platform_information
                    (id, platform_id, organization_id, properties, modified_datetime)
                SELECT id, platform_id, organization_id, properties, $5 FROM source
                RETURNING id, created_datetime, modified_datetime,
                          platform_id, organization_id, properties
            )
            SELECT inserted.id, inserted.created_datetime, inserted.modified_datetime,
                   inserted.platform_id, inserted.organization_id, inserted.properties
            FROM inserted
            JOIN source ON source.id = inserted.id
            ORDER BY source.ord
            "#,
            &ids as &[Option<Uuid>],
            &platform_ids as &[Uuid],
            &orgs as &[Option<Uuid>],
            &properties as &[serde_json::Value],
            now,
        )
        .fetch_all(&mut *tx)
        .await?;
        tx.commit().await?;
        debug!(batch_size = rows.len(), "insert_platform_information_batch");
        Ok(rows)
    }

    /// List all platform information record records.
    pub async fn list_platform_information(&self) -> Result<Vec<PlatformInformation>> {
        Ok(sqlx::query_as!(
            PlatformInformation,
            r#"
            SELECT id, created_datetime, modified_datetime,
                   platform_id, organization_id, properties
            FROM horizon_public.platform_information
            "#,
        )
        .fetch_all(&self.pool)
        .await?)
    }
}