use crate::model::AcsValue;
use bamcensus_core::{
model::identifier::{Geoid, GeoidType},
ops::agg::NumericAggregation,
};
use itertools::Itertools;
use serde_json::json;
pub fn aggregate_acs(
rows: &[(Geoid, Vec<AcsValue>)],
target: GeoidType,
agg: NumericAggregation,
) -> Result<Vec<(Geoid, Vec<AcsValue>)>, String> {
let (geoid_oks, geoid_errs): (Vec<(Geoid, &Vec<AcsValue>)>, Vec<String>) = rows
.iter()
.map(|(geoid, values)| {
let trunc_geoid = geoid.truncate_geoid_to_type(&target)?;
Ok((trunc_geoid, values))
})
.partition_result();
if !geoid_errs.is_empty() {
let msg = geoid_errs.into_iter().unique().take(5).join("\n");
return Err(format!(
"errors during aggregation. first 5 unique errors: \n{msg}"
));
}
let mut geoids_grouped = vec![];
let grouping_iter = geoid_oks.into_iter().chunk_by(|(g, _)| g.clone());
for (geoid, grouped) in &grouping_iter {
let vs = grouped.into_iter().flat_map(|(_, v)| v).collect_vec();
geoids_grouped.push((geoid, vs));
}
let reduced = geoids_grouped
.into_iter()
.map(|(geoid, values)| {
let xs = values.into_iter().chunk_by(|v| v.name.clone());
let mut agg_values = vec![];
for (name, values) in &xs {
let values = values.map(|v| {
v.value.as_f64().ok_or_else(|| format!("ACS value for {} is not numeric (found {}) but user requested aggregation", name, v.value))
})
.collect::<Result<Vec<_>, _>>()?;
let aggregated = agg.aggregate(&mut values.into_iter());
agg_values.push(AcsValue::new(name, json![aggregated]));
}
Ok((geoid, agg_values))
})
.collect::<Result<Vec<_>, String>>()?;
Ok(reduced)
}