#!/usr/bin/env python3 """ maclustr-mobile — coordinateur des nœuds mobiles (iPhone / iPad) du cluster MacLustr. Les iPhones ne peuvent pas recevoir de SSH ni tourner en arrière-plan : ils *tirent* donc leur travail. L'app MacLustr iOS (onglet « Nœud ») s'enregistre ici, envoie un battement toutes les ~20 s (batterie, thermique, réseau, mémoire) et fait du long-poll sur /api/worker/next ; chaque job exécuté est renvoyé sur /api/worker/result. Côté cluster, on soumet des jobs (HTTP fetch distribué, code JavaScript, infos système) via /api/jobs et on lit les résultats. L'agent maclustr-agentd relit /api/workers pour afficher les mobiles dans les apps MacLustr macOS / iOS. Python 3.9 stdlib uniquement (tourne avec /usr/bin/python3 sous PM2). Aucun accès réseau sortant : seulement du trafic entrant des iPhones (Tailscale ou LAN) et des clients. Variables d'environnement : MOBILE_PORT (9320) MOBILE_TOKEN (obligatoire) MOBILE_DATA (./data) MOBILE_ONLINE_S (45) battement plus vieux → worker « hors ligne » MOBILE_LEASE_GRACE_S (30) marge ajoutée au timeout d'un job avant de le re-mettre en file """ import json import os import sqlite3 import sys import threading import time import uuid from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from urllib.parse import urlparse, parse_qs 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) PORT = int(os.environ.get("MOBILE_PORT", "9320")) TOKEN = os.environ.get("MOBILE_TOKEN", "") DATA_DIR = os.path.abspath(os.environ.get("MOBILE_DATA", os.path.join(os.path.dirname(os.path.abspath(__file__)), "..", "data"))) DB_PATH = os.path.join(DATA_DIR, "mobile.db") ONLINE_S = float(os.environ.get("MOBILE_ONLINE_S", "45")) LEASE_GRACE_S = float(os.environ.get("MOBILE_LEASE_GRACE_S", "30")) MAX_WAIT_S = 30 # borne du long-poll MAX_PAYLOAD = 1024 * 1024 # 1 Mo par job soumis MAX_RESULT = 4 * 1024 * 1024 # 4 Mo par résultat conservé JOB_TYPES = ("http.fetch", "js.run", "sys.info", "ping", "compute.bench", "net.probe") DEFAULT_TIMEOUT_S = {"http.fetch": 30, "js.run": 60, "sys.info": 10, "ping": 10, "compute.bench": 60, "net.probe": 30} METRICS_RETENTION_S = 24 * 3600 METRIC_KEYS = ("battery", "charging", "thermal", "lowPower", "memUsedMb", "memTotalMb", "memAvailableMb", "diskFreeGb", "diskTotalGb", "net", "cpuLoad", "cpuUsage", "uptimeS", "concurrency", "jobsDone", "jobsFailed", "activeJobs", "screenOn", "workerUptimeS", "bytesFetched", "lastJobAt", "ip", "ssid", "brightness") HISTORY_DAYS = 7 if not TOKEN: print("MOBILE_TOKEN manquant", file=sys.stderr) sys.exit(2) started_at = time.time() state_lock = threading.Lock() queue_cv = threading.Condition(state_lock) # réveil des pollers (nouveau job) result_cv = threading.Condition(state_lock) # réveil des clients qui attendent un résultat workers = {} # name -> dict (état live : dernier battement, métriques, jobs en cours) queued = [] # liste ordonnée d'ids de jobs en file jobs = {} # id -> dict (jobs vivants : queued/leased ; les terminés vivent en base) recent = {} # id -> dict des jobs terminés récemment (cache 1 h, pour /wait et /jobs) stats = {"submitted": 0, "done": 0, "failed": 0, "expired": 0, "bytes_in": 0} # --------------------------------------------------------------------------- # Base SQLite (persistance des workers, des jobs et de leurs résultats) # --------------------------------------------------------------------------- def db(): os.makedirs(DATA_DIR, exist_ok=True) c = sqlite3.connect(DB_PATH, timeout=10) c.execute("PRAGMA journal_mode=WAL") c.execute("""CREATE TABLE IF NOT EXISTS workers ( name TEXT PRIMARY KEY, identifier TEXT, model TEXT, marketing TEXT, os TEXT, chip TEXT, ram_mb INTEGER, storage_gb REAL, app_version TEXT, first_seen REAL, last_seen REAL, jobs_done INTEGER DEFAULT 0, jobs_failed INTEGER DEFAULT 0, meta TEXT)""") for col, typ in (("alias", "TEXT"), ("pinned", "INTEGER DEFAULT 0"), ("role", "TEXT"), ("kind", "TEXT"), ("platform", "TEXT"), ("cores", "INTEGER"), ("bench", "TEXT"), ("notes", "TEXT")): try: c.execute("ALTER TABLE workers ADD COLUMN %s %s" % (col, typ)) except sqlite3.OperationalError: pass c.execute("""CREATE TABLE IF NOT EXISTS worker_metrics ( ts REAL, name TEXT, battery REAL, charging INTEGER, thermal TEXT, mem_used_mb INTEGER, mem_avail_mb INTEGER, cpu REAL, net TEXT, active INTEGER, disk_free_gb REAL)""") c.execute("CREATE INDEX IF NOT EXISTS wm_name_ts ON worker_metrics(name, ts)") c.execute("""CREATE TABLE IF NOT EXISTS jobs ( id TEXT PRIMARY KEY, type TEXT, payload TEXT, target TEXT, tag TEXT, state TEXT, attempts INTEGER DEFAULT 0, max_attempts INTEGER DEFAULT 3, timeout_s REAL, created REAL, leased_at REAL, lease_worker TEXT, lease_expires REAL, finished REAL, duration_ms INTEGER, worker TEXT, result TEXT, error TEXT, submitted_by TEXT)""") c.execute("CREATE INDEX IF NOT EXISTS jobs_state ON jobs(state, created)") c.execute("CREATE INDEX IF NOT EXISTS jobs_tag ON jobs(tag)") return c def db_upsert_worker(w): c = db() with c: c.execute("""INSERT INTO workers(name, identifier, model, marketing, os, chip, ram_mb, storage_gb, app_version, first_seen, last_seen, meta, alias, pinned, role, kind, platform, cores, bench, notes) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(name) DO UPDATE SET identifier=excluded.identifier, model=excluded.model, marketing=excluded.marketing, os=excluded.os, chip=excluded.chip, ram_mb=excluded.ram_mb, storage_gb=excluded.storage_gb, app_version=excluded.app_version, last_seen=excluded.last_seen, meta=excluded.meta, alias=excluded.alias, pinned=excluded.pinned, role=excluded.role, kind=excluded.kind, platform=excluded.platform, cores=excluded.cores, bench=excluded.bench, notes=excluded.notes""", (w["name"], w.get("identifier"), w.get("model"), w.get("marketing"), w.get("os"), w.get("chip"), w.get("ramMb"), w.get("storageGb"), w.get("appVersion"), w.get("firstSeen"), w.get("lastSeen"), json.dumps(w.get("meta") or {}), w.get("alias"), 1 if w.get("pinned") else 0, w.get("role") or "worker", w.get("kind"), w.get("platform"), w.get("cores"), json.dumps(w.get("bench")) if w.get("bench") else None, w.get("notes"))) c.close() def db_record_metrics(name, ts, m): try: c = db() with c: c.execute("INSERT INTO worker_metrics VALUES (?,?,?,?,?,?,?,?,?,?,?)", (ts, name, m.get("battery"), 1 if m.get("charging") else 0, m.get("thermal"), m.get("memUsedMb"), m.get("memAvailableMb"), m.get("cpuUsage", m.get("cpuLoad")), m.get("net"), m.get("activeJobs"), m.get("diskFreeGb"))) c.execute("DELETE FROM worker_metrics WHERE ts < ?", (ts - METRICS_RETENTION_S,)) c.close() except Exception as e: print("metrics db:", e, flush=True) def db_metrics_history(name, window_s, bucket_s): since = time.time() - window_s c = db() rows = c.execute( "SELECT (CAST(ts AS INTEGER)/%d)*%d AS b, AVG(battery), AVG(mem_used_mb), AVG(cpu), MAX(active), MAX(charging)" " FROM worker_metrics WHERE name=? AND ts>=? GROUP BY b ORDER BY b" % (bucket_s, bucket_s), (name, since)).fetchall() c.close() return [{"ts": r[0], "battery": round(r[1], 1) if r[1] is not None else None, "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, "active": r[4] or 0, "charging": bool(r[5])} for r in rows] def db_touch_worker(name, ts, done_inc=0, failed_inc=0): c = db() with c: c.execute("UPDATE workers SET last_seen=?, jobs_done=jobs_done+?, jobs_failed=jobs_failed+? WHERE name=?", (ts, done_inc, failed_inc, name)) c.close() def db_insert_job(j): c = db() with c: c.execute("""INSERT INTO jobs(id, type, payload, target, tag, state, attempts, max_attempts, timeout_s, created, submitted_by) VALUES(?,?,?,?,?,?,?,?,?,?,?)""", (j["id"], j["type"], json.dumps(j["payload"], ensure_ascii=False), j.get("target"), j.get("tag"), j["state"], j["attempts"], j["maxAttempts"], j["timeoutS"], j["created"], j.get("submittedBy"))) c.close() def db_update_job(j): c = db() with c: c.execute("""UPDATE jobs SET state=?, attempts=?, leased_at=?, lease_worker=?, lease_expires=?, finished=?, duration_ms=?, worker=?, result=?, error=? WHERE id=?""", (j["state"], j["attempts"], j.get("leasedAt"), j.get("leaseWorker"), j.get("leaseExpires"), j.get("finished"), j.get("durationMs"), j.get("worker"), json.dumps(j.get("result"), ensure_ascii=False) if j.get("result") is not None else None, j.get("error"), j["id"])) c.close() def db_get_job(job_id): c = db() cur = c.execute("SELECT * FROM jobs WHERE id=?", (job_id,)) row = cur.fetchone() cols = [d[0] for d in cur.description] if row else [] c.close() return row_to_job(dict(zip(cols, row))) if row else None def db_list_jobs(state=None, tag=None, limit=100): c = db() sql = "SELECT * FROM jobs" cond, args = [], [] if state: cond.append("state=?"); args.append(state) if tag: cond.append("tag=?"); args.append(tag) if cond: sql += " WHERE " + " AND ".join(cond) sql += " ORDER BY created DESC LIMIT ?" args.append(int(limit)) cur = c.execute(sql, args) rows = cur.fetchall() cols = [d[0] for d in cur.description] c.close() return [row_to_job(dict(zip(cols, r))) for r in rows] def db_load_workers(): c = db() cur = c.execute("SELECT * FROM workers") rows = cur.fetchall() cols = [d[0] for d in cur.description] c.close() out = {} for r in rows: d = dict(zip(cols, r)) out[d["name"]] = { "name": d["name"], "identifier": d["identifier"], "model": d["model"], "marketing": d["marketing"], "os": d["os"], "chip": d["chip"], "ramMb": d["ram_mb"], "storageGb": d["storage_gb"], "appVersion": d["app_version"], "firstSeen": d["first_seen"], "lastSeen": d["last_seen"], "jobsDone": d["jobs_done"] or 0, "jobsFailed": d["jobs_failed"] or 0, "meta": json.loads(d["meta"] or "{}"), "metrics": {}, "active": [], "alias": d.get("alias"), "pinned": bool(d.get("pinned")), "role": d.get("role") or "worker", "kind": d.get("kind") or _kind_of(d["model"]), "platform": d.get("platform") or _platform_of(d["model"]), "cores": d.get("cores"), "bench": json.loads(d["bench"]) if d.get("bench") else None, "notes": d.get("notes"), } return out def db_recover_jobs(): """Au démarrage : les jobs queued/leased reviennent en file.""" c = db() cur = c.execute("SELECT * FROM jobs WHERE state IN ('queued','leased') ORDER BY created") rows = cur.fetchall() cols = [d[0] for d in cur.description] c.close() out = [] for r in rows: j = row_to_job(dict(zip(cols, r))) j["state"] = "queued" j["leasedAt"] = j["leaseWorker"] = j["leaseExpires"] = None out.append(j) return out def db_purge(): c = db() with c: c.execute("DELETE FROM jobs WHERE state IN ('done','failed','expired','cancelled') AND created < ?", (time.time() - HISTORY_DAYS * 86400,)) c.close() def row_to_job(d): return { "id": d["id"], "type": d["type"], "payload": json.loads(d["payload"] or "{}"), "target": d["target"], "tag": d["tag"], "state": d["state"], "attempts": d["attempts"] or 0, "maxAttempts": d["max_attempts"] or 3, "timeoutS": d["timeout_s"], "created": d["created"], "leasedAt": d["leased_at"], "leaseWorker": d["lease_worker"], "leaseExpires": d["lease_expires"], "finished": d["finished"], "durationMs": d["duration_ms"], "worker": d["worker"], "result": json.loads(d["result"]) if d.get("result") else None, "error": d["error"], "submittedBy": d["submitted_by"], } # --------------------------------------------------------------------------- # Logique de file # --------------------------------------------------------------------------- def _kind_of(model): m = (model or "").lower() return "ipad" if m.startswith("ipad") else ("iphone" if m.startswith("iphone") else "mobile") def _platform_of(model): return "ipados" if _kind_of(model) == "ipad" else "ios" def worker_online(w, now=None): now = now or time.time() return (w.get("lastSeen") or 0) >= now - ONLINE_S def public_worker(w, now=None): now = now or time.time() d = {k: v for k, v in w.items() if k not in ("identifier",)} d["online"] = worker_online(w, now) d["ageS"] = round(now - (w.get("lastSeen") or 0), 1) d["activeJobs"] = len(w.get("active") or []) d["displayName"] = w.get("alias") or w["name"] d.setdefault("kind", _kind_of(w.get("model"))) d.setdefault("platform", _platform_of(w.get("model"))) d.setdefault("role", "worker") d.setdefault("pinned", False) return d def public_job(j, with_result=True): d = dict(j) if not with_result: d.pop("result", None) if isinstance(j.get("payload"), dict) and "code" in j["payload"]: d["payload"] = dict(j["payload"], code="<%d octets>" % len(j["payload"].get("code") or "")) return d def submit_job(spec, submitted_by=None): t = spec.get("type") if t not in JOB_TYPES: raise ValueError("type inconnu : %r (attendu : %s)" % (t, ", ".join(JOB_TYPES))) payload = spec.get("payload") or {} if not isinstance(payload, dict): raise ValueError("payload doit être un objet") if t == "http.fetch" and not payload.get("url"): raise ValueError("http.fetch : payload.url requis") if t == "js.run" and not payload.get("code"): raise ValueError("js.run : payload.code requis (définir function main(args) { … })") if t == "net.probe" and not payload.get("hosts"): raise ValueError("net.probe : payload.hosts (liste d'hôtes ou d'URL) requis") if len(json.dumps(payload)) > MAX_PAYLOAD: raise ValueError("payload trop gros (> 1 Mo)") timeout_s = float(spec.get("timeoutS") or DEFAULT_TIMEOUT_S[t]) timeout_s = max(2.0, min(timeout_s, 600.0)) j = { "id": spec.get("id") or uuid.uuid4().hex[:12], "type": t, "payload": payload, "target": (spec.get("target") or None), "tag": spec.get("tag"), "state": "queued", "attempts": 0, "maxAttempts": int(spec.get("maxAttempts") or 3), "timeoutS": timeout_s, "created": time.time(), "leasedAt": None, "leaseWorker": None, "leaseExpires": None, "finished": None, "durationMs": None, "worker": None, "result": None, "error": None, "submittedBy": submitted_by, } if j["target"] and j["target"] not in ("any",): with state_lock: if j["target"] not in workers: by_alias = next((n for n, w in workers.items() if (w.get("alias") or "").lower() == str(j["target"]).lower()), None) if not by_alias: raise ValueError("worker cible inconnu : %s" % j["target"]) j["target"] = by_alias db_insert_job(j) with queue_cv: jobs[j["id"]] = j queued.append(j["id"]) stats["submitted"] += 1 queue_cv.notify_all() return j def _pick_job_for(name): """Premier job en file compatible avec ce worker (cible = ce worker ou n'importe qui).""" for i, jid in enumerate(queued): j = jobs.get(jid) if not j: continue if j["target"] in (None, "any", name): del queued[i] return j return None def lease_job(name, wait_s): deadline = time.time() + wait_s with queue_cv: while True: j = _pick_job_for(name) if j: now = time.time() j["state"] = "leased" j["attempts"] += 1 j["leasedAt"] = now j["leaseWorker"] = name j["leaseExpires"] = now + j["timeoutS"] + LEASE_GRACE_S w = workers.get(name) if w is not None: w.setdefault("active", []).append(j["id"]) break remaining = deadline - time.time() if remaining <= 0: return None queue_cv.wait(remaining) db_update_job(j) return j def finish_job(name, job_id, ok, result, error, duration_ms): with result_cv: j = jobs.get(job_id) if not j: return None, "job inconnu ou déjà clos" if j["state"] != "leased" or j["leaseWorker"] != name: return None, "job non loué par %s (état %s, loué à %s)" % (name, j["state"], j["leaseWorker"]) now = time.time() if ok: j["state"] = "done" j["result"] = result stats["done"] += 1 if j["type"] == "compute.bench" and isinstance(result, dict): wb = workers.get(name) if wb is not None: wb["bench"] = {"score": result.get("score"), "singleCore": result.get("singleCore"), "multiCore": result.get("multiCore"), "ts": now} else: j["error"] = (error or "erreur inconnue")[:2000] if j["attempts"] < j["maxAttempts"] and _retryable(error): j["state"] = "queued" j["leasedAt"] = j["leaseWorker"] = j["leaseExpires"] = None queued.append(j["id"]) queue_cv.notify_all() else: j["state"] = "failed" stats["failed"] += 1 j["finished"] = now if j["state"] in ("done", "failed") else None j["durationMs"] = duration_ms j["worker"] = name w = workers.get(name) if w is not None: w["active"] = [x for x in (w.get("active") or []) if x != job_id] if j["state"] == "done": w["jobsDone"] = (w.get("jobsDone") or 0) + 1 elif j["state"] == "failed": w["jobsFailed"] = (w.get("jobsFailed") or 0) + 1 if j["state"] in ("done", "failed"): jobs.pop(job_id, None) recent[job_id] = j result_cv.notify_all() db_update_job(j) if j["state"] in ("done", "failed"): db_touch_worker(name, time.time(), 1 if j["state"] == "done" else 0, 1 if j["state"] == "failed" else 0) return j, None def _retryable(error): e = (error or "").lower() # Une erreur « métier » (JS levée par le code, HTTP 4xx) n'est pas rejouée ; réseau/timeout oui. return any(k in e for k in ("timeout", "timed out", "network", "connexion", "connection", "réseau", "cancelled", "-1001", "-1005", "-1009")) def cancel_job(job_id): with queue_cv: j = jobs.get(job_id) if not j or j["state"] != "queued": return None j["state"] = "cancelled" j["finished"] = time.time() try: queued.remove(job_id) except ValueError: pass jobs.pop(job_id, None) recent[job_id] = j db_update_job(j) return j def wait_job(job_id, timeout_s): deadline = time.time() + min(max(timeout_s, 0), MAX_WAIT_S) with result_cv: while True: j = recent.get(job_id) if j and j["state"] in ("done", "failed", "cancelled", "expired"): return j if job_id not in jobs and not j: break remaining = deadline - time.time() if remaining <= 0: return jobs.get(job_id) or j result_cv.wait(remaining) return db_get_job(job_id) def reaper_loop(): """Re-met en file les jobs dont le bail a expiré (iPhone parti en arrière-plan, app fermée…).""" last_purge = 0 while True: time.sleep(5) now = time.time() expired = [] with queue_cv: for j in list(jobs.values()): if j["state"] == "leased" and j.get("leaseExpires") and j["leaseExpires"] < now: w = workers.get(j["leaseWorker"] or "") if w is not None: w["active"] = [x for x in (w.get("active") or []) if x != j["id"]] if j["attempts"] < j["maxAttempts"]: j["state"] = "queued" j["error"] = "bail expiré sur %s" % j["leaseWorker"] j["leasedAt"] = j["leaseWorker"] = j["leaseExpires"] = None queued.append(j["id"]) else: j["state"] = "expired" j["error"] = "bail expiré %d fois" % j["attempts"] j["finished"] = now jobs.pop(j["id"], None) recent[j["id"]] = j stats["expired"] += 1 expired.append(j) if expired: queue_cv.notify_all() result_cv.notify_all() # cache des résultats récents : 1 h for jid in [k for k, v in recent.items() if (v.get("finished") or 0) < now - 3600]: recent.pop(jid, None) for j in expired: db_update_job(j) if now - last_purge > 3600: try: db_purge() except Exception as e: print("purge:", e, flush=True) last_purge = now # --------------------------------------------------------------------------- # HTTP # --------------------------------------------------------------------------- class Handler(BaseHTTPRequestHandler): server_version = "maclustr-mobile/" + VERSION protocol_version = "HTTP/1.1" def log_message(self, fmt, *args): pass def _send(self, code, payload=None): body = b"" if payload is None else json.dumps(payload, ensure_ascii=False).encode() self.send_response(code) self.send_header("Content-Type", "application/json; charset=utf-8") self.send_header("Content-Length", str(len(body))) self.send_header("Cache-Control", "no-store") self.end_headers() if body: self.wfile.write(body) def _auth(self): h = self.headers.get("Authorization", "") if h == "Bearer " + TOKEN: return True q = parse_qs(urlparse(self.path).query) if (q.get("token") or [""])[0] == TOKEN: return True self._send(401, {"error": "unauthorized"}) return False def _json(self): try: n = int(self.headers.get("Content-Length", 0)) if n > MAX_RESULT + MAX_PAYLOAD: self._send(413, {"error": "corps trop gros"}) return None raw = self.rfile.read(n) if n else b"{}" stats["bytes_in"] += n return json.loads(raw or b"{}") except Exception as e: self._send(400, {"error": "JSON invalide : %s" % e}) return None # -- GET -------------------------------------------------------------- def do_GET(self): u = urlparse(self.path) q = parse_qs(u.query) parts = [p for p in u.path.split("/") if p] now = time.time() if u.path == "/health": with state_lock: online = sum(1 for w in workers.values() if worker_online(w, now)) active = sum(len(w.get("active") or []) for w in workers.values()) nq = len(queued) return self._send(200, {"ok": True, "version": VERSION, "uptime": int(now - started_at), "workersOnline": online, "workersTotal": len(workers), "queued": nq, "active": active, "stats": stats}) if not self._auth(): return if u.path == "/api/workers": with state_lock: lst = [public_worker(w, now) for w in workers.values()] lst.sort(key=lambda w: (not w["online"], w["name"])) return self._send(200, {"workers": lst, "ts": now, "onlineS": ONLINE_S}) if len(parts) == 3 and parts[0] == "api" and parts[1] == "workers": with state_lock: w = _find_worker(parts[2]) pub = public_worker(w, now) if w else None if not pub: return self._send(404, {"error": "worker inconnu"}) 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] pub["history"] = db_metrics_history(w["name"], 3600, 60) return self._send(200, pub) if len(parts) == 4 and parts[0] == "api" and parts[1] == "workers" and parts[3] == "history": with state_lock: w = _find_worker(parts[2]) if not w: return self._send(404, {"error": "worker inconnu"}) window = {"1h": 3600, "6h": 21600, "24h": 86400}.get((q.get("window") or ["1h"])[0], 3600) bucket = {3600: 60, 21600: 300, 86400: 900}[window] return self._send(200, {"name": w["name"], "points": db_metrics_history(w["name"], window, bucket), "window": window}) if u.path == "/api/stats": with state_lock: return self._send(200, {"stats": stats, "queued": len(queued), "leased": sum(1 for j in jobs.values() if j["state"] == "leased"), "workers": {n: {"online": worker_online(w, now), "active": len(w.get("active") or [])} for n, w in workers.items()}, "uptime": int(now - started_at)}) if u.path == "/api/worker/next": name = (q.get("worker") or [""])[0] wait = min(float((q.get("wait") or ["25"])[0]), MAX_WAIT_S) with state_lock: w = workers.get(name) if w is None: return self._send(409, {"error": "worker non enregistré : appeler /api/worker/register"}) w["lastSeen"] = now j = lease_job(name, wait) if not j: return self._send(204) return self._send(200, {"id": j["id"], "type": j["type"], "payload": j["payload"], "timeoutS": j["timeoutS"], "attempt": j["attempts"], "tag": j.get("tag")}) if u.path == "/api/jobs": state = (q.get("state") or [None])[0] tag = (q.get("tag") or [None])[0] wname = (q.get("worker") or [None])[0] limit = int((q.get("limit") or ["100"])[0]) full = (q.get("full") or ["0"])[0] == "1" with state_lock: if wname: ww = _find_worker(wname) wname = ww["name"] if ww else wname 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) and (not wname or j.get("leaseWorker") == wname or j.get("worker") == wname or j.get("target") == wname)] if state in (None, "done", "failed", "expired", "cancelled"): stored = db_list_jobs(state, tag, limit if not wname else max(limit * 5, 300)) seen = {j["id"] for j in live} 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)] live.sort(key=lambda j: j["created"], reverse=True) return self._send(200, {"jobs": live[:limit], "count": len(live)}) if len(parts) == 3 and parts[0] == "api" and parts[1] == "jobs": with state_lock: j = jobs.get(parts[2]) or recent.get(parts[2]) if not j: j = db_get_job(parts[2]) if not j: return self._send(404, {"error": "job inconnu"}) return self._send(200, j) if len(parts) == 4 and parts[0] == "api" and parts[1] == "jobs" and parts[3] == "wait": timeout = float((q.get("timeout") or ["25"])[0]) j = wait_job(parts[2], timeout) if not j: return self._send(404, {"error": "job inconnu"}) return self._send(200, j) return self._send(404, {"error": "not found"}) # -- POST ------------------------------------------------------------- def do_POST(self): u = urlparse(self.path) if not self._auth(): return body = self._json() if body is None: return now = time.time() if u.path == "/api/worker/register": name = _clean_name(body.get("name")) if not name: return self._send(400, {"error": "name requis"}) with state_lock: w = workers.get(name) or {"name": name, "firstSeen": now, "jobsDone": 0, "jobsFailed": 0, "active": [], "metrics": {}} for k in ("identifier", "model", "marketing", "os", "chip", "ramMb", "storageGb", "appVersion", "cores", "kind", "platform", "role"): if body.get(k) is not None: w[k] = body[k] if body.get("alias") is not None: w["alias"] = _clean_name(body["alias"]) or None if body.get("pinned") is not None: w["pinned"] = bool(body["pinned"]) w.setdefault("kind", _kind_of(w.get("model"))) w.setdefault("platform", _platform_of(w.get("model"))) w.setdefault("role", "worker") w["meta"] = body.get("meta") or w.get("meta") or {} w["lastSeen"] = now w["active"] = [] # un (re)démarrage de l'app annule ses baux w["ip"] = self.client_address[0] workers[name] = w db_upsert_worker(w) print("register %s (%s, %s) depuis %s" % (name, w.get("marketing") or w.get("model"), w.get("os"), w["ip"]), flush=True) return self._send(200, {"ok": True, "name": name, "serverTime": now, "version": VERSION, "config": {"heartbeatS": 20, "pollWaitS": 25}}) if u.path == "/api/worker/heartbeat": name = _clean_name(body.get("name")) with state_lock: w = workers.get(name) if w is None: return self._send(409, {"error": "worker non enregistré"}) w["lastSeen"] = now w["ip"] = self.client_address[0] m = body.get("metrics") or {} if isinstance(m, dict): w["metrics"] = {k: m[k] for k in m if k in METRIC_KEYS} w["metrics"]["ts"] = now if body.get("appVersion"): w["appVersion"] = body["appVersion"] nq = len(queued) metrics_copy = dict(w.get("metrics") or {}) db_touch_worker(name, now) if metrics_copy: db_record_metrics(name, now, metrics_copy) return self._send(200, {"ok": True, "serverTime": now, "queued": nq, "config": {"heartbeatS": 20, "pollWaitS": 25}}) if u.path == "/api/worker/result": name = _clean_name(body.get("name")) job_id = body.get("jobId") result = body.get("result") if result is not None and len(json.dumps(result)) > MAX_RESULT: result = {"truncated": True, "note": "résultat > 4 Mo, tronqué côté coordinateur"} j, err = finish_job(name, job_id, bool(body.get("ok")), result, body.get("error"), body.get("durationMs")) if err: return self._send(409, {"error": err}) return self._send(200, {"ok": True, "state": j["state"]}) if u.path == "/api/jobs": specs = body.get("jobs") if isinstance(body, dict) and "jobs" in body else [body] if not isinstance(specs, list) or not specs: return self._send(400, {"error": "attendu un job ou {\"jobs\": [...]}"}) if len(specs) > 5000: return self._send(400, {"error": "max 5000 jobs par lot"}) ids, errors = [], [] who = self.client_address[0] for s in specs: try: ids.append(submit_job(s, who)["id"]) except ValueError as e: errors.append(str(e)) code = 200 if ids else 400 return self._send(code, {"ids": ids, "errors": errors, "count": len(ids)}) return self._send(404, {"error": "not found"}) def do_PATCH(self): u = urlparse(self.path) if not self._auth(): return body = self._json() if body is None: return parts = [p for p in u.path.split("/") if p] if len(parts) == 3 and parts[0] == "api" and parts[1] == "workers": with state_lock: w = _find_worker(parts[2]) if not w: return self._send(404, {"error": "worker inconnu"}) if "alias" in body: w["alias"] = _clean_name(body["alias"]) or None if "pinned" in body: w["pinned"] = bool(body["pinned"]) if "role" in body: w["role"] = str(body["role"])[:40] if "notes" in body: w["notes"] = str(body["notes"])[:500] if isinstance(body.get("meta"), dict): w["meta"] = dict(w.get("meta") or {}, **body["meta"]) snap = dict(w) db_upsert_worker(snap) return self._send(200, public_worker(snap)) return self._send(404, {"error": "not found"}) def do_DELETE(self): u = urlparse(self.path) if not self._auth(): return parts = [p for p in u.path.split("/") if p] if len(parts) == 3 and parts[0] == "api" and parts[1] == "jobs": j = cancel_job(parts[2]) if not j: return self._send(409, {"error": "job absent ou plus en file"}) return self._send(200, j) if len(parts) == 3 and parts[0] == "api" and parts[1] == "workers": with state_lock: w = workers.pop(parts[2], None) if not w: return self._send(404, {"error": "worker inconnu"}) c = db() with c: c.execute("DELETE FROM workers WHERE name=?", (parts[2],)) c.close() return self._send(200, {"ok": True}) return self._send(404, {"error": "not found"}) def _find_worker(key): """Worker par nom ou par alias (insensible à la casse). À appeler sous state_lock.""" if key in workers: return workers[key] k = (key or "").lower() return next((w for w in workers.values() if (w.get("alias") or "").lower() == k or w["name"].lower() == k), None) def _clean_name(n): n = (n or "").strip() return "".join(ch for ch in n if ch.isalnum() or ch in "-_.")[:40] class Server(ThreadingHTTPServer): daemon_threads = True allow_reuse_address = True request_queue_size = 128 def main(): global workers os.makedirs(DATA_DIR, exist_ok=True) workers = db_load_workers() for j in db_recover_jobs(): jobs[j["id"]] = j queued.append(j["id"]) db_update_job(j) threading.Thread(target=reaper_loop, daemon=True).start() srv = Server(("0.0.0.0", PORT), Handler) print("maclustr-mobile %s : %d worker(s) connus, %d job(s) repris, écoute :%d" % (VERSION, len(workers), len(queued), PORT), flush=True) srv.serve_forever() if __name__ == "__main__": main()