#!/usr/bin/env python3 """Athena Deck prototype: loopback API, owned demo process, 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 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 DemoService: def __init__(self): self.lock = threading.RLock() self.process = None self.port = None def status(self): with self.lock: running = self.process is not None and self.process.poll() is None reachable = False if running: try: with urllib.request.urlopen(f'http://127.0.0.1:{self.port}/health', timeout=.5) as response: reachable = json.load(response).get('service') == 'athena-deck-demo' except (OSError, ValueError): pass return dict(state='running' if running else 'stopped', reachable=reachable, pid=self.process.pid if running else None, port=self.port if running else None, location='Deck · isolierter Demo-Prozess') def start(self): with self.lock: if self.process is None or self.process.poll() is not None: with socket.socket() as sock: sock.bind(('127.0.0.1', 0)) sock.listen(8) self.port = sock.getsockname()[1] self.process = subprocess.Popen([sys.executable, str(ROOT/'demo.py'), str(sock.fileno())], pass_fds=(sock.fileno(),), stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) for _ in range(30): if self.status()['reachable']: break time.sleep(.05) return self.status() def stop(self): with self.lock: if self.process is not None and self.process.poll() is None: self.process.terminate() try: self.process.wait(timeout=3) except subprocess.TimeoutExpired: self.process.kill() self.process.wait() return self.status() 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.demo = DemoService() 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.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','/api/v1/demo')) or (self.command == 'POST' and self.path in ('/api/v1/demo/start','/api/v1/demo/stop')) 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 = {'/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': return self.respond(dict(name='Athena Deck', version='0.6.0', state='ready', uptime_seconds=round(time.time()-self.server.started), mode='isolated', location=os.environ.get('DECK_LOCATION', 'Athena · Debian-Server'), demo=self.server.demo.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())) if self.path == '/api/v1/demo': return self.respond(self.server.demo.status()) 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/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) if self.path not in ('/api/v1/demo/start', '/api/v1/demo/stop'): return self.respond({'error': 'Not found'}, 404) action = self.server.demo.start if self.path.endswith('/start') else self.server.demo.stop result = action() self.respond(result, 503 if self.path.endswith('/start') and not result['reachable'] else 200) 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.image_runtime.stop() server.image_tests.stop() server.runtime.stop() server.catalog.stop() server.demo.stop() server.server_close() if __name__ == '__main__': main()