Files
echo-core/src/adapters/discord_voice.py
Marius Mutu c1414616ad feat: US-010 - Extinde /voice doctor cu health-check pocket-tts
- adaugă ping GET /health la pocket-tts (URL din config tts.pockettts_url)
- adaugă verificare prezență hf_token în keyring
- checks-urile existente (libopus, voice load error) rămân neschimbate

gates rulate: tests PASS (1043 passed, 22 preexistente neschimbate), /review (backend) manual PASS
2026-07-11 11:16:06 +00:00

614 lines
24 KiB
Python

"""Discord voice slash commands (Pas 7 — CONVERGENCE wiring).
Registers the `/voice` slash command group on the existing CommandTree and
exposes an async `warmup_models()` for eager model load at bot startup.
Owns nothing in `src/voice/*` — purely the Discord-facing wiring. Defers
heavy lifting to:
- ``src.voice.pipeline.VoiceSession`` — per-guild session state machine
- ``src.voice.pipeline.EchoVoiceSink`` — discord-ext-voice-recv sink
- ``src.voice.tts_stream.TTSQueue`` / ``EchoStreamingAudioSource``
- ``src.voice._discord_voice_adapter.connect_voice``
"""
from __future__ import annotations
import asyncio
import io
import logging
import os
import re
import subprocess
import tempfile
import wave
from pathlib import Path
from typing import Optional
import discord
import httpx
from discord import app_commands
# Optional DAVE dep (mandatory at runtime when discord.py 2.7.1 is paired with
# Discord voice gateway v=8; tolerated missing in tests / dev environments).
try:
import davey
_HAS_DAVE = True
except ImportError:
_HAS_DAVE = False
from src.config import Config
from src.voice.pipeline import (
VoiceSession,
EchoVoiceSink,
_get_whisper_model,
_get_silero_vad,
)
from src.voice.tts_stream import TTSQueue, EchoStreamingAudioSource
from src.voice._discord_voice_adapter import connect_voice
log = logging.getLogger("echo-core.discord.voice")
PROJECT_ROOT = Path(__file__).resolve().parent.parent.parent
POCKET_TTS_VENV_PYTHON = PROJECT_ROOT / ".venv-pockettts" / "bin" / "python"
ADD_VOICE_SCRIPT = PROJECT_ROOT / "tools" / "pocket_tts_add_voice.py"
POCKET_TTS_SERVICE = "pocket-tts.service"
_MIN_ADDVOICE_SAMPLE_SECONDS = 3.0
_ADDVOICE_TIMEOUT_SECONDS = 300
# Per-guild voice session registry. Key = guild_id.
_voice_sessions: dict[int, VoiceSession] = {}
# Set if model warmup failed; surfaces as ephemeral error on /voice join.
_voice_load_error: Optional[str] = None
# Reference to the eager warmup task created in on_ready, so /voice join can
# await it if the user is faster than the background load.
_models_warmup_future: Optional[asyncio.Task] = None
async def warmup_models() -> None:
"""Eager model load — called from `on_ready()` as a background task.
Runs the (synchronous, blocking) model loaders on a worker thread so the
event loop stays responsive. On failure, sets `_voice_load_error` instead
of raising, so `/voice join` can degrade gracefully.
"""
global _voice_load_error
try:
if not discord.opus.is_loaded():
discord.opus.load_opus("libopus.so.0")
if _HAS_DAVE:
log.info("DAVE protocol v%d available (davey %s)",
davey.DAVE_PROTOCOL_VERSION, davey.__version__)
await asyncio.to_thread(_get_whisper_model)
await asyncio.to_thread(_get_silero_vad)
log.info("Voice models warm")
except Exception as e:
_voice_load_error = f"{type(e).__name__}: {e}"
log.error("Voice models load failed: %s", _voice_load_error)
def _get_whitelist() -> set[int]:
"""Read `voice.allowed_user_ids` from config and coerce to int set.
Re-reads config from disk to pick up any runtime edits between bot start
and /voice join.
"""
try:
raw = Config().get("voice.allowed_user_ids", [])
except Exception:
raw = []
out: set[int] = set()
for v in raw or []:
try:
out.add(int(v))
except (TypeError, ValueError):
continue
return out
def _get_default_voice() -> str:
try:
return Config().get("voice.default_voice", "M2") or "M2"
except Exception:
return "M2"
def _default_voice_for_engine(engine: str) -> str:
"""Vocea implicită asociată engine-ului (nu vocea implicită globală per-guild)."""
if engine == "pockettts":
return "alba"
return _get_default_voice()
def _systemctl_user(action: str, unit: str) -> None:
"""Best-effort `systemctl --user <action> <unit>` — nu ridică, doar loghează eșecul."""
try:
subprocess.run(
["systemctl", "--user", action, unit],
capture_output=True, text=True, timeout=15,
)
except Exception as e:
log.warning("systemctl --user %s %s failed: %s", action, unit, e)
_REGISTERED_VOICE_RE = re.compile(r"Registered voice '(.+?)' ->")
def _parse_registered_voice_name(stdout: str) -> Optional[str]:
m = _REGISTERED_VOICE_RE.search(stdout or "")
return m.group(1) if m else None
def _tts_voice_names() -> list[str]:
"""Nume de voci din tts_voices.json (import lazy din tools/tts.py, catalog live)."""
import sys as _sys
tools_dir = str(PROJECT_ROOT / "tools")
if tools_dir not in _sys.path:
_sys.path.insert(0, tools_dir)
try:
import tts as _tts_mod
return _tts_mod.list_voice_names()
except Exception:
return []
def _tts_synthesize_preview(text: str, voice: str) -> dict:
"""Import tools/tts.py (nu e package, sys.path trick) și generează un preview audio."""
import sys as _sys
tools_dir = str(PROJECT_ROOT / "tools")
if tools_dir not in _sys.path:
_sys.path.insert(0, tools_dir)
try:
import importlib
import tts as _tts_mod
importlib.reload(_tts_mod)
return _tts_mod.synthesize(text, voice=voice, lang="ro")
except Exception as e:
return {"ok": False, "error": f"{type(e).__name__}: {e}"}
def register(tree: app_commands.CommandTree, bot: discord.Client) -> app_commands.Group:
"""Build the `/voice` slash command group and return it (caller registers)."""
voice_group = app_commands.Group(
name="voice", description="Echo Core voice channel"
)
@voice_group.command(name="join", description="Echo intră în voice channel-ul tău")
async def join(interaction: discord.Interaction) -> None:
await interaction.response.defer(ephemeral=True)
if _voice_load_error:
await interaction.followup.send(
f"Voice unavailable: {_voice_load_error}", ephemeral=True
)
return
if _models_warmup_future is not None and not _models_warmup_future.done():
try:
await _models_warmup_future
except Exception as e:
await interaction.followup.send(
f"Voice unavailable: {type(e).__name__}: {e}", ephemeral=True
)
return
user = interaction.user
if not isinstance(user, discord.Member) or user.voice is None or user.voice.channel is None:
await interaction.followup.send(
"Intră într-un voice channel întâi.", ephemeral=True
)
return
channel = user.voice.channel
whitelist = _get_whitelist()
if user.id not in whitelist:
await interaction.followup.send(
"Nu ești pe whitelist voice.", ephemeral=True
)
return
# Reject double-join on the same guild.
guild_id = channel.guild.id
if guild_id in _voice_sessions:
await interaction.followup.send(
"Sunt deja în voice pe acest server. Folosește /voice leave întâi.",
ephemeral=True,
)
return
# Connect
try:
vc = await connect_voice(channel)
except Exception as e:
log.exception("connect_voice failed")
await interaction.followup.send(
f"Conectare eșuată: {type(e).__name__}: {e}", ephemeral=True
)
return
# Build TTS queue + session
ttsq = TTSQueue(voice_id=_get_default_voice(), lang="ro")
ttsq.start()
try:
session = VoiceSession(
text_channel_id=int(interaction.channel.id),
voice_channel_id=int(channel.id),
guild_id=guild_id,
voice_client=vc,
record_enabled=False,
mirror_enabled=True,
whitelist=whitelist,
ttsq=ttsq,
bot=bot,
loop=asyncio.get_running_loop(),
)
except Exception as e:
log.exception("VoiceSession construction failed")
ttsq.stop()
try:
await vc.disconnect(force=True)
except Exception:
pass
await interaction.followup.send(
f"Sesiune voice eșuată: {type(e).__name__}: {e}", ephemeral=True
)
return
_voice_sessions[guild_id] = session
# Start TTS streaming source for the entire session. Chain the
# wake-up beep via `after=` so streaming takes over when beep ends.
def _start_stream(error: Optional[Exception] = None) -> None:
if error is not None:
log.warning("Beep playback ended with error: %s", error)
try:
vc.play(EchoStreamingAudioSource(ttsq))
log.info("TTS streaming source attached")
except Exception:
log.exception("EchoStreamingAudioSource attach failed")
try:
vc.play(
discord.FFmpegPCMAudio("assets/voice/beep_200ms.wav"),
after=_start_stream,
)
except Exception:
log.warning("Beep playback skipped, starting stream directly", exc_info=True)
_start_stream()
# Attach sink
try:
bot_user_id = int(bot.user.id) if bot.user is not None else 0
sink = EchoVoiceSink(session=session, bot_user_id=bot_user_id)
vc.listen(sink)
except Exception as e:
log.exception("Sink attach failed")
_voice_sessions.pop(guild_id, None)
try:
session.cleanup("sink_attach_failed")
except Exception:
pass
await interaction.followup.send(
f"Atașare sink eșuată: {type(e).__name__}: {e}", ephemeral=True
)
return
# Presence
try:
await bot.change_presence(activity=discord.Activity(
type=discord.ActivityType.listening,
name=f"{user.display_name} în #{channel.name}",
))
except Exception:
log.warning("Presence update skipped", exc_info=True)
await interaction.followup.send(
f"În voce în #{channel.name}.", ephemeral=True
)
@voice_group.command(name="leave", description="Echo iese din voice channel")
async def leave(interaction: discord.Interaction) -> None:
await interaction.response.defer(ephemeral=True)
guild_id = interaction.guild.id if interaction.guild else None
session = _voice_sessions.pop(guild_id, None) if guild_id is not None else None
if session is None:
await interaction.followup.send(
"Nu sunt în niciun voice channel aici.", ephemeral=True
)
return
try:
session.cleanup("user_leave")
except Exception:
log.exception("session.cleanup raised")
try:
await bot.change_presence(activity=None)
except Exception:
log.warning("Presence reset skipped", exc_info=True)
await interaction.followup.send("Plecat.", ephemeral=True)
async def _voice_autocomplete(
interaction: discord.Interaction, current: str
) -> list[app_commands.Choice[str]]:
current_low = (current or "").lower()
names = [n for n in _tts_voice_names() if current_low in n.lower()]
return [app_commands.Choice(name=n, value=n) for n in names[:25]]
@voice_group.command(name="setvoice", description="Schimbă vocea Echo (din catalogul tts_voices.json)")
@app_commands.describe(voice="Voce nouă")
@app_commands.autocomplete(voice=_voice_autocomplete)
async def setvoice(
interaction: discord.Interaction,
voice: str,
) -> None:
await interaction.response.defer(ephemeral=True)
if voice not in _tts_voice_names():
await interaction.followup.send(
f"Voce necunoscută: {voice!r}. Alege din autocomplete.", ephemeral=True
)
return
new_voice = voice
# Live-swap on the active session if Echo is in voice on this guild.
guild_id = interaction.guild.id if interaction.guild else None
session = _voice_sessions.get(guild_id) if guild_id is not None else None
live_swapped = False
if session is not None and session.ttsq is not None:
session.ttsq.voice_id = new_voice
live_swapped = True
# Persist as the new default for future sessions.
try:
cfg = Config()
cfg.set("voice.default_voice", new_voice)
cfg.save()
except Exception as e:
log.warning("config save failed for new default voice: %s", e)
await interaction.followup.send(
f"Voce schimbată live ({new_voice}), dar config-ul nu s-a salvat: {e}",
ephemeral=True,
)
return
if live_swapped:
msg = f"Vocea schimbată **live** pe {new_voice}. Următoarea frază va folosi vocea nouă."
else:
msg = f"Default voce setată {new_voice}. Va intra în vigoare la următorul /voice join."
await interaction.followup.send(msg, ephemeral=True)
_ENGINE_CHOICES = [
app_commands.Choice(name="pocket-tts", value="pockettts"),
app_commands.Choice(name="Supertonic", value="supertonic"),
]
@voice_group.command(name="engine", description="Schimbă engine-ul TTS implicit (pockettts/supertonic)")
@app_commands.describe(engine="Engine TTS")
@app_commands.choices(engine=_ENGINE_CHOICES)
async def engine_cmd(
interaction: discord.Interaction,
engine: app_commands.Choice[str],
) -> None:
await interaction.response.defer(ephemeral=True)
new_engine = engine.value
try:
cfg = Config()
cfg.set("tts.default_engine", new_engine)
cfg.save()
except Exception as e:
log.warning("config save failed for tts.default_engine: %s", e)
await interaction.followup.send(
f"Eroare la salvarea engine-ului: {e}", ephemeral=True
)
return
default_voice = _default_voice_for_engine(new_engine)
await interaction.followup.send(
f"Engine TTS implicit setat pe **{new_engine}**. Voce implicită: {default_voice}.",
ephemeral=True,
)
@voice_group.command(name="addvoice", description="Adaugă o voce nouă clonată (pocket-tts) dintr-un sample WAV")
@app_commands.describe(nume="Nume voce (ex: Marius)", sample="Fișier WAV cu vocea (minim ~3s)")
async def addvoice(
interaction: discord.Interaction,
nume: str,
sample: discord.Attachment,
) -> None:
await interaction.response.defer(ephemeral=True)
nume = nume.strip()
if not nume:
await interaction.followup.send("Numele vocii nu poate fi gol.", ephemeral=True)
return
filename = sample.filename or ""
if not filename.lower().endswith(".wav"):
await interaction.followup.send(
"Sample-ul trebuie să fie un fișier .wav.", ephemeral=True
)
return
try:
content = await sample.read()
except Exception as e:
await interaction.followup.send(
f"Descărcare sample eșuată: {type(e).__name__}: {e}", ephemeral=True
)
return
try:
with wave.open(io.BytesIO(content), "rb") as wf:
duration = wf.getnframes() / float(wf.getframerate())
except (wave.Error, EOFError) as e:
await interaction.followup.send(f"Fișier WAV invalid: {e}", ephemeral=True)
return
if duration < _MIN_ADDVOICE_SAMPLE_SECONDS:
await interaction.followup.send(
f"Sample prea scurt ({duration:.1f}s) — minim {_MIN_ADDVOICE_SAMPLE_SECONDS:.0f}s.",
ephemeral=True,
)
return
if not POCKET_TTS_VENV_PYTHON.exists():
await interaction.followup.send(
f"venv pocket-tts lipsă: {POCKET_TTS_VENV_PYTHON}", ephemeral=True
)
return
await interaction.followup.send(
"Adaug voce, TTS indisponibil ~30s...", ephemeral=True
)
tmp_wav_path: Optional[Path] = None
try:
fd, tmp_name = tempfile.mkstemp(prefix="echo-addvoice-", suffix=".wav")
with open(fd, "wb") as f:
f.write(content)
tmp_wav_path = Path(tmp_name)
await asyncio.to_thread(_systemctl_user, "stop", POCKET_TTS_SERVICE)
try:
proc = await asyncio.to_thread(
subprocess.run,
[str(POCKET_TTS_VENV_PYTHON), str(ADD_VOICE_SCRIPT),
"--wav", str(tmp_wav_path), "--name", nume],
capture_output=True, text=True,
timeout=_ADDVOICE_TIMEOUT_SECONDS, cwd=str(PROJECT_ROOT),
)
except subprocess.TimeoutExpired:
await interaction.followup.send(
f"Export voce a depășit timeout-ul ({_ADDVOICE_TIMEOUT_SECONDS}s).",
ephemeral=True,
)
return
finally:
await asyncio.to_thread(_systemctl_user, "start", POCKET_TTS_SERVICE)
finally:
if tmp_wav_path is not None:
try:
tmp_wav_path.unlink(missing_ok=True)
except OSError:
pass
if proc.returncode != 0:
err = (proc.stderr or proc.stdout or "eroare necunoscută").strip()
await interaction.followup.send(
f"Export voce eșuat: {err[-500:]}", ephemeral=True
)
return
final_name = _parse_registered_voice_name(proc.stdout) or nume
preview = _tts_synthesize_preview(f"Salut, sunt vocea {final_name}.", final_name)
if preview.get("ok"):
preview_path = preview["path"]
try:
await interaction.followup.send(
f"Voce adăugată: **{final_name}**.",
file=discord.File(preview_path, filename="preview.wav"),
ephemeral=True,
)
finally:
try:
os.unlink(preview_path)
except OSError:
pass
else:
await interaction.followup.send(
f"Voce adăugată: **{final_name}** (preview audio eșuat: {preview.get('error')})",
ephemeral=True,
)
@voice_group.command(name="stop", description="Oprește audio-ul curent (golește coada TTS)")
async def stop_audio(interaction: discord.Interaction) -> None:
await interaction.response.defer(ephemeral=True)
guild_id = interaction.guild.id if interaction.guild else None
session = _voice_sessions.get(guild_id) if guild_id is not None else None
if session is None or session.ttsq is None:
await interaction.followup.send("Nu sunt în voice.", ephemeral=True)
return
try:
session.ttsq.clear()
log.info("voice stop: TTS queue cleared by user %s", interaction.user)
except Exception as e:
log.warning("voice stop: ttsq.clear failed: %s", e)
await interaction.followup.send(f"Eroare la oprire: {e}", ephemeral=True)
return
await interaction.followup.send("Audio oprit.", ephemeral=True)
@voice_group.command(name="doctor", description="Verifică voice stack")
async def doctor(interaction: discord.Interaction) -> None:
await interaction.response.defer(ephemeral=True)
checks: list[tuple[str, bool]] = []
# libopus
try:
checks.append(("libopus", bool(discord.opus.is_loaded())))
except Exception:
checks.append(("libopus", False))
# warmup
checks.append(("voice load error", _voice_load_error is None))
# hf_token în keyring
try:
from src.credential_store import get_secret
checks.append(("hf_token (keyring)", get_secret("hf_token") is not None))
except Exception:
checks.append(("hf_token (keyring)", False))
# pocket-tts /health
pockettts_ok = False
try:
pockettts_url = Config().get("tts.pockettts_url", "http://127.0.0.1:7789")
async with httpx.AsyncClient(timeout=3.0) as client:
resp = await client.get(f"{pockettts_url}/health")
pockettts_ok = resp.status_code == 200
except Exception:
pockettts_ok = False
checks.append(("pocket-tts /health", pockettts_ok))
# Build response
lines = ["**Voice doctor:**"]
for label, ok in checks:
lines.append(f"{'OK' if ok else 'FAIL'}{label}")
if _voice_load_error:
lines.append(f" details: {_voice_load_error}")
await interaction.followup.send("\n".join(lines), ephemeral=True)
# --- /voice mirror on|off ---
mirror_group = app_commands.Group(
name="mirror", description="Text mirror", parent=voice_group
)
@mirror_group.command(name="on", description="Activează text mirror în canal")
async def mirror_on(interaction: discord.Interaction) -> None:
await interaction.response.defer(ephemeral=True)
guild_id = interaction.guild.id if interaction.guild else None
s = _voice_sessions.get(guild_id) if guild_id is not None else None
if s is None:
await interaction.followup.send("Nu sunt în voice.", ephemeral=True)
return
s.mirror_enabled = True
await interaction.followup.send("Mirror ON.", ephemeral=True)
@mirror_group.command(name="off", description="Dezactivează text mirror")
async def mirror_off(interaction: discord.Interaction) -> None:
await interaction.response.defer(ephemeral=True)
guild_id = interaction.guild.id if interaction.guild else None
s = _voice_sessions.get(guild_id) if guild_id is not None else None
if s is None:
await interaction.followup.send("Nu sunt în voice.", ephemeral=True)
return
s.mirror_enabled = False
await interaction.followup.send("Mirror OFF.", ephemeral=True)
# --- /voice record on|off ---
record_group = app_commands.Group(
name="record", description="KB recording", parent=voice_group
)
@record_group.command(name="on", description="Activează înregistrare în KB")
async def record_on(interaction: discord.Interaction) -> None:
await interaction.response.defer(ephemeral=True)
guild_id = interaction.guild.id if interaction.guild else None
s = _voice_sessions.get(guild_id) if guild_id is not None else None
if s is None:
await interaction.followup.send("Nu sunt în voice.", ephemeral=True)
return
s.record_enabled = True
await interaction.followup.send("Record ON.", ephemeral=True)
@record_group.command(name="off", description="Dezactivează înregistrare")
async def record_off(interaction: discord.Interaction) -> None:
await interaction.response.defer(ephemeral=True)
guild_id = interaction.guild.id if interaction.guild else None
s = _voice_sessions.get(guild_id) if guild_id is not None else None
if s is None:
await interaction.followup.send("Nu sunt în voice.", ephemeral=True)
return
s.record_enabled = False
await interaction.followup.send("Record OFF.", ephemeral=True)
return voice_group