422 lines
22 KiB
Python
422 lines
22 KiB
Python
"""Машина состояний голосового ассистента (сердце Модуля 5).
|
||
|
||
Состояния и переходы:
|
||
|
||
idle ──talk_down──▶ listening ──talk_up──▶ thinking ──готово──▶ speaking ──конец──▶ idle
|
||
▲ │ │ │
|
||
└────── stop ────────┴───────────────────────┴────────────────────┘
|
||
|
||
Правила по ТЗ:
|
||
- если ассистент ГОВОРИТ, а отец нажал «Слушай» — речь мгновенно стихает
|
||
и начинается запись (одним нажатием);
|
||
- «Замолчи» гасит речь / отменяет запись / отбрасывает «думание»;
|
||
- «Повтори» — последняя фраза ассистента (в состоянии idle);
|
||
- удержание клавиши «Слушай» короче min_seconds отбрасывается.
|
||
|
||
Тяжёлая работа (STT→LLM→TTS) идёт в фоновом потоке; устаревшие результаты
|
||
отбрасываются по счётчику поколений _generation (если отец прервал и спросил
|
||
заново — старый ответ не прозвучит).
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import threading
|
||
import time
|
||
from pathlib import Path
|
||
from typing import Callable, Optional, Tuple
|
||
|
||
import numpy as np
|
||
|
||
from .errors import GENERAL_CUE, UNCLEAR_CUE, UNCLEAR_LIMIT, classify_error, cue_path
|
||
from .farewell import FAREWELL_REPLY, is_farewell, summarize_and_save
|
||
from .memory import MemoryStore
|
||
from .wake_word import WakeWordWatcher, matches_wake_word
|
||
|
||
|
||
def resample_for_player(audio: np.ndarray, source_rate: int, target_rate: int) -> np.ndarray:
|
||
"""Чанк в другой частоте (для склейки); edge всегда 24k — на всякий случай."""
|
||
from ..audio_io.resample import resample_to_16k
|
||
return resample_to_16k(audio, source_rate, target_rate)
|
||
|
||
|
||
class Assistant:
|
||
IDLE = "idle" # ждёт нажатия
|
||
LISTENING = "listening" # запись голоса
|
||
THINKING = "thinking" # STT → LLM → TTS
|
||
SPEAKING = "speaking" # воспроизведение ответа
|
||
|
||
def __init__(
|
||
self,
|
||
stt, brain, tts, # провайдеры модулей 1/4/3
|
||
recorder, player, # модуль 2
|
||
memory: Optional[MemoryStore] = None, # долговременная память (М5.1)
|
||
history_limit: int = 10, # реплик диалога в контексте (2×N сообщений)
|
||
on_state_change: Optional[Callable[[str, str], None]] = None, # для earcons М7
|
||
on_log: Optional[Callable[[str], None]] = None,
|
||
) -> None:
|
||
self.stt = stt
|
||
self.brain = brain
|
||
self.tts = tts
|
||
self.recorder = recorder
|
||
self.player = player
|
||
self.memory = memory or MemoryStore()
|
||
self.history_limit = history_limit
|
||
self.session_dialog: list[dict] = [] # реплики текущего разговора (для конспекта)
|
||
self.on_state_change = on_state_change or (lambda state, note: None)
|
||
self.on_log = on_log or (lambda msg: print(f" {msg}"))
|
||
# Долговременная память → в системный промпт LLM при каждом вопросе
|
||
memory_block = self.memory.as_prompt_block()
|
||
if memory_block:
|
||
self.brain.context_block = memory_block
|
||
self._log(f"Память загружена ({len(self.memory.load())} симв.)")
|
||
|
||
self.history: list[dict] = []
|
||
self._state = self.IDLE
|
||
self._lock = threading.Lock()
|
||
self._generation = 0 # инкремент при прерываниях; устаревшие воркеры молчат
|
||
self._last_speech: Optional[Tuple[np.ndarray, int]] = None
|
||
# Свободный режим (HandsFreeRecorder): событие «фраза закончена» от VAD
|
||
self._vad_mode = hasattr(recorder, "suppress") # duck-typing: HandsFreeRecorder
|
||
if self._vad_mode:
|
||
recorder.on_phrase = self._on_vad_phrase # публичный атрибут — переписываем
|
||
# СЕАНС ДИАЛОГА (пробел-тумблер): вне сеанса микрофон глушится
|
||
self._in_dialog = False
|
||
self._cues_dir = Path(__file__).resolve().parents[2] / "assets" / "earcons"
|
||
self._unclear_streak = 0 # неразборчивых фраз подряд (после UNCLEAR_LIMIT — куи)
|
||
# Wake-word: вне диалога слушаем тихо и ждём «Злата» (или своё из .env)
|
||
self.wake_word = WakeWordWatcher(on_detected=self._open_dialog)
|
||
if self._vad_mode:
|
||
recorder.suppress(False) # слушаем сразу — ждём wake-word
|
||
# Прогрев TTS-соединения и аудио-выхода в фоне (первый ответ звучит быстрее)
|
||
threading.Thread(target=self._warmup, daemon=True).start()
|
||
|
||
def _warmup(self) -> None:
|
||
"""Разогреть TTS (DNS/TLS до сервиса синтеза) и аудио-выход — беззвучно."""
|
||
try:
|
||
warm = getattr(self.tts, "warmup", None)
|
||
if warm:
|
||
warm()
|
||
except Exception:
|
||
pass
|
||
try:
|
||
import numpy as _np
|
||
self.player.play(_np.zeros(1600, dtype=_np.float32), 24000) # 67 мс тишины
|
||
except Exception:
|
||
pass
|
||
|
||
# ------------------------------------------------------------------ state
|
||
@property
|
||
def state(self) -> str:
|
||
return self._state
|
||
|
||
def _set_state(self, new_state: str, note: str = "") -> None:
|
||
with self._lock:
|
||
self._state = new_state
|
||
self.on_state_change(new_state, note)
|
||
|
||
def _log(self, msg: str) -> None:
|
||
self.on_log(msg)
|
||
|
||
# ------------------------------------------------- сеанс диалога (пробел)
|
||
def on_dialog_toggle(self) -> None:
|
||
"""ПРОБЕЛ: начать диалог или завершить его (куи-файлы — мгновенно)."""
|
||
if self._in_dialog:
|
||
self._close_dialog(play_end_cue=True)
|
||
else:
|
||
self._open_dialog()
|
||
|
||
def _open_dialog(self) -> None:
|
||
if self._state == self.SPEAKING:
|
||
self.player.stop()
|
||
if self._state == self.THINKING:
|
||
self._generation += 1 # старое «думание» отменяется
|
||
self._in_dialog = True
|
||
self._set_state(self.IDLE, f"Диалог начат (wake-word «{self.wake_word.wake_word}»)")
|
||
if self._vad_mode:
|
||
# микрофон слушает, но на время куи глушим (чтобы не услышать саму себя)
|
||
self.recorder.suppress(True)
|
||
self._log("Диалог начат (wake-word)")
|
||
self._play_cue("dialog_start.wav", unpause_after=True)
|
||
|
||
def _close_dialog(self, play_end_cue: bool) -> None:
|
||
if not self._in_dialog:
|
||
return
|
||
self._generation += 1 # всё текущее (запись/думание/речь) отменяется
|
||
self.player.stop()
|
||
if self._vad_mode:
|
||
self.recorder.suppress(False) # сбрасываем недозапись и слушаем снова
|
||
self._in_dialog = False
|
||
self._set_state(self.IDLE, f"Диалог завершён — жду «{self.wake_word.wake_word}»")
|
||
self._log("Диалог завершён — в режиме ожидания wake-word")
|
||
if play_end_cue:
|
||
self._play_cue("dialog_end.wav")
|
||
|
||
def _play_cue(self, filename: str, unpause_after: bool = False) -> None:
|
||
"""Мгновенно проиграть готовый wav-куи (без синтеза)."""
|
||
path = self._cues_dir / filename
|
||
if not path.is_file():
|
||
self._log(f"(куи-файл не найден: {filename})")
|
||
if unpause_after and self._vad_mode:
|
||
self.recorder.suppress(False)
|
||
return
|
||
|
||
def _watch() -> None:
|
||
while self.player.is_playing:
|
||
time.sleep(0.05)
|
||
if unpause_after and self._in_dialog and self._vad_mode:
|
||
self.recorder.suppress(False) # куи прозвучал — слушаем отца
|
||
|
||
self.player.play_file(str(path))
|
||
threading.Thread(target=_watch, daemon=True).start()
|
||
|
||
# ----------------------------------------------------- свободный режим
|
||
def _on_vad_phrase(self, audio: np.ndarray, samplerate: int) -> None:
|
||
"""VAD прислал законченную фразу (пауза тишины после речи).
|
||
|
||
Вне диалога: STT → если фраза начинается с wake-word («Злата») —
|
||
открыть диалог. Иначе молча игнорируем (фоновый разговор не будит).
|
||
В диалоге — обрабатываем как вопрос (перебивание работает).
|
||
"""
|
||
if not self._in_dialog:
|
||
self._generation += 1
|
||
generation = self._generation
|
||
self._set_state(self.LISTENING, "жду wake-word")
|
||
threading.Thread(
|
||
target=self._wake_check, args=(audio, generation), daemon=True
|
||
).start()
|
||
return
|
||
state = self._state
|
||
if state == self.SPEAKING:
|
||
self.player.stop()
|
||
if state == self.THINKING:
|
||
self._generation += 1
|
||
self._log("Перебиваю размышление новой фразой")
|
||
self._generation += 1
|
||
generation = self._generation
|
||
self._set_state(self.THINKING)
|
||
threading.Thread(target=self._pipeline, args=(audio, generation), daemon=True).start()
|
||
|
||
def _wake_check(self, audio: np.ndarray, generation: int) -> None:
|
||
"""Проверка фразы на wake-word (вне диалога). Дешёво: только STT."""
|
||
try:
|
||
t0 = time.perf_counter()
|
||
res = self.stt.transcribe(audio)
|
||
t_stt = time.perf_counter() - t0
|
||
text = res.text.strip()
|
||
if generation != self._generation:
|
||
return
|
||
if matches_wake_word(text, self.wake_word.wake_word):
|
||
self._log(f"Wake-word услышан ({t_stt:.1f} c): {text!r}")
|
||
self._open_dialog()
|
||
else:
|
||
# чужая речь/телевизор — молча, без реакций
|
||
self._log(f"(вне диалога фраза мимо wake-word: {text[:40]!r})")
|
||
self._set_state(self.IDLE, "жду wake-word")
|
||
except Exception as exc:
|
||
self._log(f"wake-check: {exc}")
|
||
|
||
def _speak(self, audio: np.ndarray, samplerate: int) -> None:
|
||
"""Озвучить и вернуть состояние в idle по окончании (или после прерывания)."""
|
||
generation = self._generation
|
||
self._set_state(self.SPEAKING)
|
||
# пока говорим — микрофон не слушает (не слышим сами себя из динамиков)
|
||
if self._vad_mode:
|
||
self.recorder.suppress(True)
|
||
self.player.play(audio, samplerate)
|
||
|
||
def _watch() -> None:
|
||
while self.player.is_playing:
|
||
time.sleep(0.05)
|
||
# сюда попадаем и при player.stop() — тогда state уже IDLE
|
||
if generation == self._generation and self._state == self.SPEAKING:
|
||
self._set_state(self.IDLE)
|
||
if self._vad_mode and generation == self._generation:
|
||
self.recorder.suppress(False) # снова слушаем
|
||
|
||
threading.Thread(target=_watch, daemon=True).start()
|
||
|
||
# ------------------------------------------------------------- события
|
||
def on_talk_down(self) -> None:
|
||
"""Клавиша «Слушай» нажата."""
|
||
state = self._state
|
||
|
||
if state == self.SPEAKING:
|
||
self.player.stop() # мгновенно стихаем (по ТЗ)
|
||
if state == self.THINKING:
|
||
self._generation += 1 # текущий ответ станет устаревшим
|
||
self._log("Прервал размышление — начинаю запись")
|
||
|
||
if state in (self.IDLE, self.SPEAKING, self.THINKING):
|
||
try:
|
||
self.recorder.start()
|
||
except RuntimeError as exc:
|
||
self._set_state(self.IDLE, f"микрофон недоступен: {exc}")
|
||
return
|
||
self._set_state(self.LISTENING)
|
||
self._log("Слушаю…")
|
||
# в LISTENING повторный down игнорируем (автодубли клавиши)
|
||
|
||
def on_talk_up(self) -> None:
|
||
"""Клавиша «Слушай» отпущена — запись закончена, обрабатываем."""
|
||
if self._state != self.LISTENING:
|
||
return
|
||
|
||
audio = self.recorder.stop()
|
||
if len(audio) == 0:
|
||
self._set_state(self.IDLE, "Слишком коротко — ничего не записал")
|
||
return
|
||
|
||
self._generation += 1
|
||
generation = self._generation
|
||
self._set_state(self.THINKING)
|
||
threading.Thread(target=self._pipeline, args=(audio, generation), daemon=True).start()
|
||
|
||
def on_stop(self) -> None:
|
||
"""Клавиша «Замолчи» — мгновенно тишина."""
|
||
state = self._state
|
||
self._generation += 1 # отбрасываем все текущие работы
|
||
if state == self.SPEAKING:
|
||
self.player.stop()
|
||
self._log("Прервал речь")
|
||
elif state == self.LISTENING:
|
||
self.recorder.cancel()
|
||
self._log("Запись отменена")
|
||
elif state == self.THINKING:
|
||
self._log("Размышление отменено")
|
||
self._set_state(self.IDLE)
|
||
|
||
def on_repeat(self) -> None:
|
||
"""Клавиша «Повтори» — снова озвучить последний ответ."""
|
||
if self._state != self.IDLE:
|
||
self._log("«Повтори» работает только в режиме ожидания")
|
||
return
|
||
if not self._last_speech:
|
||
self._log("Пока нечего повторять")
|
||
return
|
||
audio, samplerate = self._last_speech
|
||
self._log("Повторяю последний ответ")
|
||
self._speak(audio, samplerate)
|
||
|
||
# ------------------------------------------------------------ конвейер
|
||
def _pipeline(self, audio: np.ndarray, generation: int) -> None:
|
||
"""STT → LLM → TTS → озвучка (фоновый поток). Тайминги этапов — в лог."""
|
||
try:
|
||
t0 = time.perf_counter()
|
||
stt_res = self.stt.transcribe(audio)
|
||
t_stt = time.perf_counter() - t0
|
||
|
||
question = stt_res.text.strip()
|
||
if generation != self._generation:
|
||
return # устарело (прервали)
|
||
if not question or len(question) < 2:
|
||
self._handle_unclear("Речь не распознана")
|
||
return
|
||
self._unclear_streak = 0 # фраза распознана — счётчик в ноль
|
||
self._log(f"Вы: {question}")
|
||
|
||
# Прощание: отвечаем сразу, пишем конспект и ЗАКРЫВАЕМ сеанс диалога
|
||
if is_farewell(question) and self.session_dialog:
|
||
self._log("Прощание — записываю разговор в память")
|
||
threading.Thread(
|
||
target=summarize_and_save,
|
||
args=(self.brain, self.session_dialog + [{"role": "user", "text": question}], self.memory),
|
||
kwargs={"on_error": lambda msg: self._log(msg)},
|
||
daemon=True,
|
||
).start()
|
||
answer = FAREWELL_REPLY
|
||
answer_for_history = answer
|
||
llm = None
|
||
t_llm = 0.0
|
||
else:
|
||
t0 = time.perf_counter()
|
||
llm = self.brain.ask(question, history=self.history)
|
||
t_llm = time.perf_counter() - t0
|
||
if generation != self._generation:
|
||
return
|
||
answer = llm.cleaned_text # для ушей (без служебных маркеров)
|
||
answer_for_history = llm.raw_text # маркер [ЧАСТЬ i ИЗ n] остаётся — модель помнит часть
|
||
if not answer:
|
||
answer = "Извините, я задумалась и забыла, что хотела сказать. Спросите ещё раз."
|
||
answer_for_history = answer
|
||
self._log(f"Ассистент ({t_llm:.1f} с): {answer}")
|
||
|
||
# История и конспект-диалог пополняются одинаково для обоих путей озвучки
|
||
turn = [
|
||
{"role": "user", "text": question},
|
||
{"role": "assistant", "text": answer_for_history if not is_farewell(question) else answer},
|
||
]
|
||
self.history += turn
|
||
self.history = self.history[-self.history_limit:]
|
||
self.session_dialog += turn
|
||
|
||
# «Отвечаю» — сразу (синтез идёт далее в фоне, звук стартует по готовности)
|
||
self._log("Отвечаю...")
|
||
|
||
# Потоковая озвучка: чанки-предложения по мере готовности.
|
||
# Если синтезатор не умеет stream() — обычный путь (целиком).
|
||
streamer = getattr(self.tts, "stream", None)
|
||
if streamer is not None and generation == self._generation:
|
||
self._set_state(self.SPEAKING)
|
||
if self._vad_mode:
|
||
self.recorder.suppress(True)
|
||
self._last_speech = None # заполним по чанкам
|
||
|
||
pause = getattr(self.tts, "_pause", 0.0)
|
||
sr_prev = None
|
||
for audio, sr in streamer(answer):
|
||
if generation != self._generation: # перебили — гасим поток
|
||
self.player.stop()
|
||
break
|
||
piece = audio
|
||
if sr_prev is not None and sr != sr_prev:
|
||
piece = resample_for_player(audio, sr, sr_prev) # Rare: edge всегда 24k
|
||
sr_prev = sr
|
||
if self._last_speech is None:
|
||
self._last_speech = (piece, sr)
|
||
else:
|
||
prev, prev_sr = self._last_speech
|
||
gap = np.zeros(int(pause * prev_sr), dtype=np.float32)
|
||
self._last_speech = (np.concatenate([prev, gap, piece]), sr)
|
||
self.player.play(piece, sr)
|
||
# ждём конца чанка (если перебили — generation сместится и выйдем)
|
||
while self.player.is_playing and generation == self._generation:
|
||
time.sleep(0.05)
|
||
if generation != self._generation:
|
||
break
|
||
if generation == self._generation:
|
||
self._set_state(self.IDLE)
|
||
# Прощание: после ответа закрываем сеанс (микрофон глушится)
|
||
if is_farewell(question) and self._in_dialog:
|
||
self._close_dialog(play_end_cue=False)
|
||
elif self._vad_mode:
|
||
self.recorder.suppress(False)
|
||
return
|
||
|
||
t0 = time.perf_counter()
|
||
speech = self.tts.synthesize(answer)
|
||
t_tts = time.perf_counter() - t0
|
||
if generation != self._generation:
|
||
return
|
||
|
||
self._last_speech = (speech.audio, speech.samplerate)
|
||
|
||
self._speak(speech.audio, speech.samplerate)
|
||
t_play = getattr(self.player, "last_start_latency", 0.0)
|
||
self._log(f"тайминг: распознавание {t_stt:.1f} | модель {t_llm:.1f} | "
|
||
f"синтез {t_tts:.1f} | старт звука {t_play:.1f}")
|
||
except Exception as exc: # сеть, API, TTS — что угодно
|
||
if generation == self._generation:
|
||
self._log(f"Ошибка: {exc}")
|
||
cue = classify_error(exc)
|
||
self._set_state(self.IDLE, f"Ошибка: {cue}")
|
||
self._play_cue(cue)
|
||
|
||
# --------------------------------------------------- неразборчивая речь
|
||
def _handle_unclear(self, note: str) -> None:
|
||
"""Фраза не распознана: счётчик подряд; после UNCLEAR_LIMIT — куи и сброс."""
|
||
self._unclear_streak += 1
|
||
self._log(f"{note} (неразборчивых подряд: {self._unclear_streak})")
|
||
self._set_state(self.IDLE, note)
|
||
if self._unclear_streak >= UNCLEAR_LIMIT:
|
||
self._log("Много неразборчивых подряд — напоминаю про микрофон")
|
||
self._unclear_streak = 0
|
||
self._play_cue(UNCLEAR_CUE) |