From 1f56a602aca4cab258056b7e5e15c74bd115a1c7 Mon Sep 17 00:00:00 2001 From: Andrew Tridgell Date: Mon, 21 Sep 2026 20:46:50 +1000 Subject: [PATCH] video: relay and record lossless thermal Matroska Accept authenticated chunked Matroska publishers and preserve native FFV1 samples and frame metadata through fan-out and recording. Cache codec headers for late viewers and rotate recordings at complete cluster boundaries, retaining existing admission, viewer and quota policies. Provide a greyscale desktop thermal viewer and browser guidance, with tests for authentication, fragmented input, reconnects, slow viewers and recording rotation. --- Makefile | 4 + README.md | 63 ++++++++++ cleanup.cpp | 3 + conntdb.h | 1 + conntdb_lib.py | 2 + scripts/view_raw_thermal.py | 125 ++++++++++++++++++++ session.cpp | 6 +- tests/test_video_matroska.py | 190 ++++++++++++++++++++++++++++++ tests/webadmin/test_video_page.py | 13 ++ video.cpp | 108 +++++++++++++++-- videomkv.h | 117 ++++++++++++++++++ videorec.cpp | 2 +- videorec.h | 2 + videoview.cpp | 44 +++++-- videoview.h | 15 ++- webadmin/logs.py | 2 +- webadmin/routes_video.py | 7 +- webadmin/templates/video.html | 11 +- 18 files changed, 685 insertions(+), 30 deletions(-) create mode 100644 scripts/view_raw_thermal.py create mode 100644 tests/test_video_matroska.py create mode 100644 videomkv.h diff --git a/Makefile b/Makefile index bb2579d..21036a5 100644 --- a/Makefile +++ b/Makefile @@ -87,10 +87,14 @@ session.o: session.cpp session.h keydb.h binlog.o: binlog.cpp binlog.h session.h mavlink.h util.h cleanup.h $(MAVLINK_DIR)/protocol.h cleanup.o: cleanup.cpp cleanup.h keydb.h websocket.o: websocket.cpp websocket.h util.h +video.o videoview.o: videomkv.h + video.o: video.cpp video.h videoauth.h videots.h videostream.h videorec.h videoview.h httpreq.h videortsp.h videortmp.h conntdb.h keydb.h util.h videoauth.o: videoauth.cpp videoauth.h conntdb.h keydb.h videots.o: videots.cpp videots.h videostream.o: videostream.cpp videostream.h +video.o: videorec.h + videorec.o: videorec.cpp videorec.h session.h cleanup.h videoview.o: videoview.cpp videoview.h httpreq.h videostream.h videots.h videoauth.h keydb.h httpreq.o: httpreq.cpp httpreq.h diff --git a/README.md b/README.md index fbda681..b9e9484 100644 --- a/README.md +++ b/README.md @@ -666,3 +666,66 @@ runs one pytest invocation against exactly what you passed. SupportProxy is licensed under the GNU General Public License version 3 or later. See `COPYING.txt` for full license terms. + +### Lossless raw thermal (FFV1/Matroska) + +Updated MT11/SITL-MT11 firmware can publish its third stream to an allocated +video port using `[support_proxy] video3_port` (MAVLink `PROXY_VID3_PORT`). +Enable video and allocate three distinct ports on the proxy entry; no database +migration or additional slot flag is required. The raw stream uses the same +publish password or permitted MAVLink-session admission as ordinary video. +The camera's `RAW_STREAM_FPS` controls transmission rate. This requires both +the camera firmware and SupportProxy Matroska updates. + +Publishers use `PUT /vN.mkv?pw=...`, `Content-Type: video/x-matroska`, +`Transfer-Encoding: chunked` and `Expect: 100-continue`. An accepted publisher +receives `100 Continue` before sending the EBML header and Clusters. A wrong +or explicitly empty password fails even when session fallback is enabled. +The bounded parser accepts the APCG live profile: EBML header, unknown-sized +Segment, Info, Tracks and finite independent FFV1 Clusters (maximum 4 MiB). +HTTP chunk boundaries need not align with EBML elements. It is not a generic +upload endpoint for arbitrary Matroska layouts or inter-frame codecs. + +Viewers use `http://HOST:PORT/v3.mkv` for the camera's third stream. Existing +viewer-password (`?pw=...` or HTTP Basic), short-lived viewer tokens (`?t=...`), +WebSocket and enabled raw-TCP access policies apply. Late viewers receive the +cached Matroska header followed by the newest complete Cluster and subsequent +publisher bytes. Pixels and BlockAdditional telemetry are never transcoded. +Slow viewers are disconnected independently; a publisher disconnect closes its +viewers, who must reconnect to the next session. + +The slot's **record** option saves `.v3.mkv` files. Segments rotate on complete +Cluster boundaries with a repeated header, and participate in the existing +video quota and log-download listing. They do not use the MPEG-TS browser log +player. The live web page identifies active Matroska publishers and offers a +desktop viewer instead of its MPEG-TS player (reload after changing formats). + +MAVProxy connected to this proxy discovers the correct URL automatically: + +```text +module load camera +camera discover +camera view rawthermal +``` + +For standalone viewing, install a MAVProxy version with the raw thermal +reader, PyAV >= 18.1, numpy and OpenCV with GUI support, then run: + +```sh +python3 scripts/view_raw_thermal.py http://HOST:PORT/v3.mkv +# Or use a MAVProxy checkout: +python3 scripts/view_raw_thermal.py --mavproxy /path/to/MAVProxy http://HOST:PORT/v3.mkv +``` + +Greyscale is the default. Hover for pixel temperature; Space pauses, C selects +optional palettes, S saves native little-endian uint16 pixels and frame JSON, +and Q/Escape exits. The window readout is separate from the image. A display +palette never changes the saved samples. `--headless --frames 10` validates +reception without a display. ffplay can also display the stream, but lacks the +temperature and capture-metadata handling of the thermal viewers. The existing +browser MPEG-TS player cannot decode this FFV1/16-bit stream. + +Tests: `python3 -m pytest tests/test_video_matroska.py` covers authentication, +fragmentation, late joining, restart, slow viewers, malformed lengths and +recording rotation. The camera repository's `sitl/test_support_proxy.py +--raw-thermal --reconnect` additionally checks actual FFV1 pixels and telemetry. diff --git a/cleanup.cpp b/cleanup.cpp index 0230ca0..51527ab 100644 --- a/cleanup.cpp +++ b/cleanup.cpp @@ -128,6 +128,9 @@ static session_kind session_file_kind(const char *name) && name[n - 4] >= '1' && name[n - 4] <= '9') { return SESSION_VIDEO; } + if (n > 7 && strcmp(name+n-4, ".mkv") == 0 && + name[n-7] == '.' && name[n-6] == 'v' && name[n-5] >= '1' && name[n-5] <= '9') + return SESSION_VIDEO; return SESSION_NONE; } diff --git a/conntdb.h b/conntdb.h index 4c5ed81..43cb96f 100644 --- a/conntdb.h +++ b/conntdb.h @@ -58,6 +58,7 @@ #define CONN_APP_HTTP 3 #define CONN_APP_SRT 4 #define CONN_APP_RTMP 5 +#define CONN_APP_MATROSKA 7 #define CONN_APP_RTP 6 // bare RTP over UDP /* diff --git a/conntdb_lib.py b/conntdb_lib.py index 85e82a6..5b19ce6 100644 --- a/conntdb_lib.py +++ b/conntdb_lib.py @@ -74,6 +74,7 @@ CONN_APP_SRT = 4 CONN_APP_RTMP = 5 CONN_APP_RTP = 6 +CONN_APP_MATROSKA = 7 APP_NAMES = { CONN_APP_MAVLINK: 'mavlink', @@ -83,6 +84,7 @@ CONN_APP_SRT: 'srt', CONN_APP_RTMP: 'rtmp', CONN_APP_RTP: 'rtp', + CONN_APP_MATROSKA: 'matroska', } # Video rows occupy a conn_index range disjoint from the MAVLink ones so diff --git a/scripts/view_raw_thermal.py b/scripts/view_raw_thermal.py new file mode 100644 index 0000000..3792c26 --- /dev/null +++ b/scripts/view_raw_thermal.py @@ -0,0 +1,125 @@ +#!/usr/bin/env python3 +"""View lossless APCG thermal Matroska from SupportProxy or a local camera. + +Requires the MAVProxy raw thermal reader, PyAV >= 18.1, numpy and OpenCV. +Space pauses; C cycles greyscale/inferno/turbo; S saves native pixels and JSON. +""" +import argparse +import json +from pathlib import Path +import sys +import threading +import time + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument('uri') + parser.add_argument('--mavproxy', type=Path, help='MAVProxy checkout containing the raw thermal reader') + parser.add_argument('--output', type=Path, default=Path('thermal-captures')) + parser.add_argument('--headless', action='store_true', help='validate decoded frames without opening a window') + parser.add_argument('--frames', type=int, default=0, help='exit after this many received frames (0: unlimited)') + args = parser.parse_args() + if args.frames < 0: + parser.error('--frames must be nonnegative') + if args.mavproxy: + sys.path.insert(0, str(args.mavproxy.resolve())) + import numpy as np + from MAVProxy.modules.mavproxy_camera.thermal_stream import ThermalReader + if not args.headless: + import cv2 + stop = threading.Event() + lock = threading.Lock() + latest, error, count = None, None, 0 + + def receive(): + nonlocal latest, error, count + while not stop.is_set(): + reader = None + try: + reader = ThermalReader(args.uri) + for pixels, metadata in reader.frames(): + with lock: + latest = pixels, metadata + error = None + count += 1 + if stop.is_set() or (args.frames and count >= args.frames): + return + raise RuntimeError('stream ended') + except Exception as exc: + # Exception URLs can contain the viewer password. + with lock: + error = 'Stream unavailable (%s); reconnecting' % type(exc).__name__ + finally: + if reader: + reader.close() + stop.wait(1) + + worker = threading.Thread(target=receive, daemon=True) + worker.start() + shown = None + cursor = None + paused, palette = False, 0 + window = 'Raw Thermal' + if not args.headless: + cv2.namedWindow(window, cv2.WINDOW_AUTOSIZE) + def mouse(_event, x, y, _flags, _data): + nonlocal cursor + cursor = (x, y) if 0 <= x < 640 and 0 <= y < 512 else None + cv2.setMouseCallback(window, mouse) + try: + while True: + with lock: + frame, status, received = latest, error, count + if not paused and frame is not None: + shown = frame + if args.headless: + if args.frames and received >= args.frames: + print('Decoded %u native 16-bit frames with metadata' % received) + break + time.sleep(.02) + continue + if shown is not None: + pixels, metadata = shown + lo, hi = int(pixels.min()), int(pixels.max()) + grey = ((pixels.astype(np.float32)-lo)*(255/max(1, hi-lo))).astype(np.uint8) + if metadata['rotation_deg'] == 180: + grey = grey[::-1, ::-1] + display = cv2.cvtColor(grey, cv2.COLOR_GRAY2BGR) if palette == 0 else cv2.applyColorMap( + grey, (cv2.COLORMAP_INFERNO, cv2.COLORMAP_TURBO)[palette-1]) + display = cv2.copyMakeBorder(display, 0, 68, 0, 0, cv2.BORDER_CONSTANT) + text = 'Min %.2f C Max %.2f C%s' % ( + metadata['minimum_c'], metadata['maximum_c'], ' PAUSED' if paused else '') + pixel_text = 'Hover for temperature; Space pause, C palette, S save' + if cursor: + x, y = cursor + sx, sy = (639-x, 511-y) if metadata['rotation_deg'] == 180 else (x, y) + temp = int(pixels[sy, sx])*metadata['temperature_scale_k']+metadata['temperature_offset_k']-273.15 + pixel_text = 'Pixel (%u,%u): %.3f C' % (x, y, temp) + cv2.putText(display, text, (8, 536), cv2.FONT_HERSHEY_SIMPLEX, .55, (255,255,255), 1) + cv2.putText(display, status or pixel_text, (8, 564), cv2.FONT_HERSHEY_SIMPLEX, .5, (255,255,255), 1) + cv2.imshow(window, display) + key = cv2.waitKey(20) & 255 + if key in (27, ord('q')) or cv2.getWindowProperty(window, cv2.WND_PROP_VISIBLE) < 1: + break + if key == ord(' '): paused = not paused + if key == ord('c'): palette = (palette+1) % 3 + if key == ord('s') and shown is not None: + pixels, metadata = shown + args.output.mkdir(parents=True, exist_ok=True) + stem = args.output / ('%u_%u' % (metadata['capture_monotonic_us'], metadata['frame_id'])) + with stem.with_suffix('.bin').open('xb') as f: + f.write(pixels.astype('= args.frames: break + finally: + stop.set() + # PyAV containers are owned and closed by the reader thread. Its + # bounded network timeout lets shutdown complete without racing decode. + worker.join(timeout=6) + if not args.headless: cv2.destroyAllWindows() + + +if __name__ == '__main__': + main() diff --git a/session.cpp b/session.cpp index ab070c2..ec0368f 100644 --- a/session.cpp +++ b/session.cpp @@ -86,9 +86,9 @@ static bool basename_free(const char *dir, const char *candidate) } } for (int slot = 1; slot <= KEY_MAX_VIDEO_PORTS; slot++) { - snprintf(p, sizeof(p), "%s/%s.v%d.ts", dir, candidate, slot); - if (stat(p, &st) == 0) { - return false; + for (const char *ext : {"ts", "mkv"}) { + snprintf(p, sizeof(p), "%s/%s.v%d.%s", dir, candidate, slot, ext); + if (stat(p, &st) == 0) return false; } } return true; diff --git a/tests/test_video_matroska.py b/tests/test_video_matroska.py new file mode 100644 index 0000000..13d1ea5 --- /dev/null +++ b/tests/test_video_matroska.py @@ -0,0 +1,190 @@ +"""Matroska relay: admission, fragmented input, late joins, rotation and loss.""" +import socket +import time +from urllib.parse import quote + +import keydb_lib +from test_video_child import Proxy, _make_workdir, PORT_ENG, VPORT +from test_video_view import http_get, ws_connect, ws_read_payload, mint_token + +PASSWORD = 'raw&publish?=test' + + +def element(eid, payload): + width = max(1, (eid.bit_length()+7)//8) + n = next(n for n in range(1, 9) if len(payload) < (1 << (7*n))-1) + return eid.to_bytes(width, 'big') + ((1 << (7*n)) | len(payload)).to_bytes(n, 'big') + payload + + +# Transport fixtures deliberately contain EBML-looking bytes inside frames: +# scanning for a magic Cluster signature instead of parsing sizes would fail. +HEADER = (element(0x1A45DFA3, b'matroska') + bytes.fromhex('1853806701ffffffffffffff') + + element(0x1549A966, b'info') + element(0x1654AE6B, b'V_FFV1')) + + +def cluster(i, size=128): + return element(0x1F43B675, i.to_bytes(4, 'big') + bytes.fromhex('1f43b675') + bytes([i % 256])*size) + + +def publish(password=PASSWORD, fragmented=False): + s = socket.create_connection(('127.0.0.1', VPORT), timeout=4) + target = '/v1.mkv' + ('' if password is None else '?pw='+quote(password, safe='')) + request = ('PUT %s HTTP/1.1\r\nHost: localhost\r\nContent-Type: video/x-matroska\r\n' + 'Transfer-Encoding: chunked\r\nExpect: 100-continue\r\n\r\n' % target).encode() + if fragmented: + for byte in request: + s.sendall(bytes([byte])) + time.sleep(.001) + else: + s.sendall(request) + response = b'' + while b'\r\n\r\n' not in response: + part = s.recv(1024) + if not part: break + response += part + return s, response + + +def send(s, data, fragment=0): + chunk = ('%x\r\n' % len(data)).encode()+data+b'\r\n' + if fragment: + for offset in range(0, len(chunk), fragment): s.sendall(chunk[offset:offset+fragment]) + else: + s.sendall(chunk) + + +def receive(s, data, count): + while len(data) < count: + part = s.recv(count-len(data)) + assert part, 'unexpected disconnect' + data += part + return data + + +def start(tmp_path, viewer=False): + wd = _make_workdir(tmp_path, publish_pass=PASSWORD) + db = keydb_lib.init_db(str(wd/'keys.tdb')) + db.transaction_start() + keydb_lib.set_video_slot_flag(db, PORT_ENG, 0, 'record') + if viewer: keydb_lib.set_video_viewer_pass(db, PORT_ENG, 'viewer') + db.transaction_prepare_commit(); db.transaction_commit(); db.close() + proxy = Proxy(wd) + assert proxy.wait_for('video slot 0 listening'), proxy.log + return wd, proxy + + +def test_auth_fragmentation_late_join_and_reconnect(tmp_path): + wd, proxy = start(tmp_path, viewer=True) + try: + for wrong in ('wrong', '', None): + s, response = publish(wrong) + s.close() + assert b'403' in response + pub, response = publish(fragmented=True) + assert b'100 Continue' in response + send(pub, HEADER, fragment=1) + send(pub, cluster(1), fragment=3) + send(pub, cluster(2), fragment=7) + assert proxy.wait_for('join=ready'), proxy.log + bad, response, _ = http_get(VPORT, '/v1.mkv') + bad.close() + assert b'401' in response + view, response, body = http_get(VPORT, '/v1.mkv?pw=viewer') + assert b'video/x-matroska' in response + expected = HEADER+cluster(2) + assert receive(view, body, len(expected)) == expected + busy, response = publish() + busy.close() + assert b'409' in response + send(pub, cluster(3), fragment=11) + assert receive(view, b'', len(cluster(3))) == cluster(3) + pub.close() + assert view.recv(1024) == b'' + view.close() + pub, response = publish() + assert b'100 Continue' in response + send(pub, HEADER+cluster(10)) + time.sleep(.15) + view, response, body = http_get(VPORT, '/v1.mkv?pw=viewer') + expected = HEADER+cluster(10) + assert receive(view, body, len(expected)) == expected + view.close(); pub.close() + assert list(wd.rglob('*.mkv')) + finally: + proxy.stop() + + +def test_record_rotates_at_cluster_boundaries(tmp_path, monkeypatch): + monkeypatch.setenv('SUPPORTPROXY_VIDEO_SEGMENT_BYTES', '300') + wd, proxy = start(tmp_path) + try: + pub, response = publish() + assert b'100 Continue' in response + send(pub, HEADER) + frames = [cluster(i, 220) for i in range(12)] + for frame in frames: send(pub, frame, fragment=17) + time.sleep(.4) + pub.close() + time.sleep(.4) + files = sorted(wd.rglob('*.mkv'), key=lambda p:p.stat().st_mtime_ns) + assert len(files) >= 6 + recovered = b'' + for path in files: + data = path.read_bytes() + assert data.startswith(HEADER) + recovered += data[len(HEADER):] + assert recovered == b''.join(frames) + finally: + proxy.stop() + + +def test_bad_sizes_release_slot_and_slow_viewer_isolated(tmp_path, monkeypatch): + monkeypatch.setenv('SUPPORTPROXY_VIDEO_RING_BYTES', '65536') + _, proxy = start(tmp_path) + try: + pub, response = publish() + assert b'100' in response + send(pub, b'\x00') # invalid EBML ID + assert pub.recv(1024) == b'' + pub.close() + pub, response = publish() + assert b'100' in response + pub.sendall(b'fffffffffff\r\n') # reject without allocating payload + assert pub.recv(1024) == b'' + pub.close() + pub, response = publish() + assert b'100' in response + send(pub, HEADER+cluster(1, 12000)) + assert proxy.wait_for('join=ready'), proxy.log + slow, response, _ = http_get(VPORT, '/v1.mkv') + for i in range(2, 500): send(pub, cluster(i, 12000)) + send(pub, cluster(500, 12000)) + time.sleep(.3) + fast, response, body = http_get(VPORT, '/v1.mkv') + expected = HEADER+cluster(500, 12000) + assert receive(fast, body, len(expected)) == expected + assert proxy.wait_for('lapped|chronically behind'), proxy.log + fast.close(); slow.close(); pub.close() + finally: + proxy.stop() + + +def test_websocket_prefix_and_http_token(tmp_path): + wd, proxy = start(tmp_path, viewer=True) + try: + pub, response = publish() + assert b'100' in response + send(pub, HEADER+cluster(4)) + assert proxy.wait_for('join=ready'), proxy.log + token = mint_token(wd, PORT_ENG, 0) + http, response, body = http_get(VPORT, '/v1.mkv?t='+token) + expected = HEADER+cluster(4) + assert b'200' in response + assert receive(http, body, len(expected)) == expected + http.close() + ws, response, body = ws_connect(VPORT, '/v1?t='+token) + assert b'101' in response + assert ws_read_payload(ws, 2, body) == expected + ws.close(); pub.close() + finally: + proxy.stop() diff --git a/tests/webadmin/test_video_page.py b/tests/webadmin/test_video_page.py index 5f73587..e445ca0 100644 --- a/tests/webadmin/test_video_page.py +++ b/tests/webadmin/test_video_page.py @@ -467,3 +467,16 @@ def test_ffplay_still_probes_enough_to_find_the_stream(self, client, def test_vlc_command_caps_its_cache(self, client, keydb_path): html = self._page(client, keydb_path) assert 'network-caching' in html + + +def test_raw_thermal_offers_desktop_viewer(client, keydb_path, monkeypatch): + from types import SimpleNamespace + from webadmin import connections + import conntdb_lib + _enable_video(keydb_path, ALICE_PORT2, ports=(VPORT, VPORT+1, VPORT+2)) + monkeypatch.setattr(connections, 'list_for_port2', lambda port: [SimpleNamespace( + stream_idx=2, role=conntdb_lib.CONN_ROLE_VIDEO_PUB, app_proto=conntdb_lib.CONN_APP_MATROSKA)]) + login_as(client, ALICE_PORT1, ALICE_PASS) + html = client.get('/video/').get_data(as_text=True) + assert 'view_raw_thermal.py' in html and '/v3.mkv' in html + assert 'id="player3"' not in html and 'id="player1"' in html diff --git a/video.cpp b/video.cpp index 758691d..d5e0dd6 100644 --- a/video.cpp +++ b/video.cpp @@ -33,7 +33,7 @@ #include "videoauth.h" #include "videostream.h" #include "videorec.h" -#include "videots.h" +#include "videomkv.h" #include "videortmp.h" #include "videortsp.h" #include "videoview.h" @@ -192,9 +192,12 @@ struct Slot { */ PendingRtmp pending[VIDEO_MAX_PENDING_RTMP]; + int mkv_fd = -1; + VideoChunks chunks; + // media VideoRing ring; - TSScanner scanner; + VideoScanner scanner; VideoWriter rec; VideoViewer viewers[VIDEO_MAX_VIEWERS]; uint32_t viewer_events[VIDEO_MAX_VIEWERS] {}; // last epoll mask armed @@ -296,6 +299,9 @@ class VideoChild { bool load_entry(void); int bind_slots(void); void signal_ready(int err); + void handle_mkv(Slot &s, int idx, int fd, struct sockaddr_in &from, time_t now); + bool pump_mkv(Slot &s, int idx, time_t now); + void close_mkv(Slot &s, int idx, const char *why); void handle_udp(Slot &s, int idx); void handle_tcp(Slot &s, int idx); void tick(time_t now); @@ -524,6 +530,7 @@ void VideoChild::handle_udp(Slot &s, int idx) if (n <= 0) { return; } + if (s.mkv_fd >= 0) return; // this slot belongs to the authenticated TCP publisher if (size_t(n) > sizeof(buf)) { // MSG_TRUNC reports the real length: a truncated RTP packet // would be forwarded as a whole one, so drop it instead. @@ -717,7 +724,8 @@ void VideoChild::latch_publisher(Slot &s, int idx, time_t now) s.pub_last = now; s.pub_bytes = 0; s.ring.init(video_ring_bytes()); - s.scanner = TSScanner(); + s.scanner = VideoScanner(); + s.rec.set_matroska(false); s.had_anchor = false; s.reset_rtp(); s.recording = (video_slot_opts_of(ke_, unsigned(idx)) @@ -730,6 +738,77 @@ void VideoChild::latch_publisher(Slot &s, int idx, time_t now) last_tick_ = 0; // snapshot connections.tdb promptly } +void VideoChild::handle_mkv(Slot &s, int idx, int fd, + struct sockaddr_in &from, time_t now) +{ + uint8_t buf[2048]; + ssize_t n = recv(fd, buf, sizeof(buf), MSG_PEEK); + HttpRequest req; + if (n <= 0 || req.feed(buf, size_t(n)) != 1) { close(fd); return; } + char path[32]; snprintf(path, sizeof(path), "/v%d.mkv", idx+1); + int code = 400; + std::string pw; + bool offered = http_query_value(req.target(), "pw", pw); + auto admission = auth_.admit(ke_, uint32_t(from.sin_addr.s_addr), offered ? &pw : nullptr, + true, session_ok(idx), open_publish(idx), now); + bool valid = req.path() == path && req.header_count("Content-Type") == 1 && + req.header("Content-Type") == "video/x-matroska" && + req.header_count("Transfer-Encoding") == 1 && req.header("Transfer-Encoding") == "chunked" && + req.header_count("Content-Length") == 0 && req.header("Expect") == "100-continue"; + if (admission != VIDEO_ADMIT_OK) code = 403; + else if (s.has_pub || s.rtsp.running() || s.mkv_fd >= 0) code = 409; + else if (valid) code = 100; + if (code != 100) { + const auto response = http_simple_response(code, "publish rejected", "text/plain", "Cannot publish this stream.\n"); + (void)send(fd, response.data(), response.size(), MSG_NOSIGNAL); + close(fd); return; + } + const size_t header_size = size_t(n)-req.leftover().size(); + if (recv(fd, buf, header_size, 0) != ssize_t(header_size)) { close(fd); return; } + const char response[] = "HTTP/1.1 100 Continue\r\n\r\n"; + if (send(fd, response, sizeof(response)-1, MSG_NOSIGNAL) != sizeof(response)-1) { close(fd); return; } + s.pub_ip_be = from.sin_addr.s_addr; s.pub_port_be = from.sin_port; + latch_publisher(s, idx, now); + s.mkv_fd = fd; s.chunks = VideoChunks(); s.scanner.matroska = true; + s.rec.set_matroska(true); + s.scanner.cluster = [&s](const uint8_t *data, size_t length) { + if (!s.recording || s.rec.stopped()) return; + const time_t t = time(nullptr); + if (s.rec.rotation_due(t)) s.rec.rotate(t, true); + if (!s.rec.is_open()) { + const auto &header = s.scanner.prefix; + if (!s.rec.write(reinterpret_cast(header.data()), header.size(), t)) return; + } + s.rec.write(data, length, t); + }; + struct epoll_event ev {}; ev.events = EPOLLIN | EPOLLRDHUP; ev.data.fd = fd; + epoll_ctl(epfd_, EPOLL_CTL_ADD, fd, &ev); + printf("[%d] video slot %d Matroska publisher connected\n", port2_, idx); +} + +bool VideoChild::pump_mkv(Slot &s, int idx, time_t now) +{ + uint8_t buf[65536]; + const ssize_t n = recv(s.mkv_fd, buf, sizeof(buf), 0); + if (n < 0) return errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR; + if (!n) return false; + s.pub_last = now; + bool ok = s.chunks.feed(buf, size_t(n), [&s](const uint8_t *data, size_t length) { + s.scanner.feed(data, length, s.ring.write_pos()); + s.ring.write(data, length); s.pub_bytes += length; + }); + (void)idx; + return ok && !s.scanner.failed; +} + +void VideoChild::close_mkv(Slot &s, int idx, const char *why) +{ + epoll_ctl(epfd_, EPOLL_CTL_DEL, s.mkv_fd, nullptr); + close(s.mkv_fd); s.mkv_fd = -1; s.has_pub = false; + s.rec.close_segment(); s.recording = false; + end_stream(s, idx, why); +} + void VideoChild::handle_rtsp(Slot &s, int idx, int fd, struct sockaddr_in &from, time_t now, splice_proto_t proto) @@ -831,7 +910,7 @@ void VideoChild::handle_rtsp(Slot &s, int idx, int fd, close(fd); return; } - if (s.rtsp.running() || s.rtsp_client_fd >= 0 + if (s.mkv_fd >= 0 || s.rtsp.running() || s.rtsp_client_fd >= 0 || (s.has_pub && now - s.pub_last <= VIDEO_PUB_IDLE_S)) { log_reject(s, idx, uint32_t(from.sin_addr.s_addr), VIDEO_ADMIT_SLOT_BUSY, now); @@ -1134,7 +1213,7 @@ bool VideoChild::promote_pending(Slot &s, int idx, PendingRtmp &p, // Authorised -- but the slot may have been taken while this one was // still negotiating. - if (s.rtsp.running() || s.rtsp_client_fd >= 0 + if (s.mkv_fd >= 0 || s.rtsp.running() || s.rtsp_client_fd >= 0 || (s.has_pub && now - s.pub_last <= VIDEO_PUB_IDLE_S)) { log_reject(s, idx, p.ip_be, VIDEO_ADMIT_SLOT_BUSY, now); r.reject_publish("NetStream.Publish.Denied", @@ -1558,7 +1637,7 @@ void VideoChild::pump_viewers(Slot &s, int idx, time_t now) if (!vw.active()) { continue; } - if (vw.kind() == VVK_RTSP || vw.kind() == VVK_RTMP) { + if (vw.kind() == VVK_RTSP || vw.kind() == VVK_RTMP || vw.kind() == VVK_MKV) { // A publisher, not a viewer. Take the socket -- untouched, // since the detect phase only ever peeked -- and splice it. struct sockaddr_in from {}; @@ -1567,10 +1646,12 @@ void VideoChild::pump_viewers(Slot &s, int idx, time_t now) from.sin_port = vw.peer_port_be(); const splice_proto_t proto = vw.kind() == VVK_RTMP ? SPLICE_RTMP : SPLICE_RTSP; + const bool mkv = vw.kind() == VVK_MKV; const int fd = vw.release_fd(); epoll_ctl(epfd_, EPOLL_CTL_DEL, fd, nullptr); s.viewer_events[v] = 0; - handle_rtsp(s, idx, fd, from, now, proto); + if (mkv) handle_mkv(s, idx, fd, from, now); + else handle_rtsp(s, idx, fd, from, now, proto); continue; } if (vw.state() == VV_DETECT && vw.kind() == VVK_WS) { @@ -1629,13 +1710,13 @@ void VideoChild::write_conn_rows(time_t now) */ const bool rtp_pub = s.rtsp.running() && s.rtsp.proto() == SPLICE_RTP; - const bool tcp_pub = !rtp_pub && (s.rtmp || s.rtsp.running() + const bool tcp_pub = !rtp_pub && (s.mkv_fd >= 0 || s.rtmp || s.rtsp.running() || s.rtsp_client_fd >= 0); e.transport = tcp_pub ? CONN_TRANSPORT_TCP : CONN_TRANSPORT_UDP; e.is_user = 0; e.role = CONN_ROLE_VIDEO_PUB; e.stream_idx = uint8_t(i); - e.app_proto = s.rtmp ? CONN_APP_RTMP + e.app_proto = s.mkv_fd >= 0 ? CONN_APP_MATROSKA : s.rtmp ? CONN_APP_RTMP : rtp_pub ? CONN_APP_RTP : tcp_pub ? CONN_APP_RTSP : CONN_APP_MPEGTS; conn_write(db, e); @@ -1658,6 +1739,7 @@ void VideoChild::tick(time_t now) auth_.observe(now); for (int i = 0; i < KEY_MAX_VIDEO_PORTS; i++) { Slot &s = slots_[i]; + if (s.mkv_fd >= 0 && now-s.pub_last > VIDEO_PUB_IDLE_S) close_mkv(s, i, "publisher idle"); if (s.rtsp.running() && s.rtsp.reap()) { close_rtsp(s, i, "backend exited"); } @@ -1820,6 +1902,14 @@ void VideoChild::run(void) if (matched) { continue; } + for (int i=0; i +#include +#include + +// The APCG live profile has an unknown-size Segment and finite, independent +// FFV1 Clusters. Preserve bytes, including BlockAdditional telemetry. Bounds +// are independent of advertised EBML sizes and HTTP chunk boundaries. +class VideoScanner : public TSScanner { +public: + bool matroska = false; + bool failed = false; + std::string prefix; + std::function cluster; + void reset() { *this = VideoScanner(); } + bool join_offset(uint64_t &out) const { + if (!matroska) return TSScanner::join_offset(out); + out = anchor; + return have_anchor && !failed; + } + void feed(const uint8_t *buf, size_t n, uint64_t base) { + if (!matroska) { TSScanner::feed(buf, n, base); return; } + if (failed) return; + if (pending.empty()) offset = base; + pending.insert(pending.end(), buf, buf+n); + while (!pending.empty()) { + uint64_t id, size; size_t a, b; + int r = vint(pending.data(), pending.size(), id, a, true); + if (r < 0) { failed = true; return; } + if (!r) return; + r = vint(pending.data()+a, pending.size()-a, size, b, false); + if (r < 0) { failed = true; return; } + if (!r) return; + const bool unknown = size == ((uint64_t(1) << (7*b))-1); + if (stage == 1) { + if (id != 0x18538067 || !unknown) { failed = true; return; } + size = 0; // descend into the live Segment + } else if (unknown || size > 4*1024*1024 || (stage < 4 && size > 65536)) { + failed = true; return; + } + const size_t total = a+b+size; + if (pending.size() < total) return; + const uint64_t expected[] = {0x1a45dfa3, 0x18538067, 0x1549a966, 0x1654ae6b}; + if (stage < 4) { + if (id != expected[stage] || prefix.size()+total > 131072) { failed = true; return; } + prefix.append(reinterpret_cast(pending.data()), total); + stage++; + } else { + if (id != 0x1f43b675) { failed = true; return; } + anchor = offset; + have_anchor = true; + if (cluster) cluster(pending.data(), total); + } + pending.erase(pending.begin(), pending.begin()+total); + offset += total; + } + } +private: + unsigned stage = 0; + uint64_t offset = 0, anchor = 0; + bool have_anchor = false; + std::vector pending; + static int vint(const uint8_t *p, size_t n, uint64_t &value, size_t &len, bool id) { + if (!n) return 0; + if (!p[0]) return -1; + len = 1; uint8_t mask = 0x80; + while (!(p[0]&mask)) { len++; mask >>= 1; } + if (len > (id ? 4U : 8U)) return -1; + if (n < len) return 0; + value = id ? p[0] : (p[0]&~mask); + for (size_t i=1; i &emit) { + while (n) { + if (remaining) { + const size_t take = n < remaining ? n : remaining; + emit(p, take); p += take; n -= take; remaining -= take; + if (!remaining) ending = 2; + } else if (ending) { + if (*p++ != (ending == 2 ? '\r' : '\n')) return false; + n--; ending--; + } else { + char c = *p++; n--; + if (c == '\n') { + if (line.size() < 2 || line.back() != '\r') return false; + size_t length = 0; + for (size_t i=0; i+1= '0' && h <= '9' ? h-'0' : + h >= 'a' && h <= 'f' ? h-'a'+10 : h >= 'A' && h <= 'F' ? h-'A'+10 : 16; + if (digit > 15 || length > 4*1024*1024/16) return false; + length = length*16+digit; + } + if (!length || length > 4*1024*1024) return false; // zero chunk ends publication + remaining = length; line.clear(); + } else { + line += c; + if (line.size() > 10) return false; + } + } + } + return true; + } +private: + std::string line; + size_t remaining = 0; + unsigned ending = 0; +}; diff --git a/videorec.cpp b/videorec.cpp index b1be74a..6c2b95e 100644 --- a/videorec.cpp +++ b/videorec.cpp @@ -123,7 +123,7 @@ bool VideoWriter::open_segment(time_t now) session_unique_basename(base_dir_.c_str(), port2_, datedir, base, sizeof(base)); char path[1200]; - snprintf(path, sizeof(path), "%s/%s.v%d.ts", dir, base, slot_ + 1); + snprintf(path, sizeof(path), "%s/%s.v%d.%s", dir, base, slot_ + 1, matroska_ ? "mkv" : "ts"); const int fd = open(path, O_WRONLY | O_CREAT | O_EXCL, 0600); if (fd >= 0) { fd_ = fd; diff --git a/videorec.h b/videorec.h index 046b862..582f8ea 100644 --- a/videorec.h +++ b/videorec.h @@ -42,6 +42,7 @@ class VideoWriter { void configure(uint32_t port2, int slot, bool use_tz, float tz_offset, const char *base_dir, uint32_t quota_mb); + void set_matroska(bool value) { matroska_ = value; } bool is_open(void) const { return fd_ >= 0; } // Append. Returns false if recording has stopped (disk full); the @@ -78,6 +79,7 @@ class VideoWriter { std::string base_dir_ = "logs"; uint32_t quota_mb_ = 0; + bool matroska_ = false; int fd_ = -1; std::string path_; std::string date_dir_; diff --git a/videoview.cpp b/videoview.cpp index 49aa196..054fe9c 100644 --- a/videoview.cpp +++ b/videoview.cpp @@ -40,6 +40,7 @@ void VideoViewer::start(int fd, int port2, uint32_t peer_ip_be, last_progress_ = now; behind_since_ = 0; out_.clear(); + prefix_.clear(); prefix_sent_ = 0; out_sent_ = 0; read_pos_ = 0; streaming_ = false; @@ -84,7 +85,7 @@ void VideoViewer::fail(int code, const char *reason, const char *text) } bool VideoViewer::begin_stream(const struct KeyEntry &ke, int slot, - const VideoRing &ring, const TSScanner &scanner, + const VideoRing &ring, const VideoScanner &scanner, bool http, time_t now) { uint64_t anchor = 0; @@ -99,8 +100,10 @@ bool VideoViewer::begin_stream(const struct KeyEntry &ke, int slot, return true; } if (!ring.resident(anchor)) { - anchor = ring.oldest(); + fail(503, "stream not ready", "Join point is no longer buffered."); + return true; } + prefix_ = scanner.prefix; prefix_sent_ = 0; read_pos_ = anchor; streaming_ = true; last_progress_ = now; @@ -114,6 +117,10 @@ bool VideoViewer::begin_stream(const struct KeyEntry &ke, int slot, "Cache-Control: no-store\r\n" "Connection: close\r\n" "\r\n"; + if (scanner.matroska) { + const auto pos = out_.find("video/mp2t"); + out_.replace(pos, strlen("video/mp2t"), "video/x-matroska"); + } out_sent_ = 0; state_ = VV_RESPONDING; } else { @@ -127,7 +134,7 @@ bool VideoViewer::begin_stream(const struct KeyEntry &ke, int slot, } bool VideoViewer::on_readable(const struct KeyEntry &ke, int slot, - const VideoRing &ring, const TSScanner &scanner, + const VideoRing &ring, const VideoScanner &scanner, time_t now) { if (state_ != VV_DETECT) { @@ -254,6 +261,11 @@ bool VideoViewer::on_readable(const struct KeyEntry &ke, int slot, return true; // wait for the rest, still unconsumed } + if (peek.method() == "PUT") { + kind_ = VVK_MKV; + return true; // parent consumes only the header after authenticating + } + // A WebSocket upgrade: hand the untouched socket to WebSocket. if (!peek.header("Upgrade").empty() && peek.header("Upgrade").find("ebsocket") != std::string::npos) { @@ -281,7 +293,7 @@ bool VideoViewer::on_readable(const struct KeyEntry &ke, int slot, // Path selects the slot: /v1.ts .. /v3.ts. The connection already // arrived on this slot's port, so the path only has to agree. char want[32]; - snprintf(want, sizeof(want), "/v%d.ts", slot + 1); + snprintf(want, sizeof(want), "/v%d.%s", slot + 1, scanner.matroska ? "mkv" : "ts"); if (req_.path() != want && req_.path() != "/" && req_.path() != "/stream.ts") { fail(404, "not found", "Try /v1.ts on this port."); return true; @@ -291,7 +303,7 @@ bool VideoViewer::on_readable(const struct KeyEntry &ke, int slot, if (pw.empty()) { pw = http_basic_password(req_.header("Authorization")); } - if (!video_viewer_authorised(ke, pw)) { + if (!ws_authorise(ke, slot, req_)) { fail(401, "unauthorized", "A viewer password is required."); return true; } @@ -300,7 +312,7 @@ bool VideoViewer::on_readable(const struct KeyEntry &ke, int slot, bool VideoViewer::detect_timeout(const struct KeyEntry &ke, int slot, const VideoRing &ring, - const TSScanner &scanner, time_t now) + const VideoScanner &scanner, time_t now) { if (state_ != VV_DETECT) { return true; @@ -381,6 +393,20 @@ bool VideoViewer::on_writable(const VideoRing &ring, time_t now) return false; } + // Send codec setup before any ring bytes, through the same framing as + // the viewer. A slow joiner is subject to the same bounded buffering. + while (prefix_sent_ < prefix_.size()) { + const size_t n = std::min(size_t(16384), prefix_.size()-prefix_sent_); + const uint8_t *p = reinterpret_cast(prefix_.data()+prefix_sent_); + ssize_t w = ws_ ? ws_->send(p, n) : ::send(fd_, p, n, MSG_NOSIGNAL); + if (w < 0 && errno == EINTR) continue; + if ((w < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) || (w == 0 && ws_)) { + blocked_ = true; + return now-last_progress_ <= VIDEO_VIEWER_STUCK_S; + } + if (w <= 0) return false; + prefix_sent_ += size_t(w); last_progress_ = now; + } size_t budget = VIDEO_VIEWER_WRITE_CHUNK; while (budget > 0 && read_pos_ < ring.write_pos()) { uint8_t chunk[16384]; @@ -464,7 +490,7 @@ bool VideoViewer::ws_authorise(const struct KeyEntry &ke, int slot, } bool VideoViewer::begin_ws(const struct KeyEntry &ke, int slot, - const VideoRing &ring, const TSScanner &scanner, + const VideoRing &ring, const VideoScanner &scanner, time_t now) { if (ws_ == nullptr) { @@ -497,8 +523,10 @@ bool VideoViewer::begin_ws(const struct KeyEntry &ke, int slot, return true; // wait for a decodable start point } if (!ring.resident(anchor)) { - anchor = ring.oldest(); + fail(503, "stream not ready", "Join point is no longer buffered."); + return true; } + prefix_ = scanner.prefix; prefix_sent_ = 0; read_pos_ = anchor; streaming_ = true; state_ = VV_STREAMING; diff --git a/videoview.h b/videoview.h index 81133d6..efadb24 100644 --- a/videoview.h +++ b/videoview.h @@ -24,7 +24,7 @@ #include "websocket.h" #include "keydb.h" #include "videostream.h" -#include "videots.h" +#include "videomkv.h" // Per-slot viewer cap. The limit that bites first is egress bandwidth, // not CPU: 32 viewers of an 8 Mbit/s stream is 256 Mbit/s. @@ -55,6 +55,7 @@ enum viewer_kind { VVK_RAW, // raw TCP, no framing VVK_WS, // WebSocket (or WSS), binary frames VVK_RTSP, // RTSP: handed to the ingest splice, not served + VVK_MKV, // HTTP chunked Matroska publisher VVK_RTMP, // RTMP publish: same, a different backend }; @@ -83,7 +84,7 @@ class VideoViewer { // Readable: consume the request. Returns false if the viewer should // be dropped. bool on_readable(const struct KeyEntry &ke, int slot, - const VideoRing &ring, const TSScanner &scanner, + const VideoRing &ring, const VideoScanner &scanner, time_t now); // Writable (or just a poll tick): push bytes. Returns false when the @@ -92,7 +93,7 @@ class VideoViewer { // Called when the detect deadline passes with nothing received. bool detect_timeout(const struct KeyEntry &ke, int slot, - const VideoRing &ring, const TSScanner &scanner, + const VideoRing &ring, const VideoScanner &scanner, time_t now); /* @@ -107,7 +108,7 @@ class VideoViewer { */ // Drive a WebSocket viewer's handshake and start of stream. bool begin_ws_pump(const struct KeyEntry &ke, int slot, - const VideoRing &ring, const TSScanner &scanner, + const VideoRing &ring, const VideoScanner &scanner, time_t now) { return begin_ws(ke, slot, ring, scanner, now); @@ -143,6 +144,8 @@ class VideoViewer { time_t behind_since_ = 0; HttpRequest req_; + std::string prefix_; // codec header, also sent inside WebSocket framing + size_t prefix_sent_ = 0; std::string out_; // pending response bytes size_t out_sent_ = 0; @@ -161,7 +164,7 @@ class VideoViewer { const char *drop_reason_ = ""; bool begin_stream(const struct KeyEntry &ke, int slot, - const VideoRing &ring, const TSScanner &scanner, + const VideoRing &ring, const VideoScanner &scanner, bool http, time_t now); void fail(int code, const char *reason, const char *text); bool flush(time_t now); @@ -170,7 +173,7 @@ class VideoViewer { bool ws_authorise(const struct KeyEntry &ke, int slot, const HttpRequest &req); bool begin_ws(const struct KeyEntry &ke, int slot, const VideoRing &ring, - const TSScanner &scanner, time_t now); + const VideoScanner &scanner, time_t now); }; /* diff --git a/webadmin/logs.py b/webadmin/logs.py index 7944caf..60f62f4 100644 --- a/webadmin/logs.py +++ b/webadmin/logs.py @@ -44,7 +44,7 @@ # and the legacy sessionN names so old logs stay browsable. SESSION_RE = re.compile( r'^(session\d+|\d{4}_\d{2}_\d{2}_\d{2}:\d{2}:\d{2}(-\d+)?)' - r'\.(tlog|bin|v[1-5]\.ts)$') + r'\.(tlog|bin|v[1-5]\.(?:ts|mkv))$') # Natural-sort key: treat embedded digit runs as numbers so that # session10.tlog sorts AFTER session2.tlog (not between session1 and diff --git a/webadmin/routes_video.py b/webadmin/routes_video.py index cb50c46..88d905d 100644 --- a/webadmin/routes_video.py +++ b/webadmin/routes_video.py @@ -14,7 +14,8 @@ from .auth import current_owner, is_admin, require_login from .db import tdb_readonly -from . import videotoken +from . import videotoken, connections +import conntdb_lib bp = Blueprint('video', __name__, url_prefix='/video') @@ -89,10 +90,14 @@ def index(): if not ke.fetch(db): abort(404) + raw_slots = {c.stream_idx for c in connections.list_for_port2(port2) + if c.role == conntdb_lib.CONN_ROLE_VIDEO_PUB and + c.app_proto == conntdb_lib.CONN_APP_MATROSKA} slots = [] for slot, port in ke.active_video_ports(): slots.append({ 'slot': slot, + 'raw_thermal': slot in raw_slots, 'number': slot + 1, 'port': port, 'opts': ke.slot_opt_names(slot), diff --git a/webadmin/templates/video.html b/webadmin/templates/video.html index c665141..c8b157b 100644 --- a/webadmin/templates/video.html +++ b/webadmin/templates/video.html @@ -18,6 +18,14 @@

