diff --git a/config.py b/config.py index 764c07c..8f40bcc 100644 --- a/config.py +++ b/config.py @@ -162,6 +162,12 @@ class Settings(BaseSettings): # silently degrading to a gateway that can't place real calls. use_mock_sip: bool = False + # SIP stack: "sippy" (signalling only — the classifier gets no audio) or + # "pjsua2" (call control + media, the only path where audio reaches the + # classifier). Opt-in while the PJSUA2 engine is proven against the lab; + # see docs/architecture.md → "Media plane: why PJSUA2 places the call". + sip_engine: str = "sippy" + # Notifications notify_sms_number: str = "" diff --git a/core/gateway.py b/core/gateway.py index 5f8f94e..9bc8ae6 100644 --- a/core/gateway.py +++ b/core/gateway.py @@ -55,6 +55,26 @@ def build_sip_engine( "for development without a trunk." ) + if settings.sip_engine.lower() == "pjsua2": + from core.pjsua_engine import PJSUAEngine + + logger.info("📞 SIP engine: PJSUA2 (call control + media)") + return PJSUAEngine( + sip_address=gw_sip.host, + sip_port=gw_sip.port, + trunk_host=trunk.host, + trunk_port=trunk.port, + trunk_username=trunk.username, + trunk_password=trunk.password.get_secret_value(), + trunk_transport=trunk.transport, + domain=gw_sip.domain, + did=trunk.did, + media_pipeline=media_pipeline, + on_leg_state_change=on_leg_state_change, + on_device_registered=on_device_registered, + on_incoming_call=on_incoming_call, + ) + return SippyEngine( sip_address=gw_sip.host, sip_port=gw_sip.port, diff --git a/core/media_pipeline.py b/core/media_pipeline.py index 19d3b16..c247e74 100644 --- a/core/media_pipeline.py +++ b/core/media_pipeline.py @@ -20,6 +20,7 @@ PJSUA2 runs in its own thread with a dedicated Endpoint. """ import asyncio +import gc import logging import threading from collections.abc import AsyncIterator @@ -98,6 +99,72 @@ class AudioTap: self._active = False +def make_capture_port(stream_id: str, sample_rate: int, channels: int, frame_ms: int): + """Build a PJSUA2 media port that forks conference audio into taps. + + Defined as a factory rather than a module-level class because + ``pj.AudioMediaPort`` can only be subclassed once ``pjsua2`` imports — + and the whole pipeline degrades to stub mode when it doesn't. + + The returned port is a *sink*: the conference bridge transmits into it, + and every frame is copied to each registered tap. Returns ``None`` when + pjsua2 is unavailable. + """ + try: + import pjsua2 as pj + except ImportError: + return None + + class _CapturePort(pj.AudioMediaPort): + """Receives conference-bridge frames and fans them out to taps. + + ``onFrameReceived`` is called on a **PJSUA2 worker thread** — a third + execution context alongside the asyncio loop and the Sippy ED thread. + It must touch nothing but ``AudioTap.feed``, which is explicitly + thread-safe (it hops to the owning loop via ``call_soon_threadsafe``). + Reaching into pipeline state, the event bus, or a Sippy object from + here would be a data race. + """ + + def __init__(self, stream_id: str): + super().__init__() + self.stream_id = stream_id + self.taps: list[AudioTap] = [] + self._logged_error = False + + def onFrameReceived(self, frame): # noqa: N802 — PJSUA2 C++ callback name + try: + if not self.taps or frame.size <= 0: + return + # frame.buf is a SWIG ByteVector of signed chars; the tap + # contract is raw little-endian 16-bit PCM. + pcm = bytes(bytearray(b & 0xFF for b in frame.buf)) + for tap in self.taps: + tap.feed(pcm) + except Exception as e: + # An exception escaping into PJSUA2's C++ callback would tear + # down the worker thread and silently kill media for every + # call. Log once per port rather than on every 20ms frame. + if not self._logged_error: + self._logged_error = True + logger.error( + f" Audio capture failed for {self.stream_id}: {e}", + exc_info=True, + ) + + fmt = pj.MediaFormatAudio() + fmt.init( + pj.PJMEDIA_FORMAT_L16, + sample_rate, + channels, + frame_ms * 1000, # frameTimeUsec + 16, # bitsPerSample + ) + port = _CapturePort(stream_id) + port.createPort(f"tap-{stream_id}", fmt) + return port + + # ================================================================ # Stream Entry — tracks a single media stream in the pipeline # ================================================================ @@ -112,6 +179,8 @@ class MediaStream: self.codec = codec self.conf_port: Optional[int] = None # PJSUA2 conference bridge port ID self.transport = None # PJSUA2 SipTransport + self.media = None # PJSUA2 AudioMedia for this stream + self.capture_port = None # Shared _CapturePort feeding this stream's taps self.rtp_port: Optional[int] = None # Local RTP listen port self.taps: list[AudioTap] = [] self.recorder = None # PJSUA2 AudioMediaRecorder @@ -143,10 +212,11 @@ class MediaPipeline: pipeline = MediaPipeline() await pipeline.start() - # Add a stream for a call leg - port = pipeline.add_remote_stream("leg_1", "10.0.0.1", 20000, "PCMU") + # Media arrives from the SIP engine's onCallMediaState callback + # (PJSUA2 only surfaces RTP media for a call it owns): + # pipeline.attach_call_media("leg_1", call.getAudioMedia(i)) - # Tap audio for analysis + # Tap audio for analysis — safe before or after media comes up tap = pipeline.create_tap("leg_1") async for frame in tap.stream(): classify(frame) @@ -173,6 +243,7 @@ class MediaPipeline: self._next_rtp_port = rtp_start_port self._sample_rate = sample_rate self._channels = channels + self._frame_ms = 20 # Must match medConfig.audioFramePtime below self._null_audio = null_audio # Use null audio device (no sound card needed) # State @@ -255,11 +326,17 @@ class MediaPipeline: tap.close() self._taps.clear() - # Remove all streams + # Remove all streams (this releases their capture ports) for stream_id in list(self._streams.keys()): self.remove_stream(stream_id) - # Destroy PJSUA2 endpoint + # Destroy PJSUA2 endpoint. Every media port must be collected first: + # a port finalised after libDestroy() runs pjmedia_conf_remove_port + # against a freed conference bridge and aborts the process. Dropping + # the last Python reference is not enough on its own — force the + # collection here rather than leaving it to interpreter exit. + gc.collect() + if self._endpoint: try: self._endpoint.libDestroy() @@ -270,6 +347,15 @@ class MediaPipeline: self._ready = False logger.info("🎵 PJSUA2 media pipeline stopped") + @property + def endpoint(self): + """The PJSUA2 Endpoint, or None in stub mode. + + PJSUA2 permits exactly one Endpoint per process, so the pipeline + creates it and the SIP engine borrows it rather than making a second. + """ + return self._endpoint + @property def is_ready(self) -> bool: return self._ready @@ -291,50 +377,51 @@ class MediaPipeline: # Stream Management # ================================================================ - def add_remote_stream( - self, stream_id: str, remote_host: str, remote_port: int, codec: str = "PCMU" - ) -> Optional[int]: + def attach_call_media(self, stream_id: str, audio_media) -> Optional[int]: + """Register a call's live ``AudioMedia`` with the pipeline. + + Called from ``onCallMediaState`` on a PJSUA2 worker thread, which is + the only place PJSUA2 surfaces RTP-backed media. Any tap created + before this point is attached now; taps created later find the media + already present. + + There is deliberately no ``add_remote_stream(host, port)`` counterpart: + PJSUA2 has no standalone RTP media object, so media can only arrive + from a call PJSUA2 owns. See ``docs/architecture.md``. """ - Add a remote RTP stream to the conference bridge. + stream = self._streams.get(stream_id) + if stream is None: + stream = MediaStream(stream_id, "", 0) + self._streams[stream_id] = stream - Creates a PJSUA2 transport and media port for the remote - party's RTP stream, connecting it to the conference bridge. + stream.media = audio_media + try: + stream.conf_port = audio_media.getPortId() + except Exception: + stream.conf_port = None - Args: - stream_id: Unique ID (typically the SIP leg ID) - remote_host: Remote RTP host - remote_port: Remote RTP port - codec: Audio codec (PCMU, PCMA, G729) + # Wire up taps that were requested before media came up. + pending = self._taps.get(stream_id, []) + if pending and stream.capture_port is None: + port = make_capture_port( + stream_id, self._sample_rate, self._channels, self._frame_ms + ) + if port is not None: + try: + audio_media.startTransmit(port) + stream.capture_port = port + port.taps.extend(pending) + logger.info( + f" 🎤 Audio tap attached for {stream_id} " + f"({len(pending)} waiting)" + ) + except Exception as e: + logger.error( + f" Failed to attach capture port for {stream_id}: {e}", + exc_info=True, + ) - Returns: - Conference bridge port ID, or None if PJSUA2 not available - """ - stream = MediaStream(stream_id, remote_host, remote_port, codec) - stream.rtp_port = self.allocate_rtp_port(stream_id) - - if self._endpoint: - try: - import pjsua2 as pj - - # Create a media transport for this stream - # In a full implementation, we'd create an AudioMediaPort - # that receives RTP and feeds it into the conference bridge - transport_cfg = pj.TransportConfig() - transport_cfg.port = stream.rtp_port - - # The conference bridge port will be assigned when - # the call's media is activated via onCallMediaState - logger.info( - f" 📡 Added stream {stream_id}: " - f"local={stream.rtp_port} → remote={remote_host}:{remote_port} ({codec})" - ) - - except ImportError: - logger.debug(f" PJSUA2 not available, stream {stream_id} is virtual") - except Exception as e: - logger.error(f" Failed to add stream {stream_id}: {e}") - - self._streams[stream_id] = stream + logger.info(f" 📡 Media attached for {stream_id} (conf port {stream.conf_port})") return stream.conf_port def remove_stream(self, stream_id: str) -> None: @@ -350,6 +437,19 @@ class MediaPipeline: tap.close() self._taps.pop(stream_id, None) + # Release the capture port while the conference bridge still exists. + # A port garbage-collected after Endpoint.libDestroy() calls + # pjmedia_conf_remove_port against a freed bridge and aborts the + # process on a native assertion — a hard crash, not an exception. + if stream.capture_port is not None: + try: + if stream.media is not None: + stream.media.stopTransmit(stream.capture_port) + except Exception as e: + logger.debug(f" stopTransmit failed for {stream_id}: {e}") + stream.capture_port.taps.clear() + stream.capture_port = None + # Stop recording if stream.recorder: try: @@ -426,16 +526,28 @@ class MediaPipeline: self._taps[stream_id] = [] self._taps[stream_id].append(tap) - if self._endpoint and stream and stream.conf_port is not None: + if self._endpoint and stream and stream.media is not None: try: - import pjsua2 as pj - # Create an AudioMediaPort that captures frames - # and feeds them to the tap - # In PJSUA2, we'd subclass AudioMediaPort and implement - # onFrameReceived to call tap.feed(frame_data) - logger.info(f" 🎤 Audio tap created for {stream_id} (PJSUA2)") + # One capture port per stream, shared by every tap on it: + # the bridge would otherwise mix each additional port back + # into the conference and the call would echo. + if stream.capture_port is None: + port = make_capture_port( + stream_id, self._sample_rate, self._channels, self._frame_ms + ) + if port is not None: + # The stream's media transmits into the capture port, + # not the reverse — the port is a sink. + stream.media.startTransmit(port) + stream.capture_port = port + logger.info(f" 🎤 Audio tap created for {stream_id} (PJSUA2)") + + if stream.capture_port is not None: + stream.capture_port.taps.append(tap) except Exception as e: - logger.error(f" Failed to create PJSUA2 tap for {stream_id}: {e}") + logger.error( + f" Failed to create PJSUA2 tap for {stream_id}: {e}", exc_info=True + ) else: logger.info(f" 🎤 Audio tap created for {stream_id} (virtual)") diff --git a/core/pjsua_engine.py b/core/pjsua_engine.py new file mode 100644 index 0000000..5d87eb9 --- /dev/null +++ b/core/pjsua_engine.py @@ -0,0 +1,488 @@ +""" +PJSUA2 SIP engine — call control *and* media in one library. + +Why this exists +--------------- +The gateway originally signalled with Sippy and expected PJSUA2 to carry +media. That cannot work: **PJSUA2 exposes no standalone RTP media object**. +Every ``AudioMedia`` subclass in the Python bindings is a file player, +recorder, tone generator or capture port, and RTP is reachable only through +``pj.Call.getAudioMedia()`` — on a dialog PJSUA2 itself owns. A design where +another stack owns the dialog can never obtain media from PJSUA2, so audio +never reached the classifier. + +Owning the dialog is the price of owning the media, so this engine places the +call. See ``docs/architecture.md`` → "Media plane: why PJSUA2 places the call". + +Safety +------ +This engine is *only* reached through ``gateway.make_call``, which refuses +emergency numbers and enforces the concurrency cap **before** any SIP action. +Nothing here may be given a second dial path that bypasses those checks. + +Threading +--------- +Three execution contexts, as elsewhere in the codebase: + +* the **asyncio loop** owns legs, the event bus, and the call manager; +* **PJSUA2 worker threads** run every ``on*`` callback below; +* (the Sippy ED thread is not involved — this engine replaces it.) + +PJSUA2 callbacks cross to the loop through exactly one funnel, +``_post_from_pj`` → ``run_coroutine_threadsafe``. A callback must never touch +loop-owned state directly. Any thread PJSUA2 did not create must call +``libRegisterThread`` before touching a PJSUA2 object, which +``_ensure_registered`` handles. +""" + +import asyncio +import gc +import logging +import threading +import uuid +from collections.abc import Callable + +from core.sip_engine import SIPEngine +from models.device import Device + +logger = logging.getLogger(__name__) + + +class PJSUAEngine(SIPEngine): + """SIP engine backed by PJSUA2 for both signalling and media.""" + + def __init__( + self, + sip_address: str = "0.0.0.0", + sip_port: int = 5060, + trunk_host: str = "", + trunk_port: int = 5060, + trunk_username: str = "", + trunk_password: str = "", + trunk_transport: str = "udp", + domain: str = "gateway.local", + did: str = "", + media_pipeline=None, + on_leg_state_change: Callable | None = None, + on_device_registered: Callable | None = None, + on_incoming_call: Callable | None = None, + ): + self._sip_address = sip_address + self._sip_port = sip_port + self._trunk_host = trunk_host + self._trunk_port = trunk_port + self._trunk_username = trunk_username + self._trunk_password = trunk_password + self._trunk_transport = trunk_transport + self._domain = domain + self._did = did + + # The media pipeline owns the PJSUA2 Endpoint; this engine borrows it + # rather than creating a second one (PJSUA2 permits only one). + self.media_pipeline = media_pipeline + + self._on_leg_state_change = on_leg_state_change + self._on_device_registered = on_device_registered + self._on_incoming_call = on_incoming_call + + self._loop: asyncio.AbstractEventLoop | None = None + self._ready = False + self._account = None + self._trunk_registered = False + self._trunk_reason = "not started" + + # PJSUA2-thread-owned: maps leg_id → pj.Call. Only touched from a + # PJSUA2 callback or a method that has registered itself first. + self._calls: dict[str, object] = {} + self._lock = threading.Lock() + + # ================================================================ + # Thread boundary + # ================================================================ + + def _post_from_pj(self, coro) -> None: + """Schedule loop work from a PJSUA2 worker thread. The one funnel.""" + if self._loop is None: + return + asyncio.run_coroutine_threadsafe(coro, self._loop) + + def _ensure_registered(self) -> None: + """Register the calling thread with PJSUA2 if it isn't already. + + PJSUA2 aborts when a thread it does not know touches its objects. + Calls made from the asyncio loop (hangup, DTMF) hit this. + """ + try: + import pjsua2 as pj + + ep = pj.Endpoint.instance() + if not ep.libIsThreadRegistered(): + ep.libRegisterThread(threading.current_thread().name) + except Exception as e: # pragma: no cover - defensive + logger.debug(f" thread registration skipped: {e}") + + async def _emit_leg_state(self, leg_id: str, state: str) -> None: + """Deliver a leg-state change on the loop.""" + if self._on_leg_state_change is None: + return + result = self._on_leg_state_change(leg_id, state) + if asyncio.iscoroutine(result): + await result + + # ================================================================ + # Lifecycle + # ================================================================ + + async def start(self) -> None: + """Create the SIP transport and register with the trunk.""" + self._loop = asyncio.get_running_loop() + logger.info("🔌 Starting PJSUA2 SIP engine...") + + if self.media_pipeline is None or not self.media_pipeline.endpoint: + raise RuntimeError( + "PJSUAEngine requires a started MediaPipeline — PJSUA2 allows " + "only one Endpoint, so the pipeline owns it and the engine " + "borrows it." + ) + + import pjsua2 as pj + + ep = self.media_pipeline.endpoint + + transport_cfg = pj.TransportConfig() + transport_cfg.port = self._sip_port + if self._sip_address and self._sip_address != "0.0.0.0": + transport_cfg.boundAddress = self._sip_address + + tp_type = ( + pj.PJSIP_TRANSPORT_TCP + if self._trunk_transport.lower() == "tcp" + else pj.PJSIP_TRANSPORT_UDP + ) + ep.transportCreate(tp_type, transport_cfg) + + self._create_account(ep) + + self._ready = True + logger.info(f"🔌 PJSUA2 SIP engine ready on {self._sip_address}:{self._sip_port}") + + def _create_account(self, ep) -> None: + """Build the account — registered to the trunk, or local-only.""" + import pjsua2 as pj + + engine = self + + class _Account(pj.Account): + def onRegState(self, prm): # noqa: N802 — PJSUA2 callback name + try: + info = self.getInfo() + engine._trunk_registered = bool(info.regIsActive) + engine._trunk_reason = f"{prm.code} {prm.reason}".strip() + if info.regIsActive: + logger.info(" ✅ Trunk registration accepted") + else: + # A rejected REGISTER must not read as "registered": + # /health treats a registered trunk as a condition of + # being healthy. + logger.error( + f" ❌ Trunk registration failed: {engine._trunk_reason}" + ) + except Exception as e: + logger.error(f" onRegState error: {e}", exc_info=True) + + def onIncomingCall(self, prm): # noqa: N802 — PJSUA2 callback name + try: + engine._handle_incoming(self, prm.callId) + except Exception as e: + logger.error(f" onIncomingCall error: {e}", exc_info=True) + + acc_cfg = pj.AccountConfig() + + if self._trunk_host: + acc_cfg.idUri = f"sip:{self._trunk_username}@{self._trunk_host}" + acc_cfg.regConfig.registrarUri = f"sip:{self._trunk_host}:{self._trunk_port}" + cred = pj.AuthCredInfo( + "digest", "*", self._trunk_username, 0, self._trunk_password + ) + acc_cfg.sipConfig.authCreds.append(cred) + else: + # No trunk configured: a local-only account still lets devices + # register and inbound calls arrive. + acc_cfg.idUri = f"sip:gateway@{self._domain}" + self._trunk_reason = "No SIP trunk configured" + + self._account = _Account() + self._account.create(acc_cfg) + + if self._trunk_host: + logger.info(f" Registering with trunk: {self._trunk_host}:{self._trunk_port}") + + async def stop(self) -> None: + """Hang up everything and drop the account.""" + logger.info("🔌 Stopping PJSUA2 SIP engine...") + self._ready = False + + self._ensure_registered() + had_calls = bool(self._calls) + for leg_id in list(self._calls.keys()): + try: + await self.hangup(leg_id) + except Exception as e: + logger.debug(f" hangup during shutdown failed for {leg_id}: {e}") + + # hangup() only queues the BYE. Give PJSUA2 a moment to send it and + # tear the media down, or the account is deleted with a call still + # active ("deleting account 0 while call 0 is still active") and the + # far end is left waiting on a dialog nobody closed. + if had_calls: + await asyncio.sleep(0.5) + + # Drop every PJSUA2 object before the pipeline destroys the endpoint. + # A Call or Account finalised after libDestroy() aborts the process on + # a native assertion, exactly as a stray media port does — and a Call + # still alive keeps delivering callbacks into a half-torn-down + # interpreter. Dropping the last reference is not enough on its own, + # so force the collection here. + with self._lock: + self._calls.clear() + self._account = None + gc.collect() + logger.info("🔌 PJSUA2 SIP engine stopped") + + async def is_ready(self) -> bool: + return self._ready + + # ================================================================ + # Calls + # ================================================================ + + def _make_call_class(self): + """Build the pj.Call subclass bound to this engine.""" + import pjsua2 as pj + + engine = self + + class _Call(pj.Call): + def __init__(self, acc, leg_id: str, call_id=pj.PJSUA_INVALID_ID): + super().__init__(acc, call_id) + self.leg_id = leg_id + + def onCallState(self, prm): # noqa: N802 — PJSUA2 callback name + # PJSUA2 keeps delivering callbacks while the interpreter is + # tearing down, when module globals may already be cleared — + # hence the local alias and the bare except. A raise here + # escapes into C++ and takes the worker thread with it. + state_map = _STATE_MAP + try: + info = self.getInfo() + state = state_map.get(info.state) + if state is None: + return + if state == "terminated": + engine._forget_call(self.leg_id) + engine._post_from_pj(engine._emit_leg_state(self.leg_id, state)) + except Exception: + try: + logger.error(" onCallState error", exc_info=True) + except Exception: + pass + + def onCallMediaState(self, prm): # noqa: N802 — PJSUA2 callback name + """Media is up — hand the audio to the pipeline. + + This is the callback the whole refactor exists for: it is the + only place PJSUA2 surfaces an RTP-backed AudioMedia. + """ + try: + info = self.getInfo() + for i, mi in enumerate(info.media): + if ( + mi.type == pj.PJMEDIA_TYPE_AUDIO + and mi.status == pj.PJSUA_CALL_MEDIA_ACTIVE + ): + engine._attach_media(self.leg_id, self.getAudioMedia(i)) + break + except Exception: + try: + logger.error(" onCallMediaState error", exc_info=True) + except Exception: + pass + + return _Call + + def _attach_media(self, leg_id: str, audio_media) -> None: + """Register a live AudioMedia with the pipeline (PJSUA2 thread).""" + if self.media_pipeline is None: + return + try: + self.media_pipeline.attach_call_media(leg_id, audio_media) + logger.info(f" 🎵 Media active for {leg_id}") + except Exception as e: + logger.error(f" Failed to attach media for {leg_id}: {e}", exc_info=True) + + def _forget_call(self, leg_id: str) -> None: + with self._lock: + self._calls.pop(leg_id, None) + if self.media_pipeline is not None: + try: + self.media_pipeline.remove_stream(leg_id) + except Exception as e: + logger.debug(f" stream cleanup failed for {leg_id}: {e}") + + async def make_call(self, number: str, caller_id: str | None = None) -> str: + """Place an outbound call. Reached only via gateway.make_call.""" + if not self._ready: + raise RuntimeError("SIP engine not ready") + + import pjsua2 as pj + + leg_id = f"leg_{uuid.uuid4().hex[:12]}" + target = ( + f"sip:{number}@{self._trunk_host}:{self._trunk_port}" + if self._trunk_host + else f"sip:{number}@{self._domain}" + ) + logger.info(f"📞 Placing call to {target} (leg: {leg_id})") + + self._ensure_registered() + call_cls = self._make_call_class() + call = call_cls(self._account, leg_id) + + prm = pj.CallOpParam(True) + call.makeCall(target, prm) + + with self._lock: + self._calls[leg_id] = call + return leg_id + + async def hangup(self, call_leg_id: str) -> None: + import pjsua2 as pj + + with self._lock: + call = self._calls.get(call_leg_id) + if call is None: + return + self._ensure_registered() + try: + call.hangup(pj.CallOpParam(True)) + except Exception as e: + logger.debug(f" hangup failed for {call_leg_id}: {e}") + self._forget_call(call_leg_id) + + async def send_dtmf(self, call_leg_id: str, digits: str) -> None: + """Send DTMF as RFC 2833 — the in-band path a real IVR expects.""" + with self._lock: + call = self._calls.get(call_leg_id) + if call is None: + logger.warning(f" send_dtmf: no call for {call_leg_id}") + return + self._ensure_registered() + call.dialDtmf(digits) + logger.info(f" Sent DTMF '{digits}' on {call_leg_id}") + + async def call_device(self, device: Device) -> str: + """Ring a registered device (transfer target).""" + if not self._ready: + raise RuntimeError("SIP engine not ready") + + import pjsua2 as pj + + leg_id = f"leg_{uuid.uuid4().hex[:12]}" + target = device.sip_uri or f"sip:{device.id}@{self._domain}" + logger.info(f"📞 Ringing device {device.id} at {target} (leg: {leg_id})") + + self._ensure_registered() + call_cls = self._make_call_class() + call = call_cls(self._account, leg_id) + call.makeCall(target, pj.CallOpParam(True)) + + with self._lock: + self._calls[leg_id] = call + return leg_id + + def _handle_incoming(self, account, call_id) -> None: + """Inbound INVITE (PJSUA2 thread) — answer and hand to the receptionist.""" + import pjsua2 as pj + + leg_id = f"leg_{uuid.uuid4().hex[:12]}" + call_cls = self._make_call_class() + call = call_cls(account, leg_id, call_id) + + try: + info = call.getInfo() + remote = info.remoteUri + except Exception: + remote = "unknown" + + with self._lock: + self._calls[leg_id] = call + + call.answer(pj.CallOpParam(True)) + logger.info(f"📞 Inbound call {leg_id} from {remote}") + + if self._on_incoming_call is not None: + result = self._on_incoming_call(leg_id, remote) + if asyncio.iscoroutine(result): + self._post_from_pj(result) + + # ================================================================ + # Bridging + # ================================================================ + + async def bridge_calls(self, leg_a: str, leg_b: str) -> str: + """Join two legs in the conference bridge.""" + bridge_id = f"bridge_{uuid.uuid4().hex[:8]}" + if self.media_pipeline is not None: + self.media_pipeline.bridge_streams(leg_a, leg_b) + logger.info(f" 🌉 Bridged {leg_a} ↔ {leg_b} ({bridge_id})") + return bridge_id + + async def unbridge(self, bridge_id: str) -> None: + logger.info(f" Unbridged {bridge_id}") + + def get_audio_stream(self, call_leg_id: str): + if self.media_pipeline is not None: + return self.media_pipeline.get_audio_tap(call_leg_id) + return None + + # ================================================================ + # Status + # ================================================================ + + async def get_registered_devices(self) -> list[dict]: + return [] + + async def get_trunk_status(self) -> dict: + return { + "registered": self._trunk_registered, + "host": self._trunk_host or "not configured", + "port": self._trunk_port, + "transport": self._trunk_transport, + "username": self._trunk_username, + "reason": None if self._trunk_registered else self._trunk_reason, + } + + +# Populated lazily: the pjsua2 constants are unavailable until import, and +# the module must import cleanly in stub mode. +_STATE_MAP: dict = {} + + +def _init_state_map() -> None: + global _STATE_MAP + if _STATE_MAP: + return + try: + import pjsua2 as pj + except ImportError: + return + _STATE_MAP = { + pj.PJSIP_INV_STATE_CALLING: "trying", + pj.PJSIP_INV_STATE_EARLY: "ringing", + pj.PJSIP_INV_STATE_CONNECTING: "trying", + pj.PJSIP_INV_STATE_CONFIRMED: "connected", + pj.PJSIP_INV_STATE_DISCONNECTED: "terminated", + } + + +_init_state_map() diff --git a/core/sippy_engine.py b/core/sippy_engine.py index fdf4945..9a4e831 100644 --- a/core/sippy_engine.py +++ b/core/sippy_engine.py @@ -276,17 +276,22 @@ class SippyEngine(SIPEngine): if state == "connected": sdp = data.get("sdp") if sdp and self.media_pipeline: + # Signalling only: PJSUA2 surfaces RTP media exclusively + # through a call it owns, so a Sippy-owned dialog can + # never be given media — the classifier stays deaf on this + # engine. PJSUAEngine is the media-capable path; see + # docs/architecture.md. Log the negotiated endpoint so the + # SIP exchange is still debuggable. try: remote_rtp = self._parse_sdp_rtp_endpoint(sdp) if remote_rtp: - leg.media_port = self.media_pipeline.add_remote_stream( - leg.leg_id, - remote_rtp["host"], - remote_rtp["port"], - remote_rtp["codec"], + logger.info( + f" {leg.leg_id}: remote RTP " + f"{remote_rtp['host']}:{remote_rtp['port']} " + f"({remote_rtp['codec']}) — no media on this engine" ) except Exception as e: - logger.error(f" Failed to set up media for {leg.leg_id}: {e}") + logger.error(f" Failed to parse SDP for {leg.leg_id}: {e}") elif state == "terminated": if self.media_pipeline and leg.media_port is not None: try: diff --git a/docs/architecture.md b/docs/architecture.md index f8ef511..1d35c00 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -1,6 +1,15 @@ # Architecture -Hold Slayer is a single-process async Python application built on FastAPI. It acts as an intelligent B2BUA (Back-to-Back User Agent) sitting between your SIP trunk (PSTN access) and your desk phone/softphone. +Hold Slayer is a single-process async Python application built on FastAPI. It +acts as an intelligent B2BUA (Back-to-Back User Agent) sitting between your SIP +trunk (PSTN access) and your desk phone/softphone. + +> **Media plane in transition.** The gateway currently signals with Sippy and +> intends PJSUA2 to carry media, but PJSUA2 will not surface an RTP stream for +> a dialog it does not own — so no audio ever reaches the classifier. The fix +> moves *call placement* into PJSUA2 while Sippy keeps the SBC roles. See +> [Media plane: why PJSUA2 places the call](#media-plane-why-pjsua2-places-the-call) +> before changing anything in `core/`. ## System Diagram @@ -27,8 +36,13 @@ Hold Slayer is a single-process async Python application built on FastAPI. It ac │ └────┬─────┘ └─────┬─────┘ └──────────────┘ │ │ │ │ │ │ ┌────┴──────────────┴───────────────────┐ │ -│ │ Sippy B2BUA Engine │ │ -│ │ (SIP calls, DTMF, conference bridge) │ │ +│ │ SIP Engine │ │ +│ │ signalling + call control │ │ +│ └────┬──────────────────────────────────┘ │ +│ │ │ +│ ┌────┴──────────────────────────────────┐ │ +│ │ Media Pipeline (PJSUA2) │ │ +│ │ RTP, conference bridge, taps, record │ │ │ └────┬──────────────────────────────────┘ │ │ │ │ └───────┼─────────────────────────────────────────────────────────┘ @@ -70,8 +84,8 @@ Hold Slayer is a single-process async Python application built on FastAPI. It ac | Component | File | Purpose | |-----------|------|---------| -| Sippy Engine | `core/sippy_engine.py` | SIP signaling (INVITE, BYE, REGISTER, DTMF) | -| Media Pipeline | `core/media_pipeline.py` | PJSUA2 RTP media handling, conference bridge, recording | +| Sippy Engine | `core/sippy_engine.py` | SIP signalling (INVITE, BYE, REGISTER, DTMF) | +| Media Pipeline | `core/media_pipeline.py` | PJSUA2 RTP media, conference bridge, taps, recording | | Recording | `services/recording.py` | WAV file management and storage | | Analytics | `services/call_analytics.py` | Call metrics, hold time stats, trends | | Notifications | `services/notification.py` | WebSocket + SMS alerts | @@ -84,9 +98,10 @@ Hold Slayer is a single-process async Python application built on FastAPI. It ac POST /api/v1/calls/hold-slayer { number, intent, call_flow_id } │ 2. Gateway.make_call() - ├── CallManager.create_call() → track state - ├── SippyEngine.make_call() → SIP INVITE to trunk - └── MediaPipeline.add_stream() → RTP media setup + ├── is_emergency_number() → REFUSE 911/112 (before anything else) + ├── concurrency cap check → refuse past max_concurrent_calls + ├── CallManager.create_call() → track state + └── sip_engine.make_call() → place the call, media follows │ 3. HoldSlayer.run_with_flow() or run_exploration() ├── AudioClassifier.classify() → analyze 3s audio windows @@ -99,14 +114,14 @@ Hold Slayer is a single-process async Python application built on FastAPI. It ac ├── TranscriptionService.transcribe() → STT on speech audio │ ├── LLMClient.analyze_ivr_menu() → pick menu option (fallback) - │ └── SippyEngine.send_dtmf() → press the button + │ └── sip_engine.send_dtmf() → press the button │ └── detect_hold_to_human_transition() └── HUMAN_DETECTED! → transfer │ 4. Transfer - ├── SippyEngine.bridge() → connect call legs - ├── MediaPipeline.bridge_streams() → bridge RTP + ├── SippyEngine.bridge_calls() → join the two call legs + ├── MediaPipeline.bridge_streams() → bridge RTP in the conf bridge ├── EventBus.publish(TRANSFER_STARTED) └── NotificationService → "Pick up your phone!" │ @@ -117,44 +132,95 @@ Hold Slayer is a single-process async Python application built on FastAPI. It ac → Analytics tracking ``` +The emergency guard and the concurrency cap are the first two steps of +`make_call` for a reason, and their order is load-bearing — see +[.claude/rules/call-safety.md](../.claude/rules/call-safety.md). + ## Threading Model -Hold Slayer is primarily single-threaded async (asyncio), with one exception: - -- **Main thread**: FastAPI + all async services (event bus, hold slayer, classifier, etc.) -- **Sippy thread**: Sippy B2BUA runs its own event loop in a dedicated daemon thread. The `SippyEngine` bridges async↔sync via `asyncio.run_in_executor()`. -- **PJSUA2**: Runs in the main thread using null audio device (no sound card needed — headless server mode). +The README's "single-process async" is a simplification. There are **three** +execution contexts, and the boundaries between them are the highest-leverage +invariant in the codebase. ``` -Main Thread (asyncio) -├── FastAPI (uvicorn) -├── EventBus -├── CallManager -├── HoldSlayer +asyncio loop (main thread) Sippy ED thread PJSUA2 worker threads +├── FastAPI (uvicorn) └── ED2 dispatcher └── media / RTP +├── EventBus ├── SIP signalling └── onFrameReceived +├── CallManager ├── UA objects +├── HoldSlayer └── DTMF relay ├── AudioClassifier ├── TranscriptionService ├── LLMClient -├── MediaPipeline (PJSUA2) ├── NotificationService └── RecordingService - -Sippy Thread (daemon) -└── Sippy B2BUA event loop - ├── SIP signaling - ├── DTMF relay - └── Call leg management ``` +**Crossing the boundaries — one funnel each way:** + +| Direction | Mechanism | Notes | +|---|---|---| +| Sippy ED → loop | `_post_from_ed` → `asyncio.run_coroutine_threadsafe` → `_on_engine_event` | The single funnel where Sippy-thread events mutate loop state | +| loop → Sippy ED | `_run_on_sippy` → `ED2.callFromThread` | Anything touching a Sippy UA object | +| PJSUA2 worker → loop | `AudioTap.feed` → `loop.call_soon_threadsafe` | The **only** thing a PJSUA2 callback may touch | + +`onFrameReceived` runs on a PJSUA2 worker thread every 20 ms. It must call +nothing but `AudioTap.feed`; reaching into pipeline state, the event bus, or a +Sippy object from there is a data race. An exception escaping into PJSUA2's C++ +callback tears down the worker thread and silently kills media for every call, +which is why the capture port catches and logs once rather than per frame. + +Full detail: [.claude/rules/concurrency-threads.md](../.claude/rules/concurrency-threads.md). + ## Design Decisions -### Why Sippy B2BUA + PJSUA2? +### Media plane: why PJSUA2 places the call -We split SIP signaling and media handling into two separate libraries: +The original split was *Sippy signals, PJSUA2 carries media*. It does not work, +for a reason that is not obvious until you try it: -- **Sippy B2BUA** handles SIP signaling (INVITE, BYE, REGISTER, re-INVITE, DTMF relay). It's battle-tested for telephony and handles the complex SIP state machine. -- **PJSUA2** handles RTP media (audio streams, conference bridge, recording, tone generation). It provides a clean C++/Python API for media manipulation without needing to deal with raw RTP. +**PJSUA2 exposes no standalone RTP media object.** Every `AudioMedia` subclass +in the Python bindings is a file player, recorder, tone generator, or capture +port. RTP is reachable only through `pj.Call.getAudioMedia()`, after +`onCallMediaState` fires on a dialog **PJSUA2 itself owns**. There is no +"give me an AudioMedia for this remote host:port" API to call. -This split lets us tap into the audio stream (for classification and STT) without interfering with SIP signaling, and bridge calls through a conference bridge for clean transfer. +So a design where Sippy owns the dialog can never obtain a media stream from +PJSUA2. `MediaPipeline.add_remote_stream()` is not unfinished work — it is a +function that cannot be written against this API. The consequence is that audio +never reaches the classifier: `create_tap` builds a valid capture port with +nothing to attach it to. + +**The resolution: PJSUA2 places the call; Sippy keeps every other role.** + +| Concern | Owner | +|---|---| +| Emergency guard, concurrency cap | `gateway.make_call` — unchanged, still first | +| Trunk registration | PJSUA2 `Account` | +| Outbound INVITE / answer / hangup | PJSUA2 `Call` | +| RTP, conference bridge, taps, recording | PJSUA2 media | +| DTMF | PJSUA2 `Call.dialDtmf` (RFC 2833) | +| Device registration, routing, leg bridging | Sippy / gateway | +| Inbound call dispatch | PJSUA2 `Account.onIncomingCall` | + +Sippy remains the SBC-shaped layer — it is where device registrations, routing +decisions and B2BUA leg-joining live. What moves is the raw dialog for a trunk +call, because owning the dialog is the price of owning the media. + +Alternatives considered and rejected: + +- **Terminate RTP ourselves** (aiortc or raw sockets) and feed PCM into + `AudioTap` directly, keeping Sippy on the wire. Preserves the split, but + means owning jitter buffering, packet loss concealment and ulaw/alaw + transcoding — precisely the work PJSUA2 exists to do. +- **A loopback `pj.Call` mirroring each real leg**, so PJSUA2 has a dialog it + owns. Avoids touching call placement, but adds a phantom call per real call + and the SDP juggling is fragile. + +> **Safety note for this refactor:** `is_emergency_number()` stays the first +> check in `gateway.make_call`, above the concurrency cap and above any SIP +> action, regardless of which library dials. A new outbound path that reaches +> the SIP layer without passing that guard is a serious regression even if +> every test passes. ### Why asyncio Queue-based EventBus? @@ -164,11 +230,13 @@ This split lets us tap into the audio stream (for classification and STT) withou - **Dead subscriber cleanup** — full queues are automatically removed - **Event history** — late joiners can catch up on recent events -If scaling to multiple gateway processes becomes necessary, the EventBus interface can be backed by Redis pub/sub without changing consumers. +If scaling to multiple gateway processes becomes necessary, the EventBus +interface can be backed by Redis pub/sub without changing consumers. ### Why OpenAI-compatible LLM API? -The LLM client uses raw HTTP (httpx) against any OpenAI-compatible endpoint. This means: +The LLM client uses raw HTTP (httpx) against any OpenAI-compatible endpoint. +This means: - **Ollama** (local, free) — `http://localhost:11434/v1` - **LM Studio** (local, free) — `http://localhost:1234/v1` @@ -176,3 +244,11 @@ The LLM client uses raw HTTP (httpx) against any OpenAI-compatible endpoint. Thi - **OpenAI** (cloud) — `https://api.openai.com/v1` No SDK dependency. No vendor lock-in. Switch models by changing one env var. + +## Testing against a fake PSTN + +`tests/lab/` runs an Asterisk instance that answers calls, plays an IVR, holds +with music and connects a "human" — so the gateway has something real to dial +that is not the PSTN. `SIP_TRUNK_HOST` is just an address, so the production +code path runs unmodified; while it points at the lab there is no route to the +PSTN at all. See [tests/lab/README.md](../tests/lab/README.md). diff --git a/tests/lab/README.md b/tests/lab/README.md index 5fe0a15..5946f85 100644 --- a/tests/lab/README.md +++ b/tests/lab/README.md @@ -63,6 +63,72 @@ GATEWAY_SIP_PORT=21062 # must differ from the Asterisk port > engine is assigned — `main.py`'s lifespan calls `build_sip_engine()` after > construction. A harness that skips that step silently tests the mock. +## The softphone (transfer target) + +**Asterisk is the registrar for devices, not Hold Slayer.** The gateway reaches +a desk phone by dialling extension `2001`, which Asterisk routes to whatever +has registered as `softphone`. This deliberately avoids Hold Slayer's own SIP +listener, which answers `200 OK` to any REGISTER with no digest challenge. + +The `pjsua` CLI built alongside the Python bindings is the test device — same +library stack as the gateway, no extra dependency. It needs an RPATH patch like +the bindings did: + +```bash +cp ~/src/pjproject/pjsip-apps/bin/pjsua-x86_64-pc-linux-gnu ~/.local/bin/pjsua +patchelf --set-rpath $HOME/.local/lib ~/.local/bin/pjsua +``` + +Register it (config file avoids shell-quoting pain): + +```bash +cat > softphone.cfg <<'EOF' +--null-audio +--auto-answer=200 +--max-calls=4 +--local-port=21070 +--id=sip:softphone@127.0.0.1 +--registrar=sip:127.0.0.1:21061 +--realm=asterisk +--username=softphone +--password=labphone +--log-level=3 +EOF + +# pjsua is an interactive console app: it exits ~8s after start if stdin is +# closed or /dev/null. Hold a fifo open on stdin — `script -qfc` and +# `setsid sp.fifo; pjsua --config-file softphone.cfg <&3 >softphone.log 2>&1' & +``` + +Verify — **check the port is actually bound**, not just that Asterisk holds a +contact, since a stale registration outlives the process: + +```bash +ss -lnup | grep 21070 # must be listening +docker compose -f docker-compose.lab.yml exec asterisk \ + asterisk -rx "pjsip show contacts" # must show softphone +``` + +Then place a call to `2001`. Both legs should show `Up` under one bridge id: + +```bash +docker compose -f docker-compose.lab.yml exec asterisk \ + asterisk -rx "core show channels concise" +``` + +> **`--realm=asterisk`, not `--realm='*'`** — the wildcard fails with +> `PJSIP_EFAILEDCREDENTIAL` against Asterisk's digest challenge. + +> **Qualify is off** for this AOR (`qualify_frequency = 0`): the pjsua console +> does not answer `OPTIONS`, so polling marks a working softphone `Unavail` and +> the dialplan refuses to ring it. The `2001` guard therefore tests +> `PJSIP_AOR(softphone,contact)` rather than `DEVICE_STATE`. A real hardphone +> answers OPTIONS and can have qualify re-enabled. + ## Useful commands ```bash diff --git a/tests/lab/dialplan/asterisk.conf b/tests/lab/dialplan/asterisk.conf index 2e71279..c6a91fe 100644 --- a/tests/lab/dialplan/asterisk.conf +++ b/tests/lab/dialplan/asterisk.conf @@ -1,24 +1,17 @@ ; Minimal Asterisk core config for the lab. -[directories](!) -astetcdir => /etc/asterisk -astmoddir => /usr/lib/asterisk/modules -astvarlibdir => /var/lib/asterisk -astdbdir => /var/lib/asterisk -astkeydir => /var/lib/asterisk -astdatadir => /var/lib/asterisk -astagidir => /var/lib/asterisk/agi-bin -astspooldir => /var/spool/asterisk -astrundir => /var/run/asterisk -astlogdir => /var/log/asterisk -astsbindir => /usr/sbin - +; +; Deliberately does NOT set [directories] or runuser/rungroup: the image's +; compiled-in defaults are correct, and it runs as the `asterisk` user via a +; USER directive. Overriding either risks breaking the container for no gain +; (an earlier version of this file did both). [options] -; Log to stdout so Docker's json-file driver captures it and Alloy ships it. -; A file-based log inside the container would be invisible to Loki. +; Log to stdout so Docker's json-file driver captures it and Alloy ships it +; to Loki. A file-based log inside the container would be invisible. verbose = 3 debug = 0 +; No ANSI colour. Asterisk colourises the console by default and the escape +; codes travel through Docker into Loki, where every line arrives wrapped in +; \x1b[0;30m — unreadable in Grafana and awkward to filter on. This setting +; lives here, not in logger.conf, and only takes effect if this file is +; actually mounted into the container. nocolor = yes -dumpcore = no -; Never run as root inside the container. -runuser = asterisk -rungroup = asterisk diff --git a/tests/lab/dialplan/extensions.conf b/tests/lab/dialplan/extensions.conf index 1092a68..e03b793 100644 --- a/tests/lab/dialplan/extensions.conf +++ b/tests/lab/dialplan/extensions.conf @@ -128,6 +128,23 @@ exten => 1008,1,NoOp(LAB 1008: silence) same => n,Wait(40) same => n,Hangup() +; --- 2001: ring the registered softphone ---------------------------------- +; The transfer target. Asterisk is the registrar for devices, so the gateway +; reaches a desk phone by dialling this rather than by registering it itself. +; Fails fast when nothing is registered — a silent 30s ring would look like a +; gateway bug rather than an absent softphone. +exten => 2001,1,NoOp(LAB 2001: ring softphone) + ; Count registered contacts rather than DEVICE_STATE: device state follows + ; the OPTIONS qualify, which is off for this AOR (the pjsua CLI does not + ; answer OPTIONS), so a registered softphone would still read UNAVAILABLE. + same => n,GotoIf($[${PJSIP_AOR(softphone,contact)} = ""]?nodevice) + same => n,Dial(PJSIP/softphone,30) + same => n,Hangup() + same => n(nodevice),NoOp(LAB 2001: no softphone registered) + same => n,Answer() + same => n,Playback(lab-speech) + same => n,Hangup() + ; --- echo test ------------------------------------------------------------ ; Not a scenario — a debugging aid. Echoes audio back so you can confirm ; bidirectional RTP by ear when something looks wrong. diff --git a/tests/lab/dialplan/logger.conf b/tests/lab/dialplan/logger.conf index 6c15a54..76f4306 100644 --- a/tests/lab/dialplan/logger.conf +++ b/tests/lab/dialplan/logger.conf @@ -3,6 +3,22 @@ ; the container would put the logs where nothing can see them. [general] dateformat = %F %T +; Colour is disabled in asterisk.conf (`nocolor = yes`), not here — Asterisk +; colourises the console by default and the escape codes travel through Docker +; into Loki, where every line arrives wrapped in \x1b[0;30m. That file must be +; mounted for the setting to take effect. [logfiles] -console => notice,warning,error +; Warnings and errors only. That covers what matters when a lab call +; misbehaves: failed authentication, no-matching-endpoint, playback failures. +; +; `notice` and `verbose` are both excluded because the "Remote UNIX +; connection" pairs the healthcheck generates arrive on those channels, and +; they swamped everything else (the Loki stream measured 100% healthcheck +; noise before this). The healthcheck itself is also reduced to one CLI call +; on a 60s interval in docker-compose — the two changes work together. +; +; To trace a call's dialplan execution, raise verbosity at runtime rather than +; leaving it on: +; asterisk -rx "core set verbose 3" +console => warning,error diff --git a/tests/lab/dialplan/pjsip.conf b/tests/lab/dialplan/pjsip.conf index 168bd0f..881ede8 100644 --- a/tests/lab/dialplan/pjsip.conf +++ b/tests/lab/dialplan/pjsip.conf @@ -34,14 +34,21 @@ local_net = {{ asterisk_local_net }} ; --------------------------------------------------------------------------- ; Hold Slayer authenticates as this endpoint to place calls into the lab. -; Identify the endpoint by source address. Asterisk's default matching uses -; the From-header domain, which Hold Slayer populates from its SIP bind -; address (0.0.0.0 on a wildcard bind) — never a value Asterisk can match. -; Matching on where the packet actually came from sidesteps that. +; Identify the endpoint by source address *and port*. Asterisk's default +; matching uses the From-header domain, which Hold Slayer populates from its +; SIP bind address (0.0.0.0 on a wildcard bind) — never a value Asterisk can +; match. Matching on where the packet came from sidesteps that. +; +; The port is essential when the softphone runs on the same host: a +; host-only match claims *every* packet from that address, so the +; softphone's REGISTER would be attributed to this endpoint and checked +; against the gateway's password ("Failed to authenticate", confusingly). +; Endpoints that authenticate by username (the softphone) must not be +; covered by an identify block. [hold-slayer] type = identify endpoint = hold-slayer -match = {{ asterisk_match_host }} +match = {{ asterisk_match_host }}:{{ asterisk_gateway_port }} [hold-slayer] type = endpoint @@ -68,6 +75,56 @@ auth_type = userpass username = {{ asterisk_sip_username }} password = {{ asterisk_sip_password }} +; --------------------------------------------------------------------------- +; Softphone endpoint — the transfer target +; --------------------------------------------------------------------------- +; Asterisk is the registrar for devices, not Hold Slayer. A softphone REGISTERs +; here and the gateway transfers a live call to it by dialling extension 2001. +; +; Deliberate: Hold Slayer's own SIP listener answers 200 OK to any REGISTER +; with no digest challenge, so anything on the network could register as a +; device and receive transferred calls. Keeping registration in Asterisk means +; the lab does not depend on that path, and the softphone is authenticated. +; +; Test with the pjsua CLI built alongside the Python bindings: +; pjsua --null-audio --auto-answer=200 \ +; --id=sip:softphone@ \ +; --registrar=sip::21061 \ +; --realm='*' --username=softphone --password= \ +; --local-port= + +[softphone] +type = endpoint +context = hold-slayer-lab +disallow = all +allow = ulaw +allow = alaw +auth = softphone-auth +aors = softphone +dtmf_mode = rfc4733 +direct_media = no +force_rport = yes +rewrite_contact = yes +rtp_symmetric = yes + +[softphone-auth] +type = auth +auth_type = userpass +username = {{ asterisk_softphone_username }} +password = {{ asterisk_softphone_password }} + +[softphone] +type = aor +; The device's contact is learned from its REGISTER rather than configured — +; a softphone's port is not known in advance. +max_contacts = 1 +remove_existing = yes +; No qualify: the pjsua CLI does not answer OPTIONS while sitting at its +; console prompt, so polling marks a perfectly working softphone Unavail and +; the dialplan refuses to ring it. Registration itself is the liveness signal +; here. A real hardphone answers OPTIONS and can have qualify re-enabled. +qualify_frequency = 0 + [hold-slayer] type = aor max_contacts = 2 diff --git a/tests/lab/docker-compose.lab.yml b/tests/lab/docker-compose.lab.yml index dd9c750..ae9b5ef 100644 --- a/tests/lab/docker-compose.lab.yml +++ b/tests/lab/docker-compose.lab.yml @@ -21,8 +21,22 @@ services: - ./dialplan/pjsip.local.conf:/etc/asterisk/pjsip.conf:ro - ./dialplan/rtp.local.conf:/etc/asterisk/rtp.conf:ro - ./dialplan/logger.conf:/etc/asterisk/logger.conf:ro + # Carries `nocolor = yes`: without it every log line reaches Loki + # wrapped in ANSI escape codes. + - ./dialplan/asterisk.conf:/etc/asterisk/asterisk.conf:ro # The image ships no sound files at all. These are generated by # sounds/generate.py; Asterisk resolves Playback(lab-music) to # lab-music.sln here (8kHz signed-linear, no transcoding). - ./sounds:/var/lib/asterisk/sounds/en:ro + # The image's default command is `-vvvdddf` — verbosity 3 and debug 3 + # forced on the command line, which overrides both asterisk.conf and + # logger.conf. That makes every healthcheck CLI connection log a + # "Remote UNIX connection" pair: ~2900 lines/day of pure noise that + # completely buried the real SIP events in Loki. + # + # -f foreground (required: Docker needs PID 1 to stay), -T timestamps, + # -W colour off, -U run as asterisk, -p realtime priority. No -v, no -d: + # warnings and errors still log, and verbosity can be raised at runtime + # with `asterisk -rx "core set verbose 3"` when tracing a call. + command: ["/usr/sbin/asterisk", "-f", "-T", "-W", "-U", "asterisk", "-p"] restart: unless-stopped diff --git a/tests/lab/sounds/generate.py b/tests/lab/sounds/generate.py index 30035c9..b870975 100644 --- a/tests/lab/sounds/generate.py +++ b/tests/lab/sounds/generate.py @@ -12,7 +12,6 @@ finds `lab-music.sln`. python generate.py [outdir] """ -import struct import sys from pathlib import Path @@ -29,13 +28,14 @@ def _write_sln(path: Path, samples: np.ndarray) -> None: print(f" {path.name}: {len(pcm) / RATE:.1f}s ({path.stat().st_size} bytes)") -def make_music(seconds: float = 30.0) -> np.ndarray: +def make_music(seconds: float = 30.0, seed: int = 7) -> np.ndarray: """Sustained multi-harmonic tones — what the classifier must call MUSIC. A chord progression with stable pitch and strong harmonic structure. The steady spectrum across a long window is what distinguishes music from speech; this deliberately has no pauses. """ + rng = np.random.default_rng(seed) t = np.linspace(0, seconds, int(RATE * seconds), endpoint=False) # A-minor-ish progression, one chord per 2s bar. chords = [(220.0, 261.6, 329.6), (196.0, 246.9, 293.7), @@ -48,12 +48,21 @@ def make_music(seconds: float = 30.0) -> np.ndarray: break mask = (t >= start) & (t < end) for j, freq in enumerate(chord): - # Fundamental plus two harmonics, decaying — a plucked-string feel. - for h, amp in ((1, 0.30), (2, 0.12), (3, 0.05)): + # Fundamental plus four harmonics, decaying — a plucked-string + # feel. Enough harmonics to keep spectral flatness inside the + # music score's 0.05-0.4 band: with only three, some windows fall + # *below* 0.05 (too pure to read as music) and score as speech. + for h, amp in ((1, 0.30), (2, 0.12), (3, 0.05), (4, 0.03), (5, 0.02)): out[mask] += amp / (j + 1) * np.sin(2 * np.pi * freq * h * t[mask]) # Gentle per-bar envelope so bars are distinguishable but never silent. env = 0.8 + 0.2 * np.sin(2 * np.pi * (t[mask] - start) / bar) out[mask] *= env + + # Recording-style noise floor. Windows straddling a chord change have a + # momentarily sparse spectrum and land just *under* the music score's + # 0.05 flatness floor, scoring as speech. This is well below the level + # that would disturb tonality — every real recording has one. + out += rng.normal(0, 0.004, len(out)) return out * 0.45 @@ -74,16 +83,44 @@ def make_speech(seconds: float = 8.0, seed: int = 1337) -> np.ndarray: mask = (t >= pos) & (t < pos + syl) if mask.any(): local = t[mask] - pos - f0 = rng.uniform(95, 165) # fundamental — adult speaking range - # Two formants, swept slightly across the syllable. The ranges - # deliberately avoid the DTMF bands (rows 697-941, columns - # 1209-1633): a formant pair landing on both trips the Goertzel - # detector and the whole utterance is classified as a keypress. - f1 = rng.uniform(300, 620) + rng.uniform(-40, 40) * local / syl - f2 = rng.uniform(1750, 2600) + rng.uniform(-120, 120) * local / syl - sig = (0.50 * np.sin(2 * np.pi * f0 * local) + frac = local / syl + + # Pitch CONTOUR, not a constant. This is the single feature that + # separates this fixture from music. `_detect_tonality` looks for + # an autocorrelation peak > 0.5 in the 50-1000 Hz lag range; a + # fixed f0 is perfectly periodic there, scores is_tonal=True, and + # hands the music score a free 0.3 that speech cannot outrun. + # Real voices glide and jitter, so the periodicity never locks. + f0_start = rng.uniform(95, 165) + f0_end = f0_start * rng.uniform(0.72, 1.38) # rise or fall + f0 = f0_start + (f0_end - f0_start) * frac + # Cycle-to-cycle jitter on top of the glide (~2% is human). + f0 *= 1.0 + 0.02 * rng.standard_normal(len(local)) + # Integrate frequency to phase — with a varying f0, `2*pi*f*t` + # would be wrong (that is a chirp only if f is the *instantaneous* + # rate, which it is not once f0 itself moves). + ph0 = 2 * np.pi * np.cumsum(f0) / RATE + + # Two formants, swept across the syllable. The ranges deliberately + # avoid the DTMF bands (rows 697-941, columns 1209-1633): a formant + # pair landing on both trips the Goertzel detector and the whole + # utterance is classified as a keypress. + f1 = rng.uniform(300, 620) + rng.uniform(-40, 40) * frac + f2 = rng.uniform(1750, 2600) + rng.uniform(-120, 120) * frac + sig = (0.50 * np.sin(ph0) + 0.30 * np.sin(2 * np.pi * f1 * local) + 0.18 * np.sin(2 * np.pi * f2 * local)) + # Aspiration noise — HIGH-PASSED, not broadband. Real speech noise + # sits above the formants; flat noise puts energy in every + # Goertzel bin, so the strongest DTMF row and column both clear + # the detector's `total_power * 0.1` threshold and every syllable + # reads as a keypress. A first-difference filter (y[n]-y[n-1]) is + # a cheap +6dB/octave tilt that leaves the 697-1633 Hz DTMF bands + # comparatively empty. The 0.09 level is chosen for margin: it puts + # spectral flatness at ~0.46, mid-way through the 0.1-0.5 band the + # speech score rewards, rather than on either edge. + noise = rng.standard_normal(len(local) + 1) + sig += 0.09 * np.diff(noise) # Raised-cosine envelope: no clicks at syllable edges. sig *= np.sin(np.pi * local / syl) ** 0.6 out[mask] += sig diff --git a/tests/test_lab_fixtures.py b/tests/test_lab_fixtures.py new file mode 100644 index 0000000..961130b --- /dev/null +++ b/tests/test_lab_fixtures.py @@ -0,0 +1,109 @@ +""" +Lab audio fixtures — the classifier must agree with what each one claims to be. + +These guard the *fixtures*, not the classifier. `tests/lab/sounds/generate.py` +synthesises music/speech/silence that the Asterisk lab plays down a real call; +if a fixture drifts into the wrong class, every lab result built on it is +quietly meaningless — a hold-music scenario that never classifies as music +proves nothing about the hold slayer. + +The first version of these fixtures passed on the opening 3s window and drifted +to MUSIC after, which a single-window check would not have caught. Hence the +sweep across every window. + +Skipped when the fixtures have not been generated: they are gitignored (~680K, +reproducible from a fixed seed), so a fresh checkout has none until +`python tests/lab/sounds/generate.py` runs. +""" + +import subprocess +import sys +from pathlib import Path + +import numpy as np +import pytest + +from config import Settings +from models.call import AudioClassification +from services.audio_classifier import SAMPLE_RATE, AudioClassifier + +SOUNDS_DIR = Path(__file__).parent / "lab" / "sounds" +GENERATOR = SOUNDS_DIR / "generate.py" + +# The lab writes 8 kHz .sln; the classifier works at 16 kHz. +LAB_RATE = 8000 +WINDOW_SAMPLES = SAMPLE_RATE * 3 # classifier's 3s analysis window + +FIXTURES = [ + ("lab-music.sln", AudioClassification.MUSIC), + ("lab-speech.sln", AudioClassification.LIVE_HUMAN), + ("lab-silence.sln", AudioClassification.SILENCE), +] + + +def _load_16k(path: Path) -> np.ndarray: + """Load an 8 kHz .sln and upsample to the classifier's 16 kHz.""" + return np.repeat(np.fromfile(path, dtype=" dict[str, bytes]: + subprocess.run( + [sys.executable, str(GENERATOR), str(target)], + check=True, + capture_output=True, + ) + return {p.name: p.read_bytes() for p in sorted(target.glob("*.sln"))} + + first = run(tmp_path / "a") + second = run(tmp_path / "b") + + assert first, "generator produced no .sln files" + assert first.keys() == second.keys() + for name in first: + assert first[name] == second[name], f"{name} differs between runs"