#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""SkillFishOS - the AI cluster: which boards, how they are doing, and adding one.

    skillfish-cluster stato            telemetry of every board, as JSON
    skillfish-cluster elenco           just the list
    skillfish-cluster aggiungi <ip> <user> <password>
    skillfish-cluster togli <ip>
    skillfish-cluster avvia | ferma    the RPC node on every board
    skillfish-cluster demone [sec]     keep /run/skillfish/cluster.json fresh

WHY A CLUSTER AT ALL
One BC-250 has about 14 GB of GPU memory. A 27B model at Q6 needs 22 and simply
does not load: `failed to fit params to free device memory`, and the GPU drops
to ErrorDeviceLost. Two boards have 28 GB and it runs. Measured on the .32 and
the .40 on 12/09/2026: 7.57 tokens per second.

⚠️ AND IT IS NOT FASTER. A model that fits on one board runs SLOWER when you
split it: 35.7 tokens per second on one, 23.3 on two. Every layer boundary is
a trip over the network, and that trip costs more than the work it saves. The
cluster is for running what does not fit, nothing else. Anybody who reads this
file should know that before turning it on.

⚠️ THE RPC NODE HAS NO AUTHENTICATION. llama.cpp says so itself at startup:
"This is an experimental feature and is not secure". Whoever can reach that
port can run whatever they like on that GPU. It binds to the LAN only when
somebody asks for it, and it is not started at boot.
"""
import base64
import io
import json
import os
import shlex
import subprocess
import sys
import time

CONF = "/etc/skillfish/cluster.json"
CHIAVE = "/etc/skillfish/cluster_ed25519"
PORTA_RPC = 50052
NODO = "/opt/skillfish-rpc/ggml-rpc-server"

D = "/sys/class/drm/card0/device"


# ----------------------------------------------------------------- la lista --
def leggi():
    try:
        with open(CONF, encoding="utf-8") as f:
            d = json.load(f)
    except (OSError, ValueError):
        return {"schede": []}
    d.setdefault("schede", [])
    return d


def scrivi(d):
    os.makedirs(os.path.dirname(CONF), exist_ok=True)
    tmp = CONF + ".nuovo"
    with open(tmp, "w", encoding="utf-8") as f:
        json.dump(d, f, indent=1)
    os.replace(tmp, CONF)
    # ⚠️ 0600: qui dentro ci sono gli indirizzi e gli utenti delle altre
    # schede, e lo leggono soltanto skillfish-cluster e l'helper, che girano
    # da root. Il file che leggono tutti e' un altro, quello in /run.
    os.chmod(CONF, 0o600)


# ------------------------------------------------------------- la telemetria --
def _leggi_int(p):
    try:
        with open(p, encoding="utf-8") as f:
            return int(f.read().strip())
    except (OSError, ValueError):
        return 0


def locale():
    """I numeri di QUESTA scheda, letti dal kernel."""
    import glob
    watt = 0
    for p in glob.glob(D + "/hwmon/hwmon*/power1_average"):
        watt = _leggi_int(p) // 1000000
        break
    clock = 0
    try:
        with open(D + "/pp_dpm_sclk", encoding="utf-8") as f:
            for r in f:
                if "*" in r:
                    clock = int("".join(c for c in r.split()[1] if c.isdigit()))
                    break
    except (OSError, ValueError, IndexError):
        # il nodo non c'e' o non dice quello che ci aspettiamo: clock resta 0,
        # e uno zero nella telemetria si vede benissimo.
        pass
    temp = 0
    for p in glob.glob(D + "/hwmon/hwmon*/temp1_input"):
        temp = _leggi_int(p) // 1000
        break
    # ⚠️ Il carico NON si legge da un file del kernel: su gfx1013 l'SMU non
    # espone quella metrica e riporta 0xFFFF. Lo campiona
    # skillfish-gpu-util.service con radeontop e lo scrive qui.
    try:
        with open("/run/skillfish-gpu-util", encoding="utf-8") as f:
            carico = int(float(f.read().strip()))
    except (OSError, ValueError):
        carico = -1
    return {
        "vram_usata": _leggi_int(D + "/mem_info_vram_used"),
        "vram_totale": _leggi_int(D + "/mem_info_vram_total"),
        "gtt_usata": _leggi_int(D + "/mem_info_gtt_used"),
        "gtt_totale": _leggi_int(D + "/mem_info_gtt_total"),
        "watt": watt, "clock": clock, "temp": temp, "carico": carico,
    }


# Lo stesso identico lavoro, ma da eseguire sull'altra scheda. Si manda il
# codice invece di installare un agente: una scheda in piu' nel cluster non
# deve voler dire un pacchetto in piu' da aggiornare.
RACCOGLI = r'''
import glob, json
D = "/sys/class/drm/card0/device"
def i(p):
    try:
        return int(open(p).read().strip())
    except Exception:
        return 0
w = 0
for p in glob.glob(D + "/hwmon/hwmon*/power1_average"):
    w = i(p) // 1000000
    break
t = 0
for p in glob.glob(D + "/hwmon/hwmon*/temp1_input"):
    t = i(p) // 1000
    break
c = 0
try:
    for r in open(D + "/pp_dpm_sclk"):
        if "*" in r:
            c = int("".join(x for x in r.split()[1] if x.isdigit()))
            break
except Exception:
    pass
try:
    u = int(float(open("/run/skillfish-gpu-util").read().strip()))
except Exception:
    u = -1
print(json.dumps({
    "vram_usata": i(D + "/mem_info_vram_used"),
    "vram_totale": i(D + "/mem_info_vram_total"),
    "gtt_usata": i(D + "/mem_info_gtt_used"),
    "gtt_totale": i(D + "/mem_info_gtt_total"),
    "watt": w, "clock": c, "temp": t, "carico": u}))
'''


def ssh(ip, utente, comando, password="", t=25):
    """Un comando sull'altra scheda. Con la chiave se c'e', altrimenti con la
    password (solo la prima volta, per installare la chiave).

    ⚠️ StrictHostKeyChecking=accept-new e non "no": la prima volta si accetta,
    dalle volte dopo un cambio di chiave e' un errore, come deve essere.
    """
    # ⚠️ Le opzioni si costruiscono a COPPIE. La prima stesura toglieva
    # "BatchMode=yes" da una lista piatta e lasciava un "-o" orfano: ssh
    # rispondeva "no argument after keyword -o" e sembrava un problema di rete.
    def opz(*coppie):
        fuori = []
        for c in coppie:
            fuori += ["-o", c]
        return fuori

    comuni = ["StrictHostKeyChecking=accept-new", "ConnectTimeout=8",
              "LogLevel=ERROR"]
    if password:
        if not _ha("sshpass"):
            return 127, "", "manca sshpass"
        # Con la password: niente BatchMode (vorrebbe dire "non chiedere mai
        # niente", e la password e' proprio quello che stiamo passando) e
        # niente chiavi, cosi' non prova quelle sbagliate e finisce i tentativi.
        cmd = (["sshpass", "-p", password, "ssh"]
               + opz(*comuni, "PubkeyAuthentication=no", "PreferredAuthentications=password")
               + ["%s@%s" % (utente, ip), comando])
    else:
        cmd = (["ssh", "-i", CHIAVE] + opz(*comuni, "BatchMode=yes")
               + ["%s@%s" % (utente, ip), comando])
    try:
        r = subprocess.run(cmd, capture_output=True, text=True, timeout=t)
        return r.returncode, r.stdout, r.stderr
    except (OSError, subprocess.TimeoutExpired) as e:
        return 1, "", str(e)


def _ha(x):
    return subprocess.run(["sh", "-c", "command -v " + shlex.quote(x)],
                          capture_output=True).returncode == 0


def remota(s):
    """La telemetria di una scheda del cluster, presa via ssh."""
    codice = base64.b64encode(RACCOGLI.encode()).decode()
    rc, out, err = ssh(s["ip"], s.get("utente", "root"),
                       "echo %s | base64 -d | python3 -" % codice)
    if rc != 0:
        return {"errore": (err or "non risponde").strip()[:120]}
    try:
        return json.loads(out)
    except ValueError:
        return {"errore": "risposta non leggibile"}


def nodo_vivo(ip):
    """Il nodo RPC risponde su quella porta?"""
    import socket
    s = socket.socket()
    s.settimeout(3)
    try:
        s.connect((ip, PORTA_RPC))
        return True
    except OSError:
        return False
    finally:
        s.close()


def stato():
    d = leggi()
    schede = [{"ip": "locale", "nome": os.uname().nodename, "locale": True,
               "nodo": True, **locale()}]
    for s in d["schede"]:
        schede.append({"ip": s["ip"], "nome": s.get("nome") or s["ip"],
                       "locale": False, "nodo": nodo_vivo(s["ip"]),
                       **remota(s)})
    # Il totale: e' il numero per cui si accende un cluster, e va detto grosso.
    tot = sum((x.get("vram_totale", 0) + x.get("gtt_totale", 0)) for x in schede)
    uso = sum((x.get("vram_usata", 0) + x.get("gtt_usata", 0)) for x in schede)
    return {"schede": schede, "memoria_totale": tot, "memoria_usata": uso,
            "watt_totali": sum(x.get("watt", 0) for x in schede)}


# --------------------------------------------------------------- aggiungere --
SORGENTI = [
    # dove puo' stare llama.cpp: Unsloth lo mette nella home dell'utente
    "~/.unsloth/llama.cpp",
    "/opt/unsloth/llama.cpp",
]


def _dove_llama():
    """L'albero di llama.cpp su questa macchina, se c'e'."""
    import glob
    for base in SORGENTI:
        for d in glob.glob(os.path.expanduser(base)):
            if os.path.isdir(os.path.join(d, "tools", "rpc")):
                return d
    # e le home degli utenti veri, che e' dove finisce davvero
    for u in os.listdir("/home"):
        d = "/home/%s/.unsloth/llama.cpp" % u
        if os.path.isdir(os.path.join(d, "tools", "rpc")):
            return d
    return ""


