269 lines
13 KiB
Python
269 lines
13 KiB
Python
"""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)
|