From 0237024d1cc6127162a3d743ed2ed1524d6d6b81 Mon Sep 17 00:00:00 2001 From: Cory Hawklvelt Date: Mon, 27 Oct 2025 08:42:45 +1030 Subject: [PATCH] feat(task-assigner): add debug flag for task assignment logging Introduce `ENABLE_TASK_ASSIGNMENT_DEBUG` configuration to control verbose logging in task assignment, Redis subscription, and worker management flows. This reduces log noise by default while allowing detailed debugging when enabled. --- websocket_server/config.py | 2 ++ websocket_server/redis_utils.py | 14 ++++++++++---- websocket_server/task_assigner.py | 25 +++++++++++++++++-------- websocket_server/worker_manager.py | 11 +++++++---- 4 files changed, 36 insertions(+), 16 deletions(-) diff --git a/websocket_server/config.py b/websocket_server/config.py index 390e82a..678498f 100644 --- a/websocket_server/config.py +++ b/websocket_server/config.py @@ -22,6 +22,8 @@ ASSIGN_INTERVAL_SECONDS = 30 LIVENESS_CHECK_INTERVAL_SECONDS = 5 # How often to cross check connected_workers with ping results from Redis ENABLE_WEBSOCKET_PING_DEBUG = os.getenv('ENABLE_WEBSOCKET_PING_DEBUG', 'false').lower() == 'true' + +ENABLE_TASK_ASSIGNMENT_DEBUG = os.getenv('ENABLE_TASK_ASSIGNMENT_DEBUG', 'false').lower() == 'true' # Database configuration DATABASE_URL = "mysql://root:password@172.17.0.1:3306/theapi" diff --git a/websocket_server/redis_utils.py b/websocket_server/redis_utils.py index 9d9c6a9..1ba08dc 100644 --- a/websocket_server/redis_utils.py +++ b/websocket_server/redis_utils.py @@ -8,6 +8,8 @@ from websocket_server.task_assigner import assign_task_to_worker from websocket_server.worker_manager import notify_worker_disconnect import logging +from websocket_server.config import ENABLE_TASK_ASSIGNMENT_DEBUG + logger = logging.getLogger("websocket_server") # Shared state @@ -38,16 +40,20 @@ def redis_subscribe(worker_id): message = pubsub.get_message(timeout=1.0) if message and message["type"] in ["pmessage", "message"]: current_status = redis_client.get(f"worker_status_{worker_id}") - logger.debug(f"[{worker_id}] Current Redis status: {current_status}") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Current Redis status: {current_status}") if current_status and current_status.lower() == "idle": if redis_client.get(f"worker_queue_{worker_id}").lower() == "true": - logger.info(f"[{worker_id}] Queue is true, setting dispatch flag") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Queue is true, setting dispatch flag") worker_dispatch_flags[worker_id] = True else: - logger.info(f"[{worker_id}] Queue false, no dispatch needed") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Queue false, no dispatch needed") else: - logger.info(f"[{worker_id}] Status is not idle, skipping") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Status is not idle, skipping") logger.warning(f"[{worker_id}] Redis pub/sub loop exited, disconnecting") notify_worker_disconnect(worker_id) diff --git a/websocket_server/task_assigner.py b/websocket_server/task_assigner.py index 437b782..2c9d1e0 100644 --- a/websocket_server/task_assigner.py +++ b/websocket_server/task_assigner.py @@ -9,6 +9,7 @@ from websocket_server.models import Task import logging from websocket_server.shared_state import connected_workers, worker_lock from websocket_server.events import base +from websocket_server.config import ENABLE_TASK_ASSIGNMENT_DEBUG logger = logging.getLogger("websocket_server") @@ -21,14 +22,17 @@ def assign_task_to_worker(worker_id): max_lock_retries = 5 lock_retry_wait = 0.2 - logger.info(f"[{worker_id}] Starting task assignment process") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Starting task assignment process") try: redis_client.set(f"worker_status_{worker_id}", "busy") - logger.debug(f"[{worker_id}] Marked as 'busy' in Redis") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Marked as 'busy' in Redis") for attempt in range(max_lock_retries): - logger.debug(f"[{worker_id}] Lock attempt {attempt + 1}/{max_lock_retries}") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Lock attempt {attempt + 1}/{max_lock_retries}") with get_db_session() as session: now = datetime.utcnow() candidates = session.query(Task).filter( @@ -50,7 +54,8 @@ def assign_task_to_worker(worker_id): ready.append(task) if not ready: - logger.info(f"[{worker_id}] No ready tasks with satisfied dependencies") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] No ready tasks with satisfied dependencies") redis_client.set(f"worker_queue_{worker_id}", "False") redis_client.set(f"worker_status_{worker_id}", "idle") return @@ -63,21 +68,25 @@ def assign_task_to_worker(worker_id): ).with_for_update(skip_locked=True).limit(1).one_or_none() if selected: - logger.info(f"[{worker_id}] Task selected: {selected.id}") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Task selected: {selected.id}") break - logger.warning(f"[{worker_id}] Lock contention, retrying...") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Lock contention, retrying...") time.sleep(lock_retry_wait) lock_retry_wait *= 2 else: - logger.error(f"[{worker_id}] Failed to acquire task after {max_lock_retries} attempts") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Failed to acquire task after {max_lock_retries} attempts") redis_client.set(f"worker_status_{worker_id}", "idle") redis_client.set(f"worker_queue_{worker_id}", "False") return with worker_lock: if worker_id in connected_workers: - logger.debug(f"[{worker_id}] Dispatching task {selected.id}") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{worker_id}] Dispatching task {selected.id}") base.socketio.emit("task", { "task_id": selected.id, "worker_id": worker_id, diff --git a/websocket_server/worker_manager.py b/websocket_server/worker_manager.py index 5a1e7b4..5c1d838 100644 --- a/websocket_server/worker_manager.py +++ b/websocket_server/worker_manager.py @@ -6,7 +6,7 @@ import time import threading import logging from websocket_server.shared_state import connected_workers, worker_lock, ping_tracker -from websocket_server.config import get_redis_client,PING_INTERVAL_SECONDS, ASSIGN_INTERVAL_SECONDS, LIVENESS_CHECK_INTERVAL_SECONDS, PING_EXPIRY_SECONDS, ENABLE_WEBSOCKET_PING_DEBUG +from websocket_server.config import get_redis_client,PING_INTERVAL_SECONDS, ASSIGN_INTERVAL_SECONDS, LIVENESS_CHECK_INTERVAL_SECONDS, PING_EXPIRY_SECONDS, ENABLE_WEBSOCKET_PING_DEBUG, ENABLE_TASK_ASSIGNMENT_DEBUG from websocket_server.task_assigner import assign_task_to_worker from websocket_server.events import base @@ -66,7 +66,8 @@ def worker_dispatch_flag_check(): logger.info("Worker dispatch flag thread started") while True: for wid in list(worker_dispatch_flags.keys()): - logger.info(f"[{wid}] Dispatch flag set. Triggering task assign.") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{wid}] Dispatch flag set. Triggering task assign.") del worker_dispatch_flags[wid] assign_task_to_worker(wid) threading.Event().wait(0.01) @@ -109,9 +110,11 @@ def all_worker_watchdog(): # 2. Trigger task assignment if due if now - last_assign_time >= ASSIGN_INTERVAL_SECONDS: if worker_ids: - logger.debug("Checking for task assignment across all workers") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug("Checking for task assignment across all workers") for wid in worker_ids: - logger.info(f"[{wid}] Triggering assign check") + if ENABLE_TASK_ASSIGNMENT_DEBUG: + logger.debug(f"[{wid}] Triggering assign check") assign_task_to_worker(wid) last_assign_time = now