diff --git a/frontend/composables/state.ts b/frontend/composables/state.ts index 930f6f5..eee06ed 100644 --- a/frontend/composables/state.ts +++ b/frontend/composables/state.ts @@ -100,7 +100,7 @@ export const PREVIEW_PAGE_SIZE = 100; export const PREVIEW_PAGE_INCREMENT = PREVIEW_PAGE_SIZE; export const APP_NAME = "HydroServer Streaming Data Loader"; export const API_KEY_DOCS_URL = - "https://hydroserver2.github.io/hydroserver/tutorials/creating-your-first-orchestration-system#create-an-api-key"; + "https://hydroserver.org/user-guides/tutorials/hydroserver-101/part-3-sdl-setup#create-an-api-key"; export function emptyServerConfig(): ServerConfig { return { diff --git a/src/hydroserver.rs b/src/hydroserver.rs index ece566c..734d2f9 100644 --- a/src/hydroserver.rs +++ b/src/hydroserver.rs @@ -1,12 +1,8 @@ -use std::{ - collections::HashMap, - sync::{Arc, Mutex}, - time::Duration, -}; +use std::{collections::HashMap, sync::Mutex, time::Duration}; use chrono::{DateTime, Utc}; use reqwest::{ - header::{HeaderMap, HeaderValue, ACCEPT, AUTHORIZATION, CONTENT_TYPE}, + header::{ACCEPT, CONTENT_TYPE}, Client, Method, Response, StatusCode, }; use serde_json::{json, Value}; @@ -18,16 +14,13 @@ use crate::models::{ ServerConfig, ServerUrlValidationResponse, }; -const AUTH_ROUTE: &str = "/api/auth"; const BASE_ROUTE: &str = "/api/data"; const DATASTREAM_PAGE_SIZE: usize = 1000; +const DATASTREAM_INCLUDE: &str = "monitoringSite,observedProperty,processingLevel,unit,method"; const DATASTREAM_CACHE_TTL_SECONDS: i64 = 300; type DatastreamCacheEntry = (DateTime, Vec); type DatastreamCache = HashMap; -/// Cached username/password session tokens keyed by [`token_cache_key`], shared -/// across the short-lived [`HydroServerSession`]s spun up per request. -type TokenCache = HashMap; #[derive(Debug, Clone)] pub struct ObservationPayloadRow { @@ -38,7 +31,6 @@ pub struct ObservationPayloadRow { pub struct HydroServerService { http: Client, datastream_cache: Mutex, - session_tokens: Arc>, } impl HydroServerService { @@ -51,15 +43,11 @@ impl HydroServerService { Ok(Self { http, datastream_cache: Mutex::new(HashMap::new()), - session_tokens: Arc::new(Mutex::new(HashMap::new())), }) } - /// Build a session that shares this service's cached session tokens, so - /// username/password auth reuses a bearer token across requests instead of - /// logging in afresh for every batch. fn session(&self, server: ServerConfig) -> HydroServerSession { - HydroServerSession::new(self.http.clone(), server, self.session_tokens.clone()) + HydroServerSession::new(self.http.clone(), server) } pub async fn validate_url(&self, url: &str) -> ServerUrlValidationResponse { @@ -72,39 +60,12 @@ impl HydroServerService { }; } - let auth_probe_url = format!("{normalized_url}{AUTH_ROUTE}/app/session"); let data_probe_url = format!("{normalized_url}{BASE_ROUTE}/workspaces"); - let auth_response = self - .http - .get(&auth_probe_url) - .header(ACCEPT, "application/json") - .send() - .await; - - if let Ok(response) = auth_response { - if looks_like_hydroserver_auth_response(response).await { - let instance_name = instance_name(&normalized_url); - return ServerUrlValidationResponse { - ok: true, - message: format!("HydroServer API detected at {instance_name}."), - instance_name: Some(instance_name), - }; - } - } else if let Err(err) = auth_response { - if err.is_connect() || err.is_timeout() { - return ServerUrlValidationResponse { - ok: false, - message: "Couldn't reach that URL. Check the server URL and try again." - .to_string(), - instance_name: None, - }; - } - } - match self .http .get(&data_probe_url) + .query(&[("limit", "1")]) .header(ACCEPT, "application/json") .send() .await @@ -223,7 +184,7 @@ impl HydroServerService { }) => ConnectionTestResponse { ok: false, state: ConnectionState::Error, - message: "These credentials are invalid or do not have the permissions the loader needs. Make sure they can access workspaces, datastreams, and orchestration systems.".to_string(), + message: "These credentials are invalid or do not have the permissions the loader needs. Make sure they can access workspaces, datastreams, and observations.".to_string(), invalid_field: Some(match server.auth_type { AuthType::Apikey => "api_key".to_string(), AuthType::Userpass => "username".to_string(), @@ -295,67 +256,19 @@ impl HydroServerService { } } - if let Some(datastreams) = self - .list_datastreams_from_bootstrap(&mut session, &workspace_id) - .await? - { - self.set_cached_datastreams(&cache_key, datastreams.clone()); - return Ok(datastreams); - } - - if let Some(datastreams) = self - .list_datastreams_expanded(&mut session, &workspace_id) - .await? - { - self.set_cached_datastreams(&cache_key, datastreams.clone()); - return Ok(datastreams); - } - let datastreams = session .fetch_all_collection( &format!("{BASE_ROUTE}/datastreams"), - &[("workspace_id", workspace_id.clone())], + &[ + ("workspace_id", workspace_id.clone()), + ("include", DATASTREAM_INCLUDE.to_string()), + ], ) .await .map_err(|err| err.to_string())?; - - if datastreams.is_empty() { - return Ok(Vec::new()); - } - - let things_by_id = session - .fetch_collection_lookup(&format!("{BASE_ROUTE}/things"), &workspace_id) - .await - .unwrap_or_default(); - let observed_properties_by_id = session - .fetch_collection_lookup(&format!("{BASE_ROUTE}/observed-properties"), &workspace_id) - .await - .unwrap_or_default(); - let processing_levels_by_id = session - .fetch_collection_lookup(&format!("{BASE_ROUTE}/processing-levels"), &workspace_id) - .await - .unwrap_or_default(); - let units_by_id = session - .fetch_collection_lookup(&format!("{BASE_ROUTE}/units"), &workspace_id) - .await - .unwrap_or_default(); - let sensors_by_id = session - .fetch_collection_lookup(&format!("{BASE_ROUTE}/sensors"), &workspace_id) - .await - .unwrap_or_default(); - let summaries = datastreams .iter() - .map(|item| { - datastream_to_summary( - item, - &things_by_id, - &observed_properties_by_id, - &processing_levels_by_id, - &units_by_id, - &sensors_by_id, - ) - }) + .map(expanded_datastream_to_summary) .collect::>(); self.set_cached_datastreams(&cache_key, summaries.clone()); @@ -390,14 +303,18 @@ impl HydroServerService { &format!("{BASE_ROUTE}/datastreams/{datastream_id}"), &[ ("workspace_id", workspace_id), - ("expand_related", "true".to_string()), + ("include", DATASTREAM_INCLUDE.to_string()), ], None, ) .await .map_err(|err| err.to_string())?; - Ok(expanded_datastream_to_detail(&payload)) + let mut item = response_item(&payload) + .map_err(|err| err.to_string())? + .clone(); + attach_datastream_relations(std::slice::from_mut(&mut item), &payload); + Ok(expanded_datastream_to_detail(&item)) } /// Fetch the datastream's current `phenomenonEndTime` — the timestamp of @@ -427,7 +344,7 @@ impl HydroServerService { ) .await?; - Ok(parse_phenomenon_end_time(&payload)) + Ok(parse_phenomenon_end_time(response_item(&payload)?)) } pub(crate) async fn post_observations_batch( @@ -441,9 +358,8 @@ impl HydroServerService { } let mut session = self.session(server.clone().normalized()); - // The earlier Rust port posted ["timestamp", "value"], which does not match the - // HydroServer bulk observation schema. The API expects SensorThings field names. let body = json!({ + "datastreamId": datastream_id, "fields": ["phenomenonTime", "result"], "data": observations .iter() @@ -454,158 +370,13 @@ impl HydroServerService { session .request_void( Method::POST, - &format!("{BASE_ROUTE}/datastreams/{datastream_id}/observations/bulk-create"), + &format!("{BASE_ROUTE}/observations/bulk-create"), &[("mode", "insert".to_string())], Some(body), ) .await } - async fn list_datastreams_from_bootstrap( - &self, - session: &mut HydroServerSession, - workspace_id: &str, - ) -> Result>, String> { - let payload = match session - .request_json( - Method::GET, - &format!("{BASE_ROUTE}/datastreams/visualization-bootstrap"), - &[("workspace_id", workspace_id.to_string())], - None, - ) - .await - { - Ok(payload) => payload, - Err(_) => return Ok(None), - }; - - let Some(datastreams) = payload.get("datastreams").and_then(Value::as_array) else { - return Ok(None); - }; - let Some(things) = payload.get("things").and_then(Value::as_array) else { - return Ok(None); - }; - let Some(observed_properties) = payload - .get("observed_properties") - .or_else(|| payload.get("observedProperties")) - .and_then(Value::as_array) - else { - return Ok(None); - }; - let Some(processing_levels) = payload - .get("processing_levels") - .or_else(|| payload.get("processingLevels")) - .and_then(Value::as_array) - else { - return Ok(None); - }; - - let units_by_id = match session - .fetch_collection_lookup(&format!("{BASE_ROUTE}/units"), workspace_id) - .await - { - Ok(units) => units, - Err(_) => return Ok(None), - }; - - let things_by_id = map_items_by_id(things); - let observed_properties_by_id = map_items_by_id(observed_properties); - let processing_levels_by_id = map_items_by_id(processing_levels); - - Ok(Some( - datastreams - .iter() - .map(|datastream| { - let thing_id = string_value(datastream, &["thing_id", "thingId"]); - let observed_property_id = - string_value(datastream, &["observed_property_id", "observedPropertyId"]); - let processing_level_id = - string_value(datastream, &["processing_level_id", "processingLevelId"]); - let unit_id = string_value(datastream, &["unit_id", "unitId"]); - - DatastreamSummary { - id: string_value(datastream, &["id", "uid"]).unwrap_or_default(), - name: string_value(datastream, &["name"]) - .unwrap_or_else(|| "Unnamed datastream".to_string()), - thing_id: thing_id.clone().unwrap_or_default(), - thing_name: string_value_from_map(&things_by_id, &thing_id, &["name"]), - observed_property_name: string_value_from_map( - &observed_properties_by_id, - &observed_property_id, - &["name"], - ), - processing_level_definition: string_value_from_map( - &processing_levels_by_id, - &processing_level_id, - &["definition"], - ), - unit_name: string_value_from_map(&units_by_id, &unit_id, &["name"]), - unit_symbol: string_value_from_map(&units_by_id, &unit_id, &["symbol"]), - sampled_medium: String::new(), - sensor_name: String::new(), - result_type: String::new(), - } - }) - .collect(), - )) - } - - async fn list_datastreams_expanded( - &self, - session: &mut HydroServerSession, - workspace_id: &str, - ) -> Result>, String> { - let mut page = 1_u32; - let mut datastreams = Vec::new(); - - loop { - let response = match session - .request_response( - Method::GET, - &format!("{BASE_ROUTE}/datastreams"), - &[ - ("workspace_id", workspace_id.to_string()), - ("expand_related", "true".to_string()), - ("page", page.to_string()), - ("page_size", DATASTREAM_PAGE_SIZE.to_string()), - ], - None, - ) - .await - { - Ok(response) => response, - Err(_) => return Ok(None), - }; - - let headers = response.headers().clone(); - let payload = response - .json::() - .await - .map_err(|err| err.to_string())?; - let Some(items) = payload.as_array() else { - return Ok(None); - }; - - if items.is_empty() { - break; - } - - datastreams.extend(items.iter().map(expanded_datastream_to_summary)); - - if let Some(total_pages) = header_int(&headers, "X-Total-Pages") { - if page >= total_pages { - break; - } - } else if items.len() < DATASTREAM_PAGE_SIZE { - break; - } - - page += 1; - } - - Ok(Some(datastreams)) - } - fn cached_datastreams(&self, cache_key: &str) -> Option> { let mut cache = self.datastream_cache.lock().ok()?; let (cached_at, datastreams) = cache.get(cache_key)?.clone(); @@ -627,52 +398,11 @@ impl HydroServerService { struct HydroServerSession { http: Client, server: ServerConfig, - bearer_token: Option, - token_cache: Arc>, } impl HydroServerSession { - fn new(http: Client, server: ServerConfig, token_cache: Arc>) -> Self { - Self { - http, - server, - bearer_token: None, - token_cache, - } - } - - /// Identity a cached session token belongs to. Username/password tokens are - /// scoped to the instance URL and account; the password is deliberately not - /// part of the key, since a token minted under an old password simply 401s - /// and gets refreshed by [`request_response`]. - fn token_cache_key(&self) -> String { - format!( - "{}|{}", - self.server.url.trim_end_matches('/'), - self.server.username.trim() - ) - } - - fn cached_token(&self) -> Option { - self.token_cache - .lock() - .ok()? - .get(&self.token_cache_key()) - .cloned() - } - - fn store_token(&self, token: &str) { - if let Ok(mut cache) = self.token_cache.lock() { - cache.insert(self.token_cache_key(), token.to_string()); - } - } - - /// Drop the cached token for this account so the next request re-authenticates. - fn invalidate_token(&mut self) { - self.bearer_token = None; - if let Ok(mut cache) = self.token_cache.lock() { - cache.remove(&self.token_cache_key()); - } + fn new(http: Client, server: ServerConfig) -> Self { + Self { http, server } } async fn associated_workspace( @@ -706,61 +436,51 @@ impl HydroServerSession { .collect()) } - async fn fetch_collection_lookup( - &mut self, - path: &str, - workspace_id: &str, - ) -> Result, RequestError> { - let items = self - .fetch_all_collection(path, &[("workspace_id", workspace_id.to_string())]) - .await?; - Ok(items - .into_iter() - .filter_map(|item| { - let id = string_value(&item, &["id", "uid"])?; - Some((id, item)) - }) - .collect()) - } - async fn fetch_all_collection( &mut self, path: &str, params: &[(&str, String)], ) -> Result, RequestError> { - let mut page = 1_u32; + let mut offset = 0_u64; let mut items = Vec::new(); loop { let mut page_params = params.to_vec(); - page_params.push(("page", page.to_string())); - page_params.push(("page_size", DATASTREAM_PAGE_SIZE.to_string())); - let response = self - .request_response(Method::GET, path, &page_params, None) + page_params.push(("offset", offset.to_string())); + page_params.push(("limit", DATASTREAM_PAGE_SIZE.to_string())); + let payload = self + .request_json(Method::GET, path, &page_params, None) .await?; - let headers = response.headers().clone(); - let payload = response - .json::() - .await - .map_err(|err| RequestError::Other(err.to_string()))?; - let Some(page_items) = payload.as_array() else { - return Ok(Vec::new()); - }; - - if page_items.is_empty() { - break; + let mut page_items = payload + .get("data") + .and_then(Value::as_array) + .ok_or_else(|| { + RequestError::Other( + "HydroServer returned an invalid collection response.".to_string(), + ) + })? + .clone(); + let total_count = payload + .pointer("/meta/totalCount") + .and_then(Value::as_u64) + .ok_or_else(|| { + RequestError::Other( + "HydroServer returned invalid pagination metadata.".to_string(), + ) + })?; + + if page_items.is_empty() && offset < total_count { + return Err(RequestError::Other( + "HydroServer returned an empty page before the collection was complete." + .to_string(), + )); } - items.extend(page_items.iter().cloned()); - - if let Some(total_pages) = header_int(&headers, "X-Total-Pages") { - if page >= total_pages { - break; - } - } else if page_items.len() < DATASTREAM_PAGE_SIZE { + offset += page_items.len() as u64; + attach_datastream_relations(&mut page_items, &payload); + items.extend(page_items); + if offset >= total_count { break; } - - page += 1; } Ok(items) @@ -804,143 +524,40 @@ impl HydroServerSession { ) -> Result { let url = build_url(&self.server.url, path); - for attempt in 0..2 { - let mut request = self - .http - .request(method.clone(), &url) - .header(ACCEPT, "application/json"); - - if !params.is_empty() { - request = request.query(params); - } - - if let Some(payload) = body.clone() { - request = request - .header(CONTENT_TYPE, "application/json") - .json(&payload); - } - - request = self.apply_auth(request).await?; - - let response = match request.send().await { - Ok(response) => response, - Err(err) if err.is_connect() => return Err(RequestError::Connection), - Err(err) if err.is_timeout() => return Err(RequestError::Timeout), - Err(err) => return Err(RequestError::Other(err.to_string())), - }; - - if response.status().is_success() { - return Ok(response); - } - - if attempt == 0 - && self.server.auth_type == AuthType::Userpass - && matches!( - response.status(), - StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN - ) - { - // The token we used — possibly one cached from an earlier - // request — was rejected. Drop it from the shared cache so the - // retry (and other sessions) re-authenticate instead of reusing - // a known-bad token. - self.invalidate_token(); - continue; - } - - let status = response.status(); - let message = response_error_message(response).await; - return Err(RequestError::Http { - status: Some(status), - message, - }); + let mut request = self + .http + .request(method, &url) + .header(ACCEPT, "application/json"); + if !params.is_empty() { + request = request.query(params); } - - Err(RequestError::Other( - "HydroServer request failed after retry.".to_string(), - )) - } - - async fn apply_auth( - &mut self, - request: reqwest::RequestBuilder, - ) -> Result { - match self.server.auth_type { - AuthType::Apikey => Ok(request.header("X-API-Key", self.server.api_key.clone())), + if let Some(payload) = body { + request = request.json(&payload); + } + request = match self.server.auth_type { + AuthType::Apikey => request.header("X-API-Key", &self.server.api_key), + // HydroServer v2 removed the app/session endpoint. The data API + // accepts email/password directly through HTTP Basic authentication. AuthType::Userpass => { - let token = self.session_token().await?; - Ok(request.header( - AUTHORIZATION, - HeaderValue::from_str(&format!("Bearer {token}")) - .map_err(|err| RequestError::Other(err.to_string()))?, - )) + request.basic_auth(&self.server.username, Some(&self.server.password)) } + }; + let response = request.send().await.map_err(|err| { + if err.is_connect() { + RequestError::Connection + } else if err.is_timeout() { + RequestError::Timeout + } else { + RequestError::Other(err.to_string()) + } + })?; + if response.status().is_success() { + return Ok(response); } - } - - async fn session_token(&mut self) -> Result { - if let Some(token) = &self.bearer_token { - return Ok(token.clone()); - } - - // Reuse a token cached by an earlier request for this account before - // paying for a fresh login. A stale one is purged on the 401/403 retry. - if let Some(token) = self.cached_token() { - self.bearer_token = Some(token.clone()); - return Ok(token); - } - - let payload = json!({ - "email": self.server.username, - "password": self.server.password, - }); - - let response = self - .http - .post(build_url( - &self.server.url, - &format!("{AUTH_ROUTE}/app/session"), - )) - .header(CONTENT_TYPE, "application/json") - .header(ACCEPT, "application/json") - .json(&payload) - .send() - .await - .map_err(|err| { - if err.is_connect() { - RequestError::Connection - } else if err.is_timeout() { - RequestError::Timeout - } else { - RequestError::Other(err.to_string()) - } - })?; - - if !response.status().is_success() { - let status = response.status(); - let message = response_error_message(response).await; - return Err(RequestError::Http { - status: Some(status), - message, - }); - } - - let payload = response - .json::() - .await - .map_err(|err| RequestError::Other(err.to_string()))?; - let token = payload - .get("meta") - .and_then(|meta| meta.get("session_token")) - .and_then(Value::as_str) - .ok_or_else(|| { - RequestError::Other("Authentication failed: No access token returned.".to_string()) - })? - .to_string(); - - self.bearer_token = Some(token.clone()); - self.store_token(&token); - Ok(token) + Err(RequestError::Http { + status: Some(response.status()), + message: response_error_message(response).await, + }) } } @@ -1105,56 +722,21 @@ fn instance_name(url: &str) -> String { .unwrap_or_else(|| url.to_string()) } -async fn looks_like_hydroserver_auth_response(response: Response) -> bool { - if !matches!( - response.status(), - StatusCode::OK - | StatusCode::UNAUTHORIZED - | StatusCode::FORBIDDEN - | StatusCode::METHOD_NOT_ALLOWED - | StatusCode::UNPROCESSABLE_ENTITY - ) { - return false; - } - - let Some(payload) = response_json(response).await else { - return false; - }; - - let Some(payload) = payload.as_object() else { - return false; - }; - - payload - .get("meta") - .and_then(Value::as_object) - .map(|meta| meta.contains_key("is_authenticated")) - .unwrap_or(false) - || payload - .get("data") - .and_then(Value::as_object) - .map(|data| data.contains_key("flows")) - .unwrap_or(false) - || payload.get("detail").and_then(Value::as_array).is_some() -} - async fn looks_like_hydroserver_data_response(response: Response) -> bool { - if !matches!( - response.status(), - StatusCode::OK | StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN - ) { + if response.status() != StatusCode::OK { return false; } - let Some(payload) = response_json(response).await else { return false; }; - - payload.is_array() - || payload - .as_object() - .map(|object| object.contains_key("detail") || object.contains_key("status")) - .unwrap_or(false) + payload.get("data").and_then(Value::as_array).is_some() + && ["limit", "offset", "totalCount"].iter().all(|key| { + payload + .get("meta") + .and_then(|meta| meta.get(key)) + .and_then(Value::as_u64) + .is_some() + }) } async fn response_json(response: Response) -> Option { @@ -1176,7 +758,12 @@ async fn response_error_message(response: Response) -> String { let text = response.text().await.unwrap_or_default(); let detail = serde_json::from_str::(&text) .ok() - .and_then(|payload| payload.get("detail").cloned()) + .and_then(|payload| { + payload + .get("message") + .or_else(|| payload.get("detail")) + .cloned() + }) .map(format_error_detail); detail.unwrap_or_else(|| format!("HydroServer returned status {}.", status.as_u16())) @@ -1210,7 +797,7 @@ fn datastream_cache_key(server: &ServerConfig, workspace_id: &str) -> String { } fn expanded_datastream_to_summary(item: &Value) -> DatastreamSummary { - let thing = item.get("thing"); + let thing = item.get("monitoringSite"); let observed_property = item .get("observed_property") .or_else(|| item.get("observedProperty")); @@ -1218,21 +805,11 @@ fn expanded_datastream_to_summary(item: &Value) -> DatastreamSummary { .get("processing_level") .or_else(|| item.get("processingLevel")); let unit = item.get("unit"); - let sensor = item.get("sensor"); + let sensor = item.get("method"); - let thing_id = string_value(item, &["thing_id", "thingId"]) + let thing_id = string_value(item, &["monitoring_site_id", "monitoringSiteId"]) .or_else(|| thing.and_then(|thing| string_value(thing, &["id", "uid"]))) .unwrap_or_default(); - let observed_property_id = string_value(item, &["observed_property_id", "observedPropertyId"]) - .or_else(|| observed_property.and_then(|value| string_value(value, &["id", "uid"]))) - .unwrap_or_default(); - let processing_level_id = string_value(item, &["processing_level_id", "processingLevelId"]) - .or_else(|| processing_level.and_then(|value| string_value(value, &["id", "uid"]))) - .unwrap_or_default(); - let unit_id = string_value(item, &["unit_id", "unitId"]) - .or_else(|| unit.and_then(|value| string_value(value, &["id", "uid"]))) - .unwrap_or_default(); - let _ = (observed_property_id, processing_level_id, unit_id); DatastreamSummary { id: string_value(item, &["id", "uid"]).unwrap_or_default(), @@ -1245,7 +822,7 @@ fn expanded_datastream_to_summary(item: &Value) -> DatastreamSummary { .and_then(|value| string_value(value, &["name"])) .unwrap_or_default(), processing_level_definition: processing_level - .and_then(|value| string_value(value, &["definition"])) + .and_then(|value| string_value(value, &["name"])) .unwrap_or_default(), unit_name: unit .and_then(|value| string_value(value, &["name"])) @@ -1263,10 +840,7 @@ fn expanded_datastream_to_summary(item: &Value) -> DatastreamSummary { } fn expanded_datastream_to_detail(item: &Value) -> DatastreamDetail { - let thing = item.get("thing"); - let location = thing - .and_then(|value| value.get("location")) - .or_else(|| thing.and_then(|value| value.get("Location"))); + let thing = item.get("monitoringSite"); let observed_property = item .get("observed_property") .or_else(|| item.get("observedProperty")); @@ -1274,7 +848,7 @@ fn expanded_datastream_to_detail(item: &Value) -> DatastreamDetail { .get("processing_level") .or_else(|| item.get("processingLevel")); let unit = item.get("unit"); - let sensor = item.get("sensor"); + let sensor = item.get("method"); DatastreamDetail { id: string_value(item, &["id", "uid"]).unwrap_or_default(), @@ -1324,7 +898,7 @@ fn expanded_datastream_to_detail(item: &Value) -> DatastreamDetail { is_private: bool_value(item, &["is_private", "isPrivate"]), is_visible: bool_value(item, &["is_visible", "isVisible"]), thing: DatastreamThingDetail { - id: string_value(item, &["thing_id", "thingId"]) + id: string_value(item, &["monitoring_site_id", "monitoringSiteId"]) .or_else(|| thing.and_then(|value| string_value(value, &["id", "uid"]))) .unwrap_or_default(), name: thing @@ -1334,12 +908,10 @@ fn expanded_datastream_to_detail(item: &Value) -> DatastreamDetail { .and_then(|value| string_value(value, &["description"])) .unwrap_or_default(), sampling_feature_code: thing - .and_then(|value| { - string_value(value, &["sampling_feature_code", "samplingFeatureCode"]) - }) + .and_then(|value| string_value(value, &["code"])) .unwrap_or_default(), site_type: thing - .and_then(|value| string_value(value, &["site_type", "siteType"])) + .and_then(|value| string_value(value, &["type"])) .unwrap_or_default(), sampling_feature_type: thing .and_then(|value| { @@ -1350,28 +922,27 @@ fn expanded_datastream_to_detail(item: &Value) -> DatastreamDetail { .map(|value| bool_value(value, &["is_private", "isPrivate"])) .unwrap_or(false), location: DatastreamThingLocationDetail { - latitude: scalar_string_value(location.unwrap_or(item), &["latitude"]), - longitude: scalar_string_value(location.unwrap_or(item), &["longitude"]), - elevation_m: scalar_string_value( - location.unwrap_or(item), - &["elevation_m", "elevationM"], - ), - elevation_datum: string_value( - location.unwrap_or(item), - &["elevation_datum", "elevationDatum"], - ) - .unwrap_or_default(), - admin_area_1: string_value( - location.unwrap_or(item), - &["admin_area_1", "adminArea1"], - ) - .unwrap_or_default(), - admin_area_2: string_value( - location.unwrap_or(item), - &["admin_area_2", "adminArea2"], - ) - .unwrap_or_default(), - country: string_value(location.unwrap_or(item), &["country"]).unwrap_or_default(), + latitude: thing + .map(|value| scalar_string_value(value, &["latitude"])) + .unwrap_or_default(), + longitude: thing + .map(|value| scalar_string_value(value, &["longitude"])) + .unwrap_or_default(), + elevation_m: thing + .map(|value| scalar_string_value(value, &["elevation_m", "elevationM"])) + .unwrap_or_default(), + elevation_datum: thing + .and_then(|value| string_value(value, &["elevation_datum", "elevationDatum"])) + .unwrap_or_default(), + admin_area_1: thing + .and_then(|value| string_value(value, &["admin_area_1", "adminArea1"])) + .unwrap_or_default(), + admin_area_2: thing + .and_then(|value| string_value(value, &["admin_area_2", "adminArea2"])) + .unwrap_or_default(), + country: thing + .and_then(|value| string_value(value, &["country"])) + .unwrap_or_default(), }, }, observed_property: DatastreamObservedPropertyDetail { @@ -1412,7 +983,7 @@ fn expanded_datastream_to_detail(item: &Value) -> DatastreamDetail { .unwrap_or_default(), }, sensor: DatastreamSensorDetail { - id: string_value(item, &["sensor_id", "sensorId"]) + id: string_value(item, &["method_id", "methodId"]) .or_else(|| sensor.and_then(|value| string_value(value, &["id", "uid"]))) .unwrap_or_default(), name: sensor @@ -1422,25 +993,25 @@ fn expanded_datastream_to_detail(item: &Value) -> DatastreamDetail { .and_then(|value| string_value(value, &["description"])) .unwrap_or_default(), manufacturer: sensor - .and_then(|value| string_value(value, &["manufacturer"])) + .and_then(|value| string_value(value, &["sensorModelManufacturer"])) .unwrap_or_default(), model: sensor - .and_then(|value| string_value(value, &["model"])) + .and_then(|value| string_value(value, &["sensorModel"])) .unwrap_or_default(), method_type: sensor - .and_then(|value| string_value(value, &["method_type", "methodType"])) + .and_then(|value| string_value(value, &["type"])) .unwrap_or_default(), method_code: sensor - .and_then(|value| string_value(value, &["method_code", "methodCode"])) + .and_then(|value| string_value(value, &["code"])) .unwrap_or_default(), method_link: sensor - .and_then(|value| string_value(value, &["method_link", "methodLink"])) + .and_then(|value| string_value(value, &["definition"])) .unwrap_or_default(), encoding_type: sensor .and_then(|value| string_value(value, &["encoding_type", "encodingType"])) .unwrap_or_default(), model_link: sensor - .and_then(|value| string_value(value, &["model_link", "modelLink"])) + .and_then(|value| string_value(value, &["sensorModelDefinition"])) .unwrap_or_default(), }, processing_level: DatastreamProcessingLevelDetail { @@ -1451,52 +1022,58 @@ fn expanded_datastream_to_detail(item: &Value) -> DatastreamDetail { .and_then(|value| string_value(value, &["code"])) .unwrap_or_default(), definition: processing_level - .and_then(|value| string_value(value, &["definition"])) + .and_then(|value| string_value(value, &["name"])) .unwrap_or_default(), explanation: processing_level - .and_then(|value| string_value(value, &["explanation"])) + .and_then(|value| string_value(value, &["description"])) .unwrap_or_default(), }, } } -fn datastream_to_summary( - item: &Value, - things_by_id: &HashMap, - observed_properties_by_id: &HashMap, - processing_levels_by_id: &HashMap, - units_by_id: &HashMap, - sensors_by_id: &HashMap, -) -> DatastreamSummary { - let thing_id = string_value(item, &["thing_id", "thingId"]).unwrap_or_default(); - let observed_property_id = - string_value(item, &["observed_property_id", "observedPropertyId"]).unwrap_or_default(); - let processing_level_id = - string_value(item, &["processing_level_id", "processingLevelId"]).unwrap_or_default(); - let unit_id = string_value(item, &["unit_id", "unitId"]).unwrap_or_default(); - let sensor_id = string_value(item, &["sensor_id", "sensorId"]).unwrap_or_default(); +fn response_item(payload: &Value) -> Result<&Value, RequestError> { + payload + .get("data") + .filter(|item| item.is_object()) + .ok_or_else(|| { + RequestError::Other("HydroServer returned an invalid item response.".to_string()) + }) +} - DatastreamSummary { - id: string_value(item, &["id", "uid"]).unwrap_or_default(), - name: string_value(item, &["name"]).unwrap_or_else(|| "Unnamed datastream".to_string()), - thing_id: thing_id.clone(), - thing_name: string_value_from_map(things_by_id, &Some(thing_id), &["name"]), - observed_property_name: string_value_from_map( - observed_properties_by_id, - &Some(observed_property_id), - &["name"], +// Keep the SDL's persisted/UI models stable while resolving v2's side-loaded +// resources. Monitoring sites and methods fill the existing thing/sensor fields. +fn attach_datastream_relations(items: &mut [Value], payload: &Value) { + for (id_field, relation, bucket) in [ + ("monitoringSiteId", "monitoringSite", "monitoringSites"), + ( + "observedPropertyId", + "observedProperty", + "observedProperties", ), - processing_level_definition: string_value_from_map( - processing_levels_by_id, - &Some(processing_level_id), - &["definition"], - ), - unit_name: string_value_from_map(units_by_id, &Some(unit_id.clone()), &["name"]), - unit_symbol: string_value_from_map(units_by_id, &Some(unit_id), &["symbol"]), - sampled_medium: string_value(item, &["sampled_medium", "sampledMedium"]) - .unwrap_or_default(), - sensor_name: string_value_from_map(sensors_by_id, &Some(sensor_id), &["name"]), - result_type: string_value(item, &["result_type", "resultType"]).unwrap_or_default(), + ("processingLevelId", "processingLevel", "processingLevels"), + ("unitId", "unit", "units"), + ("methodId", "method", "methods"), + ] { + let Some(resources) = payload + .get("included") + .and_then(|included| included.get(bucket)) + .and_then(Value::as_array) + else { + continue; + }; + let by_id = map_items_by_id(resources); + for item in items.iter_mut() { + if let Some(related) = item + .get(id_field) + .and_then(Value::as_str) + .and_then(|id| by_id.get(id)) + { + let related = related.clone(); + if let Some(object) = item.as_object_mut() { + object.insert(relation.to_string(), related); + } + } + } } } @@ -1541,24 +1118,6 @@ fn bool_value(item: &Value, keys: &[&str]) -> bool { .unwrap_or(false) } -fn string_value_from_map( - items: &HashMap, - id: &Option, - keys: &[&str], -) -> String { - id.as_ref() - .and_then(|key| items.get(key)) - .and_then(|item| string_value(item, keys)) - .unwrap_or_default() -} - -fn header_int(headers: &HeaderMap, header: &str) -> Option { - headers - .get(header) - .and_then(|value| value.to_str().ok()) - .and_then(|value| value.parse::().ok()) -} - #[cfg(test)] #[path = "tests/hydroserver.rs"] mod tests; diff --git a/src/tests/hydroserver.rs b/src/tests/hydroserver.rs index 725c89c..566a2ab 100644 --- a/src/tests/hydroserver.rs +++ b/src/tests/hydroserver.rs @@ -1,40 +1,6 @@ use super::*; use chrono::TimeZone; -#[test] -fn observation_batch_payload_matches_sensorthings_schema() { - let observations = [ - ObservationPayloadRow { - phenomenon_time: Utc.with_ymd_and_hms(2026, 4, 3, 8, 0, 0).unwrap(), - result: json!(2.41), - }, - ObservationPayloadRow { - phenomenon_time: Utc.with_ymd_and_hms(2026, 4, 3, 8, 5, 0).unwrap(), - result: json!("qualitative"), - }, - ]; - - let body = json!({ - "fields": ["phenomenonTime", "result"], - "data": observations - .iter() - .map(|row| json!([row.phenomenon_time.to_rfc3339(), row.result])) - .collect::>(), - }); - - assert_eq!( - body["fields"], - json!(["phenomenonTime", "result"]), - "fields must use SensorThings naming" - ); - - let data = body["data"].as_array().expect("data should be array"); - assert_eq!(data.len(), 2); - assert_eq!(data[0][0], "2026-04-03T08:00:00+00:00"); - assert_eq!(data[0][1], 2.41); - assert_eq!(data[1][1], "qualitative"); -} - #[test] fn build_url_normalizes_correctly() { assert_eq!( @@ -160,52 +126,453 @@ fn userpass_server(url: &str, username: &str) -> ServerConfig { } } -fn session_for(cache: &Arc>, url: &str, username: &str) -> HydroServerSession { - HydroServerSession::new(Client::new(), userpass_server(url, username), cache.clone()) +// Exercise the real HTTP client against v2 response shapes, including servers +// that return fewer items per page than the requested limit. +use axum::{body::to_bytes, extract::Request, response::IntoResponse, routing::any, Json, Router}; +use std::{collections::VecDeque, sync::Arc}; +use tokio::{net::TcpListener, task::JoinHandle}; + +#[derive(Debug)] +struct RecordedRequest { + method: Method, + url: reqwest::Url, + headers: reqwest::header::HeaderMap, + body: Value, } -#[test] -fn token_cache_key_ignores_trailing_slash_and_surrounding_whitespace() { - let cache = Arc::new(Mutex::new(TokenCache::new())); - let a = session_for(&cache, "https://hydro.example/", " user@example.com "); - let b = session_for(&cache, "https://hydro.example", "user@example.com"); - assert_eq!(a.token_cache_key(), b.token_cache_key()); +struct TestServer { + url: String, + requests: Arc>>, + task: JoinHandle<()>, } -#[test] -fn cached_token_is_reused_across_sessions_for_the_same_account() { - let cache = Arc::new(Mutex::new(TokenCache::new())); +impl TestServer { + async fn spawn(responses: Vec<(u16, Value)>) -> Self { + let requests = Arc::new(Mutex::new(Vec::new())); + let responses = Arc::new(Mutex::new(VecDeque::from(responses))); + let handler = { + let requests = requests.clone(); + move |request: Request| { + let requests = requests.clone(); + let responses = responses.clone(); + async move { + let (parts, body) = request.into_parts(); + let bytes = to_bytes(body, 1024 * 1024).await.unwrap(); + requests.lock().unwrap().push(RecordedRequest { + method: parts.method, + url: reqwest::Url::parse(&format!("http://localhost{}", parts.uri)) + .unwrap(), + headers: parts.headers, + body: if bytes.is_empty() { + Value::Null + } else { + serde_json::from_slice(&bytes).unwrap() + }, + }); + let (status, body) = responses + .lock() + .unwrap() + .pop_front() + .unwrap_or((500, json!({"message": "Unexpected request"}))); + (StatusCode::from_u16(status).unwrap(), Json(body)).into_response() + } + } + }; + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}", listener.local_addr().unwrap()); + let app = Router::new().fallback(any(handler)); + let task = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + Self { + url, + requests, + task, + } + } - session_for(&cache, "https://hydro.example", "user@example.com").store_token("tok-123"); + fn config(&self) -> ServerConfig { + ServerConfig { + url: self.url.clone(), + ..server_with_workspace("workspace-1") + } + } +} - let next = session_for(&cache, "https://hydro.example", "user@example.com"); - assert_eq!(next.cached_token().as_deref(), Some("tok-123")); +impl Drop for TestServer { + fn drop(&mut self) { + self.task.abort(); + } } -#[test] -fn cached_token_is_scoped_per_account() { - let cache = Arc::new(Mutex::new(TokenCache::new())); - session_for(&cache, "https://hydro.example", "alice@example.com").store_token("alice-tok"); +fn collection(items: Value, offset: u64, total: u64) -> Value { + json!({"data": items, "meta": {"offset": offset, "limit": 1, "totalCount": total}}) +} - let bob = session_for(&cache, "https://hydro.example", "bob@example.com"); - assert_eq!(bob.cached_token(), None); +fn query(request: &RecordedRequest, name: &str) -> Option { + request + .url + .query_pairs() + .find(|(key, _)| key == name) + .map(|(_, value)| value.into_owned()) } -#[test] -fn invalidate_token_clears_the_shared_cache_for_other_sessions() { - let cache = Arc::new(Mutex::new(TokenCache::new())); +#[tokio::test] +async fn validates_v2_instance_with_no_public_workspaces() { + let server = TestServer::spawn(vec![(200, collection(json!([]), 0, 0))]).await; + let service = HydroServerService::new().unwrap(); + assert!(service.validate_url(&format!("{}/", server.url)).await.ok); + let requests = server.requests.lock().unwrap(); + assert_eq!(requests.len(), 1); + assert_eq!(requests[0].url.path(), "/api/data/workspaces"); + assert_eq!(query(&requests[0], "limit").as_deref(), Some("1")); + assert!(!requests[0].headers.contains_key("authorization")); +} - let mut session = session_for(&cache, "https://hydro.example", "user@example.com"); - session.store_token("stale"); - assert_eq!(session.cached_token().as_deref(), Some("stale")); +#[tokio::test] +async fn rejects_unrecognized_server_responses() { + for body in [ + json!([]), + json!({"data": []}), + json!({"message": "Not found"}), + ] { + let server = TestServer::spawn(vec![(200, body)]).await; + assert!( + !HydroServerService::new() + .unwrap() + .validate_url(&server.url) + .await + .ok + ); + } +} - session.invalidate_token(); +#[tokio::test] +async fn api_key_connection_reads_all_workspace_pages() { + let server = TestServer::spawn(vec![ + ( + 200, + collection(json!([{"id": "other", "name": "Other"}]), 0, 2), + ), + ( + 200, + collection( + json!([{"id": "workspace-1", "name": "Test workspace"}]), + 1, + 2, + ), + ), + ]) + .await; + let response = HydroServerService::new() + .unwrap() + .test_connection(&server.config()) + .await; + assert!(response.ok, "{}", response.message); + assert_eq!(response.workspace_count, 2); + assert_eq!(response.workspace_id.as_deref(), Some("workspace-1")); + let requests = server.requests.lock().unwrap(); + assert_eq!(requests.len(), 2); + for (index, request) in requests.iter().enumerate() { + assert_eq!(request.headers["x-api-key"], "test-key"); + assert!(!request.headers.contains_key("authorization")); + assert_eq!(query(request, "is_associated").as_deref(), Some("true")); + assert_eq!(query(request, "offset"), Some(index.to_string())); + assert_eq!(query(request, "limit").as_deref(), Some("1000")); + assert_eq!(query(request, "page"), None); + } +} - assert_eq!(session.cached_token(), None, "own cache lookup is cleared"); - let other = session_for(&cache, "https://hydro.example", "user@example.com"); +#[tokio::test] +async fn username_password_uses_basic_auth_without_session_login() { + let server = TestServer::spawn(vec![ + ( + 200, + collection( + json!([{"id": "workspace-1", "name": "Test workspace"}]), + 0, + 1, + ), + ), + (401, json!({"message": "Invalid username or password"})), + ]) + .await; + let mut config = userpass_server(&server.url, "user@example.com"); + config.workspace_name = "Test workspace".to_string(); + let service = HydroServerService::new().unwrap(); + assert!(service.test_connection(&config).await.ok); + config.password = "wrong".to_string(); + let rejected = service.test_connection(&config).await; + assert!(!rejected.ok); + assert_eq!(rejected.invalid_field.as_deref(), Some("username")); + let requests = server.requests.lock().unwrap(); + assert_eq!( + requests.len(), + 2, + "invalid credentials must not trigger a session login" + ); + assert_eq!( + requests[0].headers["authorization"], + "Basic dXNlckBleGFtcGxlLmNvbTpzZWNyZXQ=" + ); + assert_ne!( + requests[0].headers["authorization"], + requests[1].headers["authorization"] + ); + assert!(requests + .iter() + .all(|r| r.method == Method::GET && r.url.path() == "/api/data/workspaces")); +} + +fn datastream_item() -> Value { + json!({ + "id": "ds-1", "name": "Stage", "description": "River stage", "workspaceId": "workspace-1", + "monitoringSiteId": "site-1", "methodId": "method-1", "observedPropertyId": "op-1", + "processingLevelId": "pl-1", "unitId": "unit-1", "sampledMedium": "Water", + "resultType": "Time series", "observationType": "OM_Measurement", "noDataValue": -9999, + "phenomenonEndTime": "2026-09-24T12:00:00Z", "valueCount": 7, "isVisible": true + }) +} + +fn included_resources() -> Value { + json!({ + "monitoringSites": [ + {"id": "other", "name": "Wrong site"}, + {"id": "site-1", "name": "River", "code": "R1", "type": "Stream", "description": "River site", + "latitude": 41.5, "longitude": -111.5, "elevation_m": 1400, "elevationDatum": "NAVD88", + "adminArea1": "Utah", "adminArea2": "Cache", "country": "US", "isPrivate": true} + ], + "methods": [{"id": "method-1", "name": "Pressure", "code": "P1", "type": "Instrument", + "description": "Pressure transducer", "definition": "https://example.com/method", + "sensorModel": "Model A", "sensorModelManufacturer": "Manufacturer", + "sensorModelDefinition": "https://example.com/model"}], + "observedProperties": [{"id": "op-1", "name": "Stage", "type": "Depth", "code": "STAGE", "description": "Water level"}], + "processingLevels": [{"id": "pl-1", "name": "Raw", "description": "Unprocessed measurements", "code": "0", "definition": "https://example.com/level"}], + "units": [{"id": "unit-1", "name": "Meter", "symbol": "m", "type": "Length"}] + }) +} + +#[tokio::test] +async fn datastream_pages_resolve_included_metadata_and_cache_results() { + let mut first = collection(json!([datastream_item()]), 0, 2); + first["included"] = included_resources(); + let mut next_item = datastream_item(); + next_item["id"] = json!("ds-2"); + let mut second = collection(json!([next_item]), 1, 2); + second["included"] = included_resources(); + second["included"]["monitoringSites"][1]["name"] = json!("Second page site"); + let server = TestServer::spawn(vec![ + (200, first), + (200, second), + (200, collection(json!([]), 0, 0)), + ]) + .await; + let service = HydroServerService::new().unwrap(); + let config = server.config(); + let items = service.list_datastreams(&config, false).await.unwrap(); + assert_eq!(items.len(), 2); + assert_eq!(items[0].thing_id, "site-1"); + assert_eq!(items[0].thing_name, "River"); + assert_eq!(items[0].sensor_name, "Pressure"); + assert_eq!(items[0].observed_property_name, "Stage"); + assert_eq!(items[0].processing_level_definition, "Raw"); + assert_eq!(items[0].unit_name, "Meter"); + assert_eq!(items[0].unit_symbol, "m"); + assert_eq!(items[0].sampled_medium, "Water"); + assert_eq!(items[0].result_type, "Time series"); + assert_eq!(items[1].thing_name, "Second page site"); + assert_eq!( + service + .list_datastreams(&config, false) + .await + .unwrap() + .len(), + 2 + ); + assert_eq!(server.requests.lock().unwrap().len(), 2); + assert!(service + .list_datastreams(&config, true) + .await + .unwrap() + .is_empty()); + let requests = server.requests.lock().unwrap(); + assert_eq!(requests.len(), 3); + for request in requests.iter() { + assert_eq!(request.url.path(), "/api/data/datastreams"); + assert_eq!( + query(request, "include").as_deref(), + Some(DATASTREAM_INCLUDE) + ); + assert_eq!( + query(request, "workspace_id").as_deref(), + Some("workspace-1") + ); + assert_eq!(query(request, "expand_related"), None); + } + assert_eq!(query(&requests[1], "offset").as_deref(), Some("1")); +} + +#[tokio::test] +async fn datastream_detail_maps_v2_metadata_into_existing_ui_fields() { + let server = TestServer::spawn(vec![( + 200, + json!({"data": datastream_item(), "included": included_resources()}), + )]) + .await; + let detail = HydroServerService::new() + .unwrap() + .get_datastream_detail(&server.config(), "ds-1") + .await + .unwrap(); + assert_eq!(detail.id, "ds-1"); + assert_eq!(detail.no_data_value, "-9999"); + assert_eq!(detail.value_count, "7"); + assert!(detail.is_visible); + assert_eq!(detail.thing.id, "site-1"); + assert_eq!(detail.thing.name, "River"); + assert_eq!(detail.thing.sampling_feature_code, "R1"); + assert_eq!(detail.thing.site_type, "Stream"); + assert!(detail.thing.is_private); + assert_eq!(detail.thing.location.latitude, "41.5"); + assert_eq!(detail.thing.location.longitude, "-111.5"); + assert_eq!(detail.thing.location.elevation_m, "1400"); + assert_eq!(detail.thing.location.elevation_datum, "NAVD88"); + assert_eq!(detail.thing.location.admin_area_1, "Utah"); + assert_eq!(detail.thing.location.admin_area_2, "Cache"); + assert_eq!(detail.thing.location.country, "US"); + assert_eq!(detail.sensor.id, "method-1"); + assert_eq!(detail.sensor.name, "Pressure"); + assert_eq!(detail.sensor.manufacturer, "Manufacturer"); + assert_eq!(detail.sensor.model, "Model A"); + assert_eq!(detail.sensor.method_type, "Instrument"); + assert_eq!(detail.sensor.method_code, "P1"); + assert_eq!(detail.sensor.method_link, "https://example.com/method"); + assert_eq!(detail.sensor.model_link, "https://example.com/model"); + assert_eq!(detail.processing_level.definition, "Raw"); + assert_eq!( + detail.processing_level.explanation, + "Unprocessed measurements" + ); + assert_eq!(detail.unit.symbol, "m"); + assert_eq!(detail.observed_property.code, "STAGE"); + let requests = server.requests.lock().unwrap(); + assert_eq!(requests[0].url.path(), "/api/data/datastreams/ds-1"); + assert_eq!( + query(&requests[0], "include").as_deref(), + Some(DATASTREAM_INCLUDE) + ); +} + +#[tokio::test] +async fn malformed_collections_fail_instead_of_silently_hiding_datastreams() { + for payload in [json!([]), json!({"data": []}), collection(json!([]), 0, 1)] { + let server = TestServer::spawn(vec![(200, payload)]).await; + assert!(HydroServerService::new() + .unwrap() + .list_datastreams(&server.config(), false) + .await + .is_err()); + } +} + +#[test] +fn missing_monitoring_site_does_not_use_datastream_location_fields() { + let item = json!({ + "id": "ds-1", "monitoringSiteId": "site-1", + "latitude": 41.5, "longitude": -111.5, "elevation_m": 1400, + "elevationDatum": "NAVD88", "adminArea1": "Utah", + "adminArea2": "Cache", "country": "US" + }); + let detail = expanded_datastream_to_detail(&item); + assert_eq!(detail.thing.id, "site-1"); + let location = detail.thing.location; + assert!(location.latitude.is_empty()); + assert!(location.longitude.is_empty()); + assert!(location.elevation_m.is_empty()); + assert!(location.elevation_datum.is_empty()); + assert!(location.admin_area_1.is_empty()); + assert!(location.admin_area_2.is_empty()); + assert!(location.country.is_empty()); +} + +#[tokio::test] +async fn reads_conflict_watermark_from_item_envelope() { + let server = TestServer::spawn(vec![ + ( + 200, + json!({"data": {"phenomenonEndTime": "2026-09-24T12:00:00Z"}}), + ), + (200, json!({"data": {"phenomenonEndTime": null}})), + (200, json!({})), + ]) + .await; + let service = HydroServerService::new().unwrap(); + assert_eq!( + service + .fetch_phenomenon_end_time(&server.config(), "ds-1") + .await + .unwrap(), + Some(Utc.with_ymd_and_hms(2026, 9, 24, 12, 0, 0).unwrap()) + ); + assert_eq!( + service + .fetch_phenomenon_end_time(&server.config(), "ds-1") + .await + .unwrap(), + None + ); + assert!(service + .fetch_phenomenon_end_time(&server.config(), "ds-1") + .await + .is_err()); +} + +#[tokio::test] +async fn uploads_to_v2_bulk_endpoint_and_preserves_insert_conflicts() { + let server = TestServer::spawn(vec![ + (201, Value::Null), + ( + 409, + json!({"message": "Duplicate phenomenonTime found on this datastream."}), + ), + ]) + .await; + let service = HydroServerService::new().unwrap(); + let rows = vec![ + ObservationPayloadRow { + phenomenon_time: Utc.with_ymd_and_hms(2026, 9, 24, 12, 0, 0).unwrap(), + result: json!(2.41), + }, + ObservationPayloadRow { + phenomenon_time: Utc.with_ymd_and_hms(2026, 9, 24, 12, 5, 0).unwrap(), + result: Value::Null, + }, + ]; + service + .post_observations_batch(&server.config(), "ds-1", &[]) + .await + .unwrap(); + assert!(server.requests.lock().unwrap().is_empty()); + service + .post_observations_batch(&server.config(), "ds-1", &rows) + .await + .unwrap(); + let conflict = service + .post_observations_batch(&server.config(), "ds-1", &rows) + .await + .unwrap_err(); + assert!(conflict.is_conflict()); + assert!(conflict.to_string().contains("Duplicate phenomenonTime")); + let requests = server.requests.lock().unwrap(); + assert_eq!(requests.len(), 2); + assert_eq!(requests[0].method, Method::POST); + assert_eq!(requests[0].url.path(), "/api/data/observations/bulk-create"); + assert_eq!(query(&requests[0], "mode").as_deref(), Some("insert")); + assert_eq!(requests[0].headers["x-api-key"], "test-key"); assert_eq!( - other.cached_token(), - None, - "the stale token is gone for future sessions too" + requests[0].body, + json!({ + "datastreamId": "ds-1", "fields": ["phenomenonTime", "result"], + "data": [["2026-09-24T12:00:00+00:00", 2.41], ["2026-09-24T12:05:00+00:00", null]] + }) ); + assert_eq!(requests[0].body, requests[1].body); } diff --git a/src/tests/uploader.rs b/src/tests/uploader.rs index a04195e..0da6005 100644 --- a/src/tests/uploader.rs +++ b/src/tests/uploader.rs @@ -38,8 +38,8 @@ impl TestObservationServer { } /// Spawn a mock that pops `statuses` once per request. Successful responses - /// to GET requests return `get_body` (used to serve `phenomenonEndTime` for - /// conflict-reconciliation tests); successful POSTs return `{}`. + /// to GET requests wrap `get_body` in the v2 item envelope (used to serve + /// `phenomenonEndTime` for conflict reconciliation); successful POSTs return `{}`. async fn spawn_with_get_body(statuses: Vec, get_body: String) -> Self { let listener = TcpListener::bind("127.0.0.1:0") .await @@ -72,9 +72,11 @@ impl TestObservationServer { let status = statuses.lock().await.pop_front().unwrap_or(200); let payload = if status >= 400 { - json!({ "detail": "temporary outage" }).to_string() + json!({ "message": "temporary outage" }).to_string() } else if method == "GET" { - get_body.clone() + let data: serde_json::Value = + serde_json::from_str(&get_body).expect("valid item"); + json!({ "data": data }).to_string() } else { "{}".to_string() };