horizon-sdk 6.9.0

Canonical Rust data access layer for the Horizon platform
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
//! Domain model types for the Horizon platform.
//!
//! Each struct maps to a table in the `horizon_public` `PostgreSQL` schema.
//! Fields use `Option` where the database column is nullable or has a server-side default.

use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use sqlx::error::BoxDynError;
use sqlx::postgres::{PgTypeInfo, PgValueRef};
use uuid::Uuid;

/// Minimum byte length for a `PostgreSQL` `POINT` binary value.
const POINT_BYTE_LEN: usize = 16;

/// A geographic position as (longitude, latitude).
///
/// Maps to `PostgreSQL` `POINT` type with native sqlx decode support.
/// The binary wire format is two big-endian `f64` values (x, y).
#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
pub struct Position(pub f64, pub f64);

impl sqlx::Type<sqlx::Postgres> for Position {
    fn compatible(ty: &PgTypeInfo) -> bool {
        *ty == Self::type_info()
    }

    fn type_info() -> PgTypeInfo {
        PgTypeInfo::with_name("POINT")
    }
}

#[allow(
    clippy::as_conversions,
    clippy::indexing_slicing,
    clippy::std_instead_of_core,
    clippy::big_endian_bytes,
    clippy::float_arithmetic,
    reason = "byte slicing, Box cast, and BE f64 decode required for PostgreSQL POINT wire format"
)]
impl<'de> sqlx::Decode<'de, sqlx::Postgres> for Position {
    fn decode(value: PgValueRef<'de>) -> Result<Self, BoxDynError> {
        let buf = <&[u8] as sqlx::Decode<sqlx::Postgres>>::decode(value)?;
        if buf.len() < POINT_BYTE_LEN {
            return Err("POINT requires 16 bytes".into());
        }
        let (x_bytes, rest) = buf.split_at(8_usize);
        let x = f64::from_be_bytes(
            x_bytes
                .try_into()
                .map_err(|err| Box::new(err) as BoxDynError)?,
        );
        let y = f64::from_be_bytes(
            rest[..8]
                .try_into()
                .map_err(|err| Box::new(err) as BoxDynError)?,
        );
        Ok(Self(x, y))
    }
}

/// A label/annotation on a platform.
///
/// Time and value coordinates describe the labeled region. `time_coordinates`
/// pairs with `value_coordinates` element-wise (a point, line, or polygon
/// depending on length). `bearing_time_record_specification_id` and
/// `spectrogram_specification_id` are mutually exclusive — only one may be
/// set. User attribution (`created_by_user_id`, `modified_by_user_id`) is
/// managed by database triggers, not exposed here.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Annotation {
    pub bearing_time_record_specification_id: Option<Uuid>,
    pub confidence: Option<f64>,
    pub created_datetime: Option<DateTime<Utc>>,
    pub duration_seconds: Option<i32>,
    pub feed_context: Option<serde_json::Value>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub notes: Option<String>,
    pub ontology_class_id: Option<Uuid>,
    pub organization_id: Option<Uuid>,
    pub parent_annotation_id: Option<Uuid>,
    pub platform_id: Uuid,
    pub spectrogram_specification_id: Option<Uuid>,
    pub time_coordinates: Vec<DateTime<Utc>>,
    pub value_coordinates: Vec<f64>,
}

/// A specification for audio data processing.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct AudioSpecification {
    pub bit_depth: i64,
    pub channel_count: i64,
    pub channel_index: i64,
    pub created_datetime: Option<DateTime<Utc>>,
    pub encoding: String,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub organization_id: Option<Uuid>,
    pub sample_rate: i64,
}

/// A specification for beamgram data processing.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct BeamgramSpecification {
    pub center_bearing: Option<f64>,
    pub center_bin_width: Option<f64>,
    pub created_datetime: Option<DateTime<Utc>>,
    pub elevation_increment: Option<f64>,
    pub fft_sample_count: Option<i32>,
    pub id: Option<Uuid>,
    pub lower_elevation: Option<f64>,
    pub max_frequency: Option<f64>,
    pub min_frequency: Option<f64>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub name: Option<String>,
    pub normalizer: Option<String>,
    pub organization_id: Option<Uuid>,
    pub update_rate_ms: Option<i64>,
    pub upper_elevation: Option<f64>,
}