def compila_nodo(log=None):
    """Compila ggml-rpc-server su QUESTA macchina. Restituisce (percorso, motivo).

    ⚠️ Il target si chiama `ggml-rpc-server`, non `rpc-server`: quello non
    esiste e make risponde "No rule to make target", che sembra un guasto e
    invece e' solo il nome sbagliato.

    ⚠️ Servono glslc e spirv-headers, che su una SkillFishOS appena installata
    non ci sono: senza, cmake si ferma su "Could not find SPIRV-Headers" e il
    messaggio non dice quale pacchetto manchi.
    """
    def nota(x):
        if log is not None:
            log.append(x)

    src = _dove_llama()
    if not src:
        return "", "non trovo il sorgente di llama.cpp"
    pronto = os.path.join(src, "build-rpc", "bin", "ggml-rpc-server")
    if os.access(pronto, os.X_OK):
        nota("nodo gia' compilato qui")
        return pronto, ""
    nota("compilo il nodo su questa macchina")

    mancanti = [x for x in ("cmake", "glslc") if not _ha(x)]
    if not os.path.exists("/usr/share/cmake/SPIRV-Headers") and \
       not os.path.exists("/usr/lib/cmake/SPIRV-Headers"):
        mancanti.append("spirv-headers")
    if mancanti:
        nota("installo: " + " ".join(mancanti))
        pacchetti = ["cmake" if m == "cmake" else m for m in mancanti]
        r = subprocess.run(["apt-get", "install", "-y"] + pacchetti,
                           capture_output=True, text=True, timeout=900,
                           env={**os.environ, "DEBIAN_FRONTEND": "noninteractive"})
        if r.returncode != 0:
            return "", "non riesco a installare: " + " ".join(pacchetti)

    nota("preparo la compilazione")
    r = subprocess.run(
        ["cmake", "-B", "build-rpc", "-DGGML_VULKAN=ON", "-DGGML_RPC=ON",
         "-DLLAMA_BUILD_TESTS=OFF", "-DLLAMA_BUILD_EXAMPLES=OFF",
         "-DLLAMA_CURL=OFF", "-DCMAKE_BUILD_TYPE=Release"],
        cwd=src, capture_output=True, text=True, timeout=600)
    if r.returncode != 0:
        return "", (r.stderr or r.stdout).strip()[-300:]

    nota("compilo (qualche minuto)")
    r = subprocess.run(["cmake", "--build", "build-rpc",
                        "--target", "ggml-rpc-server", "-j", str(os.cpu_count() or 4)],
                       cwd=src, capture_output=True, text=True, timeout=3600)
    if not os.access(pronto, os.X_OK):
        return "", (r.stderr or r.stdout).strip()[-300:] or "compilazione non riuscita"
    return pronto, ""


