File size: 10,592 Bytes
f705ee4
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
47f9388
 
 
 
 
f705ee4
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
47f9388
 
 
 
 
 
 
 
 
 
 
f705ee4
47f9388
 
 
 
 
 
 
 
 
 
eba33d0
 
 
 
 
f705ee4
47f9388
f705ee4
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
47f9388
f705ee4
 
 
eba33d0
 
 
 
 
 
 
f705ee4
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
"""Agent de pod : un shell pilote par le Hub.

POURQUOI. Ce pod n'expose ni IP publique ni port TCP -- seulement 8080 en HTTP
-- et son `PUBLIC_KEY` est vide. L'image `vllm/vllm-openai` ne contient pas de
`sshd`, et notre `dockerStartCmd` remplace l'entrypoint : rien n'ecoute sur 22.
Aucun SSH n'est donc possible. Le seul canal disponible est sortant : HTTPS
vers le Hub. Cet agent en fait un shell.

  - il relit `cmd/<pod>.sh` toutes les 8 s ; si l'empreinte a change, il
    l'execute dans un fil separe et publie la sortie dans `etat/out-<pod>.log` ;
  - il publie un etat machine toutes les 30 s dans `etat/vie-<pod>.log` :
    GPU, RAM DU CGROUP, disque, avancement du telechargement, queue des
    journaux en cours ;
  - il lance le telechargement des poids des le demarrage, sans attendre une
    premiere commande : 125,9 Go ne s'attendent pas.

PIEGE MESURE UNE FOIS. `free -g` dans un conteneur lit /proc/meminfo, donc la
RAM de l'HOTE (1 133 Go ici) et non la limite du conteneur (125 Go). Toute
conclusion sur la faisabilite du dechargement PLE tiree de `free` est fausse.
On lit donc le cgroup, v2 puis v1.

PIEGE DE LA VIE DU POD. `volumeInGb = 0` : le disque conteneur est efface a
chaque arret. Les 125,9 Go ne sont telecharges qu'une fois par vie du pod, donc
toute l'iteration doit tenir dans une seule session -- d'ou ce shell.
"""
import hashlib
import os
import shutil
import subprocess
import threading
import time
import urllib.request

DEPOT = "patdev/k3-a40-bootstrap"
POD = os.environ.get("RUNPOD_POD_ID", "inconnu")
JETON = os.environ.get("HF_TOKEN")
TRAVAIL = "/travail"
MODELE = os.environ.get("BANC_MODEL", "RadixArk/Qwen3.8-Flash-Next-NVFP4")
# Taille annoncee du depot, pour transformer les octets en pourcentage.
ATTENDU_GO = float(os.environ.get("BANC_TAILLE_GO", "125.9"))

os.makedirs(TRAVAIL, exist_ok=True)
os.makedirs(TRAVAIL + "/sorties", exist_ok=True)

UA = ("Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 "
      "(KHTML, like Gecko) Chrome/126.0.0.0 Safari/537.36")

COURANTES = []


def sh(cmd, timeout=25):
    try:
        r = subprocess.run(["bash", "-lc", cmd], capture_output=True,
                           text=True, timeout=timeout)
        return (r.stdout + r.stderr).strip()
    except Exception as e:
        return "%s: %s" % (type(e).__name__, str(e)[:120])


def publier(chemin_local, chemin_depot):
    from huggingface_hub import HfApi
    for essai in range(3):
        try:
            HfApi().upload_file(path_or_fileobj=chemin_local,
                                path_in_repo=chemin_depot,
                                repo_id=DEPOT, token=JETON)
            return True
        except Exception as e:
            if essai == 2:
                print("publication %s : %s" % (chemin_depot, e), flush=True)
        time.sleep(3)
    return False


def lire_hub(chemin):
    """Lit un fichier du depot en contournant le cache du CDN."""
    url = ("https://huggingface.co/%s/resolve/main/%s?t=%d"
           % (DEPOT, chemin, int(time.time())))
    entetes = {"User-Agent": UA, "Cache-Control": "no-cache"}
    if JETON:
        entetes["Authorization"] = "Bearer " + JETON
    try:
        req = urllib.request.Request(url, headers=entetes)
        with urllib.request.urlopen(req, timeout=25) as r:
            return r.read().decode("utf-8", "replace")
    except Exception:
        return None


