#![allow(
clippy::absolute_paths,
clippy::expect_used,
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, and panic freely"
)]
use chrono::{TimeZone as _, Utc};
use horizon_sdk::postgres::repository::PostgresRepository;
use horizon_sdk::types::error::{EntityId, HorizonError, PostgresError};
use horizon_sdk::types::model::{
Annotation, DataRow, DataStream, DataStreamQuerySource, MetadataRow, Mission, Platform,
Position,
};
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_for_annotation(r: &PostgresRepository) -> Result<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,
};
Ok(r.insert_platform(&p).await?.id.unwrap())
}
#[tokio::test]
#[ignore = "requires live Postgres"]
async fn insert_annotation_round_trip() -> Result<(), HorizonError> {
let r = setup().await;
let platform_id = make_platform_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,
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.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 = make_platform_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,
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);
}
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 = make_platform_for_annotation(&r).await?;
let annotation = Annotation {
id: None,
created_datetime: None,
modified_datetime: None,
platform_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 rows: Vec<DataRow> = (0_i32..20_i32)
.map(|i| DataRow {
created_datetime: None,
data_stream_id,
data_type: "spectrogram".to_owned(),
datetime: base + chrono::Duration::seconds(i64::from(i)),
modified_datetime: None,
specification_id,
vector: vec![f64::from(i), f64::from(i) + 0.5_f64],
vector_end_bound: 1.0_f64,
vector_start_bound: 0.0_f64,
})
.collect();
let inserted = r.insert_data_row_batch(&rows).await?;
assert_eq!(inserted, 20_u64);
let listed = r.list_data_rows(data_stream_id).await?;
assert_eq!(listed.len(), 20);
r.delete_platform(platform.id.unwrap()).await?;
Ok(())
}
#[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 = make_platform_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,
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,
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,
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(())
}