horizon-sdk 6.7.2

Canonical Rust data access layer for the Horizon platform
Documentation
#![allow(
    clippy::big_endian_bytes,
    reason = "PostgreSQL binary COPY format specifies big-endian encoding throughout"
)]
use chrono::{DateTime, Utc};
use futures::future::try_join_all;
use uuid::Uuid;

use crate::postgres::retry::{RetryPolicy, retry};
use crate::types::error::Result;
use crate::types::model::DataRow;

use super::PostgresRepository;
use sqlx::PgPool;

impl PostgresRepository {
    /// Run a single insert pass via up to `workers` concurrent binary COPY streams:
    /// the rows are split into `workers` chunks, each `COPYed` on its own
    /// connection/backend. The streams are driven concurrently on one task, so while
    /// one is blocked on socket I/O the others make progress. Returns the total
    /// inserted.
    ///
    /// Each partition is retried independently under a reduced-budget [`RetryPolicy`] so a
    /// transient failure on one connection re-COPYs only that partition rather than
    /// bubbling up to a whole-batch retry, which would re-COPY partitions that already
    /// committed and duplicate rows.
    async fn binary_parallel_pass(
        &self,
        row_slice: &[DataRow],
        worker_count_option: Option<usize>,
    ) -> Result<u64> {
        // 8 workers was simply the optimal number found in local testing. In the future this should
        // be refactored to be an argument passed to the SDK, chose not to do that for the moment to avoid a
        // breaking change
        let worker_count = worker_count_option.unwrap_or(8);
        let part = row_slice.len().div_ceil(worker_count).max(1);
        // Tighter budget than the idempotent read default: a connection drop during a
        // partition COMMIT is ambiguous, so each retry risks a double-write. Two retries
        // (three attempts total) matches horizon-data-core's `insert_batch` policy.
        let policy = RetryPolicy::default().with_max_retries(2);
        let count_vec = try_join_all(row_slice.chunks(part).map(|chunk| {
            retry(policy, move || {
                Self::copy_binary_partition(&self.pool, chunk)
            })
        }))
        .await?;
        Ok(count_vec.into_iter().sum())
    }

    /// Run one binary `COPY ... FROM STDIN WITH (FORMAT binary)` for `rows` inside its
    /// own transaction on a single connection from `pool`, returning the number of rows
    /// `COPYed`. Builds the binary stream ourselves (file header, one tuple per row with
    /// field count plus length-prefixed binary fields, and the -1 trailer), streaming it
    /// to the `PgCopyIn` sink, then commits.
    ///
    /// Wrapping the `COPY` in a transaction makes the per-partition [`retry`] in
    /// `binary_parallel_pass` idempotent for server-side aborts: a failed attempt is
    /// dropped without `commit`, rolling back any partial work so the retry re-COPYs from
    /// a clean slate. The residual case is a connection drop during `COMMIT`, which is
    /// ambiguous (the commit may have landed); a retry can then re-COPY the partition, so
    /// `data_row` writes are at-least-once.
    #[allow(
        clippy::single_call_fn,
        reason = "Used to breakup binary_parallel_pass function for readability"
    )]
    async fn copy_binary_partition(pool: &PgPool, row_slice: &[DataRow]) -> Result<u64> {
        let mut transaction = pool.begin().await?;
        let copied = {
            let mut copy = transaction
                .copy_in_raw(
                    "COPY horizon_public.data_row \
                     (data_stream_id, datetime, vector, data_type, specification_id, \
                      vector_start_bound, vector_end_bound, created_datetime, modified_datetime) \
                     FROM STDIN WITH (FORMAT binary)",
                )
                .await?;

            // File header: 11-byte signature, int32 flags = 0, int32 header-extension length = 0.
            let mut buf: Vec<u8> = Vec::with_capacity(16 * 1024 * 1024);
            buf.extend_from_slice(b"PGCOPY\n\xFF\r\n\0");
            buf.extend_from_slice(&0_i32.to_be_bytes()); // flags
            buf.extend_from_slice(&0_i32.to_be_bytes()); // header extension area length

            for row in row_slice {
                let now = Utc::now();
                buf.extend_from_slice(&9_i16.to_be_bytes()); // field count for this tuple
                put_field(&mut buf, Some(row.data_stream_id.as_bytes()));
                put_timestamptz(&mut buf, row.datetime);
                put_float8_array_field(&mut buf, &row.vector);
                put_field(&mut buf, Some(row.data_type.as_bytes()));
                put_field(&mut buf, Some(row.specification_id.as_bytes()));
                put_field(&mut buf, Some(&row.vector_start_bound.to_be_bytes()));
                put_field(&mut buf, Some(&row.vector_end_bound.to_be_bytes()));
                put_timestamptz(&mut buf, row.created_datetime.unwrap_or(now));
                put_timestamptz(&mut buf, now);

                // Stream in ~16 MB chunks so the whole payload never sits in memory at once.
                if buf.len() >= 16 * 1024 * 1024 {
                    copy.send(buf.as_slice()).await?;
                    buf.clear();
                }
            }

            // File trailer: int16 = -1 marks end of data.
            buf.extend_from_slice(&(-1_i16).to_be_bytes());
            copy.send(buf.as_slice()).await?;
            copy.finish().await?
        };
        transaction.commit().await?;
        Ok(copied)
    }
    /// If `created_datetime` is `None`, defaults to now (allows backfilling historical data).
    pub async fn insert_data_row(&self, data_row: &DataRow) -> Result<DataRow> {
        let now = Utc::now();
        Ok(sqlx::query_as!(
            DataRow,
            r#"
            INSERT INTO horizon_public.data_row
                (data_stream_id, datetime, vector, data_type, specification_id,
                 vector_start_bound, vector_end_bound, created_datetime, modified_datetime)
            VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
            RETURNING created_datetime, modified_datetime, data_stream_id,
                      datetime, vector, data_type, specification_id,
                      vector_start_bound, vector_end_bound
            "#,
            data_row.data_stream_id,
            data_row.datetime,
            &data_row.vector,
            data_row.data_type,
            data_row.specification_id,
            data_row.vector_start_bound,
            data_row.vector_end_bound,
            data_row.created_datetime.unwrap_or(now),
            now,
        )
        .fetch_one(&self.pool)
        .await?)
    }

    pub async fn insert_data_row_batch(&self, row_slice: &[DataRow]) -> Result<u64> {
        self.binary_parallel_pass(row_slice, None).await
    }

    /// Ordered by datetime. Joins through `data_stream` for RLS scoping. Retried under
    /// the default [`RetryPolicy`] on transient Postgres errors; reads are naturally
    /// idempotent.
    pub async fn list_data_rows(&self, data_stream_id: Uuid) -> Result<Vec<DataRow>> {
        retry(RetryPolicy::default(), || async move {
            Ok(sqlx::query_as!(
                DataRow,
                r#"
                SELECT dr.created_datetime, dr.modified_datetime, dr.data_stream_id,
                       dr.datetime, dr.vector, dr.data_type, dr.specification_id,
                       dr.vector_start_bound, dr.vector_end_bound
                FROM horizon_public.data_row dr
                JOIN horizon_public.data_stream ds ON dr.data_stream_id = ds.id
                WHERE dr.data_stream_id = $1
                ORDER BY dr.datetime
                "#,
                data_stream_id,
            )
            .fetch_all(&self.pool)
            .await?)
        })
        .await
    }
}

