first commit
This commit is contained in:
@@ -0,0 +1,422 @@
|
||||
"""Машина состояний голосового ассистента (сердце Модуля 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)
|
||||
Reference in New Issue
Block a user