From 473c7c71664dd85722fd408945ed62175f0494ba Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 21 Feb 2026 15:49:40 +0400 Subject: [PATCH 1/3] Make startup logs more clear --- Cargo.lock | 65 +++++++++++++++++++++++++++ crates/fluxqueue-worker/Cargo.toml | 3 +- crates/fluxqueue-worker/src/logger.rs | 2 +- crates/fluxqueue-worker/src/main.rs | 8 ++++ crates/fluxqueue-worker/src/worker.rs | 49 ++++++++++++++++---- 5 files changed, 117 insertions(+), 10 deletions(-) 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..2db8f0e 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 { @@ -69,6 +74,8 @@ pub async fn run_worker( Ok::<_, anyhow::Error>(( shutdown, + concurrency, + counter, queue_name, executor_id, redis_client, @@ -81,10 +88,20 @@ pub async fn run_worker( 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?; + let ( + shutdown, + concurrency, + counter, + queue_name, + executor_id, + redis_client, + task_registry, + python_dispatcher, + ) = result?; executors.spawn(executor_loop( shutdown, + concurrency, + counter, queue_name, executor_id, redis_client, @@ -93,14 +110,18 @@ pub async fn run_worker( )); } + 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); executors.spawn(janitor_loop( janitor_shutdown, + concurrency, + counter, janitor_queue_name, janitor_executor_ids, save_dead_tasks, @@ -126,6 +147,8 @@ pub async fn run_worker( async fn executor_loop( mut shutdown: watch::Receiver, + concurrency: Arc, + c: Arc, queue_name: Arc, executor_id: Arc, redis_client: Arc, @@ -134,12 +157,12 @@ async fn executor_loop( ) -> Result<()> { let logger = Logger::new(format!("EXECUTOR {}", &executor_id)); - logger.info(format_args!("Started!")); + c.fetch_add(1, Ordering::SeqCst); + check_worker_is_ready(concurrency, c); loop { tokio::select! { _ = shutdown.changed() => { - logger.info(format_args!("Shutting down...")); return Ok(()) } @@ -213,6 +236,8 @@ async fn executor_loop( async fn janitor_loop( mut shutdown: watch::Receiver, + concurrency: Arc, + c: Arc, queue_name: Arc, executor_ids: Arc>>, save_dead_tasks: Arc, @@ -220,12 +245,12 @@ async fn janitor_loop( ) -> Result<()> { let logger = Logger::new("JANITOR"); - logger.info(format_args!("Started!")); + c.fetch_add(1, Ordering::SeqCst); + check_worker_is_ready(concurrency, c); loop { tokio::select! { _ = shutdown.changed() => { - logger.info(format_args!("Shutting down...")); return Ok(()) } @@ -401,6 +426,14 @@ fn generate_executor_ids(num_executors: usize) -> Arc>> { Arc::new(ids) } +fn check_worker_is_ready(concurrency: Arc, c: Arc) { + let current = c.load(Ordering::SeqCst); + if current == concurrency.load(Ordering::SeqCst) + 1 { + tracing::info!("Worker ready ({:?} executors, 1 janitor)", concurrency); + tracing::info!("{}", "-".repeat(65)); + } +} + #[cfg(test)] mod tests { use super::*; From 856edcdf61425865bd24ea96b8bd0611401a7d23 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 21 Feb 2026 15:59:43 +0400 Subject: [PATCH 2/3] Fix clippy error by moving args to a struct --- crates/fluxqueue-worker/src/worker.rs | 108 +++++++++++++------------- 1 file changed, 54 insertions(+), 54 deletions(-) diff --git a/crates/fluxqueue-worker/src/worker.rs b/crates/fluxqueue-worker/src/worker.rs index 2db8f0e..c281854 100644 --- a/crates/fluxqueue-worker/src/worker.rs +++ b/crates/fluxqueue-worker/src/worker.rs @@ -14,6 +14,19 @@ use crate::redis_client::RedisClient; use crate::task::{PythonDispatcher, TaskRegistry}; use fluxqueue_common::{Task, deserialize_raw_task_data}; +struct ExecutorContext { + queue_name: Arc, + executor_id: Arc, + redis_client: Arc, + task_registry: Arc, + python_dispatcher: Arc, +} + +struct ReadyCheck { + concurrency: Arc, + counter: Arc, +} + pub async fn run_worker( mut shutdown: watch::Receiver, concurrency: usize, @@ -72,42 +85,28 @@ 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, - concurrency, - counter, - queue_name, - executor_id, - redis_client, - task_registry, - python_dispatcher, - ) = result?; - executors.spawn(executor_loop( - shutdown, - concurrency, - counter, - 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(); @@ -118,10 +117,14 @@ pub async fn run_worker( let janitor_executor_ids = Arc::clone(&executor_ids); let save_dead_tasks = Arc::new(save_dead_tasks); - executors.spawn(janitor_loop( - janitor_shutdown, + let ready_check = ReadyCheck { concurrency, counter, + }; + + executors.spawn(janitor_loop( + janitor_shutdown, + ready_check, janitor_queue_name, janitor_executor_ids, save_dead_tasks, @@ -147,18 +150,13 @@ pub async fn run_worker( async fn executor_loop( mut shutdown: watch::Receiver, - concurrency: Arc, - c: Arc, - queue_name: Arc, - executor_id: Arc, - redis_client: Arc, - task_registry: Arc, - python_dispatcher: Arc, + ready_check: ReadyCheck, + ctx: ExecutorContext, ) -> Result<()> { - let logger = Logger::new(format!("EXECUTOR {}", &executor_id)); + let logger = Logger::new(format!("EXECUTOR {}", &ctx.executor_id)); - c.fetch_add(1, Ordering::SeqCst); - check_worker_is_ready(concurrency, c); + ready_check.counter.fetch_add(1, Ordering::SeqCst); + check_worker_is_ready(ready_check); loop { tokio::select! { @@ -166,8 +164,8 @@ async fn executor_loop( 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)) => { @@ -180,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)); } @@ -191,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)); } @@ -215,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)); } @@ -225,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; } } @@ -236,8 +234,7 @@ async fn executor_loop( async fn janitor_loop( mut shutdown: watch::Receiver, - concurrency: Arc, - c: Arc, + ready_check: ReadyCheck, queue_name: Arc, executor_ids: Arc>>, save_dead_tasks: Arc, @@ -245,8 +242,8 @@ async fn janitor_loop( ) -> Result<()> { let logger = Logger::new("JANITOR"); - c.fetch_add(1, Ordering::SeqCst); - check_worker_is_ready(concurrency, c); + ready_check.counter.fetch_add(1, Ordering::SeqCst); + check_worker_is_ready(ready_check); loop { tokio::select! { @@ -426,10 +423,13 @@ fn generate_executor_ids(num_executors: usize) -> Arc>> { Arc::new(ids) } -fn check_worker_is_ready(concurrency: Arc, c: Arc) { - let current = c.load(Ordering::SeqCst); - if current == concurrency.load(Ordering::SeqCst) + 1 { - tracing::info!("Worker ready ({:?} executors, 1 janitor)", concurrency); +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)); } } From 689340c8bc96013721dca6ef30a35d57cdb336e4 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 21 Feb 2026 16:01:01 +0400 Subject: [PATCH 3/3] Move structs below the run_worker func --- crates/fluxqueue-worker/src/worker.rs | 26 +++++++++++++------------- 1 file changed, 13 insertions(+), 13 deletions(-) diff --git a/crates/fluxqueue-worker/src/worker.rs b/crates/fluxqueue-worker/src/worker.rs index c281854..8e2ba8c 100644 --- a/crates/fluxqueue-worker/src/worker.rs +++ b/crates/fluxqueue-worker/src/worker.rs @@ -14,19 +14,6 @@ use crate::redis_client::RedisClient; use crate::task::{PythonDispatcher, TaskRegistry}; use fluxqueue_common::{Task, deserialize_raw_task_data}; -struct ExecutorContext { - queue_name: Arc, - executor_id: Arc, - redis_client: Arc, - task_registry: Arc, - python_dispatcher: Arc, -} - -struct ReadyCheck { - concurrency: Arc, - counter: Arc, -} - pub async fn run_worker( mut shutdown: watch::Receiver, concurrency: usize, @@ -148,6 +135,19 @@ pub async fn run_worker( Ok(()) } +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,