diff --git a/dave/bot.mjs b/dave/bot.mjs index 5c87c43..ad9c3d8 100644 --- a/dave/bot.mjs +++ b/dave/bot.mjs @@ -55,6 +55,10 @@ const RUN_MS = process.env.RUN_MS != null ? Number(process.env.RUN_MS) : 0; // 0 // Python STT+TTS voice-turn endpoint (run `python -m wsai --voice-server`). const VOICE_ENDPOINT = process.env.WSAI_VOICE_ENDPOINT || env.WSAI_VOICE_ENDPOINT || 'http://127.0.0.1:8787/api/voice-turn'; +// Dashboard control plane (state push + command poll), same host as voice-turn. +const API_BASE = VOICE_ENDPOINT.replace(/\/api\/voice-turn\/?$/, ''); +const REPORT_ENDPOINT = API_BASE + '/api/bot/report'; +const REPORT_INTERVAL_MS = Number(process.env.WSAI_REPORT_INTERVAL_MS || 2500); // Ignore utterances shorter than this many PCM bytes (48kHz*2ch*2B = 192000 B/s), // so key clicks / brief noise don't trigger a turn. ~0.35s. const MIN_UTTERANCE_BYTES = Number(process.env.WSAI_MIN_UTTERANCE_BYTES || 67000); @@ -109,10 +113,19 @@ async function handleUtterance(userId, pcm) { return; } const wav = Buffer.concat([wavHeader(pcm.length), pcm]); + // Resolve who spoke (Discord display name) so the dashboard log can show it. + let speaker = userId; + try { + const g = currentGuildId && client.guilds.cache.get(currentGuildId); + const m = g && (g.members.cache.get(userId) || await g.members.fetch(userId).catch(() => null)); + if (m) speaker = m.displayName || m.user.username; + } catch {} let resp; try { resp = await fetch(VOICE_ENDPOINT, { - method: 'POST', headers: { 'Content-Type': 'audio/wav' }, body: wav, + method: 'POST', + headers: { 'Content-Type': 'audio/wav', 'X-User-Name': encodeURIComponent(speaker) }, + body: wav, }); } catch (e) { log(`voice-turn POST failed (is \`python -m wsai --voice-server\` running?): ${e.message}`); @@ -147,11 +160,131 @@ const client = new Client({ const perUser = new Map(); // userId -> { opusPackets, pcmFrames } +// --- control-plane state (dashboard drives which channel we're in) --------- # +let currentGuildId = null, currentChannelId = null, currentChannelName = null; +const speakingSet = new Set(); // userIds currently speaking (for participant list) +const activeSubs = new Set(); // userIds with an in-flight receive subscription + +// Attach the bot's audio player (so it can speak) to a fresh connection. +function setupPlayer(connection) { + voicePlayer = createAudioPlayer({ behaviors: { noSubscriber: NoSubscriberBehavior.Play } }); + voicePlayer.on('error', (e) => log(`player error: ${e.message}`)); + connection.subscribe(voicePlayer); +} + +// Attach the receive path: capture each utterance and run the voice turn. +function setupReceiver(connection) { + const receiver = connection.receiver; + receiver.speaking.on('start', (userId) => { + speakingSet.add(userId); + if (userId === client.user.id || activeSubs.has(userId)) return; + activeSubs.add(userId); + if (!perUser.has(userId)) perUser.set(userId, { opusPackets: 0, pcmFrames: 0 }); + const opusStream = receiver.subscribe(userId, { + end: { behavior: EndBehaviorType.AfterSilence, duration: 800 }, + }); + const decoder = new prism.opus.Decoder({ rate: 48000, channels: 2, frameSize: 960 }); + const chunks = []; + opusStream.on('data', () => { perUser.get(userId).opusPackets++; }); + opusStream.on('error', (e) => { logThrottled(`recv:${e.message}`, `recv stream error user=${userId}: ${e.message}`); activeSubs.delete(userId); }); + opusStream.pipe(decoder); + decoder.on('data', (d) => { chunks.push(d); perUser.get(userId).pcmFrames++; }); + decoder.on('error', (e) => log(`decode error user=${userId}: ${e.message}`)); + decoder.on('end', () => { + activeSubs.delete(userId); + const pcm = Buffer.concat(chunks); + log(`utterance end user=${userId} pcm=${pcm.length}B — running voice turn`); + handleUtterance(userId, pcm).catch((e) => log(`voice turn error: ${e.message}`)); + }); + }); + receiver.speaking.on('end', (userId) => speakingSet.delete(userId)); +} + +// Join (or switch to) a voice channel on command from the dashboard. +async function joinChannel(guildId, channelId) { + const guild = await client.guilds.fetch(guildId).catch(() => null); + const channel = guild && await guild.channels.fetch(channelId).catch(() => null); + if (!channel || !channel.isVoiceBased()) { log(`join failed: ${guildId}/${channelId} not a voice channel`); return; } + if (currentGuildId && currentGuildId !== guildId) leaveChannel(); + const connection = joinVoiceChannel({ + channelId, guildId, adapterCreator: guild.voiceAdapterCreator, + selfDeaf: false, selfMute: false, + }); + connection.on('error', (e) => log(`voice connection error: ${e.message}`)); + try { + await entersState(connection, VoiceConnectionStatus.Ready, 40_000); + } catch (e) { log(`join not Ready in 40s: ${e.message}`); return; } + currentGuildId = guildId; currentChannelId = channelId; currentChannelName = channel.name; + speakingSet.clear(); activeSubs.clear(); + setupPlayer(connection); + setupReceiver(connection); + connection.on(VoiceConnectionStatus.Disconnected, () => { + log('voice: disconnected — attempting to resume…'); + Promise.race([ + entersState(connection, VoiceConnectionStatus.Signalling, 5_000), + entersState(connection, VoiceConnectionStatus.Connecting, 5_000), + ]).catch(() => { log('voice: could not resume'); leaveChannel(); }); + }); + log(`✅ joined voice: "${guild.name}" / "${channel.name}"`); +} + +function leaveChannel() { + try { getVoiceConnection(currentGuildId)?.destroy(); } catch {} + currentGuildId = currentChannelId = currentChannelName = null; + speakingSet.clear(); activeSubs.clear(); +} + +// Build the state snapshot the dashboard shows (identity, joinable servers + +// voice channels, current channel, and who is in it). +function buildState() { + const guilds = [...client.guilds.cache.values()].map((g) => ({ + id: g.id, name: g.name, + voiceChannels: [...g.channels.cache.values()] + .filter((c) => c.isVoiceBased()) + .map((c) => ({ id: c.id, name: c.name })), + })); + let members = []; + if (currentGuildId && currentChannelId) { + const ch = client.guilds.cache.get(currentGuildId)?.channels.cache.get(currentChannelId); + if (ch && ch.members) { + members = [...ch.members.values()].map((m) => ({ + id: m.id, name: m.displayName, speaking: speakingSet.has(m.id), + })); + } + } + return { + identity: { id: client.user.id, username: client.user.username, tag: client.user.tag }, + guilds, + current: { guildId: currentGuildId, channelId: currentChannelId, channelName: currentChannelName }, + members, + }; +} + +async function handleCommand(cmd) { + if (cmd.type === 'join') await joinChannel(cmd.guildId, cmd.channelId); + else if (cmd.type === 'leave') { leaveChannel(); log('left voice on command'); } +} + +// Push state to the dashboard and apply any commands it hands back. +async function reportLoop() { + let j; + try { + const r = await fetch(REPORT_ENDPOINT, { + method: 'POST', headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify(buildState()), + }); + j = await r.json(); + } catch { return; } // dashboard down: keep running, retry next tick + for (const cmd of (j?.commands || [])) { + try { await handleCommand(cmd); } catch (e) { log(`command ${cmd?.type} failed: ${e.message}`); } + } +} + let leaving = false; function leaveAndExit(code = 0) { if (leaving) return; leaving = true; - try { getVoiceConnection(GUILD_ID)?.destroy(); } catch {} + try { getVoiceConnection(currentGuildId || GUILD_ID)?.destroy(); } catch {} const summary = [...perUser.entries()].map(([u, s]) => `${u}:opus=${s.opusPackets},pcm=${s.pcmFrames}`); log(`leaving. speakers heard: ${summary.length ? summary.join(' ') : '(none)'}`); try { client.destroy(); } catch {} @@ -165,84 +298,21 @@ if (RUN_MS > 0) setTimeout(() => { log(`RUN_MS=${RUN_MS} hard ceiling elapsed client.once('clientReady', async () => { log(`logged in as ${client.user.tag} (${client.user.id})`); - let guild, channel; - try { - guild = await client.guilds.fetch(GUILD_ID); - channel = await guild.channels.fetch(CHANNEL_ID); - } catch (e) { - log(`FATAL: cannot access guild/channel — is the bot invited to guild ${GUILD_ID}? (${e.message})`); - log('run `node bot.mjs --invite` and have a server admin authorise the bot, then retry.'); - return leaveAndExit(1); + log(`voice endpoint: ${VOICE_ENDPOINT} · control: ${REPORT_ENDPOINT}`); + + // Backward-compat: if a default guild/channel is configured, auto-join it. + // Otherwise idle and wait for the dashboard to pick a channel. + if (GUILD_ID && CHANNEL_ID) { + await joinChannel(GUILD_ID, CHANNEL_ID).catch((e) => log(`initial join failed: ${e.message}`)); + } else { + log('no default channel — waiting for the dashboard to select a server/voice channel…'); } - if (!channel || !channel.isVoiceBased()) { log(`FATAL: channel ${CHANNEL_ID} is not a voice channel`); return leaveAndExit(1); } - log(`joining voice: guild="${guild.name}" channel="${channel.name}"`); - const connection = joinVoiceChannel({ - channelId: CHANNEL_ID, - guildId: GUILD_ID, - adapterCreator: guild.voiceAdapterCreator, - selfDeaf: false, // MUST be false to receive audio (the STT input path) - selfMute: false, // false so we can also speak later (M5 TTS) - }); - connection.on('error', (e) => log(`voice connection error: ${e.message}`)); - - try { - // The DAVE/MLS handshake here cycles signalling<->connecting several times - // and can take ~25s, so give it a generous ceiling before declaring failure. - await entersState(connection, VoiceConnectionStatus.Ready, 40_000); - } catch (e) { - log(`FATAL: voice connection did not become Ready in 40s (${e.message})`); - return leaveAndExit(1); - } - log(`✅ JOINED & READY. channel=${CHANNEL_ID} — staying connected, listening for speakers…`); - - // Playback path (bot speaks): one player, subscribed to the connection. - voicePlayer = createAudioPlayer({ behaviors: { noSubscriber: NoSubscriberBehavior.Play } }); - voicePlayer.on('error', (e) => log(`player error: ${e.message}`)); - connection.subscribe(voicePlayer); - log(`voice endpoint: ${VOICE_ENDPOINT}`); - - // ---------- receive path: capture each utterance and run the voice turn ---- - const receiver = connection.receiver; - const active = new Set(); // userIds with an in-flight subscription (avoid dupes) - receiver.speaking.on('start', (userId) => { - if (userId === client.user.id || active.has(userId)) return; // skip self / dupes - active.add(userId); - if (!perUser.has(userId)) perUser.set(userId, { opusPackets: 0, pcmFrames: 0 }); - log(`SPEAKING start user=${userId}`); - const opusStream = receiver.subscribe(userId, { - // End the utterance after a short silence so natural pauses don't cut words. - end: { behavior: EndBehaviorType.AfterSilence, duration: 800 }, - }); - // Decode Opus -> 48kHz stereo s16le PCM, buffered until the utterance ends. - const decoder = new prism.opus.Decoder({ rate: 48000, channels: 2, frameSize: 960 }); - const chunks = []; - opusStream.on('data', () => { perUser.get(userId).opusPackets++; }); - // A receive-stream error (e.g. a DAVE decrypt/UDP GenericFailure on one - // packet) must NOT crash the process — log it and free the slot so the - // next utterance still works. - opusStream.on('error', (e) => { logThrottled(`recv:${e.message}`, `recv stream error user=${userId}: ${e.message}`); active.delete(userId); }); - opusStream.pipe(decoder); - decoder.on('data', (d) => { - chunks.push(d); - const s = perUser.get(userId); s.pcmFrames++; - }); - decoder.on('error', (e) => log(`decode error user=${userId}: ${e.message}`)); - decoder.on('end', () => { - active.delete(userId); - const pcm = Buffer.concat(chunks); - log(`utterance end user=${userId} pcm=${pcm.length}B — running voice turn`); - handleUtterance(userId, pcm).catch((e) => log(`voice turn error: ${e.message}`)); - }); - }); - - connection.on(VoiceConnectionStatus.Disconnected, () => { - log('voice: disconnected — attempting to resume…'); - Promise.race([ - entersState(connection, VoiceConnectionStatus.Signalling, 5_000), - entersState(connection, VoiceConnectionStatus.Connecting, 5_000), - ]).catch(() => { log('voice: could not resume, leaving'); leaveAndExit(0); }); - }); + // Report state + poll commands forever. This is what powers the dashboard's + // bot info, server/voice-channel pickers, participant list, and join/leave. + reportLoop(); + const reportTimer = setInterval(reportLoop, REPORT_INTERVAL_MS); + if (typeof reportTimer.unref === 'function') reportTimer.unref(); }); client.on('error', (e) => log('client error', e.message)); diff --git a/wsai/bot_control.py b/wsai/bot_control.py new file mode 100644 index 0000000..c05d494 --- /dev/null +++ b/wsai/bot_control.py @@ -0,0 +1,61 @@ +"""Control plane between the dashboard (Python) and the Discord bot (Node). + +The bot only ever makes *outbound* HTTP (it already POSTs voice turns), so we +keep that single direction: + +* the bot PUSHES its live state here (identity, joinable guilds + voice + channels, current channel, members, whitelist/blacklist) via + ``POST /api/bot/report`` → :meth:`report`; +* the bot POLLS ``GET /api/bot/commands`` → :meth:`drain` for pending commands + (join a channel, leave) that the dashboard UI enqueued via :meth:`enqueue`. + +Thread-safe: the HTTP handler threads read/write from several threads. +""" + +from __future__ import annotations + +import threading +import time +from collections import deque +from typing import Any + +# The bot is considered offline if it hasn't reported within this window. +STALE_AFTER_S = 8.0 + + +class BotControl: + def __init__(self) -> None: + self._lock = threading.Lock() + self._state: dict[str, Any] = {"connected": False, "ts": 0.0} + self._commands: deque[dict[str, Any]] = deque() + self._cmd_id = 0 + + # -- bot -> dashboard (state push) ----------------------------------- # + def report(self, state: dict[str, Any]) -> None: + s = dict(state) + s["ts"] = time.time() + s["connected"] = True + with self._lock: + self._state = s + + def state(self) -> dict[str, Any]: + with self._lock: + s = dict(self._state) + # A report older than STALE_AFTER_S means the bot stopped polling/pushing. + if s.get("connected") and (time.time() - s.get("ts", 0.0)) > STALE_AFTER_S: + s["connected"] = False + return s + + # -- dashboard -> bot (command queue) -------------------------------- # + def enqueue(self, cmd: dict[str, Any]) -> int: + with self._lock: + self._cmd_id += 1 + cmd = {**cmd, "id": self._cmd_id} + self._commands.append(cmd) + return self._cmd_id + + def drain(self) -> list[dict[str, Any]]: + with self._lock: + cmds = list(self._commands) + self._commands.clear() + return cmds diff --git a/wsai/dashboard.py b/wsai/dashboard.py index b3d0687..6a3ad96 100644 --- a/wsai/dashboard.py +++ b/wsai/dashboard.py @@ -74,6 +74,10 @@ def _make_handler(dash: "Dashboard"): self._send(200, body, "application/json; charset=utf-8") elif path == "/api/prompt": self._handle_prompt_get() + elif path == "/api/bot/state": + self._send_json(dash.bot.state()) + elif path == "/api/bot/commands": + self._send_json({"commands": dash.bot.drain()}) elif path == "/events": self._stream_events() else: @@ -94,6 +98,10 @@ def _make_handler(dash: "Dashboard"): self._handle_log_mutate("delete") elif path == "/api/logs/edit": self._handle_log_mutate("edit") + elif path == "/api/bot/report": + self._handle_bot_report() + elif path == "/api/bot/select": + self._handle_bot_select() else: self._send(404, b"not found", "text/plain; charset=utf-8") @@ -119,8 +127,9 @@ def _make_handler(dash: "Dashboard"): self._send(400, json.dumps({"ok": False, "error": "empty upload"}).encode(), "application/json; charset=utf-8") return + speaker = urllib.parse.unquote(self.headers.get("X-User-Name", "") or "") try: - res = dash.voice_turn(raw) + res = dash.voice_turn(raw, speaker=speaker) except Exception as exc: # noqa: BLE001 log.exception("voice-turn failed") self._send(500, json.dumps({"ok": False, "error": f"{type(exc).__name__}: {exc}"}, @@ -206,6 +215,42 @@ def _make_handler(dash: "Dashboard"): }, ensure_ascii=False).encode("utf-8") self._send(200, body, "application/json; charset=utf-8") + def _send_json(self, obj, code: int = 200) -> None: + self._send(code, json.dumps(obj, ensure_ascii=False).encode("utf-8"), + "application/json; charset=utf-8") + + def _handle_bot_report(self) -> None: + """The Discord bot pushes its live state (identity, guilds, voice + channels, current channel + members, list settings).""" + raw = self._read_body() + try: + data = json.loads(raw.decode("utf-8")) if raw else {} + except (ValueError, AttributeError): + self._send_json({"ok": False, "error": "invalid JSON"}, 400) + return + dash.bot.report(data) + # Hand the bot any queued commands in the same round trip so it does + # not have to poll a second endpoint. + self._send_json({"ok": True, "commands": dash.bot.drain()}) + + def _handle_bot_select(self) -> None: + """UI picked a server/voice channel → queue a join (or leave) command.""" + raw = self._read_body() + try: + data = json.loads(raw.decode("utf-8")) if raw else {} + except (ValueError, AttributeError): + self._send_json({"ok": False, "error": "invalid JSON"}, 400) + return + guild_id = (data.get("guildId") or "").strip() + channel_id = (data.get("channelId") or "").strip() + if channel_id and guild_id: + cid = dash.bot.enqueue({"type": "join", "guildId": guild_id, "channelId": channel_id}) + monitor.log("info", f"음성채널 참여 요청 (guild={guild_id} channel={channel_id})") + else: + cid = dash.bot.enqueue({"type": "leave"}) + monitor.log("info", "음성채널 나가기 요청") + self._send_json({"ok": True, "commandId": cid}) + def _handle_log_mutate(self, action: str) -> None: """Per-line log delete/edit by event id.""" raw = self._read_body() @@ -267,12 +312,14 @@ class Dashboard: def __init__(self, monitor: Monitor, host: str = "0.0.0.0", port: int = 8787, stt=None, tts=None, brain=None, history_turns: int = 12) -> None: + from .bot_control import BotControl self.monitor = monitor self.host = host self.port = port self.stt = stt self.tts = tts self.brain = brain + self.bot = BotControl() # dashboard <-> Discord bot control plane self._history: list[tuple[str, str]] = [] self._history_turns = history_turns self._server: ThreadingHTTPServer | None = None @@ -325,7 +372,7 @@ class Dashboard: if self.tts is not None: self._submit(self.tts._ensure()) - def voice_turn(self, audio_bytes: bytes) -> dict: + def voice_turn(self, audio_bytes: bytes, speaker: str = "") -> dict: """One Discord voice turn: decode the uploaded utterance, recognise it on the GPU, think of a reply (Claude brain if wired, else echo), synthesise it on the GPU, and return {heard, reply, wav} where wav is the @@ -344,6 +391,8 @@ class Dashboard: with open(src, "wb") as f: f.write(audio_bytes) turn = self.monitor.turn(source="discord") + if speaker: + turn.speaker = speaker # who spoke (for the "누가 말했는지" log) t0 = time.monotonic() try: subprocess.run( @@ -545,6 +594,18 @@ PAGE = r""" .sttres .txt{color:var(--heard);font-weight:600;line-height:1.5} .sttres .meta{color:var(--muted);font-size:12px;margin-top:5px} .hbtn{padding:6px 12px;font-size:12.5px} + /* Bot control bar (봇 정보 · 서버/채널 선택 · 참여자) */ + .botbar{display:flex;gap:14px;align-items:center;flex-wrap:wrap;background:var(--panel); + border:1px solid var(--line);border-radius:12px;padding:10px 14px;margin:0 0 16px;font-size:13px} + .botbar label{display:flex;gap:6px;align-items:center;color:var(--muted)} + .botbar select{background:var(--panel2);border:1px solid var(--line);color:var(--fg); + border-radius:8px;padding:6px 9px;font-size:13px;max-width:230px} + .botinfo{display:flex;gap:7px;align-items:center;font-weight:600} + .parts{display:flex;gap:6px;align-items:center;flex-wrap:wrap;margin-left:auto;color:var(--muted)} + .part{display:inline-flex;gap:5px;align-items:center;background:var(--panel2);border:1px solid var(--line); + border-radius:999px;padding:3px 10px;font-size:12px} + .part.spk{border-color:#1f5236;color:#9ff0bd} + .speaker{color:var(--muted);font-size:11.5px} /* Modal / popup (reused by 프롬프트, 화이트/블랙리스트 …) */ .modal{position:fixed;inset:0;z-index:20;background:rgba(4,7,11,.66); display:flex;align-items:center;justify-content:center;padding:20px} @@ -614,6 +675,12 @@ PAGE = r"""
+
+ 봇: 연결 안 됨 + + + +