use chrono::Utc;
use tracing::debug;
use uuid::Uuid;
use crate::postgres::retry::{RetryPolicy, retry};
use crate::types::error::{HorizonError, PostgresError, Result};
use crate::types::model::{DataStream, DataStreamQuerySource};
use super::PostgresRepository;
impl PostgresRepository {
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())
}
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?)
}
#[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_data_stream_batch(
&self,
data_streams: &[DataStream],
) -> Result<Vec<DataStream>> {
if data_streams.is_empty() {
return Ok(Vec::new());
}
let now = Utc::now();
let ids: Vec<Option<Uuid>> = data_streams
.iter()
.map(|data_stream| data_stream.id)
.collect();
let platform_ids: Vec<Option<Uuid>> = data_streams
.iter()
.map(|data_stream| data_stream.platform_id)
.collect();
let orgs: Vec<Option<Uuid>> = data_streams
.iter()
.map(|data_stream| self.organization_id.or(data_stream.organization_id))
.collect();
let names: Vec<Option<String>> = data_streams
.iter()
.map(|data_stream| data_stream.name.clone())
.collect();
let query_sources: Vec<DataStreamQuerySource> = data_streams
.iter()
.map(|data_stream| data_stream.query_source)
.collect();
let mut tx = self.pool.begin().await?;
let rows = sqlx::query_as!(
DataStream,
r#"
WITH source AS (
SELECT
COALESCE(id, gen_random_uuid()) AS id,
platform_id, organization_id, name, query_source, ord
FROM UNNEST(
$1::uuid[], $2::uuid[], $3::uuid[], $4::text[],
$5::horizon_public.data_stream_query_source[]
) WITH ORDINALITY
AS batch(id, platform_id, organization_id, name, query_source, ord)
),
inserted AS (
INSERT INTO horizon_public.data_stream
(id, platform_id, organization_id, name, query_source, modified_datetime)
SELECT id, platform_id, organization_id, name, query_source, $6 FROM source
RETURNING
id, created_datetime, modified_datetime,
platform_id, organization_id, name, query_source
)
SELECT
inserted.id,
inserted.created_datetime,
inserted.modified_datetime,
inserted.platform_id,
inserted.organization_id,
inserted.name,
inserted.query_source AS "query_source!: _"
FROM inserted
JOIN source ON source.id = inserted.id
ORDER BY source.ord
"#,
&ids as &[Option<Uuid>],
&platform_ids as &[Option<Uuid>],
&orgs as &[Option<Uuid>],
&names as &[Option<String>],
&query_sources as &[DataStreamQuerySource],
now,
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
debug!(batch_size = rows.len(), "insert_data_stream_batch");
Ok(rows)
}
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
}
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
}
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(),
})
})
}
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?)
}
}