horizon-sdk 6.7.2

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

use crate::postgres::retry::{RetryPolicy, retry};
use crate::types::error::{HorizonError, PostgresError, Result};
use crate::types::model::DataStream;

use super::PostgresRepository;

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

    /// Create a new data stream.
    pub async fn insert_data_stream(&self, data_stream: &DataStream) -> Result<DataStream> {
        let now = Utc::now();
        let org = self.organization_id.or(data_stream.organization_id);
        Ok(sqlx::query_as!(
            DataStream,
            r#"
            INSERT INTO horizon_public.data_stream
                (id, platform_id, organization_id, name, query_source, modified_datetime)
            VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6)
            RETURNING
                id,
                created_datetime,
                modified_datetime,
                platform_id,
                organization_id,
                name,
                query_source AS "query_source!: _"
            "#,
            data_stream.id,
            data_stream.platform_id,
            org,
            data_stream.name,
            data_stream.query_source as _,
            now,
        )
        .fetch_one(&self.pool)
        .await?)
    }

    /// List all data streams. Retried under the default [`RetryPolicy`] on transient
    /// Postgres errors; reads are naturally idempotent.
    pub async fn list_data_streams(&self) -> Result<Vec<DataStream>> {
        retry(RetryPolicy::default(), || async move {
            Ok(sqlx::query_as!(
                DataStream,
                r#"
                SELECT
                    id,
                    created_datetime,
                    modified_datetime,
                    platform_id,
                    organization_id,
                    name,
                    query_source AS "query_source!: _"
                FROM horizon_public.data_stream
                "#,
            )
            .fetch_all(&self.pool)
            .await?)
        })
        .await
    }

    /// Read a data stream by ID. Retried under the default [`RetryPolicy`] on transient
    /// Postgres errors; reads are naturally idempotent.
    pub async fn read_data_stream(&self, id: Uuid) -> Result<Option<DataStream>> {
        retry(RetryPolicy::default(), || async move {
            Ok(sqlx::query_as!(
                DataStream,
                r#"
                SELECT
                    id,
                    created_datetime,
                    modified_datetime,
                    platform_id,
                    organization_id,
                    name,
                    query_source AS "query_source!: _"
                FROM horizon_public.data_stream
                WHERE id = $1
                "#,
                id,
            )
            .fetch_optional(&self.pool)
            .await?)
        })
        .await
    }

    /// Update an existing data stream by ID.
    ///
    /// Returns `HorizonError::Postgres(PostgresError::NotFound)` if no row
    /// with the given ID exists. Matches Python `BaseRepository.update()`
    /// semantics.
    pub async fn update_data_stream(&self, data_stream: &DataStream) -> Result<DataStream> {
        let now = Utc::now();
        let result = sqlx::query_as!(
            DataStream,
            r#"
            UPDATE horizon_public.data_stream
            SET
                platform_id = $2,
                name = $3,
                query_source = $4,
                modified_datetime = $5
            WHERE id = $1
            RETURNING
                id,
                created_datetime,
                modified_datetime,
                platform_id,
                organization_id,
                name,
                query_source AS "query_source!: _"
            "#,
            data_stream.id,
            data_stream.platform_id,
            data_stream.name,
            data_stream.query_source as _,
            now,
        )
        .fetch_optional(&self.pool)
        .await?;

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

    /// Create a new data stream or update if it already exists.
    pub async fn upsert_data_stream(&self, data_stream: &DataStream) -> Result<DataStream> {
        let now = Utc::now();
        let org = self.organization_id.or(data_stream.organization_id);
        Ok(sqlx::query_as!(
            DataStream,
            r#"
            INSERT INTO horizon_public.data_stream
                (id, platform_id, organization_id, name, query_source, modified_datetime)
            VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6)
            ON CONFLICT (id) DO UPDATE SET
                platform_id = EXCLUDED.platform_id,
                name = EXCLUDED.name,
                query_source = EXCLUDED.query_source,
                modified_datetime = EXCLUDED.modified_datetime
            RETURNING
                id,
                created_datetime,
                modified_datetime,
                platform_id,
                organization_id,
                name,
                query_source AS "query_source!: _"
            "#,
            data_stream.id,
            data_stream.platform_id,
            org,
            data_stream.name,
            data_stream.query_source as _,
            now,
        )
        .fetch_one(&self.pool)
        .await?)
    }
}