PJSUA2 media plane — audio finally reaches the classifier #8

Merged
r merged 5 commits from feat/pjsua2-media-tap into main 2026-07-29 21:34:10 +00:00
14 changed files with 1147 additions and 131 deletions

View File

@@ -162,6 +162,12 @@ class Settings(BaseSettings):
# silently degrading to a gateway that can't place real calls. # silently degrading to a gateway that can't place real calls.
use_mock_sip: bool = False 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 # Notifications
notify_sms_number: str = "" notify_sms_number: str = ""

View File

@@ -55,6 +55,26 @@ def build_sip_engine(
"for development without a trunk." "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( return SippyEngine(
sip_address=gw_sip.host, sip_address=gw_sip.host,
sip_port=gw_sip.port, sip_port=gw_sip.port,

View File

@@ -20,6 +20,7 @@ PJSUA2 runs in its own thread with a dedicated Endpoint.
""" """
import asyncio import asyncio
import gc
import logging import logging
import threading import threading
from collections.abc import AsyncIterator from collections.abc import AsyncIterator
@@ -98,6 +99,72 @@ class AudioTap:
self._active = False 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 # Stream Entry — tracks a single media stream in the pipeline
# ================================================================ # ================================================================
@@ -112,6 +179,8 @@ class MediaStream:
self.codec = codec self.codec = codec
self.conf_port: Optional[int] = None # PJSUA2 conference bridge port ID self.conf_port: Optional[int] = None # PJSUA2 conference bridge port ID
self.transport = None # PJSUA2 SipTransport 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.rtp_port: Optional[int] = None # Local RTP listen port
self.taps: list[AudioTap] = [] self.taps: list[AudioTap] = []
self.recorder = None # PJSUA2 AudioMediaRecorder self.recorder = None # PJSUA2 AudioMediaRecorder
@@ -143,10 +212,11 @@ class MediaPipeline:
pipeline = MediaPipeline() pipeline = MediaPipeline()
await pipeline.start() await pipeline.start()
# Add a stream for a call leg # Media arrives from the SIP engine's onCallMediaState callback
port = pipeline.add_remote_stream("leg_1", "10.0.0.1", 20000, "PCMU") # (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") tap = pipeline.create_tap("leg_1")
async for frame in tap.stream(): async for frame in tap.stream():
classify(frame) classify(frame)
@@ -173,6 +243,7 @@ class MediaPipeline:
self._next_rtp_port = rtp_start_port self._next_rtp_port = rtp_start_port
self._sample_rate = sample_rate self._sample_rate = sample_rate
self._channels = channels 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) self._null_audio = null_audio # Use null audio device (no sound card needed)
# State # State
@@ -255,11 +326,17 @@ class MediaPipeline:
tap.close() tap.close()
self._taps.clear() self._taps.clear()
# Remove all streams # Remove all streams (this releases their capture ports)
for stream_id in list(self._streams.keys()): for stream_id in list(self._streams.keys()):
self.remove_stream(stream_id) 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: if self._endpoint:
try: try:
self._endpoint.libDestroy() self._endpoint.libDestroy()
@@ -270,6 +347,15 @@ class MediaPipeline:
self._ready = False self._ready = False
logger.info("🎵 PJSUA2 media pipeline stopped") 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 @property
def is_ready(self) -> bool: def is_ready(self) -> bool:
return self._ready return self._ready
@@ -291,50 +377,51 @@ class MediaPipeline:
# Stream Management # Stream Management
# ================================================================ # ================================================================
def add_remote_stream( def attach_call_media(self, stream_id: str, audio_media) -> Optional[int]:
self, stream_id: str, remote_host: str, remote_port: int, codec: str = "PCMU" """Register a call's live ``AudioMedia`` with the pipeline.
) -> Optional[int]:
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 stream.media = audio_media
party's RTP stream, connecting it to the conference bridge.
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)
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: try:
import pjsua2 as pj stream.conf_port = audio_media.getPortId()
except Exception:
stream.conf_port = None
# Create a media transport for this stream # Wire up taps that were requested before media came up.
# In a full implementation, we'd create an AudioMediaPort pending = self._taps.get(stream_id, [])
# that receives RTP and feeds it into the conference bridge if pending and stream.capture_port is None:
transport_cfg = pj.TransportConfig() port = make_capture_port(
transport_cfg.port = stream.rtp_port stream_id, self._sample_rate, self._channels, self._frame_ms
)
# The conference bridge port will be assigned when if port is not None:
# the call's media is activated via onCallMediaState try:
audio_media.startTransmit(port)
stream.capture_port = port
port.taps.extend(pending)
logger.info( logger.info(
f" 📡 Added stream {stream_id}: " f" 🎤 Audio tap attached for {stream_id} "
f"local={stream.rtp_port} → remote={remote_host}:{remote_port} ({codec})" f"({len(pending)} waiting)"
)
except Exception as e:
logger.error(
f" Failed to attach capture port for {stream_id}: {e}",
exc_info=True,
) )
except ImportError: logger.info(f" 📡 Media attached for {stream_id} (conf port {stream.conf_port})")
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
return stream.conf_port return stream.conf_port
def remove_stream(self, stream_id: str) -> None: def remove_stream(self, stream_id: str) -> None:
@@ -350,6 +437,19 @@ class MediaPipeline:
tap.close() tap.close()
self._taps.pop(stream_id, None) 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 # Stop recording
if stream.recorder: if stream.recorder:
try: try:
@@ -426,16 +526,28 @@ class MediaPipeline:
self._taps[stream_id] = [] self._taps[stream_id] = []
self._taps[stream_id].append(tap) 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: try:
import pjsua2 as pj # One capture port per stream, shared by every tap on it:
# Create an AudioMediaPort that captures frames # the bridge would otherwise mix each additional port back
# and feeds them to the tap # into the conference and the call would echo.
# In PJSUA2, we'd subclass AudioMediaPort and implement if stream.capture_port is None:
# onFrameReceived to call tap.feed(frame_data) 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)") 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: 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: else:
logger.info(f" 🎤 Audio tap created for {stream_id} (virtual)") logger.info(f" 🎤 Audio tap created for {stream_id} (virtual)")

488
core/pjsua_engine.py Normal file
View File

@@ -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()

View File

@@ -276,17 +276,22 @@ class SippyEngine(SIPEngine):
if state == "connected": if state == "connected":
sdp = data.get("sdp") sdp = data.get("sdp")
if sdp and self.media_pipeline: 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: try:
remote_rtp = self._parse_sdp_rtp_endpoint(sdp) remote_rtp = self._parse_sdp_rtp_endpoint(sdp)
if remote_rtp: if remote_rtp:
leg.media_port = self.media_pipeline.add_remote_stream( logger.info(
leg.leg_id, f" {leg.leg_id}: remote RTP "
remote_rtp["host"], f"{remote_rtp['host']}:{remote_rtp['port']} "
remote_rtp["port"], f"({remote_rtp['codec']}) — no media on this engine"
remote_rtp["codec"],
) )
except Exception as e: 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": elif state == "terminated":
if self.media_pipeline and leg.media_port is not None: if self.media_pipeline and leg.media_port is not None:
try: try:

View File

@@ -1,6 +1,15 @@
# Architecture # 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 ## System Diagram
@@ -27,8 +36,13 @@ Hold Slayer is a single-process async Python application built on FastAPI. It ac
│ └────┬─────┘ └─────┬─────┘ └──────────────┘ │ │ └────┬─────┘ └─────┬─────┘ └──────────────┘ │
│ │ │ │ │ │ │ │
│ ┌────┴──────────────┴───────────────────┐ │ │ ┌────┴──────────────┴───────────────────┐ │
│ │ Sippy B2BUA Engine │ │ │ │ SIP Engine │ │
│ │ (SIP calls, DTMF, conference bridge) │ │ │ │ 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 | | Component | File | Purpose |
|-----------|------|---------| |-----------|------|---------|
| Sippy Engine | `core/sippy_engine.py` | SIP signaling (INVITE, BYE, REGISTER, DTMF) | | Sippy Engine | `core/sippy_engine.py` | SIP signalling (INVITE, BYE, REGISTER, DTMF) |
| Media Pipeline | `core/media_pipeline.py` | PJSUA2 RTP media handling, conference bridge, recording | | Media Pipeline | `core/media_pipeline.py` | PJSUA2 RTP media, conference bridge, taps, recording |
| Recording | `services/recording.py` | WAV file management and storage | | Recording | `services/recording.py` | WAV file management and storage |
| Analytics | `services/call_analytics.py` | Call metrics, hold time stats, trends | | Analytics | `services/call_analytics.py` | Call metrics, hold time stats, trends |
| Notifications | `services/notification.py` | WebSocket + SMS alerts | | 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 } POST /api/v1/calls/hold-slayer { number, intent, call_flow_id }
2. Gateway.make_call() 2. Gateway.make_call()
├── is_emergency_number() → REFUSE 911/112 (before anything else)
├── concurrency cap check → refuse past max_concurrent_calls
├── CallManager.create_call() → track state ├── CallManager.create_call() → track state
── SippyEngine.make_call() → SIP INVITE to trunk ── sip_engine.make_call() → place the call, media follows
└── MediaPipeline.add_stream() → RTP media setup
3. HoldSlayer.run_with_flow() or run_exploration() 3. HoldSlayer.run_with_flow() or run_exploration()
├── AudioClassifier.classify() → analyze 3s audio windows ├── 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 ├── TranscriptionService.transcribe() → STT on speech audio
├── LLMClient.analyze_ivr_menu() → pick menu option (fallback) ├── 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() └── detect_hold_to_human_transition()
└── HUMAN_DETECTED! → transfer └── HUMAN_DETECTED! → transfer
4. Transfer 4. Transfer
├── SippyEngine.bridge() → connect call legs ├── SippyEngine.bridge_calls() → join the two call legs
├── MediaPipeline.bridge_streams() → bridge RTP ├── MediaPipeline.bridge_streams() → bridge RTP in the conf bridge
├── EventBus.publish(TRANSFER_STARTED) ├── EventBus.publish(TRANSFER_STARTED)
└── NotificationService → "Pick up your phone!" └── 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 → 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 ## Threading Model
Hold Slayer is primarily single-threaded async (asyncio), with one exception: The README's "single-process async" is a simplification. There are **three**
execution contexts, and the boundaries between them are the highest-leverage
- **Main thread**: FastAPI + all async services (event bus, hold slayer, classifier, etc.) invariant in the codebase.
- **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).
``` ```
Main Thread (asyncio) asyncio loop (main thread) Sippy ED thread PJSUA2 worker threads
├── FastAPI (uvicorn) ├── FastAPI (uvicorn) └── ED2 dispatcher └── media / RTP
├── EventBus ├── EventBus ├── SIP signalling └── onFrameReceived
├── CallManager ├── CallManager ├── UA objects
├── HoldSlayer ├── HoldSlayer └── DTMF relay
├── AudioClassifier ├── AudioClassifier
├── TranscriptionService ├── TranscriptionService
├── LLMClient ├── LLMClient
├── MediaPipeline (PJSUA2)
├── NotificationService ├── NotificationService
└── RecordingService └── 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 ## 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 exposes no standalone RTP media object.** Every `AudioMedia` subclass
- **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. 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? ### 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 - **Dead subscriber cleanup** — full queues are automatically removed
- **Event history** — late joiners can catch up on recent events - **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? ### 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` - **Ollama** (local, free) — `http://localhost:11434/v1`
- **LM Studio** (local, free) — `http://localhost:1234/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` - **OpenAI** (cloud) — `https://api.openai.com/v1`
No SDK dependency. No vendor lock-in. Switch models by changing one env var. 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).

View File

@@ -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 > engine is assigned — `main.py`'s lifespan calls `build_sip_engine()` after
> construction. A harness that skips that step silently tests the mock. > 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 </dev/null` both look like they work (registration succeeds) and
# then the process dies, leaving a stale contact in Asterisk that routes
# INVITEs to a port nobody is listening on.
mkfifo sp.fifo
setsid sh -c 'exec 3<>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 ## Useful commands
```bash ```bash

View File

@@ -1,24 +1,17 @@
; Minimal Asterisk core config for the lab. ; Minimal Asterisk core config for the lab.
[directories](!) ;
astetcdir => /etc/asterisk ; Deliberately does NOT set [directories] or runuser/rungroup: the image's
astmoddir => /usr/lib/asterisk/modules ; compiled-in defaults are correct, and it runs as the `asterisk` user via a
astvarlibdir => /var/lib/asterisk ; USER directive. Overriding either risks breaking the container for no gain
astdbdir => /var/lib/asterisk ; (an earlier version of this file did both).
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
[options] [options]
; Log to stdout so Docker's json-file driver captures it and Alloy ships it. ; 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. ; to Loki. A file-based log inside the container would be invisible.
verbose = 3 verbose = 3
debug = 0 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 nocolor = yes
dumpcore = no
; Never run as root inside the container.
runuser = asterisk
rungroup = asterisk

View File

@@ -128,6 +128,23 @@ exten => 1008,1,NoOp(LAB 1008: silence)
same => n,Wait(40) same => n,Wait(40)
same => n,Hangup() 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 ------------------------------------------------------------ ; --- echo test ------------------------------------------------------------
; Not a scenario — a debugging aid. Echoes audio back so you can confirm ; Not a scenario — a debugging aid. Echoes audio back so you can confirm
; bidirectional RTP by ear when something looks wrong. ; bidirectional RTP by ear when something looks wrong.

View File

@@ -3,6 +3,22 @@
; the container would put the logs where nothing can see them. ; the container would put the logs where nothing can see them.
[general] [general]
dateformat = %F %T 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] [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

View File

@@ -34,14 +34,21 @@ local_net = {{ asterisk_local_net }}
; --------------------------------------------------------------------------- ; ---------------------------------------------------------------------------
; Hold Slayer authenticates as this endpoint to place calls into the lab. ; Hold Slayer authenticates as this endpoint to place calls into the lab.
; Identify the endpoint by source address. Asterisk's default matching uses ; Identify the endpoint by source address *and port*. Asterisk's default
; the From-header domain, which Hold Slayer populates from its SIP bind ; matching uses the From-header domain, which Hold Slayer populates from its
; address (0.0.0.0 on a wildcard bind) — never a value Asterisk can match. ; SIP bind address (0.0.0.0 on a wildcard bind) — never a value Asterisk can
; Matching on where the packet actually came from sidesteps that. ; 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] [hold-slayer]
type = identify type = identify
endpoint = hold-slayer endpoint = hold-slayer
match = {{ asterisk_match_host }} match = {{ asterisk_match_host }}:{{ asterisk_gateway_port }}
[hold-slayer] [hold-slayer]
type = endpoint type = endpoint
@@ -68,6 +75,56 @@ auth_type = userpass
username = {{ asterisk_sip_username }} username = {{ asterisk_sip_username }}
password = {{ asterisk_sip_password }} 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@<asterisk-host> \
; --registrar=sip:<asterisk-host>:21061 \
; --realm='*' --username=softphone --password=<pw> \
; --local-port=<free 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] [hold-slayer]
type = aor type = aor
max_contacts = 2 max_contacts = 2

View File

@@ -21,8 +21,22 @@ services:
- ./dialplan/pjsip.local.conf:/etc/asterisk/pjsip.conf:ro - ./dialplan/pjsip.local.conf:/etc/asterisk/pjsip.conf:ro
- ./dialplan/rtp.local.conf:/etc/asterisk/rtp.conf:ro - ./dialplan/rtp.local.conf:/etc/asterisk/rtp.conf:ro
- ./dialplan/logger.conf:/etc/asterisk/logger.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 # The image ships no sound files at all. These are generated by
# sounds/generate.py; Asterisk resolves Playback(lab-music) to # sounds/generate.py; Asterisk resolves Playback(lab-music) to
# lab-music.sln here (8kHz signed-linear, no transcoding). # lab-music.sln here (8kHz signed-linear, no transcoding).
- ./sounds:/var/lib/asterisk/sounds/en:ro - ./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 restart: unless-stopped

View File

@@ -12,7 +12,6 @@ finds `lab-music.sln`.
python generate.py [outdir] python generate.py [outdir]
""" """
import struct
import sys import sys
from pathlib import Path 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)") 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. """Sustained multi-harmonic tones — what the classifier must call MUSIC.
A chord progression with stable pitch and strong harmonic structure. The A chord progression with stable pitch and strong harmonic structure. The
steady spectrum across a long window is what distinguishes music from steady spectrum across a long window is what distinguishes music from
speech; this deliberately has no pauses. speech; this deliberately has no pauses.
""" """
rng = np.random.default_rng(seed)
t = np.linspace(0, seconds, int(RATE * seconds), endpoint=False) t = np.linspace(0, seconds, int(RATE * seconds), endpoint=False)
# A-minor-ish progression, one chord per 2s bar. # A-minor-ish progression, one chord per 2s bar.
chords = [(220.0, 261.6, 329.6), (196.0, 246.9, 293.7), 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 break
mask = (t >= start) & (t < end) mask = (t >= start) & (t < end)
for j, freq in enumerate(chord): for j, freq in enumerate(chord):
# Fundamental plus two harmonics, decaying — a plucked-string feel. # Fundamental plus four harmonics, decaying — a plucked-string
for h, amp in ((1, 0.30), (2, 0.12), (3, 0.05)): # 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]) 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. # 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) env = 0.8 + 0.2 * np.sin(2 * np.pi * (t[mask] - start) / bar)
out[mask] *= env 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 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) mask = (t >= pos) & (t < pos + syl)
if mask.any(): if mask.any():
local = t[mask] - pos local = t[mask] - pos
f0 = rng.uniform(95, 165) # fundamental — adult speaking range frac = local / syl
# Two formants, swept slightly across the syllable. The ranges
# deliberately avoid the DTMF bands (rows 697-941, columns # Pitch CONTOUR, not a constant. This is the single feature that
# 1209-1633): a formant pair landing on both trips the Goertzel # separates this fixture from music. `_detect_tonality` looks for
# detector and the whole utterance is classified as a keypress. # an autocorrelation peak > 0.5 in the 50-1000 Hz lag range; a
f1 = rng.uniform(300, 620) + rng.uniform(-40, 40) * local / syl # fixed f0 is perfectly periodic there, scores is_tonal=True, and
f2 = rng.uniform(1750, 2600) + rng.uniform(-120, 120) * local / syl # hands the music score a free 0.3 that speech cannot outrun.
sig = (0.50 * np.sin(2 * np.pi * f0 * local) # 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.30 * np.sin(2 * np.pi * f1 * local)
+ 0.18 * np.sin(2 * np.pi * f2 * 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. # Raised-cosine envelope: no clicks at syllable edges.
sig *= np.sin(np.pi * local / syl) ** 0.6 sig *= np.sin(np.pi * local / syl) ** 0.6
out[mask] += sig out[mask] += sig

109
tests/test_lab_fixtures.py Normal file
View File

@@ -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="<i2"), SAMPLE_RATE // LAB_RATE)
def _windows(samples: np.ndarray, step: int):
"""Yield successive analysis windows; at least one, even for short files."""
end = max(1, len(samples) - WINDOW_SAMPLES)
for offset in range(0, end, step):
yield samples[offset : offset + WINDOW_SAMPLES].astype("<i2").tobytes()
@pytest.fixture(scope="module")
def classifier():
return AudioClassifier(settings=Settings().classifier)
@pytest.mark.parametrize("filename,expected", FIXTURES)
def test_fixture_classifies_correctly_in_every_window(filename, expected, classifier):
"""Every window must classify correctly — not just the first.
Stepped at half the window length so windows overlap: a fixture that only
works on aligned boundaries would still be a trap in a live call, where
the window has no relationship to where the audio started.
"""
path = SOUNDS_DIR / filename
if not path.exists():
pytest.skip(f"{filename} not generated — run {GENERATOR}")
samples = _load_16k(path)
results = [
classifier.classify_chunk(w).audio_type
for w in _windows(samples, step=WINDOW_SAMPLES // 2)
]
wrong = [(i, r.value) for i, r in enumerate(results) if r is not expected]
assert not wrong, (
f"{filename} must classify as {expected.value} in all "
f"{len(results)} windows; wrong: {wrong}"
)
def test_generator_is_deterministic(tmp_path):
"""Same bytes on every run — the whole point of synthesising them.
Real hold music varies per call, so a classifier regression on the PSTN is
indistinguishable from noise. Fixed-seed audio makes the answer binary.
"""
if not GENERATOR.exists():
pytest.skip("generator not present")
def run(target: Path) -> 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"