/// A specification for bearing-time record data processing.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct BearingTimeRecordSpecification {
    pub amplitude_unit_mode: Option<String>,
    pub baffle_bearing_list: Option<Vec<f64>>,
    pub bearing_bin_count: Option<i32>,
    pub created_datetime: Option<DateTime<Utc>>,
    pub elevation_increment: Option<f64>,
    pub fft_sample_count: Option<i32>,
    pub focus_range: Option<String>,
    pub frequency_spacing: Option<f64>,
    pub heading_data_type: Option<String>,
    pub heading_vector_index: Option<i32>,
    pub id: Option<Uuid>,
    pub lower_elevation: Option<f64>,
    pub max_frequency: Option<f64>,
    pub max_pixel: Option<i32>,
    pub min_frequency: Option<f64>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub name: Option<String>,
    pub normalizer: Option<String>,
    pub organization_id: Option<Uuid>,
    pub update_rate_ms: Option<i64>,
    pub upper_elevation: Option<f64>,
}

/// A row of data in the Iceberg `data_row` table.
///
/// This type is used for both Postgres inserts and Arrow IPC serialization.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DataRow {
    pub created_datetime: Option<DateTime<Utc>>,
    pub data_stream_id: Uuid,
    pub data_type: String,
    pub datetime: DateTime<Utc>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub specification_id: Uuid,
    pub vector: Vec<f64>,
    pub vector_end_bound: f64,
    pub vector_start_bound: f64,
}

/// Backend that powers tile queries for a data stream.
///
/// Maps to the `horizon_public.data_stream_query_source` Postgres ENUM.
#[allow(
    clippy::exhaustive_enums,
    reason = "variants mirror a closed Postgres ENUM, so exhaustive matching downstream is desirable"
)]
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, sqlx::Type)]
#[sqlx(
    type_name = "horizon_public.data_stream_query_source",
    rename_all = "lowercase"
)]
#[serde(rename_all = "lowercase")]
pub enum DataStreamQuerySource {
    /// Query the row-store `data_row` table directly.
    #[default]
    Postgres,
    /// Query Iceberg through Trino.
    Trino,
}

/// A data stream belonging to a platform.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DataStream {
    pub created_datetime: Option<DateTime<Utc>>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub name: Option<String>,
    pub organization_id: Option<Uuid>,
    pub platform_id: Option<Uuid>,
    #[serde(default)]
    pub query_source: DataStreamQuerySource,
}

impl DataStream {
    /// Create a `DataStream` with a deterministic UUID derived from the platform ID and name.
    ///
    /// Defaults `query_source` to [`DataStreamQuerySource::Postgres`]; override the
    /// field after construction if the stream should be served from Trino.
    #[must_use]
    pub fn from_platform_and_name(platform_id: Uuid, data_stream_name: &str) -> Self {
        let name = capitalize(data_stream_name);
        Self {
            created_datetime: None,
            id: Some(data_stream_uuid(platform_id, data_stream_name)),
            modified_datetime: None,
            name: Some(name),
            organization_id: None,
            platform_id: Some(platform_id),
            query_source: DataStreamQuerySource::Postgres,
        }
    }
}

/// A row of metadata in the Iceberg `metadata_row` table.
///
/// This type is used for both Postgres inserts and Arrow IPC serialization.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct MetadataRow {
    pub altitude: Option<f64>,
    pub created_datetime: Option<DateTime<Utc>>,
    pub data_stream_id: Uuid,
    pub datetime: DateTime<Utc>,
    pub heading: Option<f64>,
    pub latitude: Option<f64>,
    pub longitude: Option<f64>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub pitch: Option<f64>,
    pub roll: Option<f64>,
    pub speed: Option<f64>,
    pub speed_over_ground: Option<f64>,
}

/// A group that can group platforms around some bounded time or purpose.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Mission {
    pub created_datetime: Option<DateTime<Utc>>,
    pub end_datetime: Option<DateTime<Utc>>,
    pub free_text: Option<String>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub name: Option<String>,
    pub organization_id: Option<Uuid>,
    pub position: Option<Position>,
    pub start_datetime: Option<DateTime<Utc>>,
}

/// A relationship that associates a platform with a mission.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct MissionPlatform {
    pub created_datetime: Option<DateTime<Utc>>,
    pub id: Option<Uuid>,
    pub mission_id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub organization_id: Option<Uuid>,
    pub platform_id: Option<Uuid>,
}

/// An ontology for organizing platforms.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Ontology {
    pub created_datetime: Option<DateTime<Utc>>,
    pub description: Option<String>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub name: Option<String>,
    pub organization_id: Option<Uuid>,
}

