From 92505a4cacaf0dff9158e475f182267f92b1559b Mon Sep 17 00:00:00 2001 From: alf Date: Fri, 19 Jun 2026 14:24:24 +0200 Subject: [PATCH] =?UTF-8?q?feat(voice):=20moteur=20d'enr=C3=B4lement=20?= =?UTF-8?q?=E2=80=94=20worker=20cv=5Fvenv=20chaud=20+=20IPC=20+=20transcri?= =?UTF-8?q?ption?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - worker.py : process persistant sous cv_venv (teacher + Whisper chauds), protocole JSON-lines. Canal stdout propre (fd1→stderr) pour ne pas corrompre l'IPC avec le bruit torch/librosa. - bridge.py : pont côté API py3.14 (sans torch) — lance/pilote le worker en subprocess, jobs sérialisés, démarrage paresseux. - transcribe.py : Whisper sur le segment PRÉPARÉ (texte == segment enrôlé), défaut "small" ("base" transcrit mal le FR, validé). - Validé live : 1 worker chaud sert ping+transcribe(fr)+enroll, réutilisé (même PID), protocole non corrompu (bruit libs → stderr). Co-Authored-By: Claude Opus 4.8 (1M context) --- kazeia_central/voice/bridge.py | 101 +++++++++++++++++++++++++++++ kazeia_central/voice/transcribe.py | 45 +++++++++++++ kazeia_central/voice/worker.py | 69 ++++++++++++++++++++ 3 files changed, 215 insertions(+) create mode 100644 kazeia_central/voice/bridge.py create mode 100644 kazeia_central/voice/transcribe.py create mode 100644 kazeia_central/voice/worker.py diff --git a/kazeia_central/voice/bridge.py b/kazeia_central/voice/bridge.py new file mode 100644 index 0000000..347e765 --- /dev/null +++ b/kazeia_central/voice/bridge.py @@ -0,0 +1,101 @@ +"""Pont API ↔ worker d'enrôlement (côté Python 3.14, SANS torch). + +L'encodeur vit dans cv_venv (torch), l'API en Python 3.14 → on ne peut pas l'importer +en-process. `VoiceBridge` lance et pilote le worker `cv_venv` (`voice/worker.py`) en +subprocess, lui parle en JSON-lines, et sérialise les jobs (un verrou : torch n'est de +toute façon pas thread-safe en inférence concurrente). Démarrage paresseux (le worker +n'est lancé qu'au 1ᵉʳ job → l'API ne paie le coût que si on enrôle). +""" + +from __future__ import annotations + +import json +import os +import subprocess +import threading +from pathlib import Path + +# Interpréteur de l'env encodeur (torch + cosyvoice + whisper + gguf). +CV_PYTHON = os.environ.get("CV_PYTHON", "/opt/Kazeia/cv_venv/bin/python") +_REPO = Path(__file__).resolve().parents[2] # /opt/Kazeia-central + + +class VoiceWorkerError(RuntimeError): + pass + + +class VoiceBridge: + """Gère un worker cv_venv chaud. Thread-safe (jobs sérialisés).""" + + def __init__(self, cv_python: str = CV_PYTHON) -> None: + self.cv_python = cv_python + self._proc: subprocess.Popen | None = None + self._lock = threading.Lock() + self._id = 0 + + def available(self) -> bool: + return os.path.exists(self.cv_python) + + def _ensure(self) -> None: + if self._proc and self._proc.poll() is None: + return + if not os.path.exists(self.cv_python): + raise VoiceWorkerError( + f"interpréteur encodeur introuvable: {self.cv_python} " + f"(provisionner cv_venv, ou définir CV_PYTHON)") + env = dict(os.environ, PYTHONPATH=str(_REPO), PYTHONUNBUFFERED="1") + self._proc = subprocess.Popen( + [self.cv_python, "-m", "kazeia_central.voice.worker"], + stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=None, # stderr → console + text=True, env=env, bufsize=1, + ) + # Attendre l'event "ready" (la 1ʳᵉ ligne du protocole). + while True: + line = self._proc.stdout.readline() + if not line: + raise VoiceWorkerError("le worker s'est arrêté avant d'être prêt") + try: + msg = json.loads(line) + except Exception: + continue + if msg.get("event") == "ready": + return + + def _call(self, method: str, params: dict | None = None) -> dict: + with self._lock: + self._ensure() + self._id += 1 + rid = self._id + assert self._proc and self._proc.stdin and self._proc.stdout + self._proc.stdin.write( + json.dumps({"id": rid, "method": method, "params": params or {}}) + "\n") + self._proc.stdin.flush() + line = self._proc.stdout.readline() + if not line: + self._proc = None + raise VoiceWorkerError(f"worker terminé sans réponse (method={method})") + resp = json.loads(line) + if not resp.get("ok"): + raise VoiceWorkerError(resp.get("error", "échec worker")) + return resp["result"] + + # ---- API publique ----------------------------------------------------- + def ping(self) -> dict: + return self._call("ping") + + def transcribe(self, wav_path: str, model: str | None = None) -> dict: + return self._call("transcribe", {"wav": wav_path, "model": model}) + + def enroll(self, wav_path: str, text: str, out_path: str) -> dict: + return self._call("enroll", {"wav": wav_path, "text": text, "out": out_path}) + + def shutdown(self) -> None: + with self._lock: + if self._proc and self._proc.poll() is None: + try: + self._proc.stdin.write(json.dumps({"id": 0, "method": "shutdown"}) + "\n") + self._proc.stdin.flush() + self._proc.wait(timeout=10) + except Exception: + self._proc.kill() + self._proc = None diff --git a/kazeia_central/voice/transcribe.py b/kazeia_central/voice/transcribe.py new file mode 100644 index 0000000..4f4e06a --- /dev/null +++ b/kazeia_central/voice/transcribe.py @@ -0,0 +1,45 @@ +"""Transcription Whisper (PC) du segment de référence — tourne sous cv_venv. + +Important : on transcrit le **segment préparé** (mono/16k/~15 s), PAS la capture +brute. Le texte du `.cvps` doit correspondre à ce qui est parlé DANS le segment +enrôlé (§3/§4) ; transcrire les 111 s d'origine donnerait un texte qui ne colle pas +aux ~15 s réellement extraits. `prepare_wav` étant déterministe, le segment transcrit +ici == le segment enrôlé. +""" + +from __future__ import annotations + +import os +import tempfile + +from .enroll import prepare_wav + +# Défaut "small" : "base" transcrit mal le FR (validé), "small" est le bon compromis +# qualité/CPU ; "medium" encore mieux mais lourd. La transcription est de toute façon +# relue par l'opérateur. Configurable via CV_WHISPER_MODEL. +CV_WHISPER_MODEL = os.environ.get("CV_WHISPER_MODEL", "small") + +_model_cache: dict[str, object] = {} + + +def _model(size: str): + if size not in _model_cache: + import whisper # openai-whisper (présent dans cv_venv) + _model_cache[size] = whisper.load_model(size) + return _model_cache[size] + + +def transcribe(wav_path: str, model: str | None = None) -> dict: + """Transcrit le segment de référence. Renvoie {text, language, model, prep}.""" + size = model or CV_WHISPER_MODEL + m = _model(size) + with tempfile.TemporaryDirectory() as td: + prep = os.path.join(td, "prep.wav") + prep_info = prepare_wav(wav_path, prep) + result = m.transcribe(prep, fp16=False) # CPU → fp16 off + return { + "text": (result.get("text") or "").strip(), + "language": result.get("language"), + "model": size, + "prep": prep_info, + } diff --git a/kazeia_central/voice/worker.py b/kazeia_central/voice/worker.py new file mode 100644 index 0000000..1f57637 --- /dev/null +++ b/kazeia_central/voice/worker.py @@ -0,0 +1,69 @@ +"""Worker d'enrôlement vocal — process persistant sous cv_venv. + +Lancé par `voice/bridge.py` (côté API py3.14) via : + cv_venv/bin/python -m kazeia_central.voice.worker +Garde le teacher CosyVoice (et le modèle Whisper) **chauds**, et sert une file de +jobs via un protocole JSON-lines sur stdin/stdout. + +Canal protocole propre : la stack torch/librosa/whisper écrit du bruit sur stdout. +On capture le **vrai** stdout (fd 1) pour le protocole AVANT de rediriger fd 1 vers +stderr — ainsi tout `print()`/sortie C de lib part sur stderr (console), et seul le +JSON du protocole circule sur le pipe lu par le bridge. Sans ça, l'IPC serait corrompu. + +Protocole : + requête (1 ligne) : {"id": , "method": "ping|transcribe|enroll|shutdown", "params": {...}} + réponse (1 ligne) : {"id": , "ok": true, "result": {...}} | {"id": , "ok": false, "error": "..."} + au démarrage : {"event": "ready"} +""" + +from __future__ import annotations + +import json +import os +import sys + +# --- canal protocole propre (avant tout import lourd) --------------------- +_proto = os.fdopen(os.dup(1), "w", buffering=1) # vrai stdout (pipe vers le bridge) +os.dup2(2, 1) # fd 1 → stderr : bruit des libs hors protocole + + +def _send(obj: dict) -> None: + _proto.write(json.dumps(obj, ensure_ascii=False) + "\n") + _proto.flush() + + +def _handle(method: str, params: dict): + if method == "ping": + return {"pong": True, "pid": os.getpid()} + if method == "transcribe": + from .transcribe import transcribe + return transcribe(params["wav"], params.get("model")) + if method == "enroll": + from .enroll import enroll + return enroll(params["wav"], params["text"], params["out"]) + raise ValueError(f"méthode inconnue: {method}") + + +def main() -> None: + _send({"event": "ready"}) + for line in sys.stdin: + line = line.strip() + if not line: + continue + try: + req = json.loads(line) + except Exception: + continue + rid = req.get("id") + method = req.get("method") + if method == "shutdown": + _send({"id": rid, "ok": True, "result": {"bye": True}}) + break + try: + _send({"id": rid, "ok": True, "result": _handle(method, req.get("params") or {})}) + except Exception as e: # un job qui échoue ne tue pas le worker + _send({"id": rid, "ok": False, "error": f"{type(e).__name__}: {e}"}) + + +if __name__ == "__main__": + main()