linnet/app/api/events.py
pyr0ball 7e14f9135e feat: Notation v0.1.x scaffold — full backend + frontend + tests
FastAPI backend (port 8522):
- Session lifecycle: POST /session/start, DELETE /session/{id}/end, GET /session/{id}
- SSE stream: GET /session/{id}/stream — per-subscriber asyncio.Queue fan-out, 15s heartbeat
- History: GET /session/{id}/history with min_confidence + limit filters
- Audio: WS /session/{id}/audio — binary PCM ingestion stub (real inference in v0.2.x)
- Export: GET /session/{id}/export — downloadable JSON session log
- ContextClassifier background task per session (CF_VOICE_MOCK=1 in dev)
- ToneEvent SSE wire format per cf-core#40 (locked field names)
- Tier gate: CFG-LNNT- prefix check, 402 for paid features

Vue 3 frontend (port 8521, Vite + UnoCSS + Pinia):
- NowPanel: affect-aware border tint, subtext, prosody flags, shift indicator
- HistoryStrip: horizontal scroll, last 8 events with affect color
- ComposeBar: start/stop session, SSE connection lifecycle
- useToneStream: EventSource composable
- useAudioCapture: AudioWorklet → Int16 PCM → WebSocket (v0.1.x stub)
- audio-processor.js: 100ms chunk accumulator in AudioWorklet thread
- Respects prefers-reduced-motion globally

26 tests passing, manage.sh, Dockerfile, compose.yml included.
2026-04-06 18:23:52 -07:00

58 lines
1.9 KiB
Python

# app/api/events.py — SSE stream endpoint
from __future__ import annotations
import asyncio
import logging
from fastapi import APIRouter, HTTPException, Request
from fastapi.responses import StreamingResponse
from app.services import session_store
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/session", tags=["events"])
@router.get("/{session_id}/stream")
async def stream_events(session_id: str, request: Request) -> StreamingResponse:
"""
SSE stream of ToneEvent annotations for a session.
Clients connect with EventSource:
const es = new EventSource(`/session/${sessionId}/stream`)
es.addEventListener('tone-event', e => { ... })
The stream runs until the client disconnects or the session ends.
"""
session = session_store.get_session(session_id)
if session is None:
raise HTTPException(status_code=404, detail=f"Session {session_id} not found")
queue = session.subscribe()
async def generator():
try:
# Heartbeat comment every 15s to keep connection alive through proxies
yield ": heartbeat\n\n"
while True:
if await request.is_disconnected():
break
try:
event = await asyncio.wait_for(queue.get(), timeout=15.0)
yield event.to_sse()
except asyncio.TimeoutError:
yield ": heartbeat\n\n"
except asyncio.CancelledError:
break
finally:
session.unsubscribe(queue)
logger.debug("SSE subscriber disconnected from session %s", session_id)
return StreamingResponse(
generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no", # disable nginx buffering
},
)