From 620d21eb22d68c911c1c151193f6155136a94abc Mon Sep 17 00:00:00 2001 From: Stackie Jia Date: Sun, 19 Jul 2026 19:28:10 +0800 Subject: [PATCH 1/3] feat(stream): expose upstream metadata headers --- docs/en/guide/fcc-setup.md | 2 + docs/en/guide/url-formats.md | 25 +++ docs/en/guide/video-snapshot.md | 2 + docs/guide/fcc-setup.md | 2 + docs/guide/url-formats.md | 25 +++ docs/guide/video-snapshot.md | 2 + e2e/helpers/mock_rtsp.py | 18 +- e2e/helpers/rtp.py | 4 +- e2e/test_config.py | 3 +- e2e/test_fcc.py | 47 ++++- e2e/test_http_proxy.py | 1 + e2e/test_multicast.py | 73 ++++++- e2e/test_rtsp_misc.py | 7 +- e2e/test_rtsp_stun.py | 33 +++ e2e/test_rtsp_transport.py | 132 +++++++++++- src/connection.c | 35 +++- src/fcc.c | 6 +- src/multicast.c | 2 +- src/rtp.c | 2 +- src/rtp_reorder.c | 2 +- src/rtsp.c | 282 ++++++++++++++++++++++---- src/rtsp.h | 1 + src/snapshot.c | 2 +- src/stream.c | 342 +++++++++++++++++++++++++++----- src/stream.h | 74 ++++++- src/worker.c | 12 ++ 26 files changed, 1012 insertions(+), 124 deletions(-) diff --git a/docs/en/guide/fcc-setup.md b/docs/en/guide/fcc-setup.md index 21039aec..72f40ebd 100644 --- a/docs/en/guide/fcc-setup.md +++ b/docs/en/guide/fcc-setup.md @@ -81,6 +81,8 @@ rtp://239.253.64.120:5140/?fcc=10.255.14.152:15970 The difference between the two protocols is very small. In most cases, `telecom` will work. In some specific network environments, the `huawei` protocol may be required. +FCC media responses report the configured protocol in `R2H-FCC-Type` and the startup result in `R2H-FCC-Status`. The status is `active` when FCC unicast first produces valid output, or `fallback` when FCC does not succeed and multicast produces the initial output. This value is a response-time snapshot; a later transition from FCC unicast to multicast does not change the header already sent. A `HEAD` request returns FCC Type only, sends no FCC request, and omits FCC Status. See [URL Formats](/en/guide/url-formats#upstream-stream-metadata-response-headers) for the complete field reference. + ### Configure FCC in M3U If you are using an M3U playlist, you can add FCC parameters to the URL of each channel: diff --git a/docs/en/guide/url-formats.md b/docs/en/guide/url-formats.md index 495cd795..47e8829a 100644 --- a/docs/en/guide/url-formats.md +++ b/docs/en/guide/url-formats.md @@ -91,6 +91,31 @@ This is caused by timezone mismatch. You need to perform timezone conversion. Tr For detailed information on time-shift parameter handling (timezone, offset), see [Time Processing Guide](/en/guide/time-processing). +## Upstream Stream Metadata Response Headers + +RTSP, standard multicast, and FCC streams return confirmed upstream metadata in `R2H-*` response headers. Fields that are unknown, invalid, or not yet confirmed when the response is sent are omitted. + +| Response Header | Meaning | +| --- | --- | +| `R2H-Upstream-Protocol` | Upstream type: `rtsp` or `multicast`. FCC streams are also reported as `multicast`. | +| `R2H-Upstream-Transport` | Transport negotiated by RTSP: `tcp-interleaved` or `udp`. | +| `R2H-Upstream-Payload` | Confirmed MPEG-TS encapsulation: `mp2t-rtp` or `mp2t-direct`. | +| `R2H-Playback-Scale` | Playback scale explicitly confirmed by the RTSP `PLAY` response. | +| `R2H-Playback-Range` | Playback range explicitly confirmed by the RTSP `PLAY` response. | +| `R2H-Media-Duration` | Duration of a finite SDP `npt` range, in seconds. | +| `R2H-FCC-Type` | Configured FCC protocol: `telecom` or `huawei`. | +| `R2H-FCC-Status` | `active` when FCC unicast produced the initial output; `fallback` when multicast produced it. | + +These response headers are a snapshot taken when the HTTP headers are sent. Changes after stream startup do not update them. FEC status is not reported, and sensitive data such as RTSP sessions, authentication information, upstream addresses and ports, and raw SDP is never exposed. + +MPEG-TS responses and successfully generated JPEG snapshots include all fields confirmed at that point. `HEAD` requests behave differently: + +- RTSP `HEAD` connects upstream and performs only `OPTIONS` and `DESCRIBE`. It can therefore return Protocol and any Payload or Duration confirmed by SDP, but it never performs `SETUP` or `PLAY`. +- Multicast and FCC `HEAD` requests do not join multicast groups or send FCC requests. They return Protocol and, when configured, FCC Type only. +- HTTP reverse proxy `HEAD` requests continue to be handled by the upstream HTTP service. + +When `cors-allow-origin` is configured, all eight fields are listed in `Access-Control-Expose-Headers` so browser JavaScript can read them through the Fetch API. + ## HTTP Reverse Proxy ```url diff --git a/docs/en/guide/video-snapshot.md b/docs/en/guide/video-snapshot.md index ae5b91f7..26ff4728 100644 --- a/docs/en/guide/video-snapshot.md +++ b/docs/en/guide/video-snapshot.md @@ -181,6 +181,8 @@ http://192.168.1.1:5140/CCTV-1?snapshot=1 curl -H "Accept: image/jpeg" http://192.168.1.1:5140/rtp/239.253.64.120:5140 ``` +When JPEG generation succeeds, the response includes both `Content-Length` and all `R2H-*` upstream stream metadata confirmed before conversion. If snapshot conversion fails and falls back to a media stream, the media response sends the metadata once; no JPEG headers are sent in advance. See [URL Formats](/en/guide/url-formats#upstream-stream-metadata-response-headers) for the field reference. + ## Troubleshooting ### Snapshot Request Returns Video Stream Instead of Image diff --git a/docs/guide/fcc-setup.md b/docs/guide/fcc-setup.md index 190dc7c0..bae0d0dd 100644 --- a/docs/guide/fcc-setup.md +++ b/docs/guide/fcc-setup.md @@ -81,6 +81,8 @@ rtp://239.253.64.120:5140/?fcc=10.255.14.152:15970 两种协议区别很小,大多数情况使用 `telecom` 就可以工作。在一些特定网络环境下可能需要 `huawei` 协议。 +FCC 媒体响应会通过 `R2H-FCC-Type` 返回实际配置的协议,并通过 `R2H-FCC-Status` 返回起播结果:FCC 单播首先产生有效输出时为 `active`,FCC 未成功且首先从组播输出时为 `fallback`。这个状态是响应发出时的快照;后续从 FCC 单播切换到组播不会改变已经发送的 Header。`HEAD` 请求只返回 FCC Type,不发送 FCC 请求,也不返回 FCC Status。完整字段说明见 [URL 格式说明](./url-formats.md#上游流-metadata-响应头)。 + ### 在 M3U 中配置 FCC 如果使用 M3U 播放列表,可以在每个频道的 URL 中添加 FCC 参数: diff --git a/docs/guide/url-formats.md b/docs/guide/url-formats.md index f4294095..df38df13 100644 --- a/docs/guide/url-formats.md +++ b/docs/guide/url-formats.md @@ -91,6 +91,31 @@ http://192.168.1.1:5140/rtsp/iptv.example.com:554/channel1?r2h-ifname=eth0 关于时移回看的参数处理(时区、偏移),详见 [时间处理说明](./time-processing.md)。 +## 上游流 Metadata 响应头 + +RTSP、普通组播和 FCC 流会在媒体响应中返回已确认的 `R2H-*` Metadata。未知、无效或在响应发出时尚未确认的字段会被省略。 + +| 响应头 | 含义 | +| --- | --- | +| `R2H-Upstream-Protocol` | 上游类型:`rtsp` 或 `multicast`。FCC 流也属于 `multicast`。 | +| `R2H-Upstream-Transport` | RTSP 实际协商的传输方式:`tcp-interleaved` 或 `udp`。 | +| `R2H-Upstream-Payload` | 已确认的 MPEG-TS 封装:`mp2t-rtp` 或 `mp2t-direct`。 | +| `R2H-Playback-Scale` | RTSP `PLAY` 响应明确确认的播放速率。 | +| `R2H-Playback-Range` | RTSP `PLAY` 响应明确确认的播放范围。 | +| `R2H-Media-Duration` | SDP 有限 `npt` 范围的时长,单位为秒。 | +| `R2H-FCC-Type` | 已配置的 FCC 协议:`telecom` 或 `huawei`。 | +| `R2H-FCC-Status` | 起播时 FCC 单播已输出为 `active`;从组播起播为 `fallback`。 | + +这些响应头是发送 HTTP Header 时的状态快照,流启动后的状态变化不会更新它们。不会返回 FEC 工作状态,也不会暴露 RTSP Session、认证信息、上游地址、端口或原始 SDP 等敏感信息。 + +MPEG-TS 响应和成功生成的 JPEG 快照都会携带当时已经确认的字段。`HEAD` 请求的行为有所不同: + +- RTSP `HEAD` 只连接上游并执行 `OPTIONS` 和 `DESCRIBE`,因此只可能返回 Protocol,以及 SDP 能确认的 Payload 和 Duration;不会执行 `SETUP` 或 `PLAY`。 +- 组播和 FCC `HEAD` 不加入组播、不发送 FCC 请求,只返回 Protocol,以及已配置时的 FCC Type。 +- HTTP 反向代理的 `HEAD` 仍由上游 HTTP 服务处理。 + +配置 `cors-allow-origin` 后,以上八个字段也会列入 `Access-Control-Expose-Headers`,网页 JavaScript 可以通过 Fetch API 读取。 + ## HTTP 反向代理 ```url diff --git a/docs/guide/video-snapshot.md b/docs/guide/video-snapshot.md index e290983e..592925b5 100644 --- a/docs/guide/video-snapshot.md +++ b/docs/guide/video-snapshot.md @@ -181,6 +181,8 @@ http://192.168.1.1:5140/CCTV-1?snapshot=1 curl -H "Accept: image/jpeg" http://192.168.1.1:5140/rtp/239.253.64.120:5140 ``` +成功生成 JPEG 时,响应会同时包含 `Content-Length` 和生成前已经确认的 `R2H-*` 上游流 Metadata。若快照转换失败并回退为媒体流,则由媒体流响应发送一次 Metadata,不会提前发送 JPEG Header。字段说明见 [URL 格式说明](./url-formats.md#上游流-metadata-响应头)。 + ## 故障排查 ### 快照请求返回视频流而不是图片 diff --git a/e2e/helpers/mock_rtsp.py b/e2e/helpers/mock_rtsp.py index d9ce8702..562237f3 100644 --- a/e2e/helpers/mock_rtsp.py +++ b/e2e/helpers/mock_rtsp.py @@ -30,6 +30,8 @@ def __init__( content_base: str | None = "auto", custom_sdp: str | None = None, options_session_id: str | None = None, + play_response_headers: list[tuple[str, str]] | None = None, + close_after_describe: bool = False, host: str = "127.0.0.1", ): """ @@ -44,6 +46,8 @@ def __init__( options_session_id: If set, OPTIONS responds with this Session ID and every subsequent request (DESCRIBE, SETUP, PLAY, ...) must echo it (simulates HMS-style servers). + close_after_describe: Close the control connection immediately after + sending the DESCRIBE response. host: Address to listen on (use "::1" for IPv6 loopback). """ self.host = host @@ -52,6 +56,8 @@ def __init__( self._content_base = content_base self._custom_sdp = custom_sdp self._options_session_id = options_session_id + self._play_response_headers = play_response_headers or [] + self._close_after_describe = close_after_describe self._server_sock: socket.socket | None = None self._thread: threading.Thread | None = None self._stop = threading.Event() @@ -191,11 +197,17 @@ def _handle(self, conn: socket.socket, addr: tuple) -> None: "Content-Length: %d\r\n\r\n%s" % (cseq, cb_header, len(sdp), sdp) ).encode() ) + if self._close_after_describe: + return elif method == "SETUP": conn.sendall(self._setup_response(cseq, transport_hdr).encode()) elif method == "PLAY": + extra_headers = "".join("%s: %s\r\n" % item for item in self._play_response_headers) conn.sendall( - ("RTSP/1.0 200 OK\r\nCSeq: %s\r\nSession: %s\r\n\r\n" % (cseq, self._session_id())).encode() + ( + "RTSP/1.0 200 OK\r\nCSeq: %s\r\nSession: %s\r\n%s\r\n" + % (cseq, self._session_id(), extra_headers) + ).encode() ) self._after_play(conn, addr) return @@ -233,6 +245,8 @@ def __init__( content_base: str | None = "auto", custom_sdp: str | None = None, options_session_id: str | None = None, + play_response_headers: list[tuple[str, str]] | None = None, + close_after_describe: bool = False, host: str = "127.0.0.1", ): super().__init__( @@ -241,6 +255,8 @@ def __init__( content_base=content_base, custom_sdp=custom_sdp, options_session_id=options_session_id, + play_response_headers=play_response_headers, + close_after_describe=close_after_describe, host=host, ) self._num_packets = num_packets diff --git a/e2e/helpers/rtp.py b/e2e/helpers/rtp.py index 0b1669ef..39c75742 100644 --- a/e2e/helpers/rtp.py +++ b/e2e/helpers/rtp.py @@ -56,6 +56,7 @@ def __init__( reorder_distance: int = 0, unique_payloads: bool = False, send_duplicates: bool = False, + encapsulate_rtp: bool = True, ): self.addr = addr self.port = port or find_free_udp_port() @@ -64,6 +65,7 @@ def __init__( self.reorder_distance = reorder_distance self.unique_payloads = unique_payloads self.send_duplicates = send_duplicates + self.encapsulate_rtp = encapsulate_rtp self._payload = _TS_NULL_PACKET * ts_per_rtp self._sock: socket.socket | None = None self._thread: threading.Thread | None = None @@ -92,7 +94,7 @@ def _loop(self) -> None: payload = _make_ts_with_marker(seq) * self.ts_per_rtp else: payload = self._payload - pkt = make_rtp_packet(seq, ts, payload=payload) + pkt = make_rtp_packet(seq, ts, payload=payload) if self.encapsulate_rtp else payload if self.reorder_distance > 1: buf.append(pkt) diff --git a/e2e/test_config.py b/e2e/test_config.py index 79dce61f..452c7274 100644 --- a/e2e/test_config.py +++ b/e2e/test_config.py @@ -195,7 +195,7 @@ def _assert_rtsp_user_agent(port: int, expected_user_agent: str): rtsp = MockRTSPServer(custom_sdp=_RTSP_DURATION_SDP) rtsp.start() try: - status, _, body = http_get( + status, headers, body = http_get( "127.0.0.1", port, "/rtsp/127.0.0.1:%d/stream?r2h-duration=1" % rtsp.port, @@ -203,6 +203,7 @@ def _assert_rtsp_user_agent(port: int, expected_user_agent: str): ) assert status == 200, f"Expected 200, got {status}" assert b'"duration"' in body, "Expected RTSP duration response" + assert not any(name.lower().startswith("r2h-") for name in headers) option_reqs = [req for req in rtsp.requests_detailed if req["method"] == "OPTIONS"] assert option_reqs, "Expected OPTIONS request" diff --git a/e2e/test_fcc.py b/e2e/test_fcc.py index dfd65668..a2a4988d 100644 --- a/e2e/test_fcc.py +++ b/e2e/test_fcc.py @@ -19,6 +19,8 @@ find_free_port, find_free_udp_port, find_free_udp_port_pair, + get_header, + http_request, stream_get, ) @@ -60,7 +62,7 @@ def test_fcc_unicast_stream(self, shared_r2h): ) fcc.start() try: - status, _, body = stream_get( + status, headers, body = stream_get( "127.0.0.1", shared_r2h.port, f"/rtp/{MCAST_ADDR}:{mcast_port}?fcc=127.0.0.1:{fcc.port}", @@ -71,6 +73,10 @@ def test_fcc_unicast_stream(self, shared_r2h): assert len(body) > 0, "Expected to receive unicast stream data" assert body[0] == 0x47, f"Expected TS sync byte 0x47, got 0x{body[0]:02x}" assert fcc.requests_received >= 1 + assert headers["r2h-upstream-protocol"] == "multicast" + assert headers["r2h-upstream-payload"] == "mp2t-rtp" + assert headers["r2h-fcc-type"] == "telecom" + assert headers["r2h-fcc-status"] == "active" finally: fcc.stop() @@ -87,7 +93,7 @@ def test_fcc_single_socket_keeps_media_when_signal_port_changes(self, shared_r2h ) fcc.start() try: - status, _, body = stream_get( + status, headers, body = stream_get( "127.0.0.1", shared_r2h.port, f"/rtp/{MCAST_ADDR}:{mcast_port}?fcc=127.0.0.1:{fcc.port}", @@ -98,6 +104,8 @@ def test_fcc_single_socket_keeps_media_when_signal_port_changes(self, shared_r2h assert len(body) > 0, "Expected to receive unicast stream data" assert body[0] == 0x47, f"Expected TS sync byte 0x47, got 0x{body[0]:02x}" assert fcc.requests_received >= 1 + assert headers["r2h-fcc-type"] == "telecom" + assert headers["r2h-fcc-status"] == "active" finally: fcc.stop() @@ -118,7 +126,7 @@ def test_fcc_redirect_to_unicast_stream(self, shared_r2h): target.start() redirect.start() try: - status, _, body = stream_get( + status, headers, body = stream_get( "127.0.0.1", shared_r2h.port, f"/rtp/{MCAST_ADDR}:{mcast_port}?fcc=127.0.0.1:{redirect.port}", @@ -214,7 +222,7 @@ def test_huawei_fcc_unicast_stream(self, shared_r2h): fcc.start() try: url = f"/rtp/{MCAST_ADDR}:{mcast_port}?fcc=127.0.0.1:{fcc.port}&fcc-type=huawei" - status, _, body = stream_get( + status, headers, body = stream_get( "127.0.0.1", shared_r2h.port, url, @@ -225,6 +233,8 @@ def test_huawei_fcc_unicast_stream(self, shared_r2h): assert len(body) > 0, "Expected to receive unicast stream data" assert body[0] == 0x47, f"Expected TS sync byte 0x47, got 0x{body[0]:02x}" assert fcc.requests_received >= 1 + assert headers["r2h-fcc-type"] == "huawei" + assert headers["r2h-fcc-status"] == "active" finally: fcc.stop() @@ -417,7 +427,7 @@ def test_fcc_timeout_falls_back_to_multicast(self, shared_r2h): # Use a port where no FCC server is listening dead_fcc_port = find_free_udp_port() try: - status, _, body = stream_get( + status, headers, body = stream_get( "127.0.0.1", shared_r2h.port, f"/rtp/{MCAST_ADDR}:{mcast_port}?fcc=127.0.0.1:{dead_fcc_port}", @@ -426,5 +436,32 @@ def test_fcc_timeout_falls_back_to_multicast(self, shared_r2h): ) assert status == 200 assert len(body) > 0, "Expected multicast fallback to deliver data" + assert headers["r2h-upstream-protocol"] == "multicast" + assert headers["r2h-upstream-payload"] == "mp2t-rtp" + assert headers["r2h-fcc-type"] == "telecom" + assert headers["r2h-fcc-status"] == "fallback" finally: sender.stop() + + def test_fcc_head_returns_static_metadata_without_request(self, shared_r2h): + mcast_port = find_free_udp_port() + fcc = MockFCCServer(mcast_addr=MCAST_ADDR, protocol="huawei", sync_after=0) + fcc.start() + try: + status, headers, body = http_request( + "127.0.0.1", + shared_r2h.port, + "HEAD", + f"/rtp/{MCAST_ADDR}:{mcast_port}?fcc=127.0.0.1:{fcc.port}&fcc-type=huawei", + timeout=3.0, + ) + assert status == 200 + assert body == b"" + assert get_header(headers, "R2H-Upstream-Protocol") == "multicast" + assert get_header(headers, "R2H-FCC-Type") == "huawei" + assert get_header(headers, "R2H-Upstream-Payload") == "" + assert get_header(headers, "R2H-FCC-Status") == "" + time.sleep(0.1) + assert fcc.requests_received == 0 + finally: + fcc.stop() diff --git a/e2e/test_http_proxy.py b/e2e/test_http_proxy.py index 9f9d2e9b..bc38f842 100644 --- a/e2e/test_http_proxy.py +++ b/e2e/test_http_proxy.py @@ -62,6 +62,7 @@ def test_proxy_200(self, shared_r2h): ) assert status == 200 assert body == b"world" + assert not any(name.lower().startswith("r2h-") for name in hdrs) finally: upstream.stop() diff --git a/e2e/test_multicast.py b/e2e/test_multicast.py index dfe286fe..c60ae5f3 100644 --- a/e2e/test_multicast.py +++ b/e2e/test_multicast.py @@ -72,6 +72,29 @@ def test_rtp_stream_returns_200(self, multicast_r2h): ) assert status == 200 assert len(body) > 0, "Expected to receive stream data" + assert hdrs["r2h-upstream-protocol"] == "multicast" + assert hdrs["r2h-upstream-payload"] == "mp2t-rtp" + assert "r2h-fec-status" not in hdrs + finally: + sender.stop() + + def test_direct_ts_metadata(self, multicast_r2h): + """Raw MPEG-TS multicast should be identified without RTP.""" + mcast_port = find_free_udp_port() + sender = MulticastSender(addr=MCAST_ADDR, port=mcast_port, pps=200, encapsulate_rtp=False) + sender.start() + try: + status, headers, body = stream_get( + "127.0.0.1", + multicast_r2h.port, + f"/rtp/{MCAST_ADDR}:{mcast_port}", + read_bytes=4096, + timeout=_MCAST_STREAM_TIMEOUT, + ) + assert status == 200 + assert body + assert headers["r2h-upstream-protocol"] == "multicast" + assert headers["r2h-upstream-payload"] == "mp2t-direct" finally: sender.stop() @@ -217,7 +240,7 @@ def test_head_rtp(self, r2h_binary): r2h = R2HProcess( r2h_binary, port, - extra_args=["-v", "4", "-m", "100"], + extra_args=["-v", "4", "-m", "100", "-O", "*"], ) try: r2h.start() @@ -232,10 +255,58 @@ def test_head_rtp(self, r2h_binary): ) assert status == 200 assert len(body) == 0 # HEAD should have empty body + assert hdrs.get("R2H-Upstream-Protocol") == "multicast" + assert "R2H-Upstream-Payload" not in hdrs + assert "R2H-FCC-Status" not in hdrs + exposed = {name.strip() for name in hdrs["Access-Control-Expose-Headers"].split(",")} + assert exposed == { + "R2H-Upstream-Protocol", + "R2H-Upstream-Transport", + "R2H-Upstream-Payload", + "R2H-Playback-Scale", + "R2H-Playback-Range", + "R2H-Media-Duration", + "R2H-FCC-Type", + "R2H-FCC-Status", + } finally: r2h.stop() +class TestSnapshotFallbackMetadata: + """A failed JPEG capture must fall back to one MPEG-TS response.""" + + def test_snapshot_fallback_uses_stream_metadata(self, r2h_binary): + port = find_free_port() + mcast_port = find_free_udp_port() + r2h = R2HProcess( + r2h_binary, + port, + extra_args=["-v", "4", "-m", "100", "-r", LOOPBACK_IF, "-S"], + ) + sender = MulticastSender(addr=MCAST_ADDR, port=mcast_port, pps=300) + try: + r2h.start() + sender.start() + status, headers, body = stream_get( + "127.0.0.1", + port, + f"/rtp/{MCAST_ADDR}:{mcast_port}", + read_bytes=4096, + timeout=_MCAST_STREAM_TIMEOUT, + headers={"Accept": "image/jpeg"}, + ) + assert status == 200 + assert body + assert headers["content-type"] == "video/mp2t" + assert headers["r2h-upstream-protocol"] == "multicast" + assert headers["r2h-upstream-payload"] == "mp2t-rtp" + assert "content-length" not in headers + finally: + sender.stop() + r2h.stop() + + # --------------------------------------------------------------------------- # RTP reorder # --------------------------------------------------------------------------- diff --git a/e2e/test_rtsp_misc.py b/e2e/test_rtsp_misc.py index 53a04c42..a17dd3d1 100644 --- a/e2e/test_rtsp_misc.py +++ b/e2e/test_rtsp_misc.py @@ -97,15 +97,16 @@ def shared_r2h(r2h_binary): class TestRTSPHeadRequest: - def test_head_rtsp(self, shared_r2h): + def test_head_rtsp_unreachable_returns_503(self, shared_r2h): + dead_port = find_free_port() status, _, body = http_request( "127.0.0.1", shared_r2h.port, "HEAD", - "/rtsp/127.0.0.1:554/test", + "/rtsp/127.0.0.1:%d/test" % dead_port, timeout=3.0, ) - assert status == 200 + assert status == 503 assert len(body) == 0 diff --git a/e2e/test_rtsp_stun.py b/e2e/test_rtsp_stun.py index a049341c..c4823449 100644 --- a/e2e/test_rtsp_stun.py +++ b/e2e/test_rtsp_stun.py @@ -20,6 +20,7 @@ MockRTSPServer, R2HProcess, find_free_port, + http_request, stream_get, ) from helpers.mock_stun import MockSTUNServer @@ -29,6 +30,38 @@ _STREAM_TIMEOUT = 20.0 +class TestHeadProbeSkipsSTUN: + def test_head_does_not_allocate_media_transport(self, r2h_binary): + stun = MockSTUNServer(mapped_port=50000) + rtsp = MockRTSPServer() + stun.start() + rtsp.start() + r2h_port = find_free_port() + r2h = R2HProcess( + r2h_binary, + r2h_port, + extra_args=["-v", "4", "-m", "100", "-N", "127.0.0.1:%d" % stun.port], + ) + r2h.start() + try: + status, _, body = http_request( + "127.0.0.1", + r2h_port, + "HEAD", + "/rtsp/127.0.0.1:%d/stream" % rtsp.port, + timeout=8.0, + ) + assert status == 200 + assert body == b"" + assert rtsp.requests_received == ["OPTIONS", "DESCRIBE"] + time.sleep(0.1) + assert stun.requests_received == 0 + finally: + r2h.stop() + rtsp.stop() + stun.stop() + + class TestSTUNMappedPort: """When STUN succeeds, the mapped port should appear in SETUP.""" diff --git a/e2e/test_rtsp_transport.py b/e2e/test_rtsp_transport.py index cadf76dd..76aafeba 100644 --- a/e2e/test_rtsp_transport.py +++ b/e2e/test_rtsp_transport.py @@ -9,9 +9,12 @@ from helpers import ( MockRTSPServer, + MockRTSPServerSilent, MockRTSPServerUDP, R2HProcess, find_free_port, + get_header, + http_request, stream_get, ) @@ -54,6 +57,130 @@ def test_tcp_stream_returns_200(self, shared_r2h): finally: rtsp.stop() + def test_tcp_stream_metadata_headers(self, shared_r2h): + sdp = ( + "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=T\r\n" + "t=0 0\r\nm=video 0 RTP/AVP 33\r\n" + "a=range:npt=10.25-40.5\r\na=control:*\r\n" + ) + rtsp = MockRTSPServer( + num_packets=500, + custom_sdp=sdp, + play_response_headers=[("sCaLe", "1.500000"), ("rAnGe", "npt=10.25-40.5")], + ) + rtsp.start() + try: + status, headers, body = stream_get( + "127.0.0.1", + shared_r2h.port, + "/rtsp/127.0.0.1:%d/stream" % rtsp.port, + read_bytes=4096, + timeout=_STREAM_TIMEOUT, + ) + assert status == 200 + assert body + assert headers["r2h-upstream-protocol"] == "rtsp" + assert headers["r2h-upstream-transport"] == "tcp-interleaved" + assert headers["r2h-upstream-payload"] == "mp2t-rtp" + assert headers["r2h-playback-scale"] == "1.5" + assert headers["r2h-playback-range"] == "npt=10.25-40.5" + assert headers["r2h-media-duration"] == "30.25" + assert "r2h-metadata-version" not in headers + assert not any(name.startswith("r2h-fec-") for name in headers) + assert "session" not in headers + assert "content-base" not in headers + finally: + rtsp.stop() + + def test_head_stops_after_describe(self, shared_r2h): + sdp = ( + "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=T\r\n" + "t=0 0\r\nm=video 0 RTP/AVP 33\r\n" + "a=range:npt=2-12.125\r\na=control:*\r\n" + ) + rtsp = MockRTSPServer(custom_sdp=sdp) + rtsp.start() + try: + status, headers, body = http_request( + "127.0.0.1", + shared_r2h.port, + "HEAD", + "/rtsp/127.0.0.1:%d/stream" % rtsp.port, + timeout=10.0, + ) + assert status == 200 + assert body == b"" + assert get_header(headers, "R2H-Upstream-Protocol") == "rtsp" + assert get_header(headers, "R2H-Upstream-Payload") == "mp2t-rtp" + assert get_header(headers, "R2H-Media-Duration") == "10.125" + assert get_header(headers, "R2H-Upstream-Transport") == "" + assert get_header(headers, "R2H-Playback-Scale") == "" + assert get_header(headers, "R2H-Playback-Range") == "" + assert rtsp.requests_received == ["OPTIONS", "DESCRIBE"] + finally: + rtsp.stop() + + def test_head_metadata_takes_precedence_over_duration_query(self, shared_r2h): + sdp = ( + "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=T\r\n" + "t=0 0\r\nm=video 0 RTP/AVP 33\r\n" + "a=range:npt=2-12.125\r\na=control:*\r\n" + ) + rtsp = MockRTSPServer(custom_sdp=sdp) + rtsp.start() + try: + status, headers, body = http_request( + "127.0.0.1", + shared_r2h.port, + "HEAD", + "/rtsp/127.0.0.1:%d/stream?r2h-duration=1" % rtsp.port, + timeout=10.0, + ) + assert status == 200 + assert body == b"" + assert get_header(headers, "Content-Type") == "video/mp2t" + assert get_header(headers, "R2H-Upstream-Protocol") == "rtsp" + assert get_header(headers, "R2H-Upstream-Payload") == "mp2t-rtp" + assert get_header(headers, "R2H-Media-Duration") == "10.125" + assert rtsp.requests_received == ["OPTIONS", "DESCRIBE"] + finally: + rtsp.stop() + + def test_head_survives_upstream_close_after_describe(self, shared_r2h): + rtsp = MockRTSPServer(close_after_describe=True) + rtsp.start() + try: + status, headers, body = http_request( + "127.0.0.1", + shared_r2h.port, + "HEAD", + "/rtsp/127.0.0.1:%d/stream" % rtsp.port, + timeout=10.0, + ) + assert status == 200 + assert body == b"" + assert get_header(headers, "R2H-Upstream-Protocol") == "rtsp" + assert get_header(headers, "R2H-Upstream-Payload") == "mp2t-rtp" + assert rtsp.requests_received == ["OPTIONS", "DESCRIBE"] + finally: + rtsp.stop() + + def test_head_probe_timeout_returns_503(self, shared_r2h): + rtsp = MockRTSPServerSilent() + rtsp.start() + try: + status, _, body = http_request( + "127.0.0.1", + shared_r2h.port, + "HEAD", + "/rtsp/127.0.0.1:%d/stream" % rtsp.port, + timeout=8.0, + ) + assert status == 503 + assert body == b"" + finally: + rtsp.stop() + def test_tcp_data_is_ts(self, shared_r2h): """Relayed data should be raw MPEG-TS (RTP headers stripped).""" rtsp = MockRTSPServer(num_packets=500) @@ -129,7 +256,7 @@ def test_udp_stream_returns_200(self, shared_r2h): rtsp = MockRTSPServerUDP() rtsp.start() try: - status, _, body = stream_get( + status, headers, body = stream_get( "127.0.0.1", shared_r2h.port, "/rtsp/127.0.0.1:%d/stream" % rtsp.port, @@ -138,6 +265,9 @@ def test_udp_stream_returns_200(self, shared_r2h): ) assert status == 200, "Expected 200 for UDP RTSP" assert len(body) > 0 + assert headers["r2h-upstream-protocol"] == "rtsp" + assert headers["r2h-upstream-transport"] == "udp" + assert headers["r2h-upstream-payload"] == "mp2t-rtp" finally: rtsp.stop() diff --git a/src/connection.c b/src/connection.c index 0f5b29ce..7f1e3348 100644 --- a/src/connection.c +++ b/src/connection.c @@ -1073,22 +1073,37 @@ int connection_route_and_start(connection_t *c) { return 0; } - /* Handle HEAD requests for RTP/RTSP streams - return success without - * connecting upstream. HTTP services forward HEAD to the upstream server - * so the real Content-Type (e.g. application/vnd.apple.mpegurl for HLS) - * is returned to the client. */ + if (c->http_req.user_agent[0]) { + service->user_agent = strdup(c->http_req.user_agent); + } + + /* HTTP services forward HEAD upstream unchanged. Multicast HEAD requests + * return only static metadata. RTSP HEAD performs an asynchronous + * OPTIONS/DESCRIBE probe without opening media resources. */ if (strcasecmp(c->http_req.method, "HEAD") == 0 && service->service_type != SERVICE_HTTP) { - logger(LOG_INFO, "HEAD request detected, returning success without upstream connection"); - send_http_headers(c, STATUS_200, "video/mp2t", NULL); + if (service->service_type == SERVICE_RTSP) { + logger(LOG_INFO, "RTSP HEAD request detected, starting metadata probe"); + if (stream_context_init_rtsp_metadata_probe(&c->stream, c, service, c->epfd) == 0) { + c->streaming = 1; + c->service = service; + c->state = CONN_STREAMING; + return 0; + } + + stream_context_cleanup(&c->stream); + http_send_503(c); + service_free(service); + return 0; + } + + logger(LOG_INFO, "Multicast HEAD request detected, returning static metadata"); + stream_metadata_init(&c->stream.metadata, service); + stream_send_http_headers(c, "video/mp2t", NULL); connection_queue_output_and_flush(c, NULL, 0); service_free(service); return 0; } - if (c->http_req.user_agent[0]) { - service->user_agent = strdup(c->http_req.user_agent); - } - /* Check if this is a snapshot request (X-Request-Snapshot, Accept: * image/jpeg, or snapshot=1) */ /* 1 = snapshot=1, 2 = X-Request-Snapshot or Accept: image/jpeg */ diff --git a/src/fcc.c b/src/fcc.c index 1bb4366e..63d1db87 100644 --- a/src/fcc.c +++ b/src/fcc.c @@ -662,7 +662,7 @@ int fcc_handle_unicast_media(stream_context_t *ctx, buffer_ref_t *buf_ref) { } /* Forward RTP payload to client (with reordering) */ - stream_process_rtp_payload(ctx, buf_ref); + stream_process_rtp_payload(ctx, buf_ref, STREAM_MEDIA_ORIGIN_FCC_UNICAST); /* Check if we should terminate FCC based on reorder's delivered sequence. * base_seq - 1 is the last sequence number successfully delivered. @@ -765,7 +765,7 @@ int fcc_handle_mcast_active(stream_context_t *ctx, buffer_ref_t *buf_ref) { while (node) { /* Queue each buffer for zero-copy send */ buffer_ref_t *next = node->send_next; - int processed_bytes = stream_process_rtp_payload(ctx, node); + int processed_bytes = stream_process_rtp_payload(ctx, node, STREAM_MEDIA_ORIGIN_FCC_MULTICAST); if (likely(processed_bytes > 0)) { flushed_bytes += (uint64_t)processed_bytes; } @@ -782,7 +782,7 @@ int fcc_handle_mcast_active(stream_context_t *ctx, buffer_ref_t *buf_ref) { /* Forward multicast data to client (true zero-copy) or capture I-frame * (snapshot) */ - stream_process_rtp_payload(ctx, buf_ref); + stream_process_rtp_payload(ctx, buf_ref, STREAM_MEDIA_ORIGIN_FCC_MULTICAST); return 0; } diff --git a/src/multicast.c b/src/multicast.c index 64896fd5..0f4488f6 100644 --- a/src/multicast.c +++ b/src/multicast.c @@ -482,7 +482,7 @@ int mcast_session_handle_event(mcast_session_t *session, stream_context_t *ctx, /* Handle based on FCC state (if FCC initialized) */ if (!ctx->fcc.initialized) { /* Direct multicast without FCC - forward to client */ - stream_process_rtp_payload(ctx, recv_buf); + stream_process_rtp_payload(ctx, recv_buf, STREAM_MEDIA_ORIGIN_MULTICAST); buffer_ref_put(recv_buf); continue; /* Read next packet */ } diff --git a/src/rtp.c b/src/rtp.c index 3cf11129..31c287f5 100644 --- a/src/rtp.c +++ b/src/rtp.c @@ -87,7 +87,7 @@ int rtp_get_payload(uint8_t *buf, int recv_len, uint8_t **payload, int *size, ui int rtp_queue_buf_direct(connection_t *conn, buffer_ref_t *buf_ref) { /* Send headers lazily on first data packet */ if (!conn->headers_sent) { - send_http_headers(conn, STATUS_200, "video/mp2t", NULL); + stream_send_http_headers(conn, "video/mp2t", NULL); } /* Queue for zero-copy send */ diff --git a/src/rtp_reorder.c b/src/rtp_reorder.c index fb0cd495..81ec69d4 100644 --- a/src/rtp_reorder.c +++ b/src/rtp_reorder.c @@ -83,7 +83,7 @@ static int deliver_raw_packet(uint8_t *data, int len, connection_t *conn, int is } /* Send headers lazily on first data packet (same as rtp_queue_buf_direct) */ if (!conn->headers_sent) { - send_http_headers(conn, STATUS_200, "video/mp2t", NULL); + stream_send_http_headers(conn, "video/mp2t", NULL); } return connection_queue_output(conn, data, len) == 0 ? len : -1; } diff --git a/src/rtsp.c b/src/rtsp.c index 4e9a111d..d9ac2f38 100644 --- a/src/rtsp.c +++ b/src/rtsp.c @@ -12,13 +12,16 @@ #include "stream.h" #include "utils.h" #include "worker.h" +#include #include +#include #include #include #include #include #include #include +#include #include #include @@ -36,6 +39,7 @@ static const char rtsp_default_user_agent[] = "rtp2httpd/" VERSION; #define RTSP_RESPONSE_ADVANCE 1 #define RTSP_RESPONSE_KEEPALIVE 2 #define RTSP_RESPONSE_DURATION 3 +#define RTSP_RESPONSE_METADATA_READY 4 #define RTSP_RESPONSE_ERROR -1 /* Helper function prototypes */ static int rtsp_prepare_request(rtsp_session_t *session, const char *method, const char *url, @@ -52,6 +56,7 @@ static void rtsp_send_udp_nat_probe(rtsp_session_t *session); static int rtsp_process_interleaved_buffer(rtsp_session_t *session, connection_t *conn); static int rtsp_handle_redirect(rtsp_session_t *session, const char *location); static void rtsp_parse_describe_sdp(rtsp_session_t *session, const char *header_start, const char *sdp_body); +static void rtsp_parse_play_metadata(rtsp_session_t *session, const char *response); static int rtsp_initiate_teardown(rtsp_session_t *session); static int rtsp_reconnect_for_teardown(rtsp_session_t *session); static void rtsp_force_cleanup(rtsp_session_t *session); @@ -280,6 +285,7 @@ void rtsp_session_init(rtsp_session_t *session) { session->use_playseek_range = 0; session->r2h_duration = 0; session->r2h_duration_value = -1; + session->metadata_probe = 0; /* Initialize transport parameters - mode will be negotiated during SETUP */ session->transport_mode = RTSP_TRANSPORT_TCP; /* Default preference */ @@ -849,8 +855,8 @@ int rtsp_connect(rtsp_session_t *session) { * This allows STUN to run in parallel with TCP connection establishment * Only do this on initial connect (not on redirect or reconnect for TEARDOWN) * Check: UDP socket not yet created and STUN not already in progress/completed */ - if (config.rtsp_stun_server && config.rtsp_stun_server[0] != '\0' && session->rtp_socket < 0 && - !session->stun.in_progress && !session->stun.completed) { + if (!session->metadata_probe && config.rtsp_stun_server && config.rtsp_stun_server[0] != '\0' && + session->rtp_socket < 0 && !session->stun.in_progress && !session->stun.completed) { if (rtsp_setup_udp_sockets(session) == 0) { if (stun_send_request(&session->stun, session->rtp_socket) == 0) { logger(LOG_DEBUG, "RTSP: Started STUN discovery before TCP connect"); @@ -927,7 +933,8 @@ int rtsp_handle_socket_event(rtsp_session_t *session, uint32_t events) { } /* Check for connection errors or hangup */ - else if (events & (POLLER_HUP | POLLER_ERR | POLLER_RDHUP)) { + else if ((events & (POLLER_HUP | POLLER_ERR | POLLER_RDHUP)) && + !(session->metadata_probe && session->awaiting_response && (events & POLLER_IN) && !(events & POLLER_ERR))) { if (events & POLLER_ERR) { int sock_error = 0; socklen_t error_len = sizeof(sock_error); @@ -996,7 +1003,16 @@ int rtsp_handle_socket_event(rtsp_session_t *session, uint32_t events) { /* Handle readable socket - try to receive response */ if (events & POLLER_IN) { if (session->awaiting_response) { - int response_result = rtsp_try_receive_response(session); + int response_result; + size_t previous_buffer_pos; + do { + previous_buffer_pos = session->response_buffer_pos; + response_result = rtsp_try_receive_response(session); + /* A closing probe socket can deliver its final response and EOF in + * one edge-triggered event. Keep draining while each recv makes + * progress so a fragmented DESCRIBE body is parsed before HUP. */ + } while (response_result == RTSP_RESPONSE_OK && session->metadata_probe && session->awaiting_response && + (events & (POLLER_HUP | POLLER_RDHUP)) && session->response_buffer_pos > previous_buffer_pos); if (response_result < 0) { logger(LOG_ERROR, "RTSP: Failed to receive response"); rtsp_session_set_state(session, RTSP_STATE_ERROR); @@ -1006,6 +1022,9 @@ int rtsp_handle_socket_event(rtsp_session_t *session, uint32_t events) { if (response_result == RTSP_RESPONSE_DURATION) { return -2; } + if (response_result == RTSP_RESPONSE_METADATA_READY) { + return -3; + } /* Re-enable POLLER_OUT for next request */ if (response_result == RTSP_RESPONSE_ADVANCE && session->epoll_fd >= 0) { @@ -1282,8 +1301,11 @@ static int rtsp_try_send_pending(rtsp_session_t *session) { * RTSP_RESPONSE_OK: Waiting for more data, or response processed internally * RTSP_RESPONSE_ERROR: recv failure, connection closed, or parse error * RTSP_RESPONSE_DURATION: Duration query completed + * RTSP_RESPONSE_METADATA_READY: HEAD metadata probe completed */ static int rtsp_try_receive_response(rtsp_session_t *session) { + int peer_closed = 0; + if (!session->awaiting_response) { return RTSP_RESPONSE_OK; } @@ -1304,12 +1326,17 @@ static int rtsp_try_receive_response(rtsp_session_t *session) { } if (received == 0) { - logger(LOG_ERROR, "RTSP: Connection closed by server"); - session->awaiting_keepalive_response = 0; - return RTSP_RESPONSE_ERROR; + if (session->response_buffer_pos == 0) { + logger(LOG_ERROR, "RTSP: Connection closed by server"); + session->awaiting_keepalive_response = 0; + return RTSP_RESPONSE_ERROR; + } + /* The peer may close immediately after its final response. Parse any + * bytes already buffered before treating EOF as an error. */ + peer_closed = 1; + } else { + session->response_buffer_pos += (size_t)received; } - - session->response_buffer_pos += (size_t)received; } /* NUL-terminate for strstr-based parsing. Clamp to buffer bounds since @@ -1329,6 +1356,11 @@ static int rtsp_try_receive_response(rtsp_session_t *session) { rtsp_parse_response_header(session, (const char *)session->response_buffer, &response_offset, &response_len); if (parse_result == 1) { + if (peer_closed) { + logger(LOG_ERROR, "RTSP: Connection closed with an incomplete response"); + session->awaiting_keepalive_response = 0; + return RTSP_RESPONSE_ERROR; + } /* Need more data - for TCP interleaved mode, drain '$'-prefixed frames * to free buffer space and forward media data while waiting for RTSP * response */ @@ -1454,11 +1486,17 @@ static int rtsp_try_receive_response(rtsp_session_t *session) { if (session->state == RTSP_STATE_AWAITING_DESCRIBE) { rtsp_parse_describe_sdp(session, (const char *)session->response_buffer + response_offset, (const char *)session->response_buffer + response_offset + response_len); + rtsp_session_set_state(session, RTSP_STATE_DESCRIBED); + /* HEAD metadata probes always stop after DESCRIBE, even when the URL also + * contains the legacy r2h-duration query parameter. */ + if (session->metadata_probe) { + session->response_buffer_pos = 0; + return RTSP_RESPONSE_METADATA_READY; + } if (session->r2h_duration) { session->response_buffer_pos = 0; return RTSP_RESPONSE_DURATION; } - rtsp_session_set_state(session, RTSP_STATE_DESCRIBED); session->response_buffer_pos = 0; return RTSP_RESPONSE_ADVANCE; } @@ -1829,7 +1867,7 @@ static int rtsp_process_interleaved_buffer(rtsp_session_t *session, connection_t session->first_media_received = 1; logger(LOG_DEBUG, "RTSP: First media packet received (TCP)"); } - int pb = stream_process_rtp_payload(&conn->stream, packet_buf); + int pb = stream_process_rtp_payload(&conn->stream, packet_buf, STREAM_MEDIA_ORIGIN_RTSP); if (pb > 0) bytes_forwarded += pb; buffer_ref_put(packet_buf); @@ -1973,7 +2011,7 @@ int rtsp_handle_udp_rtp_data(rtsp_session_t *session, connection_t *conn) { session->first_media_received = 1; logger(LOG_DEBUG, "RTSP: First media packet received (UDP)"); } - int pb = stream_process_rtp_payload(&conn->stream, rtp_buf); + int pb = stream_process_rtp_payload(&conn->stream, rtp_buf, STREAM_MEDIA_ORIGIN_RTSP); buffer_ref_put(rtp_buf); if (pb > 0) total_bytes_written += pb; @@ -2295,6 +2333,19 @@ static int rtsp_parse_response_header(rtsp_session_t *session, const char *respo /* Increment retry counter */ session->auth_retry_count++; + if (session->conn && !session->conn->stream.metadata.frozen) { + stream_metadata_t *metadata = &session->conn->stream.metadata; + if (session->state == RTSP_STATE_AWAITING_DESCRIBE) { + metadata->upstream_payload = STREAM_PAYLOAD_UNKNOWN; + metadata->media_duration_known = 0; + } else if (session->state == RTSP_STATE_AWAITING_SETUP) { + metadata->upstream_transport = STREAM_TRANSPORT_UNKNOWN; + } else if (session->state == RTSP_STATE_AWAITING_PLAY) { + metadata->playback_scale_known = 0; + metadata->playback_range[0] = '\0'; + } + } + /* Move state back to retry the same request */ if (session->state == RTSP_STATE_AWAITING_OPTIONS) { rtsp_session_set_state(session, RTSP_STATE_CONNECTED); @@ -2330,6 +2381,10 @@ static int rtsp_parse_response_header(rtsp_session_t *session, const char *respo goto cleanup; } + if (session->state == RTSP_STATE_AWAITING_PLAY) { + rtsp_parse_play_metadata(session, rtsp_start); + } + /* Parse Public header from OPTIONS response to determine supported methods */ if (session->state == RTSP_STATE_AWAITING_OPTIONS) { public_header = rtsp_find_header(rtsp_start, "Public"); @@ -2591,33 +2646,99 @@ static void rtsp_close_udp_sockets(rtsp_session_t *session, const char *reason) } static char *rtsp_find_header(const char *response, const char *header_name) { - char *header_start, *header_end; - char header_prefix[RTSP_HEADER_PREFIX_SIZE]; - int header_len; + const char *line; + size_t header_name_len; - snprintf(header_prefix, sizeof(header_prefix), "%s:", header_name); - header_start = strstr(response, header_prefix); - if (!header_start) + if (!response || !header_name) return NULL; - header_start += strlen(header_prefix); - while (*header_start == ' ') - header_start++; /* Skip whitespace */ - - header_end = strstr(header_start, "\r\n"); - if (!header_end) + line = strstr(response, "\r\n"); + if (!line) return NULL; + line += 2; + header_name_len = strlen(header_name); + + while (*line != '\0') { + const char *line_end = strstr(line, "\r\n"); + const char *colon; + const char *value_start; + const char *value_end; + size_t value_len; + char *result; + + if (!line_end || line_end == line) + return NULL; + + colon = memchr(line, ':', (size_t)(line_end - line)); + if (!colon || (size_t)(colon - line) != header_name_len || strncasecmp(line, header_name, header_name_len) != 0) { + line = line_end + 2; + continue; + } - header_len = header_end - header_start; - char *result = malloc(header_len + 1); - if (!result) { - logger(LOG_ERROR, "RTSP: Failed to allocate memory for header"); - return NULL; + value_start = colon + 1; + while (value_start < line_end && (*value_start == ' ' || *value_start == '\t')) + value_start++; + value_end = line_end; + while (value_end > value_start && (value_end[-1] == ' ' || value_end[-1] == '\t')) + value_end--; + + value_len = (size_t)(value_end - value_start); + result = malloc(value_len + 1); + if (!result) { + logger(LOG_ERROR, "RTSP: Failed to allocate memory for header"); + return NULL; + } + memcpy(result, value_start, value_len); + result[value_len] = '\0'; + return result; } - strncpy(result, header_start, header_len); - result[header_len] = '\0'; - return result; + return NULL; +} + +static void rtsp_parse_play_metadata(rtsp_session_t *session, const char *response) { + stream_metadata_t *metadata; + char *scale_header; + char *range_header; + + if (!session || !session->conn || !response) + return; + + metadata = &session->conn->stream.metadata; + if (metadata->frozen) + return; + + scale_header = rtsp_find_header(response, "Scale"); + if (scale_header) { + char *end = NULL; + double scale; + + errno = 0; + scale = strtod(scale_header, &end); + while (end && (*end == ' ' || *end == '\t')) + end++; + if (end != scale_header && end && *end == '\0' && errno != ERANGE && isfinite(scale)) { + metadata->playback_scale = scale; + metadata->playback_scale_known = 1; + } + free(scale_header); + } + + range_header = rtsp_find_header(response, "Range"); + if (range_header) { + size_t len = strlen(range_header); + int valid = len > 0 && len < sizeof(metadata->playback_range); + + for (size_t i = 0; valid && i < len; i++) { + unsigned char ch = (unsigned char)range_header[i]; + if (ch < 0x20 || ch > 0x7e) + valid = 0; + } + if (valid) { + memcpy(metadata->playback_range, range_header, len + 1); + } + free(range_header); + } } /** @@ -2650,6 +2771,56 @@ static void rtsp_resolve_relative_url(const char *base_url, const char *relative } } +static int rtsp_sdp_describes_mp2t_rtp(const char *sdp_body) { + const char *line = sdp_body; + + while (line && *line != '\0') { + const char *line_end = strpbrk(line, "\r\n"); + size_t line_len = line_end ? (size_t)(line_end - line) : strlen(line); + + if (line_len > 9 && strncasecmp(line, "a=rtpmap:", 9) == 0) { + int payload_type; + char encoding[32]; + if (sscanf(line, "a=rtpmap:%d %31[^/\r\n]", &payload_type, encoding) == 2 && strcasecmp(encoding, "MP2T") == 0) { + return 1; + } + } else if (line_len > 2 && strncasecmp(line, "m=", 2) == 0) { + char media[32]; + char transport[64]; + char payloads[128]; + int port; + + if (sscanf(line, "m=%31s %d %63s %127[^\r\n]", media, &port, transport, payloads) == 4 && + strncasecmp(transport, "RTP/", 4) == 0) { + char *cursor = payloads; + while (*cursor != '\0') { + char *end; + long payload_type; + + while (*cursor == ' ' || *cursor == '\t') + cursor++; + if (*cursor == '\0') + break; + payload_type = strtol(cursor, &end, 10); + if (end == cursor) + break; + if (payload_type == 33) + return 1; + cursor = end; + } + } + } + + if (!line_end) + break; + line = line_end; + while (*line == '\r' || *line == '\n') + line++; + } + + return 0; +} + /** * Parse DESCRIBE response SDP body in a single pass. * @@ -2668,6 +2839,7 @@ static void rtsp_resolve_relative_url(const char *base_url, const char *relative * @param sdp_body Start of the SDP body (right after \\r\\n\\r\\n) */ static void rtsp_parse_describe_sdp(rtsp_session_t *session, const char *header_start, const char *sdp_body) { + stream_metadata_t *metadata = session->conn ? &session->conn->stream.metadata : NULL; session->setup_url[0] = '\0'; char *content_base = rtsp_find_header(header_start, "Content-Base"); @@ -2676,6 +2848,10 @@ static void rtsp_parse_describe_sdp(rtsp_session_t *session, const char *header_ if (*sdp_body == '\0') goto done; + if (metadata && !metadata->frozen && rtsp_sdp_describes_mp2t_rtp(sdp_body)) { + metadata->upstream_payload = STREAM_PAYLOAD_MP2T_RTP; + } + /* --- a=control: → setup_url ----------------------------------------- */ /* Find the first media-level a=control: (after first m= line). @@ -2718,26 +2894,44 @@ static void rtsp_parse_describe_sdp(rtsp_session_t *session, const char *header_ /* --- a=range:npt= → r2h_duration_value ------------------------------ */ - if (session->r2h_duration) { + { const char *range = strstr(sdp_body, "a=range:npt="); if (range) { const char *val = range + 12; /* skip "a=range:npt=" */ const char *val_end = strpbrk(val, "\r\n"); + size_t len; + char buf[64]; + if (!val_end) val_end = val + strlen(val); + len = (size_t)(val_end - val); - size_t len = val_end - val; - char buf[64]; if (len > 0 && len < sizeof(buf)) { + char *start_end; + char *end_end; + double range_start; + double range_end; + memcpy(buf, val, len); buf[len] = '\0'; - - float range_start, range_end; - if (sscanf(buf, "%f-%f", &range_start, &range_end) == 2) { - logger(LOG_DEBUG, "RTSP: Range: %.3f-%.3f", range_start, range_end); - session->r2h_duration_value = range_end; + errno = 0; + range_start = strtod(buf, &start_end); + if (start_end != buf && *start_end == '-') { + range_end = strtod(start_end + 1, &end_end); + if (end_end != start_end + 1 && *end_end == '\0' && errno != ERANGE && isfinite(range_start) && + isfinite(range_end) && range_end >= range_start) { + logger(LOG_DEBUG, "RTSP: Range: %.3f-%.3f", range_start, range_end); + if (session->r2h_duration) + session->r2h_duration_value = (float)range_end; + if (metadata && !metadata->frozen) { + metadata->media_duration = range_end - range_start; + metadata->media_duration_known = 1; + } + } else { + logger(LOG_DEBUG, "RTSP Range: %s, cannot find finite range", buf); + } } else { - logger(LOG_DEBUG, "RTSP Range: %s, cannot find right range!", buf); + logger(LOG_DEBUG, "RTSP Range: %s, cannot find finite range", buf); } } } @@ -2843,6 +3037,8 @@ static void rtsp_parse_transport_header(rtsp_session_t *session, const char *tra if (strstr(transport, "TCP") || strstr(transport, "interleaved=")) { /* TCP transport mode */ session->transport_mode = RTSP_TRANSPORT_TCP; + if (session->conn && !session->conn->stream.metadata.frozen) + session->conn->stream.metadata.upstream_transport = STREAM_TRANSPORT_TCP_INTERLEAVED; session->keepalive_interval_ms = RTSP_KEEPALIVE_INTERVAL_MS; session->last_keepalive_ms = 0; session->keepalive_pending = 0; @@ -2864,6 +3060,8 @@ static void rtsp_parse_transport_header(rtsp_session_t *session, const char *tra } else { /* UDP transport mode */ session->transport_mode = RTSP_TRANSPORT_UDP; + if (session->conn && !session->conn->stream.metadata.frozen) + session->conn->stream.metadata.upstream_transport = STREAM_TRANSPORT_UDP; session->keepalive_interval_ms = RTSP_KEEPALIVE_INTERVAL_MS; session->last_keepalive_ms = 0; session->keepalive_pending = 0; @@ -2940,6 +3138,8 @@ static int rtsp_handle_redirect(rtsp_session_t *session, const char *location) { * session learned from the previous server so requests to the * redirected-to server don't carry a stale Session header. */ session->session_id[0] = '\0'; + if (session->conn) + stream_metadata_reset_rtsp_negotiation(&session->conn->stream.metadata); /* Close current connection and remove from poller properly */ if (session->socket >= 0) { diff --git a/src/rtsp.h b/src/rtsp.h index 96f456bb..3fc1e244 100644 --- a/src/rtsp.h +++ b/src/rtsp.h @@ -142,6 +142,7 @@ typedef struct { int use_playseek_range; int r2h_duration; float r2h_duration_value; + int metadata_probe; /* Authentication state */ char username[RTSP_CREDENTIAL_SIZE]; /* RTSP username for authentication */ diff --git a/src/snapshot.c b/src/snapshot.c index 9a91d956..a3d53ce5 100644 --- a/src/snapshot.c +++ b/src/snapshot.c @@ -463,7 +463,7 @@ int snapshot_process_packet(snapshot_context_t *ctx, int recv_len, uint8_t *buf, char content_length_header[64]; snprintf(content_length_header, sizeof(content_length_header), "Content-Length: %zu\r\n", jpeg_size); - send_http_headers(conn, STATUS_200, "image/jpeg", content_length_header); + stream_send_http_headers(conn, "image/jpeg", content_length_header); /* Queue JPEG file for non-blocking sendfile() */ if (connection_queue_file(conn, jpeg_fd, 0, jpeg_size) < 0) { diff --git a/src/stream.c b/src/stream.c index 40964694..ee13ddf4 100644 --- a/src/stream.c +++ b/src/stream.c @@ -1,6 +1,8 @@ #include "stream.h" +#include "configuration.h" #include "connection.h" #include "fcc.h" +#include "http.h" #include "http_proxy.h" #include "multicast.h" #include "rtp.h" @@ -11,14 +13,230 @@ #include "status.h" #include "utils.h" #include +#include #include #include #include +#include #include #include #include #include +#define STREAM_METADATA_HEADERS_SIZE 1024 +#define TS_PACKET_SIZE 188 +#define TS_SYNC_BYTE 0x47 + +static int stream_metadata_append_raw(char *buffer, size_t buffer_size, size_t *length, const char *value) { + size_t value_len; + + if (!buffer || !length || !value || *length >= buffer_size) + return -1; + + value_len = strlen(value); + if (value_len >= buffer_size - *length) + return -1; + + memcpy(buffer + *length, value, value_len + 1); + *length += value_len; + return 0; +} + +static int stream_metadata_append_header(char *buffer, size_t buffer_size, size_t *length, const char *name, + const char *value) { + char line[512]; + int written; + + written = snprintf(line, sizeof(line), "%s: %s\r\n", name, value); + if (written < 0 || (size_t)written >= sizeof(line)) + return -1; + return stream_metadata_append_raw(buffer, buffer_size, length, line); +} + +static void stream_metadata_format_number(double value, char *buffer, size_t buffer_size) { + char *end; + + if (!buffer || buffer_size == 0) + return; + + snprintf(buffer, buffer_size, "%.6f", value); + end = buffer + strlen(buffer); + while (end > buffer && end[-1] == '0') + *--end = '\0'; + if (end > buffer && end[-1] == '.') + *--end = '\0'; + if (strcmp(buffer, "-0") == 0) + snprintf(buffer, buffer_size, "0"); +} + +static int stream_payload_is_mpegts(const uint8_t *payload, int payload_len) { + int checked = 0; + + if (!payload || payload_len < TS_PACKET_SIZE || payload[0] != TS_SYNC_BYTE) + return 0; + + for (int offset = 0; offset + TS_PACKET_SIZE <= payload_len && checked < 3; offset += TS_PACKET_SIZE) { + if (payload[offset] != TS_SYNC_BYTE) + return 0; + checked++; + } + + return checked > 0; +} + +void stream_metadata_init(stream_metadata_t *metadata, const service_t *service) { + if (!metadata) + return; + + memset(metadata, 0, sizeof(*metadata)); + if (!service) + return; + + if (service->service_type == SERVICE_RTSP) { + metadata->upstream_protocol = STREAM_UPSTREAM_RTSP; + } else if (service->service_type == SERVICE_MRTP) { + metadata->upstream_protocol = STREAM_UPSTREAM_MULTICAST; + if (service->fcc_addr) { + metadata->fcc_type = service->fcc_type == FCC_TYPE_HUAWEI ? STREAM_FCC_TYPE_HUAWEI : STREAM_FCC_TYPE_TELECOM; + } + } +} + +void stream_metadata_reset_rtsp_negotiation(stream_metadata_t *metadata) { + if (!metadata || metadata->frozen) + return; + + metadata->upstream_transport = STREAM_TRANSPORT_UNKNOWN; + metadata->upstream_payload = STREAM_PAYLOAD_UNKNOWN; + metadata->playback_scale_known = 0; + metadata->media_duration_known = 0; + metadata->playback_range[0] = '\0'; +} + +static void stream_metadata_note_media(stream_context_t *ctx, int packet_type, const uint8_t *payload, int payload_len, + stream_media_origin_t origin) { + stream_metadata_t *metadata; + + if (!ctx || ctx->metadata.frozen || !stream_payload_is_mpegts(payload, payload_len)) + return; + + metadata = &ctx->metadata; + metadata->upstream_payload = packet_type == 1 ? STREAM_PAYLOAD_MP2T_RTP : STREAM_PAYLOAD_MP2T_DIRECT; + + if (metadata->fcc_type != STREAM_FCC_TYPE_UNKNOWN && metadata->fcc_status == STREAM_FCC_STATUS_UNKNOWN) { + if (origin == STREAM_MEDIA_ORIGIN_FCC_UNICAST) { + metadata->fcc_status = STREAM_FCC_STATUS_ACTIVE; + } else if (origin == STREAM_MEDIA_ORIGIN_FCC_MULTICAST) { + metadata->fcc_status = STREAM_FCC_STATUS_FALLBACK; + } + } +} + +void stream_send_http_headers(connection_t *conn, const char *content_type, const char *extra_headers) { + static const char exposed_headers[] = + "R2H-Upstream-Protocol, R2H-Upstream-Transport, R2H-Upstream-Payload, " + "R2H-Playback-Scale, R2H-Playback-Range, R2H-Media-Duration, R2H-FCC-Type, R2H-FCC-Status"; + stream_metadata_t *metadata; + char headers[STREAM_METADATA_HEADERS_SIZE]; + char number[64]; + size_t length = 0; + int failed = 0; + + if (!conn) { + return; + } + + metadata = &conn->stream.metadata; + headers[0] = '\0'; + + if (extra_headers && extra_headers[0]) + failed |= stream_metadata_append_raw(headers, sizeof(headers), &length, extra_headers) < 0; + + switch (metadata->upstream_protocol) { + case STREAM_UPSTREAM_RTSP: + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-Upstream-Protocol", "rtsp") < 0; + break; + case STREAM_UPSTREAM_MULTICAST: + failed |= + stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-Upstream-Protocol", "multicast") < 0; + break; + default: + break; + } + + switch (metadata->upstream_transport) { + case STREAM_TRANSPORT_TCP_INTERLEAVED: + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-Upstream-Transport", + "tcp-interleaved") < 0; + break; + case STREAM_TRANSPORT_UDP: + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-Upstream-Transport", "udp") < 0; + break; + default: + break; + } + + switch (metadata->upstream_payload) { + case STREAM_PAYLOAD_MP2T_RTP: + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-Upstream-Payload", "mp2t-rtp") < 0; + break; + case STREAM_PAYLOAD_MP2T_DIRECT: + failed |= + stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-Upstream-Payload", "mp2t-direct") < 0; + break; + default: + break; + } + + if (metadata->playback_scale_known && isfinite(metadata->playback_scale)) { + stream_metadata_format_number(metadata->playback_scale, number, sizeof(number)); + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-Playback-Scale", number) < 0; + } + if (metadata->playback_range[0]) { + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-Playback-Range", + metadata->playback_range) < 0; + } + if (metadata->media_duration_known && isfinite(metadata->media_duration)) { + stream_metadata_format_number(metadata->media_duration, number, sizeof(number)); + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-Media-Duration", number) < 0; + } + + switch (metadata->fcc_type) { + case STREAM_FCC_TYPE_TELECOM: + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-FCC-Type", "telecom") < 0; + break; + case STREAM_FCC_TYPE_HUAWEI: + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-FCC-Type", "huawei") < 0; + break; + default: + break; + } + + switch (metadata->fcc_status) { + case STREAM_FCC_STATUS_ACTIVE: + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-FCC-Status", "active") < 0; + break; + case STREAM_FCC_STATUS_FALLBACK: + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "R2H-FCC-Status", "fallback") < 0; + break; + default: + break; + } + + if (config.cors_allow_origin && config.cors_allow_origin[0]) { + failed |= stream_metadata_append_header(headers, sizeof(headers), &length, "Access-Control-Expose-Headers", + exposed_headers) < 0; + } + + if (failed) { + logger(LOG_ERROR, "Failed to build stream metadata HTTP headers"); + send_http_headers(conn, STATUS_200, content_type, extra_headers); + } else { + send_http_headers(conn, STATUS_200, content_type, headers); + } + metadata->frozen = 1; +} + void stream_on_client_drain(stream_context_t *ctx) { /* Hot path: every successful client write hits this. Bail out cheaply when * no upstream is paused (vast majority of streams). */ @@ -33,7 +251,7 @@ void stream_on_client_drain(stream_context_t *ctx) { rtsp_resume_upstream(&ctx->rtsp); } -int stream_process_rtp_payload(stream_context_t *ctx, buffer_ref_t *buf_ref) { +int stream_process_rtp_payload(stream_context_t *ctx, buffer_ref_t *buf_ref, stream_media_origin_t origin) { uint8_t *data_ptr = (uint8_t *)buf_ref->data + buf_ref->data_offset; uint8_t *payload; int payload_len; @@ -52,6 +270,8 @@ int stream_process_rtp_payload(stream_context_t *ctx, buffer_ref_t *buf_ref) { return 0; } + stream_metadata_note_media(ctx, pkt_type, payload, payload_len, origin); + if (pkt_type == 0) { /* Non-RTP packet - pass through directly (no reordering needed) */ if (ctx->snapshot.initialized) { @@ -109,6 +329,9 @@ int stream_handle_fd_event(stream_context_t *ctx, int fd, uint32_t events, int64 logger(LOG_DEBUG, "RTSP: found duration: %0.3f", ctx->rtsp.r2h_duration_value); return -2; } + if (result == -3) { + return -3; + } return -1; } return 0; @@ -147,6 +370,70 @@ int stream_handle_fd_event(stream_context_t *ctx, int fd, uint32_t events, int64 return 0; } +static int stream_init_rtsp_control(stream_context_t *ctx, service_t *service, int status_index, int metadata_probe) { + seek_parse_result_t seek_parse_result; + const char *resolved_seek_param_name = service->seek_param_name; + char resolved_rtsp_url[2048]; + + rtsp_session_init(&ctx->rtsp); + ctx->rtsp.status_index = status_index; + ctx->rtsp.epoll_fd = ctx->epoll_fd; + ctx->rtsp.conn = ctx->conn; + ctx->rtsp.metadata_probe = metadata_probe; + ctx->rtsp.upstream_ifname = get_upstream_interface_for_rtsp(service->ifname); + + if (!service->rtsp_url) { + logger(LOG_ERROR, "RTSP URL not found in service configuration"); + return -1; + } + + if (service_parse_seek_value(service->seek_param_value, service->seek_begin_offset_seconds, + service->seek_end_offset_seconds, service->user_agent, service->seek_mode, + service->seek_mode_tz_explicit, service->seek_mode_tz_offset_seconds, + service->seek_mode_window_seconds, &seek_parse_result) != 0) { + logger(LOG_ERROR, "RTSP: Failed to parse seek parameters"); + return -1; + } + + if (service_format_recent_seek_range(&seek_parse_result, ctx->rtsp.playseek_range_start, + sizeof(ctx->rtsp.playseek_range_start)) > 0) { + ctx->rtsp.use_playseek_range = 1; + resolved_seek_param_name = NULL; + } + + if (service_resolve_upstream_url(service->rtsp_url, resolved_seek_param_name, &seek_parse_result, resolved_rtsp_url, + sizeof(resolved_rtsp_url)) < 0) { + logger(LOG_ERROR, "RTSP: Failed to resolve upstream URL"); + return -1; + } + if (rtsp_parse_server_url(&ctx->rtsp, resolved_rtsp_url, NULL, NULL) < 0) { + logger(LOG_ERROR, "RTSP: Failed to parse URL"); + return -1; + } + if (rtsp_connect(&ctx->rtsp) < 0) { + logger(LOG_ERROR, "RTSP: Failed to initiate connection"); + return -1; + } + + logger(LOG_DEBUG, "RTSP: Async connection initiated, state=%d", ctx->rtsp.state); + return 0; +} + +int stream_context_init_rtsp_metadata_probe(stream_context_t *ctx, connection_t *conn, service_t *service, + int epoll_fd) { + if (!ctx || !conn || !service || service->service_type != SERVICE_RTSP) + return -1; + + memset(ctx, 0, sizeof(*ctx)); + ctx->conn = conn; + ctx->service = service; + ctx->epoll_fd = epoll_fd; + ctx->status_index = -1; + ctx->last_status_update = get_time_ms(); + stream_metadata_init(&ctx->metadata, service); + return stream_init_rtsp_control(ctx, service, -1, 1); +} + /* Initialize context for unified worker epoll (non-blocking, no own loop) */ int stream_context_init_for_worker(stream_context_t *ctx, connection_t *conn, service_t *service, int epoll_fd, int status_index, int is_snapshot) { @@ -160,6 +447,7 @@ int stream_context_init_for_worker(stream_context_t *ctx, connection_t *conn, se ctx->total_bytes_sent = 0; ctx->last_bytes_sent = 0; ctx->last_status_update = get_time_ms(); + stream_metadata_init(&ctx->metadata, service); /* Initialize media path depending on service type */ if (service->service_type == SERVICE_HTTP) { @@ -244,58 +532,8 @@ int stream_context_init_for_worker(stream_context_t *ctx, connection_t *conn, se fec_init(&ctx->fec, service->fec_port, &ctx->reorder); if (service->service_type == SERVICE_RTSP) { - /* Initialize RTSP session */ - seek_parse_result_t seek_parse_result; - const char *resolved_seek_param_name = service->seek_param_name; - - rtsp_session_init(&ctx->rtsp); - ctx->rtsp.status_index = status_index; - ctx->rtsp.epoll_fd = ctx->epoll_fd; - ctx->rtsp.conn = conn; - ctx->rtsp.upstream_ifname = get_upstream_interface_for_rtsp(service->ifname); - if (!service->rtsp_url) { - logger(LOG_ERROR, "RTSP URL not found in service configuration"); - return -1; - } - - if (service_parse_seek_value(service->seek_param_value, service->seek_begin_offset_seconds, - service->seek_end_offset_seconds, service->user_agent, service->seek_mode, - service->seek_mode_tz_explicit, service->seek_mode_tz_offset_seconds, - service->seek_mode_window_seconds, &seek_parse_result) != 0) { - logger(LOG_ERROR, "RTSP: Failed to parse seek parameters"); - return -1; - } - - if (service_format_recent_seek_range(&seek_parse_result, ctx->rtsp.playseek_range_start, - sizeof(ctx->rtsp.playseek_range_start)) > 0) { - ctx->rtsp.use_playseek_range = 1; - resolved_seek_param_name = NULL; - } - - /* Resolve URL templates / seek params, then parse. Always forward the - * parsed seek result so URL templates (${start}/${end}/...) expand even - * when use_playseek_range=1; suppression of the literal playseek query - * append on the non-template branch is already handled via the - * resolved_seek_param_name=NULL gating above. */ - char resolved_rtsp_url[2048]; - if (service_resolve_upstream_url(service->rtsp_url, resolved_seek_param_name, &seek_parse_result, - resolved_rtsp_url, sizeof(resolved_rtsp_url)) < 0) { - logger(LOG_ERROR, "RTSP: Failed to resolve upstream URL"); + if (stream_init_rtsp_control(ctx, service, status_index, 0) < 0) return -1; - } - if (rtsp_parse_server_url(&ctx->rtsp, resolved_rtsp_url, NULL, NULL) < 0) { - logger(LOG_ERROR, "RTSP: Failed to parse URL"); - return -1; - } - - if (rtsp_connect(&ctx->rtsp) < 0) { - logger(LOG_ERROR, "RTSP: Failed to initiate connection"); - return -1; - } - - /* Connection initiated - handshake will proceed asynchronously via event - * loop */ - logger(LOG_DEBUG, "RTSP: Async connection initiated, state=%d", ctx->rtsp.state); } else { /* Multicast-based services (FCC or direct multicast) */ mcast_session_init(&ctx->mcast); diff --git a/src/stream.h b/src/stream.h index c708571a..e708ec6f 100644 --- a/src/stream.h +++ b/src/stream.h @@ -18,6 +18,55 @@ * fallback to streaming */ #define SNAPSHOT_TIMEOUT_SEC 2 +#define STREAM_PLAYBACK_RANGE_SIZE 256 + +typedef enum { + STREAM_UPSTREAM_UNKNOWN = 0, + STREAM_UPSTREAM_RTSP, + STREAM_UPSTREAM_MULTICAST +} stream_upstream_protocol_t; + +typedef enum { + STREAM_TRANSPORT_UNKNOWN = 0, + STREAM_TRANSPORT_TCP_INTERLEAVED, + STREAM_TRANSPORT_UDP +} stream_upstream_transport_t; + +typedef enum { + STREAM_PAYLOAD_UNKNOWN = 0, + STREAM_PAYLOAD_MP2T_RTP, + STREAM_PAYLOAD_MP2T_DIRECT +} stream_upstream_payload_t; + +typedef enum { STREAM_FCC_TYPE_UNKNOWN = 0, STREAM_FCC_TYPE_TELECOM, STREAM_FCC_TYPE_HUAWEI } stream_fcc_type_t; + +typedef enum { + STREAM_FCC_STATUS_UNKNOWN = 0, + STREAM_FCC_STATUS_ACTIVE, + STREAM_FCC_STATUS_FALLBACK +} stream_fcc_status_t; + +typedef enum { + STREAM_MEDIA_ORIGIN_RTSP = 0, + STREAM_MEDIA_ORIGIN_MULTICAST, + STREAM_MEDIA_ORIGIN_FCC_UNICAST, + STREAM_MEDIA_ORIGIN_FCC_MULTICAST +} stream_media_origin_t; + +typedef struct stream_metadata_s { + stream_upstream_protocol_t upstream_protocol; + stream_upstream_transport_t upstream_transport; + stream_upstream_payload_t upstream_payload; + stream_fcc_type_t fcc_type; + stream_fcc_status_t fcc_status; + double playback_scale; + double media_duration; + char playback_range[STREAM_PLAYBACK_RANGE_SIZE]; + uint8_t playback_scale_known; + uint8_t media_duration_known; + uint8_t frozen; +} stream_metadata_t; + /* Stream processing context */ typedef struct stream_context_s { int epoll_fd; @@ -31,6 +80,9 @@ typedef struct stream_context_s { uint64_t last_bytes_sent; /* Bytes sent at last bandwidth calculation */ int64_t last_status_update; /* Last status update time in milliseconds */ + /* Metadata exposed as R2H-* HTTP response headers. */ + stream_metadata_t metadata; + /* FCC session for Fast Channel Change */ fcc_session_t fcc; @@ -69,6 +121,13 @@ typedef struct stream_context_s { int stream_context_init_for_worker(stream_context_t *ctx, connection_t *conn, service_t *service, int epoll_fd, int status_index, int is_snapshot); +/** + * Initialize an RTSP control-plane-only HEAD probe. The probe performs + * OPTIONS and DESCRIBE but never SETUP or PLAY and allocates no media sockets. + */ +int stream_context_init_rtsp_metadata_probe(stream_context_t *ctx, connection_t *conn, service_t *service, + int epoll_fd); + /** * Handle an event-ready fd that belongs to this stream context. * @param ctx Stream context @@ -79,6 +138,7 @@ int stream_context_init_for_worker(stream_context_t *ctx, connection_t *conn, se * 0: Success, continue processing * -1: Connection should be closed (error or graceful TEARDOWN complete) * -2: Duration query completed, send response to client + * -3: RTSP metadata probe completed, send HEAD response to client */ int stream_handle_fd_event(stream_context_t *ctx, int fd, uint32_t events, int64_t now); @@ -105,7 +165,19 @@ int stream_context_cleanup(stream_context_t *ctx); * @return bytes forwarded (>= 0) for streaming, 1 if I-frame captured for * snapshot, -1 on error */ -int stream_process_rtp_payload(stream_context_t *ctx, buffer_ref_t *buf_ref); +int stream_process_rtp_payload(stream_context_t *ctx, buffer_ref_t *buf_ref, stream_media_origin_t origin); + +/** Initialize static metadata from a parsed service. */ +void stream_metadata_init(stream_metadata_t *metadata, const service_t *service); + +/** Clear RTSP metadata learned from a server before following a redirect. */ +void stream_metadata_reset_rtsp_negotiation(stream_metadata_t *metadata); + +/** + * Send a successful HTTP response with any known stream metadata appended to + * extra_headers. extra_headers may be NULL. + */ +void stream_send_http_headers(connection_t *conn, const char *content_type, const char *extra_headers); /** * Notify that the client send queue has just been drained (some buffers diff --git a/src/worker.c b/src/worker.c index 448085e0..c5a939a2 100644 --- a/src/worker.c +++ b/src/worker.c @@ -513,6 +513,18 @@ int worker_run_event_loop(int *listen_sockets, int num_sockets, int notif_fd) { snprintf(response, sizeof(response), "{\"duration\": \"%0.3f\"}", c->stream.rtsp.r2h_duration_value); connection_queue_output_and_flush(c, (const uint8_t *)response, strlen(response)); + } else if (res == -3) { + /* A metadata probe has no RTSP media session to tear down. Close + * and unregister its control socket before queuing the HEAD + * response so a later upstream HUP cannot close the client. */ + if (stream_context_cleanup(&c->stream) != 0) { + logger(LOG_ERROR, "Worker: Metadata probe unexpectedly required async cleanup"); + http_send_503(c); + continue; + } + c->streaming = 0; + stream_send_http_headers(c, "video/mp2t", NULL); + connection_queue_output_and_flush(c, NULL, 0); } else if (!c->headers_sent && c->state != CONN_CLOSING) { /* Send 503 if headers not sent yet (no data ever arrived) */ http_send_503(c); From 6c2e22dc2f0baad1d9b2bb671ad0bb0d8ca84e5f Mon Sep 17 00:00:00 2001 From: Stackie Jia Date: Tue, 21 Jul 2026 17:41:49 +0800 Subject: [PATCH 2/3] fix(rtsp): handle NPT ranges and probe socket errors --- e2e/helpers/mock_rtsp.py | 9 +++++ e2e/test_rtsp_transport.py | 42 ++++++++++++++++++++++ src/rtsp.c | 72 ++++++++++++++++++++++++++++++++------ 3 files changed, 113 insertions(+), 10 deletions(-) diff --git a/e2e/helpers/mock_rtsp.py b/e2e/helpers/mock_rtsp.py index 562237f3..9d1f302e 100644 --- a/e2e/helpers/mock_rtsp.py +++ b/e2e/helpers/mock_rtsp.py @@ -32,6 +32,7 @@ def __init__( options_session_id: str | None = None, play_response_headers: list[tuple[str, str]] | None = None, close_after_describe: bool = False, + reset_after_describe: bool = False, host: str = "127.0.0.1", ): """ @@ -48,6 +49,8 @@ def __init__( must echo it (simulates HMS-style servers). close_after_describe: Close the control connection immediately after sending the DESCRIBE response. + reset_after_describe: Reset the control connection immediately after + sending the DESCRIBE response. host: Address to listen on (use "::1" for IPv6 loopback). """ self.host = host @@ -58,6 +61,7 @@ def __init__( self._options_session_id = options_session_id self._play_response_headers = play_response_headers or [] self._close_after_describe = close_after_describe + self._reset_after_describe = reset_after_describe self._server_sock: socket.socket | None = None self._thread: threading.Thread | None = None self._stop = threading.Event() @@ -197,6 +201,9 @@ def _handle(self, conn: socket.socket, addr: tuple) -> None: "Content-Length: %d\r\n\r\n%s" % (cseq, cb_header, len(sdp), sdp) ).encode() ) + if self._reset_after_describe: + conn.setsockopt(socket.SOL_SOCKET, socket.SO_LINGER, struct.pack("ii", 1, 0)) + return if self._close_after_describe: return elif method == "SETUP": @@ -247,6 +254,7 @@ def __init__( options_session_id: str | None = None, play_response_headers: list[tuple[str, str]] | None = None, close_after_describe: bool = False, + reset_after_describe: bool = False, host: str = "127.0.0.1", ): super().__init__( @@ -257,6 +265,7 @@ def __init__( options_session_id=options_session_id, play_response_headers=play_response_headers, close_after_describe=close_after_describe, + reset_after_describe=reset_after_describe, host=host, ) self._num_packets = num_packets diff --git a/e2e/test_rtsp_transport.py b/e2e/test_rtsp_transport.py index 76aafeba..967dbcae 100644 --- a/e2e/test_rtsp_transport.py +++ b/e2e/test_rtsp_transport.py @@ -120,6 +120,29 @@ def test_head_stops_after_describe(self, shared_r2h): finally: rtsp.stop() + def test_head_parses_clock_form_npt_duration(self, shared_r2h): + sdp = ( + "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=T\r\n" + "t=0 0\r\nm=video 0 RTP/AVP 33\r\n" + "a=range:npt=00:00:10.25-01:30:40.75\r\na=control:*\r\n" + ) + rtsp = MockRTSPServer(custom_sdp=sdp) + rtsp.start() + try: + status, headers, body = http_request( + "127.0.0.1", + shared_r2h.port, + "HEAD", + "/rtsp/127.0.0.1:%d/stream" % rtsp.port, + timeout=10.0, + ) + assert status == 200 + assert body == b"" + assert get_header(headers, "R2H-Media-Duration") == "5430.5" + assert rtsp.requests_received == ["OPTIONS", "DESCRIBE"] + finally: + rtsp.stop() + def test_head_metadata_takes_precedence_over_duration_query(self, shared_r2h): sdp = ( "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=T\r\n" @@ -165,6 +188,25 @@ def test_head_survives_upstream_close_after_describe(self, shared_r2h): finally: rtsp.stop() + def test_head_survives_upstream_reset_after_describe(self, shared_r2h): + rtsp = MockRTSPServer(reset_after_describe=True) + rtsp.start() + try: + status, headers, body = http_request( + "127.0.0.1", + shared_r2h.port, + "HEAD", + "/rtsp/127.0.0.1:%d/stream" % rtsp.port, + timeout=10.0, + ) + assert status == 200 + assert body == b"" + assert get_header(headers, "R2H-Upstream-Protocol") == "rtsp" + assert get_header(headers, "R2H-Upstream-Payload") == "mp2t-rtp" + assert rtsp.requests_received == ["OPTIONS", "DESCRIBE"] + finally: + rtsp.stop() + def test_head_probe_timeout_returns_503(self, shared_r2h): rtsp = MockRTSPServerSilent() rtsp.start() diff --git a/src/rtsp.c b/src/rtsp.c index d9ac2f38..8c2a5bd3 100644 --- a/src/rtsp.c +++ b/src/rtsp.c @@ -934,7 +934,7 @@ int rtsp_handle_socket_event(rtsp_session_t *session, uint32_t events) { /* Check for connection errors or hangup */ else if ((events & (POLLER_HUP | POLLER_ERR | POLLER_RDHUP)) && - !(session->metadata_probe && session->awaiting_response && (events & POLLER_IN) && !(events & POLLER_ERR))) { + !(session->metadata_probe && session->awaiting_response && (events & POLLER_IN))) { if (events & POLLER_ERR) { int sock_error = 0; socklen_t error_len = sizeof(sock_error); @@ -1012,7 +1012,8 @@ int rtsp_handle_socket_event(rtsp_session_t *session, uint32_t events) { * one edge-triggered event. Keep draining while each recv makes * progress so a fragmented DESCRIBE body is parsed before HUP. */ } while (response_result == RTSP_RESPONSE_OK && session->metadata_probe && session->awaiting_response && - (events & (POLLER_HUP | POLLER_RDHUP)) && session->response_buffer_pos > previous_buffer_pos); + (events & (POLLER_HUP | POLLER_ERR | POLLER_RDHUP)) && + session->response_buffer_pos > previous_buffer_pos); if (response_result < 0) { logger(LOG_ERROR, "RTSP: Failed to receive response"); rtsp_session_set_state(session, RTSP_STATE_ERROR); @@ -1026,6 +1027,15 @@ int rtsp_handle_socket_event(rtsp_session_t *session, uint32_t events) { return -3; } + /* A terminal socket event may be reported together with readable + * response data. The final DESCRIBE response above completes the probe; + * any earlier or incomplete response still requires a live connection. */ + if (session->metadata_probe && (events & (POLLER_HUP | POLLER_ERR | POLLER_RDHUP))) { + logger(LOG_ERROR, "RTSP: Metadata probe connection closed before DESCRIBE completed"); + rtsp_session_set_state(session, RTSP_STATE_ERROR); + return -1; + } + /* Re-enable POLLER_OUT for next request */ if (response_result == RTSP_RESPONSE_ADVANCE && session->epoll_fd >= 0) { if (poller_mod(session->epoll_fd, session->socket, @@ -2821,6 +2831,51 @@ static int rtsp_sdp_describes_mp2t_rtp(const char *sdp_body) { return 0; } +static int rtsp_parse_npt_time(const char *value, const char **end_out, double *seconds_out) { + char *component_end; + double first_component; + + errno = 0; + first_component = strtod(value, &component_end); + if (component_end == value || errno == ERANGE || !isfinite(first_component) || first_component < 0.0) + return -1; + + if (*component_end != ':') { + *end_out = component_end; + *seconds_out = first_component; + return 0; + } + + if (floor(first_component) != first_component) + return -1; + + const char *minutes_start = component_end + 1; + char *minutes_end; + long minutes; + + errno = 0; + minutes = strtol(minutes_start, &minutes_end, 10); + if (minutes_end == minutes_start || errno == ERANGE || minutes < 0 || minutes > 59 || *minutes_end != ':') + return -1; + + const char *seconds_start = minutes_end + 1; + char *seconds_end; + double seconds; + + errno = 0; + seconds = strtod(seconds_start, &seconds_end); + if (seconds_end == seconds_start || errno == ERANGE || !isfinite(seconds) || seconds < 0.0 || seconds >= 60.0) + return -1; + + double total = first_component * 3600.0 + (double)minutes * 60.0 + seconds; + if (!isfinite(total)) + return -1; + + *end_out = seconds_end; + *seconds_out = total; + return 0; +} + /** * Parse DESCRIBE response SDP body in a single pass. * @@ -2907,19 +2962,16 @@ static void rtsp_parse_describe_sdp(rtsp_session_t *session, const char *header_ len = (size_t)(val_end - val); if (len > 0 && len < sizeof(buf)) { - char *start_end; - char *end_end; + const char *start_end; + const char *end_end; double range_start; double range_end; memcpy(buf, val, len); buf[len] = '\0'; - errno = 0; - range_start = strtod(buf, &start_end); - if (start_end != buf && *start_end == '-') { - range_end = strtod(start_end + 1, &end_end); - if (end_end != start_end + 1 && *end_end == '\0' && errno != ERANGE && isfinite(range_start) && - isfinite(range_end) && range_end >= range_start) { + if (rtsp_parse_npt_time(buf, &start_end, &range_start) == 0 && *start_end == '-') { + if (rtsp_parse_npt_time(start_end + 1, &end_end, &range_end) == 0 && *end_end == '\0' && + range_end >= range_start) { logger(LOG_DEBUG, "RTSP: Range: %.3f-%.3f", range_start, range_end); if (session->r2h_duration) session->r2h_duration_value = (float)range_end; From 022daf8c350e98d198b37751b0784cdb0e13cba0 Mon Sep 17 00:00:00 2001 From: Stackie Jia Date: Tue, 21 Jul 2026 17:50:40 +0800 Subject: [PATCH 3/3] fix(build): link math library explicitly --- CMakeLists.txt | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/CMakeLists.txt b/CMakeLists.txt index 44de69da..32714733 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -151,7 +151,8 @@ endif() # ── Libraries ─────────────────────────────────────────────────────── find_package(Threads REQUIRED) -target_link_libraries(rtp2httpd PRIVATE Threads::Threads) +# Math functions require libm on some Unix platforms; macOS maps it to libSystem. +target_link_libraries(rtp2httpd PRIVATE Threads::Threads m) if(PLATFORM_LINUX) # rt library for clock_gettime (not needed on macOS)