def installa_nodo(ip, utente, log=None):
    """Porta il nodo RPC sulla scheda remota. Restituisce (ok, motivo)."""
    import tarfile
    import tempfile

    def nota(x):
        if log is not None:
            log.append(x)

    locale_bin = NODO if os.access(NODO, os.X_OK) else ""
    if not locale_bin:
        # ⚠️ Puo' esserci gia' compilato dentro l'albero di llama.cpp senza
        # essere stato copiato in /opt: in quel caso non si compila niente, e
        # dirlo lo stesso sarebbe raccontare un lavoro che non si e' fatto.
        locale_bin, err = compila_nodo(log)
        if not locale_bin:
            return False, err

    # il binario e le sue librerie: si spedisce la cartella intera, perche'
    # ggml-rpc-server senza le sue libggml non parte
    cartella = os.path.dirname(locale_bin)
    with tempfile.NamedTemporaryFile(suffix=".tgz", delete=False) as t:
        pacco = t.name
    with tarfile.open(pacco, "w:gz") as tar:
        for n in sorted(os.listdir(cartella)):
            if n == "ggml-rpc-server" or n.startswith("libggml"):
                tar.add(os.path.join(cartella, n), arcname=n)

    nota("copio il nodo sulla scheda")
    r = subprocess.run(
        ["scp", "-i", CHIAVE, "-o", "StrictHostKeyChecking=accept-new",
         "-o", "BatchMode=yes", "-o", "ConnectTimeout=10",
         pacco, "%s@%s:/tmp/skillfish-rpc.tgz" % (utente, ip)],
        capture_output=True, text=True, timeout=300)
    os.unlink(pacco)
    if r.returncode != 0:
        return False, (r.stderr or "copia non riuscita").strip()[:200]

    rc, out, err = ssh(ip, utente,
                       "mkdir -p /opt/skillfish-rpc && "
                       "tar xzf /tmp/skillfish-rpc.tgz -C /opt/skillfish-rpc && "
                       "rm -f /tmp/skillfish-rpc.tgz && "
                       "test -x /opt/skillfish-rpc/ggml-rpc-server && echo installato",
                       t=120)
    if rc != 0 or "installato" not in out:
        return False, (err or "non si scompatta").strip()[:200]
    nota("nodo installato")
    return True, ""


