From 06daf0533d128722807a38d9d1d3de6a8d17e25f Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 18 Mar 2026 12:12:17 +0100 Subject: [PATCH 1/7] Use hi-precision scheduler to send RTP packets --- tests/test_audio.py | 18 +++++++----------- voip/audio.py | 27 +++++++++++++++------------ 2 files changed, 22 insertions(+), 23 deletions(-) diff --git a/tests/test_audio.py b/tests/test_audio.py index 9bd9688..2449c70 100644 --- a/tests/test_audio.py +++ b/tests/test_audio.py @@ -442,27 +442,23 @@ async def test_send_rtp_audio__drops_audio_when_no_remote_addr(self, caplog): assert any("dropping audio" in r.message for r in caplog.records) async def test_send_rtp_audio__paces_packets_at_20ms_intervals(self): - """send_rtp_audio sleeps RTP_PACKET_DURATION_SECS between each packet.""" + """Packets are scheduled via sched.scheduler at fixed rpt_packet_duration intervals.""" call = make_audio_call(media=PCMU_MEDIA) remote_addr = ("10.0.0.2", 5006) call.rtp.calls = {remote_addr: call} - audio = np.zeros(320, dtype=np.float32) - sleep_calls: list[float] = [] - original_sleep = asyncio.sleep - - async def capture_sleep(delay: float) -> None: - sleep_calls.append(delay) - await original_sleep(0) with ( - patch("voip.audio.asyncio.sleep", side_effect=capture_sleep), + patch.object(call, "send_audio_scheduler") as mock_scheduler, patch.object(call, "send_packet"), ): await call.send_audio(audio) - assert len(sleep_calls) == 2 - assert all(s <= call.rpt_packet_duration.total_seconds() for s in sleep_calls) + assert mock_scheduler.enterabs.call_count == 2 + scheduled_times = [c.args[0] for c in mock_scheduler.enterabs.call_args_list] + assert scheduled_times[1] - scheduled_times[0] == pytest.approx( + call.rpt_packet_duration.total_seconds() + ) def make_echo_call(**kwargs) -> EchoCall: diff --git a/voip/audio.py b/voip/audio.py index 7652002..6578233 100644 --- a/voip/audio.py +++ b/voip/audio.py @@ -16,6 +16,7 @@ import datetime import json import logging +import sched import secrets import time from typing import ClassVar @@ -74,13 +75,15 @@ class AudioCall(RTPCall): rtp_ssrc: int = dataclasses.field( init=False, repr=False, default_factory=generate_ssrc ) - last_packet_time: float = dataclasses.field( - init=False, repr=False, default_factory=time.perf_counter - ) send_audio_lock: asyncio.Lock = dataclasses.field( default_factory=asyncio.Lock, init=False, ) + send_audio_scheduler: sched.scheduler = dataclasses.field( + default_factory=lambda: sched.scheduler(time.perf_counter, time.sleep), + init=False, + repr=False, + ) def __post_init__(self) -> None: fmt = self.media.fmt[0] @@ -216,16 +219,16 @@ async def send_audio(self, audio: np.ndarray) -> None: logger.warning("No remote RTP address for this call; dropping audio") return async with self.send_audio_lock: - for payload in self.codec.packetize(audio): - await asyncio.sleep( - max( - 0.0, - self.rpt_packet_duration.total_seconds() - - (time.perf_counter() - self.last_packet_time), - ) + loop = asyncio.get_running_loop() + t0 = time.perf_counter() + for i, payload in enumerate(self.codec.packetize(audio)): + self.send_audio_scheduler.enterabs( + t0 + i * self.rpt_packet_duration.total_seconds(), + 1, + self.send_packet, + (self.next_rtp_packet(payload), remote_addr), ) - self.send_packet(self.next_rtp_packet(payload), remote_addr) - self.last_packet_time = time.perf_counter() + await loop.run_in_executor(None, self.send_audio_scheduler.run) def audio_received(self, *, audio: np.ndarray, rms: float) -> None: """ From 4a61d5c0cc6e4c869fdc567306932070b4a0b06a Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 18 Mar 2026 14:08:06 +0100 Subject: [PATCH 2/7] Use loop.call_at --- voip/audio.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/voip/audio.py b/voip/audio.py index 6578233..be48c35 100644 --- a/voip/audio.py +++ b/voip/audio.py @@ -222,13 +222,12 @@ async def send_audio(self, audio: np.ndarray) -> None: loop = asyncio.get_running_loop() t0 = time.perf_counter() for i, payload in enumerate(self.codec.packetize(audio)): - self.send_audio_scheduler.enterabs( + loop.call_at( t0 + i * self.rpt_packet_duration.total_seconds(), - 1, self.send_packet, - (self.next_rtp_packet(payload), remote_addr), + self.next_rtp_packet(payload), + remote_addr, ) - await loop.run_in_executor(None, self.send_audio_scheduler.run) def audio_received(self, *, audio: np.ndarray, rms: float) -> None: """ From 58aa32dc4db463b8d2e4a4f8dae41617909c503d Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 18 Mar 2026 19:44:14 +0100 Subject: [PATCH 3/7] Use recrusive timerhandler loop --- tests/test_audio.py | 48 ++++++++++++++++++++++++++++++----------- voip/ai.py | 16 ++++++++++++++ voip/audio.py | 52 +++++++++++++++++++++++++++++++-------------- 3 files changed, 88 insertions(+), 28 deletions(-) diff --git a/tests/test_audio.py b/tests/test_audio.py index 2449c70..bd2954f 100644 --- a/tests/test_audio.py +++ b/tests/test_audio.py @@ -414,7 +414,7 @@ class TestSendRTPAudio: """Tests for AudioCall.send_rtp_audio.""" async def test_send_rtp_audio__sends_to_remote_addr(self): - """send_rtp_audio sends RTP packets to the caller's registered address.""" + """The first RTP packet is transmitted synchronously inside send_audio.""" call = make_audio_call(media=PCMU_MEDIA) remote_addr = ("10.0.0.1", 5004) call.rtp.calls = {remote_addr: call} @@ -442,24 +442,48 @@ async def test_send_rtp_audio__drops_audio_when_no_remote_addr(self, caplog): assert any("dropping audio" in r.message for r in caplog.records) async def test_send_rtp_audio__paces_packets_at_20ms_intervals(self): - """Packets are scheduled via sched.scheduler at fixed rpt_packet_duration intervals.""" + """Subsequent packets are scheduled at rpt_packet_duration intervals via call_later.""" call = make_audio_call(media=PCMU_MEDIA) remote_addr = ("10.0.0.2", 5006) call.rtp.calls = {remote_addr: call} - audio = np.zeros(320, dtype=np.float32) - with ( - patch.object(call, "send_audio_scheduler") as mock_scheduler, - patch.object(call, "send_packet"), - ): - await call.send_audio(audio) + with patch.object(call, "send_packet"): + await call.send_audio(np.zeros(320, dtype=np.float32)) - assert mock_scheduler.enterabs.call_count == 2 - scheduled_times = [c.args[0] for c in mock_scheduler.enterabs.call_args_list] - assert scheduled_times[1] - scheduled_times[0] == pytest.approx( - call.rpt_packet_duration.total_seconds() + loop = asyncio.get_event_loop() + assert call.outbound_handle is not None + assert call.outbound_handle.when() == pytest.approx( + loop.time() + call.rpt_packet_duration.total_seconds(), abs=0.01 ) + async def test_cancel_outbound_audio__cancels_handle_and_clears(self): + """cancel_outbound_audio cancels the pending handle and sets outbound_handle to None.""" + call = make_audio_call(media=PCMU_MEDIA) + remote_addr = ("10.0.0.1", 5004) + call.rtp.calls = {remote_addr: call} + + with patch.object(call, "send_packet"): + await call.send_audio(np.zeros(320, dtype=np.float32)) + + handle = call.outbound_handle + call.cancel_outbound_audio() + + assert handle.cancelled() + assert call.outbound_handle is None + + async def test_send_audio__preempts_pending_handle(self): + """A second send_audio cancels the pending handle from the first call.""" + call = make_audio_call(media=PCMU_MEDIA) + remote_addr = ("10.0.0.1", 5004) + call.rtp.calls = {remote_addr: call} + + with patch.object(call, "send_packet"): + await call.send_audio(np.zeros(320, dtype=np.float32)) + first_handle = call.outbound_handle + await call.send_audio(np.zeros(320, dtype=np.float32)) + + assert first_handle.cancelled() + def make_echo_call(**kwargs) -> EchoCall: """Create an EchoCall with mock rtp/sip for unit testing.""" diff --git a/voip/ai.py b/voip/ai.py index 76d9fad..fa3ecec 100644 --- a/voip/ai.py +++ b/voip/ai.py @@ -128,6 +128,9 @@ class AgentCall(TranscribeCall): _response_task: asyncio.Task | None = dataclasses.field( init=False, repr=False, default=None ) + _cancel_audio_handle: asyncio.Handle | None = dataclasses.field( + init=False, repr=False, default=None + ) def __post_init__(self) -> None: super().__post_init__() @@ -141,6 +144,7 @@ def __post_init__(self) -> None: ] def transcription_received(self, text: str) -> None: + self.cancel_outbound_audio() self._messages.append({"role": "user", "content": text}) if self._response_task is not None and not self._response_task.done(): self._response_task.cancel() @@ -167,3 +171,15 @@ async def send_speech(self, text: str) -> None: audio.numpy(), self.tts_model.sample_rate, self.codec.sample_rate_hz ) ) + + def on_audio_speech(self) -> None: + loop = asyncio.get_event_loop() + self._cancel_audio_handle = loop.call_later(0.5, self.cancel_outbound_audio) + super().on_audio_speech() + + def on_audio_silence(self) -> None: + super().on_audio_silence() + try: + self._cancel_audio_handle.cancel() + except AttributeError: + pass diff --git a/voip/audio.py b/voip/audio.py index be48c35..a26109d 100644 --- a/voip/audio.py +++ b/voip/audio.py @@ -16,9 +16,8 @@ import datetime import json import logging -import sched import secrets -import time +from collections.abc import Iterator from typing import ClassVar import numpy as np @@ -79,8 +78,8 @@ class AudioCall(RTPCall): default_factory=asyncio.Lock, init=False, ) - send_audio_scheduler: sched.scheduler = dataclasses.field( - default_factory=lambda: sched.scheduler(time.perf_counter, time.sleep), + outbound_handle: asyncio.TimerHandle | None = dataclasses.field( + default=None, init=False, repr=False, ) @@ -176,8 +175,8 @@ def decode_payload(self, payload: bytes) -> np.ndarray: def next_rtp_packet(self, payload: bytes) -> RTPPacket: packet = RTPPacket( payload_type=self.codec.payload_type, - sequence_number=self.rtp_sequence_number & 0xFFFF, - timestamp=self.rtp_timestamp & 0xFFFFFFFF, + sequence_number=self.rtp_sequence_number, + timestamp=self.rtp_timestamp, ssrc=self.rtp_ssrc, payload=payload, ) @@ -204,9 +203,37 @@ def rms(audio: np.ndarray) -> float: """ return float(np.sqrt(np.mean(np.square(audio)))) + def cancel_outbound_audio(self) -> None: + """Stop the current outbound audio while it is being sent.""" + try: + self.outbound_handle.cancel() + except AttributeError: + pass + else: + self.outbound_handle = None + + def _dispatch_next_packet( + self, + packets: Iterator[bytes], + remote_addr: tuple[str, int], + ) -> None: + loop = asyncio.get_running_loop() + try: + payload = next(packets) + except StopIteration: + self.outbound_handle = None + else: + self.send_packet(self.next_rtp_packet(payload), remote_addr) + self.outbound_handle = loop.call_later( + self.rpt_packet_duration.total_seconds(), + self._dispatch_next_packet, + packets, + remote_addr, + ) + async def send_audio(self, audio: np.ndarray) -> None: """ - Encode *audio* with the negotiated codec and transmit via RTP. + Encode `audio` with the negotiated codec and transmit via RTP. Args: audio: Float32 mono PCM at `codec.sample_rate_hz` Hz. @@ -219,15 +246,8 @@ async def send_audio(self, audio: np.ndarray) -> None: logger.warning("No remote RTP address for this call; dropping audio") return async with self.send_audio_lock: - loop = asyncio.get_running_loop() - t0 = time.perf_counter() - for i, payload in enumerate(self.codec.packetize(audio)): - loop.call_at( - t0 + i * self.rpt_packet_duration.total_seconds(), - self.send_packet, - self.next_rtp_packet(payload), - remote_addr, - ) + self.cancel_outbound_audio() + self._dispatch_next_packet(self.codec.packetize(audio), remote_addr) def audio_received(self, *, audio: np.ndarray, rms: float) -> None: """ From 7ae98b4d8e942fb55ca04963c36d63307ac370fb Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 18 Mar 2026 20:50:37 +0100 Subject: [PATCH 4/7] No dangling handlers --- voip/ai.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/voip/ai.py b/voip/ai.py index fa3ecec..0c40f06 100644 --- a/voip/ai.py +++ b/voip/ai.py @@ -174,7 +174,8 @@ async def send_speech(self, text: str) -> None: def on_audio_speech(self) -> None: loop = asyncio.get_event_loop() - self._cancel_audio_handle = loop.call_later(0.5, self.cancel_outbound_audio) + if self._cancel_audio_handle is None: + self._cancel_audio_handle = loop.call_later(0.5, self.cancel_outbound_audio) super().on_audio_speech() def on_audio_silence(self) -> None: @@ -183,3 +184,5 @@ def on_audio_silence(self) -> None: self._cancel_audio_handle.cancel() except AttributeError: pass + else: + self._cancel_audio_handle = None From f49c6c4a045f08684e7f75b9417e6032d1acec1f Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 18 Mar 2026 20:54:36 +0100 Subject: [PATCH 5/7] Apply suggestions from code review Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- voip/audio.py | 71 ++++++++++++++++++++++++++++++++++++++++++++++----- 1 file changed, 64 insertions(+), 7 deletions(-) diff --git a/voip/audio.py b/voip/audio.py index a26109d..5813d2d 100644 --- a/voip/audio.py +++ b/voip/audio.py @@ -21,6 +21,7 @@ from typing import ClassVar import numpy as np +import pytest import voip.codecs as codecs from voip.codecs import RTPCodec @@ -216,19 +217,23 @@ def _dispatch_next_packet( self, packets: Iterator[bytes], remote_addr: tuple[str, int], + next_send_at: float, ) -> None: - loop = asyncio.get_running_loop() try: payload = next(packets) except StopIteration: self.outbound_handle = None else: self.send_packet(self.next_rtp_packet(payload), remote_addr) - self.outbound_handle = loop.call_later( - self.rpt_packet_duration.total_seconds(), + duration_seconds = self.rpt_packet_duration.total_seconds() + next_deadline = next_send_at + duration_seconds + loop = asyncio.get_running_loop() + self.outbound_handle = loop.call_at( + next_deadline, self._dispatch_next_packet, packets, remote_addr, + next_deadline, ) async def send_audio(self, audio: np.ndarray) -> None: @@ -242,12 +247,23 @@ async def send_audio(self, audio: np.ndarray) -> None: (addr for addr, call in self.rtp.calls.items() if call is self), None, ) - if remote_addr is None: - logger.warning("No remote RTP address for this call; dropping audio") - return + match remote_addr: + case None: + logger.warning( + "No remote RTP address for this call; dropping audio", + ) + return + case _: + pass async with self.send_audio_lock: self.cancel_outbound_audio() - self._dispatch_next_packet(self.codec.packetize(audio), remote_addr) + loop = asyncio.get_running_loop() + next_send_at = loop.time() + self._dispatch_next_packet( + self.codec.packetize(audio), + remote_addr, + next_send_at, + ) def audio_received(self, *, audio: np.ndarray, rms: float) -> None: """ @@ -259,6 +275,47 @@ def audio_received(self, *, audio: np.ndarray, rms: float) -> None: """ +@pytest.mark.asyncio +async def test_send_audio_with_empty_packet_iterator_does_not_schedule_packets() -> None: + empty_audio = np.array([], dtype=np.float32) + + class EmptyPacketCodec: + def __init__(self) -> None: + self.payload_type = 0 + self.timestamp_increment = 160 + self.sample_rate_hz = 8000 + + def packetize(self, audio: np.ndarray) -> Iterator[bytes]: + return iter(()) + + class SingleCallRtp: + def __init__(self, call: AudioCall, remote_addr: tuple[str, int]) -> None: + self.calls = {remote_addr: call} + + call = object.__new__(AudioCall) + remote_addr = ("127.0.0.1", 4000) + codec = EmptyPacketCodec() + send_calls: list[tuple[RTPPacket, tuple[str, int]]] = [] + + def send_packet(packet: RTPPacket, addr: tuple[str, int]) -> None: + send_calls.append((packet, addr)) + + call.codec = codec + call.rtp_sequence_number = 0 + call.rtp_timestamp = 0 + call.rtp_ssrc = 1 + call.rpt_packet_duration = datetime.timedelta(milliseconds=20) + call.outbound_handle = None + call.rtp = SingleCallRtp(call, remote_addr) + call.send_audio_lock = asyncio.Lock() + call.send_packet = send_packet + + await AudioCall.send_audio(call, empty_audio) + + assert call.outbound_handle is None + assert send_calls == [] + + @dataclasses.dataclass(kw_only=True) class VoiceActivityCall(AudioCall): """ From d04e7489ee7cf7cbefdf263ae9279e87aa714ee8 Mon Sep 17 00:00:00 2001 From: "pre-commit-ci[bot]" <66853113+pre-commit-ci[bot]@users.noreply.github.com> Date: Wed, 18 Mar 2026 19:54:48 +0000 Subject: [PATCH 6/7] [pre-commit.ci] auto fixes from pre-commit.com hooks for more information, see https://pre-commit.ci --- voip/audio.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/voip/audio.py b/voip/audio.py index 5813d2d..5f20228 100644 --- a/voip/audio.py +++ b/voip/audio.py @@ -276,7 +276,9 @@ def audio_received(self, *, audio: np.ndarray, rms: float) -> None: @pytest.mark.asyncio -async def test_send_audio_with_empty_packet_iterator_does_not_schedule_packets() -> None: +async def test_send_audio_with_empty_packet_iterator_does_not_schedule_packets() -> ( + None +): empty_audio = np.array([], dtype=np.float32) class EmptyPacketCodec: From b71eb3f0086254636d4359951d00a75a8d07e7db Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 18 Mar 2026 21:18:50 +0100 Subject: [PATCH 7/7] Better emoji clearing --- voip/ai.py | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/voip/ai.py b/voip/ai.py index 0c40f06..7941346 100644 --- a/voip/ai.py +++ b/voip/ai.py @@ -12,6 +12,7 @@ import asyncio import dataclasses import logging +import re import typing import numpy as np @@ -132,6 +133,18 @@ class AgentCall(TranscribeCall): init=False, repr=False, default=None ) + emoji_pattern: typing.ClassVar[typing.Pattern[str]] = re.compile( + "[" + "\U0001f600-\U0001f64f" # emoticons + "\U0001f300-\U0001f5ff" # symbols & pictographs + "\U0001f680-\U0001f6ff" # transport & map symbols + "\U0001f1e0-\U0001f1ff" # flags (iOS) + "\U00002702-\U000027b0" + "\U000024c2-\U0001f251" + "]+", + flags=re.UNICODE, + ) + def __post_init__(self) -> None: super().__post_init__() self.tts_model = self.tts_model or TTSModel.load_model() @@ -156,7 +169,7 @@ async def respond(self) -> None: messages=self._messages, ) # clean non-ascii characters from the response for TTS processing - reply = (response.message.content or "").encode("ascii", "ignore").decode() + reply = self.emoji_pattern.sub("", response.message.content or "") self._messages.append({"role": "assistant", "content": reply}) logger.debug("Agent reply: %r", reply) await self.send_speech(reply)