#!/usr/bin/env python3 """ mlmobile — client en ligne de commande du coordinateur maclustr-mobile (iPhones du cluster). mlmobile nodes iPhones connus (en ligne, batterie, thermique, jobs) mlmobile stats file, baux, compteurs mlmobile ping [--on NOM] aller-retour vers un iPhone mlmobile info [--on NOM] infos système d'un iPhone (sys.info) mlmobile fetch URL [--on NOM] [--method M] [--header K:V]… [--body TXT] [--timeout S] [--raw] mlmobile run script.js [--args JSON] [--on NOM] [--timeout S] mlmobile map script.js items.json [--tag T] [--timeout S] [--concurrency N] 1 job par élément, résultats agrégés (JSON sur stdout) mlmobile fetch-list urls.txt [--tag T] [--out results.jsonl] 1 http.fetch par URL, réparti sur tous les iPhones mlmobile job ID | mlmobile wait ID | mlmobile jobs [--state S] [--tag T] [--limit N] | mlmobile cancel ID Le script JS doit définir `function main(args)` et retourner une valeur JSON. API disponible dans le script : ml.http(url, {method, headers, body, timeout}) -> {status, headers, body} (synchrone) ml.log(msg) ml.sleep(ms) ml.device -> {name, model, os, battery, …} Env : MLMOBILE_URL (défaut http://M4M64a.maclustr.io:9320), MLMOBILE_TOKEN. """ import argparse import json import os import sys import time import urllib.error import urllib.request URL = os.environ.get("MLMOBILE_URL", "http://M4M64a.maclustr.io:9320").rstrip("/") TOKEN = os.environ.get("MLMOBILE_TOKEN", "mlmob_4f9a1c7e2b8d6035a9e1c4b7d2f8036e5a1b9c3d") def api(method, path, body=None, timeout=40): data = json.dumps(body).encode() if body is not None else None req = urllib.request.Request(URL + path, data=data, method=method, headers={"Authorization": "Bearer " + TOKEN, "Content-Type": "application/json"}) try: with urllib.request.urlopen(req, timeout=timeout) as r: raw = r.read() return r.status, (json.loads(raw) if raw else None) except urllib.error.HTTPError as e: raw = e.read() try: return e.code, json.loads(raw) except Exception: return e.code, {"error": raw.decode(errors="replace")} except Exception as e: sys.exit("coordinateur injoignable (%s) : %s" % (URL, e)) def fmt_age(s): if s is None: return "—" if s < 90: return "%ds" % s if s < 5400: return "%dmin" % (s // 60) return "%.1fh" % (s / 3600) def cmd_nodes(a): _, d = api("GET", "/api/workers") ws = d.get("workers", []) if not ws: print("aucun iPhone enregistré (ouvrir MacLustr → Nœud → Mode nœud)") return print("%-14s %-6s %-22s %-8s %-6s %-9s %-8s %-6s %s" % ("nom", "état", "modèle", "iOS", "batt.", "therm.", "réseau", "actifs", "faits/échecs")) for w in ws: m = w.get("metrics") or {} batt = m.get("battery") batt_s = ("%d%%%s" % (round(batt * 100) if batt is not None and batt <= 1 else (batt or 0), "⚡" if m.get("charging") else "")) if batt is not None else "—" print("%-14s %-6s %-22s %-8s %-6s %-9s %-8s %-6s %d/%d vu il y a %s" % ( w["name"], "ON" if w["online"] else "off", (w.get("marketing") or w.get("model") or "")[:22], (w.get("os") or "")[:8], batt_s, (m.get("thermal") or "—")[:9], (m.get("net") or "—")[:8], w.get("activeJobs", 0), w.get("jobsDone", 0), w.get("jobsFailed", 0), fmt_age(w.get("ageS")))) def cmd_stats(a): _, d = api("GET", "/api/stats") print(json.dumps(d, indent=2, ensure_ascii=False)) def submit_and_wait(spec, timeout): code, d = api("POST", "/api/jobs", spec) if code != 200 or not d.get("ids"): sys.exit("refus : %s" % (d.get("errors") or d)) jid = d["ids"][0] return wait_for(jid, timeout) def wait_for(jid, timeout): t0 = time.time() while True: _, j = api("GET", "/api/jobs/%s/wait?timeout=25" % jid, timeout=40) if j.get("state") in ("done", "failed", "expired", "cancelled"): return j if time.time() - t0 > timeout: return j def print_job(j, raw=False): st = j.get("state") if st == "done": r = j.get("result") if raw and isinstance(r, dict) and "body" in r: sys.stdout.write(r["body"] if isinstance(r["body"], str) else json.dumps(r["body"])) sys.stdout.write("\n") else: print(json.dumps(r, indent=2, ensure_ascii=False)) print("# %s par %s en %s ms" % (j["id"], j.get("worker"), j.get("durationMs")), file=sys.stderr) else: print(json.dumps({k: j.get(k) for k in ("id", "state", "error", "worker", "attempts")}, indent=2, ensure_ascii=False)) if st != "queued": sys.exit(1) print("(toujours en file — aucun iPhone en ligne ?)", file=sys.stderr) def cmd_ping(a): j = submit_and_wait({"type": "ping", "payload": {"t": time.time()}, "target": a.on, "timeoutS": 10}, a.timeout) print_job(j) def cmd_info(a): j = submit_and_wait({"type": "sys.info", "payload": {}, "target": a.on, "timeoutS": 15}, a.timeout) print_job(j) def fetch_payload(url, a): headers = {} for h in (a.header or []): if ":" in h: k, v = h.split(":", 1) headers[k.strip()] = v.strip() p = {"url": url, "method": a.method, "timeoutS": a.timeout, "maxBytes": a.max_bytes} if headers: p["headers"] = headers if getattr(a, "body", None): p["body"] = a.body return p def cmd_fetch(a): j = submit_and_wait({"type": "http.fetch", "payload": fetch_payload(a.url, a), "target": a.on, "timeoutS": a.timeout + 5}, a.timeout + 30) print_job(j, raw=a.raw) def cmd_run(a): code = open(a.script).read() args = json.loads(a.args) if a.args else {} j = submit_and_wait({"type": "js.run", "payload": {"code": code, "args": args}, "target": a.on, "timeoutS": a.timeout}, a.timeout + 30) print_job(j) def run_batch(specs, tag, timeout, out=None): code, d = api("POST", "/api/jobs", {"jobs": specs}, timeout=120) if code != 200: sys.exit("refus : %s" % d) ids = d["ids"] if d.get("errors"): print("erreurs de soumission : %s" % d["errors"][:5], file=sys.stderr) print("%d jobs soumis (tag %s)" % (len(ids), tag), file=sys.stderr) pending = set(ids) results = {} t0 = time.time() fh = open(out, "a") if out else None while pending and time.time() - t0 < timeout: _, d = api("GET", "/api/jobs?tag=%s&limit=10000&full=1" % tag, timeout=60) for j in d.get("jobs", []): if j["id"] in pending and j["state"] in ("done", "failed", "expired", "cancelled"): pending.discard(j["id"]) results[j["id"]] = j if fh: fh.write(json.dumps(j, ensure_ascii=False) + "\n") fh.flush() done = len(results) print("\r%d/%d terminés" % (done, len(ids)), end="", file=sys.stderr) if pending: time.sleep(2) print("", file=sys.stderr) if fh: fh.close() return [results.get(i) or {"id": i, "state": "pending"} for i in ids] def cmd_map(a): code = open(a.script).read() items = json.load(open(a.items)) tag = a.tag or ("map-%d" % int(time.time())) specs = [{"type": "js.run", "payload": {"code": code, "args": it}, "tag": tag, "timeoutS": a.timeout} for it in items] res = run_batch(specs, tag, a.timeout * max(1, len(items)) / max(1, a.concurrency) + 120) print(json.dumps([{"args": items[i], "state": r.get("state"), "worker": r.get("worker"), "result": r.get("result"), "error": r.get("error")} for i, r in enumerate(res)], indent=2, ensure_ascii=False)) def cmd_fetch_list(a): urls = [u.strip() for u in open(a.urls) if u.strip() and not u.startswith("#")] tag = a.tag or ("fetch-%d" % int(time.time())) specs = [{"type": "http.fetch", "payload": fetch_payload(u, a), "tag": tag, "timeoutS": a.timeout + 5} for u in urls] res = run_batch(specs, tag, a.timeout * len(urls) / 3 + 120, out=a.out) ok = sum(1 for r in res if r.get("state") == "done") by_worker = {} for r in res: by_worker[r.get("worker") or "?"] = by_worker.get(r.get("worker") or "?", 0) + 1 print("%d/%d réussis ; répartition : %s%s" % (ok, len(urls), by_worker, " ; résultats dans %s" % a.out if a.out else "")) if not a.out: for u, r in zip(urls, res): rr = r.get("result") or {} print("%-4s %-7s %s (%s, %s o)" % (rr.get("status", "—"), r.get("worker") or "—", u, r.get("state"), rr.get("bytes", "—"))) def cmd_job(a): _, j = api("GET", "/api/jobs/%s" % a.id) print(json.dumps(j, indent=2, ensure_ascii=False)) def cmd_wait(a): print_job(wait_for(a.id, a.timeout)) def cmd_jobs(a): qs = "limit=%d" % a.limit if a.state: qs += "&state=" + a.state if a.tag: qs += "&tag=" + a.tag _, d = api("GET", "/api/jobs?" + qs) for j in d.get("jobs", []): print("%-12s %-9s %-10s %-8s %-10s %s" % (j["id"], j["state"], j["type"], j.get("worker") or "—", time.strftime("%H:%M:%S", time.localtime(j["created"])), (j.get("error") or "")[:60])) def cmd_cancel(a): code, d = api("DELETE", "/api/jobs/%s" % a.id) print(d) def main(): p = argparse.ArgumentParser(prog="mlmobile", description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) sp = p.add_subparsers(dest="cmd", required=True) sp.add_parser("nodes").set_defaults(f=cmd_nodes) sp.add_parser("stats").set_defaults(f=cmd_stats) x = sp.add_parser("ping"); x.add_argument("--on"); x.add_argument("--timeout", type=float, default=30); x.set_defaults(f=cmd_ping) x = sp.add_parser("info"); x.add_argument("--on"); x.add_argument("--timeout", type=float, default=30); x.set_defaults(f=cmd_info) x = sp.add_parser("fetch"); x.add_argument("url"); x.add_argument("--on"); x.add_argument("--method", default="GET") x.add_argument("--header", action="append"); x.add_argument("--body"); x.add_argument("--timeout", type=float, default=30) x.add_argument("--max-bytes", type=int, default=2 * 1024 * 1024); x.add_argument("--raw", action="store_true"); x.set_defaults(f=cmd_fetch) x = sp.add_parser("run"); x.add_argument("script"); x.add_argument("--args"); x.add_argument("--on"); x.add_argument("--timeout", type=float, default=60); x.set_defaults(f=cmd_run) x = sp.add_parser("map"); x.add_argument("script"); x.add_argument("items"); x.add_argument("--tag"); x.add_argument("--timeout", type=float, default=60) x.add_argument("--concurrency", type=int, default=6); x.set_defaults(f=cmd_map) x = sp.add_parser("fetch-list"); x.add_argument("urls"); x.add_argument("--tag"); x.add_argument("--out"); x.add_argument("--method", default="GET") x.add_argument("--header", action="append"); x.add_argument("--timeout", type=float, default=30); x.add_argument("--max-bytes", type=int, default=2 * 1024 * 1024); x.set_defaults(f=cmd_fetch_list) x = sp.add_parser("job"); x.add_argument("id"); x.set_defaults(f=cmd_job) x = sp.add_parser("wait"); x.add_argument("id"); x.add_argument("--timeout", type=float, default=120); x.set_defaults(f=cmd_wait) x = sp.add_parser("jobs"); x.add_argument("--state"); x.add_argument("--tag"); x.add_argument("--limit", type=int, default=50); x.set_defaults(f=cmd_jobs) x = sp.add_parser("cancel"); x.add_argument("id"); x.set_defaults(f=cmd_cancel) a = p.parse_args() a.f(a) if __name__ == "__main__": main()