Skip to main content

river_data_core/models/
streams.rs

1use serde::{Deserialize, Serialize};
2use uuid::Uuid;
3
4use crate::models::replicates::{ColumnAssignment, ReplicateSpec};
5
6#[derive(Debug, Clone, Serialize, Deserialize)]
7pub struct DataStream {
8    pub id: Uuid,
9    pub source_system: String,
10    pub source_key: String,
11    pub source_name: Option<String>,
12    pub source_path: Option<String>,
13    pub metadata: serde_json::Value,
14    pub site_parameter_id: Option<Uuid>,
15    /// Stream-level default for readings.measurement_type ('continuous' | 'spot' | 'derived').
16    /// None defers to the API's sensor-frequency resolution.
17    #[serde(default, skip_serializing_if = "Option::is_none")]
18    pub measurement_type: Option<String>,
19    pub is_active: bool,
20    pub last_data_time: Option<chrono::DateTime<chrono::Utc>>,
21    /// The authoritative replicate column-to-index mapping, present on a
22    /// register response for a stream declaring a replicate family. Absent on
23    /// list responses and on APIs that predate pinning; the same list persists
24    /// under `metadata.replicates.assignments`.
25    #[serde(default, skip_serializing_if = "Option::is_none")]
26    pub replicates: Option<Vec<ColumnAssignment>>,
27}
28
29#[derive(Debug, Serialize)]
30pub struct RegisterStreamRequest {
31    pub source_system: String,
32    pub source_key: String,
33    pub source_name: Option<String>,
34    pub source_path: Option<String>,
35    pub metadata: serde_json::Value,
36    /// Stream-level classification declared at discovery. None never clears an operator-set value.
37    #[serde(skip_serializing_if = "Option::is_none")]
38    pub measurement_type: Option<String>,
39    /// Owning sensor. Required for curve-carrying streams: the API admits a
40    /// reading's curve claim only when reading-sensor == curve-sensor.
41    #[serde(skip_serializing_if = "Option::is_none")]
42    pub sensor_id: Option<Uuid>,
43    /// Replicate-family declaration; requires `measurement_type: "spot"`.
44    #[serde(skip_serializing_if = "Option::is_none")]
45    pub replicates: Option<ReplicateSpec>,
46}
47
48#[derive(Debug, Clone, Serialize)]
49pub struct IngestReading {
50    pub time: chrono::DateTime<chrono::Utc>,
51    pub raw_value: f64,
52    #[serde(skip_serializing_if = "is_zero")]
53    pub replicate_index: i16,
54    #[serde(skip_serializing_if = "Option::is_none")]
55    pub sensor_id: Option<Uuid>,
56    #[serde(skip_serializing_if = "Option::is_none")]
57    pub calibration_id: Option<Uuid>,
58    #[serde(skip_serializing_if = "Option::is_none")]
59    pub deployment_id: Option<Uuid>,
60    /// Per-reading override ('continuous' | 'spot' | 'derived'). None resolves server-side from
61    /// the stream default, then the owning sensor's data_frequency.
62    #[serde(skip_serializing_if = "Option::is_none")]
63    pub measurement_type: Option<String>,
64    /// Standard curve the source applied to this reading.
65    #[serde(skip_serializing_if = "Option::is_none")]
66    pub standard_curve_id: Option<Uuid>,
67}
68
69impl IngestReading {
70    /// A reading at replicate 0 with no sensor attribution; the server resolves the rest.
71    pub fn new(time: chrono::DateTime<chrono::Utc>, raw_value: f64) -> Self {
72        Self {
73            time,
74            raw_value,
75            replicate_index: 0,
76            sensor_id: None,
77            calibration_id: None,
78            deployment_id: None,
79            measurement_type: None,
80            standard_curve_id: None,
81        }
82    }
83}
84
85fn is_zero(v: &i16) -> bool {
86    *v == 0
87}
88
89#[derive(Debug, Serialize)]
90pub struct IngestStatusEvent {
91    pub time: chrono::DateTime<chrono::Utc>,
92    pub value: String,
93}
94
95#[cfg(test)]
96mod tests {
97    use super::*;
98
99    #[test]
100    fn test_ingest_reading_serialization() {
101        let r = IngestReading::new(chrono::Utc::now(), 42.5);
102        let json = serde_json::to_value(&r).unwrap();
103        assert_eq!(json["raw_value"], 42.5);
104        assert!(json.get("replicate_index").is_none());
105        assert!(json.get("sensor_id").is_none());
106        assert!(json.get("measurement_type").is_none());
107    }
108
109    #[test]
110    fn test_register_stream_request() {
111        let req = RegisterStreamRequest {
112            source_system: "test_system".to_string(),
113            source_key: "source_1".to_string(),
114            source_name: Some("stream_a".to_string()),
115            source_path: None,
116            metadata: serde_json::json!({"device": "dev_001"}),
117            measurement_type: None,
118            sensor_id: None,
119            replicates: None,
120        };
121        let json = serde_json::to_value(&req).unwrap();
122        assert_eq!(json["source_system"], "test_system");
123        assert_eq!(json["metadata"]["device"], "dev_001");
124        assert!(json.get("sensor_id").is_none());
125        assert!(json.get("replicates").is_none());
126    }
127
128    #[test]
129    fn test_register_stream_request_with_replicates() {
130        let req = RegisterStreamRequest {
131            source_system: "cnet".to_string(),
132            source_key: "VAD:DOC_avg_ppb:reps".to_string(),
133            source_name: None,
134            source_path: None,
135            metadata: serde_json::json!({}),
136            measurement_type: Some("spot".to_string()),
137            sensor_id: Some(Uuid::nil()),
138            replicates: Some(crate::models::replicates::ReplicateSpec {
139                source_columns: vec!["DOC_rep_1".into(), "DOC_rep_2".into(), "DOC_rep_3".into()],
140                portal_mean_column: Some("DOC_avg_ppb".into()),
141                portal_sd_column: Some("DOC_sd_ppb".into()),
142                curve_ref_column: Some("doc_std_curve_id".into()),
143                calc: Some("calcDOCavg".into()),
144            }),
145        };
146        let json = serde_json::to_value(&req).unwrap();
147        assert_eq!(json["measurement_type"], "spot");
148        assert_eq!(json["replicates"]["source_columns"][2], "DOC_rep_3");
149        assert_eq!(json["replicates"]["curve_ref_column"], "doc_std_curve_id");
150    }
151
152    #[test]
153    fn test_data_stream_deserialization() {
154        let json = serde_json::json!({
155            "id": "550e8400-e29b-41d4-a716-446655440000",
156            "source_system": "test_system",
157            "source_key": "source_1",
158            "source_name": "stream_a",
159            "source_path": null,
160            "metadata": {},
161            "site_parameter_id": null,
162            "is_active": true,
163            "last_data_time": null
164        });
165        let stream: DataStream = serde_json::from_value(json).unwrap();
166        assert_eq!(stream.source_system, "test_system");
167        assert!(stream.is_active);
168        assert!(stream.site_parameter_id.is_none());
169        assert!(stream.replicates.is_none());
170    }
171
172    #[test]
173    fn register_response_replicates_parse() {
174        let json = serde_json::json!({
175            "id": "550e8400-e29b-41d4-a716-446655440000",
176            "source_system": "cnet",
177            "source_key": "VAD:DOC_avg_ppb:reps",
178            "source_name": null,
179            "source_path": null,
180            "metadata": {},
181            "site_parameter_id": null,
182            "is_active": true,
183            "last_data_time": null,
184            "replicates": [
185                {"column": "DOC_rep_1", "index": 0},
186                {"column": "DOC_rep_2", "index": 5, "retired": true},
187            ]
188        });
189        let stream: DataStream = serde_json::from_value(json).unwrap();
190        let assignments = stream.replicates.unwrap();
191        assert_eq!(assignments.len(), 2);
192        assert_eq!(assignments[0].index, 0);
193        assert!(!assignments[0].retired);
194        assert_eq!(assignments[1].column, "DOC_rep_2");
195        assert_eq!(assignments[1].index, 5);
196        assert!(assignments[1].retired);
197    }
198}