diff --git a/Cargo.lock b/Cargo.lock index b21c690..1c946f5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -220,6 +220,15 @@ dependencies = [ "tokio", ] +[[package]] +name = "deranged" +version = "0.5.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2163a0e204a148662b6b6816d4b5d5668a5f2f8df498ccbd5cd0e864e78fecba" +dependencies = [ + "powerfmt", +] + [[package]] name = "displaydoc" version = "0.2.5" @@ -279,6 +288,7 @@ dependencies = [ "rmp-serde", "rmpv", "serde", + "time", "tokio", "tracing", "tracing-subscriber", @@ -652,6 +662,12 @@ dependencies = [ "num-traits", ] +[[package]] +name = "num-conv" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cf97ec579c3c42f953ef76dbf8d55ac91fb219dde70e49aa4a6b7d74e9919050" + [[package]] name = "num-integer" version = "0.1.46" @@ -680,6 +696,15 @@ dependencies = [ "libc", ] +[[package]] +name = "num_threads" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c7398b9c8b70908f6371f47ed36737907c87c52af34c268fed0bf0ceb92ead9" +dependencies = [ + "libc", +] + [[package]] name = "once_cell" version = "1.21.3" @@ -719,6 +744,12 @@ dependencies = [ "zerovec", ] +[[package]] +name = "powerfmt" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" + [[package]] name = "proc-macro2" version = "1.0.106" @@ -1061,6 +1092,39 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "time" +version = "0.3.47" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "743bd48c283afc0388f9b8827b976905fb217ad9e647fae3a379a9283c4def2c" +dependencies = [ + "deranged", + "itoa", + "libc", + "num-conv", + "num_threads", + "powerfmt", + "serde_core", + "time-core", + "time-macros", +] + +[[package]] +name = "time-core" +version = "0.1.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7694e1cfe791f8d31026952abf09c69ca6f6fa4e1a1229e18988f06a04a12dca" + +[[package]] +name = "time-macros" +version = "0.2.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e70e4c5a0e0a8a4823ad65dfe1a6930e4f4d756dcd9dd7939022b5e8c501215" +dependencies = [ + "num-conv", + "time-core", +] + [[package]] name = "tinystr" version = "0.8.2" @@ -1166,6 +1230,7 @@ dependencies = [ "sharded-slab", "smallvec", "thread_local", + "time", "tracing", "tracing-core", "tracing-log", diff --git a/crates/fluxqueue-worker/Cargo.toml b/crates/fluxqueue-worker/Cargo.toml index cee1029..8c9f164 100644 --- a/crates/fluxqueue-worker/Cargo.toml +++ b/crates/fluxqueue-worker/Cargo.toml @@ -27,7 +27,8 @@ rmp-serde = "1.3.1" rmpv = { version = "1.3.1", features = ["with-serde"] } serde = { version = "1.0.228", features = ["derive"] } tracing = "0.1.44" -tracing-subscriber = { version = "0.3.22", features = ["ansi", "env-filter"] } +tracing-subscriber = { version = "0.3.22", features = ["env-filter", "fmt", "time"] } +time = { version = "0.3.47", features = ["local-offset", "macros"] } clap = { version = "4.5.56", features = ["derive", "env"] } anyhow = "1.0.100" pythonize = "0.27.0" diff --git a/crates/fluxqueue-worker/src/logger.rs b/crates/fluxqueue-worker/src/logger.rs index 7dbf8ec..fff49c4 100644 --- a/crates/fluxqueue-worker/src/logger.rs +++ b/crates/fluxqueue-worker/src/logger.rs @@ -41,7 +41,7 @@ pub fn initial_logs( info!("Redis: {}", redis_url); info!("Tasks module: {}", tasks_module_path); info!("Tasks found: {:?}", tasks); - info!("{}", "-".repeat(65)); + info!("Starting up the executors..."); } #[cfg(test)] diff --git a/crates/fluxqueue-worker/src/main.rs b/crates/fluxqueue-worker/src/main.rs index c7a8020..407d16f 100644 --- a/crates/fluxqueue-worker/src/main.rs +++ b/crates/fluxqueue-worker/src/main.rs @@ -1,6 +1,8 @@ use anyhow::Result; use clap::Parser; +use time::{UtcOffset, macros::format_description}; use tracing::{error, info, warn}; +use tracing_subscriber::fmt::time::OffsetTime; #[derive(Parser, Debug)] #[command(version, about, long_about = None)] @@ -49,9 +51,15 @@ struct Cli { #[tokio::main] async fn main() -> Result<()> { + let offset = UtcOffset::current_local_offset().unwrap(); + let timer = OffsetTime::new( + offset, + format_description!("[year]:[month]:[day]:[hour]:[minute]:[second].[subsecond digits:3]"), + ); tracing_subscriber::fmt() .with_env_filter("fluxqueue_worker=debug") .with_target(false) + .with_timer(timer) .init(); let cli = Cli::parse(); diff --git a/crates/fluxqueue-worker/src/worker.rs b/crates/fluxqueue-worker/src/worker.rs index b07622a..8e2ba8c 100644 --- a/crates/fluxqueue-worker/src/worker.rs +++ b/crates/fluxqueue-worker/src/worker.rs @@ -4,6 +4,7 @@ use pyo3::{Bound, Py, PyAny, Python}; use std::ffi::CString; use std::path::{Path, PathBuf}; use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::{Duration, Instant}; use tokio::sync::watch; use tokio::task::JoinSet; @@ -48,14 +49,18 @@ pub async fn run_worker( let queue_name = Arc::from(queue_name.to_string()); let executor_ids = generate_executor_ids(concurrency); + let atomic_concurrency = Arc::new(AtomicUsize::new(concurrency)); + let started_executors_count = Arc::new(AtomicUsize::new(0)); let mut executors = JoinSet::new(); let executor_futures: Vec<_> = (0..concurrency) .map(|i| { + let shutdown = shutdown.clone(); + let concurrency = Arc::clone(&atomic_concurrency); + let counter = Arc::clone(&started_executors_count); let redis_client = Arc::clone(&redis_client); let queue_name = Arc::clone(&queue_name); let executor_id = Arc::clone(&executor_ids[i]); - let shutdown = shutdown.clone(); let task_registry = Arc::clone(&task_registry); async move { @@ -67,40 +72,46 @@ pub async fn run_worker( redis_client.set_executor_heartbeat(&executor_id).await?; - Ok::<_, anyhow::Error>(( - shutdown, + let ready_check = ReadyCheck { + concurrency, + counter, + }; + + let executor_context = ExecutorContext { queue_name, executor_id, redis_client, task_registry, python_dispatcher, - )) + }; + + Ok::<_, anyhow::Error>((shutdown, ready_check, executor_context)) } }) .collect(); let results = futures::future::join_all(executor_futures).await; for result in results { - let (shutdown, queue_name, executor_id, redis_client, task_registry, python_dispatcher) = - result?; - executors.spawn(executor_loop( - shutdown, - queue_name, - executor_id, - redis_client, - task_registry, - python_dispatcher, - )); + let (shutdown, ready_check, executor_context) = result?; + executors.spawn(executor_loop(shutdown, ready_check, executor_context)); } + let janitor_shutdown = shutdown.clone(); + let concurrency = Arc::clone(&atomic_concurrency); + let counter = Arc::clone(&started_executors_count); let janitor_queue_name = Arc::clone(&queue_name); let janitor_redis = Arc::clone(&redis_client); let janitor_executor_ids = Arc::clone(&executor_ids); - let janitor_shutdown = shutdown.clone(); let save_dead_tasks = Arc::new(save_dead_tasks); + let ready_check = ReadyCheck { + concurrency, + counter, + }; + executors.spawn(janitor_loop( janitor_shutdown, + ready_check, janitor_queue_name, janitor_executor_ids, save_dead_tasks, @@ -124,27 +135,37 @@ pub async fn run_worker( Ok(()) } -async fn executor_loop( - mut shutdown: watch::Receiver, +struct ExecutorContext { queue_name: Arc, executor_id: Arc, redis_client: Arc, task_registry: Arc, python_dispatcher: Arc, +} + +struct ReadyCheck { + concurrency: Arc, + counter: Arc, +} + +async fn executor_loop( + mut shutdown: watch::Receiver, + ready_check: ReadyCheck, + ctx: ExecutorContext, ) -> Result<()> { - let logger = Logger::new(format!("EXECUTOR {}", &executor_id)); + let logger = Logger::new(format!("EXECUTOR {}", &ctx.executor_id)); - logger.info(format_args!("Started!")); + ready_check.counter.fetch_add(1, Ordering::SeqCst); + check_worker_is_ready(ready_check); loop { tokio::select! { _ = shutdown.changed() => { - logger.info(format_args!("Shutting down...")); return Ok(()) } - res = redis_client - .mark_task_as_processing(&queue_name, &executor_id) + res = ctx.redis_client + .mark_task_as_processing(&ctx.queue_name, &ctx.executor_id) => { match res { Ok(Some(raw_data)) => { @@ -157,10 +178,10 @@ async fn executor_loop( raw_data.len() )); - let Some(task_function) = task_registry.get(&task.name) else { + let Some(task_function) = ctx.task_registry.get(&task.name) else { logger.warn(format_args!("Task '{}' not found in registry. Skipping", &task.name)); - if let Err(e) = redis_client - .remove_from_processing(&queue_name, &executor_id, &raw_data) + if let Err(e) = ctx.redis_client + .remove_from_processing(&ctx.queue_name, &ctx.executor_id, &raw_data) .await { logger.error(format_args!("Failed to remove the task: {}", e)); } @@ -168,12 +189,12 @@ async fn executor_loop( }; let duration_start = Instant::now(); - let task_result = run_task(python_dispatcher.clone(), &task, task_function).await; + let task_result = run_task(ctx.python_dispatcher.clone(), &task, task_function).await; match task_result { Ok(_) => { - if let Err(e) = redis_client - .remove_from_processing(&queue_name, &executor_id, &raw_data) + if let Err(e) = ctx.redis_client + .remove_from_processing(&ctx.queue_name, &ctx.executor_id, &raw_data) .await { logger.error(format_args!("Failed to remove the task after successful run: {}", e)); } @@ -192,8 +213,8 @@ async fn executor_loop( duration_end.as_millis(), e )); - if let Err(err) = redis_client - .mark_as_failed(&queue_name, &executor_id, &raw_data) + if let Err(err) = ctx.redis_client + .mark_as_failed(&ctx.queue_name, &ctx.executor_id, &raw_data) .await { logger.error(format_args!("Failed to mark the task '{}' as failed: {}", &task.name, err)); } @@ -202,7 +223,7 @@ async fn executor_loop( } Ok(None) => continue, Err(e) => { - logger.error(format_args!("Worker {} redis error: {}", &executor_id, e)); + logger.error(format_args!("Worker {} redis error: {}", &ctx.executor_id, e)); tokio::time::sleep(Duration::from_millis(500)).await; } } @@ -213,6 +234,7 @@ async fn executor_loop( async fn janitor_loop( mut shutdown: watch::Receiver, + ready_check: ReadyCheck, queue_name: Arc, executor_ids: Arc>>, save_dead_tasks: Arc, @@ -220,12 +242,12 @@ async fn janitor_loop( ) -> Result<()> { let logger = Logger::new("JANITOR"); - logger.info(format_args!("Started!")); + ready_check.counter.fetch_add(1, Ordering::SeqCst); + check_worker_is_ready(ready_check); loop { tokio::select! { _ = shutdown.changed() => { - logger.info(format_args!("Shutting down...")); return Ok(()) } @@ -401,6 +423,17 @@ fn generate_executor_ids(num_executors: usize) -> Arc>> { Arc::new(ids) } +fn check_worker_is_ready(ready_check: ReadyCheck) { + let current = ready_check.counter.load(Ordering::SeqCst); + if current == ready_check.concurrency.load(Ordering::SeqCst) + 1 { + tracing::info!( + "Worker ready ({:?} executors, 1 janitor)", + ready_check.concurrency + ); + tracing::info!("{}", "-".repeat(65)); + } +} + #[cfg(test)] mod tests { use super::*;