#![allow(
clippy::absolute_paths,
clippy::default_numeric_fallback,
clippy::expect_used,
clippy::indexing_slicing,
clippy::missing_assert_message,
clippy::panic,
clippy::panic_in_result_fn,
clippy::shadow_reuse,
clippy::tests_outside_test_module,
clippy::unwrap_used,
reason = "integration tests use assertions, unwrap, expect, panic, untyped numeric literals, and indexing freely"
)]
use chrono::{TimeZone as _, Utc};
use horizon_sdk::postgres::repository::PostgresRepository;
use horizon_sdk::types::error::{EntityId, HorizonError, PostgresError};
use horizon_sdk::types::filter::{
AudioSpecificationFilter, BeamgramSpecificationFilter, SpectrogramSpecificationFilter,
};
use horizon_sdk::types::model::{
Annotation, AudioSpecification, BeamgramSpecification, BearingResolutionAlias,
BearingTimeRecordSpecification, DataRowPayload, DataStream, DataStreamQuerySource,
DirectionalSpectrogramSpecification, ElevationAlias, FocusRange, FrequencyBandAlias,
FrequencyResolutionAlias, MetadataRow, Mission, MissionPlatform, Normalizer, Ontology,
OntologyClass, Platform, PlatformAudioSpecification, PlatformInformation, PlatformKind,
Position, SpectrogramSpecification, UpdateRateAlias, VariantDataRow,
};
use sqlx::postgres::PgPoolOptions;
use uuid::Uuid;
async fn setup() -> PostgresRepository {
let url = std::env::var("DATABASE_URL").unwrap_or_else(|_| {
"postgresql://horizon_owner:insecure_password@localhost:5432/horizon?sslmode=disable"
.to_owned()
});
let pool = PgPoolOptions::new()
.max_connections(2)
.connect(&url)
.await
.expect("Failed to connect to test database");
PostgresRepository::new(pool, None)
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_platform_generates_id() -> Result<(), HorizonError> {
let r = setup().await;
let name = format!("InsertPlatform-{}", Uuid::new_v4());
let p = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(name.clone()),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let created = r.insert_platform(&p).await?;
assert!(created.id.is_some());
assert!(created.created_datetime.is_some());
assert!(created.modified_datetime.is_some());
assert_eq!(created.name.as_deref(), Some(name.as_str()));
r.delete_platform(created.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn read_platform_returns_none_for_missing() -> Result<(), HorizonError> {
let r = setup().await;
let result = r.read_platform(Uuid::new_v4()).await?;
assert!(result.is_none());
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_then_read_platform() -> Result<(), HorizonError> {
let r = setup().await;
let name = format!("ReadTest-{}", Uuid::new_v4());
let p = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(name.clone()),
kind_id: None,
free_text: Some("hello".to_owned()),
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let created = r.insert_platform(&p).await?;
let read = r.read_platform(created.id.unwrap()).await?;
assert!(read.is_some());
let read = read.unwrap();
assert_eq!(read.name, created.name);
assert_eq!(read.free_text.as_deref(), Some("hello"));
r.delete_platform(created.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn upsert_platform_creates_when_new() -> Result<(), HorizonError> {
let r = setup().await;
let id = Uuid::new_v4();
let name = format!("UpsertNew-{}", Uuid::new_v4());
let p = Platform {
id: Some(id),
created_datetime: None,
modified_datetime: None,
name: Some(name.clone()),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let result = r.upsert_platform(&p).await?;
assert_eq!(result.id, Some(id));
assert_eq!(result.name.as_deref(), Some(name.as_str()));
r.delete_platform(id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn upsert_platform_updates_when_existing() -> Result<(), HorizonError> {
let r = setup().await;
let id = Uuid::new_v4();
let p1 = Platform {
id: Some(id),
created_datetime: None,
modified_datetime: None,
name: Some(format!("Original-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
r.upsert_platform(&p1).await?;
let updated_name = format!("Updated-{}", Uuid::new_v4());
let p2 = Platform {
id: Some(id),
created_datetime: None,
modified_datetime: None,
name: Some(updated_name.clone()),
kind_id: None,
free_text: Some("new text".to_owned()),
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let result = r.upsert_platform(&p2).await?;
assert_eq!(result.name.as_deref(), Some(updated_name.as_str()));
assert_eq!(result.free_text.as_deref(), Some("new text"));
r.delete_platform(id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn create_or_update_platform_keeps_earlier_start() -> Result<(), HorizonError> {
let r = setup().await;
let id = Uuid::new_v4();
let early = chrono::Utc::now() - chrono::Duration::hours(2);
let late = chrono::Utc::now();
let p1 = Platform {
id: Some(id),
created_datetime: None,
modified_datetime: None,
name: Some(format!("DateTimeTest-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: Some(early),
end_datetime: Some(late),
};
r.create_or_update_platform(&p1).await?;
let p2 = Platform {
id: Some(id),
created_datetime: None,
modified_datetime: None,
name: Some(format!("DateTimeTest-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: Some(late),
end_datetime: Some(early),
};
let result = r.create_or_update_platform(&p2).await?;
assert_eq!(
result.start_datetime.unwrap().timestamp(),
early.timestamp()
);
assert_eq!(result.end_datetime.unwrap().timestamp(), late.timestamp());
r.delete_platform(id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn create_or_update_platform_none_then_some_datetimes() -> Result<(), HorizonError> {
let r = setup().await;
let id = Uuid::new_v4();
let now = chrono::Utc::now();
let p1 = Platform {
id: Some(id),
created_datetime: None,
modified_datetime: None,
name: Some(format!("NullDates-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let first = r.create_or_update_platform(&p1).await?;
assert!(first.start_datetime.is_none());
assert!(first.end_datetime.is_none());
let p2 = Platform {
start_datetime: Some(now),
end_datetime: Some(now),
..p1.clone()
};
let second = r.create_or_update_platform(&p2).await?;
assert!(second.start_datetime.is_some());
assert!(second.end_datetime.is_some());
r.delete_platform(id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn create_or_update_platform_some_then_none_preserves() -> Result<(), HorizonError> {
let r = setup().await;
let id = Uuid::new_v4();
let now = chrono::Utc::now();
let p1 = Platform {
id: Some(id),
created_datetime: None,
modified_datetime: None,
name: Some(format!("PreserveDates-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: Some(now),
end_datetime: Some(now),
};
r.create_or_update_platform(&p1).await?;
let p2 = Platform {
start_datetime: None,
end_datetime: None,
..p1.clone()
};
let result = r.create_or_update_platform(&p2).await?;
assert_eq!(result.start_datetime.unwrap().timestamp(), now.timestamp());
assert_eq!(result.end_datetime.unwrap().timestamp(), now.timestamp());
r.delete_platform(id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn create_or_update_platform_both_none_stays_none() -> Result<(), HorizonError> {
let r = setup().await;
let id = Uuid::new_v4();
let p = Platform {
id: Some(id),
created_datetime: None,
modified_datetime: None,
name: Some(format!("AllNone-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
r.create_or_update_platform(&p).await?;
let result = r.create_or_update_platform(&p).await?;
assert!(result.start_datetime.is_none());
assert!(result.end_datetime.is_none());
r.delete_platform(id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn delete_platform_returns_true() -> Result<(), HorizonError> {
let r = setup().await;
let p = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("ToDelete-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let created = r.insert_platform(&p).await?;
let deleted = r.delete_platform(created.id.unwrap()).await?;
assert_eq!(deleted, 1);
let read = r.read_platform(created.id.unwrap()).await?;
assert!(read.is_none());
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn delete_platform_returns_zero_for_missing() -> Result<(), HorizonError> {
let r = setup().await;
let deleted = r.delete_platform(Uuid::new_v4()).await?;
assert_eq!(deleted, 0);
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn list_platforms_returns_inserted() -> Result<(), HorizonError> {
let r = setup().await;
let name = format!("Listed-{}", Uuid::new_v4());
let p = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(name.clone()),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let created = r.insert_platform(&p).await?;
let all = r.list_platforms().await?;
assert!(
all.iter()
.any(|pl| pl.name.as_deref() == Some(name.as_str()))
);
r.delete_platform(created.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_data_stream_with_platform() -> Result<(), HorizonError> {
let r = setup().await;
let p = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("DS Parent-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let platform = r.insert_platform(&p).await?;
let ds = DataStream::from_platform_and_name(platform.id.unwrap(), "audio");
let created = r.insert_data_stream(&ds).await?;
assert!(created.id.is_some());
assert_eq!(created.platform_id, platform.id);
assert_eq!(created.name.as_deref(), Some("Audio"));
r.delete_data_stream(created.id.unwrap()).await?;
r.delete_platform(platform.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn mission_crud_lifecycle() -> Result<(), HorizonError> {
let r = setup().await;
let m = Mission {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("TestMission-{}", Uuid::new_v4())),
start_datetime: None,
end_datetime: None,
position: None,
free_text: Some("mission text".to_owned()),
organization_id: None,
};
let created = r.insert_mission(&m).await?;
assert!(created.id.is_some());
let read = r.read_mission(created.id.unwrap()).await?.unwrap();
assert_eq!(read.name, created.name);
let updated_name = format!("UpdatedMission-{}", Uuid::new_v4());
let updated = Mission {
name: Some(updated_name.clone()),
..created
};
let upserted = r.upsert_mission(&updated).await?;
assert_eq!(upserted.name.as_deref(), Some(updated_name.as_str()));
let deleted = r.delete_mission(upserted.id.unwrap()).await?;
assert_eq!(deleted, 1);
let gone = r.read_mission(upserted.id.unwrap()).await?;
assert!(gone.is_none());
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn mission_position_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let pos = Position(-122.4194_f64, 37.7749_f64);
let m = Mission {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("PositionTest-{}", Uuid::new_v4())),
start_datetime: None,
end_datetime: None,
position: Some(pos),
free_text: None,
organization_id: None,
};
let created = r.insert_mission(&m).await?;
let read = r.read_mission(created.id.unwrap()).await?.unwrap();
assert_eq!(read.position, Some(pos));
r.delete_mission(created.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn platform_position_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let pos = Position(139.6917_f64, 35.6895_f64);
let p = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("PosPlat-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: Some(pos),
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let created = r.insert_platform(&p).await?;
let read = r.read_platform(created.id.unwrap()).await?.unwrap();
assert_eq!(read.position, Some(pos));
let updated_pos = Position(0.0_f64, 0.0_f64);
let updated = Platform {
position: Some(updated_pos),
..created
};
let upserted = r.upsert_platform(&updated).await?;
assert_eq!(upserted.position, Some(updated_pos));
r.delete_platform(upserted.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn update_platform_returns_updated_row() -> Result<(), HorizonError> {
let r = setup().await;
let original_name = format!("UpdatePlat-{}", Uuid::new_v4());
let p = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(original_name.clone()),
kind_id: None,
free_text: Some("before".to_owned()),
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let created = r.insert_platform(&p).await?;
let updated_input = Platform {
free_text: Some("after".to_owned()),
..created.clone()
};
let updated = r.update_platform(&updated_input).await?;
assert_eq!(updated.id, created.id);
assert_eq!(updated.free_text.as_deref(), Some("after"));
assert_eq!(updated.created_datetime, created.created_datetime);
r.delete_platform(created.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn update_platform_returns_not_found_for_missing_id() -> Result<(), HorizonError> {
let r = setup().await;
let nonexistent = Uuid::new_v4();
let p = Platform {
id: Some(nonexistent),
created_datetime: None,
modified_datetime: None,
name: Some("Ghost".to_owned()),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let result = r.update_platform(&p).await;
let HorizonError::Postgres(PostgresError::NotFound { entity, id }) =
result.expect_err("update should fail for missing id")
else {
panic!("expected NotFound variant");
};
assert_eq!(entity, "platform");
let EntityId::Uuid(id_uuid) = id else {
panic!("expected EntityId::Uuid for Some(uuid) input");
};
assert_eq!(id_uuid, nonexistent);
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn update_platform_returns_not_found_for_none_id() -> Result<(), HorizonError> {
let r = setup().await;
let p = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some("NoIdGhost".to_owned()),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let result = r.update_platform(&p).await;
let err = result.expect_err("update should fail for None id");
assert_eq!(err.to_string(), "platform with id None not found");
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn update_data_stream_returns_updated_row() -> Result<(), HorizonError> {
let r = setup().await;
let plat_name = format!("UpdDsPlat-{}", Uuid::new_v4());
let plat = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(plat_name),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let plat = r.insert_platform(&plat).await?;
let plat_id = plat.id.expect("inserted platform has id");
let ds_original_name = format!("UpdDs-{}", Uuid::new_v4());
let ds = DataStream {
id: None,
created_datetime: None,
modified_datetime: None,
platform_id: Some(plat_id),
organization_id: None,
name: Some(ds_original_name.clone()),
query_source: DataStreamQuerySource::default(),
};
let ds = r.insert_data_stream(&ds).await?;
let new_name = format!("{ds_original_name}-renamed");
let renamed = DataStream {
name: Some(new_name.clone()),
..ds.clone()
};
let updated = r.update_data_stream(&renamed).await?;
assert_eq!(updated.id, ds.id);
assert_eq!(updated.name.as_deref(), Some(new_name.as_str()));
r.delete_data_stream(ds.id.unwrap()).await?;
r.delete_platform(plat_id).await?;
Ok(())
}
async fn make_platform_and_stream_for_annotation(
r: &PostgresRepository,
) -> Result<(Uuid, Uuid), HorizonError> {
let p = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("AnnotationTestPlatform-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let platform_id = r.insert_platform(&p).await?.id.unwrap();
let stream = DataStream::from_platform_and_name(platform_id, "annotation-test");
let data_stream_id = r.insert_data_stream(&stream).await?.id.unwrap();
Ok((platform_id, data_stream_id))
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_annotation_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let (platform_id, data_stream_id) = make_platform_and_stream_for_annotation(&r).await?;
let t0 = Utc.with_ymd_and_hms(2026, 1, 1, 0, 0, 0).unwrap();
let t1 = Utc.with_ymd_and_hms(2026, 1, 1, 0, 1, 0).unwrap();
let annotation = Annotation {
id: None,
created_datetime: None,
modified_datetime: None,
platform_id,
data_stream_id: Some(data_stream_id),
time_coordinates: vec![t0, t1],
value_coordinates: vec![100.0_f64, 200.0_f64],
ontology_class_id: None,
notes: Some("ingested via SDK integration test".to_owned()),
confidence: Some(75.0_f64),
duration_seconds: None,
bearing_time_record_specification_id: None,
spectrogram_specification_id: None,
parent_annotation_id: None,
feed_context: Some(serde_json::json!({"source": "integration-test"})),
organization_id: None,
};
let created = r.insert_annotation(&annotation).await?;
assert!(created.id.is_some());
assert!(created.created_datetime.is_some());
assert!(created.modified_datetime.is_some());
assert_eq!(created.platform_id, platform_id);
assert_eq!(created.data_stream_id, Some(data_stream_id));
assert_eq!(created.time_coordinates, vec![t0, t1]);
assert_eq!(created.value_coordinates, vec![100.0_f64, 200.0_f64]);
assert_eq!(created.confidence, Some(75.0_f64));
assert_eq!(
created.notes.as_deref(),
Some("ingested via SDK integration test")
);
r.delete_platform(platform_id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_annotation_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let (platform_id, data_stream_id) = make_platform_and_stream_for_annotation(&r).await?;
let base = Utc.with_ymd_and_hms(2026, 2, 1, 0, 0, 0).unwrap();
let annotations: Vec<Annotation> = (0..3_i64)
.map(|i| Annotation {
id: None,
created_datetime: None,
modified_datetime: None,
platform_id,
data_stream_id: Some(data_stream_id),
time_coordinates: vec![base + chrono::Duration::seconds(i * 10)],
value_coordinates: vec![f64::from(i32::try_from(i).unwrap())],
ontology_class_id: None,
notes: None,
confidence: None,
duration_seconds: None,
bearing_time_record_specification_id: None,
spectrogram_specification_id: None,
parent_annotation_id: None,
feed_context: None,
organization_id: None,
})
.collect();
let inserted = r.insert_annotation_batch(&annotations).await?;
assert_eq!(inserted.len(), 3);
for row in &inserted {
assert!(row.id.is_some());
assert_eq!(row.platform_id, platform_id);
assert_eq!(row.data_stream_id, Some(data_stream_id));
}
r.delete_platform(platform_id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_annotation_batch_empty_is_noop() -> Result<(), HorizonError> {
let r = setup().await;
let inserted = r.insert_annotation_batch(&[]).await?;
assert!(inserted.is_empty());
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_annotation_rejects_conflicting_specifications() -> Result<(), HorizonError> {
let r = setup().await;
let (platform_id, data_stream_id) = make_platform_and_stream_for_annotation(&r).await?;
let annotation = Annotation {
id: None,
created_datetime: None,
modified_datetime: None,
platform_id,
data_stream_id: Some(data_stream_id),
time_coordinates: vec![Utc::now()],
value_coordinates: vec![1.0_f64],
ontology_class_id: None,
notes: None,
confidence: None,
duration_seconds: None,
bearing_time_record_specification_id: Some(Uuid::new_v4()),
spectrogram_specification_id: Some(Uuid::new_v4()),
parent_annotation_id: None,
feed_context: None,
organization_id: None,
};
let result = r.insert_annotation(&annotation).await;
assert!(result.is_err(), "expected CHECK constraint violation");
r.delete_platform(platform_id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn update_mission_platform_returns_not_found_for_missing_id() -> Result<(), HorizonError> {
let r = setup().await;
let mp = horizon_sdk::types::model::MissionPlatform {
id: Some(Uuid::new_v4()),
created_datetime: None,
modified_datetime: None,
mission_id: Some(Uuid::new_v4()),
platform_id: Some(Uuid::new_v4()),
organization_id: None,
};
let result = r.update_mission_platform(&mp).await;
let HorizonError::Postgres(PostgresError::NotFound { entity, .. }) =
result.expect_err("update should fail for missing mission_platform")
else {
panic!("expected NotFound variant");
};
assert_eq!(entity, "mission_platform");
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_metadata_row_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let platform = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("MetadataBatchPlatform-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let platform = r.insert_platform(&platform).await?;
let ds_input = DataStream {
id: None,
created_datetime: None,
modified_datetime: None,
platform_id: platform.id,
organization_id: None,
name: Some(format!("MetadataBatchStream-{}", Uuid::new_v4())),
query_source: DataStreamQuerySource::Postgres,
};
let ds = r.insert_data_stream(&ds_input).await?;
let data_stream_id = ds.id.unwrap();
let base = Utc.with_ymd_and_hms(2026, 3, 1, 0, 0, 0).unwrap();
let rows: Vec<MetadataRow> = (0_i32..3_i32)
.map(|i| MetadataRow {
data_stream_id,
datetime: base + chrono::Duration::seconds(i64::from(i) * 10_i64),
latitude: Some(10.0_f64 + f64::from(i)),
longitude: Some(20.0_f64 + f64::from(i)),
altitude: None,
speed: None,
heading: None,
pitch: None,
roll: None,
speed_over_ground: None,
created_datetime: None,
modified_datetime: None,
})
.collect();
r.insert_metadata_row_batch(&rows).await?;
let listed = r.list_metadata_rows(data_stream_id).await?;
assert_eq!(listed.len(), 3);
r.delete_platform(platform.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_metadata_row_batch_empty_is_noop() -> Result<(), HorizonError> {
let r = setup().await;
r.insert_metadata_row_batch(&[]).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_data_row_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let platform = Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("DataRowBatchPlatform-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
};
let platform = r.insert_platform(&platform).await?;
let ds_input = DataStream {
id: None,
created_datetime: None,
modified_datetime: None,
platform_id: platform.id,
organization_id: None,
name: Some(format!("DataRowBatchStream-{}", Uuid::new_v4())),
query_source: DataStreamQuerySource::Postgres,
};
let ds = r.insert_data_stream(&ds_input).await?;
let data_stream_id = ds.id.unwrap();
let specification_id = Uuid::new_v4();
let base = Utc.with_ymd_and_hms(2026, 3, 1, 0, 0, 0).unwrap();
let typed_payload_case_list = [
DataRowPayload::Float32(1.5_f32),
DataRowPayload::Float32Array(vec![1.0_f32, 2.0_f32]),
DataRowPayload::Float64(22.5_f64),
DataRowPayload::Float64Array(vec![3.0_f64, 4.0_f64]),
DataRowPayload::Int(42_i32),
DataRowPayload::IntArray(vec![1_i32, 2_i32, 3_i32]),
DataRowPayload::String("alpha".to_owned()),
DataRowPayload::Struct(serde_json::json!({"callsign": "ALPHA"})),
DataRowPayload::Vector {
vector: vec![9.0_f64, 9.5_f64],
vector_end_bound: 1.0_f64,
vector_start_bound: 0.0_f64,
},
];
let mut rows: Vec<VariantDataRow> = typed_payload_case_list
.iter()
.enumerate()
.map(|(index, payload)| {
VariantDataRow::from_payload(
data_stream_id,
base + chrono::Duration::seconds(i64::try_from(index).unwrap()),
"spectrogram".to_owned(),
specification_id,
payload.clone(),
)
})
.collect();
let vector_batch_offset = i64::try_from(typed_payload_case_list.len()).unwrap();
rows.extend((0_i32..20_i32).map(|index| {
VariantDataRow::from_payload(
data_stream_id,
base + chrono::Duration::seconds(vector_batch_offset + i64::from(index)),
"spectrogram".to_owned(),
specification_id,
DataRowPayload::Vector {
vector: vec![f64::from(index), f64::from(index) + 0.5_f64],
vector_end_bound: 1.0_f64,
vector_start_bound: 0.0_f64,
},
)
}));
let expected_count = u64::try_from(rows.len()).unwrap();
let inserted = r.insert_variant_data_row_batch(&rows).await?;
assert_eq!(inserted, expected_count);
let listed = r.list_variant_data_rows(data_stream_id).await?;
assert_eq!(listed.len(), rows.len());
for (index, expected_payload) in typed_payload_case_list.iter().enumerate() {
assert_data_row_payload_round_trip(&listed[index], expected_payload);
}
for index in 0_i32..20_i32 {
let listed_index = typed_payload_case_list.len() + usize::try_from(index).unwrap();
assert_data_row_payload_round_trip(
&listed[listed_index],
&DataRowPayload::Vector {
vector: vec![f64::from(index), f64::from(index) + 0.5_f64],
vector_end_bound: 1.0_f64,
vector_start_bound: 0.0_f64,
},
);
}
r.delete_platform(platform.id.unwrap()).await?;
Ok(())
}
fn assert_data_row_payload_round_trip(row: &VariantDataRow, expected: &DataRowPayload) {
assert_eq!(row.variant.as_deref(), Some(expected.variant().as_str()));
assert_eq!(row.payload().as_ref(), Some(expected));
let active_arm = expected.variant().payload_arm();
assert_eq!(
row.vector.is_some(),
active_arm == "vector",
"vector arm mismatch for variant {}",
expected.variant().as_str()
);
assert_eq!(
row.payload_int.is_some(),
active_arm == "payload_int",
"payload_int arm mismatch for variant {}",
expected.variant().as_str()
);
assert_eq!(
row.payload_int_array.is_some(),
active_arm == "payload_int_array",
"payload_int_array arm mismatch for variant {}",
expected.variant().as_str()
);
assert_eq!(
row.payload_float32.is_some(),
active_arm == "payload_float32",
"payload_float32 arm mismatch for variant {}",
expected.variant().as_str()
);
assert_eq!(
row.payload_float32_array.is_some(),
active_arm == "payload_float32_array",
"payload_float32_array arm mismatch for variant {}",
expected.variant().as_str()
);
assert_eq!(
row.payload_float64.is_some(),
active_arm == "payload_float64",
"payload_float64 arm mismatch for variant {}",
expected.variant().as_str()
);
assert_eq!(
row.payload_float64_array.is_some(),
active_arm == "payload_float64_array",
"payload_float64_array arm mismatch for variant {}",
expected.variant().as_str()
);
assert_eq!(
row.payload_string.is_some(),
active_arm == "payload_string",
"payload_string arm mismatch for variant {}",
expected.variant().as_str()
);
assert_eq!(
row.payload_struct.is_some(),
active_arm == "payload_struct",
"payload_struct arm mismatch for variant {}",
expected.variant().as_str()
);
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_data_row_batch_empty_is_noop() -> Result<(), HorizonError> {
let r = setup().await;
let inserted = r.insert_data_row_batch(&[]).await?;
assert_eq!(inserted, 0_u64);
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_annotation_batch_rolls_back_on_failure() -> Result<(), HorizonError> {
let r = setup().await;
let (platform_id, data_stream_id) = make_platform_and_stream_for_annotation(&r).await?;
let base = Utc.with_ymd_and_hms(2026, 5, 1, 0, 0, 0).unwrap();
let unique_first_id = Uuid::new_v4();
let conflicting_id = Uuid::new_v4();
let pre = Annotation {
id: Some(conflicting_id),
created_datetime: None,
modified_datetime: None,
platform_id,
data_stream_id: Some(data_stream_id),
time_coordinates: vec![base],
value_coordinates: vec![0.0_f64],
ontology_class_id: None,
notes: None,
confidence: None,
duration_seconds: None,
bearing_time_record_specification_id: None,
spectrogram_specification_id: None,
parent_annotation_id: None,
feed_context: None,
organization_id: None,
};
r.insert_annotation(&pre).await?;
let batch = vec![
Annotation {
id: Some(unique_first_id),
created_datetime: None,
modified_datetime: None,
platform_id,
data_stream_id: Some(data_stream_id),
time_coordinates: vec![base + chrono::Duration::seconds(10)],
value_coordinates: vec![1.0_f64],
ontology_class_id: None,
notes: None,
confidence: None,
duration_seconds: None,
bearing_time_record_specification_id: None,
spectrogram_specification_id: None,
parent_annotation_id: None,
feed_context: None,
organization_id: None,
},
Annotation {
id: Some(conflicting_id),
created_datetime: None,
modified_datetime: None,
platform_id,
data_stream_id: Some(data_stream_id),
time_coordinates: vec![base + chrono::Duration::seconds(20)],
value_coordinates: vec![2.0_f64],
ontology_class_id: None,
notes: None,
confidence: None,
duration_seconds: None,
bearing_time_record_specification_id: None,
spectrogram_specification_id: None,
parent_annotation_id: None,
feed_context: None,
organization_id: None,
},
];
let result = r.insert_annotation_batch(&batch).await;
assert!(result.is_err(), "expected unique-constraint violation");
let pool_url = std::env::var("DATABASE_URL").unwrap_or_else(|_| {
"postgresql://horizon_owner:insecure_password@localhost:5432/horizon?sslmode=disable"
.to_owned()
});
let pool = sqlx::postgres::PgPoolOptions::new()
.max_connections(1)
.connect(&pool_url)
.await
.expect("Failed to connect for rollback verification");
let leftover: Option<(Uuid,)> =
sqlx::query_as("SELECT id FROM horizon_public.annotation WHERE id = $1")
.bind(unique_first_id)
.fetch_optional(&pool)
.await
.expect("rollback verification query failed");
assert!(
leftover.is_none(),
"transaction should have rolled back the first batch row"
);
r.delete_platform(platform_id).await?;
Ok(())
}
fn make_platform(name_suffix: &str) -> Platform {
Platform {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("{name_suffix}-{}", Uuid::new_v4())),
kind_id: None,
free_text: None,
position: None,
organization_id: None,
start_datetime: None,
end_datetime: None,
}
}
fn make_bearing_time_record_specification(name_suffix: &str) -> BearingTimeRecordSpecification {
BearingTimeRecordSpecification {
id: None,
amplitude_unit_mode: None,
baffle_bearing_list: Some(vec![0.0_f64, 1.0_f64]),
bearing_bin_count: Some(360_i32),
elevation_increment: Some(0.1_f64),
focus_range: None,
frequency_spacing: Some(1.0_f64),
heading_data_type: None,
heading_vector_index: Some(1_i32),
lower_elevation: Some(0.0_f64),
max_frequency: Some(10_000.0_f64),
max_pixel: Some(1024_i32),
min_frequency: Some(100.0_f64),
name: Some(format!("{name_suffix}-{}", Uuid::new_v4())),
normalizer: Some("none".to_owned()),
fft_sample_count: Some(1024_i32),
organization_id: None,
upper_elevation: Some(90.0_f64),
update_rate_ms: Some(1000_i64),
created_datetime: None,
modified_datetime: None,
}
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_platform_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let batch: Vec<Platform> = (0..5_i32)
.map(|i| make_platform(&format!("BatchPlatform-{i}")))
.collect();
let inserted = r.insert_platform_batch(&batch).await?;
assert_eq!(inserted.len(), 5);
for (input, row) in batch.iter().zip(&inserted) {
assert!(row.id.is_some());
assert!(row.created_datetime.is_some());
assert_eq!(row.name, input.name);
r.delete_platform(row.id.unwrap()).await?;
}
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_platform_batch_empty_is_noop() -> Result<(), HorizonError> {
let r = setup().await;
let inserted = r.insert_platform_batch(&[]).await?;
assert!(inserted.is_empty());
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_platform_batch_rolls_back_on_failure() -> Result<(), HorizonError> {
let r = setup().await;
let conflicting_id = Uuid::new_v4();
let pre = Platform {
id: Some(conflicting_id),
..make_platform("RollbackPre")
};
r.insert_platform(&pre).await?;
let unique_first_id = Uuid::new_v4();
let batch = vec![
Platform {
id: Some(unique_first_id),
..make_platform("RollbackFirst")
},
Platform {
id: Some(conflicting_id),
..make_platform("RollbackCollide")
},
];
let result = r.insert_platform_batch(&batch).await;
assert!(result.is_err(), "expected unique-constraint violation");
let leftover = r.read_platform(unique_first_id).await?;
assert!(
leftover.is_none(),
"transaction should have rolled back the first row"
);
r.delete_platform(conflicting_id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_data_stream_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let platform = r.insert_platform(&make_platform("DataStreamBatch")).await?;
let platform_id = platform.id.unwrap();
let batch: Vec<DataStream> = (0..3_i32)
.map(|i| DataStream {
id: None,
created_datetime: None,
modified_datetime: None,
platform_id: Some(platform_id),
organization_id: None,
name: Some(format!("Stream-{i}")),
query_source: DataStreamQuerySource::Postgres,
})
.collect();
let inserted = r.insert_data_stream_batch(&batch).await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert!(row.id.is_some());
assert_eq!(row.platform_id, Some(platform_id));
assert_eq!(row.name, input.name);
}
r.delete_platform(platform_id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_data_stream_batch_empty_is_noop() -> Result<(), HorizonError> {
let r = setup().await;
let inserted = r.insert_data_stream_batch(&[]).await?;
assert!(inserted.is_empty());
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_mission_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let batch: Vec<Mission> = (0..3_i32)
.map(|i| Mission {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("BatchMission-{i}-{}", Uuid::new_v4())),
start_datetime: Some(Utc.with_ymd_and_hms(2026, 1, 1, 0, 0, 0).unwrap()),
end_datetime: None,
free_text: None,
organization_id: None,
position: Some(Position(f64::from(i), f64::from(i + 10))),
})
.collect();
let inserted = r.insert_mission_batch(&batch).await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.name, input.name);
assert_eq!(row.position, input.position);
r.delete_mission(row.id.expect("mission id")).await?;
}
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_mission_batch_empty_is_noop() -> Result<(), HorizonError> {
let r = setup().await;
let inserted = r.insert_mission_batch(&[]).await?;
assert!(inserted.is_empty());
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_ontology_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let batch: Vec<Ontology> = (0..3_i32)
.map(|i| Ontology {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("BatchOntology-{i}-{}", Uuid::new_v4())),
description: Some(format!("description {i}")),
organization_id: None,
})
.collect();
let inserted = r.insert_ontology_batch(&batch).await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.name, input.name);
assert_eq!(row.description, input.description);
r.delete_ontology(row.id.expect("ontology id")).await?;
}
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_ontology_class_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let ontology = r
.insert_ontology(&Ontology {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("BatchOntologyClassRoot-{}", Uuid::new_v4())),
description: None,
organization_id: None,
})
.await?;
let ontology_id = ontology.id.unwrap();
let batch: Vec<OntologyClass> = (0..3_i32)
.map(|i| OntologyClass {
id: None,
created_datetime: None,
modified_datetime: None,
ontology_id,
parent_id: None,
name: Some(format!("Class-{i}")),
description: None,
order: Some(i),
organization_id: None,
relationship_type: None,
})
.collect();
let inserted = r.insert_ontology_class_batch(&batch).await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.name, input.name);
assert_eq!(row.order, input.order);
}
r.delete_ontology(ontology_id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_mission_platform_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let platform = r
.insert_platform(&make_platform("MissionPlatformBatch"))
.await?;
let mission = r
.insert_mission(&Mission {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("MPBatch-{}", Uuid::new_v4())),
start_datetime: None,
end_datetime: None,
free_text: None,
organization_id: None,
position: None,
})
.await?;
let mp_batch = vec![MissionPlatform {
id: None,
created_datetime: None,
modified_datetime: None,
mission_id: mission.id,
platform_id: platform.id,
organization_id: None,
}];
let inserted = r.insert_mission_platform_batch(&mp_batch).await?;
assert_eq!(inserted.len(), 1);
assert_eq!(inserted[0].mission_id, mission.id);
assert_eq!(inserted[0].platform_id, platform.id);
r.delete_mission(mission.id.unwrap()).await?;
r.delete_platform(platform.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_audio_specification_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let batch: Vec<AudioSpecification> = (0..3_i32)
.map(|i| AudioSpecification {
id: None,
created_datetime: None,
modified_datetime: None,
sample_rate: 48_000_i64,
bit_depth: 16_i64,
channel_count: 4_i64,
channel_index: i64::from(i),
encoding: "pcm_float".to_owned(),
organization_id: None,
})
.collect();
let inserted = r.insert_audio_specification_batch(&batch).await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.channel_index, input.channel_index);
assert_eq!(row.encoding, input.encoding);
r.delete_audio_specification(row.id.unwrap()).await?;
}
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_spectrogram_specification_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let batch: Vec<SpectrogramSpecification> = (0..3_i32)
.map(|i| SpectrogramSpecification {
id: None,
amplitude_unit_mode: None,
channel: Some(i),
frequency_spacing: Some(1.0_f64),
name: Some(format!("SpecBatch-{i}")),
fft_sample_count: Some(1024_i32),
fft_sample_overlap_count: Some(512_i32),
frequency_bin_count: Some(513_i32),
created_datetime: None,
modified_datetime: None,
normalizer: Some("split_window".to_owned()),
organization_id: None,
channel_role: None,
})
.collect();
let inserted = r.insert_spectrogram_specification_batch(&batch).await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.channel, input.channel);
assert_eq!(row.name, input.name);
assert_eq!(row.frequency_bin_count, input.frequency_bin_count);
assert_eq!(row.normalizer, input.normalizer);
r.delete_spectrogram_specification(row.id.unwrap()).await?;
}
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn spectrogram_specification_crud_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let spec = SpectrogramSpecification {
amplitude_unit_mode: None,
channel: Some(0_i32),
channel_role: None,
created_datetime: None,
fft_sample_count: Some(1024_i32),
fft_sample_overlap_count: Some(512_i32),
frequency_bin_count: Some(513_i32),
frequency_spacing: Some(1.0_f64),
id: None,
modified_datetime: None,
name: Some(format!("SpecCrud-{}", Uuid::new_v4())),
normalizer: Some("split_window".to_owned()),
organization_id: None,
};
let created = r.insert_spectrogram_specification(&spec).await?;
assert_eq!(created.frequency_bin_count, Some(513_i32));
assert_eq!(created.normalizer.as_deref(), Some("split_window"));
let read = r
.read_spectrogram_specification(created.id.unwrap())
.await?
.unwrap();
assert_eq!(read.frequency_bin_count, created.frequency_bin_count);
assert_eq!(read.normalizer, created.normalizer);
let listed = r
.list_spectrogram_specifications(&SpectrogramSpecificationFilter::default())
.await?;
let found = listed
.iter()
.find(|specification| specification.id == created.id)
.expect("created specification should appear in the list");
assert_eq!(found.frequency_bin_count, created.frequency_bin_count);
assert_eq!(found.normalizer, created.normalizer);
let updated = r
.update_spectrogram_specification(&SpectrogramSpecification {
frequency_bin_count: Some(257_i32),
normalizer: Some("none".to_owned()),
..created
})
.await?;
assert_eq!(updated.frequency_bin_count, Some(257_i32));
assert_eq!(updated.normalizer.as_deref(), Some("none"));
let upserted = r
.upsert_spectrogram_specification(&SpectrogramSpecification {
frequency_bin_count: Some(129_i32),
normalizer: None,
..updated
})
.await?;
assert_eq!(upserted.frequency_bin_count, Some(129_i32));
assert_eq!(upserted.normalizer, None);
r.delete_spectrogram_specification(upserted.id.unwrap())
.await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_directional_spectrogram_specification_batch_round_trip() -> Result<(), HorizonError>
{
let r = setup().await;
let parent = r
.insert_spectrogram_specification(&SpectrogramSpecification {
id: None,
amplitude_unit_mode: None,
channel: Some(0),
channel_role: None,
frequency_spacing: Some(1.0_f64),
name: Some(format!("DirectionalParent-{}", Uuid::new_v4())),
fft_sample_count: Some(1024_i32),
fft_sample_overlap_count: Some(512_i32),
frequency_bin_count: None,
created_datetime: None,
modified_datetime: None,
normalizer: None,
organization_id: None,
})
.await?;
let parent_id = parent.id.unwrap();
let batch: Vec<DirectionalSpectrogramSpecification> = (0..3_i32)
.map(|i| DirectionalSpectrogramSpecification {
id: None,
bearing_leading_rows: Some(i64::from(i)),
bearing_lagging_rows: Some(i64::from(i + 1)),
spectrogram_specification_id: Some(parent_id),
created_datetime: None,
modified_datetime: None,
organization_id: None,
})
.collect();
let inserted = r
.insert_directional_spectrogram_specification_batch(&batch)
.await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.bearing_leading_rows, input.bearing_leading_rows);
assert_eq!(row.bearing_lagging_rows, input.bearing_lagging_rows);
r.delete_directional_spectrogram_specification(row.id.unwrap())
.await?;
}
r.delete_spectrogram_specification(parent_id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_beamgram_specification_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let batch: Vec<BeamgramSpecification> = (0..3_i32)
.map(|i| BeamgramSpecification {
id: None,
center_bearing: Some(f64::from(i)),
center_bin_width: Some(1.0_f64),
elevation_increment: Some(0.1_f64),
lower_elevation: Some(0.0_f64),
upper_elevation: Some(90.0_f64),
max_frequency: Some(10_000.0_f64),
min_frequency: Some(100.0_f64),
name: Some(format!("BeamBatch-{i}")),
fft_sample_count: Some(1024_i32),
organization_id: None,
update_rate_ms: Some(1000_i64),
normalizer: Some("none".to_owned()),
created_datetime: None,
modified_datetime: None,
})
.collect();
let inserted = r.insert_beamgram_specification_batch(&batch).await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.center_bearing, input.center_bearing);
assert_eq!(row.name, input.name);
assert_eq!(row.update_rate_ms, Some(1000_i64));
r.delete_beamgram_specification(row.id.unwrap()).await?;
}
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_bearing_time_record_specification_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let batch: Vec<BearingTimeRecordSpecification> = (0..3_i32)
.map(|i| BearingTimeRecordSpecification {
id: None,
amplitude_unit_mode: None,
baffle_bearing_list: Some(vec![f64::from(i), f64::from(i) + 1.0_f64]),
bearing_bin_count: Some(360_i32),
elevation_increment: Some(0.1_f64),
focus_range: None,
frequency_spacing: Some(1.0_f64),
heading_data_type: None,
heading_vector_index: Some(1_i32),
lower_elevation: Some(0.0_f64),
max_frequency: Some(10_000.0_f64),
max_pixel: Some(1024_i32),
min_frequency: Some(100.0_f64),
name: Some(format!("BTRBatch-{i}")),
normalizer: Some("none".to_owned()),
fft_sample_count: Some(1024_i32),
organization_id: None,
upper_elevation: Some(90.0_f64),
update_rate_ms: Some(1000_i64),
created_datetime: None,
modified_datetime: None,
})
.collect();
let inserted = r
.insert_bearing_time_record_specification_batch(&batch)
.await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.name, input.name);
assert_eq!(row.baffle_bearing_list, input.baffle_bearing_list);
r.delete_bearing_time_record_specification(row.id.unwrap())
.await?;
}
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_bearing_time_record_specification_batch_rolls_back_on_failure()
-> Result<(), HorizonError> {
let r = setup().await;
let conflicting_id = Uuid::new_v4();
let pre = BearingTimeRecordSpecification {
id: Some(conflicting_id),
..make_bearing_time_record_specification("RollbackPre")
};
r.insert_bearing_time_record_specification(&pre).await?;
let unique_first_id = Uuid::new_v4();
let batch = vec![
BearingTimeRecordSpecification {
id: Some(unique_first_id),
..make_bearing_time_record_specification("RollbackFirst")
},
BearingTimeRecordSpecification {
id: Some(conflicting_id),
..make_bearing_time_record_specification("RollbackCollide")
},
];
let result = r
.insert_bearing_time_record_specification_batch(&batch)
.await;
assert!(result.is_err(), "expected unique-constraint violation");
let leftover = r
.read_bearing_time_record_specification(unique_first_id)
.await?;
assert!(
leftover.is_none(),
"transaction should have rolled back the first row"
);
r.delete_bearing_time_record_specification(conflicting_id)
.await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_platform_audio_specification_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let platform = r.insert_platform(&make_platform("PASBatch")).await?;
let spec = r
.insert_audio_specification(&AudioSpecification {
id: None,
created_datetime: None,
modified_datetime: None,
sample_rate: 48_000_i64,
bit_depth: 16_i64,
channel_count: 1_i64,
channel_index: 0_i64,
encoding: "pcm_float".to_owned(),
organization_id: None,
})
.await?;
let batch = vec![PlatformAudioSpecification {
id: None,
created_datetime: None,
modified_datetime: None,
platform_id: platform.id.unwrap(),
audio_specification_id: spec.id.unwrap(),
organization_id: None,
}];
let inserted = r.insert_platform_audio_specification_batch(&batch).await?;
assert_eq!(inserted.len(), 1);
assert_eq!(inserted[0].platform_id, platform.id.unwrap());
assert_eq!(inserted[0].audio_specification_id, spec.id.unwrap());
r.delete_platform_audio_specification(inserted[0].id.unwrap())
.await?;
r.delete_audio_specification(spec.id.unwrap()).await?;
r.delete_platform(platform.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_platform_kind_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let batch: Vec<PlatformKind> = (0..3_i32)
.map(|i| PlatformKind {
id: None,
created_datetime: None,
modified_datetime: None,
name: Some(format!("KindBatch-{i}-{}", Uuid::new_v4())),
image_url: None,
short_description: None,
long_description: None,
organization_id: None,
})
.collect();
let inserted = r.insert_platform_kind_batch(&batch).await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.name, input.name);
r.delete_platform_kind(row.id.unwrap()).await?;
}
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_platform_information_batch_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let platform = r.insert_platform(&make_platform("InfoBatch")).await?;
let platform_id = platform.id.unwrap();
let batch: Vec<PlatformInformation> = (0..3_i32)
.map(|i| PlatformInformation {
id: None,
created_datetime: None,
modified_datetime: None,
platform_id,
organization_id: None,
properties: serde_json::json!({"index": i}),
})
.collect();
let inserted = r.insert_platform_information_batch(&batch).await?;
assert_eq!(inserted.len(), 3);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.properties, input.properties);
}
r.delete_platform(platform_id).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_platform_batch_large_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let batch: Vec<Platform> = (0..100_i32)
.map(|i| make_platform(&format!("Large-{i}")))
.collect();
let inserted = r.insert_platform_batch(&batch).await?;
assert_eq!(inserted.len(), 100);
for (input, row) in batch.iter().zip(&inserted) {
assert_eq!(row.name, input.name);
r.delete_platform(row.id.unwrap()).await?;
}
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_batches_empty_input_are_noops() -> Result<(), HorizonError> {
let r = setup().await;
assert!(r.insert_ontology_batch(&[]).await?.is_empty());
assert!(r.insert_ontology_class_batch(&[]).await?.is_empty());
assert!(r.insert_mission_platform_batch(&[]).await?.is_empty());
assert!(r.insert_audio_specification_batch(&[]).await?.is_empty());
assert!(
r.insert_spectrogram_specification_batch(&[])
.await?
.is_empty()
);
assert!(
r.insert_directional_spectrogram_specification_batch(&[])
.await?
.is_empty()
);
assert!(r.insert_beamgram_specification_batch(&[]).await?.is_empty());
assert!(
r.insert_bearing_time_record_specification_batch(&[])
.await?
.is_empty()
);
assert!(
r.insert_platform_audio_specification_batch(&[])
.await?
.is_empty()
);
assert!(
r.insert_platform_beamgram_specification_batch(&[])
.await?
.is_empty()
);
assert!(
r.insert_platform_bearing_time_record_specification_batch(&[])
.await?
.is_empty()
);
assert!(
r.insert_platform_spectrogram_specification_batch(&[])
.await?
.is_empty()
);
assert!(r.insert_platform_kind_batch(&[]).await?.is_empty());
assert!(r.insert_platform_information_batch(&[]).await?.is_empty());
Ok(())
}
async fn any_organization_id() -> Uuid {
let url = std::env::var("DATABASE_URL").unwrap_or_else(|_| {
"postgresql://horizon_owner:insecure_password@localhost:5432/horizon?sslmode=disable"
.to_owned()
});
let pool = PgPoolOptions::new()
.max_connections(1)
.connect(&url)
.await
.expect("Failed to connect to test database");
sqlx::query_scalar::<_, Uuid>("SELECT id FROM horizon_public.organization LIMIT 1")
.fetch_one(&pool)
.await
.expect("an organization must exist for org-scoped alias tests")
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn bearing_resolution_alias_crud_lifecycle() -> Result<(), HorizonError> {
let r = setup().await;
let alias = BearingResolutionAlias {
bearing_bin_count: Some(180),
created_datetime: None,
description: Some("initial".to_owned()),
id: None,
modified_datetime: None,
name: Some(format!("BearingRes-{}", Uuid::new_v4())),
organization_id: None,
};
let created = r.insert_bearing_resolution_alias(&alias).await?;
assert!(created.id.is_some());
assert_eq!(created.bearing_bin_count, Some(180));
let read = r
.read_bearing_resolution_alias(created.id.unwrap())
.await?
.unwrap();
assert_eq!(read.name, created.name);
let updated = BearingResolutionAlias {
bearing_bin_count: Some(360),
..created
};
let upserted = r.upsert_bearing_resolution_alias(&updated).await?;
assert_eq!(upserted.bearing_bin_count, Some(360));
let deleted = r
.delete_bearing_resolution_alias(upserted.id.unwrap())
.await?;
assert_eq!(deleted, 1);
assert!(
r.read_bearing_resolution_alias(upserted.id.unwrap())
.await?
.is_none()
);
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn elevation_alias_crud_lifecycle() -> Result<(), HorizonError> {
let r = setup().await;
let alias = ElevationAlias {
created_datetime: None,
description: Some("initial".to_owned()),
elevation_increment: Some(1.0),
id: None,
lower_elevation: Some(-10.0),
modified_datetime: None,
name: Some(format!("Elev-{}", Uuid::new_v4())),
organization_id: None,
upper_elevation: Some(20.0),
};
let created = r.insert_elevation_alias(&alias).await?;
assert!(created.id.is_some());
let read = r.read_elevation_alias(created.id.unwrap()).await?.unwrap();
assert_eq!(read.lower_elevation, Some(-10.0));
let updated = ElevationAlias {
upper_elevation: Some(30.0),
..created
};
let upserted = r.upsert_elevation_alias(&updated).await?;
assert_eq!(upserted.upper_elevation, Some(30.0));
let deleted = r.delete_elevation_alias(upserted.id.unwrap()).await?;
assert_eq!(deleted, 1);
assert!(
r.read_elevation_alias(upserted.id.unwrap())
.await?
.is_none()
);
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn focus_range_crud_lifecycle() -> Result<(), HorizonError> {
let r = setup().await;
let focus_range = FocusRange {
created_datetime: None,
description: Some("initial".to_owned()),
id: None,
modified_datetime: None,
name: Some(format!("Focus-{}", Uuid::new_v4())),
organization_id: None,
};
let created = r.insert_focus_range(&focus_range).await?;
assert!(created.id.is_some());
let read = r.read_focus_range(created.id.unwrap()).await?.unwrap();
assert_eq!(read.name, created.name);
let updated_name = format!("UpdatedFocus-{}", Uuid::new_v4());
let updated = FocusRange {
name: Some(updated_name.clone()),
..created
};
let upserted = r.upsert_focus_range(&updated).await?;
assert_eq!(upserted.name.as_deref(), Some(updated_name.as_str()));
let deleted = r.delete_focus_range(upserted.id.unwrap()).await?;
assert_eq!(deleted, 1);
assert!(r.read_focus_range(upserted.id.unwrap()).await?.is_none());
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn frequency_band_alias_crud_lifecycle() -> Result<(), HorizonError> {
let r = setup().await;
let alias = FrequencyBandAlias {
created_datetime: None,
description: Some("initial".to_owned()),
id: None,
lower_frequency: Some(100.0),
modified_datetime: None,
name: Some(format!("FBand-{}", Uuid::new_v4())),
organization_id: None,
upper_frequency: Some(200.0),
};
let created = r.insert_frequency_band_alias(&alias).await?;
assert!(created.id.is_some());
let read = r
.read_frequency_band_alias(created.id.unwrap())
.await?
.unwrap();
assert_eq!(read.upper_frequency, Some(200.0));
let updated = FrequencyBandAlias {
upper_frequency: Some(300.0),
..created
};
let upserted = r.upsert_frequency_band_alias(&updated).await?;
assert_eq!(upserted.upper_frequency, Some(300.0));
let deleted = r.delete_frequency_band_alias(upserted.id.unwrap()).await?;
assert_eq!(deleted, 1);
assert!(
r.read_frequency_band_alias(upserted.id.unwrap())
.await?
.is_none()
);
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn frequency_resolution_alias_crud_lifecycle() -> Result<(), HorizonError> {
let r = setup().await;
let alias = FrequencyResolutionAlias {
created_datetime: None,
description: Some("initial".to_owned()),
frequency_bin_count: Some(256),
id: None,
modified_datetime: None,
name: Some(format!("FRes-{}", Uuid::new_v4())),
organization_id: None,
};
let created = r.insert_frequency_resolution_alias(&alias).await?;
assert!(created.id.is_some());
assert_eq!(created.frequency_bin_count, Some(256));
let read = r
.read_frequency_resolution_alias(created.id.unwrap())
.await?
.unwrap();
assert_eq!(read.name, created.name);
let updated = FrequencyResolutionAlias {
frequency_bin_count: Some(512),
..created
};
let upserted = r.upsert_frequency_resolution_alias(&updated).await?;
assert_eq!(upserted.frequency_bin_count, Some(512));
let deleted = r
.delete_frequency_resolution_alias(upserted.id.unwrap())
.await?;
assert_eq!(deleted, 1);
assert!(
r.read_frequency_resolution_alias(upserted.id.unwrap())
.await?
.is_none()
);
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn update_rate_alias_crud_lifecycle() -> Result<(), HorizonError> {
let r = setup().await;
let org = any_organization_id().await;
let alias = UpdateRateAlias {
created_datetime: None,
description: Some("initial".to_owned()),
id: None,
modified_datetime: None,
name: Some(format!("UpdRate-{}", Uuid::new_v4())),
organization_id: Some(org),
update_rate_ms: Some(1000),
};
let created = r.insert_update_rate_alias(&alias).await?;
assert!(created.id.is_some());
assert_eq!(created.update_rate_ms, Some(1000));
let read = r
.read_update_rate_alias(created.id.unwrap())
.await?
.unwrap();
assert_eq!(read.update_rate_ms, Some(1000));
let updated = UpdateRateAlias {
update_rate_ms: Some(2000),
..created
};
let upserted = r.upsert_update_rate_alias(&updated).await?;
assert_eq!(upserted.update_rate_ms, Some(2000));
let deleted = r.delete_update_rate_alias(upserted.id.unwrap()).await?;
assert_eq!(deleted, 1);
assert!(
r.read_update_rate_alias(upserted.id.unwrap())
.await?
.is_none()
);
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn normalizer_crud_and_archive_lifecycle() -> Result<(), HorizonError> {
let r = setup().await;
let org = any_organization_id().await;
let normalizer = Normalizer {
created_datetime: None,
description: Some("initial".to_owned()),
id: None,
is_archived: None,
modified_datetime: None,
name: Some(format!("Norm-{}", Uuid::new_v4())),
organization_id: Some(org),
};
let created = r.insert_normalizer(&normalizer).await?;
assert!(created.id.is_some());
assert_eq!(created.is_archived, Some(false));
let read = r.read_normalizer(created.id.unwrap()).await?.unwrap();
assert_eq!(read.name, created.name);
let updated = Normalizer {
description: Some("updated".to_owned()),
..created
};
let updated_res = r.update_normalizer(&updated).await?;
assert_eq!(updated_res.description.as_deref(), Some("updated"));
let archived = r.archive_normalizer(updated_res.id.unwrap()).await?;
assert_eq!(archived.is_archived, Some(true));
let reread = r.read_normalizer(archived.id.unwrap()).await?.unwrap();
assert_eq!(reread.is_archived, Some(true));
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn list_audio_specifications_applies_filters() -> Result<(), HorizonError> {
let r = setup().await;
let sample_rate = i64::from(Uuid::new_v4().as_fields().0.max(1));
let created = r
.insert_audio_specification(&AudioSpecification {
id: None,
created_datetime: None,
modified_datetime: None,
sample_rate,
bit_depth: 24_i64,
channel_count: 2_i64,
channel_index: 0_i64,
encoding: "pcm_float".to_owned(),
organization_id: None,
})
.await?;
let matched = r
.list_audio_specifications(&AudioSpecificationFilter {
sample_rate: Some(sample_rate),
..Default::default()
})
.await?;
assert_eq!(matched.len(), 1);
assert_eq!(matched[0].id, created.id);
let combined = r
.list_audio_specifications(&AudioSpecificationFilter {
bit_depth: Some(24_i64),
encoding: Some("pcm_float".to_owned()),
sample_rate: Some(sample_rate),
..Default::default()
})
.await?;
assert_eq!(combined.len(), 1);
assert_eq!(combined[0].id, created.id);
let conflicting = r
.list_audio_specifications(&AudioSpecificationFilter {
bit_depth: Some(7_i64),
sample_rate: Some(sample_rate),
..Default::default()
})
.await?;
assert!(conflicting.is_empty());
let unfiltered = r
.list_audio_specifications(&AudioSpecificationFilter::default())
.await?;
assert!(unfiltered.iter().any(|row| row.id == created.id));
r.delete_audio_specification(created.id.unwrap()).await?;
Ok(())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn list_beamgram_specifications_filters_on_update_rate() -> Result<(), HorizonError> {
let r = setup().await;
let name = format!("Beamgram-{}", Uuid::new_v4());
let created = r
.insert_beamgram_specification(&BeamgramSpecification {
center_bearing: Some(90.0),
center_bin_width: Some(1.5),
created_datetime: None,
elevation_increment: Some(1.0),
fft_sample_count: Some(1024),
id: None,
lower_elevation: Some(-10.0),
max_frequency: Some(2000.0),
min_frequency: Some(0.0),
modified_datetime: None,
name: Some(name.clone()),
normalizer: None,
organization_id: None,
update_rate_ms: Some(250_i64),
upper_elevation: Some(10.0),
})
.await?;
let matched = r
.list_beamgram_specifications(&BeamgramSpecificationFilter {
name: Some(name.clone()),
update_rate_ms: Some(250_i64),
..Default::default()
})
.await?;
assert_eq!(matched.len(), 1);
assert_eq!(matched[0].id, created.id);
let wrong_rate = r
.list_beamgram_specifications(&BeamgramSpecificationFilter {
name: Some(name),
update_rate_ms: Some(500_i64),
..Default::default()
})
.await?;
assert!(wrong_rate.is_empty());
r.delete_beamgram_specification(created.id.unwrap()).await?;
Ok(())
}