use std::collections::HashMap;
use std::path::Path;
use futures::StreamExt;
use tracing::{debug, info, warn};
use super::streaming::ManifestStreamer;
use super::types::{DatasetSummary, ManifestConfig};
use crate::app::models::{FileInfo, QualityControlVersion};
use crate::errors::{ManifestError, ManifestResult};
pub async fn collect_datasets_and_years<P: AsRef<Path>>(
manifest_path: P,
) -> ManifestResult<HashMap<String, DatasetSummary>> {
let config = ManifestConfig::default();
let manifest_version =
ManifestStreamer::extract_manifest_version_from_path(manifest_path.as_ref());
let mut streamer = ManifestStreamer::with_config_and_version(config, manifest_version);
let mut stream = streamer.stream(manifest_path).await?;
let mut datasets: HashMap<String, DatasetSummary> = HashMap::new();
while let Some(result) = stream.next().await {
let file_info = result?;
let dataset_info = &file_info.dataset_info;
let dataset_name = dataset_info.dataset_name.clone();
let entry = datasets
.entry(dataset_name.clone())
.or_insert_with(|| DatasetSummary {
name: dataset_name,
versions: Vec::new(),
counties: Vec::new(),
quality_versions: Vec::new(),
years: Vec::new(),
file_count: 0,
example_file: None,
});
let should_count = match &dataset_info.quality_version {
Some(qv) => qv == &QualityControlVersion::V1, None => true, };
if should_count {
entry.file_count += 1;
if entry.example_file.is_none() {
entry.example_file = Some(file_info.relative_path.clone());
}
}
let version = &dataset_info.version;
if !entry.versions.contains(version) {
entry.versions.push(version.clone());
}
if let Some(ref county) = dataset_info.county {
if !entry.counties.contains(county) {
entry.counties.push(county.clone());
}
}
if let Some(ref qv) = dataset_info.quality_version {
if !entry.quality_versions.contains(qv) {
entry.quality_versions.push(qv.clone());
}
}
if let Some(ref year) = dataset_info.year {
if !entry.years.contains(year) {
entry.years.push(year.clone());
}
}
}
for summary in datasets.values_mut() {
summary.versions.sort();
summary.counties.sort();
summary.years.sort();
summary.quality_versions.sort_by_key(|qv| match qv {
QualityControlVersion::V0 => 0,
QualityControlVersion::V1 => 1,
});
}
debug!("Discovered {} datasets from manifest", datasets.len());
Ok(datasets)
}
pub async fn filter_manifest_files<P: AsRef<Path>>(
manifest_path: P,
dataset_name: Option<&str>,
county: Option<&str>,
quality_version: &QualityControlVersion,
) -> ManifestResult<Vec<FileInfo>> {
let config = ManifestConfig::default();
let manifest_version =
ManifestStreamer::extract_manifest_version_from_path(manifest_path.as_ref());
let mut streamer = ManifestStreamer::with_config_and_version(config, manifest_version);
let mut stream = streamer.stream(manifest_path).await?;
let mut filtered_files = Vec::new();
while let Some(result) = stream.next().await {
let file_info = result?;
let dataset_info = &file_info.dataset_info;
if let Some(filter_dataset) = dataset_name {
if dataset_info.dataset_name != filter_dataset {
continue;
}
}
if let Some(filter_county) = county {
if dataset_info.county.as_ref() != Some(&filter_county.to_string()) {
continue;
}
}
if let Some(ref file_qv) = dataset_info.quality_version {
if file_qv != quality_version {
continue;
}
}
filtered_files.push(file_info);
}
info!(
"Filtered to {} files matching criteria",
filtered_files.len()
);
Ok(filtered_files)
}
pub async fn filter_manifest_stream<P: AsRef<Path>>(
manifest_path: P,
dataset_name: Option<&str>,
county: Option<&str>,
quality_version: &QualityControlVersion,
) -> ManifestResult<impl futures::Stream<Item = FileInfo>> {
let files = filter_manifest_files(manifest_path, dataset_name, county, quality_version).await?;
let file_stream = futures::stream::iter(files);
Ok(file_stream)
}
pub async fn fill_queue_from_manifest<P: AsRef<Path>>(
manifest_path: P,
queue: &crate::app::queue::WorkQueue,
dataset_name: Option<&str>,
county: Option<&str>,
quality_version: &QualityControlVersion,
limit: Option<usize>,
) -> ManifestResult<usize> {
let config = ManifestConfig::default();
let manifest_version =
ManifestStreamer::extract_manifest_version_from_path(manifest_path.as_ref());
let mut streamer = ManifestStreamer::with_config_and_version(config, manifest_version);
let mut stream = streamer.stream(manifest_path).await?;
let mut added_count = 0;
let mut processed_count = 0;
info!("Starting pull-based manifest processing");
while let Some(result) = stream.next().await {
if let Some(limit) = limit {
if processed_count >= limit {
info!("Reached limit of {} files", limit);
break;
}
}
let file_info = match result {
Ok(file_info) => file_info,
Err(e) => {
warn!("Skipping invalid entry: {}", e);
continue;
}
};
let dataset_info = &file_info.dataset_info;
if let Some(filter_dataset) = dataset_name {
if dataset_info.dataset_name != filter_dataset {
continue;
}
}
if let Some(filter_county) = county {
if dataset_info.county.as_ref() != Some(&filter_county.to_string()) {
continue;
}
}
if let Some(ref file_qv) = dataset_info.quality_version {
if file_qv != quality_version {
continue;
}
}
processed_count += 1;
while !queue.has_capacity().await {
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
}
if queue.add_work(file_info).await.map_err(|e| {
ManifestError::Io(std::io::Error::new(
std::io::ErrorKind::Other,
e.to_string(),
))
})? {
added_count += 1;
}
if processed_count % 10000 == 0 {
info!(
"Processed {} files, added {} to queue",
processed_count, added_count
);
}
}
info!(
"Manifest processing complete: processed {}, added {}",
processed_count, added_count
);
Ok(added_count)
}
pub async fn get_selection_options<P: AsRef<Path>>(
manifest_path: P,
) -> ManifestResult<(Vec<String>, Vec<String>)> {
use std::collections::HashSet;
let datasets_map = collect_datasets_and_years(manifest_path).await?;
let mut all_datasets: Vec<String> = datasets_map.keys().cloned().collect();
all_datasets.sort();
let mut all_versions: HashSet<String> = HashSet::new();
for summary in datasets_map.values() {
all_versions.extend(summary.versions.iter().cloned());
}
let mut versions_vec: Vec<String> = all_versions.into_iter().collect();
versions_vec.sort();
Ok((all_datasets, versions_vec))
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write;
use tempfile::NamedTempFile;
async fn create_test_manifest(content: &str) -> NamedTempFile {
let mut file = NamedTempFile::new().unwrap();
file.write_all(content.as_bytes()).unwrap();
file.flush().unwrap();
file
}
#[tokio::test]
async fn test_collect_datasets_version_years() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202507/devon/01381_twist/qc-version-1/midas-open_uk-daily-temperature-obs_dv-202507_devon_01381_twist_qcv-1_1993.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202507/devon/01382_twist/qc-version-1/midas-open_uk-daily-temperature-obs_dv-202507_devon_01382_twist_qcv-1_1994.csv"#;
let manifest_file = create_test_manifest(content).await;
let datasets = collect_datasets_and_years(manifest_file.path())
.await
.unwrap();
assert_eq!(datasets.len(), 1);
let temp_dataset = datasets.get("uk-daily-temperature-obs").unwrap();
assert_eq!(temp_dataset.versions, vec!["202507".to_string()]);
assert_eq!(temp_dataset.file_count, 2);
}
#[tokio::test]
async fn test_filter_by_dataset() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202507/devon/01381_twist/qc-version-1/midas-open_uk-daily-temperature-obs_dv-202507_devon_01381_twist_qcv-1_1993.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01382_twist/qc-version-1/midas-open_uk-daily-temperature-obs_dv-202407_devon_01382_twist_qcv-1_1994.csv
ef4718f5cb7b83d0f7bb24a3a598b3a7 ./data/uk-daily-temperature-obs/dataset-version-202507/devon/01383_twist/qc-version-1/midas-open_uk-daily-temperature-obs_dv-202507_devon_01383_twist_qcv-1_1995.csv"#;
let manifest_file = create_test_manifest(content).await;
let files = filter_manifest_files(
manifest_file.path(),
Some("uk-daily-temperature-obs"),
None,
&QualityControlVersion::V1,
)
.await
.unwrap();
assert_eq!(files.len(), 3);
let mut versions: Vec<String> = files
.iter()
.map(|f| f.dataset_info.version.clone())
.collect();
versions.sort();
versions.dedup();
assert!(versions.contains(&"202407".to_string()));
assert!(versions.contains(&"202507".to_string()));
}
#[tokio::test]
async fn test_dataset_analysis_and_filtering_consistency() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202507/devon/01381_twist/qc-version-1/midas-open_uk-daily-temperature-obs_dv-202507_devon_01381_twist_qcv-1_1993.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01382_twist/qc-version-1/midas-open_uk-daily-temperature-obs_dv-202407_devon_01382_twist_qcv-1_1994.csv
ef4718f5cb7b83d0f7bb24a3a598b3a7 ./data/uk-daily-rain-obs/dataset-version-202507/durham/01892_stanhope/qc-version-0/midas-open_uk-daily-rain-obs_dv-202507_durham_01892_stanhope_qcv-0_1995.csv"#;
let manifest_file = create_test_manifest(content).await;
let datasets = collect_datasets_and_years(manifest_file.path())
.await
.unwrap();
let temp_dataset = datasets.get("uk-daily-temperature-obs").unwrap();
assert_eq!(temp_dataset.versions, vec!["202407", "202507"]);
let rain_dataset = datasets.get("uk-daily-rain-obs").unwrap();
assert_eq!(rain_dataset.versions, vec!["202507"]);
let temp_files = filter_manifest_files(
manifest_file.path(),
Some("uk-daily-temperature-obs"),
None,
&QualityControlVersion::V1,
)
.await
.unwrap();
assert_eq!(temp_files.len(), 2, "Should find both temperature files");
for file in &temp_files {
assert_eq!(file.dataset_info.dataset_name, "uk-daily-temperature-obs");
assert!(temp_dataset.versions.contains(&file.dataset_info.version));
}
let rain_files = filter_manifest_files(
manifest_file.path(),
Some("uk-daily-rain-obs"),
None,
&QualityControlVersion::V0, )
.await
.unwrap();
assert_eq!(rain_files.len(), 1, "Should find one rain file");
assert_eq!(rain_files[0].dataset_info.dataset_name, "uk-daily-rain-obs");
assert_eq!(rain_files[0].dataset_info.version, "202507");
}
}