Video — {{ entry.port1 }}/{{ port2 }} {{ entry.name }}

Slot {{ s.number }} — port {{ s.port }}

+ {% if s.raw_thermal %} +

Lossless raw thermal (FFV1, 16-bit). This browser player cannot decode + the stream. Use camera view rawthermal in MAVProxy when + connected through this proxy, or the desktop thermal viewer:

+
python3 scripts/view_raw_thermal.py http://{{ host }}:{{ s.port }}/v{{ s.number }}.mkv{% if s.needs_password %}?pw=YOUR_VIEWER_PASSWORD{% endif %}
+

The desktop viewer defaults to greyscale, shows pixel temperatures, + and saves native samples with their capture metadata.

+ {% else %} {# One element for the life of the page. Picture-in-Picture is bound to the element, so replacing it on reconnect tears PIP down -- see the teardown note in the script. #} @@ -45,6 +53,7 @@

Slot {{ s.number }} — port {{ s.port }}

vlc --network-caching=200 --live-caching=200 \
     http://{{ host }}:{{ s.port }}/v{{ s.number }}.ts{% if s.needs_password %}?pw=YOUR_VIEWER_PASSWORD{% endif %}
+ {% endif %}
{% endfor %} @@ -52,7 +61,7 @@

Slot {{ s.number }} — port {{ s.port }}