/// A class within an ontology.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OntologyClass {
    pub created_datetime: Option<DateTime<Utc>>,
    pub description: Option<String>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub name: Option<String>,
    pub ontology_id: Uuid,
    pub order: Option<i32>,
    pub organization_id: Option<Uuid>,
    pub parent_id: Option<Uuid>,
    pub relationship_type: Option<String>,
}

/// An instance of a platform.
///
/// The `position` field maps to `PostgreSQL`'s `POINT` type as `(longitude, latitude)`.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Platform {
    pub created_datetime: Option<DateTime<Utc>>,
    pub end_datetime: Option<DateTime<Utc>>,
    pub free_text: Option<String>,
    pub id: Option<Uuid>,
    pub kind_id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub name: Option<String>,
    pub organization_id: Option<Uuid>,
    pub position: Option<Position>,
    pub start_datetime: Option<DateTime<Utc>>,
}

/// A relationship that associates a platform with an audio specification.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PlatformAudioSpecification {
    pub audio_specification_id: Uuid,
    pub created_datetime: Option<DateTime<Utc>>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub organization_id: Option<Uuid>,
    pub platform_id: Uuid,
}

/// A relationship that associates a platform with a beamgram specification.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PlatformBeamgramSpecification {
    pub beamgram_specification_id: Option<Uuid>,
    pub created_datetime: Option<DateTime<Utc>>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub organization_id: Option<Uuid>,
    pub platform_id: Option<Uuid>,
}

/// A relationship that associates a platform with a bearing-time record specification.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PlatformBearingTimeRecordSpecification {
    pub bearing_time_record_specification_id: Option<Uuid>,
    pub created_datetime: Option<DateTime<Utc>>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub organization_id: Option<Uuid>,
    pub platform_id: Option<Uuid>,
}

/// Key-value information stored for a platform (append-only).
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PlatformInformation {
    pub created_datetime: Option<DateTime<Utc>>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub organization_id: Option<Uuid>,
    pub platform_id: Uuid,
    pub properties: serde_json::Value,
}

/// A kind of platform (e.g. buoy, vessel, aircraft).
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PlatformKind {
    pub created_datetime: Option<DateTime<Utc>>,
    pub id: Option<Uuid>,
    pub image_url: Option<String>,
    pub long_description: Option<String>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub name: Option<String>,
    pub organization_id: Option<Uuid>,
    pub short_description: Option<String>,
}

/// A relationship that associates a platform with a spectrogram specification.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PlatformSpectrogramSpecification {
    pub created_datetime: Option<DateTime<Utc>>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub organization_id: Option<Uuid>,
    pub platform_id: Option<Uuid>,
    pub spectrogram_specification_id: Option<Uuid>,
}

/// A directional spectrogram specification: leading/lagging bearing row offsets
/// used to derive directional hue on top of a spectrogram.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DirectionalSpectrogramSpecification {
    pub bearing_lagging_rows: Option<i64>,
    pub bearing_leading_rows: Option<i64>,
    pub created_datetime: Option<DateTime<Utc>>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub organization_id: Option<Uuid>,
    pub spectrogram_specification_id: Option<Uuid>,
}

/// A specification for spectrogram data processing.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SpectrogramSpecification {
    pub amplitude_unit_mode: Option<String>,
    pub channel: Option<i32>,
    pub channel_role: Option<String>,
    pub created_datetime: Option<DateTime<Utc>>,
    pub fft_sample_count: Option<i32>,
    pub fft_sample_overlap_count: Option<i32>,
    pub frequency_spacing: Option<f64>,
    pub id: Option<Uuid>,
    pub modified_datetime: Option<DateTime<Utc>>,
    pub name: Option<String>,
    pub organization_id: Option<Uuid>,
}

/// Capitalize a string: uppercase first character, lowercase rest.
///
/// Matches Python's `str.capitalize()` for cross-SDK parity.
#[must_use]
pub fn capitalize(s: &str) -> String {
    let mut chars = s.chars();
    chars.next().map_or_else(String::new, |first| {
        let upper: String = first.to_uppercase().collect();
        format!("{upper}{}", chars.as_str().to_lowercase())
    })
}

/// Compute the deterministic UUID for a data stream from its platform and name.
///
/// Applies [`capitalize`] to the name for Python SDK parity, then hashes
/// the `"{platform_id}:{capitalized_name}"` string via [`name_to_uuid`].
#[must_use]
pub fn data_stream_uuid(platform_id: Uuid, data_stream_name: &str) -> Uuid {
    let name = capitalize(data_stream_name);
    name_to_uuid(&format!("{platform_id}:{name}"))
}

