"""Separate authenticated OpenAI-compatible API; only explicitly enabled Deck profiles.""" import base64 import hashlib import json import os from pathlib import Path import secrets import re import socket import threading import time from http.server import BaseHTTPRequestHandler,ThreadingHTTPServer from api_compat import normalize_chat,CompatibilityError from inference import InferenceError from stt import read_upload class APIError(ValueError): def __init__(self,message,status=400,code='invalid_request_error'): super().__init__(message);self.status=status;self.code=code class Endpoint: def __init__(self,root,profiles,worker,scheduler,images,credentials,management_port): self.root=Path(root);self.profiles=profiles;self.worker=worker;self.scheduler=scheduler;self.images=images;self.tts=None;self.stt=None;self.video=None;self.credentials=credentials self.management_port=management_port;self.lock=threading.RLock();self.http=None;self.thread=None;self.state='stopped';self.error=None;self.inflight=0 self.allowed_ports=[int(p) for p in os.environ.get('DECK_API_PORTS','').split(',') if p] self.config=dict(port=self.allowed_ports[0] if self.allowed_ports else 8120,enabled_profiles=[],autostart=False) path=self.root/'endpoint.json' if path.exists():self.config.update(json.loads(path.read_text())) image_ids={p['id'] for p in self.profiles.status()['profiles'] if p['kind']=='image'} selected=[i for i in self.config['enabled_profiles'] if i in image_ids] if len(selected)>1: self.config['enabled_profiles']=[i for i in self.config['enabled_profiles'] if i not in image_ids or i==selected[0]];self.persist() def persist(self): self.root.mkdir(parents=True,exist_ok=True,mode=0o700) path=self.root/'endpoint.tmp';path.write_text(json.dumps(self.config));path.replace(self.root/'endpoint.json') def rows(self): rows=self.profiles.status()['profiles'] with self.lock:enabled=set(self.config['enabled_profiles']) return [dict(p,enabled=p['id'] in enabled) for p in rows] def status(self): rows=self.rows();worker=self.worker.status();job=self.images.status()['job'];counts={} for key,kind in [('llm','chat'),('image','image'),('tts','audio'),('stt','stt')]: subset=[p for p in rows if p['kind']==kind] counts[key]=dict(enabled=sum(p['enabled'] for p in subset),available=sum(p['enabled'] and p['runnable'] for p in subset),loaded=bool(worker['state']=='ready' and kind=='chat') if kind=='chat' else bool(self.tts and self.tts.status()['job'] and self.tts.status()['job']['state']=='running') if kind=='audio' else bool(self.stt and self.stt.status()['job'] and self.stt.status()['job']['state']=='running') if kind=='stt' else bool(kind=='image' and job and job['state']=='running'),supported=kind in ('chat','image') or (kind=='audio' and self.tts is not None) or (kind=='stt' and self.stt is not None)) scheduler_status=self.scheduler.status() with self.lock:return dict(video=self.video.status() if self.video else None,state=self.state,reachable=bool(self.thread and self.thread.is_alive() and self.state=='running'),port=self.config['port'],bind=os.environ.get('DECK_API_BIND','127.0.0.1'),base_url=f"http://127.0.0.1:{self.config['port']}/v1",allowed_ports=self.allowed_ports,error=self.error,counts=counts,worker=worker,scheduler=scheduler_status,profiles=[dict(id=p['id'],name=p['name'],kind=p['kind'],enabled=p['enabled'],runnable=p['runnable'],blockers=p['blockers']) for p in rows],active_requests=self.inflight) def configure(self,data): if set(data)!={'port'} or type(data['port']) is not int or not 1024<=data['port']<=65535:raise ValueError('Port zwischen 1024 und 65535 erforderlich.') port=data['port'] with self.lock: if self.state!='stopped':raise ValueError('Endpunkt zuerst vollständig stoppen.') if port==self.management_port:raise ValueError('Der Verwaltungsport ist bereits belegt.') if self.allowed_ports and port not in self.allowed_ports:raise ValueError('Dieser Port ist im Docker-Installer nicht freigegeben.') try: with socket.socket() as sock:sock.bind((os.environ.get('DECK_API_BIND','127.0.0.1'),port)) except OSError:raise ValueError('Port ist bereits belegt.') from None self.config['port']=port;self.persist() return self.status() def enable(self,data): if set(data)!={'id','enabled'} or type(data['enabled']) is not bool:raise ValueError('Profil-ID und Aktivierung erforderlich.') rows=self.rows() row=next((p for p in rows if p['id']==data['id']),None) if not row:raise ValueError('Profil nicht gefunden.') if row['kind']=='video':raise ValueError('Videoprofil in der Übersicht auswählen; Video nutzt eine exklusive Auswahl.') if data['enabled'] and not row['runnable']:raise ValueError('Profil nicht ausführbar: '+' '.join(row['blockers'])) with self.lock: enabled=set(self.config['enabled_profiles']) if data['enabled']: if row['kind']=='image':enabled.difference_update(p['id'] for p in rows if p['kind']=='image') enabled.add(data['id']) else:enabled.discard(data['id']) self.config['enabled_profiles']=sorted(enabled);self.persist() return self.status() def start(self): with self.lock: if self.state=='running':return {'state':'running','port':self.config['port']} if self.state!='stopped':raise ValueError('Endpunkt wird noch gestoppt.') record=self.credentials.read() if not record or not record.get('api_token_hash'):raise ValueError('Zuerst unter Zugang & API einen API-Token einrichten.') if self.allowed_ports and self.config['port'] not in self.allowed_ports:raise ValueError('Gespeicherter Port ist nicht im Docker-Installer freigegeben.') if self.config['port']==self.management_port:raise ValueError('Verwaltungsport kann nicht als API-Port verwendet werden.') try:http=APIHTTPServer((os.environ.get('DECK_API_BIND','127.0.0.1'),self.config['port']),APIHandler) except OSError:raise ValueError('API-Port bereits belegt; kein anderer Dienst wurde verändert.') from None http.endpoint=self;self.http=http;self.state='running';self.error=None self.thread=threading.Thread(target=http.serve_forever,daemon=True);self.thread.start() self.config['autostart']=True;self.persist() return self.status() def stop(self): with self.lock: self.config['autostart']=False;self.persist() if self.state=='running': self.state='stopping';threading.Thread(target=self._drain,daemon=True).start() return self.status() def _drain(self): http=self.http if http:http.shutdown();http.server_close() while True: if self.video: vs=self.video.status() if vs['state']=='switching':time.sleep(.1);continue if vs['mode']=='video':self.video.switch('llm');time.sleep(.1);continue with self.lock:pending=self.inflight if not pending and self.scheduler.unload_idle():break time.sleep(.1) with self.lock:self.http=None;self.thread=None;self.state='stopped' def close(self): # Process shutdown does not change the user's autostart preference. with self.lock:self.state='stopping';http=self.http if http:http.shutdown();http.server_close() if self.tts:self.tts.stop() if self.stt:self.stt.stop() self.images.stop();self.worker.stop() def allowed(self): with self.lock:return self.state=='running' def authenticate(self,header): if not header.startswith('Bearer ') or len(header)>300:return False record=self.credentials.read() return bool(record and record.get('api_token_hash') and secrets.compare_digest(hashlib.sha256(header[7:].encode()).hexdigest(),record['api_token_hash'])) def find_profile(self,name,kind): if not isinstance(name,str):raise APIError('model muss den API-Namen eines aktivierten Profils enthalten.') row=next((p for p in self.rows() if (p['name']==name or (kind=='image' and name=='athena-image')) and p['kind']==kind and p['enabled']),None) if not row:raise APIError('Modellprofil nicht aktiviert oder unbekannt.',404,'model_not_found') if not row['runnable']:raise APIError('Profil derzeit nicht ausführbar: '+' '.join(row['blockers']),503,'model_unavailable') return row def model_list(self,kind="chat"): return dict(object='list',data=[dict(id='athena-image' if kind=='image' else p['name'],object='model',created=int(p['updated_at']),owned_by='athena-deck') for p in self.rows() if p['enabled'] and p['runnable'] and p['kind']==kind]) class APIHTTPServer(ThreadingHTTPServer): daemon_threads=True def __init__(self,*args,**kwargs): self.admission=threading.BoundedSemaphore(32);super().__init__(*args,**kwargs) def process_request(self,request,address): if not self.admission.acquire(blocking=False):self.shutdown_request(request);return try:super().process_request(request,address) except Exception:self.admission.release();raise def process_request_thread(self,*args): try:super().process_request_thread(*args) finally:self.admission.release() def handle_error(self,*args):pass # No request bodies, prompts, or traces in logs. class APIHandler(BaseHTTPRequestHandler): protocol_version='HTTP/1.1' def setup(self):super().setup();self.connection.settimeout(15);self.sent=False def log_message(self,*args):pass def compatibility_headers(self): info=getattr(self,'compatibility',{}) if info: self.send_header('X-Athena-Reasoning-Requested',info['requested']) self.send_header('X-Athena-Reasoning-Effective',info['effective']) self.send_header('X-Athena-Reasoning-Semantics',info['semantics']) def send(self,payload,status=200): body=json.dumps(payload).encode();self.sent=True self.send_response(status);self.compatibility_headers();self.send_header('Content-Type','application/json');self.send_header('Content-Length',str(len(body)));self.send_header('Cache-Control','no-store');self.send_header('Connection','close');self.end_headers();self.wfile.write(body);self.close_connection=True def failure(self,exc): if self.sent:return self.send({'error':{'message':str(exc),'type':getattr(exc,'code','server_error'),'param':None,'code':getattr(exc,'code','worker_unavailable')}},getattr(exc,'status',503)) def video_route(self,ep): if not self.path.startswith('/v1/videos'):return False if not ep.video:raise APIError('Video-Worker nicht eingerichtet.',503) if self.command=='GET' and self.path=='/v1/videos/models': ident=ep.video.selected self.send({'object':'list','data':[{'id':'athena-video','object':'model','owned_by':'athena-deck'}] if any(p['id']==ident and p['runnable'] for p in ep.rows()) else []});return True if self.command=='POST' and self.path=='/v1/videos': if self.headers.get('Transfer-Encoding'):raise APIError('Chunked Upload nicht unterstützt.') try:length=int(self.headers.get('Content-Length','0')) except ValueError:raise APIError('Ungültige Länge.') if not 01:raise APIError('Zunächst ein Bild pro Anfrage unterstützt.') # Re-resolve after waiting so edits/disable cannot silently launch a stale profile. key=('chat',profile['id'],profile['revision']) def allowed():return ep.allowed() and any(p['id']==profile['id'] and p['revision']==profile['revision'] and p['enabled'] for p in ep.rows()) with ep.scheduler.lease(key,profile['parameters']['slots'],lambda:ep.worker.ensure(profile),allowed=allowed): body=dict(data) for field in ('temperature','top_p','top_k'):body.setdefault(field,profile['parameters'][field]) conn,key=ep.worker.connect() try: conn.request('POST','/v1/chat/completions',body=json.dumps(body).encode(),headers={'Content-Type':'application/json','Authorization':'Bearer '+key}) response=conn.getresponse() if response.status!=200: response.read(65536);raise APIError('llama.cpp hat die Anfrage abgelehnt. Kontext und Anfrageparameter prüfen.',response.status if 400<=response.status<600 else 502,'upstream_error') if not data.get('stream'): raw=response.read(16*1024*1024+1) if len(raw)>16*1024*1024:raise InferenceError('Modellantwort überschreitet 16 MiB.') return self.send(json.loads(raw)) self.sent=True;self.connection.settimeout(30) self.send_response(200);self.compatibility_headers();self.send_header('Content-Type','text/event-stream');self.send_header('Cache-Control','no-store');self.send_header('Connection','close');self.end_headers() deadline=time.monotonic()+600 while time.monotonic()