use chrono::Utc;
use tracing::debug;
use uuid::Uuid;
use crate::types::error::{HorizonError, PostgresError, Result};
use crate::types::filter::BeamgramSpecificationFilter;
use crate::types::model::BeamgramSpecification;
use sqlx::{Postgres, QueryBuilder};
use super::PostgresRepository;
impl PostgresRepository {
pub async fn bulk_insert_beamgram_specification(
&self,
specs: &[BeamgramSpecification],
) -> Result<Vec<Uuid>> {
const COLUMN_COUNT: usize = 14;
#[allow(
clippy::decimal_literal_representation,
reason = "Integer representation is more readable for this applications"
)]
const MAX_PARAMETER_COUNT: usize = 65_535;
#[allow(
clippy::integer_division,
clippy::integer_division_remainder_used,
reason = "Can only perform whole writes at once"
)]
let max_writes = MAX_PARAMETER_COUNT / COLUMN_COUNT;
let now = Utc::now();
let mut id_vector = Vec::with_capacity(specs.len());
for chunk in specs.chunks(max_writes) {
let mut query_builder: QueryBuilder<Postgres> = QueryBuilder::new(
"INSERT INTO horizon_public.beamgram_specification \
(id, center_bearing, center_bin_width, min_frequency, max_frequency, \
fft_sample_count, organization_id, update_rate, normalizer, name, \
elevation_increment, lower_elevation, upper_elevation, modified_datetime) ",
);
query_builder.push_values(chunk, |mut b, spec| {
b.push("COALESCE(")
.push_bind_unseparated(spec.id)
.push_unseparated(", gen_random_uuid())")
.push_bind(spec.center_bearing)
.push_bind(spec.center_bin_width)
.push_bind(spec.min_frequency)
.push_bind(spec.max_frequency)
.push_bind(spec.fft_sample_count)
.push_bind(self.organization_id.or(spec.organization_id))
.push_bind(spec.update_rate_ms)
.push_unseparated("::bigint * interval '1 millisecond'")
.push_bind(spec.normalizer.clone())
.push_bind(spec.name.clone())
.push_bind(spec.elevation_increment)
.push_bind(spec.lower_elevation)
.push_bind(spec.upper_elevation)
.push_bind(now);
});
query_builder.push(" RETURNING id");
let id_batch: Vec<Uuid> = query_builder
.build_query_scalar()
.fetch_all(&self.pool)
.await?;
id_vector.extend(id_batch);
}
Ok(id_vector)
}
pub async fn delete_beamgram_specification(&self, id: Uuid) -> Result<u64> {
let result = sqlx::query!(
"DELETE FROM horizon_public.beamgram_specification WHERE id = $1",
id
)
.execute(&self.pool)
.await?;
Ok(result.rows_affected())
}
pub async fn insert_beamgram_specification(
&self,
spec: &BeamgramSpecification,
) -> Result<BeamgramSpecification> {
let now = Utc::now();
let org = self.organization_id.or(spec.organization_id);
Ok(sqlx::query_as!(
BeamgramSpecification,
r#"
INSERT INTO horizon_public.beamgram_specification
(id, center_bearing, center_bin_width, elevation_increment,
lower_elevation, upper_elevation, max_frequency,
min_frequency, name, fft_sample_count, organization_id,
update_rate, normalizer, modified_datetime)
VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8, $9, $10, $11,
$12::bigint * interval '1 millisecond', $13, $14)
RETURNING id, created_datetime, modified_datetime, center_bearing,
center_bin_width, elevation_increment, lower_elevation,
upper_elevation, max_frequency, min_frequency,
name, fft_sample_count, organization_id,
(EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
normalizer
"#,
spec.id,
spec.center_bearing,
spec.center_bin_width,
spec.elevation_increment,
spec.lower_elevation,
spec.upper_elevation,
spec.max_frequency,
spec.min_frequency,
spec.name,
spec.fft_sample_count,
org,
spec.update_rate_ms,
spec.normalizer,
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)"
)]
#[allow(
clippy::too_many_lines,
reason = "13-column batch insert: the per-column slice collection plus the order-preserving CTE query is inherently long and reads more clearly inline than split across helpers"
)]
pub async fn insert_beamgram_specification_batch(
&self,
specs: &[BeamgramSpecification],
) -> Result<Vec<BeamgramSpecification>> {
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 center_bearings: Vec<Option<f64>> =
specs.iter().map(|spec| spec.center_bearing).collect();
let center_bin_widths: Vec<Option<f64>> =
specs.iter().map(|spec| spec.center_bin_width).collect();
let elevation_increments: Vec<Option<f64>> =
specs.iter().map(|spec| spec.elevation_increment).collect();
let lower_elevations: Vec<Option<f64>> =
specs.iter().map(|spec| spec.lower_elevation).collect();
let upper_elevations: Vec<Option<f64>> =
specs.iter().map(|spec| spec.upper_elevation).collect();
let max_frequencies: Vec<Option<f64>> =
specs.iter().map(|spec| spec.max_frequency).collect();
let min_frequencies: Vec<Option<f64>> =
specs.iter().map(|spec| spec.min_frequency).collect();
let names: Vec<Option<String>> = specs.iter().map(|spec| spec.name.clone()).collect();
let fft_sample_counts: Vec<Option<i32>> =
specs.iter().map(|spec| spec.fft_sample_count).collect();
let orgs: Vec<Option<Uuid>> = specs
.iter()
.map(|spec| self.organization_id.or(spec.organization_id))
.collect();
let update_rate_ms_list: Vec<Option<i64>> =
specs.iter().map(|spec| spec.update_rate_ms).collect();
let normalizers: Vec<Option<String>> =
specs.iter().map(|spec| spec.normalizer.clone()).collect();
let mut tx = self.pool.begin().await?;
let rows = sqlx::query_as!(
BeamgramSpecification,
r#"
WITH source AS (
SELECT
COALESCE(id, gen_random_uuid()) AS id,
center_bearing, center_bin_width, elevation_increment,
lower_elevation, upper_elevation, max_frequency,
min_frequency, name, fft_sample_count, organization_id,
update_rate_ms, normalizer, ord
FROM UNNEST(
$1::uuid[], $2::float8[], $3::float8[], $4::float8[],
$5::float8[], $6::float8[], $7::float8[], $8::float8[],
$9::text[], $10::integer[], $11::uuid[], $12::bigint[],
$13::text[]
) WITH ORDINALITY
AS batch(id, center_bearing, center_bin_width, elevation_increment,
lower_elevation, upper_elevation, max_frequency,
min_frequency, name, fft_sample_count, organization_id,
update_rate_ms, normalizer, ord)
),
inserted AS (
INSERT INTO horizon_public.beamgram_specification
(id, center_bearing, center_bin_width, elevation_increment,
lower_elevation, upper_elevation, max_frequency,
min_frequency, name, fft_sample_count, organization_id,
update_rate, normalizer, modified_datetime)
SELECT
id, center_bearing, center_bin_width, elevation_increment,
lower_elevation, upper_elevation, max_frequency,
min_frequency, name, fft_sample_count, organization_id,
update_rate_ms * interval '1 millisecond',
normalizer, $14
FROM source
RETURNING id, created_datetime, modified_datetime, center_bearing,
center_bin_width, elevation_increment, lower_elevation,
upper_elevation, max_frequency, min_frequency,
name, fft_sample_count, organization_id, update_rate, normalizer
)
SELECT inserted.id, inserted.created_datetime, inserted.modified_datetime,
inserted.center_bearing, inserted.center_bin_width,
inserted.elevation_increment, inserted.lower_elevation,
inserted.upper_elevation, inserted.max_frequency, inserted.min_frequency,
inserted.name, inserted.fft_sample_count, inserted.organization_id,
(EXTRACT(EPOCH FROM inserted.update_rate) * 1000)::bigint as "update_rate_ms: i64",
inserted.normalizer
FROM inserted
JOIN source ON source.id = inserted.id
ORDER BY source.ord
"#,
&ids as &[Option<Uuid>],
¢er_bearings as &[Option<f64>],
¢er_bin_widths as &[Option<f64>],
&elevation_increments as &[Option<f64>],
&lower_elevations as &[Option<f64>],
&upper_elevations as &[Option<f64>],
&max_frequencies as &[Option<f64>],
&min_frequencies as &[Option<f64>],
&names as &[Option<String>],
&fft_sample_counts as &[Option<i32>],
&orgs as &[Option<Uuid>],
&update_rate_ms_list as &[Option<i64>],
&normalizers as &[Option<String>],
now,
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
debug!(
batch_size = rows.len(),
"insert_beamgram_specification_batch"
);
Ok(rows)
}
pub async fn list_beamgram_specifications(
&self,
filter: &BeamgramSpecificationFilter,
) -> Result<Vec<BeamgramSpecification>> {
Ok(sqlx::query_as!(
BeamgramSpecification,
r#"
SELECT id, created_datetime, modified_datetime, center_bearing,
center_bin_width, elevation_increment, lower_elevation,
upper_elevation, max_frequency, min_frequency,
name, fft_sample_count, organization_id,
(EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
normalizer
FROM horizon_public.beamgram_specification
WHERE ($1::float8 IS NULL OR center_bearing = $1)
AND ($2::float8 IS NULL OR center_bin_width = $2)
AND ($3::float8 IS NULL OR elevation_increment = $3)
AND ($4::integer IS NULL OR fft_sample_count = $4)
AND ($5::float8 IS NULL OR lower_elevation = $5)
AND ($6::float8 IS NULL OR max_frequency = $6)
AND ($7::float8 IS NULL OR min_frequency = $7)
AND ($8::text IS NULL OR name = $8)
AND ($9::text IS NULL OR normalizer = $9)
AND ($10::uuid IS NULL OR organization_id = $10)
AND ($11::bigint IS NULL OR update_rate = $11 * interval '1 millisecond')
AND ($12::float8 IS NULL OR upper_elevation = $12)
"#,
filter.center_bearing,
filter.center_bin_width,
filter.elevation_increment,
filter.fft_sample_count,
filter.lower_elevation,
filter.max_frequency,
filter.min_frequency,
filter.name,
filter.normalizer,
filter.organization_id,
filter.update_rate_ms,
filter.upper_elevation,
)
.fetch_all(&self.pool)
.await?)
}
pub async fn read_beamgram_specification(
&self,
id: Uuid,
) -> Result<Option<BeamgramSpecification>> {
Ok(sqlx::query_as!(
BeamgramSpecification,
r#"
SELECT id, created_datetime, modified_datetime, center_bearing,
center_bin_width, elevation_increment, lower_elevation,
upper_elevation, max_frequency, min_frequency,
name, fft_sample_count, organization_id,
(EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
normalizer
FROM horizon_public.beamgram_specification WHERE id = $1
"#,
id,
)
.fetch_optional(&self.pool)
.await?)
}
pub async fn update_beamgram_specification(
&self,
spec: &BeamgramSpecification,
) -> Result<BeamgramSpecification> {
let now = Utc::now();
let org = self.organization_id.or(spec.organization_id);
let result = sqlx::query_as!(
BeamgramSpecification,
r#"
UPDATE horizon_public.beamgram_specification
SET
center_bearing = $2,
center_bin_width = $3,
elevation_increment = $4,
lower_elevation = $5,
upper_elevation = $6,
max_frequency = $7,
min_frequency = $8,
name = $9,
fft_sample_count = $10,
organization_id = $11,
update_rate = $12::bigint * interval '1 millisecond',
normalizer = $13,
modified_datetime = $14
WHERE id = $1
RETURNING id, created_datetime, modified_datetime, center_bearing,
center_bin_width, elevation_increment, lower_elevation,
upper_elevation, max_frequency, min_frequency,
name, fft_sample_count, organization_id,
(EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
normalizer
"#,
spec.id,
spec.center_bearing,
spec.center_bin_width,
spec.elevation_increment,
spec.lower_elevation,
spec.upper_elevation,
spec.max_frequency,
spec.min_frequency,
spec.name,
spec.fft_sample_count,
org,
spec.update_rate_ms,
spec.normalizer,
now,
)
.fetch_optional(&self.pool)
.await?;
result.ok_or_else(|| {
HorizonError::Postgres(PostgresError::NotFound {
entity: "beamgram_specification".to_owned(),
id: spec.id.into(),
})
})
}
pub async fn upsert_beamgram_specification(
&self,
spec: &BeamgramSpecification,
) -> Result<BeamgramSpecification> {
let now = Utc::now();
let org = self.organization_id.or(spec.organization_id);
Ok(sqlx::query_as!(
BeamgramSpecification,
r#"
INSERT INTO horizon_public.beamgram_specification
(id, center_bearing, center_bin_width, elevation_increment,
lower_elevation, upper_elevation, max_frequency,
min_frequency, name, fft_sample_count, organization_id,
update_rate, normalizer, modified_datetime)
VALUES (COALESCE($1, gen_random_uuid()), $2, $3, $4, $5, $6, $7, $8, $9, $10, $11,
$12::bigint * interval '1 millisecond', $13, $14)
ON CONFLICT (id) DO UPDATE SET
center_bearing = EXCLUDED.center_bearing,
center_bin_width = EXCLUDED.center_bin_width,
elevation_increment = EXCLUDED.elevation_increment,
lower_elevation = EXCLUDED.lower_elevation,
upper_elevation = EXCLUDED.upper_elevation,
max_frequency = EXCLUDED.max_frequency,
min_frequency = EXCLUDED.min_frequency,
name = EXCLUDED.name, fft_sample_count = EXCLUDED.fft_sample_count,
organization_id = EXCLUDED.organization_id,
update_rate = EXCLUDED.update_rate,
normalizer = EXCLUDED.normalizer,
modified_datetime = EXCLUDED.modified_datetime
RETURNING id, created_datetime, modified_datetime, center_bearing,
center_bin_width, elevation_increment, lower_elevation,
upper_elevation, max_frequency, min_frequency,
name, fft_sample_count, organization_id,
(EXTRACT(EPOCH FROM update_rate) * 1000)::bigint as "update_rate_ms: i64",
normalizer
"#,
spec.id,
spec.center_bearing,
spec.center_bin_width,
spec.elevation_increment,
spec.lower_elevation,
spec.upper_elevation,
spec.max_frequency,
spec.min_frequency,
spec.name,
spec.fft_sample_count,
org,
spec.update_rate_ms,
spec.normalizer,
now,
)
.fetch_one(&self.pool)
.await?)
}
}