Add isolated Qwen3 ASR setup test UI and transcription API
This commit is contained in:
@@ -0,0 +1,158 @@
|
||||
"""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 0<w.getnframes()<=120*16000:raise ValueError('WAV muss PCM16, mono, 16 kHz und höchstens 120 Sekunden lang sein.')
|
||||
if len(w.readframes(w.getnframes()))!=w.getnframes()*2:raise ValueError('WAV-Datei ist unvollständig.')
|
||||
except (wave.Error,EOFError):raise ValueError('Ungültige WAV-Datei.') from None
|
||||
|
||||
def read_upload(handler):
|
||||
handler.connection.settimeout(30)
|
||||
if handler.headers.get('Transfer-Encoding'):raise ValueError('Chunked Upload wird nicht unterstützt.')
|
||||
try:length=int(handler.headers.get('Content-Length','0'))
|
||||
except ValueError:raise ValueError('Ungültige Upload-Länge.') from None
|
||||
content_type=handler.headers.get('Content-Type','')
|
||||
if not 0<length<=MAX_AUDIO+65536 or not content_type.lower().startswith('multipart/form-data;'):raise ValueError('Multipart-Upload bis 8 MiB erforderlich.')
|
||||
raw=handler.rfile.read(length)
|
||||
if len(raw)!=length:raise ValueError('Upload unvollständig.')
|
||||
message=BytesParser(policy=policy.default).parsebytes(b'Content-Type: '+content_type.encode()+b'\r\nMIME-Version: 1.0\r\n\r\n'+raw)
|
||||
if not message.is_multipart() or message.defects:raise ValueError('Ungültiger Multipart-Upload.')
|
||||
fields={};audio=None
|
||||
for part in message.iter_parts():
|
||||
name=part.get_param('name',header='content-disposition');data=part.get_payload(decode=True)
|
||||
if not isinstance(data,bytes):raise ValueError('Ungültiges Upload-Feld.')
|
||||
if name=='file':
|
||||
if audio is not None:raise ValueError('Nur eine Audiodatei erlaubt.')
|
||||
audio=data
|
||||
else:
|
||||
if name not in ('model','profile_id','language','response_format') or name in fields or len(data)>256: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.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.')
|
||||
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)
|
||||
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.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.')
|
||||
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)
|
||||
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
|
||||
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'
|
||||
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('<asr_text>')[-1].replace('<|endoftext|>','').strip()
|
||||
result=dict(state='complete',phase='Transkription fertig · Modell entladen',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
|
||||
with self.lock:
|
||||
if self.cancel.is_set():result=dict(state='cancelled',phase='Transkription abgebrochen')
|
||||
self.process=None;self.job.update(result)
|
||||
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
|
||||
return dict(cancellation_requested=True)
|
||||
Reference in New Issue
Block a user