spb/maclustr-mobile
Public
Python 100%
1#!/usr/bin/env python32"""3maclustr-mobile — coordinateur des nœuds mobiles (iPhone / iPad) du cluster MacLustr.45Les iPhones ne peuvent pas recevoir de SSH ni tourner en arrière-plan : ils *tirent* donc6leur travail. L'app MacLustr iOS (onglet « Nœud ») s'enregistre ici, envoie un battement7toutes les ~20 s (batterie, thermique, réseau, mémoire) et fait du long-poll sur8/api/worker/next ; chaque job exécuté est renvoyé sur /api/worker/result.910Côté cluster, on soumet des jobs (HTTP fetch distribué, code JavaScript, infos système)11via /api/jobs et on lit les résultats. L'agent maclustr-agentd relit /api/workers pour12afficher les mobiles dans les apps MacLustr macOS / iOS.1314Python 3.9 stdlib uniquement (tourne avec /usr/bin/python3 sous PM2). Aucun accès réseau15sortant : seulement du trafic entrant des iPhones (Tailscale ou LAN) et des clients.1617Variables 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 file21"""22import json23import os24import sqlite325import sys26import threading27import time28import uuid29from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer30from urllib.parse import urlparse, parse_qs3132VERSION = "1.1.0" # 1.1.0 : nœuds mobiles de plein droit (alias, pinned, historique métriques, détail, PATCH, compute.bench / net.probe)33PORT = int(os.environ.get("MOBILE_PORT", "9320"))34TOKEN = os.environ.get("MOBILE_TOKEN", "")35DATA_DIR = os.path.abspath(os.environ.get("MOBILE_DATA", os.path.join(os.path.dirname(os.path.abspath(__file__)), "..", "data")))36DB_PATH = os.path.join(DATA_DIR, "mobile.db")37ONLINE_S = float(os.environ.get("MOBILE_ONLINE_S", "45"))38LEASE_GRACE_S = float(os.environ.get("MOBILE_LEASE_GRACE_S", "30"))39MAX_WAIT_S = 30 # borne du long-poll40MAX_PAYLOAD = 1024 * 1024 # 1 Mo par job soumis41MAX_RESULT = 4 * 1024 * 1024 # 4 Mo par résultat conservé42JOB_TYPES = ("http.fetch", "js.run", "sys.info", "ping", "compute.bench", "net.probe")43DEFAULT_TIMEOUT_S = {"http.fetch": 30, "js.run": 60, "sys.info": 10, "ping": 10, "compute.bench": 60, "net.probe": 30}44METRICS_RETENTION_S = 24 * 360045METRIC_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")48HISTORY_DAYS = 74950if not TOKEN:51 print("MOBILE_TOKEN manquant", file=sys.stderr)52 sys.exit(2)5354started_at = time.time()55state_lock = threading.Lock()56queue_cv = threading.Condition(state_lock) # réveil des pollers (nouveau job)57result_cv = threading.Condition(state_lock) # réveil des clients qui attendent un résultat5859workers = {} # name -> dict (état live : dernier battement, métriques, jobs en cours)60queued = [] # liste ordonnée d'ids de jobs en file61jobs = {} # id -> dict (jobs vivants : queued/leased ; les terminés vivent en base)62recent = {} # id -> dict des jobs terminés récemment (cache 1 h, pour /wait et /jobs)63stats = {"submitted": 0, "done": 0, "failed": 0, "expired": 0, "bytes_in": 0}646566# ---------------------------------------------------------------------------67# Base SQLite (persistance des workers, des jobs et de leurs résultats)68# ---------------------------------------------------------------------------6970def 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 pass85 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 c979899def 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()115116117def 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)128129130def db_metrics_history(name, window_s, bucket_s):131 since = time.time() - window_s132 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]140141142def 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()148149150def 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()158159160def 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()170171172def 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 None179180181def 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]198199200def 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 out220221222def 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"] = None234 out.append(j)235 return out236237238def 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()244245246def 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 }256257258# ---------------------------------------------------------------------------259# Logique de file260# ---------------------------------------------------------------------------261262def _kind_of(model):263 m = (model or "").lower()264 return "ipad" if m.startswith("ipad") else ("iphone" if m.startswith("iphone") else "mobile")265266267def _platform_of(model):268 return "ipados" if _kind_of(model) == "ipad" else "ios"269270271def worker_online(w, now=None):272 now = now or time.time()273 return (w.get("lastSeen") or 0) >= now - ONLINE_S274275276def 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 d288289290def 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 d297298299def 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_alias331 db_insert_job(j)332 with queue_cv:333 jobs[j["id"]] = j334 queued.append(j["id"])335 stats["submitted"] += 1336 queue_cv.notify_all()337 return j338339340def _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 continue346 if j["target"] in (None, "any", name):347 del queued[i]348 return j349 return None350351352def lease_job(name, wait_s):353 deadline = time.time() + wait_s354 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"] += 1361 j["leasedAt"] = now362 j["leaseWorker"] = name363 j["leaseExpires"] = now + j["timeoutS"] + LEASE_GRACE_S364 w = workers.get(name)365 if w is not None:366 w.setdefault("active", []).append(j["id"])367 break368 remaining = deadline - time.time()369 if remaining <= 0:370 return None371 queue_cv.wait(remaining)372 db_update_job(j)373 return j374375376def 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"] = result387 stats["done"] += 1388 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"] = None398 queued.append(j["id"])399 queue_cv.notify_all()400 else:401 j["state"] = "failed"402 stats["failed"] += 1403 j["finished"] = now if j["state"] in ("done", "failed") else None404 j["durationMs"] = duration_ms405 j["worker"] = name406 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) + 1411 elif j["state"] == "failed":412 w["jobsFailed"] = (w.get("jobsFailed") or 0) + 1413 if j["state"] in ("done", "failed"):414 jobs.pop(job_id, None)415 recent[job_id] = j416 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, None421422423def _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"))427428429def cancel_job(job_id):430 with queue_cv:431 j = jobs.get(job_id)432 if not j or j["state"] != "queued":433 return None434 j["state"] = "cancelled"435 j["finished"] = time.time()436 try:437 queued.remove(job_id)438 except ValueError:439 pass440 jobs.pop(job_id, None)441 recent[job_id] = j442 db_update_job(j)443 return j444445446def 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 j453 if job_id not in jobs and not j:454 break455 remaining = deadline - time.time()456 if remaining <= 0:457 return jobs.get(job_id) or j458 result_cv.wait(remaining)459 return db_get_job(job_id)460461462def 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 = 0465 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"] = None479 queued.append(j["id"])480 else:481 j["state"] = "expired"482 j["error"] = "bail expiré %d fois" % j["attempts"]483 j["finished"] = now484 jobs.pop(j["id"], None)485 recent[j["id"]] = j486 stats["expired"] += 1487 expired.append(j)488 if expired:489 queue_cv.notify_all()490 result_cv.notify_all()491 # cache des résultats récents : 1 h492 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 = now502503504# ---------------------------------------------------------------------------505# HTTP506# ---------------------------------------------------------------------------507508class Handler(BaseHTTPRequestHandler):509 server_version = "maclustr-mobile/" + VERSION510 protocol_version = "HTTP/1.1"511512 def log_message(self, fmt, *args):513 pass514515 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)524525 def _auth(self):526 h = self.headers.get("Authorization", "")527 if h == "Bearer " + TOKEN:528 return True529 q = parse_qs(urlparse(self.path).query)530 if (q.get("token") or [""])[0] == TOKEN:531 return True532 self._send(401, {"error": "unauthorized"})533 return False534535 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 None541 raw = self.rfile.read(n) if n else b"{}"542 stats["bytes_in"] += n543 return json.loads(raw or b"{}")544 except Exception as e:545 self._send(400, {"error": "JSON invalide : %s" % e})546 return None547548 # -- 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 return564 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 None573 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"] = now599 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 wname614 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"})637638 # -- POST -------------------------------------------------------------639 def do_POST(self):640 u = urlparse(self.path)641 if not self._auth():642 return643 body = self._json()644 if body is None:645 return646 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 None658 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"] = now665 w["active"] = [] # un (re)démarrage de l'app annule ses baux666 w["ip"] = self.client_address[0]667 workers[name] = w668 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"] = now679 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"] = now684 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 400716 return self._send(code, {"ids": ids, "errors": errors, "count": len(ids)})717 return self._send(404, {"error": "not found"})718719 def do_PATCH(self):720 u = urlparse(self.path)721 if not self._auth():722 return723 body = self._json()724 if body is None:725 return726 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 None734 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"})746747 def do_DELETE(self):748 u = urlparse(self.path)749 if not self._auth():750 return751 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"})768769770def _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)776777778def _clean_name(n):779 n = (n or "").strip()780 return "".join(ch for ch in n if ch.isalnum() or ch in "-_.")[:40]781782783class Server(ThreadingHTTPServer):784 daemon_threads = True785 allow_reuse_address = True786 request_queue_size = 128787788789def main():790 global workers791 os.makedirs(DATA_DIR, exist_ok=True)792 workers = db_load_workers()793 for j in db_recover_jobs():794 jobs[j["id"]] = j795 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()801802803if __name__ == "__main__":804 main()805