"""Bounded per-camera encoded fanout. Readers never consume each other's data.""" import json import re import struct import threading import time from collections import deque MAX_PACKET = 4 * 1024 * 1024 MAX_BYTES = 8 * 1024 * 1024 MAX_ENTRIES = 64 MAX_WIRE = MAX_PACKET + 1028 MAX_PARAMETERS = 65536 START_CODE = re.compile(b"\x00\x00(?:\x00)?\x01") class Parameters: """Retain only codec setup, so a late viewer can decode the next keyframe. Parameter sets describe this preview generation, never a previous session. Media history and camera commands are not retained or replayed. """ def __init__(self): self.streams = {} def packet(self, header, data): codec = header["codec"] required = (7, 8) if codec == 0 else (32, 33, 34) keyframes = (5,) if codec == 0 else (16, 17, 18, 19, 20, 21) saved = self.streams.setdefault((header["stream_index"], codec), {}) markers = START_CODE.finditer(data) marker = next(markers, None) present, keyframe = set(), False while marker is not None: following = next(markers, None) if marker.end() >= len(data): break kind = data[marker.end()] & 31 if codec == 0 else (data[marker.end()] >> 1) & 63 keyframe |= kind in keyframes if kind in required: end = following.start() if following is not None else len(data) value = data[marker.start() : end] if len(value) > MAX_PARAMETERS: saved.clear() return data saved[kind] = value present.add(kind) marker = following if keyframe and all(kind in saved for kind in required): prefix = b"".join(saved[kind] for kind in required if kind not in present) if len(prefix) + len(data) <= MAX_PACKET: return prefix + data return data def encode(value): if value is None: return b"" header, data = value raw = json.dumps(header, allow_nan=False, separators=(",", ":")).encode() if len(raw) > 1024 or not 0 < len(data) <= MAX_PACKET: raise ValueError("Invalid video envelope") return struct.pack("!I", len(raw)) + raw + data def decode(raw): if not raw: return None if len(raw) < 5 or len(raw) > MAX_WIRE: raise ValueError("Invalid video envelope") size = struct.unpack("!I", raw[:4])[0] if not 0 < size <= 1024 or len(raw) <= 4 + size: raise ValueError("Invalid video envelope") header, data = json.loads(raw[4 : 4 + size]), raw[4 + size :] if ( not isinstance(header, dict) or any( type(header.get(key)) is not int or header[key] < 0 for key in ("cursor", "generation", "stream_index", "codec") ) or header["stream_index"] not in (0, 1) or header["codec"] not in (0, 1) or type(header.get("gap")) is not bool or not 0 < len(data) <= MAX_PACKET ): raise ValueError("Invalid video metadata") return header, data class Feed: def __init__(self, camera): self.camera = camera self.lock = threading.Condition() self.source_lock = threading.Lock() self.queue = deque() self.bytes = self.cursor = self.generation = 0 self.source_generation = None self.parameters = Parameters() self.closed = threading.Event() self.thread = threading.Thread(target=self.pump, daemon=True) def clear(self): with self.lock: self.queue.clear() self.bytes = 0 self.parameters = Parameters() self.generation += 1 self.lock.notify_all() def append(self, value): source, data = value if not 0 < len(data) <= MAX_PACKET: self.clear() return with self.lock: if self.source_generation != source["generation"]: self.queue.clear() self.bytes = 0 self.source_generation = source["generation"] self.generation += 1 self.parameters = Parameters() data = self.parameters.packet(source, data) while self.queue and ( self.bytes + len(data) > MAX_BYTES or len(self.queue) >= MAX_ENTRIES ): _, removed, _ = self.queue.popleft() self.bytes -= len(removed) self.cursor += 1 header = dict(source, cursor=self.cursor, generation=self.generation) self.queue.append((header, data, time.monotonic())) self.bytes += len(data) self.lock.notify_all() def read(self, cursor, wait=0.1): if type(cursor) is not int or not 0 <= cursor < 2**64: raise ValueError("Invalid video cursor") deadline = time.monotonic() + min(max(wait, 0), 0.2) with self.lock: while True: for header, data, observed in self.queue: if header["cursor"] > cursor and time.monotonic() - observed < 2: return dict( header, gap=bool(cursor and header["cursor"] != cursor + 1) ), data remaining = deadline - time.monotonic() if remaining <= 0 or self.closed.is_set(): return None self.lock.wait(remaining) def reader(self): with self.lock: cursor = self.cursor def read(): nonlocal cursor value = self.read(cursor, wait=0) if value is not None: cursor = value[0]["cursor"] return value return read def change(self, action): # No packet read from the previous preview may be appended after a # START/STOP boundary and mistaken for evidence of the new preview. with self.source_lock: self.clear() try: return action() finally: self.clear() def pump(self): while not self.closed.is_set(): try: with self.source_lock: value = self.camera.video() if value is not None: self.append(value) continue except (OSError, RuntimeError, ValueError): self.clear() self.closed.wait(0.01) def close(self): self.closed.set() self.clear() if self.thread.is_alive(): self.thread.join(timeout=1)