refactor: move webui out of cog, use shared-state pending queue

- drop bundled Flask webui (now lives in scrapyard-website TtsToy)
- _take_pending_state_updates: pop pending commands/posts + refresh user
  info under cross-process lock, matching the webui-side producer
This commit is contained in:
owen
2026-09-14 11:45:47 -05:00
parent 9b992d6dbe
commit b40e7f6cad
15 changed files with 137 additions and 2865 deletions
+137 -267
View File
@@ -12,6 +12,11 @@ from pathlib import Path
import lavalink
from typing import Optional
try:
import fcntl as _fcntl
except ImportError: # pragma: no cover - Windows
_fcntl = None
import discord
from discord.ext import tasks
import requests
@@ -384,7 +389,6 @@ class TtsToy(Cog):
self._health_server: HTTPServer | None = None
self._health_thread: threading.Thread | None = None
self._webui_process = None # subprocess for the Flask web UI
self.config: Config = Config.get_conf(
self, identifier=0x0A0A0A0A, force_registration=True
@@ -437,7 +441,7 @@ class TtsToy(Cog):
self.clear_old_tts.stop()
if self._webui_state_processor.is_running():
self._webui_state_processor.stop()
await self._stop_webui()
self._kill_leftover_webui()
await self._stop_dectalk_server()
await self._stop_morshu_server()
await self._stop_vox_server()
@@ -455,6 +459,20 @@ class TtsToy(Cog):
if self.tts_processors:
await asyncio.gather(*self.tts_processors.values(), return_exceptions=True)
def _kill_leftover_webui(self):
"""Terminate any web UI subprocess left over from before the web UI
moved out of the bot. It ran from <cog_dir>/webui/app.py, which no
longer exists — so a matching process is guaranteed stale."""
try:
pattern = str(Path(__file__).resolve().parent / "webui" / "app.py")
result = subprocess.run(
["pkill", "-f", pattern], capture_output=True, text=True
)
if result.returncode == 0:
log.info("[WebUI] Killed leftover web UI subprocess")
except Exception as e:
log.warning(f"[WebUI] Could not kill leftover web UI subprocess: {e}")
# -----------------------------------------
async def _get_minimax_settings(self):
return (
@@ -3088,40 +3106,31 @@ class TtsToy(Cog):
@ttstoy_group.command(name="webuistatus")
@commands.is_owner()
async def webui_status_cmd(self, ctx: commands.Context):
"""Check the status of the TtsToy web UI."""
"""Check the status of the TtsToy web UI (now a scrapyard service)."""
webui_url = self._get_webui_url()
process_running = self._webui_process and self._webui_process.poll() is None
pid = self._webui_process.pid if process_running else None
# Try to hit the login page
reachable = False
try:
r = await asyncio.to_thread(requests.get, f"{webui_url}/login", timeout=3)
r = await asyncio.to_thread(requests.get, f"{webui_url.rstrip('/')}/login", timeout=3)
reachable = r.status_code == 200
except Exception:
pass
lines = [
"**TtsToy Web UI Status**",
f"Process running: {'✅ Yes' if process_running else '❌ No'}{f' (pid {pid})' if pid else ''}",
f"Reachable at URL: {'✅ Yes' if reachable else '❌ No'}",
f"URL: {webui_url}",
"The web UI now runs as a scrapyard service (not a bot subprocess).",
f"Restart with: `sudo systemctl restart scrapyard`",
]
if not process_running:
lines.append(f"\nRestart with: `{ctx.clean_prefix}ttstoy webuirestart`")
await ctx.send("\n".join(lines))
@ttstoy_group.command(name="webuirestart")
@commands.is_owner()
async def webui_restart_cmd(self, ctx: commands.Context):
"""Restart the TtsToy web UI."""
await ctx.send("🔄 Restarting web UI...")
await self._stop_webui()
success = await self._start_webui()
if success:
await ctx.send(f"✅ Web UI restarted at {self._get_webui_url()}")
else:
await ctx.send("❌ Failed to start web UI. Check bot logs.")
"""Restart the TtsToy web UI (now a scrapyard service)."""
await ctx.send("🔄 The web UI is now a scrapyard service.\nRestart it with: `sudo systemctl restart scrapyard`")
# -------------------------------------------------------------------------
# Background task — poll shared state for pending commands and posts
@@ -3131,38 +3140,13 @@ class TtsToy(Cog):
async def _webui_state_processor(self):
"""Poll the shared state file and execute pending commands/posts from the web UI."""
try:
state = self._read_shared_state()
changed = False
pending_cmds, pending_posts = await asyncio.to_thread(self._take_pending_state_updates)
# --- Process pending commands (voice/mode changes from web UI) ---
pending_cmds = state.pop("pending_commands", [])
if pending_cmds:
changed = True
for cmd in pending_cmds:
await self._handle_webui_command(cmd)
for cmd in pending_cmds:
await self._handle_webui_command(cmd)
# --- Process pending posts (TTS audio to post to Discord channel) ---
pending_posts = state.pop("pending_posts", [])
if pending_posts:
changed = True
for post in pending_posts:
await self._handle_webui_post(post)
# --- Write updated user info so the profile page stays fresh ---
users_info = state.get("users", {})
for guild in self.bot.guilds:
for member in guild.members:
uid = str(member.id)
if uid not in users_info:
users_info[uid] = {}
users_info[uid]["display_name"] = member.display_name
users_info[uid]["username"] = str(member)
users_info[uid]["joined_at"] = str(member.joined_at)[:10] if member.joined_at else "unknown"
users_info[uid]["roles"] = ", ".join(r.name for r in member.roles if r.name != "@everyone")
state["users"] = users_info
if changed or True: # always write to keep user info fresh
self._write_shared_state(state)
for post in pending_posts:
await self._handle_webui_post(post)
# Reload emoji→SFX mapping in case it was changed via the web UI
_reload_emoji_sfx_map()
@@ -3170,6 +3154,40 @@ class TtsToy(Cog):
except Exception as e:
log.exception("[WebUI] State processor error: %s", e)
def _take_pending_state_updates(self):
"""Under a cross-process lock: pop pending commands/posts and refresh user info.
Returns (pending_cmds, pending_posts). Runs in a worker thread — the
lock matches the one the web UI uses when appending, so no updates are
lost between the read and the write.
"""
lock_file = Path(os.environ.get("TTSTOY_SHARED_STATE_LOCK", "/tmp/ttstoy_webui_state.lock"))
with open(lock_file, "w") as lf:
if _fcntl is not None:
_fcntl.flock(lf, _fcntl.LOCK_EX)
try:
state = self._read_shared_state()
pending_cmds = state.pop("pending_commands", [])
pending_posts = state.pop("pending_posts", [])
# Refresh user info so the profile page stays current
users_info = state.get("users", {})
for guild in self.bot.guilds:
for member in guild.members:
uid = str(member.id)
info = users_info.setdefault(uid, {})
info["display_name"] = member.display_name
info["username"] = str(member)
info["joined_at"] = str(member.joined_at)[:10] if member.joined_at else "unknown"
info["roles"] = ", ".join(r.name for r in member.roles if r.name != "@everyone")
state["users"] = users_info
self._write_shared_state(state)
finally:
if _fcntl is not None:
_fcntl.flock(lf, _fcntl.LOCK_UN)
return pending_cmds, pending_posts
async def _handle_webui_command(self, cmd: dict):
"""Execute a command queued by the web UI."""
user_id = int(cmd.get("user_id", 0))
@@ -3182,10 +3200,7 @@ class TtsToy(Cog):
return
try:
if action == "speak_in_vc":
await self._webui_speak_in_vc(user, payload)
elif action == "set_voice":
if action == "set_voice":
voice = payload.get("voice", "")
await self.config.user(user).minimax_voice.set(voice)
log.info(f"[WebUI] Set voice for {user} → {voice}")
@@ -3245,159 +3260,40 @@ class TtsToy(Cog):
except Exception as e:
log.exception(f"[WebUI] Error handling command {action} for user {user_id}: {e}")
async def _webui_speak_in_vc(self, user: discord.User, payload: dict):
"""Generate TTS from a web UI request and play it in the user's current VC."""
job_id = payload.get("job_id")
text = payload.get("text", "")
discord_name = payload.get("user", str(user))
mode_override = payload.get("mode_override") # per-request mode from web UI
voice_override = payload.get("voice_override") # per-request voice from web UI
skip_vc = payload.get("skip_vc", False)
skip_post = payload.get("skip_post", False)
guild_id_str = payload.get("guild_id", "")
log.info(f"[WebUI] speak_in_vc: user={user} job={job_id} guild={guild_id_str} mode_override={mode_override} voice_override={voice_override} skip_vc={skip_vc} skip_post={skip_post} text={text[:60]!r}")
webui_url = self._get_webui_internal_url()
secret = self._get_internal_secret()
def _update_job(status, **kwargs):
try:
requests.post(
f"{webui_url.rstrip('/')}/api/tts/{job_id}/update",
headers={"X-Internal-Secret": secret, "Content-Type": "application/json"},
json={"status": status, **kwargs},
timeout=5,
)
log.debug(f"[WebUI] Job {job_id} updated → {status}")
except Exception as e:
log.warning(f"[WebUI] Could not update job {job_id}: {e}")
# Find which guild/VC the user is currently in
# Prefer the guild they logged in from
member_vc = None
member_guild = None
target_guild_id = int(guild_id_str) if guild_id_str else None
if target_guild_id:
guild = self.bot.get_guild(target_guild_id)
if guild:
member_guild = guild
member = guild.get_member(user.id)
if member and member.voice and member.voice.channel:
member_vc = member.voice.channel
# Fallback: search all guilds for VC presence
if not member_vc:
for guild in self.bot.guilds:
member = guild.get_member(user.id)
if member and member.voice and member.voice.channel:
member_vc = member.voice.channel
if not member_guild:
member_guild = guild
break
if not member_vc:
log.info(f"[WebUI] {user} is not in a VC — will generate audio only")
log.info(f"[WebUI] Target VC: {member_vc} in {member_guild}")
# Gather TTS settings — respect web UI overrides
mode = mode_override if mode_override else await self._get_tts_mode()
chatterbox_url = await self._get_chatterbox_api_url()
dectalk_url = await self._get_dectalk_api_url()
api_key, model, _ = await self._get_minimax_settings()
# Voice: use override if provided, otherwise user's saved voice
if voice_override:
voice_id = voice_override
else:
voice_id = await self._get_effective_minimax_voice(user)
sfx_volume = await self.config.sfx_volume()
vox_pack = await self.config.vox_pack()
cb_temp = cb_exag = cb_speed = None
cb_vol_db = 0.0
if mode == "chatterbox":
cb_temp = (await self.config.user(user).chatterbox_temperature_offsets()).get(voice_id)
cb_exag = (await self.config.user(user).chatterbox_exaggeration_offsets()).get(voice_id)
cb_vol_db = (await self.config.user(user).chatterbox_volume_offsets()).get(voice_id, 0.0)
cb_speed = (await self.config.user(user).chatterbox_speed_offsets()).get(voice_id)
log.debug(f"[WebUI] TTS params: mode={mode} voice={voice_id} temp={cb_temp} exag={cb_exag}")
self.tts_storage.mkdir(parents=True, exist_ok=True)
audio_path = str(self.tts_storage / f"webui_{job_id}.mp3")
all_cb_voices = await self._get_all_chatterbox_voices()
try:
await asyncio.to_thread(
self._save_prompt_audio,
api_key, model, voice_id, text, audio_path,
sfx_volume, mode, dectalk_url, vox_pack,
chatterbox_url, cb_temp, cb_exag, cb_vol_db, cb_speed,
all_cb_voices,
)
log.info(f"[WebUI] Audio generated: {audio_path}")
except Exception as e:
log.exception(f"[WebUI] TTS generation failed for job {job_id}: {e}")
await asyncio.to_thread(_update_job, "error", error=str(e))
return
# Queue for VC playback only if user is in a voice channel and skip_vc is not set
if member_vc and member_guild and not skip_vc:
try:
audio_cog: Optional[Audio] = self.bot.get_cog("Audio")
if audio_cog:
text_channel = member_guild.system_channel or next(
(c for c in member_guild.text_channels
if c.permissions_for(member_guild.me).send_messages), None
)
if text_channel:
if not member_guild.me.voice or member_guild.me.voice.channel != member_vc:
await member_vc.connect()
class _FakeCtx:
guild = member_guild
channel = text_channel
author = member_guild.get_member(user.id) or user
interaction = None
async def react_quietly(self, *a, **kw): pass
queue = self._get_or_create_tts_queue(member_guild.id)
await queue.put((_FakeCtx(), audio_path))
log.info(f"[WebUI] Queued audio for VC playback: job={job_id}")
except Exception as e:
log.exception(f"[WebUI] VC playback failed for job {job_id}: {e}")
# Post transcript to Discord channel and mark job done
# Resolve the best guild ID for posting: login guild > VC guild > empty
post_guild_id = guild_id_str or (str(member_guild.id) if member_guild else "")
if not skip_post:
await self._handle_webui_post({
"job_id": job_id, "path": audio_path,
"text": text, "user": discord_name, "engine": mode,
"guild_id": post_guild_id,
})
await asyncio.to_thread(_update_job, "done", path=audio_path)
async def _handle_webui_post(self, post: dict):
"""Post a TTS audio file to the configured bot channel."""
"""Deliver a pre-generated TTS file queued by the web UI.
If `play_vc` is set and the user is in a voice channel, the audio is
queued for VC playback (via the standard ttstoy queue). Either way it is
also posted as a message to the configured TTS channel.
"""
try:
audio_path = post.get("path")
text = post.get("text", "")
user = post.get("user", "Unknown")
engine = post.get("engine", "tts")
guild_id_str = post.get("guild_id", "")
user_id_str = post.get("user_id", "")
play_vc = bool(post.get("play_vc", False))
if not audio_path or not os.path.exists(audio_path):
log.warning(f"[WebUI] Audio file not found: {audio_path}")
return
# Find the designated TTS channel in the user's login guild
target_guild_id = int(guild_id_str) if guild_id_str else None
# --- VC playback (optional) ---
if play_vc and user_id_str:
await self._play_webui_audio_in_vc(user_id_str, audio_path, target_guild_id)
# --- Post transcript to the configured TTS channel ---
channel = await self._get_webui_tts_channel(target_guild_id)
if not channel:
log.warning("[WebUI] No TTS channel configured. Set TTSTOY_WEBUI_CHANNEL_ID env var.")
return
filename = f"webui-{engine}-tts.mp3"
suffix = Path(audio_path).suffix or ".mp3"
filename = f"webui-{engine}-tts{suffix}"
content = f"🗣 **{discord.utils.escape_markdown(user)}** via web UI [{engine}]\n> {discord.utils.escape_markdown(text[:200])}"
await channel.send(content=content, file=discord.File(audio_path, filename=filename))
log.info(f"[WebUI] Posted TTS from {user} to channel {channel.id}")
@@ -3405,6 +3301,61 @@ class TtsToy(Cog):
except Exception as e:
log.exception(f"[WebUI] Error posting TTS: {e}")
async def _play_webui_audio_in_vc(self, user_id_str: str, audio_path: str, target_guild_id: int = None):
"""Queue a pre-generated audio file for VC playback in the user's current VC."""
try:
user_id = int(user_id_str)
member_vc = None
member_guild = None
if target_guild_id:
guild = self.bot.get_guild(target_guild_id)
if guild:
member_guild = guild
member = guild.get_member(user_id)
if member and member.voice and member.voice.channel:
member_vc = member.voice.channel
if not member_vc:
for guild in self.bot.guilds:
member = guild.get_member(user_id)
if member and member.voice and member.voice.channel:
member_vc = member.voice.channel
if not member_guild:
member_guild = guild
break
if not member_vc or not member_guild:
log.info(f"[WebUI] {user_id} is not in a VC — skipping playback")
return
member = member_guild.get_member(user_id)
audio_cog: Optional[Audio] = self.bot.get_cog("Audio")
if not audio_cog:
log.warning("[WebUI] Audio cog not loaded — skipping VC playback")
return
text_channel = member_guild.system_channel or next(
(c for c in member_guild.text_channels
if c.permissions_for(member_guild.me).send_messages), None
)
if not text_channel:
log.warning("[WebUI] No text channel available for VC playback")
return
if not member_guild.me.voice or member_guild.me.voice.channel != member_vc:
await member_vc.connect()
class _FakeCtx:
guild = member_guild
channel = text_channel
author = member or user_id
interaction = None
async def react_quietly(self, *a, **kw): pass
queue = self._get_or_create_tts_queue(member_guild.id)
await queue.put((_FakeCtx(), audio_path))
log.info(f"[WebUI] Queued audio for VC playback (user={user_id} guild={member_guild.id})")
except Exception as e:
log.exception(f"[WebUI] VC playback failed: {e}")
async def _get_webui_tts_channel(self, guild_id: int = None) -> Optional[discord.TextChannel]:
"""Get the channel to post web UI TTS messages to.
@@ -3468,84 +3419,6 @@ class TtsToy(Cog):
await self.config.guild(ctx.guild).webui_tts_channel_id.set(channel.id)
await ctx.send(f"✅ Web UI TTS audio will be posted to {channel.mention} in this server")
# -------------------------------------------------------------------------
# Web UI subprocess management
# -------------------------------------------------------------------------
def _webui_app_path(self) -> Path:
return Path(__file__).resolve().parent / "webui" / "app.py"
async def _start_webui(self):
"""Launch the Flask web UI as a subprocess."""
if self._webui_process and self._webui_process.poll() is None:
log.info("[WebUI] Already running")
return True
app_path = self._webui_app_path()
if not app_path.exists():
log.error(f"[WebUI] app.py not found at {app_path}")
return False
env = os.environ.copy()
env.setdefault("WEBUI_PORT", "8098")
env.setdefault("WEBUI_URL", f"http://localhost:{env['WEBUI_PORT']}")
env.setdefault("TTSTOY_SHARED_STATE", str(self.WEBUI_SHARED_STATE))
if "BOT_INSTANCE" not in env:
try:
from redbot.core.data_manager import instance_name
env["BOT_INSTANCE"] = instance_name()
except Exception:
env.setdefault("BOT_INSTANCE", "redbot")
try:
import sys
python = sys.executable
req_file = app_path.parent / "requirements.txt"
if req_file.exists():
subprocess.run(
[python, "-m", "pip", "install", "-q", "-r", str(req_file)],
check=True, capture_output=True,
)
self._webui_process = subprocess.Popen(
[python, str(app_path)],
env=env,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
)
await asyncio.sleep(2)
if self._webui_process.poll() is not None:
out = self._webui_process.stdout.read(500) if self._webui_process.stdout else ""
log.error(f"[WebUI] Process exited immediately: {out}")
return False
log.info(f"[WebUI] Started on port {env['WEBUI_PORT']} (pid {self._webui_process.pid})")
# Pipe subprocess stdout → RedBot logger in a background thread
def _pipe_logs(proc):
for line in proc.stdout:
line = line.rstrip()
if line:
log.info(f"[WebUI] {line}")
threading.Thread(target=_pipe_logs, args=(self._webui_process,), daemon=True).start()
return True
except Exception as e:
log.exception(f"[WebUI] Failed to start: {e}")
return False
async def _stop_webui(self):
"""Stop the Flask web UI subprocess."""
if self._webui_process and self._webui_process.poll() is None:
try:
self._webui_process.terminate()
await asyncio.sleep(1)
if self._webui_process.poll() is None:
self._webui_process.kill()
log.info("[WebUI] Stopped")
except Exception as e:
log.exception(f"[WebUI] Error stopping: {e}")
self._webui_process = None
# -------------------------------------------------------------------------
# Cog lifecycle
# -------------------------------------------------------------------------
@@ -3576,9 +3449,6 @@ class TtsToy(Cog):
except Exception:
log.exception("Failed to start TtsToy health server")
# Start web UI subprocess
await self._start_webui()
# Start web UI state processor
self._webui_state_processor.start()
log.info("[WebUI] State processor started")