From 183097ba49e55a64edc7bbb098ae1203ec9d3ef8 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 1 May 2026 19:37:03 +0000 Subject: [PATCH] Buffer ICE candidates until setRemoteDescription completes The "stuck on CONNECTING" symptom: client UI never transitions to CONNECTED even with STUN configured. Root cause was hiding in a warning we'd been ignoring - RTCPeerConnection addIceCandidate called without remote description aiortc's addIceCandidate REJECTS the call when no remote description is set yet, just emits that warning, and the candidate is gone. The publisher was spawning the answer task and per-candidate tasks as independent asyncio tasks; the candidate tasks raced ahead of the answer task's setRemoteDescription and got dropped on the floor. Every successful ICE we've seen so far was a NAT-punch lottery on the remaining candidates that happened to arrive after the answer settled. Fix: per-Viewer asyncio.Event (`remote_set`) flips when setRemoteDescription completes. _apply_ice_candidate awaits it (10 s cap, in case the answer never lands) before calling addIceCandidate. Candidates now flow correctly regardless of arrival order. Also done while we're in this neighborhood: - STUN_URL / TURN_URL / TURN_USERNAME / TURN_PASSWORD env-overridable. Default STUN stays at stun.l.google.com:19302; TURN empty. If a symmetric NAT or CGNAT prevents direct punching, the operator can point at a free public TURN (openrelay.metered.ca) without code changes. Tests: 11 publisher cases pass including a new regression that emits ICE BEFORE the answer and asserts addIceCandidate is held until the answer arrives. ruff + mypy strict + 219 tests overall. https://claude.ai/code/session_016DuCjRJK995Tj9aDhhB9at --- .env.example | 10 ++++ src/ts6_stream_bot/config.py | 19 ++++++ .../pipeline/stream_publisher.py | 50 +++++++++++++--- tests/test_stream_publisher.py | 59 +++++++++++++++++++ 4 files changed, 129 insertions(+), 9 deletions(-) diff --git a/.env.example b/.env.example index 73cb67c..3b4de2d 100644 --- a/.env.example +++ b/.env.example @@ -46,3 +46,13 @@ STREAM_BITRATE=4608 STREAM_ACCESSIBILITY=0 STREAM_MODE=1 STREAM_VIEWER_LIMIT=-1 + +# --- WebRTC ICE servers ---------------------------------------------------- +# STUN exposes the bot's public-NAT'd IP as a server-reflexive candidate; +# TURN relays media when direct NAT punching fails. The defaults work for +# most home setups. If your host is behind symmetric / CGNAT and viewers +# can't connect, point TURN_URL at a free relay (or your own coturn). +STUN_URL=stun:stun.l.google.com:19302 +# TURN_URL=turn:openrelay.metered.ca:80 +# TURN_USERNAME=openrelayproject +# TURN_PASSWORD=openrelayproject diff --git a/src/ts6_stream_bot/config.py b/src/ts6_stream_bot/config.py index 1ec7ba0..3d616c4 100644 --- a/src/ts6_stream_bot/config.py +++ b/src/ts6_stream_bot/config.py @@ -87,6 +87,25 @@ class Settings(BaseSettings): description="Max viewers; -1 = unlimited (TS3 convention).", ) + # --- WebRTC ICE ------------------------------------------------------- + # STUN exposes the bot's public-NAT'd address as a server-reflexive + # candidate. TURN relays media through a third-party server when + # direct NAT punching fails - only needed if your host's NAT is too + # restrictive (CGNAT etc.). + STUN_URL: str = Field( + default="stun:stun.l.google.com:19302", + description="STUN server URL (use 'stun:host:port'). Empty = no STUN.", + ) + TURN_URL: str = Field( + default="", + description=( + "Optional TURN server URL (e.g. 'turn:turn.example.com:3478'). " + "Required when both peers are behind symmetric NAT." + ), + ) + TURN_USERNAME: str = Field(default="", description="TURN auth username.") + TURN_PASSWORD: str = Field(default="", description="TURN auth password.") + @field_validator("BOT_API_KEY") @classmethod def _reject_insecure_api_key(cls, v: str) -> str: diff --git a/src/ts6_stream_bot/pipeline/stream_publisher.py b/src/ts6_stream_bot/pipeline/stream_publisher.py index 63e979b..3590aa0 100644 --- a/src/ts6_stream_bot/pipeline/stream_publisher.py +++ b/src/ts6_stream_bot/pipeline/stream_publisher.py @@ -36,6 +36,7 @@ from aiortc.contrib.media import MediaRelay from aiortc.sdp import candidate_from_sdp +from ts6_stream_bot.config import settings from ts6_stream_bot.pipeline.stream_signaling import ( SignalingMessage, SignalingType, @@ -52,6 +53,11 @@ class _Viewer: clid: int pc: RTCPeerConnection joined_at: float + # remote_set fires once setRemoteDescription has completed for this + # viewer. ICE candidates arriving before that event hold off in + # `_apply_ice_candidate` instead of being silently rejected by + # aiortc with "addIceCandidate called without remote description". + remote_set: asyncio.Event = field(default_factory=asyncio.Event) @dataclass(slots=True) @@ -262,16 +268,23 @@ async def _handle_viewer_join(self, viewer_clid: int, stream_id: str) -> None: with contextlib.suppress(Exception): await existing.pc.close() - # STUN gives the bot a server-reflexive (srflx) candidate so its own - # NAT/Docker-bridge address is publishable. Without it, NAT punching - # depended entirely on the viewer's side reaching us, and the live - # trace showed the first connection attempt timing out on - # ice=checking before completing on the retry. Google's public STUN - # is the conventional default. - pc = RTCPeerConnection( - configuration=RTCConfiguration( - iceServers=[RTCIceServer(urls="stun:stun.l.google.com:19302")], + # STUN exposes the bot's public-NAT'd address as a server-reflexive + # candidate. TURN relays media when direct NAT punching fails - + # both are env-overridable in case the operator's network needs + # something other than Google's public STUN / no-TURN default. + ice_servers: list[RTCIceServer] = [] + if settings.STUN_URL: + ice_servers.append(RTCIceServer(urls=settings.STUN_URL)) + if settings.TURN_URL: + ice_servers.append( + RTCIceServer( + urls=settings.TURN_URL, + username=settings.TURN_USERNAME or None, + credential=settings.TURN_PASSWORD or None, + ) ) + pc = RTCPeerConnection( + configuration=RTCConfiguration(iceServers=ice_servers) if ice_servers else None ) # Surface aiortc's own connection-state lifecycle so we can tell @@ -340,6 +353,11 @@ async def _apply_answer(self, viewer_clid: int, answer_sdp: str) -> None: await viewer.pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer") ) + # Unblock any ICE candidates that arrived first - aiortc rejects + # addIceCandidate calls before setRemoteDescription, and the + # signaling layer has no ordering guarantee between the answer + # task and the per-candidate tasks once they're both spawned. + viewer.remote_set.set() log.info("stream_publisher.viewer_connected", clid=viewer_clid) except Exception as exc: log.exception("stream_publisher.set_answer_failed", clid=viewer_clid, error=str(exc)) @@ -351,6 +369,20 @@ async def _apply_ice_candidate( viewer = self._viewers.get(viewer_clid) if viewer is None: return + # Wait until setRemoteDescription has completed; aiortc otherwise + # rejects the call with "addIceCandidate called without remote + # description" and the candidate is lost. The early candidates are + # often the only ones that include public srflx info, so losing them + # turned every previous ICE attempt into a NAT-punch lottery. + try: + await asyncio.wait_for(viewer.remote_set.wait(), timeout=10.0) + except TimeoutError: + log.warning( + "stream_publisher.ice_wait_timeout", + clid=viewer_clid, + detail="remote description never arrived", + ) + return try: candidate = candidate_from_sdp(candidate_sdp) candidate.sdpMid = sdp_mid diff --git a/tests/test_stream_publisher.py b/tests/test_stream_publisher.py index 934038e..5269058 100644 --- a/tests/test_stream_publisher.py +++ b/tests/test_stream_publisher.py @@ -324,6 +324,18 @@ async def test_ice_candidate_is_forwarded_to_pc(wired, monkeypatch) -> None: if pcs: break + # ICE candidates now wait for the answer to be applied (so aiortc's + # "addIceCandidate called without remote description" reject doesn't + # eat them). Send the answer first - that flips the per-viewer event + # which unblocks queued ICE work. + sig.on_signaling_message( # type: ignore[misc] + SignalingMessage(type=SignalingType.ANSWER, sdp="ANSWER_SDP", clid=42, raw="") + ) + for _ in range(20): + await asyncio.sleep(0.01) + if pcs[0].setRemoteDescription.await_count > 0: + break + sig.on_signaling_message( # type: ignore[misc] SignalingMessage( type=SignalingType.ICE_CANDIDATE, @@ -344,6 +356,53 @@ async def test_ice_candidate_is_forwarded_to_pc(wired, monkeypatch) -> None: assert fake_candidate.sdpMLineIndex == 0 +@pytest.mark.asyncio +async def test_ice_candidate_arriving_before_answer_is_buffered(wired, monkeypatch) -> None: + """Regression for the live deploy where every ICE candidate arrived + before the answer task finished and aiortc dropped the call. Now + candidates wait for the answer via the per-viewer remote_set event.""" + publisher, client, sig, _capture, pcs = wired + publisher._stream_id = "s1" + fake_candidate = MagicMock() + monkeypatch.setattr( + "ts6_stream_bot.pipeline.stream_publisher.candidate_from_sdp", + lambda s: fake_candidate, + ) + + _emit(client, "notifyjoinstreamrequest id=s1 clid=42") + for _ in range(20): + await asyncio.sleep(0.01) + if pcs: + break + + # ICE arrives FIRST (the racy ordering from the live trace). + sig.on_signaling_message( # type: ignore[misc] + SignalingMessage( + type=SignalingType.ICE_CANDIDATE, + candidate="candidate:abc", + sdp_mid="0", + sdp_mline_index=0, + clid=42, + raw="", + ) + ) + # Give the task a chance to run; it should NOT call addIceCandidate yet. + for _ in range(5): + await asyncio.sleep(0.01) + assert pcs[0].addIceCandidate.await_count == 0 + + # Now the answer lands and unblocks the queued candidate. + sig.on_signaling_message( # type: ignore[misc] + SignalingMessage(type=SignalingType.ANSWER, sdp="ANSWER_SDP", clid=42, raw="") + ) + for _ in range(40): + await asyncio.sleep(0.01) + if pcs[0].addIceCandidate.await_count > 0: + break + + assert pcs[0].addIceCandidate.await_count == 1 + + # --- leave + stop --------------------------------------------------------