use chrono::Utc;
use tracing::debug;
use uuid::Uuid;
use crate::types::error::Result;
use crate::types::filter::{
PlatformAudioSpecificationFilter, PlatformBeamgramSpecificationFilter,
PlatformBearingTimeRecordSpecificationFilter, PlatformSpectrogramSpecificationFilter,
};
use crate::types::model::{
PlatformAudioSpecification, PlatformBeamgramSpecification,
PlatformBearingTimeRecordSpecification, PlatformSpectrogramSpecification,
};
use super::PostgresRepository;
impl PostgresRepository {
pub async fn delete_platform_audio_specification(&self, id: Uuid) -> Result<u64> {
let result = sqlx::query!(
"DELETE FROM horizon_public.platform_audio_specification WHERE id = $1",
id
)
.execute(&self.pool)
.await?;
Ok(result.rows_affected())
}
pub async fn delete_platform_beamgram_specification(&self, id: Uuid) -> Result<u64> {
let result = sqlx::query!(
"DELETE FROM horizon_public.platform_beamgram_specification WHERE id = $1",
id
)
.execute(&self.pool)
.await?;
Ok(result.rows_affected())
}
pub async fn delete_platform_bearing_time_record_specification(&self, id: Uuid) -> Result<u64> {
let result = sqlx::query!(
"DELETE FROM horizon_public.platform_bearing_time_record_specification WHERE id = $1",
id
)
.execute(&self.pool)
.await?;
Ok(result.rows_affected())
}
pub async fn delete_platform_spectrogram_specification(&self, id: Uuid) -> Result<u64> {
let result = sqlx::query!(
"DELETE FROM horizon_public.platform_spectrogram_specification WHERE id = $1",
id
)
.execute(&self.pool)
.await?;
Ok(result.rows_affected())
}
pub async fn insert_platform_audio_specification(
&self,
j: &PlatformAudioSpecification,
) -> Result<PlatformAudioSpecification> {
let now = Utc::now();
let org = self.organization_id.or(j.organization_id);
Ok(sqlx::query_as!(
PlatformAudioSpecification,
r#"
INSERT INTO horizon_public.platform_audio_specification
(id, platform_id, audio_specification_id, organization_id, modified_datetime)
VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5)
RETURNING id, created_datetime, modified_datetime,
platform_id, audio_specification_id, organization_id
"#,
j.id,
j.platform_id,
j.audio_specification_id,
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_platform_audio_specification_batch(
&self,
links: &[PlatformAudioSpecification],
) -> Result<Vec<PlatformAudioSpecification>> {
if links.is_empty() {
return Ok(Vec::new());
}
let now = Utc::now();
let ids: Vec<Option<Uuid>> = links.iter().map(|link| link.id).collect();
let platform_ids: Vec<Uuid> = links.iter().map(|link| link.platform_id).collect();
let spec_ids: Vec<Uuid> = links
.iter()
.map(|link| link.audio_specification_id)
.collect();
let orgs: Vec<Option<Uuid>> = links
.iter()
.map(|link| self.organization_id.or(link.organization_id))
.collect();
let mut tx = self.pool.begin().await?;
let rows = sqlx::query_as!(
PlatformAudioSpecification,
r#"
WITH source AS (
SELECT
COALESCE(id, gen_random_uuid()) AS id,
platform_id, audio_specification_id, organization_id, ord
FROM UNNEST($1::uuid[], $2::uuid[], $3::uuid[], $4::uuid[]) WITH ORDINALITY
AS batch(id, platform_id, audio_specification_id, organization_id, ord)
),
inserted AS (
INSERT INTO horizon_public.platform_audio_specification
(id, platform_id, audio_specification_id, organization_id, modified_datetime)
SELECT id, platform_id, audio_specification_id, organization_id, $5 FROM source
RETURNING id, created_datetime, modified_datetime,
platform_id, audio_specification_id, organization_id
)
SELECT inserted.id, inserted.created_datetime, inserted.modified_datetime,
inserted.platform_id, inserted.audio_specification_id,
inserted.organization_id
FROM inserted
JOIN source ON source.id = inserted.id
ORDER BY source.ord
"#,
&ids as &[Option<Uuid>],
&platform_ids as &[Uuid],
&spec_ids as &[Uuid],
&orgs as &[Option<Uuid>],
now,
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
debug!(
batch_size = rows.len(),
"insert_platform_audio_specification_batch"
);
Ok(rows)
}
pub async fn insert_platform_beamgram_specification(
&self,
j: &PlatformBeamgramSpecification,
) -> Result<PlatformBeamgramSpecification> {
let now = Utc::now();
let org = self.organization_id.or(j.organization_id);
Ok(sqlx::query_as!(
PlatformBeamgramSpecification,
r#"
INSERT INTO horizon_public.platform_beamgram_specification
(id, platform_id, beamgram_specification_id, organization_id, modified_datetime)
VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5)
RETURNING id, created_datetime, modified_datetime,
platform_id, beamgram_specification_id, organization_id
"#,
j.id,
j.platform_id,
j.beamgram_specification_id,
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_platform_beamgram_specification_batch(
&self,
links: &[PlatformBeamgramSpecification],
) -> Result<Vec<PlatformBeamgramSpecification>> {
if links.is_empty() {
return Ok(Vec::new());
}
let now = Utc::now();
let ids: Vec<Option<Uuid>> = links.iter().map(|link| link.id).collect();
let platform_ids: Vec<Option<Uuid>> = links.iter().map(|link| link.platform_id).collect();
let spec_ids: Vec<Option<Uuid>> = links
.iter()
.map(|link| link.beamgram_specification_id)
.collect();
let orgs: Vec<Option<Uuid>> = links
.iter()
.map(|link| self.organization_id.or(link.organization_id))
.collect();
let mut tx = self.pool.begin().await?;
let rows = sqlx::query_as!(
PlatformBeamgramSpecification,
r#"
WITH source AS (
SELECT
COALESCE(id, gen_random_uuid()) AS id,
platform_id, beamgram_specification_id, organization_id, ord
FROM UNNEST($1::uuid[], $2::uuid[], $3::uuid[], $4::uuid[]) WITH ORDINALITY
AS batch(id, platform_id, beamgram_specification_id, organization_id, ord)
),
inserted AS (
INSERT INTO horizon_public.platform_beamgram_specification
(id, platform_id, beamgram_specification_id, organization_id, modified_datetime)
SELECT id, platform_id, beamgram_specification_id, organization_id, $5 FROM source
RETURNING id, created_datetime, modified_datetime,
platform_id, beamgram_specification_id, organization_id
)
SELECT inserted.id, inserted.created_datetime, inserted.modified_datetime,
inserted.platform_id, inserted.beamgram_specification_id,
inserted.organization_id
FROM inserted
JOIN source ON source.id = inserted.id
ORDER BY source.ord
"#,
&ids as &[Option<Uuid>],
&platform_ids as &[Option<Uuid>],
&spec_ids as &[Option<Uuid>],
&orgs as &[Option<Uuid>],
now,
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
debug!(
batch_size = rows.len(),
"insert_platform_beamgram_specification_batch"
);
Ok(rows)
}
pub async fn insert_platform_bearing_time_record_specification(
&self,
j: &PlatformBearingTimeRecordSpecification,
) -> Result<PlatformBearingTimeRecordSpecification> {
let now = Utc::now();
let org = self.organization_id.or(j.organization_id);
Ok(sqlx::query_as!(
PlatformBearingTimeRecordSpecification,
r#"
INSERT INTO horizon_public.platform_bearing_time_record_specification
(id, platform_id, bearing_time_record_specification_id,
organization_id, modified_datetime)
VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5)
RETURNING id, created_datetime, modified_datetime,
platform_id, bearing_time_record_specification_id, organization_id
"#,
j.id,
j.platform_id,
j.bearing_time_record_specification_id,
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_platform_bearing_time_record_specification_batch(
&self,
links: &[PlatformBearingTimeRecordSpecification],
) -> Result<Vec<PlatformBearingTimeRecordSpecification>> {
if links.is_empty() {
return Ok(Vec::new());
}
let now = Utc::now();
let ids: Vec<Option<Uuid>> = links.iter().map(|link| link.id).collect();
let platform_ids: Vec<Option<Uuid>> = links.iter().map(|link| link.platform_id).collect();
let spec_ids: Vec<Option<Uuid>> = links
.iter()
.map(|link| link.bearing_time_record_specification_id)
.collect();
let orgs: Vec<Option<Uuid>> = links
.iter()
.map(|link| self.organization_id.or(link.organization_id))
.collect();
let mut tx = self.pool.begin().await?;
let rows = sqlx::query_as!(
PlatformBearingTimeRecordSpecification,
r#"
WITH source AS (
SELECT
COALESCE(id, gen_random_uuid()) AS id,
platform_id, bearing_time_record_specification_id, organization_id, ord
FROM UNNEST($1::uuid[], $2::uuid[], $3::uuid[], $4::uuid[]) WITH ORDINALITY
AS batch(id, platform_id, bearing_time_record_specification_id,
organization_id, ord)
),
inserted AS (
INSERT INTO horizon_public.platform_bearing_time_record_specification
(id, platform_id, bearing_time_record_specification_id,
organization_id, modified_datetime)
SELECT
id, platform_id, bearing_time_record_specification_id, organization_id, $5
FROM source
RETURNING id, created_datetime, modified_datetime,
platform_id, bearing_time_record_specification_id, organization_id
)
SELECT inserted.id, inserted.created_datetime, inserted.modified_datetime,
inserted.platform_id, inserted.bearing_time_record_specification_id,
inserted.organization_id
FROM inserted
JOIN source ON source.id = inserted.id
ORDER BY source.ord
"#,
&ids as &[Option<Uuid>],
&platform_ids as &[Option<Uuid>],
&spec_ids as &[Option<Uuid>],
&orgs as &[Option<Uuid>],
now,
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
debug!(
batch_size = rows.len(),
"insert_platform_bearing_time_record_specification_batch"
);
Ok(rows)
}
pub async fn insert_platform_spectrogram_specification(
&self,
j: &PlatformSpectrogramSpecification,
) -> Result<PlatformSpectrogramSpecification> {
let now = Utc::now();
let org = self.organization_id.or(j.organization_id);
Ok(sqlx::query_as!(
PlatformSpectrogramSpecification,
r#"
INSERT INTO horizon_public.platform_spectrogram_specification
(id, platform_id, spectrogram_specification_id,
organization_id, modified_datetime)
VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5)
RETURNING id, created_datetime, modified_datetime,
platform_id, spectrogram_specification_id, organization_id
"#,
j.id,
j.platform_id,
j.spectrogram_specification_id,
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_platform_spectrogram_specification_batch(
&self,
links: &[PlatformSpectrogramSpecification],
) -> Result<Vec<PlatformSpectrogramSpecification>> {
if links.is_empty() {
return Ok(Vec::new());
}
let now = Utc::now();
let ids: Vec<Option<Uuid>> = links.iter().map(|link| link.id).collect();
let platform_ids: Vec<Option<Uuid>> = links.iter().map(|link| link.platform_id).collect();
let spec_ids: Vec<Option<Uuid>> = links
.iter()
.map(|link| link.spectrogram_specification_id)
.collect();
let orgs: Vec<Option<Uuid>> = links
.iter()
.map(|link| self.organization_id.or(link.organization_id))
.collect();
let mut tx = self.pool.begin().await?;
let rows = sqlx::query_as!(
PlatformSpectrogramSpecification,
r#"
WITH source AS (
SELECT
COALESCE(id, gen_random_uuid()) AS id,
platform_id, spectrogram_specification_id, organization_id, ord
FROM UNNEST($1::uuid[], $2::uuid[], $3::uuid[], $4::uuid[]) WITH ORDINALITY
AS batch(id, platform_id, spectrogram_specification_id,
organization_id, ord)
),
inserted AS (
INSERT INTO horizon_public.platform_spectrogram_specification
(id, platform_id, spectrogram_specification_id,
organization_id, modified_datetime)
SELECT
id, platform_id, spectrogram_specification_id, organization_id, $5
FROM source
RETURNING id, created_datetime, modified_datetime,
platform_id, spectrogram_specification_id, organization_id
)
SELECT inserted.id, inserted.created_datetime, inserted.modified_datetime,
inserted.platform_id, inserted.spectrogram_specification_id,
inserted.organization_id
FROM inserted
JOIN source ON source.id = inserted.id
ORDER BY source.ord
"#,
&ids as &[Option<Uuid>],
&platform_ids as &[Option<Uuid>],
&spec_ids as &[Option<Uuid>],
&orgs as &[Option<Uuid>],
now,
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
debug!(
batch_size = rows.len(),
"insert_platform_spectrogram_specification_batch"
);
Ok(rows)
}
pub async fn list_platform_audio_specifications(
&self,
filter: &PlatformAudioSpecificationFilter,
) -> Result<Vec<PlatformAudioSpecification>> {
Ok(sqlx::query_as!(
PlatformAudioSpecification,
r#"
SELECT id, created_datetime, modified_datetime,
platform_id, audio_specification_id, organization_id
FROM horizon_public.platform_audio_specification
WHERE ($1::uuid IS NULL OR audio_specification_id = $1)
AND ($2::uuid IS NULL OR organization_id = $2)
AND ($3::uuid IS NULL OR platform_id = $3)
"#,
filter.audio_specification_id,
filter.organization_id,
filter.platform_id,
)
.fetch_all(&self.pool)
.await?)
}
pub async fn list_platform_beamgram_specifications(
&self,
filter: &PlatformBeamgramSpecificationFilter,
) -> Result<Vec<PlatformBeamgramSpecification>> {
Ok(sqlx::query_as!(
PlatformBeamgramSpecification,
r#"
SELECT id, created_datetime, modified_datetime,
platform_id, beamgram_specification_id, organization_id
FROM horizon_public.platform_beamgram_specification
WHERE ($1::uuid IS NULL OR beamgram_specification_id = $1)
AND ($2::uuid IS NULL OR organization_id = $2)
AND ($3::uuid IS NULL OR platform_id = $3)
"#,
filter.beamgram_specification_id,
filter.organization_id,
filter.platform_id,
)
.fetch_all(&self.pool)
.await?)
}
pub async fn list_platform_bearing_time_record_specifications(
&self,
filter: &PlatformBearingTimeRecordSpecificationFilter,
) -> Result<Vec<PlatformBearingTimeRecordSpecification>> {
Ok(sqlx::query_as!(
PlatformBearingTimeRecordSpecification,
r#"
SELECT id, created_datetime, modified_datetime,
platform_id, bearing_time_record_specification_id, organization_id
FROM horizon_public.platform_bearing_time_record_specification
WHERE ($1::uuid IS NULL OR bearing_time_record_specification_id = $1)
AND ($2::uuid IS NULL OR organization_id = $2)
AND ($3::uuid IS NULL OR platform_id = $3)
"#,
filter.bearing_time_record_specification_id,
filter.organization_id,
filter.platform_id,
)
.fetch_all(&self.pool)
.await?)
}
pub async fn list_platform_spectrogram_specifications(
&self,
filter: &PlatformSpectrogramSpecificationFilter,
) -> Result<Vec<PlatformSpectrogramSpecification>> {
Ok(sqlx::query_as!(
PlatformSpectrogramSpecification,
r#"
SELECT id, created_datetime, modified_datetime,
platform_id, spectrogram_specification_id, organization_id
FROM horizon_public.platform_spectrogram_specification
WHERE ($1::uuid IS NULL OR organization_id = $1)
AND ($2::uuid IS NULL OR platform_id = $2)
AND ($3::uuid IS NULL OR spectrogram_specification_id = $3)
"#,
filter.organization_id,
filter.platform_id,
filter.spectrogram_specification_id,
)
.fetch_all(&self.pool)
.await?)
}
}