Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
19 changes: 19 additions & 0 deletions src/ts6_stream_bot/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
50 changes: 41 additions & 9 deletions src/ts6_stream_bot/pipeline/stream_publisher.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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))
Expand All @@ -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
Expand Down
59 changes: 59 additions & 0 deletions tests/test_stream_publisher.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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 --------------------------------------------------------


Expand Down
Loading