diff --git a/tests/test_monitor.py b/tests/test_monitor.py new file mode 100644 index 0000000..0c9c930 --- /dev/null +++ b/tests/test_monitor.py @@ -0,0 +1,95 @@ +"""The monitor must record step-by-step turns (heard / thought / answered, +per-step timing, ok vs error) and stream them to subscribers — that data is +exactly what the status dashboard renders.""" + +import asyncio +import json + +from wsai.backends.mock import MockBrain, MockSTT, MockTTS +from wsai.monitor import Monitor +from wsai.pipeline import Pipeline + + +def test_monitor_records_turn_with_timed_steps(): + async def go(): + mon = Monitor() + + pipe = Pipeline( + brain=MockBrain(), + stt=MockSTT(script=["안녕"], interval=0.01), + tts=MockTTS(), + monitor=mon, + ) + await asyncio.wait_for(pipe.run(), timeout=5) + return mon + + mon = asyncio.run(go()) + snap = mon.snapshot() + + assert snap["status"]["turns_total"] == 1 + assert snap["status"]["running"] is False # cleaned up after run + assert len(snap["turns"]) == 1 + + turn = snap["turns"][0] + assert turn["heard"] == "안녕" # what it heard + assert turn["reply"] # what it answered + assert turn["status"] == "ok" # it worked + assert turn["total_ms"] >= 0 + # step-by-step: every stage is named and timed + names = [s["name"] for s in turn["steps"]] + assert names == ["화면 맥락", "두뇌(생각)", "응답(TTS/전송)"] + assert all(s["ok"] is True for s in turn["steps"]) + assert all(s["ms"] >= 0 for s in turn["steps"]) + + +def test_monitor_marks_errors(): + class BoomBrain(MockBrain): + async def respond(self, user_text, screen, history): + raise RuntimeError("boom") + + async def go(): + mon = Monitor() + pipe = Pipeline( + brain=BoomBrain(), + stt=MockSTT(script=["안녕"], interval=0.01), + tts=MockTTS(), + monitor=mon, + ) + try: + await asyncio.wait_for(pipe.run(), timeout=5) + except BaseException: + pass # TaskGroup re-raises; we only care about recorded telemetry + return mon + + mon = asyncio.run(go()) + snap = mon.snapshot() + + turn = snap["turns"][0] + assert turn["status"] == "error" + brain_step = next(s for s in turn["steps"] if s["name"] == "두뇌(생각)") + assert brain_step["ok"] is False + assert "boom" in brain_step["error"] + assert snap["status"]["errors_total"] >= 1 + + +def test_subscriber_receives_live_turn_events(): + async def go(): + mon = Monitor() + q = mon.subscribe() + pipe = Pipeline( + brain=MockBrain(), + stt=MockSTT(script=["안녕"], interval=0.01), + tts=MockTTS(), + monitor=mon, + ) + await asyncio.wait_for(pipe.run(), timeout=5) + return q + + q = asyncio.run(go()) + events = [] + while not q.empty(): + events.append(json.loads(q.get_nowait())) + + types = {e["type"] for e in events} + assert "turn" in types # live turn updates were pushed + assert "status" in types # listening/running status changes were pushed diff --git a/wsai/__main__.py b/wsai/__main__.py index c144ac6..6b31399 100644 --- a/wsai/__main__.py +++ b/wsai/__main__.py @@ -2,11 +2,12 @@ python -m wsai # mock pipeline (no deps, no keys) — runs a demo python -m wsai --voice # eyes-free voice loop demo (STT -> Brain -> TTS) + python -m wsai --dashboard # live status website + a continuous voice demo python -m wsai --live # capture this screen + Claude eyes/brain python -m wsai --env # build from WSAI_* environment variables The mock run is bounded (a few frames + a scripted conversation) so it exits on -its own; --live/--env run until Ctrl-C. +its own; --dashboard/--live/--env run until Ctrl-C. """ from __future__ import annotations @@ -14,32 +15,53 @@ from __future__ import annotations import argparse import asyncio import logging +import os +import socket from .config import Settings from .factory import build +from .monitor import Monitor -async def _run(settings: Settings, demo: bool) -> None: +def _lan_ip() -> str: + try: + s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + s.connect(("8.8.8.8", 80)) + ip = s.getsockname()[0] + s.close() + return ip + except OSError: + return "127.0.0.1" + + +async def _run(settings: Settings, demo: bool, monitor: Monitor | None) -> None: if demo: # Bounded demo so CI / a quick check terminates. from .backends.mock import MockFrameSource, MockSTT - from .pipeline import Pipeline - pipe = build(settings) + pipe = build(settings, monitor=monitor) if pipe.source is not None: # keep eyes-free configs eyes-free pipe.source = MockFrameSource(interval=0.3, limit=4) if pipe.stt is not None: - pipe.stt = MockSTT(interval=0.4) + # When serving the dashboard, keep talking forever so there's always + # something live to watch; otherwise stay bounded so the demo exits. + loop = monitor is not None + pipe.stt = MockSTT(interval=2.0 if loop else 0.4, loop=loop) await pipe.run() return - await build(settings).run() + await build(settings, monitor=monitor).run() def main() -> None: ap = argparse.ArgumentParser(prog="wsai") ap.add_argument("--voice", action="store_true", help="eyes-free voice loop (STT -> Brain -> TTS)") + ap.add_argument("--dashboard", action="store_true", help="serve the live status website (voice demo loops)") ap.add_argument("--live", action="store_true", help="capture screen + Claude backends") ap.add_argument("--env", action="store_true", help="build from WSAI_* env vars") + ap.add_argument("--port", type=int, default=int(os.environ.get("WSAI_DASHBOARD_PORT", "8787")), + help="dashboard port (default 8787, or WSAI_DASHBOARD_PORT)") + ap.add_argument("--host", default=os.environ.get("WSAI_DASHBOARD_HOST", "0.0.0.0"), + help="dashboard bind host (default 0.0.0.0)") ap.add_argument("-v", "--verbose", action="store_true") args = ap.parse_args() @@ -48,7 +70,11 @@ def main() -> None: format="%(levelname)s %(name)s: %(message)s", ) - if args.voice: + if args.dashboard: + # Default to the eyes-free voice preset for the demo; env can override. + settings = Settings.from_env() if args.env else Settings.voice() + demo = not args.env + elif args.voice: settings, demo = Settings.voice(), True elif args.live: settings, demo = Settings.live(), False @@ -57,10 +83,25 @@ def main() -> None: else: settings, demo = Settings.mock(), True + monitor: Monitor | None = None + dash = None + if args.dashboard: + from .dashboard import Dashboard + + monitor = Monitor() + dash = Dashboard(monitor, host=args.host, port=args.port) + dash.start() + shown = args.host if args.host not in ("0.0.0.0", "") else _lan_ip() + print(f"\n 실시간 상태 사이트: http://{shown}:{args.port}") + print(f" (로컬: http://127.0.0.1:{args.port} )\n") + try: - asyncio.run(_run(settings, demo)) + asyncio.run(_run(settings, demo, monitor)) except KeyboardInterrupt: pass + finally: + if dash is not None: + dash.stop() if __name__ == "__main__": diff --git a/wsai/backends/mock.py b/wsai/backends/mock.py index 1731cae..13bd097 100644 --- a/wsai/backends/mock.py +++ b/wsai/backends/mock.py @@ -63,18 +63,27 @@ class MockVision: class MockSTT: """Feeds a scripted set of user utterances, then goes quiet.""" - def __init__(self, script: list[str] | None = None, interval: float = 2.0) -> None: + def __init__( + self, + script: list[str] | None = None, + interval: float = 2.0, + loop: bool = False, + ) -> None: self.script = script or [ "지금 화면에 뭐 보여?", "저 에러 왜 나는 거야?", "고마워", ] self.interval = interval + self.loop = loop async def utterances(self) -> AsyncIterator[Utterance]: - for line in self.script: - await asyncio.sleep(self.interval) - yield Utterance(text=line, ts=time.monotonic(), source="voice") + while True: + for line in self.script: + await asyncio.sleep(self.interval) + yield Utterance(text=line, ts=time.monotonic(), source="voice") + if not self.loop: + return async def aclose(self) -> None: return diff --git a/wsai/dashboard.py b/wsai/dashboard.py new file mode 100644 index 0000000..6f09397 --- /dev/null +++ b/wsai/dashboard.py @@ -0,0 +1,328 @@ +"""Live status website for the voice loop. + +Serves a single self-contained page plus a Server-Sent-Events stream so you can +open a browser and watch, step by step: is it listening, what it heard, what it +thought/answered, how long each stage took, and whether anything errored. + +Pure stdlib (``http.server``). Runs in a background thread so it never blocks +the asyncio pipeline. + +Endpoints: + GET / -> the dashboard HTML + GET /api/state -> JSON snapshot (initial load / fallback polling) + GET /events -> text/event-stream live push +""" + +from __future__ import annotations + +import json +import logging +import queue +import threading +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer + +from .monitor import Monitor + +log = logging.getLogger("wsai.dashboard") + + +def _make_handler(monitor: Monitor): + class Handler(BaseHTTPRequestHandler): + # Quiet: don't spam the console with one line per request. + def log_message(self, *args) -> None: # noqa: D401 + return + + def _send(self, code: int, body: bytes, ctype: str) -> None: + self.send_response(code) + self.send_header("Content-Type", ctype) + self.send_header("Content-Length", str(len(body))) + self.send_header("Cache-Control", "no-store") + self.end_headers() + self.wfile.write(body) + + def do_GET(self) -> None: # noqa: N802 + path = self.path.split("?", 1)[0] + if path == "/" or path == "/index.html": + self._send(200, PAGE.encode("utf-8"), "text/html; charset=utf-8") + elif path == "/api/state": + body = json.dumps(monitor.snapshot(), ensure_ascii=False).encode("utf-8") + self._send(200, body, "application/json; charset=utf-8") + elif path == "/events": + self._stream_events() + else: + self._send(404, b"not found", "text/plain; charset=utf-8") + + def _stream_events(self) -> None: + self.send_response(200) + self.send_header("Content-Type", "text/event-stream; charset=utf-8") + self.send_header("Cache-Control", "no-store") + self.send_header("Connection", "keep-alive") + self.end_headers() + q = monitor.subscribe() + try: + # Prime the client with a full snapshot so it renders instantly. + first = json.dumps( + {"type": "snapshot", "snapshot": monitor.snapshot()}, + ensure_ascii=False, + ) + self.wfile.write(f"data: {first}\n\n".encode("utf-8")) + self.wfile.flush() + while True: + try: + data = q.get(timeout=15) + except queue.Empty: + # Heartbeat keeps proxies / the browser from timing out. + self.wfile.write(b": ping\n\n") + self.wfile.flush() + continue + self.wfile.write(f"data: {data}\n\n".encode("utf-8")) + self.wfile.flush() + except (BrokenPipeError, ConnectionResetError): + pass + finally: + monitor.unsubscribe(q) + + return Handler + + +class Dashboard: + """Owns the HTTP server thread.""" + + def __init__(self, monitor: Monitor, host: str = "0.0.0.0", port: int = 8787) -> None: + self.monitor = monitor + self.host = host + self.port = port + self._server: ThreadingHTTPServer | None = None + self._thread: threading.Thread | None = None + + def start(self) -> None: + handler = _make_handler(self.monitor) + self._server = ThreadingHTTPServer((self.host, self.port), handler) + self._server.daemon_threads = True + self._thread = threading.Thread( + target=self._server.serve_forever, name="wsai-dashboard", daemon=True + ) + self._thread.start() + log.info("dashboard on http://%s:%d", self.host, self.port) + + def stop(self) -> None: + if self._server is not None: + self._server.shutdown() + self._server.server_close() + self._server = None + + +# --------------------------------------------------------------------------- # +# The page. One file, no external assets, so it works offline / behind a LAN. +# --------------------------------------------------------------------------- # +PAGE = r""" + + + + +watch_sceen_ai · 실시간 상태 + + + +
+
+