def aggiungi(ip, utente, password):
    """Verifica una scheda, ci mette la chiave e prepara il nodo.

    I passi sono raccontati uno per uno perche' la finestra li mostra: se
    qualcosa non va, chi guarda deve sapere DOVE si e' fermato, non solo che
    e' fallito.
    """
    passi = []

    def passo(nome, ok, nota=""):
        passi.append({"passo": nome, "ok": bool(ok), "nota": nota})
        return ok

    rc, out, err = ssh(ip, utente, "echo vivo", password=password)
    if not passo("collegamento", rc == 0 and "vivo" in out, err.strip()[:160]):
        return {"ok": False, "passi": passi}

    # ⚠️ La chiave, subito. Da qui in poi la password non serve piu' e non
    # viene salvata: tenere in chiaro la password di root di un'altra macchina
    # dentro un file di configurazione sarebbe la cosa peggiore di tutte.
    if not os.path.exists(CHIAVE):
        subprocess.run(["ssh-keygen", "-q", "-t", "ed25519", "-N", "",
                        "-C", "skillfish-cluster", "-f", CHIAVE],
                       capture_output=True, timeout=60)
        os.chmod(CHIAVE, 0o600)
    with io.open(CHIAVE + ".pub", encoding="utf-8") as f:
        pub = f.read().strip()
    rc, _, err = ssh(ip, utente,
                     "mkdir -p ~/.ssh && chmod 700 ~/.ssh && "
                     "grep -qxF %s ~/.ssh/authorized_keys 2>/dev/null || "
                     "echo %s >> ~/.ssh/authorized_keys" % (shlex.quote(pub), shlex.quote(pub)),
                     password=password)
    if not passo("chiave", rc == 0, err.strip()[:160]):
        return {"ok": False, "passi": passi}

    rc, out, _ = ssh(ip, utente, "cat /sys/devices/virtual/dmi/id/product_name 2>/dev/null")
    modello = out.strip()
    passo("scheda", True, modello or "sconosciuta")

    rc, out, _ = ssh(ip, utente, "test -x %s && echo si || echo no" % NODO)
    ha_nodo = "si" in out
    if ha_nodo:
        passo("nodo RPC", True, "presente")
    else:
        # ⚠️ Si installa copiando il binario di QUESTA macchina, non
        # compilandolo la': il protocollo RPC e' versionato e le due parti
        # devono parlare la stessa versione. Compilare sull'altra scheda da un
        # sorgente magari piu' vecchio e' esattamente l'errore che il
        # 12/09/2026 ha prodotto "RPC handshake failed".
        righe = []
        ok, err = installa_nodo(ip, utente, righe)
        ha_nodo = ok
        passo("nodo RPC", ok, "; ".join(righe) if ok else err)

    d = leggi()
    d["schede"] = [s for s in d["schede"] if s["ip"] != ip]
    d["schede"].append({"ip": ip, "utente": utente, "nome": modello or ip,
                        "aggiunta": time.strftime("%Y-%m-%dT%H:%M:%S")})
    scrivi(d)
    passo("salvata", True, "")
    return {"ok": True, "passi": passi, "ha_nodo": ha_nodo}


