pub mod background_tasks;
pub mod completion;
pub mod config;
pub mod progress;
pub mod signals;
pub mod stats;
#[cfg(test)]
pub mod tests;
use std::sync::Arc;
use std::time::Instant;
use tokio::sync::{mpsc, RwLock};
use tracing::{error, info};
use crate::app::cache::CacheManager;
use crate::app::client::CedaClient;
use crate::app::queue::WorkQueue;
use crate::app::worker::WorkerPool;
use crate::errors::DownloadResult;
pub use background_tasks::BackgroundTaskManager;
pub use completion::{CompletionDetector, CompletionStatus};
pub use config::CoordinatorConfig;
pub use progress::{ProgressAggregator, ProgressMonitor, RateCalculator};
pub use signals::{create_shutdown_channel, wait_for_shutdown_signal, SignalHandler};
pub use stats::{DownloadStats, SessionResult};
pub struct Coordinator {
config: CoordinatorConfig,
queue: Arc<WorkQueue>,
cache: Arc<CacheManager>,
client: Arc<CedaClient>,
stats: Arc<RwLock<DownloadStats>>,
}
impl Coordinator {
pub fn new(
config: CoordinatorConfig,
queue: Arc<WorkQueue>,
cache: Arc<CacheManager>,
client: Arc<CedaClient>,
) -> Self {
let stats = Arc::new(RwLock::new(DownloadStats::default()));
Self {
config,
queue,
cache,
client,
stats,
}
}
pub fn new_with_expected_files(
config: CoordinatorConfig,
queue: Arc<WorkQueue>,
cache: Arc<CacheManager>,
client: Arc<CedaClient>,
expected_files: usize,
) -> Self {
let stats = DownloadStats::new_with_expected_files(expected_files);
let stats = Arc::new(RwLock::new(stats));
Self {
config,
queue,
cache,
client,
stats,
}
}
pub async fn run_downloads(&mut self) -> DownloadResult<SessionResult> {
let session_start = Instant::now();
info!(
"Starting download coordination with {} workers",
self.config.worker_count
);
if let Err(e) = self.config.validate() {
error!("Invalid coordinator configuration: {}", e);
return Ok(SessionResult::failed(
self.stats.read().await.clone(),
session_start.elapsed(),
vec![format!("Configuration validation failed: {}", e)],
));
}
let (shutdown_tx, _) = create_shutdown_channel();
let signal_handler = SignalHandler::new(shutdown_tx.clone());
let mut shutdown_signal = signal_handler.setup();
{
let mut stats = self.stats.write().await;
stats.session_start = chrono::Utc::now();
if stats.total_files == 0 {
stats.total_files = self.queue.stats().await.total_added as usize;
}
stats.active_workers = self.config.worker_count;
}
let mut worker_pool = match self.create_worker_pool().await {
Ok(pool) => pool,
Err(e) => {
error!("Failed to create worker pool: {}", e);
return Ok(SessionResult::failed(
self.stats.read().await.clone(),
session_start.elapsed(),
vec![format!("Worker pool creation failed: {}", e)],
));
}
};
let (progress_tx, progress_rx) = mpsc::channel(self.config.progress_batch_size);
let progress_monitor = ProgressMonitor::new(
self.stats.clone(),
self.queue.clone(),
self.config.progress_update_interval,
self.config.verbose_logging,
);
let progress_handle =
progress_monitor.start_monitoring(progress_rx, shutdown_tx.subscribe());
if let Err(e) = worker_pool.start(progress_tx).await {
error!("Failed to start worker pool: {}", e);
return Ok(SessionResult::failed(
self.stats.read().await.clone(),
session_start.elapsed(),
vec![format!("Worker pool start failed: {}", e)],
));
}
let mut background_tasks = BackgroundTaskManager::new();
background_tasks.start_cleanup_task(
self.cache.clone(),
self.queue.clone(),
shutdown_tx.subscribe(),
);
background_tasks.start_periodic_logging_task(self.queue.clone(), shutdown_tx.subscribe());
background_tasks.start_timeout_monitoring_task(self.queue.clone(), shutdown_tx.subscribe());
let expected_files = {
let stats = self.stats.read().await;
stats.total_files as u64
};
let completion_detector = CompletionDetector::new(self.queue.clone(), expected_files);
let completion_result = tokio::select! {
_ = &mut shutdown_signal => {
info!("Shutdown signal received, initiating graceful shutdown");
self.handle_shutdown().await
}
_ = completion_detector.wait_for_completion() => {
info!("All downloads completed naturally");
self.handle_completion().await
}
};
let _ = shutdown_tx.send(());
background_tasks.shutdown_all().await;
let _ = progress_handle.await;
let shutdown_errors = match tokio::time::timeout(
self.config.shutdown_timeout,
worker_pool.shutdown(),
)
.await
{
Ok(Ok(())) => Vec::new(),
Ok(Err(e)) => {
error!("Worker pool shutdown error: {}", e);
vec![format!("Worker pool shutdown error: {}", e)]
}
Err(_) => {
error!(
"Worker pool shutdown timed out after {:?}",
self.config.shutdown_timeout
);
vec!["Worker pool shutdown timed out".to_string()]
}
};
let final_stats = {
let mut stats = self.stats.write().await;
stats.session_duration = session_start.elapsed();
stats.active_workers = 0;
stats.clone()
};
info!(
"Download session completed in {:?}",
session_start.elapsed()
);
info!(
"Final stats: {} completed, {} failed, {} total",
final_stats.files_completed, final_stats.files_failed, final_stats.total_files
);
Ok(if completion_result.is_ok() && shutdown_errors.is_empty() {
SessionResult::success(final_stats, session_start.elapsed())
} else {
SessionResult::failed(final_stats, session_start.elapsed(), shutdown_errors)
})
}
pub async fn get_stats(&self) -> DownloadStats {
self.stats.read().await.clone()
}
pub async fn shutdown(&self) -> DownloadResult<()> {
info!("Shutdown requested via API");
Ok(())
}
async fn create_worker_pool(&self) -> DownloadResult<WorkerPool> {
info!(
"Creating worker pool with {} workers",
self.config.worker_count
);
let mut worker_config = self.config.worker_config.clone();
worker_config.worker_count = self.config.worker_count;
let pool = WorkerPool::new(
worker_config,
self.queue.clone(),
self.cache.clone(),
self.client.clone(),
);
Ok(pool)
}
async fn handle_shutdown(&mut self) -> DownloadResult<()> {
info!("Initiating graceful shutdown...");
{
let mut stats = self.stats.write().await;
stats.update_duration();
}
Ok(())
}
async fn handle_completion(&mut self) -> DownloadResult<()> {
info!("All downloads completed successfully");
let queue_stats = self.queue.stats().await;
{
let mut stats = self.stats.write().await;
stats.files_completed = queue_stats.completed_count as usize;
stats.files_failed = queue_stats.failed_count as usize;
stats.files_in_progress = 0;
stats.active_workers = 0;
stats.update_duration();
}
Ok(())
}
}