watch_sceen_ai · 실시간 상태

+
STT → 두뇌 → TTS 음성 루프를 단계별로 관찰
+
+
연결 대기
+
+
0대화 수
+
0오류
+
0초가동시간
+
·연결
+
+
+
+
+
+
아직 대화가 없습니다. 사용자가 말하면 여기에 단계별로 나타납니다.
+
+

이벤트 / 오류 로그

+
+
+
+ + + +""" diff --git a/wsai/factory.py b/wsai/factory.py index c3fac74..5c268d4 100644 --- a/wsai/factory.py +++ b/wsai/factory.py @@ -4,18 +4,32 @@ concrete class each config name maps to, so adding a backend = one line here.""" from __future__ import annotations from .config import Settings +from .monitor import Monitor from .pipeline import Pipeline -def build(settings: Settings) -> Pipeline: - return Pipeline( +def build(settings: Settings, monitor: Monitor | None = None) -> Pipeline: + pipe = Pipeline( source=_source(settings), vision=_vision(settings), brain=_brain(settings), stt=_stt(settings), tts=_tts(settings), text_channel=_text(settings), + monitor=monitor, ) + if monitor is not None: + monitor.set_components( + { + "source": settings.source or "none", + "vision": settings.vision or "none", + "stt": settings.stt or "none", + "brain": settings.brain, + "tts": settings.tts or "none", + "text": settings.text or "none", + } + ) + return pipe def _source(s: Settings): diff --git a/wsai/monitor.py b/wsai/monitor.py new file mode 100644 index 0000000..1382469 --- /dev/null +++ b/wsai/monitor.py @@ -0,0 +1,226 @@ +"""Telemetry hub for the live status dashboard. + +The pipeline is a chain of steps (heard -> screen context -> brain -> speak). +This module records, for every conversation turn, *what happened at each step* +and *how long it took*, plus a rolling status header and any errors. The +dashboard (``wsai/dashboard.py``) reads a snapshot and subscribes for live +push updates. + +Design notes: +* Pure stdlib, no deps — matches the project's "core has no third-party deps". +* Thread-safe. The pipeline mutates it from the asyncio loop; the HTTP server + reads/subscribes from its own threads. A single lock guards everything. +* A Monitor with zero subscribers is essentially free, so the pipeline can + always hold one (no separate no-op path). +""" + +from __future__ import annotations + +import json +import queue +import threading +import time +from collections import deque +from typing import Any + + +def _now_wall() -> float: + # Wall-clock seconds for human-readable timestamps on the page. + return time.time() + + +def _now_mono() -> float: + # Monotonic seconds for measuring durations (immune to clock jumps). + return time.monotonic() + + +class Step: + """One timed stage inside a turn (e.g. "두뇌"). Used as an async context + manager so it can wrap an ``await`` and record ok/error + elapsed ms.""" + + def __init__(self, turn: "Turn", name: str) -> None: + self.turn = turn + self.name = name + self.ok: bool | None = None + self.ms: float = 0.0 + self.detail: str = "" + self.error: str = "" + self._t0 = 0.0 + + async def __aenter__(self) -> "Step": + self._t0 = _now_mono() + self.turn._steps.append(self) + self.turn._touch() + return self + + async def __aexit__(self, exc_type, exc, tb) -> bool: + self.ms = (_now_mono() - self._t0) * 1000.0 + if exc is not None: + self.ok = False + self.error = f"{exc_type.__name__}: {exc}" + else: + self.ok = True + self.turn._touch() + return False # never swallow: the pipeline/TaskGroup must still see it + + def to_dict(self) -> dict[str, Any]: + return { + "name": self.name, + "ok": self.ok, + "ms": round(self.ms, 1), + "detail": self.detail, + "error": self.error, + } + + +class Turn: + """One user utterance and everything the AI did in response.""" + + def __init__(self, monitor: "Monitor", turn_id: int, source: str) -> None: + self._monitor = monitor + self.id = turn_id + self.source = source + self.wall = _now_wall() + self._t0 = _now_mono() + self.heard_text = "" + self.reply_text = "" + self.status = "active" # active | ok | error + self.error = "" + self.total_ms = 0.0 + self._steps: list[Step] = [] + + # -- recording API (called from the pipeline) ------------------------- # + def heard(self, text: str) -> None: + self.heard_text = text + self._touch() + + def replied(self, text: str) -> None: + self.reply_text = text + self._touch() + + def step(self, name: str) -> Step: + return Step(self, name) + + def finish(self, error: str = "") -> None: + self.total_ms = (_now_mono() - self._t0) * 1000.0 + if error: + self.status = "error" + self.error = error + elif any(s.ok is False for s in self._steps): + self.status = "error" + else: + self.status = "ok" + self._touch() + + # -- internal --------------------------------------------------------- # + def _touch(self) -> None: + self._monitor._publish(self) + + def to_dict(self) -> dict[str, Any]: + return { + "id": self.id, + "source": self.source, + "wall": self.wall, + "heard": self.heard_text, + "reply": self.reply_text, + "status": self.status, + "error": self.error, + "total_ms": round(self.total_ms, 1), + "steps": [s.to_dict() for s in self._steps], + } + + +class Monitor: + """Rolling record of turns + status, with a pub/sub for live updates.""" + + def __init__(self, keep: int = 60) -> None: + self._turns: deque[Turn] = deque(maxlen=keep) + self._events: deque[dict[str, Any]] = deque(maxlen=200) + self._status: dict[str, Any] = { + "running": False, + "listening": False, + "started_wall": _now_wall(), + "components": {}, + "turns_total": 0, + "errors_total": 0, + } + self._lock = threading.Lock() + self._subs: list["queue.Queue[str]"] = [] + self._id = 0 + + # -- status ----------------------------------------------------------- # + def set_status(self, **kw: Any) -> None: + with self._lock: + self._status.update(kw) + self._broadcast({"type": "status", "status": self.status_snapshot()}) + + def set_components(self, components: dict[str, Any]) -> None: + with self._lock: + self._status["components"] = components + self._broadcast({"type": "status", "status": self.status_snapshot()}) + + def status_snapshot(self) -> dict[str, Any]: + with self._lock: + s = dict(self._status) + s["uptime_s"] = round(_now_wall() - s["started_wall"], 1) + return s + + def log(self, level: str, message: str) -> None: + """A free-form lifecycle/error line (startup, disconnect, crash…).""" + evt = {"type": "log", "level": level, "message": message, "wall": _now_wall()} + with self._lock: + self._events.append(evt) + if level == "error": + self._status["errors_total"] += 1 + self._broadcast(evt) + + # -- turns ------------------------------------------------------------ # + def turn(self, source: str = "voice") -> Turn: + with self._lock: + self._id += 1 + self._status["turns_total"] += 1 + t = Turn(self, self._id, source) + self._turns.append(t) + self._publish(t) + return t + + def _publish(self, t: Turn) -> None: + # Recompute error total lazily on error transitions. + if t.status == "error": + with self._lock: + # errors_total counts turns that ended in error at most once + pass + self._broadcast({"type": "turn", "turn": t.to_dict()}) + + # -- snapshot / subscribe (read side, HTTP threads) ------------------- # + def snapshot(self) -> dict[str, Any]: + with self._lock: + turns = [t.to_dict() for t in self._turns] + events = list(self._events) + return { + "status": self.status_snapshot(), + "turns": turns, + "events": events, + } + + def subscribe(self) -> "queue.Queue[str]": + q: "queue.Queue[str]" = queue.Queue(maxsize=256) + with self._lock: + self._subs.append(q) + return q + + def unsubscribe(self, q: "queue.Queue[str]") -> None: + with self._lock: + if q in self._subs: + self._subs.remove(q) + + def _broadcast(self, event: dict[str, Any]) -> None: + data = json.dumps(event, ensure_ascii=False) + with self._lock: + subs = list(self._subs) + for q in subs: + try: + q.put_nowait(data) + except queue.Full: + # Slow client: drop it rather than block the pipeline. + self.unsubscribe(q) diff --git a/wsai/pipeline.py b/wsai/pipeline.py index e5e1ca0..cc8841c 100644 --- a/wsai/pipeline.py +++ b/wsai/pipeline.py @@ -15,6 +15,7 @@ from .interfaces import ( Utterance, VisionBackend, ) +from .monitor import Monitor from .state import SharedScreenContext log = logging.getLogger("wsai.pipeline") @@ -41,6 +42,7 @@ class Pipeline: tts: TextToSpeech | None = None, text_channel: TextChannel | None = None, history_turns: int = 12, + monitor: Monitor | None = None, ) -> None: self.source = source self.vision = vision @@ -48,6 +50,7 @@ class Pipeline: self.stt = stt self.tts = tts self.text_channel = text_channel + self.monitor = monitor self.context = SharedScreenContext() self._history: list[tuple[str, str]] = [] self._history_turns = history_turns @@ -59,18 +62,43 @@ class Pipeline: async for frame in self.source.frames(): try: obs = await self.vision.describe(frame) - except Exception: # a single bad frame must not kill the loop + except Exception as exc: # a single bad frame must not kill the loop log.exception("vision.describe failed") + if self.monitor is not None: + self.monitor.log("error", f"화면 이해 실패: {exc}") continue await self.context.update(obs) log.debug("screen: %s", obs.text[:120]) # -- conversation ------------------------------------------------------ # async def _handle(self, utt: Utterance) -> None: - screen = await self.context.latest() - reply = await self.brain.respond(utt.text, screen, self._history) - self._remember(utt.text, reply.text) - await self._emit(reply) + if self.monitor is None: + screen = await self.context.latest() + reply = await self.brain.respond(utt.text, screen, self._history) + self._remember(utt.text, reply.text) + await self._emit(reply) + return + + # Same work, but each stage is timed and streamed to the dashboard so a + # viewer can see what was heard, what the brain answered, how long each + # step took, and whether anything errored. + turn = self.monitor.turn(source=utt.source) + turn.heard(utt.text) + try: + async with turn.step("화면 맥락"): + screen = await self.context.latest() + async with turn.step("두뇌(생각)"): + reply = await self.brain.respond(utt.text, screen, self._history) + turn.replied(reply.text) + self._remember(utt.text, reply.text) + async with turn.step("응답(TTS/전송)"): + await self._emit(reply) + except Exception as exc: + turn.finish(error=f"{type(exc).__name__}: {exc}") + self.monitor.log("error", f"대화 #{turn.id} 실패: {exc}") + raise + else: + turn.finish() def _remember(self, user: str, ai: str) -> None: self._history.append((user, ai)) @@ -91,8 +119,15 @@ class Pipeline: async def _listen_voice(self) -> None: if self.stt is None: return - async for utt in self.stt.utterances(): - await self._handle(utt) + if self.monitor is not None: + self.monitor.set_status(listening=True) + self.monitor.log("info", "음성 수신 시작 — 발화 대기 중") + try: + async for utt in self.stt.utterances(): + await self._handle(utt) + finally: + if self.monitor is not None: + self.monitor.set_status(listening=False) async def _listen_text(self) -> None: if self.text_channel is None: @@ -107,12 +142,23 @@ class Pipeline: # failing loop propagated while the siblings kept running detached, and # aclose() in the finally then closed a source/stt out from under a # still-live loop (close-during-use). + if self.monitor is not None: + self.monitor.set_status(running=True) + self.monitor.log("info", "파이프라인 시작") try: async with asyncio.TaskGroup() as tg: tg.create_task(self._perceive()) tg.create_task(self._listen_voice()) tg.create_task(self._listen_text()) + except* Exception as eg: + if self.monitor is not None: + for exc in eg.exceptions: + self.monitor.log("error", f"루프 예외: {type(exc).__name__}: {exc}") + raise finally: + if self.monitor is not None: + self.monitor.set_status(running=False, listening=False) + self.monitor.log("info", "파이프라인 종료") await self.aclose() async def aclose(self) -> None: