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