/// Append one COPY-binary field: a 4-byte big-endian length prefix followed by
/// the field's binary bytes. `None` encodes a SQL NULL as a length of -1 with no
/// bytes following.
#[allow(
    clippy::cast_possible_truncation,
    clippy::cast_possible_wrap,
    clippy::as_conversions,
    reason = "Field lengths are bounded by PostgreSQL protocol limits, well within i32 range"
)]
fn put_field(buf: &mut Vec<u8>, bytes: Option<&[u8]>) {
    match bytes {
        Some(b) => {
            buf.extend_from_slice(&(b.len() as i32).to_be_bytes());
            buf.extend_from_slice(b);
        }
        None => buf.extend_from_slice(&(-1_i32).to_be_bytes()),
    }
}

/// Encode `values` as a 1-D `float8[]` (no NULLs) directly into `buf`, including
/// the COPY field length prefix. Avoids a temporary allocation compared to building
/// a `Vec<u8>` and then copying it in via `put_field`.
#[allow(
    clippy::single_call_fn,
    reason = "Re-usable function that could be used for other array fields in the future"
)]
#[allow(
    clippy::cast_possible_truncation,
    clippy::cast_possible_wrap,
    clippy::as_conversions,
    clippy::arithmetic_side_effects,
    reason = "Vector lengths are bounded by application constraints, well within i32 range"
)]
fn put_float8_array_field(buf: &mut Vec<u8>, values: &[f64]) {
    const FLOAT8_OID: i32 = 701;
    // 20-byte array header + 12 bytes per element (4-byte length + 8-byte f64).
    let array_len = (20 + values.len() * 12) as i32;
    buf.extend_from_slice(&array_len.to_be_bytes()); // COPY field length
    buf.extend_from_slice(&1_i32.to_be_bytes()); // number of dimensions
    buf.extend_from_slice(&0_i32.to_be_bytes()); // flags (no NULLs)
    buf.extend_from_slice(&FLOAT8_OID.to_be_bytes()); // element OID
    buf.extend_from_slice(&(values.len() as i32).to_be_bytes()); // dim length
    buf.extend_from_slice(&1_i32.to_be_bytes()); // lower bound (1-based)
    for &v in values {
        buf.extend_from_slice(&8_i32.to_be_bytes()); // element byte length
        buf.extend_from_slice(&v.to_be_bytes());
    }
}

/// Append a non-null `timestamptz` field: micros since the PG epoch (2000-01-01)
/// as a big-endian `i64`, length-prefixed like any other field.
#[allow(
    clippy::arithmetic_side_effects,
    reason = "Timestamps are always after the PG epoch (2000-01-01); subtraction cannot underflow"
)]
fn put_timestamptz(buf: &mut Vec<u8>, dt: DateTime<Utc>) {
    /// Microseconds between the Unix epoch (1970-01-01) and the `PostgreSQL` epoch
    /// (2000-01-01) — timestamptz binary is micros since the PG epoch.
    const PG_EPOCH_OFFSET_MICROS: i64 = 946_684_800_000_000;
    let micros = dt.timestamp_micros() - PG_EPOCH_OFFSET_MICROS;
    put_field(buf, Some(&micros.to_be_bytes()));
}