Files
DCCONSTRUCTIONS a3c15e11e9 Add packaged Insta360 X4 integration and recover paired Node channels
Discover independent camera instances and prepare their versioned runtime
from Node or remote Core. Add isolated SDK workers, camera controls, raw
dual-fisheye WebRTC preview, and shared action/region loading states.

Recover existing Node bindings over known Tailscale addresses after a Core
LAN address change. Preserve identities and trust, pin both peers, migrate
endpoints with revision checks, and require real heartbeats for online status.
Fix the Python client certificate profile for Go X509 verification.

Pin Design Guideline 8c53f73 and retain installer/build/acceptance history.
Node 0.8.19 is installed; X4 0.1.3-3 is bundled but hardware activation is pending.

Validation: qualified DG/Node builds and Go race tests; 31 fleet tests;
Python-to-Go certificate interoperability and live tailnet recovery with five
fresh heartbeats; prior 38 X4 tests and bounded remote WebRTC acceptance.
Clean-OS, replug/power autonomy, local X4 video and long-run stability remain open.
2026-09-10 09:21:24 +03:00

414 lines
18 KiB
Python

"""Real codecs/WebRTC over loopback; synthetic frames, no CameraSDK or device."""
import asyncio
import os
import re
import tempfile
import threading
import unittest
from fractions import Fraction
from pathlib import Path
from unittest.mock import patch
from aiortc import RTCConfiguration, RTCPeerConnection, RTCSessionDescription
from check_runtime import ROOT, VideoCamera, synthetic_video
from runtime import broker, frames, http, media
from runtime.verification import Control
def packet(number, data=b"synthetic", generation=1):
return {
"stream_index": 0,
"codec": 0,
"generation": generation,
"sequence": number,
"timestamp": number,
}, data
class FrameTests(unittest.TestCase):
def test_repeated_preview_commands_preserve_generation_and_readers(self):
feed = frames.Feed(None)
camera = VideoCamera([], preview=1)
control = Control(camera, feed)
feed.append(packet(1))
before = feed.read(0, 0)
self.assertEqual(control.call("preview.start", {})["state"], "complete")
self.assertEqual(feed.read(0, 0), before)
control.call("preview.stop", {})
self.assertIsNone(feed.read(0, 0))
generation = feed.generation
control.call("preview.stop", {})
self.assertEqual(feed.generation, generation)
control.call("preview.start", {})
self.assertGreater(feed.generation, generation)
def test_uncertain_redundant_start_does_not_release_existing_decoder(self):
value, events = broker.Broker(), []
identifier = "instax4_" + "a" * 32
session = {"device_id": identifier, "session_id": "synthetic"}
class Engine:
def capture(self, key, path, start):
events.append(start)
value.media = Engine()
def request(path, route, command=None, **kwargs):
if route == "/snapshot":
return {
"id": identifier,
"session_id": "synthetic",
"online": True,
"status": {"preview": 1},
}
return {"state": "unknown"}
with patch.object(broker, "request", request):
result = value.capture_operation(
identifier, {"session": session, "action_id": "preview.start"}
)
self.assertEqual(result["state"], "unknown")
self.assertEqual(events, [True])
def test_broker_primes_before_sdk_start_and_releases_after_sdk_stop(self):
value = broker.Broker()
events = []
identifier = "instax4_" + "a" * 32
session = {"device_id": identifier, "session_id": "synthetic"}
class Engine:
def capture(self, key, path, start):
events.append("prime" if start else "release")
value.media = Engine()
def request(path, route, command=None, **kwargs):
if route == "/snapshot":
return {"id": identifier, "session_id": "synthetic", "online": True}
events.append(command["action_id"])
return {"state": "complete"}
with patch.object(broker, "request", request):
for action in ("preview.start", "preview.stop"):
value.capture_operation(identifier, {"session": session, "action_id": action})
self.assertEqual(events, ["prime", "preview.start", "preview.stop", "release"])
def test_late_decoder_gets_parameters_after_initial_packets_leave_queue(self):
import av
encoder = av.CodecContext.create("libx264", "w")
encoder.width, encoder.height, encoder.pix_fmt = 32, 16, "yuv420p"
encoder.time_base = Fraction(1, 30)
encoder.options = {
"preset": "ultrafast",
"tune": "zerolatency",
"x264-params": "keyint=6:min-keyint=6:scenecut=0",
}
feed = frames.Feed(None)
for number in range(96):
frame = av.VideoFrame(32, 16, "yuv420p")
frame.pts = number
for plane in frame.planes:
plane.update(bytes([64]) * plane.buffer_size)
for encoded in encoder.encode(frame):
data = bytes(encoded)
if number:
# Model a camera which sends SPS/PPS only at preview start.
units = re.split(b"\x00\x00(?:\x00)?\x01", data)
data = b"".join(
b"\x00\x00\x00\x01" + unit
for unit in units
if unit and unit[0] & 31 not in (7, 8)
)
feed.append(packet(number, data))
self.assertEqual(len(feed.queue), frames.MAX_ENTRIES)
self.assertGreater(feed.queue[0][0]["sequence"], 0)
decoder, decoded, cursor = av.CodecContext.create("h264", "r"), [], 0
while (value := feed.read(cursor, 0)) is not None:
header, data = value
cursor = header["cursor"]
try:
for encoded in decoder.parse(data):
decoded.extend(decoder.decode(encoded))
except av.error.FFmpegError:
pass
self.assertGreaterEqual(len(decoded), 2)
self.assertEqual((decoded[-1].width, decoded[-1].height), (32, 16))
def test_parameter_cache_cannot_cross_stream_device_or_preview_generation(self):
config = b"\x00\x00\x01\x67sps\x00\x00\x01\x68pps"
keyframe = b"\x00\x00\x01\x65frame"
first, second = frames.Feed(None), frames.Feed(None)
first.append(packet(1, config))
first.append(packet(2, keyframe))
self.assertEqual(first.queue[-1][1], config + keyframe)
second.append(packet(1, keyframe))
self.assertEqual(second.queue[-1][1], keyframe)
header, data = packet(3, keyframe)
header["stream_index"] = 1
first.append((header, data))
self.assertEqual(first.queue[-1][1], keyframe)
first.append(packet(4, keyframe, generation=2))
self.assertEqual(first.queue[-1][1], keyframe)
first.append(packet(5, config, generation=2))
first.clear()
first.append(packet(6, keyframe, generation=2))
self.assertEqual(first.queue[-1][1], keyframe)
def test_parameter_cache_is_bounded_and_hevc_parameters_are_distinct(self):
feed = frames.Feed(None)
huge = b"\x00\x00\x01\x67" + b"x" * frames.MAX_PARAMETERS
feed.append(packet(1, huge))
self.assertFalse(any(feed.parameters.streams.values()))
config = b"".join(b"\x00\x00\x01" + bytes([kind << 1, 1, 2]) for kind in (32, 33, 34))
keyframe = b"\x00\x00\x01\x26\x01frame"
header, _ = packet(2)
header["codec"] = 1
feed.append((header, config))
feed.append((header, keyframe))
self.assertEqual(feed.queue[-1][1], config + keyframe)
def test_independent_readers_and_devices_preserve_packet_identity(self):
first, second = frames.Feed(None), frames.Feed(None)
a, b = first.reader(), first.reader()
first.append(packet(1, b"first"))
second.append(packet(1, b"second"))
self.assertEqual(a()[1], b"first")
self.assertEqual(b()[1], b"first")
self.assertEqual(second.read(0, 0)[1], b"second")
self.assertIsNone(a())
def test_overrun_and_preview_boundary_are_visible_to_decoder(self):
feed = frames.Feed(None)
for n in range(100):
feed.append(packet(n))
value = feed.read(1, 0)
self.assertTrue(value[0]["gap"])
self.assertLessEqual(len(feed.queue), frames.MAX_ENTRIES)
old = value[0]["generation"]
feed.change(lambda: None)
feed.append(packet(101))
self.assertGreater(feed.read(0, 0)[0]["generation"], old)
def test_wire_bounds_and_types_reject_corrupted_or_unknown_packets(self):
feed = frames.Feed(None)
feed.append(packet(1))
value = feed.read(0, 0)
self.assertEqual(frames.decode(frames.encode(value)), value)
for data in (b"x", b"\xff" * 8, b"\0\0\0\1x", frames.encode(value)[:5]):
with self.assertRaises((ValueError, UnicodeError)):
frames.decode(data)
with self.assertRaises(ValueError):
feed.read(True, 0)
def test_public_ice_candidates_are_rejected_before_network_work(self):
with self.assertRaises(ValueError):
media.admit_sdp("v=0\r\na=candidate:1 1 udp 1 203.0.113.6 1234 typ host\r\n")
media.admit_sdp("v=0\r\na=candidate:1 1 udp 1 127.0.0.1 1234 typ host\r\n")
class MediaTests(unittest.IsolatedAsyncioTestCase):
async def test_primed_decoder_survives_late_viewer_and_peer_close(self):
with tempfile.TemporaryDirectory(dir=ROOT / "build") as temporary:
descriptor = os.open(temporary, os.O_RDONLY)
path = Path(f"/proc/self/fd/{descriptor}/video.sock")
feed, stop = frames.Feed(None), threading.Event()
packets = synthetic_video(220) # One initial keyframe; >64-packet queue.
def publish():
for number, data in enumerate(packets):
if stop.is_set():
break
feed.append(packet(number, data))
stop.wait(1 / 30)
def dispatch(method, route, value, headers):
if route == "/snapshot":
return {"session_id": "synthetic", "online": True, "status": {"preview": 1}}
if method != "POST" or route != "/video" or value["session_id"] != "synthetic":
raise ValueError("Unexpected source command")
return http.Binary(frames.encode(feed.read(value["cursor"])))
server = http.Server(path, dispatch, {os.geteuid()})
serving = threading.Thread(target=server.serve_forever, daemon=True)
producer = threading.Thread(target=publish, daemon=True)
serving.start()
peers, key = media.Peers(), ("synthetic-camera", "synthetic")
receiver = RTCPeerConnection(RTCConfiguration(iceServers=[]))
received = asyncio.get_running_loop().create_future()
tasks = []
@receiver.on("track")
def track(value):
async def observe():
first, second = await value.recv(), await value.recv()
if not received.done():
received.set_result((first, second))
tasks.append(asyncio.create_task(observe()))
try:
await peers.prime(key, path)
source = peers.sources[key]
producer.start()
await asyncio.sleep(2.4)
self.assertGreater(feed.queue[0][0]["sequence"], 0)
self.assertGreater(source.current()[0], 2)
with patch.object(media, "_host_addresses", return_value=["127.0.0.1"]):
receiver.addTransceiver("video", direction="recvonly")
receiver.createDataChannel("sensor")
await receiver.setLocalDescription(await receiver.createOffer())
answer = await peers.offer(
key, path, {"sdp": receiver.localDescription.sdp, "layer": "preview"}
)
self.assertNotIn("VP8/", answer["sdp"])
await receiver.setRemoteDescription(
RTCSessionDescription(sdp=answer["sdp"], type="answer")
)
first, second = await asyncio.wait_for(received, 4)
self.assertEqual((first.width, first.height), (32, 16))
self.assertGreater(second.pts, first.pts)
# A second UI can send START from its stale idle snapshot.
# No new camera keyframe is produced for this idempotent call.
control = Control(VideoCamera([], preview=1), feed)
control.call("preview.start", {})
verified = await peers.verify(key, timeout=1)
self.assertEqual(verified["state"], "complete")
self.assertGreaterEqual(verified["result"]["streams"][0]["frames"], 2)
await peers.close(key, answer["peer_id"])
self.assertIs(peers.sources[key], source)
self.assertFalse(source.closed.is_set())
finally:
await receiver.close()
await peers.release(key)
for task in tasks:
task.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
stop.set()
if producer.is_alive():
producer.join(2)
feed.close()
server.shutdown()
server.server_close()
serving.join(2)
os.close(descriptor)
self.assertFalse(peers.items)
self.assertFalse(peers.sources)
self.assertFalse(peers.primed)
async def test_verification_never_accepts_retained_frames_or_another_camera(self):
peers = media.Peers()
key = ("synthetic", "one")
class Retained:
closed = threading.Event()
def current(self):
return 50, object(), media.time.monotonic()
retained = Retained()
peers.sources[key] = retained
self.assertEqual((await peers.verify(key, timeout=0.05))["state"], "error")
self.assertEqual((await peers.verify(("another", "one")))["state"], "error")
retained.closed.set()
self.assertEqual((await peers.verify(key))["state"], "error")
async def test_synthetic_frames_cross_private_ipc_and_real_webrtc_with_cleanup(self):
self.loop_slow_callback_duration = 1
with tempfile.TemporaryDirectory(dir=ROOT / "build") as temporary:
directory = Path(temporary)
descriptor = os.open(directory, os.O_RDONLY)
path = Path(f"/proc/self/fd/{descriptor}/video.sock")
feed = frames.Feed(None)
stop = threading.Event()
packets = synthetic_video()
def publish():
number = 0
while not stop.is_set():
feed.append(
packet(
number,
packets[number % len(packets)],
generation=number // len(packets) + 1,
)
)
number += 1
stop.wait(1 / 30)
def dispatch(method, route, value, headers):
if method != "POST" or route != "/video" or value["session_id"] != "synthetic":
raise ValueError("Wrong camera/session")
return http.Binary(frames.encode(feed.read(value["cursor"])))
server = http.Server(path, dispatch, {os.geteuid()})
server_thread = threading.Thread(target=server.serve_forever, daemon=True)
producer = threading.Thread(target=publish, daemon=True)
server_thread.start()
producer.start()
peers = media.Peers()
receiver = RTCPeerConnection(RTCConfiguration(iceServers=[]))
received = asyncio.get_running_loop().create_future()
tasks = []
identifier = None
key = ("synthetic-camera", "synthetic")
@receiver.on("track")
def track(value):
async def observe():
first = await value.recv()
second = await value.recv()
if not received.done():
received.set_result((first, second))
tasks.append(asyncio.create_task(observe()))
try:
# Both peers use only 127.0.0.1. No STUN, LAN candidates or
# camera packet is admitted to this synthetic qualification.
with patch.object(media, "_host_addresses", return_value=["127.0.0.1"]):
receiver.addTransceiver("video", direction="recvonly")
channel = receiver.createDataChannel("sensor")
@channel.on("open")
def opened():
channel.send("keepalive")
await receiver.setLocalDescription(await receiver.createOffer())
answer = await peers.offer(
key, path, {"sdp": receiver.localDescription.sdp, "layer": "preview"}
)
identifier = answer["peer_id"]
await receiver.setRemoteDescription(
RTCSessionDescription(sdp=answer["sdp"], type="answer")
)
first, second = await asyncio.wait_for(received, 12)
self.assertEqual((first.width, first.height), (32, 16))
self.assertGreater(second.pts, first.pts)
with self.assertRaises(ValueError):
await peers.close(("another-camera", "synthetic"), identifier)
self.assertIn(identifier, peers.items)
finally:
await receiver.close()
if identifier:
await peers.close(key, identifier)
for task in tasks:
task.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
stop.set()
feed.close()
server.shutdown()
server.server_close()
server_thread.join(2)
producer.join(2)
os.close(descriptor)
self.assertFalse(peers.items)
self.assertFalse(peers.sources)
self.assertFalse(server_thread.is_alive())
self.assertFalse(producer.is_alive())
if __name__ == "__main__":
unittest.main(verbosity=2)