def togli(ip):
    d = leggi()
    prima = len(d["schede"])
    d["schede"] = [s for s in d["schede"] if s["ip"] != ip]
    scrivi(d)
    return {"ok": len(d["schede"]) < prima}


def avvia(accendi=True):
    """Accende (o spegne) il nodo RPC su tutte le schede del cluster."""
    fuori = []
    for s in leggi()["schede"]:
        if accendi:
            c = ("cd /opt/skillfish-rpc && LD_LIBRARY_PATH=/opt/skillfish-rpc "
                 "setsid nohup ./ggml-rpc-server -H 0.0.0.0 -p %d --cache "
                 "</dev/null >/tmp/skillfish-rpc.log 2>&1 & sleep 3; echo avviato"
                 % PORTA_RPC)
        else:
            # ⚠️ pkill -x sul NOME, non -f sul motivo: con -f il modello
            # ucciderebbe anche la shell che porta il comando, perche' il
            # motivo compare nella sua riga.
            c = "pkill -x ggml-rpc-server; echo fermato"
        rc, out, err = ssh(s["ip"], s.get("utente", "root"), c, t=40)
        fuori.append({"ip": s["ip"], "ok": rc == 0,
                      "nota": (out or err).strip()[:100]})
    return {"ok": all(x["ok"] for x in fuori), "schede": fuori}


STATO_PUB = "/run/skillfish/cluster.json"


def demone(ogni=5):
    """Raccoglie la telemetria e la lascia dove la puo' leggere chiunque.

    ⚠️ PERCHE' UN SERVIZIO E NON UNA CHIAMATA DALLA FINESTRA.
    La telemetria delle altre schede passa da ssh con la chiave del cluster,
    che sta in /etc/skillfish ed e' leggibile solo da root. Chiedendola dalla
    finestra si passava dall'helper, cioe' da pkexec: il grafico restava vuoto
    finche' l'utente non digitava la password per qualche altro motivo, e
    nessuno collegava le due cose. Provato sulla .32 il 12/09/2026: il
    riquadro diceva "0 schede" con due schede accese e funzionanti.
    Adesso la raccolta la fa root una volta sola, e il risultato sta in /run,
    che leggono tutti: la finestra, il Remote Manager, chiunque.

    ⚠️ E si scrive di fianco e poi si sposta, perche' chi legge non deve mai
    trovare mezzo file.
    """
    os.makedirs(os.path.dirname(STATO_PUB), exist_ok=True)
    while True:
        try:
            d = stato()
        except Exception as e:                       # una scheda che sparisce
            d = {"schede": [], "errore": str(e)[:200]}
        d["aggiornato"] = int(time.time())
        tmp = STATO_PUB + ".nuovo"
        try:
            with open(tmp, "w", encoding="utf-8") as f:
                json.dump(d, f)
            os.replace(tmp, STATO_PUB)
            # ⚠️ 0644 VOLUTO: questo file lo legge il Control Center come
            # utente normale, ed e' tutto il punto di averlo (vedi il commento
            # in cima). Dentro non c'e' niente di segreto: nomi, gradi, watt.
            os.chmod(STATO_PUB, 0o644)
        except OSError:
            # un giro perso non si recupera e non si segnala: fra un secondo
            # ce n'e' un altro, e i dati vecchi li marca gia' chi legge.
            pass
        time.sleep(ogni)


def main():
    a = sys.argv[1:] or ["stato"]
    c = a[0]
    if c == "stato":
        print(json.dumps(stato(), indent=1))
    elif c == "elenco":
        print(json.dumps(leggi(), indent=1))
    elif c == "aggiungi" and len(a) >= 4:
        print(json.dumps(aggiungi(a[1], a[2], a[3]), indent=1))
    elif c == "togli" and len(a) >= 2:
        print(json.dumps(togli(a[1]), indent=1))
    elif c == "demone":
        demone(int(a[1]) if len(a) > 1 else 5)
    elif c in ("avvia", "ferma"):
        print(json.dumps(avvia(c == "avvia"), indent=1))
    else:
        print(__doc__.strip())
        return 2
    return 0


if __name__ == "__main__":
    sys.exit(main())
