codemaivanngu commited on
Commit
fbf0ff1
·
verified ·
1 Parent(s): 45af830

Expose standalone telemetry helper for source 5bc0d84

Browse files
Files changed (1) hide show
  1. eval-telemetry-5bc0d84.py +116 -0
eval-telemetry-5bc0d84.py ADDED
@@ -0,0 +1,116 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ #!/usr/bin/env python3
2
+ """External eval telemetry: no queue/source mutation, secrets or answer text."""
3
+ import argparse
4
+ import csv
5
+ import datetime
6
+ import hashlib
7
+ import json
8
+ import os
9
+ from pathlib import Path
10
+ import re
11
+ import subprocess
12
+ import time
13
+
14
+
15
+ class Counter:
16
+ def __init__(self): self.files={}
17
+ def count(self,path):
18
+ stat=path.stat();identity=(stat.st_dev,stat.st_ino)
19
+ old=self.files.get(str(path),(identity,0,0,0))
20
+ _,offset,count,mtime=old
21
+ if old[0]!=identity or stat.st_size<offset or (stat.st_size==offset and stat.st_mtime_ns!=mtime): offset=count=0
22
+ with path.open('rb') as f:
23
+ f.seek(offset)
24
+ while chunk:=f.read(1024*1024): count+=chunk.count(b'\n')
25
+ offset=f.tell()
26
+ self.files[str(path)]=(identity,offset,count,stat.st_mtime_ns)
27
+ return count
28
+
29
+
30
+ def read(path):
31
+ try: return Path(path).read_text().strip()
32
+ except OSError: return None
33
+
34
+
35
+ def workers(plan):
36
+ found={}
37
+ table={}
38
+ for p in Path('/proc').iterdir():
39
+ if not p.name.isdigit(): continue
40
+ try:
41
+ fields=(p/'stat').read_text().rsplit(')',1)[1].split()
42
+ table[int(p.name)]={'ppid':int(fields[1]),'ticks':int(fields[11])+int(fields[12]),'start':fields[19],
43
+ 'rss_bytes':int(fields[21])*os.sysconf('SC_PAGE_SIZE')}
44
+ argv=(p/'cmdline').read_bytes().split(b'\0')
45
+ if b'worker' in argv and b'--plan' in argv and argv[argv.index(b'--plan')+1].decode()==str(plan): found[int(p.name)]={}
46
+ except (OSError,ValueError,IndexError,UnicodeError): continue
47
+ for pid in found:
48
+ family={pid}
49
+ while True:
50
+ larger=family|{k for k,v in table.items() if v['ppid'] in family}
51
+ if larger==family: break
52
+ family=larger
53
+ found[pid]={'processes':{str(k):table[k] for k in family if k in table},'descendants':len(family)-1}
54
+ return found
55
+
56
+
57
+ def gpu():
58
+ try:
59
+ raw=subprocess.check_output(['nvidia-smi','--query-gpu=index,utilization.gpu,utilization.memory,memory.used,memory.total,power.draw','--format=csv,noheader,nounits'],text=True,timeout=5)
60
+ keys=('gpu','util_pct','memory_util_pct','memory_used_mib','memory_total_mib','power_w')
61
+ return [dict(zip(keys,[v.strip() for v in row])) for row in csv.reader(raw.splitlines())]
62
+ except (OSError,subprocess.SubprocessError) as e: return {'error':type(e).__name__}
63
+
64
+
65
+ def sample(root,plan,counter):
66
+ cells=[]
67
+ for cell in sorted((root/'cells').glob('*/*/*')):
68
+ if not cell.is_dir(): continue
69
+ responses=counter.count(cell/'responses.jsonl') if (cell/'responses.jsonl').exists() else 0
70
+ scores=counter.count(cell/'scores.jsonl') if (cell/'scores.jsonl').exists() else 0
71
+ cells.append({'cell':str(cell.relative_to(root/'cells')),'responses':responses,'scores':scores,
72
+ 'waiting_for_score':responses-scores,'complete':(cell/'metrics.json').exists()})
73
+ server=[]
74
+ for log in root.glob('*-server.log'):
75
+ with log.open('rb') as f:
76
+ f.seek(max(0,log.stat().st_size-32768));lines=f.read().decode(errors='replace').splitlines()
77
+ for line in reversed(lines):
78
+ match=re.search(r'Decode batch, #running-req: (\d+).*gen throughput \(token/s\): ([\d.]+), #queue-req: (\d+)',line)
79
+ if match:
80
+ server.append({'log':log.name,'log_mtime':log.stat().st_mtime,'last_decode_time':line.split(']')[0].lstrip('['),
81
+ 'running_requests':int(match[1]),'tokens_per_second':float(match[2]),'queued_requests':int(match[3])});break
82
+ mem={}
83
+ for line in (read('/proc/meminfo') or '').splitlines():
84
+ k,v=line.split(':',1)
85
+ if k in ('MemTotal','MemAvailable','MemFree','Cached','SwapFree'): mem[k]=v.strip()
86
+ limits={name:read('/sys/fs/cgroup/'+name) for name in ('memory.current','memory.max','cpu.max','memory/memory.usage_in_bytes','memory/memory.limit_in_bytes','cpu/cpu.cfs_quota_us','cpu/cpu.cfs_period_us')}
87
+ return {'utc':datetime.datetime.now(datetime.timezone.utc).isoformat(),'epoch':time.time(),'gpu':gpu(),
88
+ 'cells':cells,'server_last_decode':server,'workers':workers(plan),'meminfo':mem,'cgroup':limits,
89
+ 'loadavg':read('/proc/loadavg'),'pressure':{k:read('/proc/pressure/'+k) for k in ('cpu','memory','io')},
90
+ 'cpu_stat':(read('/proc/stat') or '').splitlines()[0]}
91
+
92
+
93
+ def main():
94
+ p=argparse.ArgumentParser(description=__doc__);p.add_argument('--plan',type=Path,required=True)
95
+ p.add_argument('--out',type=Path,required=True);p.add_argument('--seconds',type=int,default=600);p.add_argument('--interval',type=float,default=5)
96
+ a=p.parse_args()
97
+ if not 5<=a.seconds<=3600 or not 2<=a.interval<=60: p.error('seconds 5..3600; interval 2..60')
98
+ plan=a.plan.resolve(strict=True);counter=Counter();last=None
99
+ with a.out.open('x') as output:
100
+ output.write(json.dumps({'type':'metadata','plan':str(plan),'plan_sha256':hashlib.sha256(plan.read_bytes()).hexdigest(),
101
+ 'cpu_affinity':len(os.sched_getaffinity(0)),'clock_ticks':os.sysconf('SC_CLK_TCK'),'scope':'External samples; row counts can race live writes, RSS double-counts shared pages; no answer text captured'})+'\n');output.flush()
102
+ end=time.monotonic()+a.seconds
103
+ while time.monotonic()<end:
104
+ try:
105
+ row=sample(plan.parent,plan,counter)
106
+ totals={k:sum(c[k] for c in row['cells']) for k in ('responses','scores','waiting_for_score')}
107
+ row['totals']=totals
108
+ if last:
109
+ dt=row['epoch']-last['epoch'];row['rates_per_second']={k:(totals[k]-last['totals'][k])/dt for k in ('responses','scores')}
110
+ output.write(json.dumps(row)+'\n');output.flush();last=row
111
+ except (OSError,ValueError) as e:
112
+ output.write(json.dumps({'type':'sample_error','error':type(e).__name__,'epoch':time.time()})+'\n');output.flush()
113
+ time.sleep(min(a.interval,max(0,end-time.monotonic())))
114
+ print('TELEMETRY_DONE='+str(a.out))
115
+
116
+ if __name__=='__main__': main()