8월 7일까지 워킹트리에만 남아 있던 미커밋 작업을 커밋한다. 여러 사본 폴더(worktree·clone)에 흩어져 있던 중간 스냅샷을 정리하기 전에 원본을 git 이력으로 고정하는 것이 목적이다. - contracts/routes/services: measurement, outcome_trajectory, rupture_repair, deliberate_practice, calibration_transfer, supervision_research, multimodal_alliance, continuous_improvement 계열 신규 모듈과 테스트 - infra/db/init: 07~16 마이그레이션(측정 기반~calibration transfer 실행) - apps/web: 세션 리뷰 카드·관리 화면·E2E 스펙 추가 - docs/ops: G0~G8 라이브 통합·배포·롤백 증거 문서와 evidence JSON/PNG - scripts: smoke·ledger·릴리스 에이전트·NAS 프리뷰 운영 스크립트 engine.public 로그 .bak과 apps/web/test-results 산출물은 커밋에서 제외했다.
271 lines
9.8 KiB
Python
271 lines
9.8 KiB
Python
"""Metadata-only high-water contracts for the G7 voice runtime."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import unittest
|
|
from pathlib import Path
|
|
from unittest.mock import AsyncMock, patch
|
|
|
|
from .routes import admin as admin_routes
|
|
from .routes import voice as voice_routes
|
|
from .services import voice as voice_service_module
|
|
from .services.voice import DeepgramStreamingSession
|
|
from .services.voice_runtime import (
|
|
VOICE_AUDIO_BUFFER_MAX_BYTES,
|
|
VOICE_STREAMING_EVENT_QUEUE_MAX_ITEMS,
|
|
VOICE_UVICORN_WS_MAX_QUEUE,
|
|
VoiceRuntimeMetrics,
|
|
)
|
|
from .test_voice_ws import (
|
|
SESSION_ID,
|
|
VOICE_PRESET,
|
|
FakeWebSocket,
|
|
_binary,
|
|
_control,
|
|
_principal,
|
|
)
|
|
|
|
|
|
class VoiceRuntimeMetricsTests(unittest.TestCase):
|
|
def test_connection_buffer_and_queue_high_water_survive_cleanup(self) -> None:
|
|
metrics = VoiceRuntimeMetrics()
|
|
first = metrics.websocket_opened()
|
|
second = metrics.websocket_opened()
|
|
|
|
metrics.audio_chunk_received(
|
|
first,
|
|
current_buffer_bytes=640,
|
|
chunk_bytes=640,
|
|
)
|
|
metrics.audio_chunk_received(
|
|
second,
|
|
current_buffer_bytes=384,
|
|
chunk_bytes=384,
|
|
)
|
|
metrics.streaming_queue_observed(first, queue_items=7)
|
|
metrics.streaming_queue_observed(
|
|
second,
|
|
queue_items=VOICE_STREAMING_EVENT_QUEUE_MAX_ITEMS,
|
|
saturated=True,
|
|
wait_seconds=0.125,
|
|
)
|
|
metrics.audio_overflow_rejected(second)
|
|
|
|
during = metrics.snapshot()
|
|
self.assertEqual(2, during.counters.active_websockets)
|
|
self.assertEqual(2, during.counters.websocket_high_water)
|
|
self.assertEqual(640, during.counters.route_audio_buffer_bytes)
|
|
self.assertEqual(1024, during.counters.route_audio_buffer_high_water_bytes)
|
|
self.assertEqual(1024, during.counters.audio_bytes_received_total)
|
|
self.assertEqual(1, during.counters.audio_overflow_rejections_total)
|
|
self.assertEqual(
|
|
7 + VOICE_STREAMING_EVENT_QUEUE_MAX_ITEMS,
|
|
during.counters.streaming_event_queue_high_water_items,
|
|
)
|
|
self.assertEqual(1, during.counters.streaming_event_queue_saturation_total)
|
|
self.assertEqual(
|
|
0.125,
|
|
during.counters.streaming_event_queue_wait_seconds_total,
|
|
)
|
|
|
|
metrics.audio_buffer_cleared(first)
|
|
metrics.streaming_queue_observed(first, queue_items=0)
|
|
metrics.streaming_queue_observed(second, queue_items=0)
|
|
metrics.websocket_closed(first)
|
|
metrics.websocket_closed(second)
|
|
after = metrics.snapshot()
|
|
self.assertEqual(0, after.counters.active_websockets)
|
|
self.assertEqual(0, after.counters.route_audio_buffer_bytes)
|
|
self.assertEqual(0, after.counters.streaming_event_queue_items)
|
|
self.assertEqual(2, after.counters.websocket_high_water)
|
|
self.assertFalse(after.reset_supported)
|
|
|
|
def test_streaming_provider_lifecycle_is_bounded_and_counted(self) -> None:
|
|
metrics = VoiceRuntimeMetrics()
|
|
metrics.streaming_provider_opened()
|
|
metrics.streaming_provider_opened()
|
|
metrics.streaming_provider_closed(outcome="finalized")
|
|
metrics.streaming_provider_closed(outcome="aborted")
|
|
metrics.provider_fallback()
|
|
metrics.websocket_error()
|
|
|
|
snapshot = metrics.snapshot()
|
|
self.assertEqual(0, snapshot.counters.active_streaming_provider_sessions)
|
|
self.assertEqual(2, snapshot.counters.streaming_provider_session_high_water)
|
|
self.assertEqual(2, snapshot.counters.streaming_provider_sessions_opened_total)
|
|
self.assertEqual(1, snapshot.counters.provider_finalize_total)
|
|
self.assertEqual(1, snapshot.counters.provider_abort_total)
|
|
self.assertEqual(1, snapshot.counters.provider_fallback_total)
|
|
self.assertEqual(1, snapshot.counters.websocket_error_total)
|
|
self.assertEqual("single_api_worker", snapshot.scope)
|
|
self.assertEqual(
|
|
"metadata_only_no_audio_transcript_or_session_ids",
|
|
snapshot.privacy_boundary,
|
|
)
|
|
|
|
def test_runtime_limits_match_the_production_route_and_docker_cmd(self) -> None:
|
|
snapshot = VoiceRuntimeMetrics().snapshot()
|
|
self.assertEqual(
|
|
VOICE_AUDIO_BUFFER_MAX_BYTES,
|
|
snapshot.limits.max_utterance_audio_bytes,
|
|
)
|
|
self.assertEqual(
|
|
VOICE_STREAMING_EVENT_QUEUE_MAX_ITEMS,
|
|
snapshot.limits.streaming_event_queue_max_items,
|
|
)
|
|
self.assertEqual(
|
|
VOICE_UVICORN_WS_MAX_QUEUE,
|
|
snapshot.limits.uvicorn_ws_max_queue,
|
|
)
|
|
dockerfile = (Path(__file__).resolve().parents[1] / "Dockerfile").read_text(
|
|
encoding="utf-8"
|
|
)
|
|
self.assertIn(
|
|
'"--ws", "websockets", "--ws-max-queue", "4"',
|
|
dockerfile,
|
|
)
|
|
|
|
def test_admin_snapshot_is_structured_and_metadata_only(self) -> None:
|
|
response = asyncio.run(admin_routes.admin_voice_runtime(principal=None)) # type: ignore[arg-type]
|
|
payload = response.model_dump(mode="json")
|
|
self.assertEqual("vignette.voice-runtime.v1", payload["schema_version"])
|
|
keys: set[str] = set()
|
|
|
|
def collect_keys(value: object) -> None:
|
|
if isinstance(value, dict):
|
|
keys.update(str(key) for key in value)
|
|
for child in value.values():
|
|
collect_keys(child)
|
|
elif isinstance(value, list):
|
|
for child in value:
|
|
collect_keys(child)
|
|
|
|
collect_keys(payload)
|
|
for forbidden in (
|
|
"session_id",
|
|
"transcript_text",
|
|
"raw_audio",
|
|
"provider_payload",
|
|
):
|
|
self.assertNotIn(forbidden, keys)
|
|
|
|
|
|
class _StreamingSocket:
|
|
def __init__(self) -> None:
|
|
self._frames = iter(
|
|
[
|
|
'{"type":"Results","is_final":true,"speech_final":true,'
|
|
'"channel":{"alternatives":[{"transcript":"ok","words":[]}]}}'
|
|
]
|
|
)
|
|
self.sent: list[object] = []
|
|
self.close = AsyncMock()
|
|
|
|
def __aiter__(self):
|
|
return self
|
|
|
|
async def __anext__(self) -> str:
|
|
try:
|
|
return next(self._frames)
|
|
except StopIteration as exc:
|
|
raise StopAsyncIteration from exc
|
|
|
|
async def send(self, payload: object) -> None:
|
|
self.sent.append(payload)
|
|
|
|
|
|
class DeepgramRuntimeLifecycleTests(unittest.IsolatedAsyncioTestCase):
|
|
async def test_finalize_and_abort_release_global_active_session(self) -> None:
|
|
baseline = voice_service_module.voice_runtime_metrics.snapshot().counters
|
|
finalized = DeepgramStreamingSession(
|
|
_StreamingSocket(),
|
|
model="nova-3",
|
|
language="ko",
|
|
on_event=AsyncMock(),
|
|
keepalive_seconds=60,
|
|
finalize_timeout_seconds=1,
|
|
)
|
|
await finalized.finish()
|
|
after_finalize = voice_service_module.voice_runtime_metrics.snapshot().counters
|
|
self.assertEqual(
|
|
baseline.active_streaming_provider_sessions,
|
|
after_finalize.active_streaming_provider_sessions,
|
|
)
|
|
self.assertEqual(
|
|
baseline.provider_finalize_total + 1,
|
|
after_finalize.provider_finalize_total,
|
|
)
|
|
|
|
aborted = DeepgramStreamingSession(
|
|
_StreamingSocket(),
|
|
model="nova-3",
|
|
language="ko",
|
|
on_event=AsyncMock(),
|
|
keepalive_seconds=60,
|
|
finalize_timeout_seconds=1,
|
|
)
|
|
await aborted.abort()
|
|
after_abort = voice_service_module.voice_runtime_metrics.snapshot().counters
|
|
self.assertEqual(
|
|
baseline.active_streaming_provider_sessions,
|
|
after_abort.active_streaming_provider_sessions,
|
|
)
|
|
self.assertEqual(
|
|
after_finalize.provider_abort_total + 1,
|
|
after_abort.provider_abort_total,
|
|
)
|
|
|
|
|
|
class VoiceRouteRuntimeIntegrationTests(unittest.IsolatedAsyncioTestCase):
|
|
async def test_route_records_audio_high_water_and_cleans_connection(self) -> None:
|
|
metrics = VoiceRuntimeMetrics()
|
|
websocket = FakeWebSocket(
|
|
[
|
|
_control({"type": "audio_start", "format": "webm"}),
|
|
_binary(b"first"),
|
|
_binary(b"second"),
|
|
_control({"type": "audio_end", "format": "webm"}),
|
|
_control({"type": "close"}),
|
|
]
|
|
)
|
|
|
|
with (
|
|
patch.object(voice_routes, "voice_runtime_metrics", metrics),
|
|
patch.object(
|
|
voice_routes,
|
|
"_principal_from_websocket",
|
|
AsyncMock(return_value=_principal()),
|
|
),
|
|
patch.object(
|
|
voice_routes,
|
|
"_bind_session",
|
|
AsyncMock(
|
|
return_value=(SESSION_ID, VOICE_PRESET, None, {})
|
|
),
|
|
),
|
|
patch.object(
|
|
voice_routes.voice_service,
|
|
"is_available",
|
|
return_value=True,
|
|
),
|
|
patch.object(
|
|
voice_routes.voice_service,
|
|
"can_stream_audio",
|
|
return_value=False,
|
|
),
|
|
patch.object(voice_routes, "_handle_utterance", AsyncMock()),
|
|
):
|
|
await voice_routes.voice_ws(websocket) # type: ignore[arg-type]
|
|
|
|
snapshot = metrics.snapshot().counters
|
|
self.assertEqual(1, snapshot.websockets_opened_total)
|
|
self.assertEqual(1, snapshot.websocket_high_water)
|
|
self.assertEqual(0, snapshot.active_websockets)
|
|
self.assertEqual(11, snapshot.audio_bytes_received_total)
|
|
self.assertEqual(11, snapshot.route_audio_buffer_high_water_bytes)
|
|
self.assertEqual(0, snapshot.route_audio_buffer_bytes)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|