File size: 1,520 Bytes
310a607 75e1db7 310a607 85aed5e 310a607 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 | from typing import Dict, Any, Set
import time
workers: Dict[str, Dict[str, Any]] = {}
tasks: Dict[str, Dict[str, Any]] = {}
completed_task_ids: Set[str] = set()
results: list = []
training_active = False
training_target = ""
training_total_tasks = 0
training_completed_tasks = 0
task_counter = 0
result_save_counter = 0
RESULT_SAVE_BATCH = 10
def get_online_workers():
from config import HEARTBEAT_TIMEOUT
now = time.time()
return [wid for wid, w in workers.items() if now - w.get("last_heartbeat", 0) < HEARTBEAT_TIMEOUT]
def get_worker_by_id(worker_id: str):
return workers.get(worker_id)
def add_worker(worker_id: str, data: Dict):
workers[worker_id] = data
def update_worker(worker_id: str, **kwargs):
if worker_id in workers:
workers[worker_id].update(kwargs)
def remove_worker(worker_id: str):
if worker_id in workers:
del workers[worker_id]
def get_task(task_id: str):
return tasks.get(task_id)
def add_task(task_id: str, task_data: Dict):
tasks[task_id] = task_data
def update_task(task_id: str, **kwargs):
if task_id in tasks:
tasks[task_id].update(kwargs)
def remove_task(task_id: str):
if task_id in tasks:
del tasks[task_id]
def get_pending_tasks():
return [t for t in tasks.values() if t.get("status") == "pending"]
def get_assigned_tasks():
return [t for t in tasks.values() if t.get("status") == "assigned"]
def get_completed_tasks():
return [t for t in tasks.values() if t.get("status") == "completed"] |