use cobre_core::{EntityId, entities::PumpingStation};
use serde::Deserialize;
use std::collections::HashSet;
use std::path::Path;
use super::parse_operational_start_date;
use crate::LoadError;
#[derive(Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub(crate) struct RawPumpingFile {
#[serde(rename = "$schema")]
_schema: Option<String>,
pumping_stations: Vec<RawPumpingStation>,
}
#[derive(Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub(crate) struct RawPumpingStation {
id: i32,
name: String,
operational_start_date: String,
bus_id: i32,
source_hydro_id: i32,
destination_hydro_id: i32,
#[serde(default)]
entry_stage_id: Option<i32>,
#[serde(default)]
exit_stage_id: Option<i32>,
consumption_mw_per_m3s: f64,
flow: RawPumpingFlow,
}
#[derive(Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub(crate) struct RawPumpingFlow {
min_m3s: f64,
max_m3s: f64,
}
pub fn parse_pumping_stations(path: &Path) -> Result<Vec<PumpingStation>, LoadError> {
let raw_text = std::fs::read_to_string(path).map_err(|e| LoadError::io(path, e))?;
let raw: RawPumpingFile =
serde_json::from_str(&raw_text).map_err(|e| LoadError::parse(path, e.to_string()))?;
validate_raw_pumping(&raw, path)?;
convert_pumping(raw, path)
}
fn validate_raw_pumping(raw: &RawPumpingFile, path: &Path) -> Result<(), LoadError> {
validate_no_duplicate_pumping_ids(&raw.pumping_stations, path)?;
for (i, station) in raw.pumping_stations.iter().enumerate() {
validate_consumption(station.consumption_mw_per_m3s, i, path)?;
validate_flow_bounds(&station.flow, i, path)?;
}
Ok(())
}
fn validate_no_duplicate_pumping_ids(
stations: &[RawPumpingStation],
path: &Path,
) -> Result<(), LoadError> {
let mut seen: HashSet<i32> = HashSet::new();
for (i, station) in stations.iter().enumerate() {
if !seen.insert(station.id) {
return Err(LoadError::SchemaError {
path: path.to_path_buf(),
field: format!("pumping_stations[{i}].id"),
message: format!("duplicate id {} in pumping_stations array", station.id),
});
}
}
Ok(())
}
fn validate_consumption(
consumption_mw_per_m3s: f64,
station_index: usize,
path: &Path,
) -> Result<(), LoadError> {
if consumption_mw_per_m3s < 0.0 {
return Err(LoadError::SchemaError {
path: path.to_path_buf(),
field: format!("pumping_stations[{station_index}].consumption_mw_per_m3s"),
message: format!("consumption_mw_per_m3s must be >= 0.0, got {consumption_mw_per_m3s}"),
});
}
Ok(())
}
fn validate_flow_bounds(
flow: &RawPumpingFlow,
station_index: usize,
path: &Path,
) -> Result<(), LoadError> {
if flow.min_m3s < 0.0 {
return Err(LoadError::SchemaError {
path: path.to_path_buf(),
field: format!("pumping_stations[{station_index}].flow.min_m3s"),
message: format!("flow.min_m3s must be >= 0.0, got {}", flow.min_m3s),
});
}
if flow.max_m3s < 0.0 {
return Err(LoadError::SchemaError {
path: path.to_path_buf(),
field: format!("pumping_stations[{station_index}].flow.max_m3s"),
message: format!("flow.max_m3s must be >= 0.0, got {}", flow.max_m3s),
});
}
if flow.max_m3s < flow.min_m3s {
return Err(LoadError::SchemaError {
path: path.to_path_buf(),
field: format!("pumping_stations[{station_index}].flow.max_m3s"),
message: format!(
"flow.max_m3s ({}) must be >= flow.min_m3s ({})",
flow.max_m3s, flow.min_m3s
),
});
}
Ok(())
}
fn convert_pumping(raw: RawPumpingFile, path: &Path) -> Result<Vec<PumpingStation>, LoadError> {
let mut stations: Vec<PumpingStation> = raw
.pumping_stations
.into_iter()
.enumerate()
.map(|(i, raw_station)| {
let operational_start_date = parse_operational_start_date(
&raw_station.operational_start_date,
path,
&format!("pumping_stations[{i}].operational_start_date"),
)?;
Ok(PumpingStation {
id: EntityId(raw_station.id),
name: raw_station.name,
operational_start_date,
bus_id: EntityId(raw_station.bus_id),
source_hydro_id: EntityId(raw_station.source_hydro_id),
destination_hydro_id: EntityId(raw_station.destination_hydro_id),
entry_stage_id: raw_station.entry_stage_id,
exit_stage_id: raw_station.exit_stage_id,
consumption_mw_per_m3s: raw_station.consumption_mw_per_m3s,
min_flow_m3s: raw_station.flow.min_m3s,
max_flow_m3s: raw_station.flow.max_m3s,
})
})
.collect::<Result<_, LoadError>>()?;
stations.sort_by_key(|s| s.id.0);
Ok(stations)
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::panic, clippy::too_many_lines)]
mod tests {
use super::*;
use std::io::Write;
use tempfile::NamedTempFile;
fn write_json(content: &str) -> NamedTempFile {
let mut f = NamedTempFile::new().unwrap();
f.write_all(content.as_bytes()).unwrap();
f
}
#[test]
fn test_parse_valid_pumping_stations() {
let json = r#"{
"$schema": "https://raw.githubusercontent.com/cobre-rs/cobre/refs/heads/main/schemas/pumping_stations.schema.json",
"pumping_stations": [
{
"id": 0,
"name": "Bombeamento Serra da Mesa",
"operational_start_date": "2024-01-01",
"bus_id": 10,
"source_hydro_id": 3,
"destination_hydro_id": 5,
"consumption_mw_per_m3s": 0.5,
"flow": { "min_m3s": 0.0, "max_m3s": 150.0 }
},
{
"id": 1,
"name": "Bombeamento Cana Brava",
"operational_start_date": "2024-01-01",
"bus_id": 11,
"source_hydro_id": 4,
"destination_hydro_id": 6,
"entry_stage_id": 6,
"exit_stage_id": 120,
"consumption_mw_per_m3s": 0.8,
"flow": { "min_m3s": 10.0, "max_m3s": 200.0 }
}
]
}"#;
let f = write_json(json);
let stations = parse_pumping_stations(f.path()).unwrap();
assert_eq!(stations.len(), 2);
assert_eq!(stations[0].id, EntityId(0));
assert_eq!(stations[0].name, "Bombeamento Serra da Mesa");
assert_eq!(stations[0].bus_id, EntityId(10));
assert_eq!(stations[0].source_hydro_id, EntityId(3));
assert_eq!(stations[0].destination_hydro_id, EntityId(5));
assert_eq!(stations[0].entry_stage_id, None);
assert_eq!(stations[0].exit_stage_id, None);
assert!((stations[0].consumption_mw_per_m3s - 0.5).abs() < f64::EPSILON);
assert!((stations[0].min_flow_m3s - 0.0).abs() < f64::EPSILON);
assert!((stations[0].max_flow_m3s - 150.0).abs() < f64::EPSILON);
assert_eq!(stations[1].id, EntityId(1));
assert_eq!(stations[1].entry_stage_id, Some(6));
assert_eq!(stations[1].exit_stage_id, Some(120));
assert!((stations[1].consumption_mw_per_m3s - 0.8).abs() < f64::EPSILON);
assert!((stations[1].min_flow_m3s - 10.0).abs() < f64::EPSILON);
assert!((stations[1].max_flow_m3s - 200.0).abs() < f64::EPSILON);
}
#[test]
fn test_duplicate_pumping_station_id() {
let json = r#"{
"pumping_stations": [
{
"id": 2, "name": "Alpha", "operational_start_date": "2024-01-01", "bus_id": 0,
"source_hydro_id": 1, "destination_hydro_id": 2,
"consumption_mw_per_m3s": 0.5,
"flow": { "min_m3s": 0.0, "max_m3s": 100.0 }
},
{
"id": 2, "name": "Beta", "operational_start_date": "2024-01-01", "bus_id": 1,
"source_hydro_id": 2, "destination_hydro_id": 3,
"consumption_mw_per_m3s": 0.6,
"flow": { "min_m3s": 0.0, "max_m3s": 200.0 }
}
]
}"#;
let f = write_json(json);
let err = parse_pumping_stations(f.path()).unwrap_err();
match &err {
LoadError::SchemaError { field, message, .. } => {
assert!(
field.contains("pumping_stations[1].id"),
"field should contain 'pumping_stations[1].id', got: {field}"
);
assert!(
message.contains("duplicate"),
"message should contain 'duplicate', got: {message}"
);
}
other => panic!("expected SchemaError, got: {other:?}"),
}
}
#[test]
fn test_negative_consumption() {
let json = r#"{
"pumping_stations": [
{
"id": 0, "name": "Bad", "operational_start_date": "2024-01-01", "bus_id": 0,
"source_hydro_id": 1, "destination_hydro_id": 2,
"consumption_mw_per_m3s": -0.5,
"flow": { "min_m3s": 0.0, "max_m3s": 100.0 }
}
]
}"#;
let f = write_json(json);
let err = parse_pumping_stations(f.path()).unwrap_err();
match &err {
LoadError::SchemaError { field, message, .. } => {
assert!(
field.contains("consumption_mw_per_m3s"),
"field should contain 'consumption_mw_per_m3s', got: {field}"
);
assert!(
message.contains(">= 0.0"),
"message should mention >= 0.0, got: {message}"
);
}
other => panic!("expected SchemaError, got: {other:?}"),
}
}
#[test]
fn test_negative_flow_min() {
let json = r#"{
"pumping_stations": [
{
"id": 0, "name": "Bad", "operational_start_date": "2024-01-01", "bus_id": 0,
"source_hydro_id": 1, "destination_hydro_id": 2,
"consumption_mw_per_m3s": 0.5,
"flow": { "min_m3s": -10.0, "max_m3s": 100.0 }
}
]
}"#;
let f = write_json(json);
let err = parse_pumping_stations(f.path()).unwrap_err();
match &err {
LoadError::SchemaError { field, message, .. } => {
assert!(
field.contains("flow.min_m3s"),
"field should contain 'flow.min_m3s', got: {field}"
);
assert!(
message.contains(">= 0.0"),
"message should mention >= 0.0, got: {message}"
);
}
other => panic!("expected SchemaError, got: {other:?}"),
}
}
#[test]
fn test_max_flow_less_than_min_flow() {
let json = r#"{
"pumping_stations": [
{
"id": 0, "name": "Bad", "operational_start_date": "2024-01-01", "bus_id": 0,
"source_hydro_id": 1, "destination_hydro_id": 2,
"consumption_mw_per_m3s": 0.5,
"flow": { "min_m3s": 200.0, "max_m3s": 100.0 }
}
]
}"#;
let f = write_json(json);
let err = parse_pumping_stations(f.path()).unwrap_err();
match &err {
LoadError::SchemaError { field, message, .. } => {
assert!(
field.contains("flow.max_m3s"),
"field should contain 'flow.max_m3s', got: {field}"
);
assert!(
message.contains("min_m3s"),
"message should mention min_m3s, got: {message}"
);
}
other => panic!("expected SchemaError, got: {other:?}"),
}
}
#[test]
fn test_declaration_order_invariance() {
let json_forward = r#"{
"pumping_stations": [
{
"id": 0, "name": "Alpha", "operational_start_date": "2024-01-01", "bus_id": 0,
"source_hydro_id": 1, "destination_hydro_id": 2,
"consumption_mw_per_m3s": 0.5,
"flow": { "min_m3s": 0.0, "max_m3s": 100.0 }
},
{
"id": 1, "name": "Beta", "operational_start_date": "2024-01-01", "bus_id": 1,
"source_hydro_id": 2, "destination_hydro_id": 3,
"consumption_mw_per_m3s": 0.8,
"flow": { "min_m3s": 0.0, "max_m3s": 200.0 }
}
]
}"#;
let json_reversed = r#"{
"pumping_stations": [
{
"id": 1, "name": "Beta", "operational_start_date": "2024-01-01", "bus_id": 1,
"source_hydro_id": 2, "destination_hydro_id": 3,
"consumption_mw_per_m3s": 0.8,
"flow": { "min_m3s": 0.0, "max_m3s": 200.0 }
},
{
"id": 0, "name": "Alpha", "operational_start_date": "2024-01-01", "bus_id": 0,
"source_hydro_id": 1, "destination_hydro_id": 2,
"consumption_mw_per_m3s": 0.5,
"flow": { "min_m3s": 0.0, "max_m3s": 100.0 }
}
]
}"#;
let f1 = write_json(json_forward);
let f2 = write_json(json_reversed);
let stations1 = parse_pumping_stations(f1.path()).unwrap();
let stations2 = parse_pumping_stations(f2.path()).unwrap();
assert_eq!(
stations1, stations2,
"results must be identical regardless of input ordering"
);
assert_eq!(stations1[0].id, EntityId(0));
assert_eq!(stations1[1].id, EntityId(1));
}
#[test]
fn test_file_not_found() {
let path = Path::new("/nonexistent/system/pumping_stations.json");
let err = parse_pumping_stations(path).unwrap_err();
match &err {
LoadError::IoError { path: p, .. } => {
assert_eq!(p, path);
}
other => panic!("expected IoError, got: {other:?}"),
}
}
#[test]
fn test_invalid_json() {
let f = write_json(r#"{"pumping_stations": [not valid json}}"#);
let err = parse_pumping_stations(f.path()).unwrap_err();
assert!(
matches!(err, LoadError::ParseError { .. }),
"expected ParseError for invalid JSON, got: {err:?}"
);
}
#[test]
fn test_empty_pumping_stations_array() {
let json = r#"{ "pumping_stations": [] }"#;
let f = write_json(json);
let stations = parse_pumping_stations(f.path()).unwrap();
assert!(stations.is_empty());
}
#[test]
fn test_min_equals_max_flow_is_valid() {
let json = r#"{
"pumping_stations": [
{
"id": 0, "name": "Alpha", "operational_start_date": "2024-01-01", "bus_id": 0,
"source_hydro_id": 1, "destination_hydro_id": 2,
"consumption_mw_per_m3s": 0.5,
"flow": { "min_m3s": 100.0, "max_m3s": 100.0 }
}
]
}"#;
let f = write_json(json);
let result = parse_pumping_stations(f.path());
assert!(
result.is_ok(),
"min_m3s == max_m3s should be valid, got: {result:?}"
);
}
}