|
1 |
+#!/usr/bin/env python3 |
|
2 |
+""" |
|
3 |
+maclustr-mobile — coordinateur des nœuds mobiles (iPhone / iPad) du cluster MacLustr. |
|
4 |
+ |
|
5 |
+Les iPhones ne peuvent pas recevoir de SSH ni tourner en arrière-plan : ils *tirent* donc |
|
6 |
+leur travail. L'app MacLustr iOS (onglet « Nœud ») s'enregistre ici, envoie un battement |
|
7 |
+toutes les ~20 s (batterie, thermique, réseau, mémoire) et fait du long-poll sur |
|
8 |
+/api/worker/next ; chaque job exécuté est renvoyé sur /api/worker/result. |
|
9 |
+ |
|
10 |
+Côté cluster, on soumet des jobs (HTTP fetch distribué, code JavaScript, infos système) |
|
11 |
+via /api/jobs et on lit les résultats. L'agent maclustr-agentd relit /api/workers pour |
|
12 |
+afficher les mobiles dans les apps MacLustr macOS / iOS. |
|
13 |
+ |
|
14 |
+Python 3.9 stdlib uniquement (tourne avec /usr/bin/python3 sous PM2). Aucun accès réseau |
|
15 |
+sortant : seulement du trafic entrant des iPhones (Tailscale ou LAN) et des clients. |
|
16 |
+ |
|
17 |
+Variables d'environnement : |
|
18 |
+ MOBILE_PORT (9320) MOBILE_TOKEN (obligatoire) MOBILE_DATA (./data) |
|
19 |
+ MOBILE_ONLINE_S (45) battement plus vieux → worker « hors ligne » |
|
20 |
+ MOBILE_LEASE_GRACE_S (30) marge ajoutée au timeout d'un job avant de le re-mettre en file |
|
21 |
+""" |
|
22 |
+import json |
|
23 |
+import os |
|
24 |
+import sqlite3 |
|
25 |
+import sys |
|
26 |
+import threading |
|
27 |
+import time |
|
28 |
+import uuid |
|
29 |
+from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer |
|
30 |
+from urllib.parse import urlparse, parse_qs |
|
31 |
+ |
|
32 |
+VERSION = "1.1.0" # 1.1.0 : nœuds mobiles de plein droit (alias, pinned, historique métriques, détail, PATCH, compute.bench / net.probe) |
|
33 |
+PORT = int(os.environ.get("MOBILE_PORT", "9320")) |
|
34 |
+TOKEN = os.environ.get("MOBILE_TOKEN", "") |
|
35 |
+DATA_DIR = os.path.abspath(os.environ.get("MOBILE_DATA", os.path.join(os.path.dirname(os.path.abspath(__file__)), "..", "data"))) |
|
36 |
+DB_PATH = os.path.join(DATA_DIR, "mobile.db") |
|
37 |
+ONLINE_S = float(os.environ.get("MOBILE_ONLINE_S", "45")) |
|
38 |
+LEASE_GRACE_S = float(os.environ.get("MOBILE_LEASE_GRACE_S", "30")) |
|
39 |
+MAX_WAIT_S = 30 # borne du long-poll |
|
40 |
+MAX_PAYLOAD = 1024 * 1024 # 1 Mo par job soumis |
|
41 |
+MAX_RESULT = 4 * 1024 * 1024 # 4 Mo par résultat conservé |
|
42 |
+JOB_TYPES = ("http.fetch", "js.run", "sys.info", "ping", "compute.bench", "net.probe") |
|
43 |
+DEFAULT_TIMEOUT_S = {"http.fetch": 30, "js.run": 60, "sys.info": 10, "ping": 10, "compute.bench": 60, "net.probe": 30} |
|
44 |
+METRICS_RETENTION_S = 24 * 3600 |
|
45 |
+METRIC_KEYS = ("battery", "charging", "thermal", "lowPower", "memUsedMb", "memTotalMb", "memAvailableMb", |
|
46 |
+ "diskFreeGb", "diskTotalGb", "net", "cpuLoad", "cpuUsage", "uptimeS", "concurrency", "jobsDone", "jobsFailed", |
|
47 |
+ "activeJobs", "screenOn", "workerUptimeS", "bytesFetched", "lastJobAt", "ip", "ssid", "brightness") |
|
48 |
+HISTORY_DAYS = 7 |
|
49 |
+ |
|
50 |
+if not TOKEN: |
|
51 |
+ print("MOBILE_TOKEN manquant", file=sys.stderr) |
|
52 |
+ sys.exit(2) |
|
53 |
+ |
|
54 |
+started_at = time.time() |
|
55 |
+state_lock = threading.Lock() |
|
56 |
+queue_cv = threading.Condition(state_lock) # réveil des pollers (nouveau job) |
|
57 |
+result_cv = threading.Condition(state_lock) # réveil des clients qui attendent un résultat |
|
58 |
+ |
|
59 |
+workers = {} # name -> dict (état live : dernier battement, métriques, jobs en cours) |
|
60 |
+queued = [] # liste ordonnée d'ids de jobs en file |
|
61 |
+jobs = {} # id -> dict (jobs vivants : queued/leased ; les terminés vivent en base) |
|
62 |
+recent = {} # id -> dict des jobs terminés récemment (cache 1 h, pour /wait et /jobs) |
|
63 |
+stats = {"submitted": 0, "done": 0, "failed": 0, "expired": 0, "bytes_in": 0} |
|
64 |
+ |
|
65 |
+ |
|
66 |
+# --------------------------------------------------------------------------- |
|
67 |
+# Base SQLite (persistance des workers, des jobs et de leurs résultats) |
|
68 |
+# --------------------------------------------------------------------------- |
|
69 |
+ |
|
70 |
+def db(): |
|
71 |
+ os.makedirs(DATA_DIR, exist_ok=True) |
|
72 |
+ c = sqlite3.connect(DB_PATH, timeout=10) |
|
73 |
+ c.execute("PRAGMA journal_mode=WAL") |
|
74 |
+ c.execute("""CREATE TABLE IF NOT EXISTS workers ( |
|
75 |
+ name TEXT PRIMARY KEY, identifier TEXT, model TEXT, marketing TEXT, os TEXT, chip TEXT, |
|
76 |
+ ram_mb INTEGER, storage_gb REAL, app_version TEXT, |
|
77 |
+ first_seen REAL, last_seen REAL, jobs_done INTEGER DEFAULT 0, jobs_failed INTEGER DEFAULT 0, |
|
78 |
+ meta TEXT)""") |
|
79 |
+ for col, typ in (("alias", "TEXT"), ("pinned", "INTEGER DEFAULT 0"), ("role", "TEXT"), ("kind", "TEXT"), |
|
80 |
+ ("platform", "TEXT"), ("cores", "INTEGER"), ("bench", "TEXT"), ("notes", "TEXT")): |
|
81 |
+ try: |
|
82 |
+ c.execute("ALTER TABLE workers ADD COLUMN %s %s" % (col, typ)) |
|
83 |
+ except sqlite3.OperationalError: |
|
84 |
+ pass |
|
85 |
+ c.execute("""CREATE TABLE IF NOT EXISTS worker_metrics ( |
|
86 |
+ ts REAL, name TEXT, battery REAL, charging INTEGER, thermal TEXT, mem_used_mb INTEGER, mem_avail_mb INTEGER, |
|
87 |
+ cpu REAL, net TEXT, active INTEGER, disk_free_gb REAL)""") |
|
88 |
+ c.execute("CREATE INDEX IF NOT EXISTS wm_name_ts ON worker_metrics(name, ts)") |
|
89 |
+ c.execute("""CREATE TABLE IF NOT EXISTS jobs ( |
|
90 |
+ id TEXT PRIMARY KEY, type TEXT, payload TEXT, target TEXT, tag TEXT, state TEXT, |
|
91 |
+ attempts INTEGER DEFAULT 0, max_attempts INTEGER DEFAULT 3, timeout_s REAL, |
|
92 |
+ created REAL, leased_at REAL, lease_worker TEXT, lease_expires REAL, |
|
93 |
+ finished REAL, duration_ms INTEGER, worker TEXT, result TEXT, error TEXT, submitted_by TEXT)""") |
|
94 |
+ c.execute("CREATE INDEX IF NOT EXISTS jobs_state ON jobs(state, created)") |
|
95 |
+ c.execute("CREATE INDEX IF NOT EXISTS jobs_tag ON jobs(tag)") |
|
96 |
+ return c |
|
97 |
+ |
|
98 |
+ |
|
99 |
+def db_upsert_worker(w): |
|
100 |
+ c = db() |
|
101 |
+ with c: |
|
102 |
+ c.execute("""INSERT INTO workers(name, identifier, model, marketing, os, chip, ram_mb, storage_gb, app_version, |
|
103 |
+ first_seen, last_seen, meta, alias, pinned, role, kind, platform, cores, bench, notes) |
|
104 |
+ VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) |
|
105 |
+ ON CONFLICT(name) DO UPDATE SET identifier=excluded.identifier, model=excluded.model, |
|
106 |
+ marketing=excluded.marketing, os=excluded.os, chip=excluded.chip, ram_mb=excluded.ram_mb, |
|
107 |
+ storage_gb=excluded.storage_gb, app_version=excluded.app_version, last_seen=excluded.last_seen, |
|
108 |
+ meta=excluded.meta, alias=excluded.alias, pinned=excluded.pinned, role=excluded.role, |
|
109 |
+ kind=excluded.kind, platform=excluded.platform, cores=excluded.cores, bench=excluded.bench, notes=excluded.notes""", |
|
110 |
+ (w["name"], w.get("identifier"), w.get("model"), w.get("marketing"), w.get("os"), w.get("chip"), |
|
111 |
+ w.get("ramMb"), w.get("storageGb"), w.get("appVersion"), w.get("firstSeen"), w.get("lastSeen"), |
|
112 |
+ json.dumps(w.get("meta") or {}), w.get("alias"), 1 if w.get("pinned") else 0, w.get("role") or "worker", |
|
113 |
+ w.get("kind"), w.get("platform"), w.get("cores"), json.dumps(w.get("bench")) if w.get("bench") else None, w.get("notes"))) |
|
114 |
+ c.close() |
|
115 |
+ |
|
116 |
+ |
|
117 |
+def db_record_metrics(name, ts, m): |
|
118 |
+ try: |
|
119 |
+ c = db() |
|
120 |
+ with c: |
|
121 |
+ c.execute("INSERT INTO worker_metrics VALUES (?,?,?,?,?,?,?,?,?,?,?)", |
|
122 |
+ (ts, name, m.get("battery"), 1 if m.get("charging") else 0, m.get("thermal"), m.get("memUsedMb"), |
|
123 |
+ m.get("memAvailableMb"), m.get("cpuUsage", m.get("cpuLoad")), m.get("net"), m.get("activeJobs"), m.get("diskFreeGb"))) |
|
124 |
+ c.execute("DELETE FROM worker_metrics WHERE ts < ?", (ts - METRICS_RETENTION_S,)) |
|
125 |
+ c.close() |
|
126 |
+ except Exception as e: |
|
127 |
+ print("metrics db:", e, flush=True) |
|
128 |
+ |
|
129 |
+ |
|
130 |
+def db_metrics_history(name, window_s, bucket_s): |
|
131 |
+ since = time.time() - window_s |
|
132 |
+ c = db() |
|
133 |
+ rows = c.execute( |
|
134 |
+ "SELECT (CAST(ts AS INTEGER)/%d)*%d AS b, AVG(battery), AVG(mem_used_mb), AVG(cpu), MAX(active), MAX(charging)" |
|
135 |
+ " FROM worker_metrics WHERE name=? AND ts>=? GROUP BY b ORDER BY b" % (bucket_s, bucket_s), (name, since)).fetchall() |
|
136 |
+ c.close() |
|
137 |
+ return [{"ts": r[0], "battery": round(r[1], 1) if r[1] is not None else None, |
|
138 |
+ "memUsedMb": round(r[2]) if r[2] is not None else None, "cpu": round(r[3], 1) if r[3] is not None else None, |
|
139 |
+ "active": r[4] or 0, "charging": bool(r[5])} for r in rows] |
|
140 |
+ |
|
141 |
+ |
|
142 |
+def db_touch_worker(name, ts, done_inc=0, failed_inc=0): |
|
143 |
+ c = db() |
|
144 |
+ with c: |
|
145 |
+ c.execute("UPDATE workers SET last_seen=?, jobs_done=jobs_done+?, jobs_failed=jobs_failed+? WHERE name=?", |
|
146 |
+ (ts, done_inc, failed_inc, name)) |
|
147 |
+ c.close() |
|
148 |
+ |
|
149 |
+ |
|
150 |
+def db_insert_job(j): |
|
151 |
+ c = db() |
|
152 |
+ with c: |
|
153 |
+ c.execute("""INSERT INTO jobs(id, type, payload, target, tag, state, attempts, max_attempts, timeout_s, created, submitted_by) |
|
154 |
+ VALUES(?,?,?,?,?,?,?,?,?,?,?)""", |
|
155 |
+ (j["id"], j["type"], json.dumps(j["payload"], ensure_ascii=False), j.get("target"), j.get("tag"), |
|
156 |
+ j["state"], j["attempts"], j["maxAttempts"], j["timeoutS"], j["created"], j.get("submittedBy"))) |
|
157 |
+ c.close() |
|
158 |
+ |
|
159 |
+ |
|
160 |
+def db_update_job(j): |
|
161 |
+ c = db() |
|
162 |
+ with c: |
|
163 |
+ c.execute("""UPDATE jobs SET state=?, attempts=?, leased_at=?, lease_worker=?, lease_expires=?, finished=?, |
|
164 |
+ duration_ms=?, worker=?, result=?, error=? WHERE id=?""", |
|
165 |
+ (j["state"], j["attempts"], j.get("leasedAt"), j.get("leaseWorker"), j.get("leaseExpires"), |
|
166 |
+ j.get("finished"), j.get("durationMs"), j.get("worker"), |
|
167 |
+ json.dumps(j.get("result"), ensure_ascii=False) if j.get("result") is not None else None, |
|
168 |
+ j.get("error"), j["id"])) |
|
169 |
+ c.close() |
|
170 |
+ |
|
171 |
+ |
|
172 |
+def db_get_job(job_id): |
|
173 |
+ c = db() |
|
174 |
+ cur = c.execute("SELECT * FROM jobs WHERE id=?", (job_id,)) |
|
175 |
+ row = cur.fetchone() |
|
176 |
+ cols = [d[0] for d in cur.description] if row else [] |
|
177 |
+ c.close() |
|
178 |
+ return row_to_job(dict(zip(cols, row))) if row else None |
|
179 |
+ |
|
180 |
+ |
|
181 |
+def db_list_jobs(state=None, tag=None, limit=100): |
|
182 |
+ c = db() |
|
183 |
+ sql = "SELECT * FROM jobs" |
|
184 |
+ cond, args = [], [] |
|
185 |
+ if state: |
|
186 |
+ cond.append("state=?"); args.append(state) |
|
187 |
+ if tag: |
|
188 |
+ cond.append("tag=?"); args.append(tag) |
|
189 |
+ if cond: |
|
190 |
+ sql += " WHERE " + " AND ".join(cond) |
|
191 |
+ sql += " ORDER BY created DESC LIMIT ?" |
|
192 |
+ args.append(int(limit)) |
|
193 |
+ cur = c.execute(sql, args) |
|
194 |
+ rows = cur.fetchall() |
|
195 |
+ cols = [d[0] for d in cur.description] |
|
196 |
+ c.close() |
|
197 |
+ return [row_to_job(dict(zip(cols, r))) for r in rows] |
|
198 |
+ |
|
199 |
+ |
|
200 |
+def db_load_workers(): |
|
201 |
+ c = db() |
|
202 |
+ cur = c.execute("SELECT * FROM workers") |
|
203 |
+ rows = cur.fetchall() |
|
204 |
+ cols = [d[0] for d in cur.description] |
|
205 |
+ c.close() |
|
206 |
+ out = {} |
|
207 |
+ for r in rows: |
|
208 |
+ d = dict(zip(cols, r)) |
|
209 |
+ out[d["name"]] = { |
|
210 |
+ "name": d["name"], "identifier": d["identifier"], "model": d["model"], "marketing": d["marketing"], |
|
211 |
+ "os": d["os"], "chip": d["chip"], "ramMb": d["ram_mb"], "storageGb": d["storage_gb"], |
|
212 |
+ "appVersion": d["app_version"], "firstSeen": d["first_seen"], "lastSeen": d["last_seen"], |
|
213 |
+ "jobsDone": d["jobs_done"] or 0, "jobsFailed": d["jobs_failed"] or 0, |
|
214 |
+ "meta": json.loads(d["meta"] or "{}"), "metrics": {}, "active": [], |
|
215 |
+ "alias": d.get("alias"), "pinned": bool(d.get("pinned")), "role": d.get("role") or "worker", |
|
216 |
+ "kind": d.get("kind") or _kind_of(d["model"]), "platform": d.get("platform") or _platform_of(d["model"]), |
|
217 |
+ "cores": d.get("cores"), "bench": json.loads(d["bench"]) if d.get("bench") else None, "notes": d.get("notes"), |
|
218 |
+ } |
|
219 |
+ return out |
|
220 |
+ |
|
221 |
+ |
|
222 |
+def db_recover_jobs(): |
|
223 |
+ """Au démarrage : les jobs queued/leased reviennent en file.""" |
|
224 |
+ c = db() |
|
225 |
+ cur = c.execute("SELECT * FROM jobs WHERE state IN ('queued','leased') ORDER BY created") |
|
226 |
+ rows = cur.fetchall() |
|
227 |
+ cols = [d[0] for d in cur.description] |
|
228 |
+ c.close() |
|
229 |
+ out = [] |
|
230 |
+ for r in rows: |
|
231 |
+ j = row_to_job(dict(zip(cols, r))) |
|
232 |
+ j["state"] = "queued" |
|
233 |
+ j["leasedAt"] = j["leaseWorker"] = j["leaseExpires"] = None |
|
234 |
+ out.append(j) |
|
235 |
+ return out |
|
236 |
+ |
|
237 |
+ |
|
238 |
+def db_purge(): |
|
239 |
+ c = db() |
|
240 |
+ with c: |
|
241 |
+ c.execute("DELETE FROM jobs WHERE state IN ('done','failed','expired','cancelled') AND created < ?", |
|
242 |
+ (time.time() - HISTORY_DAYS * 86400,)) |
|
243 |
+ c.close() |
|
244 |
+ |
|
245 |
+ |
|
246 |
+def row_to_job(d): |
|
247 |
+ return { |
|
248 |
+ "id": d["id"], "type": d["type"], "payload": json.loads(d["payload"] or "{}"), "target": d["target"], |
|
249 |
+ "tag": d["tag"], "state": d["state"], "attempts": d["attempts"] or 0, "maxAttempts": d["max_attempts"] or 3, |
|
250 |
+ "timeoutS": d["timeout_s"], "created": d["created"], "leasedAt": d["leased_at"], |
|
251 |
+ "leaseWorker": d["lease_worker"], "leaseExpires": d["lease_expires"], "finished": d["finished"], |
|
252 |
+ "durationMs": d["duration_ms"], "worker": d["worker"], |
|
253 |
+ "result": json.loads(d["result"]) if d.get("result") else None, "error": d["error"], |
|
254 |
+ "submittedBy": d["submitted_by"], |
|
255 |
+ } |
|
256 |
+ |
|
257 |
+ |
|
258 |
+# --------------------------------------------------------------------------- |
|
259 |
+# Logique de file |
|
260 |
+# --------------------------------------------------------------------------- |
|
261 |
+ |
|
262 |
+def _kind_of(model): |
|
263 |
+ m = (model or "").lower() |
|
264 |
+ return "ipad" if m.startswith("ipad") else ("iphone" if m.startswith("iphone") else "mobile") |
|
265 |
+ |
|
266 |
+ |
|
267 |
+def _platform_of(model): |
|
268 |
+ return "ipados" if _kind_of(model) == "ipad" else "ios" |
|
269 |
+ |
|
270 |
+ |
|
271 |
+def worker_online(w, now=None): |
|
272 |
+ now = now or time.time() |
|
273 |
+ return (w.get("lastSeen") or 0) >= now - ONLINE_S |
|
274 |
+ |
|
275 |
+ |
|
276 |
+def public_worker(w, now=None): |
|
277 |
+ now = now or time.time() |
|
278 |
+ d = {k: v for k, v in w.items() if k not in ("identifier",)} |
|
279 |
+ d["online"] = worker_online(w, now) |
|
280 |
+ d["ageS"] = round(now - (w.get("lastSeen") or 0), 1) |
|
281 |
+ d["activeJobs"] = len(w.get("active") or []) |
|
282 |
+ d["displayName"] = w.get("alias") or w["name"] |
|
283 |
+ d.setdefault("kind", _kind_of(w.get("model"))) |
|
284 |
+ d.setdefault("platform", _platform_of(w.get("model"))) |
|
285 |
+ d.setdefault("role", "worker") |
|
286 |
+ d.setdefault("pinned", False) |
|
287 |
+ return d |
|
288 |
+ |
|
289 |
+ |
|
290 |
+def public_job(j, with_result=True): |
|
291 |
+ d = dict(j) |
|
292 |
+ if not with_result: |
|
293 |
+ d.pop("result", None) |
|
294 |
+ if isinstance(j.get("payload"), dict) and "code" in j["payload"]: |
|
295 |
+ d["payload"] = dict(j["payload"], code="<%d octets>" % len(j["payload"].get("code") or "")) |
|
296 |
+ return d |
|
297 |
+ |
|
298 |
+ |
|
299 |
+def submit_job(spec, submitted_by=None): |
|
300 |
+ t = spec.get("type") |
|
301 |
+ if t not in JOB_TYPES: |
|
302 |
+ raise ValueError("type inconnu : %r (attendu : %s)" % (t, ", ".join(JOB_TYPES))) |
|
303 |
+ payload = spec.get("payload") or {} |
|
304 |
+ if not isinstance(payload, dict): |
|
305 |
+ raise ValueError("payload doit être un objet") |
|
306 |
+ if t == "http.fetch" and not payload.get("url"): |
|
307 |
+ raise ValueError("http.fetch : payload.url requis") |
|
308 |
+ if t == "js.run" and not payload.get("code"): |
|
309 |
+ raise ValueError("js.run : payload.code requis (définir function main(args) { … })") |
|
310 |
+ if t == "net.probe" and not payload.get("hosts"): |
|
311 |
+ raise ValueError("net.probe : payload.hosts (liste d'hôtes ou d'URL) requis") |
|
312 |
+ if len(json.dumps(payload)) > MAX_PAYLOAD: |
|
313 |
+ raise ValueError("payload trop gros (> 1 Mo)") |
|
314 |
+ timeout_s = float(spec.get("timeoutS") or DEFAULT_TIMEOUT_S[t]) |
|
315 |
+ timeout_s = max(2.0, min(timeout_s, 600.0)) |
|
316 |
+ j = { |
|
317 |
+ "id": spec.get("id") or uuid.uuid4().hex[:12], "type": t, "payload": payload, |
|
318 |
+ "target": (spec.get("target") or None), "tag": spec.get("tag"), "state": "queued", |
|
319 |
+ "attempts": 0, "maxAttempts": int(spec.get("maxAttempts") or 3), "timeoutS": timeout_s, |
|
320 |
+ "created": time.time(), "leasedAt": None, "leaseWorker": None, "leaseExpires": None, |
|
321 |
+ "finished": None, "durationMs": None, "worker": None, "result": None, "error": None, |
|
322 |
+ "submittedBy": submitted_by, |
|
323 |
+ } |
|
324 |
+ if j["target"] and j["target"] not in ("any",): |
|
325 |
+ with state_lock: |
|
326 |
+ if j["target"] not in workers: |
|
327 |
+ by_alias = next((n for n, w in workers.items() if (w.get("alias") or "").lower() == str(j["target"]).lower()), None) |
|
328 |
+ if not by_alias: |
|
329 |
+ raise ValueError("worker cible inconnu : %s" % j["target"]) |
|
330 |
+ j["target"] = by_alias |
|
331 |
+ db_insert_job(j) |
|
332 |
+ with queue_cv: |
|
333 |
+ jobs[j["id"]] = j |
|
334 |
+ queued.append(j["id"]) |
|
335 |
+ stats["submitted"] += 1 |
|
336 |
+ queue_cv.notify_all() |
|
337 |
+ return j |
|
338 |
+ |
|
339 |
+ |
|
340 |
+def _pick_job_for(name): |
|
341 |
+ """Premier job en file compatible avec ce worker (cible = ce worker ou n'importe qui).""" |
|
342 |
+ for i, jid in enumerate(queued): |
|
343 |
+ j = jobs.get(jid) |
|
344 |
+ if not j: |
|
345 |
+ continue |
|
346 |
+ if j["target"] in (None, "any", name): |
|
347 |
+ del queued[i] |
|
348 |
+ return j |
|
349 |
+ return None |
|
350 |
+ |
|
351 |
+ |
|
352 |
+def lease_job(name, wait_s): |
|
353 |
+ deadline = time.time() + wait_s |
|
354 |
+ with queue_cv: |
|
355 |
+ while True: |
|
356 |
+ j = _pick_job_for(name) |
|
357 |
+ if j: |
|
358 |
+ now = time.time() |
|
359 |
+ j["state"] = "leased" |
|
360 |
+ j["attempts"] += 1 |
|
361 |
+ j["leasedAt"] = now |
|
362 |
+ j["leaseWorker"] = name |
|
363 |
+ j["leaseExpires"] = now + j["timeoutS"] + LEASE_GRACE_S |
|
364 |
+ w = workers.get(name) |
|
365 |
+ if w is not None: |
|
366 |
+ w.setdefault("active", []).append(j["id"]) |
|
367 |
+ break |
|
368 |
+ remaining = deadline - time.time() |
|
369 |
+ if remaining <= 0: |
|
370 |
+ return None |
|
371 |
+ queue_cv.wait(remaining) |
|
372 |
+ db_update_job(j) |
|
373 |
+ return j |
|
374 |
+ |
|
375 |
+ |
|
376 |
+def finish_job(name, job_id, ok, result, error, duration_ms): |
|
377 |
+ with result_cv: |
|
378 |
+ j = jobs.get(job_id) |
|
379 |
+ if not j: |
|
380 |
+ return None, "job inconnu ou déjà clos" |
|
381 |
+ if j["state"] != "leased" or j["leaseWorker"] != name: |
|
382 |
+ return None, "job non loué par %s (état %s, loué à %s)" % (name, j["state"], j["leaseWorker"]) |
|
383 |
+ now = time.time() |
|
384 |
+ if ok: |
|
385 |
+ j["state"] = "done" |
|
386 |
+ j["result"] = result |
|
387 |
+ stats["done"] += 1 |
|
388 |
+ if j["type"] == "compute.bench" and isinstance(result, dict): |
|
389 |
+ wb = workers.get(name) |
|
390 |
+ if wb is not None: |
|
391 |
+ wb["bench"] = {"score": result.get("score"), "singleCore": result.get("singleCore"), |
|
392 |
+ "multiCore": result.get("multiCore"), "ts": now} |
|
393 |
+ else: |
|
394 |
+ j["error"] = (error or "erreur inconnue")[:2000] |
|
395 |
+ if j["attempts"] < j["maxAttempts"] and _retryable(error): |
|
396 |
+ j["state"] = "queued" |
|
397 |
+ j["leasedAt"] = j["leaseWorker"] = j["leaseExpires"] = None |
|
398 |
+ queued.append(j["id"]) |
|
399 |
+ queue_cv.notify_all() |
|
400 |
+ else: |
|
401 |
+ j["state"] = "failed" |
|
402 |
+ stats["failed"] += 1 |
|
403 |
+ j["finished"] = now if j["state"] in ("done", "failed") else None |
|
404 |
+ j["durationMs"] = duration_ms |
|
405 |
+ j["worker"] = name |
|
406 |
+ w = workers.get(name) |
|
407 |
+ if w is not None: |
|
408 |
+ w["active"] = [x for x in (w.get("active") or []) if x != job_id] |
|
409 |
+ if j["state"] == "done": |
|
410 |
+ w["jobsDone"] = (w.get("jobsDone") or 0) + 1 |
|
411 |
+ elif j["state"] == "failed": |
|
412 |
+ w["jobsFailed"] = (w.get("jobsFailed") or 0) + 1 |
|
413 |
+ if j["state"] in ("done", "failed"): |
|
414 |
+ jobs.pop(job_id, None) |
|
415 |
+ recent[job_id] = j |
|
416 |
+ result_cv.notify_all() |
|
417 |
+ db_update_job(j) |
|
418 |
+ if j["state"] in ("done", "failed"): |
|
419 |
+ db_touch_worker(name, time.time(), 1 if j["state"] == "done" else 0, 1 if j["state"] == "failed" else 0) |
|
420 |
+ return j, None |
|
421 |
+ |
|
422 |
+ |
|
423 |
+def _retryable(error): |
|
424 |
+ e = (error or "").lower() |
|
425 |
+ # Une erreur « métier » (JS levée par le code, HTTP 4xx) n'est pas rejouée ; réseau/timeout oui. |
|
426 |
+ return any(k in e for k in ("timeout", "timed out", "network", "connexion", "connection", "réseau", "cancelled", "-1001", "-1005", "-1009")) |
|
427 |
+ |
|
428 |
+ |
|
429 |
+def cancel_job(job_id): |
|
430 |
+ with queue_cv: |
|
431 |
+ j = jobs.get(job_id) |
|
432 |
+ if not j or j["state"] != "queued": |
|
433 |
+ return None |
|
434 |
+ j["state"] = "cancelled" |
|
435 |
+ j["finished"] = time.time() |
|
436 |
+ try: |
|
437 |
+ queued.remove(job_id) |
|
438 |
+ except ValueError: |
|
439 |
+ pass |
|
440 |
+ jobs.pop(job_id, None) |
|
441 |
+ recent[job_id] = j |
|
442 |
+ db_update_job(j) |
|
443 |
+ return j |
|
444 |
+ |
|
445 |
+ |
|
446 |
+def wait_job(job_id, timeout_s): |
|
447 |
+ deadline = time.time() + min(max(timeout_s, 0), MAX_WAIT_S) |
|
448 |
+ with result_cv: |
|
449 |
+ while True: |
|
450 |
+ j = recent.get(job_id) |
|
451 |
+ if j and j["state"] in ("done", "failed", "cancelled", "expired"): |
|
452 |
+ return j |
|
453 |
+ if job_id not in jobs and not j: |
|
454 |
+ break |
|
455 |
+ remaining = deadline - time.time() |
|
456 |
+ if remaining <= 0: |
|
457 |
+ return jobs.get(job_id) or j |
|
458 |
+ result_cv.wait(remaining) |
|
459 |
+ return db_get_job(job_id) |
|
460 |
+ |
|
461 |
+ |
|
462 |
+def reaper_loop(): |
|
463 |
+ """Re-met en file les jobs dont le bail a expiré (iPhone parti en arrière-plan, app fermée…).""" |
|
464 |
+ last_purge = 0 |
|
465 |
+ while True: |
|
466 |
+ time.sleep(5) |
|
467 |
+ now = time.time() |
|
468 |
+ expired = [] |
|
469 |
+ with queue_cv: |
|
470 |
+ for j in list(jobs.values()): |
|
471 |
+ if j["state"] == "leased" and j.get("leaseExpires") and j["leaseExpires"] < now: |
|
472 |
+ w = workers.get(j["leaseWorker"] or "") |
|
473 |
+ if w is not None: |
|
474 |
+ w["active"] = [x for x in (w.get("active") or []) if x != j["id"]] |
|
475 |
+ if j["attempts"] < j["maxAttempts"]: |
|
476 |
+ j["state"] = "queued" |
|
477 |
+ j["error"] = "bail expiré sur %s" % j["leaseWorker"] |
|
478 |
+ j["leasedAt"] = j["leaseWorker"] = j["leaseExpires"] = None |
|
479 |
+ queued.append(j["id"]) |
|
480 |
+ else: |
|
481 |
+ j["state"] = "expired" |
|
482 |
+ j["error"] = "bail expiré %d fois" % j["attempts"] |
|
483 |
+ j["finished"] = now |
|
484 |
+ jobs.pop(j["id"], None) |
|
485 |
+ recent[j["id"]] = j |
|
486 |
+ stats["expired"] += 1 |
|
487 |
+ expired.append(j) |
|
488 |
+ if expired: |
|
489 |
+ queue_cv.notify_all() |
|
490 |
+ result_cv.notify_all() |
|
491 |
+ # cache des résultats récents : 1 h |
|
492 |
+ for jid in [k for k, v in recent.items() if (v.get("finished") or 0) < now - 3600]: |
|
493 |
+ recent.pop(jid, None) |
|
494 |
+ for j in expired: |
|
495 |
+ db_update_job(j) |
|
496 |
+ if now - last_purge > 3600: |
|
497 |
+ try: |
|
498 |
+ db_purge() |
|
499 |
+ except Exception as e: |
|
500 |
+ print("purge:", e, flush=True) |
|
501 |
+ last_purge = now |
|
502 |
+ |
|
503 |
+ |
|
504 |
+# --------------------------------------------------------------------------- |
|
505 |
+# HTTP |
|
506 |
+# --------------------------------------------------------------------------- |
|
507 |
+ |
|
508 |
+class Handler(BaseHTTPRequestHandler): |
|
509 |
+ server_version = "maclustr-mobile/" + VERSION |
|
510 |
+ protocol_version = "HTTP/1.1" |
|
511 |
+ |
|
512 |
+ def log_message(self, fmt, *args): |
|
513 |
+ pass |
|
514 |
+ |
|
515 |
+ def _send(self, code, payload=None): |
|
516 |
+ body = b"" if payload is None else json.dumps(payload, ensure_ascii=False).encode() |
|
517 |
+ self.send_response(code) |
|
518 |
+ self.send_header("Content-Type", "application/json; charset=utf-8") |
|
519 |
+ self.send_header("Content-Length", str(len(body))) |
|
520 |
+ self.send_header("Cache-Control", "no-store") |
|
521 |
+ self.end_headers() |
|
522 |
+ if body: |
|
523 |
+ self.wfile.write(body) |
|
524 |
+ |
|
525 |
+ def _auth(self): |
|
526 |
+ h = self.headers.get("Authorization", "") |
|
527 |
+ if h == "Bearer " + TOKEN: |
|
528 |
+ return True |
|
529 |
+ q = parse_qs(urlparse(self.path).query) |
|
530 |
+ if (q.get("token") or [""])[0] == TOKEN: |
|
531 |
+ return True |
|
532 |
+ self._send(401, {"error": "unauthorized"}) |
|
533 |
+ return False |
|
534 |
+ |
|
535 |
+ def _json(self): |
|
536 |
+ try: |
|
537 |
+ n = int(self.headers.get("Content-Length", 0)) |
|
538 |
+ if n > MAX_RESULT + MAX_PAYLOAD: |
|
539 |
+ self._send(413, {"error": "corps trop gros"}) |
|
540 |
+ return None |
|
541 |
+ raw = self.rfile.read(n) if n else b"{}" |
|
542 |
+ stats["bytes_in"] += n |
|
543 |
+ return json.loads(raw or b"{}") |
|
544 |
+ except Exception as e: |
|
545 |
+ self._send(400, {"error": "JSON invalide : %s" % e}) |
|
546 |
+ return None |
|
547 |
+ |
|
548 |
+ # -- GET -------------------------------------------------------------- |
|
549 |
+ def do_GET(self): |
|
550 |
+ u = urlparse(self.path) |
|
551 |
+ q = parse_qs(u.query) |
|
552 |
+ parts = [p for p in u.path.split("/") if p] |
|
553 |
+ now = time.time() |
|
554 |
+ if u.path == "/health": |
|
555 |
+ with state_lock: |
|
556 |
+ online = sum(1 for w in workers.values() if worker_online(w, now)) |
|
557 |
+ active = sum(len(w.get("active") or []) for w in workers.values()) |
|
558 |
+ nq = len(queued) |
|
559 |
+ return self._send(200, {"ok": True, "version": VERSION, "uptime": int(now - started_at), |
|
560 |
+ "workersOnline": online, "workersTotal": len(workers), |
|
561 |
+ "queued": nq, "active": active, "stats": stats}) |
|
562 |
+ if not self._auth(): |
|
563 |
+ return |
|
564 |
+ if u.path == "/api/workers": |
|
565 |
+ with state_lock: |
|
566 |
+ lst = [public_worker(w, now) for w in workers.values()] |
|
567 |
+ lst.sort(key=lambda w: (not w["online"], w["name"])) |
|
568 |
+ return self._send(200, {"workers": lst, "ts": now, "onlineS": ONLINE_S}) |
|
569 |
+ if len(parts) == 3 and parts[0] == "api" and parts[1] == "workers": |
|
570 |
+ with state_lock: |
|
571 |
+ w = _find_worker(parts[2]) |
|
572 |
+ pub = public_worker(w, now) if w else None |
|
573 |
+ if not pub: |
|
574 |
+ return self._send(404, {"error": "worker inconnu"}) |
|
575 |
+ pub["recentJobs"] = [public_job(j, False) for j in db_list_jobs(None, None, 300) if j.get("worker") == w["name"] or j.get("leaseWorker") == w["name"]][:25] |
|
576 |
+ pub["history"] = db_metrics_history(w["name"], 3600, 60) |
|
577 |
+ return self._send(200, pub) |
|
578 |
+ if len(parts) == 4 and parts[0] == "api" and parts[1] == "workers" and parts[3] == "history": |
|
579 |
+ with state_lock: |
|
580 |
+ w = _find_worker(parts[2]) |
|
581 |
+ if not w: |
|
582 |
+ return self._send(404, {"error": "worker inconnu"}) |
|
583 |
+ window = {"1h": 3600, "6h": 21600, "24h": 86400}.get((q.get("window") or ["1h"])[0], 3600) |
|
584 |
+ bucket = {3600: 60, 21600: 300, 86400: 900}[window] |
|
585 |
+ return self._send(200, {"name": w["name"], "points": db_metrics_history(w["name"], window, bucket), "window": window}) |
|
586 |
+ if u.path == "/api/stats": |
|
587 |
+ with state_lock: |
|
588 |
+ return self._send(200, {"stats": stats, "queued": len(queued), "leased": sum(1 for j in jobs.values() if j["state"] == "leased"), |
|
589 |
+ "workers": {n: {"online": worker_online(w, now), "active": len(w.get("active") or [])} for n, w in workers.items()}, |
|
590 |
+ "uptime": int(now - started_at)}) |
|
591 |
+ if u.path == "/api/worker/next": |
|
592 |
+ name = (q.get("worker") or [""])[0] |
|
593 |
+ wait = min(float((q.get("wait") or ["25"])[0]), MAX_WAIT_S) |
|
594 |
+ with state_lock: |
|
595 |
+ w = workers.get(name) |
|
596 |
+ if w is None: |
|
597 |
+ return self._send(409, {"error": "worker non enregistré : appeler /api/worker/register"}) |
|
598 |
+ w["lastSeen"] = now |
|
599 |
+ j = lease_job(name, wait) |
|
600 |
+ if not j: |
|
601 |
+ return self._send(204) |
|
602 |
+ return self._send(200, {"id": j["id"], "type": j["type"], "payload": j["payload"], |
|
603 |
+ "timeoutS": j["timeoutS"], "attempt": j["attempts"], "tag": j.get("tag")}) |
|
604 |
+ if u.path == "/api/jobs": |
|
605 |
+ state = (q.get("state") or [None])[0] |
|
606 |
+ tag = (q.get("tag") or [None])[0] |
|
607 |
+ wname = (q.get("worker") or [None])[0] |
|
608 |
+ limit = int((q.get("limit") or ["100"])[0]) |
|
609 |
+ full = (q.get("full") or ["0"])[0] == "1" |
|
610 |
+ with state_lock: |
|
611 |
+ if wname: |
|
612 |
+ ww = _find_worker(wname) |
|
613 |
+ wname = ww["name"] if ww else wname |
|
614 |
+ live = [public_job(j, full) for j in jobs.values() if (not state or j["state"] == state) and (not tag or j.get("tag") == tag) |
|
615 |
+ and (not wname or j.get("leaseWorker") == wname or j.get("worker") == wname or j.get("target") == wname)] |
|
616 |
+ if state in (None, "done", "failed", "expired", "cancelled"): |
|
617 |
+ stored = db_list_jobs(state, tag, limit if not wname else max(limit * 5, 300)) |
|
618 |
+ seen = {j["id"] for j in live} |
|
619 |
+ live += [public_job(j, full) for j in stored if j["id"] not in seen and (not wname or j.get("worker") == wname or j.get("target") == wname)] |
|
620 |
+ live.sort(key=lambda j: j["created"], reverse=True) |
|
621 |
+ return self._send(200, {"jobs": live[:limit], "count": len(live)}) |
|
622 |
+ if len(parts) == 3 and parts[0] == "api" and parts[1] == "jobs": |
|
623 |
+ with state_lock: |
|
624 |
+ j = jobs.get(parts[2]) or recent.get(parts[2]) |
|
625 |
+ if not j: |
|
626 |
+ j = db_get_job(parts[2]) |
|
627 |
+ if not j: |
|
628 |
+ return self._send(404, {"error": "job inconnu"}) |
|
629 |
+ return self._send(200, j) |
|
630 |
+ if len(parts) == 4 and parts[0] == "api" and parts[1] == "jobs" and parts[3] == "wait": |
|
631 |
+ timeout = float((q.get("timeout") or ["25"])[0]) |
|
632 |
+ j = wait_job(parts[2], timeout) |
|
633 |
+ if not j: |
|
634 |
+ return self._send(404, {"error": "job inconnu"}) |
|
635 |
+ return self._send(200, j) |
|
636 |
+ return self._send(404, {"error": "not found"}) |
|
637 |
+ |
|
638 |
+ # -- POST ------------------------------------------------------------- |
|
639 |
+ def do_POST(self): |
|
640 |
+ u = urlparse(self.path) |
|
641 |
+ if not self._auth(): |
|
642 |
+ return |
|
643 |
+ body = self._json() |
|
644 |
+ if body is None: |
|
645 |
+ return |
|
646 |
+ now = time.time() |
|
647 |
+ if u.path == "/api/worker/register": |
|
648 |
+ name = _clean_name(body.get("name")) |
|
649 |
+ if not name: |
|
650 |
+ return self._send(400, {"error": "name requis"}) |
|
651 |
+ with state_lock: |
|
652 |
+ w = workers.get(name) or {"name": name, "firstSeen": now, "jobsDone": 0, "jobsFailed": 0, "active": [], "metrics": {}} |
|
653 |
+ for k in ("identifier", "model", "marketing", "os", "chip", "ramMb", "storageGb", "appVersion", "cores", "kind", "platform", "role"): |
|
654 |
+ if body.get(k) is not None: |
|
655 |
+ w[k] = body[k] |
|
656 |
+ if body.get("alias") is not None: |
|
657 |
+ w["alias"] = _clean_name(body["alias"]) or None |
|
658 |
+ if body.get("pinned") is not None: |
|
659 |
+ w["pinned"] = bool(body["pinned"]) |
|
660 |
+ w.setdefault("kind", _kind_of(w.get("model"))) |
|
661 |
+ w.setdefault("platform", _platform_of(w.get("model"))) |
|
662 |
+ w.setdefault("role", "worker") |
|
663 |
+ w["meta"] = body.get("meta") or w.get("meta") or {} |
|
664 |
+ w["lastSeen"] = now |
|
665 |
+ w["active"] = [] # un (re)démarrage de l'app annule ses baux |
|
666 |
+ w["ip"] = self.client_address[0] |
|
667 |
+ workers[name] = w |
|
668 |
+ db_upsert_worker(w) |
|
669 |
+ print("register %s (%s, %s) depuis %s" % (name, w.get("marketing") or w.get("model"), w.get("os"), w["ip"]), flush=True) |
|
670 |
+ return self._send(200, {"ok": True, "name": name, "serverTime": now, "version": VERSION, |
|
671 |
+ "config": {"heartbeatS": 20, "pollWaitS": 25}}) |
|
672 |
+ if u.path == "/api/worker/heartbeat": |
|
673 |
+ name = _clean_name(body.get("name")) |
|
674 |
+ with state_lock: |
|
675 |
+ w = workers.get(name) |
|
676 |
+ if w is None: |
|
677 |
+ return self._send(409, {"error": "worker non enregistré"}) |
|
678 |
+ w["lastSeen"] = now |
|
679 |
+ w["ip"] = self.client_address[0] |
|
680 |
+ m = body.get("metrics") or {} |
|
681 |
+ if isinstance(m, dict): |
|
682 |
+ w["metrics"] = {k: m[k] for k in m if k in METRIC_KEYS} |
|
683 |
+ w["metrics"]["ts"] = now |
|
684 |
+ if body.get("appVersion"): |
|
685 |
+ w["appVersion"] = body["appVersion"] |
|
686 |
+ nq = len(queued) |
|
687 |
+ metrics_copy = dict(w.get("metrics") or {}) |
|
688 |
+ db_touch_worker(name, now) |
|
689 |
+ if metrics_copy: |
|
690 |
+ db_record_metrics(name, now, metrics_copy) |
|
691 |
+ return self._send(200, {"ok": True, "serverTime": now, "queued": nq, "config": {"heartbeatS": 20, "pollWaitS": 25}}) |
|
692 |
+ if u.path == "/api/worker/result": |
|
693 |
+ name = _clean_name(body.get("name")) |
|
694 |
+ job_id = body.get("jobId") |
|
695 |
+ result = body.get("result") |
|
696 |
+ if result is not None and len(json.dumps(result)) > MAX_RESULT: |
|
697 |
+ result = {"truncated": True, "note": "résultat > 4 Mo, tronqué côté coordinateur"} |
|
698 |
+ j, err = finish_job(name, job_id, bool(body.get("ok")), result, body.get("error"), body.get("durationMs")) |
|
699 |
+ if err: |
|
700 |
+ return self._send(409, {"error": err}) |
|
701 |
+ return self._send(200, {"ok": True, "state": j["state"]}) |
|
702 |
+ if u.path == "/api/jobs": |
|
703 |
+ specs = body.get("jobs") if isinstance(body, dict) and "jobs" in body else [body] |
|
704 |
+ if not isinstance(specs, list) or not specs: |
|
705 |
+ return self._send(400, {"error": "attendu un job ou {\"jobs\": [...]}"}) |
|
706 |
+ if len(specs) > 5000: |
|
707 |
+ return self._send(400, {"error": "max 5000 jobs par lot"}) |
|
708 |
+ ids, errors = [], [] |
|
709 |
+ who = self.client_address[0] |
|
710 |
+ for s in specs: |
|
711 |
+ try: |
|
712 |
+ ids.append(submit_job(s, who)["id"]) |
|
713 |
+ except ValueError as e: |
|
714 |
+ errors.append(str(e)) |
|
715 |
+ code = 200 if ids else 400 |
|
716 |
+ return self._send(code, {"ids": ids, "errors": errors, "count": len(ids)}) |
|
717 |
+ return self._send(404, {"error": "not found"}) |
|
718 |
+ |
|
719 |
+ def do_PATCH(self): |
|
720 |
+ u = urlparse(self.path) |
|
721 |
+ if not self._auth(): |
|
722 |
+ return |
|
723 |
+ body = self._json() |
|
724 |
+ if body is None: |
|
725 |
+ return |
|
726 |
+ parts = [p for p in u.path.split("/") if p] |
|
727 |
+ if len(parts) == 3 and parts[0] == "api" and parts[1] == "workers": |
|
728 |
+ with state_lock: |
|
729 |
+ w = _find_worker(parts[2]) |
|
730 |
+ if not w: |
|
731 |
+ return self._send(404, {"error": "worker inconnu"}) |
|
732 |
+ if "alias" in body: |
|
733 |
+ w["alias"] = _clean_name(body["alias"]) or None |
|
734 |
+ if "pinned" in body: |
|
735 |
+ w["pinned"] = bool(body["pinned"]) |
|
736 |
+ if "role" in body: |
|
737 |
+ w["role"] = str(body["role"])[:40] |
|
738 |
+ if "notes" in body: |
|
739 |
+ w["notes"] = str(body["notes"])[:500] |
|
740 |
+ if isinstance(body.get("meta"), dict): |
|
741 |
+ w["meta"] = dict(w.get("meta") or {}, **body["meta"]) |
|
742 |
+ snap = dict(w) |
|
743 |
+ db_upsert_worker(snap) |
|
744 |
+ return self._send(200, public_worker(snap)) |
|
745 |
+ return self._send(404, {"error": "not found"}) |
|
746 |
+ |
|
747 |
+ def do_DELETE(self): |
|
748 |
+ u = urlparse(self.path) |
|
749 |
+ if not self._auth(): |
|
750 |
+ return |
|
751 |
+ parts = [p for p in u.path.split("/") if p] |
|
752 |
+ if len(parts) == 3 and parts[0] == "api" and parts[1] == "jobs": |
|
753 |
+ j = cancel_job(parts[2]) |
|
754 |
+ if not j: |
|
755 |
+ return self._send(409, {"error": "job absent ou plus en file"}) |
|
756 |
+ return self._send(200, j) |
|
757 |
+ if len(parts) == 3 and parts[0] == "api" and parts[1] == "workers": |
|
758 |
+ with state_lock: |
|
759 |
+ w = workers.pop(parts[2], None) |
|
760 |
+ if not w: |
|
761 |
+ return self._send(404, {"error": "worker inconnu"}) |
|
762 |
+ c = db() |
|
763 |
+ with c: |
|
764 |
+ c.execute("DELETE FROM workers WHERE name=?", (parts[2],)) |
|
765 |
+ c.close() |
|
766 |
+ return self._send(200, {"ok": True}) |
|
767 |
+ return self._send(404, {"error": "not found"}) |
|
768 |
+ |
|
769 |
+ |
|
770 |
+def _find_worker(key): |
|
771 |
+ """Worker par nom ou par alias (insensible à la casse). À appeler sous state_lock.""" |
|
772 |
+ if key in workers: |
|
773 |
+ return workers[key] |
|
774 |
+ k = (key or "").lower() |
|
775 |
+ return next((w for w in workers.values() if (w.get("alias") or "").lower() == k or w["name"].lower() == k), None) |
|
776 |
+ |
|
777 |
+ |
|
778 |
+def _clean_name(n): |
|
779 |
+ n = (n or "").strip() |
|
780 |
+ return "".join(ch for ch in n if ch.isalnum() or ch in "-_.")[:40] |
|
781 |
+ |
|
782 |
+ |
|
783 |
+class Server(ThreadingHTTPServer): |
|
784 |
+ daemon_threads = True |
|
785 |
+ allow_reuse_address = True |
|
786 |
+ request_queue_size = 128 |
|
787 |
+ |
|
788 |
+ |
|
789 |
+def main(): |
|
790 |
+ global workers |
|
791 |
+ os.makedirs(DATA_DIR, exist_ok=True) |
|
792 |
+ workers = db_load_workers() |
|
793 |
+ for j in db_recover_jobs(): |
|
794 |
+ jobs[j["id"]] = j |
|
795 |
+ queued.append(j["id"]) |
|
796 |
+ db_update_job(j) |
|
797 |
+ threading.Thread(target=reaper_loop, daemon=True).start() |
|
798 |
+ srv = Server(("0.0.0.0", PORT), Handler) |
|
799 |
+ print("maclustr-mobile %s : %d worker(s) connus, %d job(s) repris, écoute :%d" % (VERSION, len(workers), len(queued), PORT), flush=True) |
|
800 |
+ srv.serve_forever() |
|
801 |
+ |
|
802 |
+ |
|
803 |
+if __name__ == "__main__": |
|
804 |
+ main() |
|
805 |
|