"""Orchestration de l'enrôlement vocal (le workflow complet). Compose : provider (`/voices` → wav_path), transfert adb, bridge (worker ov_venv), store chiffré. Tourne sous l'API py3.14 ; le calcul lourd est délégué au worker. Workflow (VOICE_ENROLLMENT_SPEC §1, décisions 2026-06-19) : connexion → lister voix sans .ovsp → pull WAV → ARCHIVER (store chiffré) → transcrire (ASR OmniVoice, relecture opérateur) → enrôler (.ovsp) → pousser sur la tablette → SUPPRIMER le WAV de la tablette (gardé sur Kazeia-central). Découpé pour permettre la relecture de la transcription : `prepare` (pull+archive+ transcribe) puis `enroll_deploy` (enroll+push+delete). Le WAV est archivé DÈS prepare → jamais de suppression device sans copie PC. """ from __future__ import annotations import os import posixpath import tempfile from typing import Any from ..adb import Adb from ..provider import ProviderClient from ..store import Store from . import admin_ingest, transfer from .bridge import VoiceBridge def _provider_voices(adb: Adb, serial: str): return ProviderClient(adb, serial).voices() def _find_wav_path(adb: Adb, serial: str, voice_id: str) -> str | None: for v in _provider_voices(adb, serial): if v.id == voice_id: return v.wav_path return None def _status(deployed: bool, rec: dict) -> str: if deployed: return "deployed" if rec.get("enrolled_at"): return "enrolled" if rec.get("archived_at"): return "archived" return "pending" def _profiles_using(adb: Adb, serial: str) -> tuple[dict[str, list[str]], set[str]]: """Renvoie ({voice_id: [profil_id,...]}, {profils présents}). Croise profile.voice_id.""" used: dict[str, list[str]] = {} present: set[str] = set() try: for p in ProviderClient(adb, serial).profiles(): present.add(p.id) if p.voice_id: used.setdefault(p.voice_id, []).append(p.id) except Exception: pass return used, present def list_voices(adb: Adb, store: Store, serial: str) -> list[dict[str, Any]]: """Inventaire voix de la tablette croisé avec l'archive store (statut, verrou, profils utilisateurs) par voix.""" rows = _provider_voices(adb, serial) try: deployed = transfer.deployed_ovsp(adb, serial, _patient_omnivoice_dir(adb, serial)) except Exception: deployed = set() recs = {r["voice_id"]: r for r in store.voices(serial)} if store.is_unlocked else {} used_by, _ = _profiles_using(adb, serial) out = [] for v in rows: rec = recs.get(v.id, {}) is_dep = v.id in deployed locked = rec.get("locked_profile_id") out.append({ "voice_id": v.id, "wav_path": v.wav_path, "wav_on_device": v.wav_exists, "deployed_on_device": is_dep, "archived": bool(rec.get("archived_at")), "enrolled": bool(rec.get("enrolled_at")), "wav_deleted": bool(rec.get("wav_deleted_at")), "transcription": rec.get("transcription"), "language": rec.get("language"), "locked_profile_id": locked, # exclusif à ce profil, ou None = généraliste "used_by_profiles": used_by.get(v.id, []), # profils qui sélectionnent cette voix "status": _status(is_dep, rec), }) return out def sync_device(adb: Adb, store: Store, bridge: VoiceBridge, serial: str, *, delete_source: bool, now: int, actor: str | None = None) -> list[dict[str, Any]]: """Enrôle + déploie les voix en attente d'une tablette (rapport par voix). Verrou-conscient : ne déploie pas une voix exclusive sur une tablette dépourvue du profil propriétaire. Transcrit automatiquement si pas encore relu (la relecture fine reste possible voix par voix via prepare/enroll).""" _, present = _profiles_using(adb, serial) report: list[dict[str, Any]] = [] for v in list_voices(adb, store, serial): vid = v["voice_id"] if v["deployed_on_device"]: continue lock = v["locked_profile_id"] if lock and lock not in present: report.append({"voice_id": vid, "ok": False, "skipped": "verrou_profil_absent", "profile": lock}) continue try: text = v["transcription"] if not text: text = prepare(adb, store, bridge, serial, vid, now=now, actor=actor)["transcription"] res = enroll_deploy(adb, store, bridge, serial, vid, text, delete_source=delete_source, now=now, actor=actor) report.append({"voice_id": vid, "ok": True, **res}) except Exception as e: report.append({"voice_id": vid, "ok": False, "error": str(e)}) return report def _ensure_wav_archived(adb: Adb, store: Store, serial: str, voice_id: str, *, now: int) -> bytes: """Garantit que le WAV est archivé (chiffré) côté PC ; le renvoie déchiffré.""" data = store.voice_wav_bytes(serial, voice_id) if data is not None: return data wav_path = _find_wav_path(adb, serial, voice_id) if not wav_path: raise ValueError(f"voix introuvable sur la tablette: {voice_id}") with tempfile.TemporaryDirectory() as td: local = os.path.join(td, f"{voice_id}.wav") transfer.pull_wav(adb, serial, wav_path, local) data = open(local, "rb").read() store.archive_voice_wav(serial, voice_id, data, source_wav_path=wav_path, now=now) return data def prepare(adb: Adb, store: Store, bridge: VoiceBridge, serial: str, voice_id: str, *, now: int, actor: str | None = None) -> dict[str, Any]: """Pull + archive (chiffré) + transcription (ASR OmniVoice). Renvoie le texte à relire.""" data = _ensure_wav_archived(adb, store, serial, voice_id, now=now) with tempfile.TemporaryDirectory() as td: local = os.path.join(td, f"{voice_id}.wav") open(local, "wb").write(data) tr = bridge.transcribe(local) store.set_voice_transcription(serial, voice_id, tr["text"], tr.get("language"), now=now) store.audit("voice_prepare", actor=actor, target=f"{serial}/{voice_id}", now=now) return {"voice_id": voice_id, "transcription": tr["text"], "language": tr.get("language"), "model": tr.get("model")} def enroll_deploy(adb: Adb, store: Store, bridge: VoiceBridge, serial: str, voice_id: str, transcription: str, *, delete_source: bool, now: int, actor: str | None = None) -> dict[str, Any]: """Enrôle (.ovsp) → pousse sur la tablette → (option) supprime le WAV device. Le WAV reste archivé côté PC → suppression device réversible par ré-enrôlement.""" data = _ensure_wav_archived(adb, store, serial, voice_id, now=now) wav_path = _find_wav_path(adb, serial, voice_id) ov = _patient_omnivoice_dir(adb, serial) with tempfile.TemporaryDirectory() as td: local = os.path.join(td, f"{voice_id}.wav") open(local, "wb").write(data) ovsp = os.path.join(td, f"{voice_id}.ovsp") info = bridge.enroll(local, transcription, ovsp) store.set_voice_transcription(serial, voice_id, transcription, None, now=now) store.mark_voice_enrolled(serial, voice_id, info["bytes"], now=now) remote = transfer.push_ovsp(adb, serial, ovsp, ov, voice_id) store.mark_voice_deployed(serial, voice_id, now=now) store.audit("voice_enroll_deploy", actor=actor, target=f"{serial}/{voice_id}", detail={"ovsp_bytes": info["bytes"]}, now=now) result = {"voice_id": voice_id, "ovsp_bytes": info["bytes"], "deployed_to": remote, "wav_deleted": False} if delete_source and wav_path: transfer.delete_device_wav(adb, serial, wav_path) store.mark_voice_wav_deleted(serial, voice_id, now=now) store.audit("voice_wav_deleted", actor=actor, target=f"{serial}/{voice_id}", now=now) result["wav_deleted"] = True return result # ===== Ingestion depuis le stockage ADMIN (manifestes record-time, §4.2) ===== # Source de vérité des voix fraîches : `Android/data/com.kazeia.admin/files/voix/` # (scope/owner/reference_text ensemble), pas le provider patient `/voices`. # MODELS_DIR candidats côté patient (dev legacy / prod scoped-storage). _CANDIDATE_MODELS_DIRS = ( "/data/local/tmp/kazeia/models", "/sdcard/Android/data/com.kazeia/files/kazeia/models", ) def _patient_omnivoice_dir(adb: Adb, serial: str) -> str: """Dossier `.ovsp` cible côté patient (`/omnivoice/voices`). Les manifestes admin ne le donnent pas, et post-migration `/voices.wav_path` est VIDE pour une voix déployée → on dérive en cascade : override env → wav_path legacy (marqueur MODELS_DIR) → sonde des dossiers modèles connus sur le device.""" env = os.environ.get("KAZEIA_OMNIVOICE_DIR") if env: return env for r in _provider_voices(adb, serial): if r.wav_path and transfer._WAV_MARKER in r.wav_path: return transfer.omnivoice_dir_for(r.wav_path) for base in _CANDIDATE_MODELS_DIRS: if base in adb.shell(f"ls -d {base} 2>/dev/null", serial=serial): return posixpath.join(base, "omnivoice", "voices") raise ValueError( "dossier omnivoice indéterminable (ni /voices wav_path, ni dossier modèles connu ; " "définir KAZEIA_OMNIVOICE_DIR)") def _admin_status(man: dict, rec: dict, deployed: bool) -> str: if man["scope"] == "pending": return "pending" # ni enrôlée ni déployable tant que non résolue if deployed: return "deployed" if rec.get("enrolled_at"): return "enrolled" if rec.get("archived_at"): return "archived" return "to_enroll" def list_admin_voices(adb: Adb, store: Store, serial: str) -> list[dict[str, Any]]: """Voix enregistrées via l'admin (manifestes) croisées avec l'archive store et l'état de déploiement. C'est la liste à piloter depuis l'UI.""" manifests = admin_ingest.list_manifests(adb, serial) recs = {r["voice_id"]: r for r in store.voices(serial)} if store.is_unlocked else {} try: deployed = transfer.deployed_ovsp(adb, serial, _patient_omnivoice_dir(adb, serial)) except Exception: deployed = set() out = [] for man in manifests: rec = recs.get(man["id"], {}) is_dep = man["id"] in deployed out.append({ "voice_id": man["id"], "name": man["name"], "scope": man["scope"], "owner_profile_id": man["owner_profile_id"], "owner_name": man["owner_name"], "reference_text": man["reference_text"], "duration_seconds": man["duration_seconds"], "deployed_on_device": is_dep, "archived": bool(rec.get("archived_at")), "enrolled": bool(rec.get("enrolled_at")), "status": _admin_status(man, rec, is_dep), }) return out def ingest_voice(adb: Adb, store: Store, bridge: VoiceBridge, serial: str, manifest: dict, *, delete_source: bool, now: int, actor: str | None = None) -> dict[str, Any]: """Ingère UNE voix admin : pull WAV → archive chiffrée → métadonnées (scope/owner/ consentement) → verrou si exclusive → enrôle (ASR du 16 s, §6) → déploie selon scope. `pending` = archivée mais NI enrôlée NI déployée (cadre à fixer à l'attribution).""" vid = manifest["id"] scope = manifest["scope"] owner = manifest.get("owner_profile_id") # 1) Archiver le WAV (depuis le stockage admin) s'il ne l'est pas déjà. data = store.voice_wav_bytes(serial, vid) if data is None: with tempfile.TemporaryDirectory() as td: local = os.path.join(td, f"{vid}.wav") transfer.pull_wav(adb, serial, manifest["wav_device_path"], local) data = open(local, "rb").read() store.archive_voice_wav(serial, vid, data, source_wav_path=manifest["wav_device_path"], now=now) # 2) Métadonnées manifeste (scope/propriétaire/consentement chiffré) + verrou. store.set_voice_manifest(serial, vid, scope=scope, owner_name=manifest.get("owner_name"), consent_text=manifest.get("reference_text"), now=now) if scope == "exclusive" and owner: store.lock_voice(serial, vid, owner, now=now) store.audit("voice_admin_ingest", actor=actor, target=f"{serial}/{vid}", detail={"scope": scope, "owner": owner}, now=now) # 3) pending → on s'arrête (archivée, à résoudre plus tard). if scope == "pending": return {"voice_id": vid, "scope": scope, "deployed": False, "reason": "pending"} # 4) exclusive → le profil propriétaire doit être présent sur la tablette. if scope == "exclusive" and owner: _, present = _profiles_using(adb, serial) if owner not in present: return {"voice_id": vid, "scope": scope, "deployed": False, "reason": "verrou_profil_absent", "profile": owner} # 5) Enrôler (texte vide → ASR auto du segment 16 s, §6) + déployer sur la tablette. ov = _patient_omnivoice_dir(adb, serial) with tempfile.TemporaryDirectory() as td: local = os.path.join(td, f"{vid}.wav") open(local, "wb").write(data) ovsp = os.path.join(td, f"{vid}.ovsp") info = bridge.enroll(local, "", ovsp) store.set_voice_transcription(serial, vid, info.get("ref_text", ""), None, now=now) store.mark_voice_enrolled(serial, vid, info["bytes"], now=now) remote = transfer.push_ovsp(adb, serial, ovsp, ov, vid) store.mark_voice_deployed(serial, vid, now=now) store.audit("voice_enroll_deploy", actor=actor, target=f"{serial}/{vid}", detail={"ovsp_bytes": info["bytes"], "scope": scope}, now=now) result = {"voice_id": vid, "scope": scope, "deployed": True, "deployed_to": remote, "ovsp_bytes": info["bytes"], "wav_deleted": False} if delete_source: admin_ingest.delete_admin_voice(adb, serial, manifest) store.mark_voice_wav_deleted(serial, vid, now=now) store.audit("voice_wav_deleted", actor=actor, target=f"{serial}/{vid}", now=now) result["wav_deleted"] = True return result def sync_admin(adb: Adb, store: Store, bridge: VoiceBridge, serial: str, *, delete_source: bool, now: int, actor: str | None = None) -> list[dict[str, Any]]: """Ingère toutes les voix admin en attente d'une tablette (rapport par voix).""" report: list[dict[str, Any]] = [] for man in admin_ingest.list_manifests(adb, serial): try: report.append({"ok": True, **ingest_voice(adb, store, bridge, serial, man, delete_source=delete_source, now=now, actor=actor)}) except Exception as e: report.append({"ok": False, "voice_id": man.get("id"), "error": str(e)}) return report