From 90a97567f37b08d78e1c3d1b86cfca271448aa91 Mon Sep 17 00:00:00 2001 From: Mikei386 <44135113+Mikei386@users.noreply.github.com> Date: Tue, 29 Sep 2026 20:37:29 +0200 Subject: [PATCH] Keep TTS and STT workers resident between requests --- ENDPOINT.md | 4 +-- endpoint.py | 8 ++++- inference.py | 27 +++++++++++----- server.py | 8 +++-- stt-ui.js | 2 +- stt.py | 69 ++++++++++++++++++++++++----------------- test_audio_residency.py | 59 +++++++++++++++++++++++++++++++++++ tts-ui.js | 2 +- tts_test.py | 61 ++++++++++++++++++++++-------------- tts_worker.py | 35 +++++++++++++-------- video.py | 2 +- 11 files changed, 197 insertions(+), 80 deletions(-) create mode 100644 test_audio_residency.py diff --git a/ENDPOINT.md b/ENDPOINT.md index e9c03d1..acdb222 100644 --- a/ENDPOINT.md +++ b/ENDPOINT.md @@ -170,11 +170,11 @@ Profilfreigaben: mehrere Chatprofile, höchstens ein Bildprofil. Die Auswahl ein ### Sprachausgabe -`POST /v1/audio/speech` verwendet ein ausdrücklich freigegebenes TTS-Profil als `model`, `input` (1–1000 Zeichen), optional `voice` (Standard Ryan), `language` (Standard German), `speed` (0,25–4; sonst Profilwert) und `response_format: "wav"`. Antwort: WAV-Datei, kein JSON. MP3 und Streaming sind noch nicht unterstützt. TTS teilt sich die Reservierung mit Chat und Bildern; fremde GPU-Prozesse werden nicht beendet. Das Modell wird nach jeder Ausgabe entladen. +`POST /v1/audio/speech` verwendet ein ausdrücklich freigegebenes TTS-Profil als `model`, `input` (1–1000 Zeichen), optional `voice` (Standard Ryan), `language` (Standard German), `speed` (0,25–4; sonst Profilwert) und `response_format: "wav"`. Antwort: WAV-Datei, kein JSON. MP3 und Streaming sind noch nicht unterstützt. TTS lädt beim ersten Auftrag auf die RTX 3060 und hält den eigenen Worker für weitere Anfragen bereit. Ein geladenes Chatprofil kann parallel bestehen, sofern die RTX 3060 mindestens 7 GiB für TTS frei hat. Ein exklusiver Bild- oder Videomodus, ein inkompatibler GPU-Profilwechsel oder das Stoppen des Endpunkts entlädt TTS. Fremde GPU-Prozesse werden nicht beendet. ### Spracherkennung -`POST /v1/audio/transcriptions` benötigt Bearer-Token und ein explizit freigegebenes STT-Profil. Multipart-Felder: `model` (API-Profilname), `file` (PCM16-WAV, mono, 16 kHz, maximal 120 Sekunden / 8 MiB), optional `language` (`de`, `en`, `auto`, Standard `de`), `response_format` (`json`). Antwort: `{"text":"…"}`. Andere Formate, Chunked-Uploads, Zeitstempel und Streaming werden abgelehnt. Ein STT-Auftrag zur Zeit; CPU-Worker wird danach beendet. Keine Nutzung des alten Routers. +`POST /v1/audio/transcriptions` benötigt Bearer-Token und ein explizit freigegebenes STT-Profil. Multipart-Felder: `model` (API-Profilname), `file` (PCM16-WAV, mono, 16 kHz, maximal 120 Sekunden / 8 MiB), optional `language` (`de`, `en`, `auto`, Standard `de`), `response_format` (`json`). Antwort: `{"text":"…"}`. Andere Formate, Chunked-Uploads, Zeitstempel und Streaming werden abgelehnt. Ein STT-Auftrag zur Zeit; der CPU-Worker bleibt danach für weitere Aufnahmen geladen. Er wird bei Endpunkt-Stopp, einem neuen Build oder Fehler entladen. Keine Nutzung des alten Routers. ### Modelllisten für Clients diff --git a/endpoint.py b/endpoint.py index 9edf681..dd2c385 100644 --- a/endpoint.py +++ b/endpoint.py @@ -42,7 +42,7 @@ class Endpoint: 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)) + 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().get('loaded')) if kind=='audio' else bool(self.stt and self.stt.status().get('loaded')) 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): @@ -70,6 +70,12 @@ class Endpoint: enabled.add(data['id']) else:enabled.discard(data['id']) self.config['enabled_profiles']=sorted(enabled);self.persist() + if not data['enabled'] and row['kind']=='audio' and self.tts: + job=self.tts.status().get('job') or {} + if job.get('profile_id')==row['id'] and job.get('state')!='running':self.tts.stop() + if not data['enabled'] and row['kind']=='stt' and self.stt: + job=self.stt.status().get('job') or {} + if job.get('profile_id')==row['id'] and job.get('state')!='running':self.stt.stop() return self.status() def start(self): with self.lock: diff --git a/inference.py b/inference.py index 3e82cc6..107ef43 100644 --- a/inference.py +++ b/inference.py @@ -19,9 +19,9 @@ class InferenceError(ValueError): class Scheduler: """FIFO admission; same-profile requests share configured slots, switches drain.""" def __init__(self, worker): - self.worker=worker;self.cv=threading.Condition();self.queue=[];self.key=None;self.active=0;self.transition=False;self.gpu_mode="llm" + self.worker=worker;self.evict_tts=lambda:None;self.cv=threading.Condition();self.queue=[];self.key=None;self.active=0;self.tts_active=0;self.transition=False;self.gpu_mode="llm" def status(self): - with self.cv:return dict(active_requests=self.active,waiting_requests=len(self.queue),switching=self.transition) + with self.cv:return dict(active_requests=self.active+self.tts_active,waiting_requests=len(self.queue),switching=self.transition) @contextlib.contextmanager def lease(self, key, slots=1, prepare=None, timeout=600, allowed=lambda:True): ticket=object();deadline=time.monotonic()+timeout;claimed=False @@ -32,7 +32,7 @@ class Scheduler: while True: if self.gpu_mode!="llm":raise InferenceError("Video-Modus aktiv oder Moduswechsel läuft; GPU-Aufträge sind gesperrt.") if not allowed():raise InferenceError('Endpunkt wird gestoppt oder Profil ist nicht mehr aktiviert.') - if self.queue[0] is ticket and not self.transition and (not self.active or (self.key==key and self.active=deadline:raise InferenceError('GPU-Auftrag läuft; TTS später erneut versuchen.') + self.cv.wait(min(1,max(.01,deadline-time.monotonic()))) + def release(): + with self.cv:self.tts_active-=1;self.cv.notify_all() + return release def unload_idle(self): with self.cv: - if self.active:return False + if self.active or self.tts_active:return False self.transition=True - try:self.worker.stop();self.key=None + try:self.worker.stop();self.evict_tts();self.key=None finally:self.transition=False;self.cv.notify_all() return True diff --git a/server.py b/server.py index 1c17da4..405b9a5 100644 --- a/server.py +++ b/server.py @@ -94,7 +94,8 @@ class Server(ThreadingHTTPServer): self.scheduler=Scheduler(self.worker) self.profiles.chat_blockers=self.worker.blockers self.image_tests.acquire=self.scheduler.image_reservation - self.tts_tests.acquire=self.scheduler.image_reservation + self.tts_tests.acquire=self.scheduler.tts_reservation + self.scheduler.evict_tts=self.tts_tests.unload_idle self.stt=STT(self.profiles,self.catalog,self.runtime) self.profiles.stt_blockers=self.stt.blockers self.endpoint=Endpoint(self.catalog.root.parent,self.profiles,self.worker,self.scheduler,self.image_tests,self.credentials,self.server_port) @@ -471,7 +472,10 @@ class Handler(BaseHTTPRequestHandler): if tts_job and tts_job.get("state")=="running" and tts_job.get("profile_id")==data.get("id"):raise ValueError("Dieses TTS-Profil wird gerade ausgeführt. Zuerst den Auftrag beenden.") 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)) + result=self.server.profiles.delete(data) + if stt_job and stt_job.get('profile_id')==data.get('id'):self.server.stt.stop() + if tts_job and tts_job.get('profile_id')==data.get('id'):self.server.tts_tests.stop() + return self.respond(result) 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': diff --git a/stt-ui.js b/stt-ui.js index aac22f1..9570cd8 100644 --- a/stt-ui.js +++ b/stt-ui.js @@ -4,7 +4,7 @@ window.STTUI=(()=>{ function html(setup=false){return `

