diff --git a/src/db.rs b/src/db.rs index c49547f..5b399e0 100644 --- a/src/db.rs +++ b/src/db.rs @@ -303,7 +303,7 @@ impl Database { error ) VALUES ( - TO_TIMESTAMP($1::DOUBLE PRECISION / 1000.0), + TO_TIMESTAMP(($1::BIGINT)::DOUBLE PRECISION / 1000.0), $1, $2, $3, diff --git a/src/redis_rtsp.rs b/src/redis_rtsp.rs index c058496..f773d23 100644 --- a/src/redis_rtsp.rs +++ b/src/redis_rtsp.rs @@ -43,42 +43,60 @@ pub async fn run_redis_coordinator( let result_queue = result_queue_name(&config.redis_queue_name); let total_key = total_ips_key(&config.redis_queue_name); let completed_key = completed_ips_key(&config.redis_queue_name); - clear_queue( + let should_resume = should_resume_existing_run( &config.redis_url, &config.redis_queue_name, &result_queue, - &done_key, &total_key, &completed_key, + ips.len(), ) .await?; - set_redis_key(&config.redis_url, &total_key, ips.len()).await?; - set_redis_key(&config.redis_url, &completed_key, 0).await?; - enqueue_ips( - &config.redis_url, - &config.redis_queue_name, - &ips, - config.redis_queue_block_size.max(1), - config.redis_worker_count.max(1), - ) - .await?; - mark_done(&config.redis_url, &done_key).await?; - if ui.is_enabled() { - ui.print_banner(); - ui.render(&ScreenState { - current_ip: format!("fila Redis: {}", config.redis_queue_name), - tested: ips.len(), - total: ips.len(), - }); - ui.finish(); - } else { + if should_resume { println!( - "Fila Redis alimentada com {} IPs em {}. Workers esperados: {}", - ips.len(), + "Retomando fila Redis existente em {}. Total esperado: {} IPs", config.redis_queue_name, - config.redis_worker_count.max(1) + ips.len() ); + } else { + clear_queue( + &config.redis_url, + &config.redis_queue_name, + &result_queue, + &done_key, + &total_key, + &completed_key, + ) + .await?; + set_redis_key(&config.redis_url, &total_key, ips.len()).await?; + set_redis_key(&config.redis_url, &completed_key, 0).await?; + enqueue_ips( + &config.redis_url, + &config.redis_queue_name, + &ips, + config.redis_queue_block_size.max(1), + config.redis_worker_count.max(1), + ) + .await?; + mark_done(&config.redis_url, &done_key).await?; + + if ui.is_enabled() { + ui.print_banner(); + ui.render(&ScreenState { + current_ip: format!("fila Redis: {}", config.redis_queue_name), + tested: ips.len(), + total: ips.len(), + }); + ui.finish(); + } else { + println!( + "Fila Redis alimentada com {} IPs em {}. Workers esperados: {}", + ips.len(), + config.redis_queue_name, + config.redis_worker_count.max(1) + ); + } } consume_result_queue( @@ -708,6 +726,28 @@ async fn read_global_progress( }) } +async fn should_resume_existing_run( + redis_url: &str, + queue_name: &str, + result_queue: &str, + total_key: &str, + completed_key: &str, + expected_total: usize, +) -> Result { + let existing_total = redis_usize(redis_url, total_key).await?.unwrap_or_default(); + if existing_total != expected_total { + return Ok(false); + } + + let completed = redis_usize(redis_url, completed_key) + .await? + .unwrap_or_default(); + let queued_work = redis_list_len(redis_url, queue_name).await?; + let queued_results = redis_list_len(redis_url, result_queue).await?; + + Ok(queued_work > 0 || queued_results > 0 || (completed > 0 && completed < expected_total)) +} + fn total_ips_key(queue_name: &str) -> String { format!("{queue_name}:total_ips") } @@ -730,8 +770,8 @@ async fn push_fetch_result( Ok(()) } -async fn pop_fetch_result(redis_url: &str, result_queue: &str) -> Result> { - let output = redis_cli(redis_url, &["--raw", "LPOP", result_queue]).await?; +async fn peek_fetch_result(redis_url: &str, result_queue: &str) -> Result> { + let output = redis_cli(redis_url, &["--raw", "LINDEX", result_queue, "0"]).await?; let value = output.trim_end_matches('\n'); if value.is_empty() { return Ok(None); @@ -739,6 +779,11 @@ async fn pop_fetch_result(redis_url: &str, result_queue: &str) -> Result Result<()> { + let _ = redis_cli(redis_url, &["LPOP", result_queue]).await?; + Ok(()) +} + async fn consume_result_queue( redis_url: &str, queue_name: &str, @@ -750,7 +795,7 @@ async fn consume_result_queue( let mut inserted = 0usize; loop { let mut drained_any = false; - while let Some(result) = pop_fetch_result(redis_url, result_queue).await? { + while let Some(result) = peek_fetch_result(redis_url, result_queue).await? { database .insert_rtsp_fetch_result( result.checked_at, @@ -763,6 +808,7 @@ async fn consume_result_queue( result.error.as_deref(), ) .await?; + ack_fetch_result(redis_url, result_queue).await?; inserted += 1; drained_any = true; }