# --------------------------------------------------------------- etat machine
def octets_repertoire(d):
    total = 0
    for racine, _, fichiers in os.walk(d):
        for f in fichiers:
            try:
                total += os.stat(os.path.join(racine, f),
                                 follow_symlinks=False).st_size
            except OSError:
                pass
    return total


def ram_cgroup():
    """La limite REELLE du conteneur, pas celle de l'hote."""
    paires = (("/sys/fs/cgroup/memory.max", "/sys/fs/cgroup/memory.current"),
              ("/sys/fs/cgroup/memory/memory.limit_in_bytes",
               "/sys/fs/cgroup/memory/memory.usage_in_bytes"))
    for lim, use in paires:
        try:
            t = open(lim).read().strip()
            u = int(open(use).read().strip())
            if t == "max":
                return None, u
            t = int(t)
            if t > (1 << 50):          # « pas de limite » deguise en immense
                return None, u
            return t, u
        except Exception:
            continue
    return None, None


CACHE = os.environ.get("HF_HOME") or os.path.expanduser("~/.cache/huggingface")

# Deux mesures successives de la taille du cache donnent un debit, donc un
# temps restant. Sans ca, « 41 Go sur 125,9 » ne dit pas si l'on attend cinq
# minutes ou une heure -- et c'est exactement la question qui coute de l'argent.
DERNIERE = {"t": None, "octets": None, "debit": None}


def etat():
    l = []
    l.append("[VIE] %s pod=%s" % (time.strftime("%H:%M:%S", time.gmtime()), POD))
    l.append(sh("nvidia-smi --query-gpu=name,memory.used,memory.total,"
                "utilization.gpu --format=csv,noheader"))
    lim, use = ram_cgroup()
    if lim:
        l.append("RAM conteneur : %.1f / %.1f Go  (cgroup)"
                 % (use / 1e9, lim / 1e9))
    elif use is not None:
        l.append("RAM conteneur : %.1f Go utilises, aucune limite cgroup lisible"
                 % (use / 1e9))
    try:
        du = shutil.disk_usage("/")
        l.append("disque / : %.0f Go libres sur %.0f" % (du.free / 1e9,
                                                         du.total / 1e9))
    except Exception:
        pass
    try:
        octets = octets_repertoire(CACHE)
        go = octets / 1e9
        maintenant = time.time()
        if DERNIERE["t"] and maintenant > DERNIERE["t"] + 5:
            brut = ((octets - DERNIERE["octets"]) / 1e9
                    / ((maintenant - DERNIERE["t"]) / 60.0))
            # Lissage : le debit instantane saute trop pour etre lisible.
            DERNIERE["debit"] = (brut if DERNIERE["debit"] is None
                                 else 0.6 * DERNIERE["debit"] + 0.4 * brut)
        DERNIERE["t"], DERNIERE["octets"] = maintenant, octets
        ligne = ("poids telecharges : %.1f / %.1f Go  (%.0f %%)"
                 % (go, ATTENDU_GO, 100 * go / ATTENDU_GO))
        d = DERNIERE["debit"]
        if d and d > 0.05:
            ligne += "  %.1f Go/min, reste ~%.0f min" % (d, (ATTENDU_GO - go) / d)
        elif d is not None:
            ligne += "  (a l'arret : %.2f Go/min)" % d
        l.append(ligne)
    except Exception as e:
        l.append("poids : mesure impossible (%s)" % type(e).__name__)
    l.append(sh("ps -eo comm,pcpu,rss --sort=-rss --no-headers | head -5 | "
                "awk '{printf \"%s %s%% %.1fGo | \", $1, $2, $3/1048576}'"))
    # `/tmp/vllm.log` et `/tmp/carte.log` sont ceux du banc lance par l'ancien
    # bootstrap : l'agent peut etre greffe sur un conteneur deja en marche.
    for nom in ("dl.log", "vllm.log", "banc.log",
                "/tmp/vllm.log", "/tmp/carte.log"):
        p = nom if nom.startswith("/") else os.path.join(TRAVAIL, nom)
        if os.path.exists(p):
            q = sh("tail -c 2500 %s | tr '\\r' '\\n' | grep -v '^$' | tail -6" % p)
            if q:
                l.append("--- %s ---\n%s" % (nom, q))
    vivantes = [c for c in COURANTES if c["fil"].is_alive()]
    for c in vivantes:
        l.append("commande en cours : %s (depuis %d s)"
                 % (c["nom"], time.time() - c["t0"]))
    return "\n".join(x for x in l if x)


