use std::collections::HashSet;
use std::path::{Path, PathBuf};
use std::pin::Pin;
use std::task::{Context, Poll};
use futures::stream::{Stream, StreamExt};
use serde::{Deserialize, Serialize};
use tokio::fs::File;
use tokio::io::{AsyncBufReadExt, BufReader, Lines};
use tracing::{debug, error, info, warn};
use crate::app::hash::Md5Hash;
use crate::app::models::{parse_manifest_line, FileInfo};
use crate::constants::workers;
use crate::errors::{ManifestError, ManifestResult};
#[derive(Debug, Clone, Default)]
pub struct ManifestStats {
pub lines_processed: usize,
pub valid_entries: usize,
pub invalid_lines: usize,
pub duplicate_hashes: usize,
pub empty_lines: usize,
}
impl ManifestStats {
pub fn success_rate(&self) -> f64 {
if self.lines_processed == 0 {
0.0
} else {
(self.valid_entries as f64 / self.lines_processed as f64) * 100.0
}
}
pub fn total_skipped(&self) -> usize {
self.invalid_lines + self.duplicate_hashes + self.empty_lines
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ManifestConfig {
pub destination_root: PathBuf,
pub max_tracked_hashes: usize,
pub allow_duplicates: bool,
pub progress_batch_size: usize,
}
impl Default for ManifestConfig {
fn default() -> Self {
let destination_root = dirs::config_dir()
.map(|dir| dir.join("midas-fetcher").join("cache"))
.unwrap_or_else(|| PathBuf::from("./cache"));
Self {
destination_root,
max_tracked_hashes: 1_000_000, allow_duplicates: false,
progress_batch_size: workers::MANIFEST_BATCH_SIZE,
}
}
}
pub struct ManifestStreamer {
config: ManifestConfig,
seen_hashes: HashSet<Md5Hash>,
stats: ManifestStats,
current_line: usize,
}
impl ManifestStreamer {
pub fn new() -> Self {
Self::with_config(ManifestConfig::default())
}
pub fn with_config(config: ManifestConfig) -> Self {
Self {
config,
seen_hashes: HashSet::new(),
stats: ManifestStats::default(),
current_line: 0,
}
}
pub async fn stream<P: AsRef<Path>>(
&mut self,
manifest_path: P,
) -> ManifestResult<impl Stream<Item = ManifestResult<FileInfo>> + '_> {
let file = File::open(manifest_path.as_ref()).await?;
let reader = BufReader::new(file);
let lines = reader.lines();
info!(
"Starting manifest streaming from: {}",
manifest_path.as_ref().display()
);
Ok(FileInfoStream {
lines,
streamer: self,
})
}
fn process_line(&mut self, line: String) -> Option<ManifestResult<FileInfo>> {
self.current_line += 1;
self.stats.lines_processed += 1;
if line.trim().is_empty() {
self.stats.empty_lines += 1;
return None;
}
let (hash, path) = match parse_manifest_line(&line) {
Ok((hash, path)) => (hash, path),
Err(mut e) => {
if let ManifestError::InvalidFormat { ref mut line, .. } = e {
*line = self.current_line;
}
self.stats.invalid_lines += 1;
warn!(
"Skipping malformed line {}: {}",
self.current_line,
line.trim()
);
return Some(Err(e));
}
};
if !self.config.allow_duplicates {
if self.seen_hashes.contains(&hash) {
self.stats.duplicate_hashes += 1;
debug!("Skipping duplicate hash: {}", hash);
return None;
}
if self.seen_hashes.len() >= self.config.max_tracked_hashes {
warn!(
"Reached maximum tracked hashes limit ({}), clearing cache",
self.config.max_tracked_hashes
);
self.seen_hashes.clear();
}
self.seen_hashes.insert(hash);
}
match FileInfo::new(hash, path, &self.config.destination_root) {
Ok(file_info) => {
self.stats.valid_entries += 1;
if self.stats.valid_entries % self.config.progress_batch_size == 0 {
debug!(
"Processed {} valid entries from {} lines ({}% success rate)",
self.stats.valid_entries,
self.stats.lines_processed,
self.stats.success_rate()
);
}
Some(Ok(file_info))
}
Err(e) => {
self.stats.invalid_lines += 1;
warn!(
"Failed to create FileInfo for line {}: {}",
self.current_line, e
);
Some(Err(e))
}
}
}
pub fn stats(&self) -> &ManifestStats {
&self.stats
}
pub fn reset(&mut self) {
self.seen_hashes.clear();
self.stats = ManifestStats::default();
self.current_line = 0;
}
pub fn estimated_memory_usage(&self) -> usize {
self.seen_hashes.len() * 48 + std::mem::size_of::<Self>()
}
}
impl Default for ManifestStreamer {
fn default() -> Self {
Self::new()
}
}
struct FileInfoStream<'a> {
lines: Lines<BufReader<File>>,
streamer: &'a mut ManifestStreamer,
}
impl Stream for FileInfoStream<'_> {
type Item = ManifestResult<FileInfo>;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let this = self.get_mut();
match Pin::new(&mut this.lines).poll_next_line(cx) {
Poll::Ready(Ok(Some(line))) => {
match this.streamer.process_line(line) {
Some(result) => Poll::Ready(Some(result)),
None => {
cx.waker().wake_by_ref();
Poll::Pending
}
}
}
Poll::Ready(Ok(None)) => {
let stats = this.streamer.stats();
info!(
"Manifest processing completed: {} valid entries from {} lines ({:.1}% success rate)",
stats.valid_entries,
stats.lines_processed,
stats.success_rate()
);
if stats.duplicate_hashes > 0 {
info!("Skipped {} duplicate entries", stats.duplicate_hashes);
}
if stats.invalid_lines > 0 {
warn!("Encountered {} invalid lines", stats.invalid_lines);
}
Poll::Ready(None)
}
Poll::Ready(Err(e)) => {
error!("Error reading manifest file: {}", e);
Poll::Ready(Some(Err(ManifestError::Io(e))))
}
Poll::Pending => Poll::Pending,
}
}
}
pub async fn collect_all_files<P: AsRef<Path>>(
manifest_path: P,
config: ManifestConfig,
) -> ManifestResult<Vec<FileInfo>> {
let mut streamer = ManifestStreamer::with_config(config);
let mut stream = streamer.stream(manifest_path).await?;
let mut files = Vec::new();
while let Some(result) = stream.next().await {
match result {
Ok(file_info) => files.push(file_info),
Err(e) => {
warn!("Skipping invalid entry: {}", e);
}
}
}
Ok(files)
}
pub async fn validate_manifest<P: AsRef<Path>>(
manifest_path: P,
sample_size: usize,
) -> ManifestResult<ManifestStats> {
let mut streamer = ManifestStreamer::new();
let mut stream = streamer.stream(manifest_path).await?;
let mut processed = 0;
while let Some(result) = stream.next().await {
match result {
Ok(_) => {
}
Err(e) => {
debug!("Validation error: {}", e);
}
}
processed += 1;
if sample_size > 0 && processed >= sample_size {
break;
}
}
drop(stream);
Ok(streamer.stats().clone())
}
#[cfg(test)]
mod tests {
use super::*;
use futures::StreamExt;
use std::io::Write;
use tempfile::{NamedTempFile, TempDir};
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_manifest_streaming_basic() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/test.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01382_another/test2.csv"#;
let manifest_file = create_test_manifest(content).await;
let temp_dir = TempDir::new().unwrap();
let config = ManifestConfig {
destination_root: temp_dir.path().to_path_buf(),
..Default::default()
};
let mut streamer = ManifestStreamer::with_config(config);
let files = {
let mut stream = streamer.stream(manifest_file.path()).await.unwrap();
let mut files = Vec::new();
while let Some(result) = stream.next().await {
files.push(result.unwrap());
}
files
};
assert_eq!(files.len(), 2);
assert_eq!(files[0].hash.to_hex(), "50c9d1c465f3cbff652be1509c2e2a4e");
assert_eq!(files[1].hash.to_hex(), "9734faa872681f96b144f60d29d52011");
let stats = streamer.stats();
assert_eq!(stats.valid_entries, 2);
assert_eq!(stats.lines_processed, 2);
assert_eq!(stats.invalid_lines, 0);
}
#[tokio::test]
async fn test_duplicate_detection() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/test1.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01382_twist/test2.csv
50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01383_twist/test3.csv
ef4718f5cb7b83d0f7bb24a3a598b3a7 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01384_twist/test4.csv"#;
let manifest_file = create_test_manifest(content).await;
let temp_dir = TempDir::new().unwrap();
let config = ManifestConfig {
destination_root: temp_dir.path().to_path_buf(),
allow_duplicates: false,
..Default::default()
};
let mut streamer = ManifestStreamer::with_config(config);
let files = {
let mut stream = streamer.stream(manifest_file.path()).await.unwrap();
let mut files = Vec::new();
while let Some(result) = stream.next().await {
files.push(result.unwrap());
}
files
};
assert_eq!(files.len(), 3);
let stats = streamer.stats();
assert_eq!(stats.valid_entries, 3);
assert_eq!(stats.duplicate_hashes, 1);
assert_eq!(stats.lines_processed, 4);
}
#[tokio::test]
async fn test_malformed_lines() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/test1.csv
invalid_hash ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01382_twist/test2.csv
50c9d1c465f3cbff652be1509c2e2a4e data/uk-daily-temperature-obs/dataset-version-202407/devon/01383_twist/test3.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01384_twist/test4.csv"#;
let manifest_file = create_test_manifest(content).await;
let temp_dir = TempDir::new().unwrap();
let config = ManifestConfig {
destination_root: temp_dir.path().to_path_buf(),
..Default::default()
};
let mut streamer = ManifestStreamer::with_config(config);
let (valid_files, errors) = {
let mut stream = streamer.stream(manifest_file.path()).await.unwrap();
let mut valid_files = 0;
let mut errors = 0;
while let Some(result) = stream.next().await {
match result {
Ok(_) => valid_files += 1,
Err(_) => errors += 1,
}
}
(valid_files, errors)
};
assert_eq!(valid_files, 2); assert_eq!(errors, 2);
let stats = streamer.stats();
assert_eq!(stats.valid_entries, 2);
assert_eq!(stats.invalid_lines, 2);
assert_eq!(stats.empty_lines, 1);
assert_eq!(stats.lines_processed, 5);
}
#[tokio::test]
async fn test_collect_all_files() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/test1.csv
invalid_hash ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01382_twist/test2.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01383_twist/test3.csv"#;
let manifest_file = create_test_manifest(content).await;
let temp_dir = TempDir::new().unwrap();
let config = ManifestConfig {
destination_root: temp_dir.path().to_path_buf(),
..Default::default()
};
let files = collect_all_files(manifest_file.path(), config)
.await
.unwrap();
assert_eq!(files.len(), 2);
assert_eq!(files[0].hash.to_hex(), "50c9d1c465f3cbff652be1509c2e2a4e");
assert_eq!(files[1].hash.to_hex(), "9734faa872681f96b144f60d29d52011");
}
#[tokio::test]
async fn test_validate_manifest() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/test1.csv
invalid_hash ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01382_twist/test2.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01383_twist/test3.csv"#;
let manifest_file = create_test_manifest(content).await;
let stats = validate_manifest(manifest_file.path(), 0).await.unwrap();
assert_eq!(stats.lines_processed, 3);
assert_eq!(stats.valid_entries, 2);
assert_eq!(stats.invalid_lines, 1);
assert!(stats.success_rate() > 60.0);
}
#[tokio::test]
async fn test_memory_limits() {
let temp_dir = TempDir::new().unwrap();
let config = ManifestConfig {
destination_root: temp_dir.path().to_path_buf(),
max_tracked_hashes: 2, allow_duplicates: false,
..Default::default()
};
let mut streamer = ManifestStreamer::with_config(config);
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/test1.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01382_twist/test2.csv
ef4718f5cb7b83d0f7bb24a3a598b3a7 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01383_twist/test3.csv
3b71d64ef33dbd6f76497d3ebd4ab976 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01384_twist/test4.csv"#;
let manifest_file = create_test_manifest(content).await;
let files = {
let mut stream = streamer.stream(manifest_file.path()).await.unwrap();
let mut files = Vec::new();
while let Some(result) = stream.next().await {
files.push(result.unwrap());
}
files
};
assert_eq!(files.len(), 4);
assert!(streamer.estimated_memory_usage() < 1024); }
#[tokio::test]
async fn test_real_manifest_sample() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/midas-open_uk-daily-temperature-obs_dv-202407_devon_01381_twist_capability.csv
9734faa872681f96b144f60d29d52011 ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/qc-version-1/midas-open_uk-daily-temperature-obs_dv-202407_devon_01381_twist_qcv-1_1992.csv"#;
let manifest_file = create_test_manifest(content).await;
let temp_dir = TempDir::new().unwrap();
let config = ManifestConfig {
destination_root: temp_dir.path().to_path_buf(),
..Default::default()
};
let files = collect_all_files(manifest_file.path(), config)
.await
.unwrap();
assert_eq!(files.len(), 2);
let capability_file = &files[0];
assert_eq!(
capability_file.hash.to_hex(),
"50c9d1c465f3cbff652be1509c2e2a4e"
);
assert_eq!(
capability_file.dataset_info.dataset_name,
"uk-daily-temperature-obs"
);
assert_eq!(capability_file.dataset_info.version, "202407");
assert_eq!(
capability_file.dataset_info.county,
Some("devon".to_string())
);
assert_eq!(
capability_file.dataset_info.station_id,
Some("01381".to_string())
);
assert_eq!(
capability_file.dataset_info.station_name,
Some("twist".to_string())
);
assert_eq!(
capability_file.dataset_info.file_type,
Some("capability".to_string())
);
let data_file = &files[1];
assert_eq!(data_file.hash.to_hex(), "9734faa872681f96b144f60d29d52011");
assert_eq!(data_file.dataset_info.year, Some("1992".to_string()));
assert_eq!(
data_file.dataset_info.quality_version,
Some(crate::app::models::QualityControlVersion::V1)
);
assert_eq!(data_file.dataset_info.file_type, Some("data".to_string()));
}
#[tokio::test]
async fn test_comprehensive_fileinfo_generation() {
let content = r#"50c9d1c465f3cbff652be1509c2e2a4e ./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/midas-open_uk-daily-temperature-obs_dv-202407_devon_01381_twist_capability.csv"#;
let manifest_file = create_test_manifest(content).await;
let temp_dir = TempDir::new().unwrap();
let config = ManifestConfig {
destination_root: temp_dir.path().to_path_buf(),
..Default::default()
};
let files = collect_all_files(manifest_file.path(), config)
.await
.unwrap();
assert_eq!(files.len(), 1);
let file_info = &files[0];
assert_eq!(file_info.hash.to_hex(), "50c9d1c465f3cbff652be1509c2e2a4e");
assert_eq!(
file_info.relative_path,
"./data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/midas-open_uk-daily-temperature-obs_dv-202407_devon_01381_twist_capability.csv"
);
assert_eq!(
file_info.file_name,
"midas-open_uk-daily-temperature-obs_dv-202407_devon_01381_twist_capability.csv"
);
let dataset = &file_info.dataset_info;
assert_eq!(dataset.dataset_name, "uk-daily-temperature-obs");
assert_eq!(dataset.version, "202407");
assert_eq!(dataset.county, Some("devon".to_string()));
assert_eq!(dataset.station_id, Some("01381".to_string()));
assert_eq!(dataset.station_name, Some("twist".to_string()));
assert_eq!(dataset.quality_version, None); assert_eq!(dataset.year, None); assert_eq!(dataset.file_type, Some("capability".to_string()));
let expected_dest = temp_dir.path().join("data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/midas-open_uk-daily-temperature-obs_dv-202407_devon_01381_twist_capability.csv");
assert_eq!(file_info.destination_path, expected_dest);
assert_eq!(file_info.retry_count, 0);
assert!(file_info.last_attempt.is_none());
assert!(file_info.estimated_size.is_none());
let download_url = file_info.download_url("https://data.ceda.ac.uk");
assert_eq!(
download_url,
"https://data.ceda.ac.uk/badc/ukmo-midas-open/data/uk-daily-temperature-obs/dataset-version-202407/devon/01381_twist/midas-open_uk-daily-temperature-obs_dv-202407_devon_01381_twist_capability.csv"
);
let display_name = dataset.display_name();
assert_eq!(display_name, "uk-daily-temperature-obs-v202407-devon-twist");
println!("✅ Generated complete FileInfo from manifest line:");
println!(" Hash: {}", file_info.hash);
println!(" Dataset: {}", dataset.dataset_name);
println!(" County: {:?}", dataset.county);
println!(
" Station: {} ({})",
dataset.station_name.as_ref().unwrap(),
dataset.station_id.as_ref().unwrap()
);
println!(" File Type: {:?}", dataset.file_type);
println!(" Download URL: {}", download_url);
}
#[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_version_year() {
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,
&crate::app::models::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_year_selection_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,
&crate::app::models::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,
&crate::app::models::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");
println!("✅ Dataset collection and filtering are consistent!");
println!(
" Temperature dataset versions: {:?}",
temp_dataset.versions
);
println!(" Rain dataset versions: {:?}", rain_dataset.versions);
}
}
#[derive(Debug, Clone)]
pub struct DatasetSummary {
pub name: String,
pub versions: Vec<String>,
pub counties: Vec<String>,
pub quality_versions: Vec<crate::app::models::QualityControlVersion>,
pub years: Vec<String>,
pub file_count: usize,
pub example_file: Option<String>,
}
impl DatasetSummary {
pub fn latest_version(&self) -> Option<&String> {
self.versions.iter().max()
}
pub fn latest_year(&self) -> Option<&String> {
self.latest_version()
}
pub fn has_version(&self, version: &str) -> bool {
self.versions.contains(&version.to_string())
}
pub fn has_year(&self, year: &str) -> bool {
self.has_version(year)
}
pub fn has_county(&self, county: &str) -> bool {
self.counties.iter().any(|c| c == county)
}
pub fn has_quality_version(&self, qv: &crate::app::models::QualityControlVersion) -> bool {
self.quality_versions.contains(qv)
}
pub fn earliest_year(&self) -> Option<&String> {
self.years.iter().min()
}
pub fn latest_data_year(&self) -> Option<&String> {
self.years.iter().max()
}
pub fn year_range(&self) -> String {
if self.years.is_empty() {
return "N/A".to_string();
}
let earliest = self.earliest_year().unwrap();
let latest = self.latest_data_year().unwrap();
if earliest == latest {
earliest.clone()
} else {
format!("{}-{}", earliest, latest)
}
}
}
pub async fn collect_datasets_and_years<P: AsRef<Path>>(
manifest_path: P,
) -> ManifestResult<std::collections::HashMap<String, DatasetSummary>> {
use futures::StreamExt;
use std::collections::HashMap;
let config = ManifestConfig::default();
let mut streamer = ManifestStreamer::with_config(config);
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,
});
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 {
crate::app::models::QualityControlVersion::V0 => 0,
crate::app::models::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: &crate::app::models::QualityControlVersion,
) -> ManifestResult<Vec<FileInfo>> {
use futures::StreamExt;
let config = ManifestConfig::default();
let mut streamer = ManifestStreamer::with_config(config);
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: &crate::app::models::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: &crate::app::models::QualityControlVersion,
limit: Option<usize>,
) -> ManifestResult<usize> {
use futures::StreamExt;
let config = ManifestConfig::default();
let mut streamer = ManifestStreamer::with_config(config);
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))
}