Files
3cloud-backend/websocket_server/worker_manager.py
T

91 lines
3.2 KiB
Python

# websocket_server/worker_manager.py
import threading
import requests
import logging
# from websocket_server.redis_utils import redis_subscribe, thread_stop_flags, worker_dispatch_flags
from websocket_server.task_assigner import assign_task_to_worker
from websocket_server.config import get_redis_client
from websocket_server.shared_state import connected_workers, worker_lock
logger = logging.getLogger("websocket_server")
# API server endpoint
api_server_url = "http://172.17.0.1:5000/api"
# Tracks threads for each worker
worker_threads = {}
def notify_worker_online(worker_id):
"""
Inform the API server that a worker has connected.
"""
logger.debug(f"[{worker_id}] Notifying API server: online")
payload = {"status": "online"}
headers = {"Content-Type": "application/json"}
try:
response = requests.put(f"{api_server_url}/workload_hosts/{worker_id}", json=payload, headers=headers)
logger.debug(response.text)
logger.info(f"[{worker_id}] API server acknowledged online state")
except Exception as e:
logger.error(f"[{worker_id}] Failed to notify API server of online status: {e}")
def notify_worker_disconnect(worker_id):
"""
Inform the API server that a worker has disconnected.
"""
logger.debug(f"[{worker_id}] Notifying API server: offline")
payload = {"status": "offline"}
headers = {"Content-Type": "application/json"}
try:
response = requests.put(f"{api_server_url}/workload_hosts/{worker_id}", json=payload, headers=headers)
logger.debug(response.text)
logger.info(f"[{worker_id}] API server acknowledged offline state")
except Exception as e:
logger.error(f"[{worker_id}] Failed to notify API server of disconnect: {e}")
def worker_dispatch_flag_check():
from websocket_server.redis_utils import redis_subscribe, thread_stop_flags, worker_dispatch_flags
"""
Continuously checks for workers with dispatch flags and triggers task assignment.
"""
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.")
del worker_dispatch_flags[wid]
assign_task_to_worker(wid)
threading.Event().wait(0.01)
def all_worker_watchdog():
"""
Periodically iterate over connected workers and ensure they are polled for task assignments.
Prevents idle starvation or missed Redis events.
"""
logger.info("Watchdog thread started")
while True:
for wid in list(connected_workers.keys()):
logger.info(f"[{wid}] Watchdog triggering assign check")
assign_task_to_worker(wid)
threading.Event().wait(30)
def start_worker_thread(worker_id):
from websocket_server.redis_utils import redis_subscribe, thread_stop_flags, worker_dispatch_flags
"""
Launch Redis subscription thread for the given worker.
"""
logger.info(f"[{worker_id}] Starting Redis listener thread")
thread_stop_flags[worker_id] = False
sub_thread = threading.Thread(target=redis_subscribe, args=(worker_id,), daemon=True)
worker_threads[worker_id] = sub_thread
sub_thread.start()