Add Deck dashboard with Athena telemetry history
This commit is contained in:
@@ -0,0 +1,268 @@
|
||||
"""Read-only Deck dashboard, continuing the old Athena telemetry history."""
|
||||
from __future__ import annotations
|
||||
|
||||
import http.client
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import urllib.parse
|
||||
from pathlib import Path
|
||||
import threading
|
||||
import time
|
||||
|
||||
from dashboard_history import HISTORY_INTERVAL, HistoryStore
|
||||
|
||||
METRICS = {
|
||||
'prompt_tokens_total', 'prompt_tokens_cached_total', 'prompt_seconds_total',
|
||||
'tokens_predicted_total', 'tokens_predicted_seconds_total', 'n_decode_total',
|
||||
'n_tokens_max', 'spec_decode_num_draft_tokens_total',
|
||||
'spec_decode_num_accepted_tokens_total', 'spec_decode_num_drafts_total',
|
||||
'prompt_tokens_seconds', 'predicted_tokens_seconds', 'requests_processing',
|
||||
'requests_deferred', 'n_busy_slots_per_decode',
|
||||
}
|
||||
|
||||
|
||||
def metrics(raw):
|
||||
values = {}
|
||||
for line in raw.splitlines():
|
||||
if not line.startswith('llamacpp:') or ' ' not in line:
|
||||
continue
|
||||
key, value = line.rsplit(None, 1)
|
||||
key = key.removeprefix('llamacpp:')
|
||||
if key not in METRICS or '{' in key:
|
||||
continue
|
||||
try:
|
||||
number = float(value)
|
||||
if number == number and abs(number) != float('inf'):
|
||||
values[key] = int(number) if number.is_integer() else number
|
||||
except ValueError:
|
||||
pass
|
||||
return values
|
||||
|
||||
|
||||
def slot_data(raw):
|
||||
result = []
|
||||
if not isinstance(raw, list):
|
||||
return result
|
||||
for slot in raw:
|
||||
if not isinstance(slot, dict):
|
||||
continue
|
||||
token = (slot.get('next_token') or [{}])[0]
|
||||
params = slot.get('params') or {}
|
||||
prompt = int(slot.get('n_prompt_tokens') or 0)
|
||||
decoded = int(token.get('n_decoded') or 0)
|
||||
context = int(slot.get('n_ctx') or 0)
|
||||
result.append(dict(id=slot.get('id'), task_id=slot.get('id_task'),
|
||||
processing=bool(slot.get('is_processing')),
|
||||
speculative=bool(slot.get('speculative')), n_ctx=context,
|
||||
prompt_tokens=prompt,
|
||||
prompt_processed=int(slot.get('n_prompt_tokens_processed') or 0),
|
||||
prompt_cached=int(slot.get('n_prompt_tokens_cache') or 0),
|
||||
decoded_tokens=decoded,
|
||||
context_used=min(context, prompt+decoded) if context else prompt+decoded,
|
||||
remaining_generation=token.get('n_remain'),
|
||||
max_tokens=params.get('max_tokens', params.get('n_predict')),
|
||||
temperature=params.get('temperature'), stream=params.get('stream')))
|
||||
return result
|
||||
|
||||
|
||||
class Dashboard:
|
||||
def __init__(self, server, state_dir):
|
||||
self.server = server
|
||||
self.started = time.time()
|
||||
self.history = HistoryStore(Path(state_dir)/'dashboard-history.sqlite3')
|
||||
self.lock = threading.Lock()
|
||||
self.cached = None
|
||||
self.checked = 0
|
||||
self.model_checked = 0
|
||||
self.model_cache = []
|
||||
self.events = []
|
||||
self.last_state = None
|
||||
self.closed = threading.Event()
|
||||
self.thread = threading.Thread(target=self._collect_loop, name='deck-dashboard-history', daemon=True)
|
||||
self.thread.start()
|
||||
|
||||
def close(self):
|
||||
self.closed.set()
|
||||
self.thread.join(timeout=5)
|
||||
if not self.thread.is_alive():self.history.close()
|
||||
|
||||
@staticmethod
|
||||
def backup_file(name):
|
||||
if not re.fullmatch(r'athena-portable-[A-Za-z0-9_.-]+\.tar\.zst\.age', name):
|
||||
raise ValueError('Ungültiger Backupname.')
|
||||
path=Path('/host/backups')/name
|
||||
if not path.is_file() or path.is_symlink():
|
||||
raise ValueError('Backup nicht verfügbar.')
|
||||
return path
|
||||
|
||||
def backups(self):
|
||||
folder=Path('/host/backups')
|
||||
try:
|
||||
paths=sorted(folder.glob('athena-portable-*.tar.zst.age'),key=lambda p:p.stat().st_mtime,reverse=True)[:5]
|
||||
except OSError:
|
||||
paths=[]
|
||||
items=[]
|
||||
for path in paths:
|
||||
try:
|
||||
self.backup_file(path.name)
|
||||
stat=path.stat()
|
||||
checksum=(path.parent/(path.name+'.sha256')).read_text().split()[0]
|
||||
if not re.fullmatch(r'[a-fA-F0-9]{64}',checksum):checksum=None
|
||||
items.append(dict(name=path.name,size=stat.st_size,modified=stat.st_mtime,
|
||||
sha256=checksum,download_url='/api/v1/dashboard/backups/download/'+urllib.parse.quote(path.name)))
|
||||
except (OSError,IndexError,ValueError):
|
||||
continue
|
||||
return dict(backups=items,available=folder.is_dir())
|
||||
|
||||
def model_files(self):
|
||||
now=time.monotonic()
|
||||
with self.lock:
|
||||
if now-self.model_checked<30:return list(self.model_cache)
|
||||
files=[]
|
||||
root=Path('/host/models')
|
||||
if root.is_dir():
|
||||
try:
|
||||
for path in sorted(root.rglob('*.gguf'))[:100]:
|
||||
if path.is_symlink():continue
|
||||
stat=path.stat()
|
||||
files.append(dict(name=path.name,relative_path=str(path.relative_to(root)),
|
||||
size=stat.st_size,modified=stat.st_mtime))
|
||||
except OSError:
|
||||
pass
|
||||
for entry in self.server.catalog.status()['entries']:
|
||||
if entry.get('file','').lower().endswith('.gguf'):
|
||||
relative=f"Deck/{entry.get('repo','Deck')}/{entry['file']}"
|
||||
if not any(f['relative_path']==relative for f in files):
|
||||
files.append(dict(name=entry['file'],relative_path=relative,
|
||||
size=entry.get('size'),modified=entry.get('downloaded_at')))
|
||||
if not root.is_dir():
|
||||
for path in Path('/reference-models').glob('*.gguf'):
|
||||
try:
|
||||
stat=path.stat()
|
||||
if not any(f['name']==path.name for f in files):
|
||||
files.append(dict(name=path.name,relative_path='Referenzmodell/'+path.name,
|
||||
size=stat.st_size,modified=stat.st_mtime))
|
||||
except OSError:pass
|
||||
with self.lock:
|
||||
self.model_cache=files
|
||||
self.model_checked=now
|
||||
return list(files)
|
||||
|
||||
def _read_worker(self, path):
|
||||
conn, key = self.server.worker.connect()
|
||||
conn.timeout = 2
|
||||
try:
|
||||
conn.request('GET', path, headers={'Authorization': 'Bearer '+key})
|
||||
response = conn.getresponse()
|
||||
body = response.read(2*1024*1024+1)
|
||||
if response.status != 200 or len(body) > 2*1024*1024:
|
||||
raise ValueError(f'{path}: HTTP {response.status}')
|
||||
return body
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
def telemetry(self, ready):
|
||||
result = dict(available=False, slots=[], metrics={}, props={})
|
||||
if not ready:
|
||||
return result
|
||||
for path, key, parse in (('/slots', 'slots', lambda raw: slot_data(json.loads(raw))),
|
||||
('/metrics', 'metrics', lambda raw: metrics(raw.decode('utf-8', 'replace'))),
|
||||
('/props', 'props', lambda raw: json.loads(raw))):
|
||||
try:
|
||||
data = parse(self._read_worker(path))
|
||||
if key == 'props':
|
||||
data = dict(total_slots=data.get('total_slots'), model_alias=data.get('model_alias'),
|
||||
model_ftype=data.get('model_ftype'), modalities=data.get('modalities') or {},
|
||||
default_context=(data.get('default_generation_settings') or {}).get('n_ctx'))
|
||||
result[key] = data
|
||||
except (OSError, ValueError, TypeError, KeyError, http.client.HTTPException):
|
||||
pass
|
||||
result['available'] = bool(result['slots'] or result['metrics'])
|
||||
return result
|
||||
|
||||
def snapshot(self):
|
||||
now = time.monotonic()
|
||||
with self.lock:
|
||||
if self.cached is not None and now-self.checked < 1:
|
||||
return self.cached
|
||||
server = self.server
|
||||
h = server.hardware.snapshot()
|
||||
worker = server.worker.status()
|
||||
scheduler = server.scheduler.status()
|
||||
ready = worker['state'] == 'ready'
|
||||
telemetry = self.telemetry(ready)
|
||||
with server.worker.lock:
|
||||
profile = server.worker.profile if ready else None
|
||||
pid = server.worker.process.pid if ready and server.worker.process else None
|
||||
params = (profile or {}).get('parameters') or {}
|
||||
model = (profile or {}).get('model') or {}
|
||||
model_file = model.get('file')
|
||||
gpu_names = {g['uuid']: f"GPU {g['index']}" for g in h.get('gpus', [])}
|
||||
runtime = dict(pid=pid, model_file=model_file, alias=(profile or {}).get('name'),
|
||||
context_size=params.get('context'), batch_size=params.get('batch'),
|
||||
ubatch_size=params.get('ubatch'), parallel=params.get('slots'),
|
||||
threads=params.get('threads'), threads_batch=params.get('threads_batch'),
|
||||
device=', '.join(gpu_names.get(g, g) for g in worker.get('gpus') or []),
|
||||
tensor_split=params.get('tensor_split'), cache_k='q4_0', cache_v='q4_0',
|
||||
flash_attention=True if ready else None, prompt_cache=params.get('cache_prompt'),
|
||||
kv_unified=True if ready else None, mtp_draft_tokens=params.get('mtp_draft_tokens'),
|
||||
reasoning_budget=params.get('reasoning_budget'))
|
||||
image = server.image_tests.status().get('job') or {}
|
||||
image_active = image.get('state') == 'running'
|
||||
image_model = image.get('profile_name') if image_active else None
|
||||
profile_name = worker.get('profile_name') if ready else None
|
||||
profiles = {p['name']: p['parameters'].get('context') for p in server.profiles.status()['profiles'] if p['kind'] == 'chat'}
|
||||
router = dict(current_profile=profile_name or 'nicht geladen', switching='ja' if scheduler['switching'] else None,
|
||||
upstream=dict(model=model_file, ctx=params.get('context'), reachable=ready),
|
||||
qwen=dict(available=ready, active_chats=scheduler['active_requests']),
|
||||
llama_telemetry=telemetry, profiles=profiles,
|
||||
uptime_seconds=time.time()-server.started,
|
||||
image=dict(phase='running' if image_active else 'idle', model=image_model,
|
||||
worker='Deck-Bildlaufzeit' if image_active else None,
|
||||
model_loaded=image_active, last_seconds=image.get('duration_seconds')),
|
||||
tts=dict(ready=server.tts_tests.status().get('loaded',False),engine='Deck-TTS'),
|
||||
stt=dict(reachable=server.stt.status().get('loaded',False)))
|
||||
cpu = h.get('cpu') or {}
|
||||
ram = h.get('ram') or {}
|
||||
cpu_data = dict(usage_percent=cpu.get('percent'), logical_cpus=cpu.get('logical_cpus'),
|
||||
load=cpu.get('load') or [], host_uptime_seconds=cpu.get('host_uptime_seconds'),
|
||||
memory=dict(total=ram.get('total_bytes'), used=ram.get('used_bytes')),
|
||||
disk_data=cpu.get('disk_data') or {}, network=cpu.get('network') or {})
|
||||
gpus = [dict(index=g.get('index'), name=g.get('name'), uuid=g.get('uuid'),
|
||||
gpu_percent=g.get('percent'), memory_total_mib=g.get('total_mib'),
|
||||
memory_used_mib=g.get('used_mib'), temperature_c=g.get('temperature_c'),
|
||||
memory_controller_percent=g.get('memory_controller_percent'),
|
||||
memory_free_mib=g.get('free_mib'), power_w=g.get('power_w'),
|
||||
power_limit_w=g.get('power_limit_w'),graphics_clock_mhz=g.get('graphics_clock_mhz'),
|
||||
memory_clock_mhz=g.get('memory_clock_mhz'),fan_percent=g.get('fan_percent'),
|
||||
pstate=g.get('pstate')) for g in h.get('gpus') or []]
|
||||
files = self.model_files()
|
||||
state = (profile_name, model_file, worker['state'])
|
||||
with self.lock:
|
||||
if self.last_state is not None and self.last_state != state:
|
||||
self.events.insert(0, dict(timestamp=time.time(), name='Deck-Modell',
|
||||
**{'from': self.last_state[0] or self.last_state[2],
|
||||
'to': profile_name or worker['state']}))
|
||||
self.events = self.events[:50]
|
||||
self.last_state = state
|
||||
result = dict(timestamp=time.time(),dashboard_uptime_seconds=time.time()-self.started,
|
||||
cpu=cpu_data,gpus=gpus,gpu_processes=h.get('gpu_processes') or [],
|
||||
llama_runtime=runtime,router=router,models=files,
|
||||
model_summary=dict(count=len(files),total_size=sum(f.get('size') or 0 for f in files)),
|
||||
events=list(self.events),errors={'hardware':'; '.join(h.get('errors') or []) or None})
|
||||
self.cached = result
|
||||
self.checked = time.monotonic()
|
||||
return result
|
||||
|
||||
def _collect_loop(self):
|
||||
next_compaction = 0
|
||||
while not self.closed.is_set():
|
||||
try:
|
||||
self.history.record(self.snapshot())
|
||||
if time.time() >= next_compaction:
|
||||
self.history.compact()
|
||||
next_compaction = time.time()+3600
|
||||
except Exception:
|
||||
pass # Telemetry must never interrupt inference or expose request contents.
|
||||
self.closed.wait(HISTORY_INTERVAL)
|
||||
Reference in New Issue
Block a user