patdev commited on
Commit
f705ee4
·
verified ·
1 Parent(s): c15d840

agent de pod : shell pilote par le Hub

Browse files
Files changed (1) hide show
  1. agent_pod.py +239 -0
agent_pod.py ADDED
@@ -0,0 +1,239 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """Agent de pod : un shell pilote par le Hub.
2
+
3
+ POURQUOI. Ce pod n'expose ni IP publique ni port TCP -- seulement 8080 en HTTP
4
+ -- et son `PUBLIC_KEY` est vide. L'image `vllm/vllm-openai` ne contient pas de
5
+ `sshd`, et notre `dockerStartCmd` remplace l'entrypoint : rien n'ecoute sur 22.
6
+ Aucun SSH n'est donc possible. Le seul canal disponible est sortant : HTTPS
7
+ vers le Hub. Cet agent en fait un shell.
8
+
9
+ - il relit `cmd/<pod>.sh` toutes les 8 s ; si l'empreinte a change, il
10
+ l'execute dans un fil separe et publie la sortie dans `etat/out-<pod>.log` ;
11
+ - il publie un etat machine toutes les 30 s dans `etat/vie-<pod>.log` :
12
+ GPU, RAM DU CGROUP, disque, avancement du telechargement, queue des
13
+ journaux en cours ;
14
+ - il lance le telechargement des poids des le demarrage, sans attendre une
15
+ premiere commande : 125,9 Go ne s'attendent pas.
16
+
17
+ PIEGE MESURE UNE FOIS. `free -g` dans un conteneur lit /proc/meminfo, donc la
18
+ RAM de l'HOTE (1 133 Go ici) et non la limite du conteneur (125 Go). Toute
19
+ conclusion sur la faisabilite du dechargement PLE tiree de `free` est fausse.
20
+ On lit donc le cgroup, v2 puis v1.
21
+
22
+ PIEGE DE LA VIE DU POD. `volumeInGb = 0` : le disque conteneur est efface a
23
+ chaque arret. Les 125,9 Go ne sont telecharges qu'une fois par vie du pod, donc
24
+ toute l'iteration doit tenir dans une seule session -- d'ou ce shell.
25
+ """
26
+ import hashlib
27
+ import os
28
+ import shutil
29
+ import subprocess
30
+ import threading
31
+ import time
32
+ import urllib.request
33
+
34
+ DEPOT = "patdev/k3-a40-bootstrap"
35
+ POD = os.environ.get("RUNPOD_POD_ID", "inconnu")
36
+ JETON = os.environ.get("HF_TOKEN")
37
+ TRAVAIL = "/travail"
38
+ MODELE = os.environ.get("BANC_MODEL", "RadixArk/Qwen3.8-Flash-Next-NVFP4")
39
+ # Taille annoncee du depot, pour transformer les octets en pourcentage.
40
+ ATTENDU_GO = float(os.environ.get("BANC_TAILLE_GO", "125.9"))
41
+
42
+ os.makedirs(TRAVAIL, exist_ok=True)
43
+ os.makedirs(TRAVAIL + "/sorties", exist_ok=True)
44
+
45
+ UA = ("Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 "
46
+ "(KHTML, like Gecko) Chrome/126.0.0.0 Safari/537.36")
47
+
48
+ COURANTES = []
49
+
50
+
51
+ def sh(cmd, timeout=25):
52
+ try:
53
+ r = subprocess.run(["bash", "-lc", cmd], capture_output=True,
54
+ text=True, timeout=timeout)
55
+ return (r.stdout + r.stderr).strip()
56
+ except Exception as e:
57
+ return "%s: %s" % (type(e).__name__, str(e)[:120])
58
+
59
+
60
+ def publier(chemin_local, chemin_depot):
61
+ from huggingface_hub import HfApi
62
+ for essai in range(3):
63
+ try:
64
+ HfApi().upload_file(path_or_fileobj=chemin_local,
65
+ path_in_repo=chemin_depot,
66
+ repo_id=DEPOT, token=JETON)
67
+ return True
68
+ except Exception as e:
69
+ if essai == 2:
70
+ print("publication %s : %s" % (chemin_depot, e), flush=True)
71
+ time.sleep(3)
72
+ return False
73
+
74
+
75
+ def lire_hub(chemin):
76
+ """Lit un fichier du depot en contournant le cache du CDN."""
77
+ url = ("https://huggingface.co/%s/resolve/main/%s?t=%d"
78
+ % (DEPOT, chemin, int(time.time())))
79
+ entetes = {"User-Agent": UA, "Cache-Control": "no-cache"}
80
+ if JETON:
81
+ entetes["Authorization"] = "Bearer " + JETON
82
+ try:
83
+ req = urllib.request.Request(url, headers=entetes)
84
+ with urllib.request.urlopen(req, timeout=25) as r:
85
+ return r.read().decode("utf-8", "replace")
86
+ except Exception:
87
+ return None
88
+
89
+
90
+ # --------------------------------------------------------------- etat machine
91
+ def octets_repertoire(d):
92
+ total = 0
93
+ for racine, _, fichiers in os.walk(d):
94
+ for f in fichiers:
95
+ try:
96
+ total += os.stat(os.path.join(racine, f),
97
+ follow_symlinks=False).st_size
98
+ except OSError:
99
+ pass
100
+ return total
101
+
102
+
103
+ def ram_cgroup():
104
+ """La limite REELLE du conteneur, pas celle de l'hote."""
105
+ paires = (("/sys/fs/cgroup/memory.max", "/sys/fs/cgroup/memory.current"),
106
+ ("/sys/fs/cgroup/memory/memory.limit_in_bytes",
107
+ "/sys/fs/cgroup/memory/memory.usage_in_bytes"))
108
+ for lim, use in paires:
109
+ try:
110
+ t = open(lim).read().strip()
111
+ u = int(open(use).read().strip())
112
+ if t == "max":
113
+ return None, u
114
+ t = int(t)
115
+ if t > (1 << 50): # « pas de limite » deguise en immense
116
+ return None, u
117
+ return t, u
118
+ except Exception:
119
+ continue
120
+ return None, None
121
+
122
+
123
+ CACHE = os.environ.get("HF_HOME") or os.path.expanduser("~/.cache/huggingface")
124
+
125
+
126
+ def etat():
127
+ l = []
128
+ l.append("[VIE] %s pod=%s" % (time.strftime("%H:%M:%S", time.gmtime()), POD))
129
+ l.append(sh("nvidia-smi --query-gpu=name,memory.used,memory.total,"
130
+ "utilization.gpu --format=csv,noheader"))
131
+ lim, use = ram_cgroup()
132
+ if lim:
133
+ l.append("RAM conteneur : %.1f / %.1f Go (cgroup)"
134
+ % (use / 1e9, lim / 1e9))
135
+ elif use is not None:
136
+ l.append("RAM conteneur : %.1f Go utilises, aucune limite cgroup lisible"
137
+ % (use / 1e9))
138
+ try:
139
+ du = shutil.disk_usage("/")
140
+ l.append("disque / : %.0f Go libres sur %.0f" % (du.free / 1e9,
141
+ du.total / 1e9))
142
+ except Exception:
143
+ pass
144
+ try:
145
+ go = octets_repertoire(CACHE) / 1e9
146
+ l.append("poids telecharges : %.1f / %.1f Go (%.0f %%)"
147
+ % (go, ATTENDU_GO, 100 * go / ATTENDU_GO))
148
+ except Exception:
149
+ pass
150
+ for nom in ("dl.log", "vllm.log", "banc.log"):
151
+ p = os.path.join(TRAVAIL, nom)
152
+ if os.path.exists(p):
153
+ q = sh("tail -c 1500 %s | tr '\\r' '\\n' | grep -v '^$' | tail -4" % p)
154
+ if q:
155
+ l.append("--- %s ---\n%s" % (nom, q))
156
+ vivantes = [c for c in COURANTES if c["fil"].is_alive()]
157
+ for c in vivantes:
158
+ l.append("commande en cours : %s (depuis %d s)"
159
+ % (c["nom"], time.time() - c["t0"]))
160
+ return "\n".join(x for x in l if x)
161
+
162
+
163
+ # ------------------------------------------------------------- telechargement
164
+ CODE_DL = """
165
+ import os, time
166
+ from huggingface_hub import snapshot_download
167
+ t0 = time.time()
168
+ for essai in range(6):
169
+ try:
170
+ p = snapshot_download(%r, max_workers=8,
171
+ token=os.environ.get('HF_TOKEN'),
172
+ ignore_patterns=['*.pth', 'original/*'])
173
+ print('TELECHARGE en %%.0f s -> %%s' %% (time.time() - t0, p), flush=True)
174
+ break
175
+ except Exception as e:
176
+ print('essai', essai, type(e).__name__, str(e)[:200], flush=True)
177
+ time.sleep(10)
178
+ else:
179
+ print('TELECHARGEMENT ECHOUE', flush=True)
180
+ """
181
+
182
+
183
+ def telecharger():
184
+ chemin = os.path.join(TRAVAIL, "dl.py")
185
+ with open(chemin, "w") as f:
186
+ f.write(CODE_DL % MODELE)
187
+ with open(os.path.join(TRAVAIL, "dl.log"), "w") as f:
188
+ subprocess.run(["python3", chemin], stdout=f, stderr=subprocess.STDOUT)
189
+
190
+
191
+ # ------------------------------------------------------------------ commandes
192
+ def executer(nom, script):
193
+ chemin = os.path.join(TRAVAIL, "sorties", nom + ".log")
194
+ with open(chemin, "w") as f:
195
+ f.write("$ %s\n%s\n%s\n%s\n" % (nom, "-" * 70, script, "-" * 70))
196
+ f.flush()
197
+ p = subprocess.Popen(["bash", "-lc", script], stdout=f,
198
+ stderr=subprocess.STDOUT, cwd=TRAVAIL)
199
+ code = p.wait()
200
+ f.write("\n%s\n[FIN] code=%d\n" % ("-" * 70, code))
201
+ publier(chemin, "etat/out-%s.log" % POD)
202
+ print("commande %s terminee (code %d)" % (nom, code), flush=True)
203
+
204
+
205
+ def boucle_commandes():
206
+ vue = None
207
+ n = 0
208
+ while True:
209
+ texte = lire_hub("cmd/%s.sh" % POD)
210
+ if texte is not None and texte.strip():
211
+ h = hashlib.sha256(texte.encode()).hexdigest()[:12]
212
+ if h != vue:
213
+ vue = h
214
+ n += 1
215
+ nom = "%03d-%s" % (n, h)
216
+ print("commande recue %s" % nom, flush=True)
217
+ fil = threading.Thread(target=executer, args=(nom, texte),
218
+ daemon=True)
219
+ fil.start()
220
+ COURANTES.append({"nom": nom, "fil": fil, "t0": time.time()})
221
+ time.sleep(8)
222
+
223
+
224
+ def boucle_etat():
225
+ p = os.path.join(TRAVAIL, "vie.log")
226
+ while True:
227
+ try:
228
+ with open(p, "w") as f:
229
+ f.write(etat() + "\n")
230
+ publier(p, "etat/vie-%s.log" % POD)
231
+ except Exception as e:
232
+ print("etat:", e, flush=True)
233
+ time.sleep(30)
234
+
235
+
236
+ print("[AGENT] pod=%s modele=%s" % (POD, MODELE), flush=True)
237
+ threading.Thread(target=telecharger, daemon=True).start()
238
+ threading.Thread(target=boucle_commandes, daemon=True).start()
239
+ boucle_etat()