"""Owned CPU-only Qwen3-ASR worker. Audio and transcript live only in memory.""" import io,json,os,secrets,signal,socket,subprocess,threading,time,uuid,wave,urllib.request from email import policy from email.parser import BytesParser from pathlib import Path from image_test import cgroup_headroom REPO='ggml-org/Qwen3-ASR-0.6B-GGUF' REVISION='928ab958557df9aa2ef1c93e0e83c7ad0933fae2' MODEL='Qwen3-ASR-0.6B-Q8_0.gguf' PROJECTOR_REPO='ggml-org/Qwen3-ASR-0.6B-GGUF' PROJECTOR_REVISION='928ab958557df9aa2ef1c93e0e83c7ad0933fae2' PROJECTOR='mmproj-Qwen3-ASR-0.6B-Q8_0.gguf' MAX_AUDIO=8*1024*1024 def supported(m):return m.get('repo')==REPO and m.get('revision')==REVISION and m.get('file')==MODEL def validate_wav(audio): if not isinstance(audio,bytes) or not 44<=len(audio)<=MAX_AUDIO:raise ValueError('WAV-Datei bis 8 MiB erforderlich.') try: with wave.open(io.BytesIO(audio)) as w: if w.getnchannels()!=1 or w.getsampwidth()!=2 or w.getframerate()!=16000 or not 0256:raise ValueError('Ungültiges oder doppeltes Feld.') try:fields[name]=data.decode('utf-8') except UnicodeError:raise ValueError('Ungültiges Textfeld.') from None validate_wav(audio) if fields.get('response_format','json')!='json':raise ValueError('Aktuell wird response_format=json unterstützt.') return fields,audio 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.port=None;self.key=None;self.binary_path=None;self.cancel=threading.Event();self.policy='auto';self.warming=False 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.') directory=self.runtime.root/ident;binary=directory/'build/bin/llama-server' if not binary.is_file():raise ValueError('Aktive llama.cpp-Laufzeit fehlt.') return binary def projector(self): for x in self.catalog.status()['entries']: if x['repo']==PROJECTOR_REPO and x['revision']==PROJECTOR_REVISION and x['file']==PROJECTOR:return self.catalog.entry(x['id']) raise ValueError('Audio-Projektor fehlt. Unter STT → Einrichten herunterladen.') def blockers(self,p): if not supported(p['model']):return ['Diese STT-Variante ist noch nicht angebunden. Unterstützt wird Qwen3-ASR 0.6B Q8_0 aus dem geprüften Repository.'] errors=[] for f in [self.build,self.projector]: try:f() except ValueError as exc:errors.append(str(exc)) return errors def status(self): try:self.build();installed=True except ValueError:installed=False try:self.projector();projector=True 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,loaded=bool(self.process and self.process.poll() is None),policy=self.policy,warming=self.warming) def model_entry(self): for x in self.catalog.status()['entries']: if supported(x):return self.catalog.entry(x['id']) raise ValueError('Passende llama.cpp-Modellvariante fehlt. Unter STT → Einrichten herunterladen.') def assign(self,ident,revision): with self.lock: if self.job and self.job['state']=='running':raise ValueError('Zuerst den STT-Auftrag beenden.') p=next((p for p in self.profiles.status()['profiles'] if p['id']==ident and p['kind']=='stt'),None) if not p:raise ValueError('STT-Profil nicht gefunden.') return self.profiles.save(dict(id=ident,revision=revision,name=p['name'],kind='stt',model_id=self.model_entry()['id'],parameters=p['parameters'])) def setup(self): data=self.catalog.files(REPO) if data['revision']!=REVISION:raise ValueError('Quellversion hat sich geändert. Das Komponentenrezept muss zuerst geprüft werden.') queued=[] for filename in (MODEL,PROJECTOR): if any(x['repo']==REPO and x['revision']==REVISION and x['file']==filename for x in self.catalog.status()['entries']):continue if any(x.get('repo')==REPO and x.get('file')==filename and x['state'] in ('downloading','queued') for x in self.catalog.status().get('downloads',[])):continue queued.append(self.catalog.start(REPO,filename,REVISION,'stt')) return dict(queued=len(queued)) def start(self,profile_id,audio,language='de'): validate_wav(audio) if language not in ('de','en','auto'):raise ValueError('Sprache muss de, en oder auto sein.') with self.lock: if self.warming:raise ValueError('STT wird gerade vorgeladen. Gleich erneut versuchen.') 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.') 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='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: port,key=self._ensure(binary,model,projector) base=f'http://127.0.0.1:{port}';headers={'Authorization':'Bearer '+key};deadline=time.monotonic()+90 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' for name,value in [('model','deck-stt')]+([] if language=='auto' else [('language',language)]):body+=f'--{boundary}\r\nContent-Disposition: form-data; name="{name}"\r\n\r\n{value}\r\n'.encode() body+=f'--{boundary}--\r\n'.encode();headers['Content-Type']='multipart/form-data; boundary='+boundary with urllib.request.urlopen(urllib.request.Request(base+'/v1/audio/transcriptions',data=body,headers=headers),timeout=180) as response: 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' if self.policy=='per_request' else '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 result['state']!='complete' or self.policy=='per_request':self.unload_idle() with self.lock: if self.cancel.is_set():result=dict(state='cancelled',phase='Transkription abgebrochen') self.job.update(result) def _ensure(self,binary,model,projector): process=None try: with self.lock: if self.cancel.is_set():raise InterruptedError() 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 return port,key except Exception: self.unload_idle();raise def warm(self,profile_id): with self.lock: if self.job and self.job.get('state')=='running':return False if self.process and self.process.poll() is None:return True if self.warming:return False p=next((p for p in self.profiles.status()['profiles'] if p['id']==profile_id and p['kind']=='stt' and p['runnable']),None) if not p:raise ValueError('Das ausgewählte STT-Profil ist nicht ausführbar.') self.warming=True;self.cancel.clear() try: headroom=cgroup_headroom() if headroom is not None and headroom<4*1024**3:raise ValueError('Mindestens 4 GiB freier Deck-RAM werden benötigt.') self._ensure(self.build(),self.catalog.root/p['model_id']/'model.gguf',self.catalog.root/self.projector()['id']/'model.gguf');return True finally: with self.lock:self.warming=False 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() self.unload_idle() return dict(cancellation_requested=True)