#!/usr/bin/env python3 """Athena Deck prototype: loopback API, owned model processes, read-only telemetry.""" import argparse import os import hashlib import secrets from http.cookies import SimpleCookie, CookieError from auth import CredentialStore, verify_password from catalog import Catalog from profiles import Profiles from capacity import assess, overview from runtime import Runtime from image_runtime import ImageRuntime from image_test import ImageTests from docker_support import DockerSupport from inference import LlamaWorker,Scheduler from endpoint import Endpoint from chat_test import ChatTests from auto_test import AutoTests from urllib.parse import urlsplit, parse_qs from network.client import NetworkClient from network.config import parse_config, ConfigError import json import signal import socket import subprocess import sys import threading import time import urllib.request from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path ROOT = Path(__file__).resolve().parent class HardwareProvider: def __init__(self): self.lock = threading.Lock() self.cached = None self.checked = 0 def snapshot(self): if os.environ.get('DECK_LOCAL_HARDWARE', '1') == '1': from collect_hardware import collect with self.lock: if self.cached is None or time.monotonic()-self.checked >= 5: self.cached = dict(collect(), available=True) self.checked = time.monotonic() return self.cached with self.lock: if time.monotonic() - self.checked < 5 and self.cached is not None: return self.cached try: result = subprocess.run(['ssh', '-i', str(Path.home()/'.ssh/athena_key'), '-o', 'BatchMode=yes', '-o', 'ConnectTimeout=4', '-o', 'StrictHostKeyChecking=yes', 'root@192.168.1.212', 'python3 -'], input=(ROOT/'collect_hardware.py').read_text(), capture_output=True, text=True, timeout=10, check=True) self.cached = json.loads(result.stdout) self.cached['available'] = True except (OSError, subprocess.SubprocessError, ValueError): self.cached = dict(available=False, sampled_at=None, cpu={}, ram={}, gpus=[], errors=['Athena per SSH nicht erreichbar oder Messung fehlgeschlagen.']) self.checked = time.monotonic() return self.cached class Server(ThreadingHTTPServer): daemon_threads = True def __init__(self, port, state_dir=None): super().__init__((os.environ.get('DECK_BIND_HOST','127.0.0.1'), port), Handler) self.catalog = Catalog(Path(state_dir or os.environ.get("DECK_STATE_DIR", ROOT/".state"))/"models") self.runtime = Runtime(Path(state_dir or os.environ.get("DECK_STATE_DIR", ROOT/".state"))/"runtime") self.profiles = Profiles(self.catalog.root.parent/"profiles.json",self.catalog) self.image_tests = ImageTests(self.catalog.root.parent/"image-tests",self.profiles) self.image_runtime = ImageRuntime(self.catalog.root.parent/"image-runtime") self.image_tests.runtime = self.image_runtime self.profiles.image_runtime_ready=lambda:self.image_tests.status()["runtime_installed"] self.docker = DockerSupport() self.hardware = HardwareProvider() self.started = time.time() self.network = NetworkClient() if os.environ.get('DECK_AUTH_RPC') == '1': from network.rpc import request self.credentials = CredentialStore(rpc=request) else: self.credentials = CredentialStore(Path(state_dir or os.environ.get('DECK_STATE_DIR', ROOT/'.state'))/'auth.json') if os.environ.get('DECK_REQUIRE_SETUP') == '1' and self.credentials.read() is None: self.server_close() raise RuntimeError('Serverzugang muss vor dem Start eingerichtet werden.') self.worker=LlamaWorker(self.catalog.root.parent/'llama-worker',self.catalog,self.runtime) self.scheduler=Scheduler(self.worker) self.profiles.chat_blockers=self.worker.blockers self.image_tests.acquire=self.scheduler.image_reservation self.endpoint=Endpoint(self.catalog.root.parent,self.profiles,self.worker,self.scheduler,self.image_tests,self.credentials,self.server_port) if self.endpoint.config['autostart']: try:self.endpoint.start() except ValueError as exc:self.endpoint.error=str(exc) self.chat_tests=ChatTests(self.profiles,self.worker,self.scheduler) self.auto_tests=AutoTests(self.catalog.root.parent/'auto-tests.json',self.profiles,self.worker,self.scheduler) self.sessions = {} self.login_attempts = [] self.auth_lock = threading.Lock() class Handler(BaseHTTPRequestHandler): def log_message(self, *args): pass def respond(self, data, status=200, mime='application/json'): body = json.dumps(data).encode() if mime == 'application/json' else data self.send_response(status) self.send_header('Content-Type', mime) self.send_header('Content-Length', str(len(body))) self.send_header('Cache-Control', 'no-store') self.send_header('X-Content-Type-Options', 'nosniff') self.send_header('Content-Security-Policy', "default-src 'self'; script-src 'self'; style-src 'self'; frame-ancestors 'none'") self.end_headers() self.wfile.write(body) def ingress(self): token = os.environ.get('DECK_PROXY_TOKEN', '') if token and secrets.compare_digest(self.headers.get('X-Deck-Proxy', ''), token): return self.headers.get('X-Deck-Ingress', 'management') return 'management' def api_token_authenticated(self): header = self.headers.get('Authorization', '') if not header.startswith('Bearer ') or len(header)>300: return False record = self.server.credentials.read() if not record or not record['api_token_hash']: return False actual = hashlib.sha256(header[7:].encode()).hexdigest() return secrets.compare_digest(actual,record['api_token_hash']) def token_route(self): return self.command == 'GET' and self.path in ('/api/v1/status','/api/v1/hardware') def session_token(self): try: return SimpleCookie(self.headers.get('Cookie',''))['deck_session'].value except (KeyError, ValueError, CookieError): return None def authenticated(self, session_only=False): if self.headers.get('Authorization'): return not session_only and self.token_route() and self.api_token_authenticated() record = self.server.credentials.read() if not record: return False with self.server.auth_lock: session = self.server.sessions.get(self.session_token()) return bool(session and session['expires']>time.monotonic() and secrets.compare_digest(session['revision'],record['password']['hash'])) def allow_password_attempt(self): with self.server.auth_lock: now = time.monotonic() self.server.login_attempts = [t for t in self.server.login_attempts if now-t<60] if len(self.server.login_attempts)>=10: self.respond({'error':'Zu viele Versuche. Bitte eine Minute warten.'},429) return False self.server.login_attempts.append(now) return True def read_json(self): self.connection.settimeout(10) if self.headers.get('Transfer-Encoding'): raise ValueError('Chunked Uploads werden nicht unterstützt.') length = int(self.headers.get('Content-Length', '0')) if not 0 < length <= 32768 or self.headers.get('Content-Type') != 'application/json': raise ValueError('JSON-Anfrage erwartet, maximal 32 KiB.') raw = self.rfile.read(length) try: value = json.loads(raw) except (ValueError, UnicodeDecodeError): raise ValueError('Ungültige JSON-Anfrage.') from None if not isinstance(value, dict): raise ValueError('JSON-Objekt erwartet.') return value def login(self): if not self.allow_password_attempt(): return try: data = self.read_json() with self.server.credentials.lock: record = self.server.credentials.read() if not record or not verify_password(data.get('password'),record['password']): return self.respond({'error':'Anmeldung fehlgeschlagen.'},401) token = secrets.token_urlsafe(32) now = time.monotonic() with self.server.auth_lock: self.server.sessions = {k:v for k,v in self.server.sessions.items() if v['expires']>now} if len(self.server.sessions)>=64: self.server.sessions.pop(next(iter(self.server.sessions))) self.server.sessions[token] = dict(expires=now+28800,revision=record['password']['hash']) body = b'{"authenticated":true}' self.send_response(200) self.send_header('Content-Type','application/json') self.send_header('Cache-Control','no-store') self.send_header('Content-Length',str(len(body))) self.send_header('Set-Cookie',f'deck_session={token}; HttpOnly; SameSite=Strict; Path=/; Max-Age=28800') self.end_headers() self.wfile.write(body) except (ValueError, OSError): self.respond({'error':'Ungültige Anmeldung.'},400) def access_status(self): record = self.server.credentials.read() signed_in = self.authenticated(session_only=True) value = dict(initialized=record is not None, authenticated=signed_in) if signed_in: value.update(api_token_configured=bool(record['api_token_hash']), password_changed_at=record['password_changed_at'], token_changed_at=record['token_changed_at'], account='Administrator', location=os.environ.get('DECK_LOCATION', 'Athena · Debian-Server')) return self.respond(value) def change_credentials(self, kind): if not self.allow_password_attempt(): return try: data = self.read_json() key = 'new_password' if kind == 'password' else 'new_token' if set(data) != {'current_password',key}: raise ValueError('Ungültige Zugangsdatenanfrage.') self.server.credentials.change(kind,data['current_password'],data[key]) if kind == 'password': with self.server.auth_lock: self.server.sessions.clear() return self.respond({'changed':True,'login_required':kind=='password'}) except ValueError as exc: return self.respond({'error':str(exc)},400) except OSError: return self.respond({'error':'Zugangsdaten konnten nicht gespeichert werden.'},503) def valid_host(self): if os.environ.get('DECK_PROXY_TOKEN'): return self.ingress() in ('lan','tunnel') allowed = {f'127.0.0.1:{self.server.server_port}', f'localhost:{self.server.server_port}'} allowed.update(filter(None,os.environ.get('DECK_ALLOWED_HOSTS','').split(','))) return self.headers.get('Host') in allowed def do_GET(self): if not self.valid_host(): return self.respond({'error': 'Host rejected'}, 403) if self.path == '/api/v1/auth/status': return self.access_status() if self.path == '/login.js': return self.respond((ROOT/'login.js').read_bytes(), mime='text/javascript') if not self.authenticated() and self.path not in ('/', '/style.css'): return self.respond({'error':'Anmeldung erforderlich.'},401) if not self.authenticated() and self.path == '/': return self.respond((ROOT/'login.html').read_bytes(), mime='text/html; charset=utf-8') routes = {'/auto-test-ui.js':('auto-test-ui.js','text/javascript'),'/chat-test-ui.js':('chat-test-ui.js','text/javascript'),'/endpoint-ui.js':('endpoint-ui.js','text/javascript'),'/docker-ui.js': ('docker-ui.js','text/javascript'), '/': ('index.html', 'text/html; charset=utf-8'), '/app.js': ('app.js', 'text/javascript'), '/style.css': ('style.css', 'text/css'), '/network-ui.js': ('network-ui.js', 'text/javascript'), '/access-ui.js': ('access-ui.js', 'text/javascript'), '/studio.js': ('studio.js', 'text/javascript'), '/catalog-ui.js': ('catalog-ui.js','text/javascript'), '/runtime-ui.js': ('runtime-ui.js','text/javascript'), '/profiles-ui.js': ('profiles-ui.js','text/javascript'), '/image-test-ui.js': ('image-test-ui.js','text/javascript')} if self.path in routes: name, mime = routes[self.path] return self.respond((ROOT/name).read_bytes(), mime=mime) if self.path == '/api/v1/status': endpoint=self.server.endpoint.status() public_endpoint={key:endpoint[key] for key in ('state','port','counts','active_requests')} return self.respond(dict(name='Athena Deck', version='0.7.0', state='ready', uptime_seconds=round(time.time()-self.server.started), mode='isolated', location=os.environ.get('DECK_LOCATION', 'Athena · Debian-Server'), endpoint=public_endpoint)) if self.path == '/api/v1/auto-tests':return self.respond(self.server.auto_tests.status()) if self.path == '/api/v1/chat-tests':return self.respond(self.server.chat_tests.status()) if self.path == '/api/v1/endpoint':return self.respond(self.server.endpoint.status()) if self.path == '/api/v1/docker':return self.respond(self.server.docker.status()) if self.path == '/api/v1/image-runtime':return self.respond(self.server.image_runtime.status()) if self.path == '/api/v1/image-tests':return self.respond(self.server.image_tests.status()) if urlsplit(self.path).path == '/api/v1/image-tests/image': try:return self.respond(self.server.image_tests.image(parse_qs(urlsplit(self.path).query).get('id',[''])[0]),mime='image/png') except (OSError,ValueError):return self.respond({'error':'Bild nicht verfügbar.'},404) if urlsplit(self.path).path == '/api/v1/profiles/components': try:return self.respond(self.server.profiles.components(parse_qs(urlsplit(self.path).query).get('model_id',[''])[0])) except ValueError as exc:return self.respond({'error':str(exc)},400) except Exception:return self.respond({'error':'Komponentenquelle nicht erreichbar. Bitte erneut versuchen.'},503) if self.path == '/api/v1/profiles': return self.respond(self.server.profiles.status()) if self.path == '/api/v1/hardware': return self.respond(self.server.hardware.snapshot()) if self.path.startswith('/api/v1/runtime'): actions={'/api/v1/runtime':self.server.runtime.status,'/api/v1/runtime/prerequisites':self.server.runtime.prerequisites,'/api/v1/runtime/releases':self.server.runtime.releases,'/api/v1/runtime/log':self.server.runtime.log,'/api/v1/runtime/references':self.server.runtime.references} if self.path in actions: try:return self.respond(actions[self.path]()) except Exception:return self.respond({'error':'Laufzeitabfrage fehlgeschlagen; Verbindung oder Werkzeuge prüfen.'},503) if self.path.startswith('/api/v1/catalog'): try: path=urlsplit(self.path);q=parse_qs(path.query) if path.path=='/api/v1/catalog/search': base=q.get('base_only',['false'])[0] if base not in ('true','false'):raise ValueError('Ungültiger Basismodellfilter.') return self.respond(self.server.catalog.search(q.get('q',[''])[0],q.get('kind',['chat'])[0],q.get('sort',['downloads'])[0],base=='true')) if path.path=='/api/v1/catalog/assessment': data=self.server.catalog.files(q.get('repo',[''])[0]) return self.respond(overview(data['files'],self.server.hardware.snapshot())) if path.path=='/api/v1/catalog/files': data=dict(self.server.catalog.files(q.get('repo',[''])[0]));hardware=self.server.hardware.snapshot() data['files']=[dict(f,assessment=assess(f['size'],hardware,f['name'])) for f in data['files']] return self.respond(data) if path.path=='/api/v1/catalog':return self.respond(self.server.catalog.status()) except Exception as exc: return self.respond({'error':str(exc) if isinstance(exc,ValueError) else 'Katalog nicht erreichbar.'},400) if self.path == '/api/v1/network': return self.respond(self.server.network.status(self.ingress())) return self.respond({'error': 'Not found'}, 404) def do_POST(self): origin = self.headers.get('Origin') bearer = self.token_route() and self.api_token_authenticated() if not self.valid_host() or (origin and origin != 'http://' + self.headers.get('Host', '')) or (not bearer and self.headers.get('X-Athena-Deck') != '1'): return self.respond({'error':'Same-origin control required'},403) if self.path == '/api/v1/auth/setup': try: data = self.read_json() if set(data) != {'password','api_token'}: raise ValueError('Kennwort und API-Token werden benötigt.') # Container installs always arrive provisioned. Never offer remote claim-on-first-use. if os.environ.get('DECK_AUTH_RPC') or os.environ.get('DECK_REQUIRE_SETUP') == '1': return self.respond({'error':'Serverzugang wird bei der Installation eingerichtet.'},403) self.server.credentials.setup(data['password'],data['api_token']) return self.respond({'initialized':True}) except ValueError as exc: return self.respond({'error':str(exc)},400) except OSError: return self.respond({'error':'Einrichtung konnte nicht gespeichert werden.'},503) if self.path == '/api/v1/login': return self.login() if not self.authenticated(): return self.respond({'error':'Anmeldung oder gültiger API-Token erforderlich.'},401) if self.path in ('/api/v1/auth/password','/api/v1/auth/token'): return self.change_credentials('password' if self.path.endswith('/password') else 'token') if self.path == '/api/v1/logout': with self.server.auth_lock: self.server.sessions.pop(self.session_token(),None) return self.respond({'logged_out':True}) if self.path == '/api/v1/docker/install': try: if self.read_json()!={'confirm':True}:raise ValueError('Die Docker-Erstinstallation muss ausdrücklich bestätigt werden.') return self.respond(self.server.docker.install()) except (OSError,ValueError) as exc:return self.respond({'error':str(exc) if isinstance(exc,ValueError) else 'Docker-Systemhelfer nicht erreichbar.'},400) if self.path.startswith('/api/v1/runtime/'): try: data=self.read_json();action=self.path.rsplit('/',1)[1] if action=='build' and set(data)=={'revision','backend','jobs'}:return self.respond(self.server.runtime.start(**data)) if action=='fit' and set(data)=={'profile','context','slots'}:return self.respond(self.server.runtime.fit(**data)) if action=='cancel' and not data:return self.respond(self.server.runtime.stop()) if action=='activate' and set(data)=={'build_id'}:return self.respond(self.server.runtime.activate(**data)) if action=='rollback' and not data:return self.respond(self.server.runtime.rollback()) raise ValueError('Ungültige Laufzeitaktion.') except ValueError as exc:return self.respond({'error':str(exc)},400) except (OSError, subprocess.SubprocessError):return self.respond({'error':'Laufzeitaktion fehlgeschlagen; Speicher und Werkzeuge prüfen.'},503) if self.path in ('/api/v1/auto-tests/start','/api/v1/auto-tests/cancel','/api/v1/auto-tests/save'): try: data=self.read_json() if self.path.endswith('/start'):return self.respond(self.server.auto_tests.start(data)) if self.path.endswith('/save'):return self.respond(self.server.auto_tests.save(data)) if data:raise ValueError('Keine Parameter erwartet.') return self.respond(self.server.auto_tests.stop()) except ValueError as exc:return self.respond({'error':str(exc)},400) if self.path in ('/api/v1/chat-tests/start','/api/v1/chat-tests/cancel','/api/v1/chat-tests/unload'): try: data=self.read_json() if self.path.endswith('/start'):return self.respond(self.server.chat_tests.start(data)) if data:raise ValueError('Keine Parameter erwartet.') return self.respond(self.server.chat_tests.stop() if self.path.endswith('/cancel') else self.server.chat_tests.unload()) except ValueError as exc:return self.respond({'error':str(exc)},400) if self.path in ('/api/v1/endpoint/start','/api/v1/endpoint/stop','/api/v1/endpoint/config','/api/v1/endpoint/profile'): try: data=self.read_json();ep=self.server.endpoint if self.path.endswith('/config'):return self.respond(ep.configure(data)) if self.path.endswith('/profile'):return self.respond(ep.enable(data)) if data:raise ValueError('Keine Parameter erwartet.') return self.respond(ep.start() if self.path.endswith('/start') else ep.stop()) except ValueError as exc:return self.respond({'error':str(exc)},400) except OSError:return self.respond({'error':'Endpunkt konnte nicht konfiguriert werden.'},503) if self.path in ('/api/v1/image-runtime/install','/api/v1/image-runtime/cancel'): try: if self.read_json():raise ValueError('Keine Parameter erwartet.') with self.server.image_tests.lock: if self.path.endswith('/cancel'):return self.respond(self.server.image_runtime.stop()) job=self.server.image_tests.job if job and job.get('state')=='running':raise ValueError('Zuerst den laufenden Bildauftrag beenden.') return self.respond(self.server.image_runtime.start()) except ValueError as exc:return self.respond({'error':str(exc)},400) except OSError:return self.respond({'error':'Installation konnte nicht vorbereitet werden.'},503) if self.path in ('/api/v1/image-tests/start','/api/v1/image-tests/cancel'): try: data=self.read_json() if self.path.endswith('/cancel') and not data:return self.respond(self.server.image_tests.stop()) if self.path.endswith('/start') and set(data)=={'profile_id','prompt'}:return self.respond(self.server.image_tests.start(**data)) raise ValueError('Ungültige Bildtest-Anfrage.') except ValueError as exc:return self.respond({'error':str(exc)},400) except Exception:return self.respond({'error':'Bildtest konnte nicht vorbereitet werden; Ressourcen prüfen.'},503) if self.path == '/api/v1/profiles/components': try:return self.respond(self.server.profiles.assign(self.read_json())) except ValueError as exc:return self.respond({'error':str(exc)},400) except OSError:return self.respond({'error':'Zuordnung konnte nicht gespeichert werden.'},503) if self.path == '/api/v1/profiles/delete': try: data=self.read_json() with self.server.image_tests.lock: job=self.server.image_tests.job if job and job.get('state')=='running' and job.get('profile_id')==data.get('id'):raise ValueError('Dieses Profil wird gerade ausgeführt. Zuerst den Auftrag beenden.') return self.respond(self.server.profiles.delete(data)) except ValueError as exc:return self.respond({'error':str(exc)},400) except OSError:return self.respond({'error':'Profil konnte nicht gelöscht werden.'},503) if self.path == '/api/v1/profiles/save': try:return self.respond(self.server.profiles.save(self.read_json())) except ValueError as exc:return self.respond({'error':str(exc)},400) except OSError:return self.respond({'error':'Profil konnte nicht gespeichert werden.'},503) if self.path in ('/api/v1/catalog/download','/api/v1/catalog/cancel','/api/v1/catalog/dismiss'): try: data=self.read_json() if self.path.endswith('/dismiss') and set(data)=={'id'}:return self.respond(self.server.catalog.dismiss(data['id'])) if self.path.endswith('/cancel') and not data:return self.respond(self.server.catalog.stop()) if set(data)!={'repo','filename','revision','kind'}:raise ValueError('Ungültige Downloadparameter.') return self.respond(self.server.catalog.start(**data)) except Exception as exc: return self.respond({'error':str(exc) if isinstance(exc,ValueError) else 'Download konnte nicht vorbereitet werden.'},400) if self.path.startswith('/api/v1/network/'): action = self.path.rsplit('/',1)[1] if action not in ('install','import','connect','disconnect','mode','confirm','cancel','delete'): return self.respond({'error':'Not found'},404) try: data = self.read_json() if action == 'install': result = self.server.network.install(data.get('password'),data.get('api_token')) else: fields = {'import': {'config'}, 'mode': {'mode'}, 'confirm': {'trial_id'}}.get(action,set()) if set(data) != fields: raise ValueError('Ungültige Netzwerkparameter.') if action == 'import': parse_config(data.get('config')) result = self.server.network.call(action, data, self.ingress()) self.server.network.cached = None return self.respond(result) except (ValueError, OSError, subprocess.SubprocessError) as exc: message = str(exc) if isinstance(exc, ValueError) else 'Netzwerkmodul nicht erreichbar.' return self.respond({'error':message},400) return self.respond({'error':'Not found'},404) def main(): parser = argparse.ArgumentParser() parser.add_argument('--port', type=int, default=8108) args = parser.parse_args() server = Server(args.port) # Binding fails rather than taking over an occupied port. def shutdown(*_): threading.Thread(target=server.shutdown, daemon=True).start() signal.signal(signal.SIGTERM, shutdown) signal.signal(signal.SIGINT, shutdown) print(f'Athena Deck: http://127.0.0.1:{server.server_port}', flush=True) try: server.serve_forever() finally: server.auto_tests.stop() server.chat_tests.stop() server.endpoint.close() server.image_runtime.stop() server.image_tests.stop() server.runtime.stop() server.catalog.stop() server.server_close() if __name__ == '__main__': main()