# ------------------------------------------------------------- telechargement
CODE_DL = """
import os, time
from huggingface_hub import snapshot_download
t0 = time.time()
for essai in range(6):
    try:
        p = snapshot_download(%r, max_workers=8,
                              token=os.environ.get('HF_TOKEN'),
                              ignore_patterns=['*.pth', 'original/*'])
        print('TELECHARGE en %%.0f s -> %%s' %% (time.time() - t0, p), flush=True)
        break
    except Exception as e:
        print('essai', essai, type(e).__name__, str(e)[:200], flush=True)
        time.sleep(10)
else:
    print('TELECHARGEMENT ECHOUE', flush=True)
"""


def telecharger():
    chemin = os.path.join(TRAVAIL, "dl.py")
    with open(chemin, "w") as f:
        f.write(CODE_DL % MODELE)
    with open(os.path.join(TRAVAIL, "dl.log"), "w") as f:
        subprocess.run(["python3", chemin], stdout=f, stderr=subprocess.STDOUT)


# ------------------------------------------------------------------ commandes
def executer(nom, script):
    chemin = os.path.join(TRAVAIL, "sorties", nom + ".log")
    with open(chemin, "w") as f:
        f.write("$ %s\n%s\n%s\n%s\n" % (nom, "-" * 70, script, "-" * 70))
        f.flush()
        p = subprocess.Popen(["bash", "-lc", script], stdout=f,
                             stderr=subprocess.STDOUT, cwd=TRAVAIL)
        code = p.wait()
        f.write("\n%s\n[FIN] code=%d\n" % ("-" * 70, code))
    publier(chemin, "etat/out-%s.log" % POD)
    print("commande %s terminee (code %d)" % (nom, code), flush=True)


def boucle_commandes():
    vue = None
    n = 0
    while True:
        texte = lire_hub("cmd/%s.sh" % POD)
        if texte is not None and texte.strip():
            h = hashlib.sha256(texte.encode()).hexdigest()[:12]
            if h != vue:
                vue = h
                n += 1
                nom = "%03d-%s" % (n, h)
                print("commande recue %s" % nom, flush=True)
                fil = threading.Thread(target=executer, args=(nom, texte),
                                       daemon=True)
                fil.start()
                COURANTES.append({"nom": nom, "fil": fil, "t0": time.time()})
        time.sleep(8)


def boucle_etat():
    p = os.path.join(TRAVAIL, "vie.log")
    while True:
        try:
            with open(p, "w") as f:
                f.write(etat() + "\n")
            publier(p, "etat/vie-%s.log" % POD)
        except Exception as e:
            print("etat:", e, flush=True)
        time.sleep(20)


print("[AGENT] pod=%s modele=%s" % (POD, MODELE), flush=True)
# Quand l'agent est greffe sur un conteneur ou vLLM telecharge deja, relancer
# `snapshot_download` ferait ecrire deux processus dans les memes blobs du
# cache. On s'abstient : `BANC_SANS_DL=1`.
if os.environ.get("BANC_SANS_DL") == "1":
    print("[AGENT] telechargement laisse au processus deja en place", flush=True)
else:
    threading.Thread(target=telecharger, daemon=True).start()
threading.Thread(target=boucle_commandes, daemon=True).start()
boucle_etat()