spb/maclustr-mobile
Public
Python 100%
1#!/usr/bin/env python32"""3mlmobile — client en ligne de commande du coordinateur maclustr-mobile (iPhones du cluster).45 mlmobile nodes iPhones connus (en ligne, batterie, thermique, jobs)6 mlmobile stats file, baux, compteurs7 mlmobile ping [--on NOM] aller-retour vers un iPhone8 mlmobile info [--on NOM] infos système d'un iPhone (sys.info)9 mlmobile fetch URL [--on NOM] [--method M] [--header K:V]… [--body TXT] [--timeout S] [--raw]10 mlmobile run script.js [--args JSON] [--on NOM] [--timeout S]11 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)12 mlmobile fetch-list urls.txt [--tag T] [--out results.jsonl] 1 http.fetch par URL, réparti sur tous les iPhones13 mlmobile job ID | mlmobile wait ID | mlmobile jobs [--state S] [--tag T] [--limit N] | mlmobile cancel ID1415Le script JS doit définir `function main(args)` et retourner une valeur JSON. API disponible dans le script :16 ml.http(url, {method, headers, body, timeout}) -> {status, headers, body} (synchrone)17 ml.log(msg) ml.sleep(ms) ml.device -> {name, model, os, battery, …}1819Env : MLMOBILE_URL (défaut http://M4M64a.maclustr.io:9320), MLMOBILE_TOKEN.20"""21import argparse22import json23import os24import sys25import time26import urllib.error27import urllib.request2829URL = os.environ.get("MLMOBILE_URL", "http://M4M64a.maclustr.io:9320").rstrip("/")30TOKEN = os.environ.get("MLMOBILE_TOKEN", "mlmob_4f9a1c7e2b8d6035a9e1c4b7d2f8036e5a1b9c3d")313233def api(method, path, body=None, timeout=40):34 data = json.dumps(body).encode() if body is not None else None35 req = urllib.request.Request(URL + path, data=data, method=method,36 headers={"Authorization": "Bearer " + TOKEN, "Content-Type": "application/json"})37 try:38 with urllib.request.urlopen(req, timeout=timeout) as r:39 raw = r.read()40 return r.status, (json.loads(raw) if raw else None)41 except urllib.error.HTTPError as e:42 raw = e.read()43 try:44 return e.code, json.loads(raw)45 except Exception:46 return e.code, {"error": raw.decode(errors="replace")}47 except Exception as e:48 sys.exit("coordinateur injoignable (%s) : %s" % (URL, e))495051def fmt_age(s):52 if s is None:53 return "—"54 if s < 90:55 return "%ds" % s56 if s < 5400:57 return "%dmin" % (s // 60)58 return "%.1fh" % (s / 3600)596061def cmd_nodes(a):62 _, d = api("GET", "/api/workers")63 ws = d.get("workers", [])64 if not ws:65 print("aucun iPhone enregistré (ouvrir MacLustr → Nœud → Mode nœud)")66 return67 print("%-14s %-6s %-22s %-8s %-6s %-9s %-8s %-6s %s" % ("nom", "état", "modèle", "iOS", "batt.", "therm.", "réseau", "actifs", "faits/échecs"))68 for w in ws:69 m = w.get("metrics") or {}70 batt = m.get("battery")71 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 "—"72 print("%-14s %-6s %-22s %-8s %-6s %-9s %-8s %-6s %d/%d vu il y a %s" % (73 w["name"], "ON" if w["online"] else "off", (w.get("marketing") or w.get("model") or "")[:22], (w.get("os") or "")[:8],74 batt_s, (m.get("thermal") or "—")[:9], (m.get("net") or "—")[:8], w.get("activeJobs", 0),75 w.get("jobsDone", 0), w.get("jobsFailed", 0), fmt_age(w.get("ageS"))))767778def cmd_stats(a):79 _, d = api("GET", "/api/stats")80 print(json.dumps(d, indent=2, ensure_ascii=False))818283def submit_and_wait(spec, timeout):84 code, d = api("POST", "/api/jobs", spec)85 if code != 200 or not d.get("ids"):86 sys.exit("refus : %s" % (d.get("errors") or d))87 jid = d["ids"][0]88 return wait_for(jid, timeout)899091def wait_for(jid, timeout):92 t0 = time.time()93 while True:94 _, j = api("GET", "/api/jobs/%s/wait?timeout=25" % jid, timeout=40)95 if j.get("state") in ("done", "failed", "expired", "cancelled"):96 return j97 if time.time() - t0 > timeout:98 return j99100101def print_job(j, raw=False):102 st = j.get("state")103 if st == "done":104 r = j.get("result")105 if raw and isinstance(r, dict) and "body" in r:106 sys.stdout.write(r["body"] if isinstance(r["body"], str) else json.dumps(r["body"]))107 sys.stdout.write("\n")108 else:109 print(json.dumps(r, indent=2, ensure_ascii=False))110 print("# %s par %s en %s ms" % (j["id"], j.get("worker"), j.get("durationMs")), file=sys.stderr)111 else:112 print(json.dumps({k: j.get(k) for k in ("id", "state", "error", "worker", "attempts")}, indent=2, ensure_ascii=False))113 if st != "queued":114 sys.exit(1)115 print("(toujours en file — aucun iPhone en ligne ?)", file=sys.stderr)116117118def cmd_ping(a):119 j = submit_and_wait({"type": "ping", "payload": {"t": time.time()}, "target": a.on, "timeoutS": 10}, a.timeout)120 print_job(j)121122123def cmd_info(a):124 j = submit_and_wait({"type": "sys.info", "payload": {}, "target": a.on, "timeoutS": 15}, a.timeout)125 print_job(j)126127128def fetch_payload(url, a):129 headers = {}130 for h in (a.header or []):131 if ":" in h:132 k, v = h.split(":", 1)133 headers[k.strip()] = v.strip()134 p = {"url": url, "method": a.method, "timeoutS": a.timeout, "maxBytes": a.max_bytes}135 if headers:136 p["headers"] = headers137 if getattr(a, "body", None):138 p["body"] = a.body139 return p140141142def cmd_fetch(a):143 j = submit_and_wait({"type": "http.fetch", "payload": fetch_payload(a.url, a), "target": a.on, "timeoutS": a.timeout + 5}, a.timeout + 30)144 print_job(j, raw=a.raw)145146147def cmd_run(a):148 code = open(a.script).read()149 args = json.loads(a.args) if a.args else {}150 j = submit_and_wait({"type": "js.run", "payload": {"code": code, "args": args}, "target": a.on, "timeoutS": a.timeout}, a.timeout + 30)151 print_job(j)152153154def run_batch(specs, tag, timeout, out=None):155 code, d = api("POST", "/api/jobs", {"jobs": specs}, timeout=120)156 if code != 200:157 sys.exit("refus : %s" % d)158 ids = d["ids"]159 if d.get("errors"):160 print("erreurs de soumission : %s" % d["errors"][:5], file=sys.stderr)161 print("%d jobs soumis (tag %s)" % (len(ids), tag), file=sys.stderr)162 pending = set(ids)163 results = {}164 t0 = time.time()165 fh = open(out, "a") if out else None166 while pending and time.time() - t0 < timeout:167 _, d = api("GET", "/api/jobs?tag=%s&limit=10000&full=1" % tag, timeout=60)168 for j in d.get("jobs", []):169 if j["id"] in pending and j["state"] in ("done", "failed", "expired", "cancelled"):170 pending.discard(j["id"])171 results[j["id"]] = j172 if fh:173 fh.write(json.dumps(j, ensure_ascii=False) + "\n")174 fh.flush()175 done = len(results)176 print("\r%d/%d terminés" % (done, len(ids)), end="", file=sys.stderr)177 if pending:178 time.sleep(2)179 print("", file=sys.stderr)180 if fh:181 fh.close()182 return [results.get(i) or {"id": i, "state": "pending"} for i in ids]183184185def cmd_map(a):186 code = open(a.script).read()187 items = json.load(open(a.items))188 tag = a.tag or ("map-%d" % int(time.time()))189 specs = [{"type": "js.run", "payload": {"code": code, "args": it}, "tag": tag, "timeoutS": a.timeout} for it in items]190 res = run_batch(specs, tag, a.timeout * max(1, len(items)) / max(1, a.concurrency) + 120)191 print(json.dumps([{"args": items[i], "state": r.get("state"), "worker": r.get("worker"),192 "result": r.get("result"), "error": r.get("error")} for i, r in enumerate(res)], indent=2, ensure_ascii=False))193194195def cmd_fetch_list(a):196 urls = [u.strip() for u in open(a.urls) if u.strip() and not u.startswith("#")]197 tag = a.tag or ("fetch-%d" % int(time.time()))198 specs = [{"type": "http.fetch", "payload": fetch_payload(u, a), "tag": tag, "timeoutS": a.timeout + 5} for u in urls]199 res = run_batch(specs, tag, a.timeout * len(urls) / 3 + 120, out=a.out)200 ok = sum(1 for r in res if r.get("state") == "done")201 by_worker = {}202 for r in res:203 by_worker[r.get("worker") or "?"] = by_worker.get(r.get("worker") or "?", 0) + 1204 print("%d/%d réussis ; répartition : %s%s" % (ok, len(urls), by_worker, " ; résultats dans %s" % a.out if a.out else ""))205 if not a.out:206 for u, r in zip(urls, res):207 rr = r.get("result") or {}208 print("%-4s %-7s %s (%s, %s o)" % (rr.get("status", "—"), r.get("worker") or "—", u, r.get("state"), rr.get("bytes", "—")))209210211def cmd_job(a):212 _, j = api("GET", "/api/jobs/%s" % a.id)213 print(json.dumps(j, indent=2, ensure_ascii=False))214215216def cmd_wait(a):217 print_job(wait_for(a.id, a.timeout))218219220def cmd_jobs(a):221 qs = "limit=%d" % a.limit222 if a.state:223 qs += "&state=" + a.state224 if a.tag:225 qs += "&tag=" + a.tag226 _, d = api("GET", "/api/jobs?" + qs)227 for j in d.get("jobs", []):228 print("%-12s %-9s %-10s %-8s %-10s %s" % (j["id"], j["state"], j["type"], j.get("worker") or "—",229 time.strftime("%H:%M:%S", time.localtime(j["created"])),230 (j.get("error") or "")[:60]))231232233def cmd_cancel(a):234 code, d = api("DELETE", "/api/jobs/%s" % a.id)235 print(d)236237238def main():239 p = argparse.ArgumentParser(prog="mlmobile", description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)240 sp = p.add_subparsers(dest="cmd", required=True)241 sp.add_parser("nodes").set_defaults(f=cmd_nodes)242 sp.add_parser("stats").set_defaults(f=cmd_stats)243 x = sp.add_parser("ping"); x.add_argument("--on"); x.add_argument("--timeout", type=float, default=30); x.set_defaults(f=cmd_ping)244 x = sp.add_parser("info"); x.add_argument("--on"); x.add_argument("--timeout", type=float, default=30); x.set_defaults(f=cmd_info)245 x = sp.add_parser("fetch"); x.add_argument("url"); x.add_argument("--on"); x.add_argument("--method", default="GET")246 x.add_argument("--header", action="append"); x.add_argument("--body"); x.add_argument("--timeout", type=float, default=30)247 x.add_argument("--max-bytes", type=int, default=2 * 1024 * 1024); x.add_argument("--raw", action="store_true"); x.set_defaults(f=cmd_fetch)248 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)249 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)250 x.add_argument("--concurrency", type=int, default=6); x.set_defaults(f=cmd_map)251 x = sp.add_parser("fetch-list"); x.add_argument("urls"); x.add_argument("--tag"); x.add_argument("--out"); x.add_argument("--method", default="GET")252 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)253 x = sp.add_parser("job"); x.add_argument("id"); x.set_defaults(f=cmd_job)254 x = sp.add_parser("wait"); x.add_argument("id"); x.add_argument("--timeout", type=float, default=120); x.set_defaults(f=cmd_wait)255 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)256 x = sp.add_parser("cancel"); x.add_argument("id"); x.set_defaults(f=cmd_cancel)257 a = p.parse_args()258 a.f(a)259260261if __name__ == "__main__":262 main()263