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 #[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 #[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 #[serde(skip_serializing_if = "Option::is_none")]
38 pub measurement_type: Option<String>,
39 #[serde(skip_serializing_if = "Option::is_none")]
42 pub sensor_id: Option<Uuid>,
43 #[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 #[serde(skip_serializing_if = "Option::is_none")]
63 pub measurement_type: Option<String>,
64 #[serde(skip_serializing_if = "Option::is_none")]
66 pub standard_curve_id: Option<Uuid>,
67}
68
69impl IngestReading {
70 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}