"""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)