| 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"] |