${setup?'Spracherkennung einrichten':'Audio in Text umwandeln'}

${setup?'

Qwen3-ASR 0.6B Q8_0 nutzt die vorhandene llama.cpp-Laufzeit auf der CPU. Installiert werden die geprüfte Modellvariante von ggml-org (rund 805 MB) und der Audio-Projektor (214 MB). Andere GGUF-Konvertierungen können inkompatibel sein. Bestehende Dienste werden nicht verwendet oder verändert.

llama.cpp installieren / Build auswählen →

Der Download erscheint unter STT → Downloads. Nach Abschluss ein neues Profil aus der Bibliothek anlegen oder ein bestehendes Profil unten ausdrücklich umstellen. Alte Dateien bleiben erhalten.

Vorhandenes Profil umstellen

Ersetzt die Modellzuordnung durch die geprüfte ggml-org-Variante. Profilname bleibt erhalten.

':'

Transkript

Die Audiodatei wird im Browser in Mono-WAV mit 16 kHz umgewandelt. Unterstützte Eingabeformate hängen vom Browser ab. Audio und Transkript werden nicht dauerhaft gespeichert; das letzte Ergebnis bleibt bis zum nächsten Auftrag im Arbeitsspeicher.

'}
`;} async function wav(file){if(file.size>20*1024**2)throw Error('Datei ist größer als 20 MiB.');const context=new AudioContext();let decoded;try{decoded=await context.decodeAudioData(await file.arrayBuffer());}finally{await context.close();}if(!decoded.duration||decoded.duration>120)throw Error('Maximal 120 Sekunden Audio erlaubt.');const offline=new OfflineAudioContext(1,Math.ceil(decoded.duration*16000),16000),source=offline.createBufferSource();source.buffer=decoded;source.connect(offline.destination);source.start();const samples=(await offline.startRendering()).getChannelData(0),buffer=new ArrayBuffer(44+samples.length*2),v=new DataView(buffer);const str=(at,s)=>{for(let i=0;iv.setInt16(44+i*2,Math.max(-1,Math.min(1,s))*(s<0?32768:32767),true));return new Blob([buffer],{type:'audio/wav'});} function bind(setup=false){const root=document.querySelector('#stt-panel');if(!root)return;const el=id=>root.querySelector('#'+id);let busy=false,profiles=[]; - async function refresh(){try{const s=await api('stt');if(!root.isConnected)return;el('stt-status').textContent=setup?`llama.cpp: ${s.installed?'vorhanden':'fehlt'} · Modell: ${s.model?'vorhanden':'fehlt'} · Audio-Projektor: ${s.projector?'vorhanden':'fehlt'}`:s.job?.phase||'Profil und Audiodatei auswählen.';if(setup){el('stt-setup').disabled=busy||(s.projector&&s.model);el('stt-assign').disabled=busy||!s.model||!profiles.length;}else{el('stt-start').disabled=busy||s.job?.state==='running'||!profiles.length;el('stt-cancel').disabled=s.job?.state!=='running';el('stt-result').textContent=s.job?.text||'';}}catch(err){if(root.isConnected)el('stt-error').textContent=err.message;}finally{if(root.isConnected)setTimeout(refresh,2000);}} + async function refresh(){try{const s=await api('stt');if(!root.isConnected)return;el('stt-status').textContent=setup?`llama.cpp: ${s.installed?'vorhanden':'fehlt'} · Modell: ${s.model?'vorhanden':'fehlt'} · Audio-Projektor: ${s.projector?'vorhanden':'fehlt'}`:(s.job?.phase||'Profil und Audiodatei auswählen.')+` · CPU-Modell ${s.loaded?'geladen':'nicht geladen'}`;if(setup){el('stt-setup').disabled=busy||(s.projector&&s.model);el('stt-assign').disabled=busy||!s.model||!profiles.length;}else{el('stt-start').disabled=busy||s.job?.state==='running'||!profiles.length;el('stt-cancel').disabled=s.job?.state!=='running';el('stt-result').textContent=s.job?.text||'';}}catch(err){if(root.isConnected)el('stt-error').textContent=err.message;}finally{if(root.isConnected)setTimeout(refresh,2000);}} if(setup){api('profiles').then(p=>{if(!root.isConnected)return;profiles=p.profiles.filter(p=>p.kind==='stt');el('stt-profile').innerHTML=profiles.map(p=>``).join('');}).catch(err=>{if(root.isConnected)el('stt-error').textContent=err.message;});el('stt-assign').onclick=async()=>{const p=profiles.find(p=>p.id===el('stt-profile').value);if(!p)return;busy=true;try{await api('stt/assign',{id:p.id,revision:p.revision});el('stt-error').textContent='Profil umgestellt. Jetzt STT → Testen öffnen.';p.revision++;}catch(err){el('stt-error').textContent=err.message;}finally{busy=false;}};el('stt-setup').onclick=async()=>{busy=true;try{await api('stt/setup',{});el('stt-error').textContent='Download ist eingeplant. Fortschritt unter Downloads.';}catch(err){el('stt-error').textContent=err.message;}finally{busy=false;}};} else{api('profiles').then(p=>{if(!root.isConnected)return;profiles=p.profiles.filter(p=>p.kind==='stt'&&p.runnable);el('stt-profile').innerHTML=profiles.map(p=>``).join('')||'';}).catch(err=>{if(root.isConnected)el('stt-error').textContent=err.message;});el('stt-form').onsubmit=async ev=>{ev.preventDefault();busy=true;el('stt-start').disabled=true;el('stt-error').textContent='';try{const body=new FormData();body.set('file',await wav(el('stt-file').files[0]),'audio.wav');body.set('profile_id',el('stt-profile').value);body.set('language',el('stt-language').value);if(root.isConnected)await api('stt/start',body);}catch(err){if(root.isConnected)el('stt-error').textContent=err.message;}finally{busy=false;}};el('stt-cancel').onclick=()=>api('stt/cancel',{}).catch(err=>{if(root.isConnected)el('stt-error').textContent=err.message;});}refresh(); } diff --git a/stt.py b/stt.py index da80fc8..051032f 100644 --- a/stt.py +++ b/stt.py @@ -50,7 +50,7 @@ def read_upload(handler): class STT: def __init__(self,profiles,catalog,runtime): - self.profiles=profiles;self.catalog=catalog;self.runtime=runtime;self.lock=threading.RLock();self.job=None;self.process=None;self.cancel=threading.Event() + self.profiles=profiles;self.catalog=catalog;self.runtime=runtime;self.lock=threading.RLock();self.job=None;self.process=None;self.port=None;self.key=None;self.binary_path=None;self.cancel=threading.Event() def build(self): state=self.runtime.status();ident=state.get('active') if not isinstance(ident,str) or len(ident)!=32 or any(c not in '0123456789abcdef' for c in ident):raise ValueError('Zuerst unter Laufzeiten → llama.cpp einen Build erstellen und aktivieren.') @@ -75,7 +75,7 @@ class STT: except ValueError:projector=False try:self.model_entry();model=True except ValueError:model=False - with self.lock:return dict(model=model,job=dict(self.job) if self.job else None,installed=installed,projector=projector,repo=REPO,revision=REVISION) + with self.lock:return dict(model=model,job=dict(self.job) if self.job else None,installed=installed,projector=projector,repo=REPO,revision=REVISION,loaded=bool(self.process and self.process.poll() is None)) def model_entry(self): for x in self.catalog.status()['entries']: if supported(x):return self.catalog.entry(x['id']) @@ -102,30 +102,39 @@ class STT: if self.job and self.job['state']=='running':raise ValueError('Ein STT-Auftrag läuft bereits.') p=next((p for p in self.profiles.status()['profiles'] if p['id']==profile_id and p['kind']=='stt'),None) if not p or not p['runnable']:raise ValueError('STT-Profil nicht ausführbar. Zuerst Einrichten öffnen.') - headroom=cgroup_headroom() - if headroom is not None and headroom<4*1024**3:raise ValueError('Mindestens 4 GiB freier Deck-RAM werden benötigt.') + if not self.process or self.process.poll() is not None: + headroom=cgroup_headroom() + if headroom is not None and headroom<4*1024**3:raise ValueError('Mindestens 4 GiB freier Deck-RAM werden benötigt.') binary=self.build();model=self.catalog.root/p['model_id']/'model.gguf';projector=self.catalog.root/self.projector()['id']/'model.gguf' - self.cancel.clear();ident=uuid.uuid4().hex;self.job=dict(id=ident,state='running',phase='Spracherkennung lädt auf der CPU',profile_id=profile_id) + self.cancel.clear();ident=uuid.uuid4().hex;self.job=dict(id=ident,state='running',phase='Vorhandene Spracherkennung wird verwendet' if self.process and self.process.poll() is None else 'Spracherkennung lädt auf der CPU',profile_id=profile_id) threading.Thread(target=self._run,args=(binary,model,projector,audio,language),daemon=True).start();return dict(self.job) def _run(self,binary,model,projector,audio,language): process=None;result=dict(state='failed',phase='Spracherkennung fehlgeschlagen. Laufzeit und Speicher prüfen.') try: - key=secrets.token_hex(24) - with socket.socket() as s:s.bind(('127.0.0.1',0));port=s.getsockname()[1] - env=dict(os.environ,CUDA_VISIBLE_DEVICES='',OMP_NUM_THREADS='2') - args=[str(binary),'--model',str(model),'--mmproj',str(projector),'--no-mmproj-offload','--n-gpu-layers','0','--ctx-size','4096','--threads','2','--parallel','1','--host','127.0.0.1','--port',str(port),'--no-ui','--fit','off','--alias','deck-stt','--api-key',key] with self.lock: if self.cancel.is_set():raise InterruptedError() - process=subprocess.Popen(args,env=env,stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL,start_new_session=True);self.process=process + process=self.process if self.process and self.process.poll() is None and self.binary_path==str(binary) else None + if process is None: + self.unload_idle() + key=secrets.token_hex(24) + with socket.socket() as s:s.bind(('127.0.0.1',0));port=s.getsockname()[1] + env=dict(os.environ,CUDA_VISIBLE_DEVICES='',OMP_NUM_THREADS='2') + args=[str(binary),'--model',str(model),'--mmproj',str(projector),'--no-mmproj-offload','--n-gpu-layers','0','--ctx-size','4096','--threads','2','--parallel','1','--host','127.0.0.1','--port',str(port),'--no-ui','--fit','off','--alias','deck-stt','--api-key',key] + with self.lock: + if self.cancel.is_set():raise InterruptedError() + process=subprocess.Popen(args,env=env,stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL,start_new_session=True);self.process=process;self.port=port;self.key=key;self.binary_path=str(binary) + deadline=time.monotonic()+90 + while True: + if self.cancel.wait(.3):raise InterruptedError() + if process.poll() is not None:raise ValueError('llama.cpp konnte das ASR-Modell nicht laden. Build-Unterstützung und freien RAM prüfen.') + if time.monotonic()>deadline:raise ValueError('Zeitlimit beim Laden des ASR-Modells.') + try: + with urllib.request.urlopen(urllib.request.Request(f'http://127.0.0.1:{port}/health',headers={'Authorization':'Bearer '+key}),timeout=1) as response: + if response.status==200:break + except OSError:continue + else: + with self.lock:port=self.port;key=self.key base=f'http://127.0.0.1:{port}';headers={'Authorization':'Bearer '+key};deadline=time.monotonic()+90 - while True: - if self.cancel.wait(.3):raise InterruptedError() - if process.poll() is not None:raise ValueError('llama.cpp konnte das ASR-Modell nicht laden. Build-Unterstützung und freien RAM prüfen.') - if time.monotonic()>deadline:raise ValueError('Zeitlimit beim Laden des ASR-Modells.') - try: - with urllib.request.urlopen(urllib.request.Request(base+'/health',headers=headers),timeout=1) as response: - if response.status==200:break - except OSError:continue with self.lock:self.job['phase']='Audio wird auf der CPU transkribiert' boundary='deck-'+uuid.uuid4().hex body=f'--{boundary}\r\nContent-Disposition: form-data; name="file"; filename="audio.wav"\r\nContent-Type: audio/wav\r\n\r\n'.encode()+audio+b'\r\n' @@ -135,24 +144,26 @@ class STT: payload=json.loads(response.read(1024*1024)) if not isinstance(payload.get('text'),str):raise ValueError('Die Laufzeit hat kein Transkript geliefert.') text=payload['text'].split('')[-1].replace('<|endoftext|>','').strip() - result=dict(state='complete',phase='Transkription fertig · Modell entladen',text=text) + result=dict(state='complete',phase='Transkription fertig · Modell bleibt geladen',text=text) except InterruptedError:result=dict(state='cancelled',phase='Transkription abgebrochen') except ValueError as exc:result=dict(state='failed',phase=str(exc)) except Exception:pass finally: - if process and process.poll() is None: - try: - os.killpg(process.pid,signal.SIGTERM) - try:process.wait(timeout=5) - except subprocess.TimeoutExpired:os.killpg(process.pid,signal.SIGKILL);process.wait() - except ProcessLookupError:pass + if result['state']!='complete':self.unload_idle() with self.lock: if self.cancel.is_set():result=dict(state='cancelled',phase='Transkription abgebrochen') - self.process=None;self.job.update(result) + self.job.update(result) + def unload_idle(self): + with self.lock: + process=self.process;self.process=None;self.port=None;self.key=None;self.binary_path=None + if process and process.poll() is None: + try: + os.killpg(process.pid,signal.SIGTERM) + try:process.wait(timeout=5) + except subprocess.TimeoutExpired:os.killpg(process.pid,signal.SIGKILL);process.wait() + except ProcessLookupError:pass def stop(self): with self.lock: self.cancel.set() - if self.process and self.process.poll() is None: - try:os.killpg(self.process.pid,signal.SIGTERM) - except ProcessLookupError:pass + self.unload_idle() return dict(cancellation_requested=True) diff --git a/test_audio_residency.py b/test_audio_residency.py new file mode 100644 index 0000000..e866314 --- /dev/null +++ b/test_audio_residency.py @@ -0,0 +1,59 @@ +"""Synthetic process tests: audio workers survive a second request and stop cleanly.""" +import json +import sys +import tempfile +import time +import unittest +from pathlib import Path +from types import SimpleNamespace +from unittest.mock import patch +from inference import Scheduler +from stt import STT +from test_stt import audio +from tts_test import TTSTests + +def done(manager): + for _ in range(200): + job=manager.status()['job'] + if job and job['state']!='running':return job + time.sleep(.05) + raise AssertionError('synthetic worker timed out') + +class ResidencyTests(unittest.TestCase): + def test_tts_reuses_process_and_gpu_switch_eviction(self): + with tempfile.TemporaryDirectory() as d: + script=Path(d)/'worker.py' + script.write_text('import json,sys,pathlib\nprint(json.dumps({"ready":True}),flush=True)\nfor line in sys.stdin:\n q=json.loads(line);p=pathlib.Path(q["output"]);(p/"result.wav").write_bytes(b"synthetic wav");(p/"result.json").write_text(json.dumps({"ok":True}))\n') + runtime=SimpleNamespace(paths=lambda:(Path(sys.executable),Path(d)/'model'),status=lambda:{'installed':True}) + profile=dict(id='p',kind='audio',runnable=True,parameters={'speed':1}) + manager=TTSTests(Path(d)/'jobs',SimpleNamespace(status=lambda:{'profiles':[profile]}),runtime) + worker=SimpleNamespace(stop=lambda:None) + scheduler=Scheduler(worker);scheduler.evict_tts=manager.unload_idle + manager.acquire=scheduler.tts_reservation + def gpu():return [dict(uuid='synthetic-3060',name='RTX 3060',processes=int(bool(manager.process and manager.process.poll() is None)),free_mib=11000)] + try: + with patch('tts_test.WORKER',script),patch('tts_test.probe',side_effect=gpu),patch('tts_test.cgroup_headroom',return_value=None): + with scheduler.lease(('chat','medium')):pass + manager.start('p','first');self.assertEqual(done(manager)['state'],'complete');pid=manager.process.pid + manager.start('p','second');self.assertEqual(done(manager)['state'],'complete');self.assertEqual(manager.process.pid,pid) + with scheduler.lease(('chat','medium')):self.assertTrue(manager.status()['loaded']) + with scheduler.lease(('image',)):self.assertFalse(manager.status()['loaded']) + finally:manager.stop() + + def test_stt_reuses_cpu_server_until_stop(self): + with tempfile.TemporaryDirectory() as d: + script=Path(d)/'server.py' + script.write_text('#!'+sys.executable+'\nimport argparse,json\nfrom http.server import BaseHTTPRequestHandler,HTTPServer\np=argparse.ArgumentParser();p.add_argument("--port",type=int);a,_=p.parse_known_args()\nclass H(BaseHTTPRequestHandler):\n def log_message(self,*x):pass\n def do_GET(self):\n self.send_response(200);self.end_headers();self.wfile.write(b"ok")\n def do_POST(self):\n self.rfile.read(int(self.headers["Content-Length"]));b=json.dumps({"text":"synthetic transcript"}).encode();self.send_response(200);self.send_header("Content-Length",str(len(b)));self.end_headers();self.wfile.write(b)\nHTTPServer(("127.0.0.1",a.port),H).serve_forever()\n') + script.chmod(0o700) + profile=dict(id='p',kind='stt',runnable=True,model_id='m') + catalog=SimpleNamespace(root=Path(d),status=lambda:{'entries':[]}) + manager=STT(SimpleNamespace(status=lambda:{'profiles':[profile]}),catalog,SimpleNamespace(status=lambda:{'active':None})) + try: + with patch.object(manager,'build',return_value=script),patch.object(manager,'projector',return_value={'id':'projector'}),patch('stt.cgroup_headroom',return_value=None): + manager.start('p',audio());self.assertEqual(done(manager)['state'],'complete');pid=manager.process.pid + manager.start('p',audio());self.assertEqual(done(manager)['state'],'complete');self.assertEqual(manager.process.pid,pid) + self.assertTrue(manager.status()['loaded']) + finally:manager.stop() + self.assertFalse(manager.status()['loaded']) + +if __name__=='__main__':unittest.main() diff --git a/tts-ui.js b/tts-ui.js index af6836a..804fa50 100644 --- a/tts-ui.js +++ b/tts-ui.js @@ -3,7 +3,7 @@ window.TTSUI=(()=>{ async function api(path='',data){const r=await fetch('/api/v1/'+path,(data===undefined)?{}:{method:'POST',headers:{'Content-Type':'application/json','X-Athena-Deck':'1'},body:JSON.stringify(data)});const x=await r.json();if(!r.ok)throw Error(x.error||'TTS-Anfrage fehlgeschlagen');return x;} function html(setup=false){return `

${setup?'Qwen3-TTS einrichten':'Sprache ausprobieren'}

${setup?'

Die Einrichtung installiert eine eigene CUDA-Laufzeit und das offizielle Qwen3-TTS 1.7B CustomVoice mit eingebauten Stimmen. Dafür werden mehrere GB heruntergeladen. MLX-Dateien sind für Apple-Geräte; vorhandene Dateien werden nicht gelöscht.

Vorhandenes TTS-Profil verbinden

Ersetzt im ausgewählten Profil die Modellzuordnung durch die eingerichtete CUDA-Variante. Der Profilname bleibt erhalten.

':'

Ausführung auf der RTX 3060. Andere Deck-Modellaufträge werden koordiniert; fremde Dienste werden nicht gestoppt. Text wird nicht dauerhaft gespeichert. Ergebnisdateien bleiben auf Athena.

'}
`;} function bind(setup=false){const root=document.querySelector('#tts-panel');if(!root)return;const el=id=>root.querySelector('#'+id);let state,profiles=[],busy=false,lastResult='';const error=x=>{if(root.isConnected)el('tts-error').textContent=x;}; - async function refresh(){try{state=await api('tts');if(!root.isConnected)return;const rt=state.runtime,j=state.job;el('tts-status').textContent=setup?(rt.job?.phase||(rt.installed?'TTS ist eingerichtet.':'Noch nicht eingerichtet.')):(j?.phase||'Profil wählen und Text eingeben.'); + async function refresh(){try{state=await api('tts');if(!root.isConnected)return;const rt=state.runtime,j=state.job;el('tts-status').textContent=setup?(rt.job?.phase||(rt.installed?'TTS ist eingerichtet.':'Noch nicht eingerichtet.')):(j?.phase||'Profil wählen und Text eingeben.')+` · RTX 3060 ${state.loaded?'geladen':'nicht geladen'}`; if(setup){el('tts-install').disabled=busy||rt.installed||rt.job?.state==='running';el('tts-install-cancel').disabled=rt.job?.state!=='running';el('tts-assign').disabled=busy||!rt.installed||!profiles.length;} else{el('tts-generate').disabled=busy||j?.state==='running'||!profiles.length||!rt.installed;el('tts-cancel').disabled=j?.state!=='running';if(j?.state==='complete'&&lastResult!==j.id){lastResult=j.id;el('tts-result').innerHTML=`

Ergebnis

WAV herunterladen

`;}} }catch(err){error(err.message);}finally{if(root.isConnected)setTimeout(refresh,2000);}} diff --git a/tts_test.py b/tts_test.py index 1ecada4..3e9de5e 100644 --- a/tts_test.py +++ b/tts_test.py @@ -1,13 +1,14 @@ """Serial own TTS workers with shared model reservation, WAV results and cancellation.""" -import json,os,signal,subprocess,threading,time,uuid,hashlib +import json,os,signal,subprocess,threading,time,uuid,hashlib,select from pathlib import Path from tts_runtime import REPO,REVISION,SPEAKERS,LANGUAGES from image_test import probe,cgroup_headroom +WORKER=Path(__file__).with_name('tts_worker.py') class TTSTests: def __init__(self,root,profiles,runtime): - self.root=Path(root);self.profiles=profiles;self.runtime=runtime;self.lock=threading.RLock();self.job=None;self.process=None;self.cancel=threading.Event();self.acquire=lambda wait=False:lambda:None + self.root=Path(root);self.profiles=profiles;self.runtime=runtime;self.lock=threading.RLock();self.job=None;self.process=None;self.gpu_uuid=None;self.cancel=threading.Event();self.acquire=lambda wait=False:lambda:None def status(self): - with self.lock:return dict(job=dict(self.job) if self.job else None,runtime=self.runtime.status()) + with self.lock:return dict(job=dict(self.job) if self.job else None,runtime=self.runtime.status(),loaded=bool(self.process and self.process.poll() is None),gpu_uuid=self.gpu_uuid) def blockers(self,p): if p['model']['repo']!=REPO or p['model']['file']!='model.safetensors' or p['model'].get('revision')!=REVISION:return ['Für diese TTS-Variante fehlt die Anbindung in Deck. Unter TTS → Einrichten kannst du ausdrücklich auf die unterstützte Variante mit eingebauten Stimmen wechseln. Die ursprüngliche Datei bleibt erhalten.'] return [] if self.runtime.status()['installed'] else ['TTS-Laufzeit unter TTS → Einrichten installieren.'] @@ -27,53 +28,67 @@ class TTSTests: if self.job and self.job['state']=='running':raise ValueError('TTS-Auftrag läuft bereits.') p=next((p for p in self.profiles.status()['profiles'] if p['id']==profile_id and p['kind']=='audio'),None) if not p or not p['runnable']:raise ValueError('TTS-Profil noch nicht eingerichtet.') - gpu=next((g for g in probe() if 'RTX 3060' in g['name'] and not g['processes'] and g['free_mib']>7000),None) - if not gpu:raise ValueError('Die RTX 3060 muss frei sein und mindestens 7 GiB freien Speicher haben.') + resident=self.process and self.process.poll() is None + gpu=next((g for g in probe() if g['uuid']==self.gpu_uuid),None) if resident else next((g for g in probe() if 'RTX 3060' in g['name'] and not g['processes'] and g['free_mib']>7000),None) + if not gpu or resident and gpu['processes']>1:raise ValueError('Die RTX 3060 ist nicht verfügbar oder wurde anderweitig belegt.') headroom=cgroup_headroom() if headroom is not None and headroom<8*1024**3:raise ValueError('Zu wenig freier Deck-RAM für TTS.') self.cancel.clear();ident=uuid.uuid4().hex;self.root.mkdir(parents=True,exist_ok=True,mode=0o700) - self.job=dict(id=ident,state='running',phase='Sprachmodell wird geladen',profile_id=profile_id,started_at=time.time()) + self.job=dict(id=ident,state='running',phase='Vorhandenes Sprachmodell wird verwendet' if resident else 'Sprachmodell wird geladen',profile_id=profile_id,started_at=time.time()) threading.Thread(target=self._run,args=(ident,text,speaker,language,speed if speed is not None else p['parameters']['speed'],gpu,release),daemon=True).start();return dict(self.job) except Exception:release();raise def _run(self,ident,text,speaker,language,speed,gpu,release): directory=self.root/ident;process=None;state='failed';phase='TTS-Ausführung fehlgeschlagen' try: directory.mkdir(mode=0o700);python,model=self.runtime.paths() - env={k:v for k,v in os.environ.items() if not k.startswith(('HF_','LLAMA_'))};env.update(CUDA_VISIBLE_DEVICES=gpu['uuid'],HF_HUB_OFFLINE='1',TRANSFORMERS_OFFLINE='1',OMP_NUM_THREADS='2',HOME=str(directory)) with self.lock: if self.cancel.is_set():raise InterruptedError() - process=subprocess.Popen([str(python),str(Path(__file__).with_name('tts_worker.py'))],stdin=subprocess.PIPE,stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL,start_new_session=True,env=env);self.process=process;self.job['phase']='Sprache wird auf der RTX 3060 erzeugt' - process.stdin.write(json.dumps(dict(model=str(model),output=str(directory),text=text,speaker=speaker,language=language,speed=speed)).encode());process.stdin.close() + process=self.process if self.process and self.process.poll() is None else None + new_process=process is None + if process is None: + env={k:v for k,v in os.environ.items() if not k.startswith(('HF_','LLAMA_'))};env.update(CUDA_VISIBLE_DEVICES=gpu['uuid'],HF_HUB_OFFLINE='1',TRANSFORMERS_OFFLINE='1',OMP_NUM_THREADS='2',HOME=str(self.root)) + process=subprocess.Popen([str(python),str(WORKER),str(model)],stdin=subprocess.PIPE,stdout=subprocess.PIPE,stderr=subprocess.DEVNULL,start_new_session=True,env=env,text=True,bufsize=1);self.process=process;self.gpu_uuid=gpu['uuid'] + if new_process: + if not select.select([process.stdout],[],[],180)[0]:raise ValueError('Zeitlimit beim Laden des TTS-Modells.') + ready=json.loads(process.stdout.readline()) + if not ready.get('ready'):raise ValueError(ready.get('error','TTS-Modell konnte nicht geladen werden.')) + with self.lock:self.job['phase']='Sprache wird auf der RTX 3060 erzeugt' + process.stdin.write(json.dumps(dict(output=str(directory),text=text,speaker=speaker,language=language,speed=speed))+'\n');process.stdin.flush() end=time.monotonic()+600 - while process.poll() is None: + while not (directory/'result.json').is_file(): if self.cancel.wait(1):raise InterruptedError() if time.monotonic()>end:raise ValueError('Zeitlimit der Sprachgenerierung erreicht.') + if process.poll() is not None:raise ValueError('TTS-Worker wurde unerwartet beendet.') current=next((g for g in probe() if g['uuid']==gpu['uuid']),None) if not current or current['processes']>1:raise ValueError('Die reservierte GPU wurde anderweitig belegt.') - result=json.loads((directory/'result.json').read_text()) if (directory/'result.json').exists() else dict(ok=False,error='Worker beendet: möglicher Speichermangel oder Laufzeitfehler.') + result=json.loads((directory/'result.json').read_text()) if self.cancel.is_set():raise InterruptedError() if not result['ok']:raise ValueError(result['error']) wav=directory/'result.wav' if not wav.exists() or wav.stat().st_size>32*1024**2:raise ValueError('Ungültige Audioausgabe.') - state='complete';phase='Sprache fertig · Modell entladen' + state='complete';phase='Sprache fertig · Modell bleibt geladen' except InterruptedError:state='cancelled';phase='Sprachgenerierung abgebrochen' except Exception as exc:phase=str(exc) if isinstance(exc,ValueError) else 'Sprachgenerierung fehlgeschlagen. Laufzeit und Ressourcen prüfen.' finally: - if process and process.poll() is None: - try: - os.killpg(process.pid,signal.SIGTERM) - try:process.wait(timeout=5) - except subprocess.TimeoutExpired:os.killpg(process.pid,signal.SIGKILL);process.wait() - except ProcessLookupError:pass + if state!='complete':self.unload_idle() if state=='complete':(directory/'complete').touch() - with self.lock:self.process=None;self.job.update(state=state,phase=phase,finished_at=time.time()) + with self.lock:self.job.update(state=state,phase=phase,finished_at=time.time()) release() + def unload_idle(self): + with self.lock: + process=self.process;self.process=None;self.gpu_uuid=None + if process and process.poll() is None: + try: + os.killpg(process.pid,signal.SIGTERM) + try:process.wait(timeout=5) + except subprocess.TimeoutExpired:os.killpg(process.pid,signal.SIGKILL);process.wait() + except ProcessLookupError:pass + if process: + if process.stdin:process.stdin.close() + if process.stdout:process.stdout.close() def stop(self): self.cancel.set() - with self.lock:process=self.process - if process and process.poll() is None: - try:os.killpg(process.pid,signal.SIGTERM) - except ProcessLookupError:pass + self.unload_idle() return dict(cancellation_requested=True) def audio(self,ident): if not isinstance(ident,str) or len(ident)!=32 or any(c not in '0123456789abcdef' for c in ident) or not (self.root/ident/'complete').exists():raise ValueError('Audio noch nicht verfügbar.') diff --git a/tts_worker.py b/tts_worker.py index 088f223..ee1758f 100644 --- a/tts_worker.py +++ b/tts_worker.py @@ -1,21 +1,30 @@ -"""One isolated offline Qwen-TTS request; stdin is never persisted.""" +"""Persistent offline Qwen-TTS worker; line-delimited stdin, no prompt logging.""" import json,sys from pathlib import Path def main(): - request=json.load(sys.stdin);root=Path(request['output']) + model_path=sys.argv[1] try: import torch,soundfile as sf from qwen_tts import Qwen3TTSModel - model=Qwen3TTSModel.from_pretrained(request['model'],device_map='cuda:0',dtype=torch.bfloat16,attn_implementation='sdpa',local_files_only=True) - waves,rate=model.generate_custom_voice(text=request['text'],language=request['language'],speaker=request['speaker'],max_new_tokens=1024) - wave=waves[0] - if request['speed']!=1: - import librosa - wave=librosa.effects.time_stretch(wave,rate=request['speed']) - sf.write(root/'result.wav',wave,rate,subtype='PCM_16') - (root/'result.json').write_text(json.dumps(dict(ok=True,sample_rate=rate,duration=len(wave)/rate))) + model=Qwen3TTSModel.from_pretrained(model_path,device_map='cuda:0',dtype=torch.bfloat16,attn_implementation='sdpa',local_files_only=True) except Exception as exc: - (root/'result.json').write_text(json.dumps(dict(ok=False,error='CUDA-Speicher reicht nicht aus.' if 'outofmemory' in type(exc).__name__.lower() or 'out of memory' in str(exc).lower() else 'TTS-Ausführung fehlgeschlagen: '+type(exc).__name__))) - sys.exit(1) -if __name__=='__main__':main() + print(json.dumps(dict(ready=False,error='CUDA-Speicher reicht nicht aus.' if 'outofmemory' in type(exc).__name__.lower() or 'out of memory' in str(exc).lower() else 'TTS-Modell konnte nicht geladen werden.')),flush=True) + return 1 + print(json.dumps(dict(ready=True)),flush=True) + for line in sys.stdin: + root=None + try: + request=json.loads(line);root=Path(request['output']) + waves,rate=model.generate_custom_voice(text=request['text'],language=request['language'],speaker=request['speaker'],max_new_tokens=1024) + wave=waves[0] + if request['speed']!=1: + import librosa + wave=librosa.effects.time_stretch(wave,rate=request['speed']) + sf.write(root/'result.wav',wave,rate,subtype='PCM_16') + (root/'result.json').write_text(json.dumps(dict(ok=True,sample_rate=rate,duration=len(wave)/rate))) + except Exception as exc: + if root is not None: + (root/'result.json').write_text(json.dumps(dict(ok=False,error='CUDA-Speicher reicht nicht aus.' if 'outofmemory' in type(exc).__name__.lower() or 'out of memory' in str(exc).lower() else 'TTS-Ausführung fehlgeschlagen: '+type(exc).__name__))) + return 0 +if __name__=='__main__':sys.exit(main()) diff --git a/video.py b/video.py index 8b31301..3177887 100644 --- a/video.py +++ b/video.py @@ -47,7 +47,7 @@ class Video: self.stop_owned() deadline=time.monotonic()+60 with self.scheduler.cv: - while self.scheduler.active or self.scheduler.transition: + while self.scheduler.active or self.scheduler.tts_active or self.scheduler.transition: if time.monotonic()>deadline:raise ValueError('Deck-Aufträge konnten noch nicht beendet werden.') self.scheduler.cv.wait(.2) self.scheduler.key=None