/// Deterministic UUID generation for data idempotency.
///
/// Uses UUID v5 (SHA-1 namespace) with a fixed namespace to produce
/// the same UUID for the same input string across invocations.
#[must_use]
pub fn name_to_uuid(name: &str) -> Uuid {
    const NAMESPACE: Uuid = Uuid::from_bytes([
        0x13, 0xdb, 0x1b, 0xb5, 0x27, 0xac, 0x4b, 0xf7, 0x81, 0x86, 0x10, 0x43, 0x34, 0x7c, 0x42,
        0x47,
    ]);
    let normalized = name.trim().to_lowercase();
    Uuid::new_v5(&NAMESPACE, normalized.as_bytes())
}

#[cfg(test)]
#[allow(clippy::unwrap_used, reason = "tests use unwrap for brevity")]
mod tests {
    use super::*;

    #[test]
    fn capitalize_matches_python() {
        assert_eq!(capitalize("hydrophone"), "Hydrophone");
        assert_eq!(capitalize("Hydrophone"), "Hydrophone");
        assert_eq!(capitalize("hYDROPHONE"), "Hydrophone");
        assert_eq!(capitalize("HYDROPHONE"), "Hydrophone");
        assert_eq!(capitalize(""), "");
        assert_eq!(capitalize("a"), "A");
    }

    #[test]
    fn data_stream_from_platform_and_name() {
        let platform_id = Uuid::new_v4();
        let ds = DataStream::from_platform_and_name(platform_id, "audio");
        assert_eq!(ds.platform_id, Some(platform_id));
        assert_eq!(ds.name.as_deref(), Some("Audio"));
        assert!(ds.id.is_some());

        let ds2 = DataStream::from_platform_and_name(platform_id, "audio");
        assert_eq!(ds.id, ds2.id);
    }

    #[test]
    fn data_stream_uuid_is_deterministic() {
        let platform_id = Uuid::new_v4();
        let a = data_stream_uuid(platform_id, "audio");
        let b = data_stream_uuid(platform_id, "audio");
        assert_eq!(a, b);
    }

    #[test]
    fn data_stream_uuid_matches_from_platform_and_name() {
        let platform_id = Uuid::new_v4();
        let ds = DataStream::from_platform_and_name(platform_id, "audio");
        assert_eq!(ds.id, Some(data_stream_uuid(platform_id, "audio")));
    }

    #[test]
    fn mission_serde_round_trip() {
        let mission = Mission {
            created_datetime: None,
            end_datetime: None,
            free_text: None,
            id: Some(Uuid::new_v4()),
            modified_datetime: None,
            name: Some("Test Mission".to_owned()),
            organization_id: None,
            position: Some(Position(1.5_f64, 2.5_f64)),
            start_datetime: Some(Utc::now()),
        };
        let json = serde_json::to_string(&mission).unwrap();
        let deserialized: Mission = serde_json::from_str(&json).unwrap();
        assert_eq!(mission, deserialized);
    }

    #[test]
    fn name_to_uuid_is_deterministic() {
        let a = name_to_uuid("test");
        let b = name_to_uuid("test");
        assert_eq!(a, b);
    }

    #[test]
    fn name_to_uuid_matches_python_sdk() {
        let expected = Uuid::parse_str("e32be881-c656-5b8e-8dd7-889066ac1c83").unwrap();
        assert_eq!(name_to_uuid("Test Platform"), expected);

        let expected2 = Uuid::parse_str("eaeebe45-0eaa-59e4-9b79-7ceda69789e2").unwrap();
        assert_eq!(name_to_uuid("hello world"), expected2);
    }

    #[test]
    fn name_to_uuid_normalizes_whitespace_and_case() {
        let a = name_to_uuid("  Hello World  ");
        let b = name_to_uuid("hello world");
        assert_eq!(a, b);
    }

    #[test]
    fn platform_information_serde_with_json_properties() {
        let info = PlatformInformation {
            created_datetime: None,
            id: None,
            modified_datetime: None,
            organization_id: None,
            platform_id: Uuid::new_v4(),
            properties: serde_json::json!({"key": "value", "nested": {"a": 1_i32}}),
        };
        let json = serde_json::to_string(&info).unwrap();
        let deserialized: PlatformInformation = serde_json::from_str(&json).unwrap();
        assert_eq!(info, deserialized);
    }
}