From 77cc2ffeac603ea481000893ecbb62be85e0ee63 Mon Sep 17 00:00:00 2001 From: stefpi <19478336+stefpi@users.noreply.github.com> Date: Fri, 31 Jul 2026 15:00:56 -0700 Subject: [PATCH 1/4] fix hang --- teleoprtc/stream.py | 37 ++++++++++++++++++++++++++++++++++--- tests/test_stream.py | 14 ++++++++++++++ 2 files changed, 48 insertions(+), 3 deletions(-) diff --git a/teleoprtc/stream.py b/teleoprtc/stream.py index 37d793e..04c989f 100644 --- a/teleoprtc/stream.py +++ b/teleoprtc/stream.py @@ -43,6 +43,14 @@ class RTCSessionDescription: class WebRTCBaseStream(abc.ABC): + # libdatachannel resets a DataChannel's callbacks on an RTC worker. Destroying + # the Python wrapper concurrently resets them again while holding the GIL, + # which can deadlock with that worker acquiring the GIL to release callbacks. + # Retain closed wrappers for the daemon lifetime; the underlying channels are + # still closed by PeerConnection.close(). + _retained_messaging_channels: List[DataChannel] = [] + _retain_messaging_channel_on_close = False + def __init__(self, consumed_camera_types: List[str], consume_audio: bool, @@ -181,13 +189,32 @@ def on_message(message: Union[bytes, str]): for handler in list(self.incoming_message_handlers): self._call_soon_threadsafe(handler, message) + def on_open(): + self._set_event(self.messaging_channel_ready_event) + + def on_closed(): + self._set_event(self.connection_stopped_event) + channel.on_message(on_message) - channel.on_open(lambda: self._set_event(self.messaging_channel_ready_event)) - channel.on_closed(lambda: self._set_event(self.connection_stopped_event)) + channel.on_open(on_open) + channel.on_closed(on_closed) if channel.is_open(): self._set_event(self.messaging_channel_ready_event) self._on_after_media() + def _retain_messaging_channel(self) -> None: + if self.messaging_channel is None: + return + + # No native callback can be running before a remote description is set. + if not self.messaging_channel_ready_event.is_set() and self.peer_connection.remote_description() is None: + self.messaging_channel = None + return + + if self._retain_messaging_channel_on_close: + self._retained_messaging_channels.append(self.messaging_channel) + self.messaging_channel = None + def _on_connectionstatechange(self, state: PeerConnection.State): self._log_debug("connection state is %s", state) if state in (PeerConnection.State.Connected, PeerConnection.State.Failed): @@ -355,8 +382,8 @@ async def stop(self): with contextlib.suppress(asyncio.CancelledError): await task self._sender_tasks.clear() + self._retain_messaging_channel() self.peer_connection.close() - self.messaging_channel = None self.incoming_camera_tracks.clear() self.incoming_audio_tracks.clear() self._consumer_tracks.clear() @@ -398,6 +425,10 @@ async def start(self) -> RTCSessionDescription: class WebRTCAnswerStream(WebRTCBaseStream): + # Incoming DataChannels are closed from an RTC worker when the remote offerer + # disconnects, which is the path affected by the callback/GIL deadlock. + _retain_messaging_channel_on_close = True + def __init__(self, session: RTCSessionDescription, *args, **kwargs): super().__init__(*args, **kwargs) self.session = session diff --git a/tests/test_stream.py b/tests/test_stream.py index ca74e4e..f2d52be 100755 --- a/tests/test_stream.py +++ b/tests/test_stream.py @@ -62,6 +62,20 @@ async def test_offer_stream_sdp_channel(self): info = parse_info_from_offer(capture.offer.sdp) assert info.incoming_datachannel + async def test_answer_side_messaging_channel_wrapper_retained_after_stop(self): + stream = WebRTCOfferBuilder(OfferCapture()).stream() + stream._retain_messaging_channel_on_close = True + channel = stream.peer_connection.create_data_channel("data") + stream.messaging_channel = channel + stream.messaging_channel_ready_event.set() + + try: + await stream.stop() + assert stream.messaging_channel is None + assert stream._retained_messaging_channels[-1] is channel + finally: + stream._retained_messaging_channels.pop() + @pytest.mark.asyncio class TestAnswerStream: From d09659ab8c0e250e2d3360a7a299c36fd4420874 Mon Sep 17 00:00:00 2001 From: stefpi <19478336+stefpi@users.noreply.github.com> Date: Mon, 3 Aug 2026 16:31:36 -0700 Subject: [PATCH 2/4] clean --- teleoprtc/stream.py | 9 ++------- 1 file changed, 2 insertions(+), 7 deletions(-) diff --git a/teleoprtc/stream.py b/teleoprtc/stream.py index 04c989f..cc4f009 100644 --- a/teleoprtc/stream.py +++ b/teleoprtc/stream.py @@ -43,11 +43,8 @@ class RTCSessionDescription: class WebRTCBaseStream(abc.ABC): - # libdatachannel resets a DataChannel's callbacks on an RTC worker. Destroying - # the Python wrapper concurrently resets them again while holding the GIL, - # which can deadlock with that worker acquiring the GIL to release callbacks. - # Retain closed wrappers for the daemon lifetime; the underlying channels are - # still closed by PeerConnection.close(). + # destorying wrapper on close can cause deadlock + # TODO: upstream a fix to this _retained_messaging_channels: List[DataChannel] = [] _retain_messaging_channel_on_close = False @@ -425,8 +422,6 @@ async def start(self) -> RTCSessionDescription: class WebRTCAnswerStream(WebRTCBaseStream): - # Incoming DataChannels are closed from an RTC worker when the remote offerer - # disconnects, which is the path affected by the callback/GIL deadlock. _retain_messaging_channel_on_close = True def __init__(self, session: RTCSessionDescription, *args, **kwargs): From 1598f491500f9545a792c4e603fb6ddcbc09de4f Mon Sep 17 00:00:00 2001 From: stefpi <19478336+stefpi@users.noreply.github.com> Date: Wed, 5 Aug 2026 13:51:40 -0700 Subject: [PATCH 3/4] audio support --- teleoprtc/stream.py | 120 +++++++++++++++++++++++++++---- tests/test_stream.py | 163 ++++++++++++++++++++++++++++++++++++++++++- 2 files changed, 270 insertions(+), 13 deletions(-) diff --git a/teleoprtc/stream.py b/teleoprtc/stream.py index cc4f009..2a3c1cb 100644 --- a/teleoprtc/stream.py +++ b/teleoprtc/stream.py @@ -14,9 +14,12 @@ H264RtpPacketizer, IceServer, NalUnit, + OpusRtpPacketizer, + OpusRtpDepacketizer, PeerConnection, PliHandler, RtcpNackResponder, + RtcpReceivingSession, RtcpSrReporter, RtpPacketizationConfig, Track, @@ -76,8 +79,10 @@ def __init__(self, self.messaging_channel: Optional[DataChannel] = None self.incoming_message_handlers: List[MessageHandler] = [] self._consumer_tracks: List[Track] = [] + self._offered_tracks: Dict[str, Track] = {} + self._incoming_audio_handlers: List[Any] = [] self._sender_tasks: List[asyncio.Task] = [] - self._track_state: List[Tuple[Track, TiciVideoStreamTrack, RtpPacketizationConfig]] = [] + self._track_state: List[Tuple[Track, Any, RtpPacketizationConfig]] = [] self._receiver_reports: Dict[str, RtcpReceiverReport] = {} self._receiver_report_tracks: Dict[str, Tuple[Track, int]] = {} @@ -125,9 +130,17 @@ def _add_consumer_transceivers(self): media = Description.Audio("audio", Description.Direction.RecvOnly) media.add_opus_codec(111) track = self.peer_connection.add_track(media) + self._set_incoming_audio_handlers(track) self._consumer_tracks.append(track) self.incoming_audio_tracks.append(track) + def _set_incoming_audio_handlers(self, track: Track) -> None: + depacketizer = OpusRtpDepacketizer() + rtcp_session = RtcpReceivingSession() + depacketizer.add_to_chain(rtcp_session) + track.set_media_handler(depacketizer) + self._incoming_audio_handlers.extend((depacketizer, rtcp_session)) + def _find_offer_video(self, remote_sdp: str, used_mids: set[str]) -> Tuple[str, int]: desc = Description(remote_sdp, Description.Type.Offer) for i in range(desc.media_count()): @@ -141,6 +154,31 @@ def _find_offer_video(self, remote_sdp: str, used_mids: set[str]) -> Tuple[str, return media.mid(), payload_type raise ValueError("Remote SDP does not offer H264 video") + def _find_offer_audio(self, remote_sdp: str, used_mids: set[str]) -> Tuple[str, int, Description.Direction]: + desc = Description(remote_sdp, Description.Type.Offer) + for i in range(desc.media_count()): + media = desc.media(i) + if media is None or media.type() != "audio" or media.mid() in used_mids: + continue + if media.direction() not in (Description.Direction.RecvOnly, Description.Direction.SendRecv): + continue + for payload_type in media.payload_types(): + with contextlib.suppress(ValueError): + rtp_map = media.rtp_map(payload_type) + if rtp_map is not None and rtp_map.format.upper() == "OPUS": + return media.mid(), payload_type, media.direction() + raise ValueError("Remote SDP does not offer Opus audio") + + def _find_track_h264(self, track: Track) -> Tuple[int, str, Optional[str]]: + media = track.description() + for payload_type in media.payload_types(): + with contextlib.suppress(ValueError): + rtp_map = media.rtp_map(payload_type) + if rtp_map.format.upper() == "H264": + profile = rtp_map.fmtps[0] if rtp_map.fmtps else None + return payload_type, rtp_map.format, profile + raise ValueError("Track does not offer H264 video") + def _make_video_media(self, track: TiciVideoStreamTrack, remote_sdp: str, used_mids: set[str]) -> Tuple[Description.Video, int, int, str]: mid, payload_type = self._find_offer_video(remote_sdp, used_mids) used_mids.add(mid) @@ -156,7 +194,13 @@ def _add_producer_tracks(self, remote_sdp: Optional[str] = None): used_mids: set[str] = set() for track in self.outgoing_video_tracks: media, ssrc, payload_type, cname = self._make_video_media(track, remote_sdp or "", used_mids) - rtc_track = self.peer_connection.add_track(media) + rtc_track = self._offered_tracks.pop(media.mid(), None) + if rtc_track is None: + rtc_track = self.peer_connection.add_track(media) + else: + offered_media = rtc_track.description() + offered_media.add_ssrc(ssrc, cname, f"stream-{random.getrandbits(32):08x}", track.id) + rtc_track.set_description(offered_media) rtp_config = RtpPacketizationConfig(ssrc, cname, payload_type, H264RtpPacketizer.CLOCK_RATE) rtp_config.start_timestamp = random.randint(0, 0xFFFFFFFF) @@ -174,8 +218,49 @@ def _add_producer_tracks(self, remote_sdp: Optional[str] = None): self._receiver_report_tracks[camera_type] = (rtc_track, ssrc) self._track_state.append((rtc_track, track, rtp_config)) - if self.outgoing_audio_tracks: - raise NotImplementedError("Audio producer tracks are not implemented with libdatachannel") + for track in self.outgoing_audio_tracks: + if remote_sdp is None: + mid, payload_type = "audio", 111 + direction = Description.Direction.SendRecv if self.expected_incoming_audio else Description.Direction.SendOnly + else: + mid, payload_type, offered_direction = self._find_offer_audio(remote_sdp, used_mids) + direction = Description.Direction.SendRecv if offered_direction == Description.Direction.SendRecv else Description.Direction.SendOnly + used_mids.add(mid) + + ssrc = random.randint(1, 0xFFFFFFFF) + cname = f"teleoprtc-{random.getrandbits(32):08x}" + stream_id = f"stream-{random.getrandbits(32):08x}" + media = Description.Audio(mid, direction) + media.add_opus_codec(payload_type) + media.add_ssrc(ssrc, cname, stream_id, track.id) + rtc_track = self._offered_tracks.pop(mid, None) + if rtc_track is not None: + rtc_track.set_description(media) + elif remote_sdp is None and self.expected_incoming_audio: + rtc_track = self.incoming_audio_tracks[0] + rtc_track.set_description(media) + else: + rtc_track = self.peer_connection.add_track(media) + + rtp_config = RtpPacketizationConfig(ssrc, cname, payload_type, OpusRtpPacketizer.DEFAULT_CLOCK_RATE) + rtp_config.start_timestamp = random.randint(0, 0xFFFFFFFF) + rtp_config.timestamp = rtp_config.start_timestamp + rtp_config.sequence_number = random.randint(0, 0xFFFF) + + packetizer = OpusRtpPacketizer(rtp_config) + packetizer.add_to_chain(RtcpSrReporter(rtp_config)) + packetizer.add_to_chain(RtcpNackResponder()) + rtc_track.set_media_handler(packetizer) + self._track_state.append((rtc_track, track, rtp_config)) + + for mid, rtc_track in self._offered_tracks.items(): + if rtc_track.description().type() == "video": + # libdatachannel creates local tracks for every remote recvonly section. + # Keep compatibility video MIDs valid, but do not advertise empty streams. + payload_type, codec, profile = self._find_track_h264(rtc_track) + media = Description.Video(mid, Description.Direction.Inactive) + media.add_video_codec(payload_type, codec, profile) + rtc_track.set_description(media) def _add_messaging_channel(self, channel: Optional[DataChannel] = None): if channel is None: @@ -226,14 +311,19 @@ def _on_gatheringstatechange(self, state: PeerConnection.GatheringState): def _on_incoming_track(self, track: Track): self._log_debug("got track: %s", track.mid()) - try: - camera_type, _ = parse_video_track_id(track.mid()) - except ValueError: - camera_type = track.mid() - if camera_type in self.expected_incoming_camera_types: - self.incoming_camera_tracks[camera_type] = track - elif self.expected_incoming_audio: + # An offer-created track may be the same transceiver used by an outgoing + # producer. Reusing it avoids adding a duplicate track with the same MID. + self._offered_tracks[track.mid()] = track + if track.description().type() == "audio" and self.expected_incoming_audio: + self._set_incoming_audio_handlers(track) self.incoming_audio_tracks.append(track) + elif track.description().type() == "video": + try: + camera_type, _ = parse_video_track_id(track.mid()) + except ValueError: + camera_type = track.mid() + if camera_type in self.expected_incoming_camera_types: + self.incoming_camera_tracks[camera_type] = track self._on_after_media() def _on_incoming_datachannel(self, channel: DataChannel): @@ -312,7 +402,7 @@ async def _wait_for_gathering_complete(self): if self.peer_connection.gathering_state() != PeerConnection.GatheringState.Complete: await self.gathering_complete_event.wait() - async def _send_track_loop(self, rtc_track: Track, producer_track: TiciVideoStreamTrack, rtp_config: RtpPacketizationConfig): + async def _send_track_loop(self, rtc_track: Track, producer_track: Any, rtp_config: RtpPacketizationConfig): while True: if not rtc_track.is_open(): await asyncio.sleep(0.01) @@ -325,6 +415,9 @@ async def _send_track_loop(self, rtc_track: Track, producer_track: TiciVideoStre continue pts = int(packet.pts or 0) + time_base = getattr(packet, "time_base", None) + if time_base is not None: + pts = int(pts * time_base * rtp_config.clock_rate) timestamp = (rtp_config.start_timestamp + pts) & 0xFFFFFFFF rtc_track.send_frame(data, FrameInfo(timestamp)) except asyncio.CancelledError: @@ -384,6 +477,8 @@ async def stop(self): self.incoming_camera_tracks.clear() self.incoming_audio_tracks.clear() self._consumer_tracks.clear() + self._offered_tracks.clear() + self._incoming_audio_handlers.clear() self._track_state.clear() self._receiver_reports.clear() self._receiver_report_tracks.clear() @@ -403,6 +498,7 @@ async def start(self) -> RTCSessionDescription: self._add_consumer_transceivers() if self.should_add_data_channel: self._add_messaging_channel() + self._add_producer_tracks() self.peer_connection.set_local_description(Description.Type.Offer) await self._wait_for_gathering_complete() diff --git a/tests/test_stream.py b/tests/test_stream.py index f2d52be..1fd09f8 100755 --- a/tests/test_stream.py +++ b/tests/test_stream.py @@ -4,7 +4,7 @@ import pytest -from libdatachannel import Description +from libdatachannel import Description, OpusRtpDepacketizer, RtcpReceivingSession from teleoprtc.builder import WebRTCOfferBuilder, WebRTCAnswerBuilder from teleoprtc.info import parse_info_from_offer @@ -27,6 +27,14 @@ async def recv(self): raise NotImplementedError() +class DummyOpusAudioStreamTrack: + kind = "audio" + id = "audio-track" + + async def recv(self): + raise NotImplementedError() + + @pytest.mark.asyncio class TestOfferStream: async def test_offer_stream_sdp_recvonly_audio(self): @@ -46,6 +54,40 @@ async def test_offer_stream_sdp_recvonly_audio(self): assert info.expected_audio_track assert not info.incoming_audio_track + async def test_offer_stream_sdp_sendonly_audio(self): + capture = OfferCapture() + builder = WebRTCOfferBuilder(capture) + builder.add_audio_stream(DummyOpusAudioStreamTrack()) + stream = builder.stream() + + try: + with contextlib.suppress(Exception): + await stream.start() + finally: + await stream.stop() + + info = parse_info_from_offer(capture.offer.sdp) + assert not info.expected_audio_track + assert info.incoming_audio_track + + async def test_offer_stream_sdp_sendrecv_audio(self): + capture = OfferCapture() + builder = WebRTCOfferBuilder(capture) + builder.offer_to_receive_audio_stream() + builder.add_audio_stream(DummyOpusAudioStreamTrack()) + stream = builder.stream() + + try: + with contextlib.suppress(Exception): + await stream.start() + finally: + await stream.stop() + + desc = Description(capture.offer.sdp, Description.Type.Offer) + audio = next(desc.media(i) for i in range(desc.media_count()) if desc.media(i).type() == "audio") + assert audio.direction() == Description.Direction.SendRecv + assert len(audio.get_ssrcs()) == 1 + async def test_offer_stream_sdp_channel(self): capture = OfferCapture() builder = WebRTCOfferBuilder(capture) @@ -79,6 +121,125 @@ async def test_answer_side_messaging_channel_wrapper_retained_after_stop(self): @pytest.mark.asyncio class TestAnswerStream: + async def test_video_answer_when_receiving_audio(self): + offer_sdp = """v=0 +o=- 3910274679 3910274679 IN IP4 0.0.0.0 +s=- +t=0 0 +a=group:BUNDLE 0 1 2 +a=msid-semantic:WMS * +m=video 9 UDP/TLS/RTP/SAVPF 103 +c=IN IP4 0.0.0.0 +a=recvonly +a=mid:0 +a=rtcp-mux +a=rtpmap:103 H264/90000 +a=rtcp-fb:103 nack +a=rtcp-fb:103 nack pli +a=fmtp:103 level-asymmetry-allowed=1;packetization-mode=1;profile-level-id=42001f +a=ice-ufrag:1234 +a=ice-pwd:1234 +a=fingerprint:sha-256 15:F3:F0:23:67:44:EE:2C:AA:8C:D9:50:95:26:42:7C:67:EA:1F:D2:92:C5:97:01:7B:2E:57:C9:A3:13:00:4A +a=setup:actpass +m=video 9 UDP/TLS/RTP/SAVPF 103 +c=IN IP4 0.0.0.0 +a=recvonly +a=mid:1 +a=rtcp-mux +a=rtpmap:103 H264/90000 +a=rtcp-fb:103 nack +a=rtcp-fb:103 nack pli +a=fmtp:103 level-asymmetry-allowed=1;packetization-mode=1;profile-level-id=42001f +a=ice-ufrag:1234 +a=ice-pwd:1234 +a=fingerprint:sha-256 15:F3:F0:23:67:44:EE:2C:AA:8C:D9:50:95:26:42:7C:67:EA:1F:D2:92:C5:97:01:7B:2E:57:C9:A3:13:00:4A +a=setup:actpass +m=audio 9 UDP/TLS/RTP/SAVPF 111 +c=IN IP4 0.0.0.0 +a=sendonly +a=mid:2 +a=rtcp-mux +a=rtpmap:111 opus/48000/2 +a=ice-ufrag:1234 +a=ice-pwd:1234 +a=fingerprint:sha-256 15:F3:F0:23:67:44:EE:2C:AA:8C:D9:50:95:26:42:7C:67:EA:1F:D2:92:C5:97:01:7B:2E:57:C9:A3:13:00:4A +a=setup:actpass""" + builder = WebRTCAnswerBuilder(offer_sdp) + builder.offer_to_receive_audio_stream() + builder.add_video_stream("road", DummyH264VideoStreamTrack("road", 0.05)) + stream = builder.stream() + try: + answer = await stream.start() + desc = Description(answer.sdp, Description.Type.Answer) + videos = [desc.media(i) for i in range(desc.media_count()) if desc.media(i).type() == "video"] + audio = next(desc.media(i) for i in range(desc.media_count()) if desc.media(i).type() == "audio") + assert [video.direction() for video in videos] == [Description.Direction.SendOnly, Description.Direction.Inactive] + assert len(videos[0].get_ssrcs()) == 1 + assert len(videos[1].get_ssrcs()) == 0 + assert audio.direction() == Description.Direction.RecvOnly + assert stream.has_incoming_audio_track() + handler = stream.get_incoming_audio_track().get_media_handler() + assert isinstance(handler, OpusRtpDepacketizer) + assert isinstance(handler.next(), RtcpReceivingSession) + finally: + await stream.stop() + + async def test_receive_opus_audio_track(self): + offer_sdp = """v=0 +o=- 1 1 IN IP4 0.0.0.0 +s=- +t=0 0 +a=group:BUNDLE audio +m=audio 9 UDP/TLS/RTP/SAVPF 111 +c=IN IP4 0.0.0.0 +a=sendonly +a=mid:audio +a=rtpmap:111 opus/48000/2 +a=ice-ufrag:1234 +a=ice-pwd:1234 +a=fingerprint:sha-256 15:F3:F0:23:67:44:EE:2C:AA:8C:D9:50:95:26:42:7C:67:EA:1F:D2:92:C5:97:01:7B:2E:57:C9:A3:13:00:4A +a=setup:actpass""" + builder = WebRTCAnswerBuilder(offer_sdp) + builder.offer_to_receive_audio_stream() + stream = builder.stream() + try: + answer = await stream.start() + desc = Description(answer.sdp, Description.Type.Answer) + audio = next(desc.media(i) for i in range(desc.media_count()) if desc.media(i).type() == "audio") + assert audio.direction() == Description.Direction.RecvOnly + assert stream.has_incoming_audio_track() + finally: + await stream.stop() + + async def test_opus_audio_track(self): + offer_sdp = """v=0 +o=- 1 1 IN IP4 0.0.0.0 +s=- +t=0 0 +a=group:BUNDLE audio +m=audio 9 UDP/TLS/RTP/SAVPF 111 0 +c=IN IP4 0.0.0.0 +a=recvonly +a=mid:audio +a=rtpmap:111 opus/48000/2 +a=rtpmap:0 PCMU/8000 +a=ice-ufrag:1234 +a=ice-pwd:1234 +a=fingerprint:sha-256 15:F3:F0:23:67:44:EE:2C:AA:8C:D9:50:95:26:42:7C:67:EA:1F:D2:92:C5:97:01:7B:2E:57:C9:A3:13:00:4A +a=setup:actpass""" + builder = WebRTCAnswerBuilder(offer_sdp) + builder.add_audio_stream(DummyOpusAudioStreamTrack()) + stream = builder.stream() + try: + answer = await stream.start() + desc = Description(answer.sdp, Description.Type.Answer) + audio = next(desc.media(i) for i in range(desc.media_count()) if desc.media(i).type() == "audio") + assert audio.direction() == Description.Direction.SendOnly + assert audio.rtp_map(audio.payload_types()[0]).format.lower() == "opus" + assert len(audio.get_ssrcs()) == 1 + finally: + await stream.stop() + async def test_codec_preference(self): offer_sdp = """v=0 o=- 3910274679 3910274679 IN IP4 0.0.0.0 From b5631eff9dae4d5dee0b485bfcc0d64342bae942 Mon Sep 17 00:00:00 2001 From: stefpi <19478336+stefpi@users.noreply.github.com> Date: Thu, 6 Aug 2026 18:27:30 -0700 Subject: [PATCH 4/4] send-recv --- teleoprtc/stream.py | 14 +++++++++++++- 1 file changed, 13 insertions(+), 1 deletion(-) diff --git a/teleoprtc/stream.py b/teleoprtc/stream.py index 2a3c1cb..80c9779 100644 --- a/teleoprtc/stream.py +++ b/teleoprtc/stream.py @@ -250,7 +250,19 @@ def _add_producer_tracks(self, remote_sdp: Optional[str] = None): packetizer = OpusRtpPacketizer(rtp_config) packetizer.add_to_chain(RtcpSrReporter(rtp_config)) packetizer.add_to_chain(RtcpNackResponder()) - rtc_track.set_media_handler(packetizer) + if direction == Description.Direction.SendRecv: + # MediaHandler runs outgoing chains front-to-back and incoming chains + # back-to-front. Put receive handlers first so incoming RTP is handled + # by sender RTCP handlers while still packetized, then depacketized as + # the final operation before delivery to Track.receive(). + depacketizer = OpusRtpDepacketizer() + rtcp_session = RtcpReceivingSession() + depacketizer.add_to_chain(rtcp_session) + depacketizer.add_to_chain(packetizer) + rtc_track.set_media_handler(depacketizer) + self._incoming_audio_handlers.extend((depacketizer, rtcp_session)) + else: + rtc_track.set_media_handler(packetizer) self._track_state.append((rtc_track, track, rtp_config)) for mid, rtc_track in self._offered_tracks.items():