SPB Git forge
2commits 1branches 0releases
196.0 KBsize
maindefault branch
1 h agolast push
Python 100%
35.9 KB · 805 lines python
Raw Blame History
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