"""Потоковая озвучка ответа: фразы синтезируются и проигрываются по мере генерации текста.""" from __future__ import annotations import logging import queue import threading from collections.abc import Callable from typing import Protocol import numpy as np from PySide6.QtCore import QObject, Signal, Slot from agr_assistent.tts.base import TTSEngine from agr_assistent.tts.text import SpeechTextStream log = logging.getLogger(__name__) _END = object() # конец реплики _LOAD = object() # предзагрузка модели _NO_GENERATION = -1 # ошибка, не привязанная к реплике _IDLE_POLL_SECONDS = 0.25 class Player(Protocol): def play( self, audio: np.ndarray, sample_rate: int, should_continue: Callable[[], bool] ) -> None: ... def finish(self) -> None: ... def abort(self) -> None: ... class Speaker(QObject): """Живёт в главном потоке; синтез и воспроизведение — в двух фоновых потоках. Синтез идёт параллельно с воспроизведением, поэтому следующая фраза обычно готова к моменту, когда доиграла предыдущая. """ playback_started = Signal() finished = Signal() # реплика доиграна или остановлена error_occurred = Signal(str) # Мост из фоновых потоков в главный; int — номер реплики _worker_started = Signal(int) _worker_done = Signal(int) _worker_failed = Signal(int, str) def __init__( self, engine: TTSEngine, player: Player, *, enabled: bool, parent: QObject | None = None, ) -> None: super().__init__(parent) self._engine = engine self._player = player self._enabled = enabled # Увеличивается при каждой новой реплике и остановке: устаревшие фразы отбрасываются self._generation = 0 self._active = False self._playing = False self._text: SpeechTextStream | None = None self._synthesis_queue: queue.Queue[tuple[int, object]] = queue.Queue() self._audio_queue: queue.Queue[tuple[int, object]] = queue.Queue() self._worker_started.connect(self._on_worker_started) self._worker_done.connect(self._on_worker_done) self._worker_failed.connect(self._on_worker_failed) threading.Thread(target=self._synthesis_loop, name="tts-synthesis", daemon=True).start() threading.Thread(target=self._playback_loop, name="tts-playback", daemon=True).start() @property def enabled(self) -> bool: return self._enabled @property def is_active(self) -> bool: """Реплика начата и ещё не доиграна.""" return self._active @property def is_playing(self) -> bool: return self._playing def set_enabled(self, enabled: bool) -> None: self._enabled = enabled if enabled: self.warm_up() else: self.stop() def warm_up(self) -> None: """Загружает модель заранее, чтобы первый ответ не ждал её.""" self._synthesis_queue.put((_NO_GENERATION, _LOAD)) def begin(self) -> None: self.stop() if not self._enabled: return self._generation += 1 self._active = True self._text = SpeechTextStream() def feed(self, chunk: str) -> None: if self._text is not None: for sentence in self._text.feed(chunk): self._synthesis_queue.put((self._generation, sentence)) def end(self) -> None: if self._text is not None: for sentence in self._text.flush(): self._synthesis_queue.put((self._generation, sentence)) self._synthesis_queue.put((self._generation, _END)) self._text = None def stop(self) -> None: if not self._active: return self._generation += 1 self._finish() def _finish(self) -> None: self._text = None self._active = False self._playing = False self.finished.emit() def _synthesis_loop(self) -> None: """Фоновый поток: текст -> звук.""" while True: generation, item = self._synthesis_queue.get() if item is _LOAD: try: self._engine.load() except Exception as exc: log.exception("Не удалось загрузить модель синтеза речи") self._worker_failed.emit(_NO_GENERATION, str(exc)) continue if generation != self._generation: continue if item is _END: self._audio_queue.put((generation, _END)) continue try: audio = self._engine.synthesize(str(item)) except Exception as exc: log.exception("Сбой синтеза речи") self._worker_failed.emit(generation, f"Ошибка синтеза речи: {exc}") continue self._audio_queue.put((generation, audio)) def _playback_loop(self) -> None: """Фоновый поток: воспроизведение звука по порядку.""" started_generation: int | None = None while True: try: generation, item = self._audio_queue.get(timeout=_IDLE_POLL_SECONDS) except queue.Empty: # Реплику остановили между фразами — освобождаем аудиоустройство if started_generation is not None and started_generation != self._generation: self._player.abort() started_generation = None continue if generation != self._generation: continue try: if item is _END: self._player.finish() started_generation = None self._worker_done.emit(generation) continue if started_generation != generation: started_generation = generation self._worker_started.emit(generation) self._player.play( item, # type: ignore[arg-type] self._engine.sample_rate, lambda: generation == self._generation, ) except Exception as exc: log.exception("Сбой воспроизведения") self._player.abort() started_generation = None self._worker_failed.emit(generation, f"Ошибка воспроизведения звука: {exc}") @Slot(int) def _on_worker_started(self, generation: int) -> None: if generation == self._generation and self._active: self._playing = True self.playback_started.emit() @Slot(int) def _on_worker_done(self, generation: int) -> None: if generation == self._generation and self._active: self._finish() @Slot(int, str) def _on_worker_failed(self, generation: int, message: str) -> None: if generation == self._generation: self.stop() elif generation != _NO_GENERATION: return self.error_occurred.emit(message)