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