use anyhow::Result;
use crate::config::RuntimeConfig;
use crate::results::Phase2Result;
use crate::setup::OutputSink;
use crate::utils::{BATCHED_NO_VALID_FILES_DETAIL_CAP, no_valid_files_error, zarr_paths};
use super::chunking::{PathBatchMode, calculate_adaptive_chunking, determine_batch_size};
use super::phases::{mining::phase2_mining, scanning::phase1_scan};
pub fn run_pipeline(
input_paths: &[String],
output_dir: Option<&str>,
config: &RuntimeConfig,
output_sink: &OutputSink,
) -> Result<(Vec<(String, String)>, Phase2Result)> {
let input_paths: Vec<String> = zarr_paths::filter_zarr_input_paths(input_paths.iter().cloned());
let input_slice = input_paths.as_slice();
let skip_file_write = output_dir.is_none();
match determine_batch_size(input_slice) {
PathBatchMode::Single => run_pipeline_single(
input_slice,
output_dir,
config,
output_sink,
skip_file_write,
),
PathBatchMode::Batched { batch_size } => run_pipeline_batched(
input_slice,
output_dir,
config,
output_sink,
skip_file_write,
batch_size,
),
}
}
fn run_pipeline_single(
input_paths: &[String],
output_dir: Option<&str>,
config: &RuntimeConfig,
output_sink: &OutputSink,
skip_file_write: bool,
) -> Result<(Vec<(String, String)>, Phase2Result)> {
let phase1 = phase1_scan(input_paths, output_dir, config);
let tasks = phase1.tasks;
if tasks.is_empty() {
return Err(no_valid_files_error(&phase1.failed, None));
}
let adaptive = calculate_adaptive_chunking(&tasks, config.max_workers, config);
let phase2 = phase2_mining(&tasks, config, &adaptive, skip_file_write, output_sink);
Ok((phase1.failed, phase2))
}
fn run_pipeline_batched(
input_paths: &[String],
output_dir: Option<&str>,
config: &RuntimeConfig,
output_sink: &OutputSink,
skip_file_write: bool,
batch_size: usize,
) -> Result<(Vec<(String, String)>, Phase2Result)> {
let mut all_phase1_failed = Vec::new();
let mut all_phase2_failed = Vec::new();
let mut all_outputs = Vec::new();
let mut any_batch_had_tasks = false;
for chunk in input_paths.chunks(batch_size) {
let phase1 = phase1_scan(chunk, output_dir, config);
if phase1.tasks.is_empty() {
all_phase1_failed.extend(phase1.failed);
continue;
}
any_batch_had_tasks = true;
let adaptive = calculate_adaptive_chunking(&phase1.tasks, config.max_workers, config);
let phase2 = phase2_mining(
&phase1.tasks,
config,
&adaptive,
skip_file_write,
output_sink,
);
all_phase1_failed.extend(phase1.failed);
all_phase2_failed.extend(phase2.failed);
if matches!(output_sink, OutputSink::Collect) {
all_outputs.extend(phase2.outputs);
}
}
if !any_batch_had_tasks {
return Err(no_valid_files_error(
&all_phase1_failed,
Some(BATCHED_NO_VALID_FILES_DETAIL_CAP),
));
}
Ok((
all_phase1_failed,
Phase2Result {
outputs: all_outputs,
failed: all_phase2_failed,
},
))
}