use chrono::Utc;
use tracing::debug;
use uuid::Uuid;
use crate::types::error::{HorizonError, PostgresError, Result};
use crate::types::filter::AudioSpecificationFilter;
use crate::types::model::AudioSpecification;
use super::PostgresRepository;
impl PostgresRepository {
pub async fn delete_audio_specification(&self, id: Uuid) -> Result<u64> {
let result = sqlx::query!(
"DELETE FROM horizon_public.audio_specification WHERE id = $1",
id
)
.execute(&self.pool)
.await?;
Ok(result.rows_affected())
}
pub async fn insert_audio_specification(
&self,
spec: &AudioSpecification,
) -> Result<AudioSpecification> {
let now = Utc::now();
let org = self.organization_id.or(spec.organization_id);
Ok(sqlx::query_as!(
AudioSpecification,
r#"
INSERT INTO horizon_public.audio_specification
(id, sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id, modified_datetime)
VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8)
RETURNING id, created_datetime, modified_datetime,
sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id
"#,
spec.id,
spec.sample_rate,
spec.bit_depth,
spec.channel_count,
spec.channel_index,
spec.encoding,
org,
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_audio_specification_batch(
&self,
specs: &[AudioSpecification],
) -> Result<Vec<AudioSpecification>> {
if specs.is_empty() {
return Ok(Vec::new());
}
let now = Utc::now();
let ids: Vec<Option<Uuid>> = specs.iter().map(|spec| spec.id).collect();
let sample_rates: Vec<i64> = specs.iter().map(|spec| spec.sample_rate).collect();
let bit_depths: Vec<i64> = specs.iter().map(|spec| spec.bit_depth).collect();
let channel_counts: Vec<i64> = specs.iter().map(|spec| spec.channel_count).collect();
let channel_indices: Vec<i64> = specs.iter().map(|spec| spec.channel_index).collect();
let encodings: Vec<String> = specs.iter().map(|spec| spec.encoding.clone()).collect();
let orgs: Vec<Option<Uuid>> = specs
.iter()
.map(|spec| self.organization_id.or(spec.organization_id))
.collect();
let mut tx = self.pool.begin().await?;
let rows = sqlx::query_as!(
AudioSpecification,
r#"
WITH source AS (
SELECT
COALESCE(id, gen_random_uuid()) AS id,
sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id, ord
FROM UNNEST(
$1::uuid[], $2::bigint[], $3::bigint[], $4::bigint[], $5::bigint[],
$6::text[], $7::uuid[]
) WITH ORDINALITY
AS batch(id, sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id, ord)
),
inserted AS (
INSERT INTO horizon_public.audio_specification
(id, sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id, modified_datetime)
SELECT
id, sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id, $8
FROM source
RETURNING id, created_datetime, modified_datetime,
sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id
)
SELECT inserted.id, inserted.created_datetime, inserted.modified_datetime,
inserted.sample_rate, inserted.bit_depth, inserted.channel_count,
inserted.channel_index, inserted.encoding, inserted.organization_id
FROM inserted
JOIN source ON source.id = inserted.id
ORDER BY source.ord
"#,
&ids as &[Option<Uuid>],
&sample_rates as &[i64],
&bit_depths as &[i64],
&channel_counts as &[i64],
&channel_indices as &[i64],
&encodings as &[String],
&orgs as &[Option<Uuid>],
now,
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
debug!(batch_size = rows.len(), "insert_audio_specification_batch");
Ok(rows)
}
pub async fn list_audio_specifications(
&self,
filter: &AudioSpecificationFilter,
) -> Result<Vec<AudioSpecification>> {
Ok(sqlx::query_as!(
AudioSpecification,
r#"
SELECT id, created_datetime, modified_datetime,
sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id
FROM horizon_public.audio_specification
WHERE ($1::bigint IS NULL OR sample_rate = $1)
AND ($2::bigint IS NULL OR bit_depth = $2)
AND ($3::bigint IS NULL OR channel_count = $3)
AND ($4::bigint IS NULL OR channel_index = $4)
AND ($5::text IS NULL OR encoding = $5)
AND ($6::uuid IS NULL OR organization_id = $6)
"#,
filter.sample_rate,
filter.bit_depth,
filter.channel_count,
filter.channel_index,
filter.encoding,
filter.organization_id,
)
.fetch_all(&self.pool)
.await?)
}
pub async fn read_audio_specification(&self, id: Uuid) -> Result<Option<AudioSpecification>> {
Ok(sqlx::query_as!(
AudioSpecification,
r#"
SELECT id, created_datetime, modified_datetime,
sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id
FROM horizon_public.audio_specification WHERE id = $1
"#,
id,
)
.fetch_optional(&self.pool)
.await?)
}
pub async fn update_audio_specification(
&self,
spec: &AudioSpecification,
) -> Result<AudioSpecification> {
let now = Utc::now();
let org = self.organization_id.or(spec.organization_id);
let result = sqlx::query_as!(
AudioSpecification,
r#"
UPDATE horizon_public.audio_specification
SET
sample_rate = $2,
bit_depth = $3,
channel_count = $4,
channel_index = $5,
encoding = $6,
organization_id = $7,
modified_datetime = $8
WHERE id = $1
RETURNING id, created_datetime, modified_datetime,
sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id
"#,
spec.id,
spec.sample_rate,
spec.bit_depth,
spec.channel_count,
spec.channel_index,
spec.encoding,
org,
now,
)
.fetch_optional(&self.pool)
.await?;
result.ok_or_else(|| {
HorizonError::Postgres(PostgresError::NotFound {
entity: "audio_specification".to_owned(),
id: spec.id.into(),
})
})
}
pub async fn upsert_audio_specification(
&self,
spec: &AudioSpecification,
) -> Result<AudioSpecification> {
let now = Utc::now();
let org = self.organization_id.or(spec.organization_id);
Ok(sqlx::query_as!(
AudioSpecification,
r#"
INSERT INTO horizon_public.audio_specification
(id, sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id, modified_datetime)
VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8)
ON CONFLICT (id) DO UPDATE SET
sample_rate = EXCLUDED.sample_rate,
bit_depth = EXCLUDED.bit_depth,
channel_count = EXCLUDED.channel_count,
channel_index = EXCLUDED.channel_index,
encoding = EXCLUDED.encoding,
organization_id = EXCLUDED.organization_id,
modified_datetime = EXCLUDED.modified_datetime
RETURNING id, created_datetime, modified_datetime,
sample_rate, bit_depth, channel_count, channel_index,
encoding, organization_id
"#,
spec.id,
spec.sample_rate,
spec.bit_depth,
spec.channel_count,
spec.channel_index,
spec.encoding,
org,
now,
)
.fetch_one(&self.pool)
.await?)
}
}