Files
Jetson/bridge/bridge.py
T
AngePierreandClaude Opus 5 c167c836e6 Portable Ithaca capture node: one install for every Jetson
Extract the install that was applied by hand to the Jetson Nano into a single
checkout that builds itself for whatever Jetson it lands on.

The two nodes on the rig share no JetPack — tegra210 caps at 4, tegra234 needs
5+ — so nothing binary is portable between them. The source and the procedure
are: install.sh detects the platform (L4T / JetPack / SoC), installs deps,
builds the OrbbecSDK and the preview server from source locally, and generates
a per-user config and systemd unit. The tested bridge.py ships verbatim; all
machine-specific paths live in the generated config, so it stays unmodified —
its one hardcoded path is made home-relative.

The retired GStreamer/RTSP target is dropped: the direct route encodes no video.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-10 15:57:15 +02:00

1432 lines
65 KiB
Python

#!/usr/bin/env python3
# Copyright (c) French Touch Factory. All Rights Reserved.
#
# Ithaca WebSocket bridge — Jetson client.
#
# The Jetson is always a *client* in the Ithaca mesh (never hosts the relay):
# it UDP-discovers the Windows PC currently running the Node.js relay, opens a
# plain WebSocket to it, announces its camera(s) via CameraUpsert, then reacts
# to Record/Stop/RecordSettings events by driving the local
# ob_stream_depth_color_k4arec_pipeline binary. See the protocol notes this
# was built from (message envelope, enum wire values, folder convention) —
# ask the project owner for the architecture report if this file needs
# updating; the wire format itself is not documented anywhere in this repo.
#
# Python 3.6 compatible on purpose (that's what ships on this Jetson's
# Ubuntu 18.04 — no asyncio.run(), no f-strings assumed available elsewhere).
import asyncio
import json
import logging
import os
import re
import shutil
import signal
import socket
import subprocess
import sys
import time
import websockets
# ---- Wire protocol -----------------------------------------------------------
class MessageType:
REGISTER = 0
COMMAND = 1
RESPONSE = 2
EVENT = 3
ERROR = 4
class ActionType:
GET_CLIENT = 0
RECORD = 10
PREVIEW = 11
STOP = 12
PROGRESS = 16
STATUT = 17
CLIENT_DISCONNECTED = 20
CLIENT_CONNECTED = 21
CAMERA_UPSERT = 22
CAMERA_REMOVE = 23
TAKE_NAME = 24
RECORD_SETTINGS = 25
GATHER_TAKES = 26
AUTO_GATHER = 27
DELETE_AFTER_UPLOAD = 28
# CamEnum: 0=Orbbec, 1=Kinect, 2=RealSense, 3=Unknown — Femto Bolt/Mega are Orbbec hardware.
CAM_TYPE_ORBBEC = 0
# SyncModeEnum
SYNC_MODE_WIRE = {"standalone": 0, "master": 1, "subordinate": 2}
SYNC_MODE_NAME = {v: k for k, v in SYNC_MODE_WIRE.items()}
# DepthModeEnum -> our binary's --depth-mode strings
DEPTH_MODE_MAP = {
0: "NFOV_2X2BINNED", # NFOV_Binned_320x288
1: "NFOV_UNBINNED", # NFOV_Unbinned_640x576
2: "WFOV_2X2BINNED", # WFOV_Binned_512x512
3: "WFOV_UNBINNED", # WFOV_Unbinned_1024x1024
}
# ColorModeEnum -> our binary's --color-mode strings (5 = 3072p 4:3, the exe
# doesn't support it either — Windows falls back to 1536p, so do we).
COLOR_MODE_MAP = {
0: "720p",
1: "1080p",
2: "1440p",
3: "1536p",
4: "2160p",
5: "1536p",
}
# AlignModeEnum -> our binary's --align strings. Documentation seule desormais :
# aucun de ces modes n'est plus selectionne ici, le noeud streame toujours en raw
# (voir _stream_options). Les alignements sont calcules par le viewer.
ALIGN_MODE_MAP = {
0: "raw",
1: "d2c",
2: "c2d",
# Diagnostic : la profondeur brute ET la sortie alignee du SDK dans le meme
# frameset, pour superposer les deux nuages. A servi a valider notre transform
# contre celle du SDK ; plus demande depuis.
3: "compare",
}
DISCOVERY_MSG = b"DISCOVER_WS_SERVER"
DISCOVERY_PORT = 41234
WS_PORT = 8765
# USB vendor 0x2bc5 = Orbbec. Product IDs seen in this project's own camera
# detection scripts (universal_camera_recorder.sh).
KNOWN_ORBBEC_PIDS = {
"066b": "Femto Bolt",
"0669": "Femto Mega",
"0660": "Femto Mega",
}
DEFAULT_CONFIG = {
"server_ip": None, # null = auto-discover via UDP broadcast
"working_directory": os.path.expanduser("~/IthacaRecordings"),
"recorder_bin": os.path.expanduser(
"~/OrbbecSDK_v2_jetson/build/linux_arm64/bin/ob_stream_depth_color_k4arec_pipeline"
),
"sync_mode": "standalone", # standalone | master | subordinate
"camera_serial": None, # null = announce every detected camera; a serial (or
# list of serials) restricts us to those
# Pousser brightness/contrast/saturation/sharpness/gain au capteur. Actif :
# l'UI est la source de verite, elle doit refleter ce qui est applique.
# Passer a false pour laisser le capteur sur ses propres defauts.
"apply_image_controls": True,
"log_level": "INFO",
}
# Configuration par camera, persistee a cote de config.json.
#
# C'est NOUS qui possedons les cameras branchees ici, donc c'est a nous de
# memoriser leurs reglages : CameraPersistence, cote Unity, ne sauvegarde que les
# cameras locales et jamais les distantes ("qui appartiennent a d'autres postes").
# Sans store local, tout reglage fait dans l'UI etait perdu a la reconnexion et
# notre re-annonce periodique remettait l'ancienne valeur.
#
# Les valeurs sont stockees telles qu'elles circulent sur le fil (entiers/booleens
# de CameraPayload) : la re-annonce les renvoie verbatim, sans conversion qui
# pourrait deriver.
# La composition servie, relue par le serveur de preview une fois par seconde :
# c'est ce qui permet a une case de fermer un capteur sans relancer le serveur, donc
# sans couper les autres cameras du noeud une dizaine de secondes.
STREAM_MAP_FILE = "stream_map.txt"
CAM_SETTINGS_FILE = "cam_settings.json"
LEGACY_CAM_INDEX_FILE = "cam_index.json"
# Memes champs que CameraPersistence.ApplySaved : l'identite materielle (serial,
# modele) et l'etat disque n'en font pas partie, ce ne sont pas des reglages.
PERSISTED_FIELDS = ("camIndex", "syncMode", "exposureAuto", "whiteBalanceAuto",
"exposure", "whiteBalance", "brightness", "contrast",
"saturation", "sharpness", "gain", "orientationQuarters",
"disabled")
# Quarts de tour autour de l'axe optique. Hors de MANUAL_CONTROLS a dessein : ce
# n'est pas un reglage du capteur, rien n'est passe au recorder. On le memorise
# parce que c'est un fait du montage, qui doit survivre a un redemarrage du noeud.
ORIENTATION_QUARTERS = 4
# Controles image. Tous memorises et renvoyes verbatim dans la re-annonce, pour
# que l'UI ne soit jamais remise a zero par notre propre message. 0 vaut "jamais
# transmis" et l'option n'est alors pas passee au recorder.
# Note : le payload ne porte un booleen "auto" que pour l'exposition et la
# balance des blancs ; pour les cinq autres, un reglage voulu ne se distingue pas
# d'un champ jamais initialise (l'UI clampe un champ vide vers le minimum de sa
# plage, d'ou les contrast=1 / saturation=1 / sharpness=1 / gain=1 observes sur
# une camera vierge). Un booleen par controle dans CameraPayload leverait
# l'ambiguite.
MANUAL_CONTROLS = ("exposure", "whiteBalance", "brightness", "contrast",
"saturation", "sharpness", "gain")
# Ecriture du store au plus une fois par cet intervalle (secondes).
SAVE_MIN_INTERVAL = 10.0
USB_DEVICES = "/sys/bus/usb/devices"
# Preview streaming: one server for the whole node, serving every camera on one
# TCP connection — colour as the sensor's own JPEG and depth as its own 16-bit
# millimetres, both halves of a frameset in the same message. Overridable
# through "stream_server" in config.json.
# Repli home-relatif, comme recorder_bin ci-dessus. En pratique jamais utilise :
# install.sh ecrit "stream_server" dans config.json avec le chemin reel du build.
DEFAULT_STREAM_SERVER = os.path.expanduser(
"~/ithaca_preview/build/ithaca_rgbd_server")
def load_config(path):
cfg = dict(DEFAULT_CONFIG)
if os.path.exists(path):
with open(path) as f:
cfg.update(json.load(f))
return cfg
def _sysfs_read(path):
try:
with open(path) as f:
return f.read().strip()
except (IOError, OSError):
return ""
def detect_cameras():
"""Enumerate connected Orbbec cameras from sysfs.
Deliberately NOT `lsusb -v`: that opens the device and issues control
transfers to dump its descriptors, and this runs every 3s while a recording
may be streaming. Reading sysfs touches no USB endpoint at all.
Returns cameras sorted by USB port path, so the ordering is a property of
the wiring rather than of plug order (`lsusb` is ordered by device number,
which is handed out at plug time and therefore shuffles across reboots).
"""
cams = []
try:
entries = sorted(os.listdir(USB_DEVICES))
except OSError:
return cams
for entry in entries:
dev = os.path.join(USB_DEVICES, entry)
if _sysfs_read(os.path.join(dev, "idVendor")).lower() != "2bc5":
continue
pid = _sysfs_read(os.path.join(dev, "idProduct")).lower()
if pid not in KNOWN_ORBBEC_PIDS:
continue
serial = _sysfs_read(os.path.join(dev, "serial"))
if not serial:
continue
cams.append({
"serial": serial,
"model": KNOWN_ORBBEC_PIDS[pid],
"port": entry,
})
return cams
class CamSettingsStore:
"""serial -> reglages persistes (forme du fil)."""
def __init__(self, path, legacy_index_path=None):
self.path = path
self.cams = {}
# Serials currently on the USB bus, refreshed by the bridge. Only these
# take part in an index swap (see update()).
self.present = set()
self._last_save = 0.0
self._dirty = False
if os.path.exists(path):
try:
with open(path) as f:
loaded = json.load(f)
if isinstance(loaded, dict):
for serial, settings in loaded.items():
if isinstance(settings, dict):
self.cams[serial] = settings
except (ValueError, IOError, OSError):
logging.warning("Lecture de %s impossible — on repart d'un store vide.", path)
elif legacy_index_path and os.path.exists(legacy_index_path):
# Reprise de l'ancien fichier, qui ne contenait que les index : on ne
# veut pas que la renumerotation des cameras deja en service reparte
# de zero a la mise a jour.
try:
with open(legacy_index_path) as f:
for serial, idx in json.load(f).items():
self.cams[serial] = {"camIndex": int(idx)}
logging.info("Reprise de %s : %d caméra(s) migrée(s) vers %s.",
legacy_index_path, len(self.cams), os.path.basename(path))
self._save()
except (ValueError, IOError, OSError, AttributeError):
logging.warning("Reprise de %s impossible — store vide.", legacy_index_path)
def settings(self, serial):
"""Reglages connus pour cette camera (dict vide si inconnue)."""
return dict(self.cams.get(serial, {}))
def next_index(self, serial):
"""Index de la camera, attribue a la premiere detection.
Volontairement monotone : debrancher une camera ne doit pas renumeroter
les autres, sinon les MKV changeraient de suffixe d'une prise a l'autre."""
cur = self.cams.setdefault(serial, {})
if cur.get("camIndex") is None:
used = [c["camIndex"] for c in self.cams.values()
if isinstance(c.get("camIndex"), int)]
cur["camIndex"] = (max(used) + 1) if used else 0
self._save()
logging.info("camIndex %d attribué à la nouvelle caméra %s",
cur["camIndex"], serial)
return cur["camIndex"]
def update(self, serial, changes):
"""Applique et persiste des reglages. Renvoie les serials touches.
L'index reste unique : il forme le suffixe du nom de fichier
("<take>_<camIndex>.mkv"), donc deux cameras qui le partageraient
ecriraient dans le meme MKV. Si l'index demande est deja pris, on echange
avec son detenteur — l'intention naturelle quand on renumerote.
On n'echange qu'avec une camera PRESENTE. Echanger avec une absente
deplace un index que personne ne voit, et Unity, dont le modele garde
encore son ancienne valeur, la renvoie aussitot : les deux index
s'echangent alors en boucle. Comme l'index nomme le chemin RTSP publie
("cam<N>_color"), le flux changeait de nom entre deux relances et Unity
demandait un chemin inexistant (404). Une absente cede simplement son
index ; il sera reattribue par next_index a son retour."""
cur = self.cams.setdefault(serial, {})
touched = set()
for field, value in changes.items():
if field not in PERSISTED_FIELDS or cur.get(field) == value:
continue
if field == "camIndex":
previous = cur.get("camIndex")
holder = next((s for s, c in self.cams.items()
if s != serial and c.get("camIndex") == value
and s in self.present), None)
if holder is not None and previous is not None:
self.cams[holder]["camIndex"] = previous
touched.add(holder)
logging.info("camIndex : %s <- %d (échangé avec %s, qui prend %d)",
serial, value, holder, previous)
else:
logging.info("camIndex : %s <- %d", serial, value)
else:
logging.info("%s : %s = %s (était %s)", serial, field, value, cur.get(field))
cur[field] = value
touched.add(serial)
if touched:
self._save()
return touched
def _save(self, force=False):
"""Ecrit le store, au plus une fois par SAVE_MIN_INTERVAL.
Amortissement necessaire : deplacer un curseur dans l'UI produit une rafale
de CameraUpsert, et l'auto-exposition peut en emettre en continu. Sans
cela on reecrivait le fichier plusieurs fois par minute sur une eMMC.
L'etat en memoire est toujours a jour ; seule sa mise sur disque attend.
flush() force l'ecriture des changements en attente."""
self._dirty = True
now = time.time()
if not force and (now - self._last_save) < SAVE_MIN_INTERVAL:
return
try:
tmp = self.path + ".tmp"
with open(tmp, "w") as f:
json.dump(self.cams, f, indent=2, sort_keys=True)
os.rename(tmp, self.path) # atomique : jamais de fichier tronque
self._last_save = now
self._dirty = False
except (IOError, OSError):
logging.exception("Persistance impossible vers %s", self.path)
def flush(self):
if self._dirty:
self._save(force=True)
def free_space(path):
"""(libre, capacite, % libre) en Go pour le volume portant `path`.
La capacite est envoyee telle quelle plutot que laissee a deduire de
libre/pourcentage : les deux sont des entiers arrondis, et Unity affiche
"<libre> / <total> GB" (CustomPCFoldout.SetStorage) — une capacite reconstruite
par division derivait de plusieurs Go."""
os.makedirs(path, exist_ok=True)
usage = shutil.disk_usage(path)
free_gb = int(usage.free / (1024 ** 3))
total_gb = int(usage.total / (1024 ** 3))
pct = int(usage.free * 100 / usage.total) if usage.total else 0
return free_gb, total_gb, pct
def sanitize_take_name(name):
name = re.sub(r'[<>:"/\\|?*]', "_", (name or "").strip())
return name or "take"
# ---- Rapatriement (Gather) ---------------------------------------------------
# Taille de lecture pour l'upload. Sert aussi de granularite a l'avancement :
# pysmb rappelle notre lecteur a chaque bloc.
GATHER_CHUNK = 1024 * 1024
def norm_rel(rel):
"""Chemin relatif normalise pour COMPARAISON avec la liste du serveur.
Le serveur enumere sous Windows et envoie donc des antislash
("cap\\mkv\\x.mkv") ; nous produisons des slash. WSClient.NormalizeRel fait
la meme conversion et compare sans tenir compte de la casse
(StringComparer.OrdinalIgnoreCase) : sans ca, aucun fichier ne serait jamais
reconnu comme deja present et on re-uploaderait tout a chaque fois."""
return rel.replace("/", "\\").lstrip("\\").lower()
class ProgressReader(object):
"""Fichier en lecture qui rapporte l'avancement.
pysmb.storeFile() lit par blocs dans l'objet qu'on lui donne : en
l'enveloppant, on obtient un avancement a l'octet — la meme granularite que
SMBUploader.OnProgressChanged cote Windows — sans dependre d'une API de
progression que pysmb n'expose pas."""
def __init__(self, path, on_bytes):
self._f = open(path, "rb")
self._on_bytes = on_bytes
def read(self, size=-1):
data = self._f.read(size)
if data:
self._on_bytes(len(data))
return data
def close(self):
self._f.close()
def __enter__(self):
return self
def __exit__(self, *exc):
self.close()
class Bridge:
def __init__(self, cfg, cam_settings_store):
self.cfg = cfg
self.cam_settings = cam_settings_store
self.my_ip = None
self.ws = None
# serial -> {"serial", "model", "port", "camIndex"}. Keyed by serial to
# match how the mesh identifies a camera (CamerasBridge.cs keeps the very
# same dict); several cameras share one computerHostname, which is only
# the node they are displayed under.
self.cameras = {}
# Sane defaults until the mesh sends a RecordSettings event.
self.record_settings = {
"Duration": 0,
"DurationBool": True,
"DepthMode": 1,
"ColorMode": 1,
"AlignMode": 0,
"FpsMode": 30,
"SyncDelay": 0,
}
self.procs = {} # serial -> Popen, one recorder per camera
# One preview server for the whole node, not one per camera: it opens
# every camera itself and stamps each frame with its camIndex, so the
# viewer needs a single connection and the cameras of one node cannot
# drift apart.
self.stream_proc = None
self.stream_serials = set() # cameras the server was started with
# Ce qui etait PHYSIQUEMENT sur le bus au demarrage du serveur. Sert a
# distinguer un branchement a chaud (qui doit relancer) d'un simple
# changement de case (qui ne doit pas).
self.stream_present = set()
# Le mode d alignement ecrit dans le fichier de composition, pour savoir
# quand il change. None tant que rien n a ete ecrit.
self.stream_align = None
self.stream_options = None # and the settings it was started with
# Preview is an intent, not just a process: with every camera unplugged
# there is nothing running, and we must still know to start streaming
# again when one comes back.
self.preview_active = False
self.current_take = None
self._warned_no_camera = False
self._gathering = False # un seul rapatriement a la fois
def _allowed(self, serial):
"""cfg["camera_serial"]: None = all; a string or a list = whitelist."""
want = self.cfg.get("camera_serial")
if not want:
return True
if isinstance(want, (list, tuple, set)):
return serial in want
return serial == want
# ---- connection lifecycle ------------------------------------------------
async def run(self):
while True:
try:
await self.connect_and_serve()
except (websockets.ConnectionClosed, OSError) as e:
logging.warning("Connection lost (%s) — retrying in 3s.", e)
except Exception:
logging.exception("Unexpected error — retrying in 3s.")
self.ws = None
# Unity est parti : on eteint les cameras. Sans ca le noeud continuait
# d'ouvrir ses capteurs et de streamer pour personne, indefiniment — la
# preview d'Unity ne redemande rien d'elle-meme a la reconnexion.
#
# Un enregistrement en cours n'est pas concerne : stop_streams ne
# connait que le serveur de preview, et il n'y en a pas pendant un
# record. Le retour se fait tout seul : Unity, s'il est toujours en
# preview, renvoie l'ordre a l'arrivee du noeud (PreviewSyncBridge).
if self.preview_active or self.stream_proc is not None:
logging.info("Liaison avec Unity perdue — extinction des caméras.")
self.preview_active = False
await self.stop_streams()
await asyncio.sleep(3)
async def discover_server(self):
if self.cfg["server_ip"]:
return self.cfg["server_ip"]
loop = asyncio.get_event_loop()
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_BROADCAST, 1)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
sock.settimeout(1.5) # blocking with a timeout; driven via run_in_executor below
sock.bind(("", 0))
logging.info("Discovering Ithaca server via UDP broadcast on :%d ...", DISCOVERY_PORT)
try:
while True:
sock.sendto(DISCOVERY_MSG, ("255.255.255.255", DISCOVERY_PORT))
try:
# loop.sock_recv() doesn't return the sender address on
# Python 3.6 (no sock_recvfrom() until 3.11) — run the
# blocking recvfrom() (with sock.settimeout above) in the
# executor instead.
data, addr = await loop.run_in_executor(None, sock.recvfrom, 1024)
except socket.timeout:
continue
text = data.decode("utf-8", "replace")
if not text.startswith("WS_SERVER:"):
continue
parts = text.split(":")
ip = parts[2] if len(parts) > 2 and parts[2] else addr[0]
logging.info("Found Ithaca server at %s (from %s)", ip, addr[0])
return ip
finally:
sock.close()
def determine_my_ip(self, server_ip):
"""LAN IP the server will see us connect from — used as computerHostname."""
s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
try:
s.connect((server_ip, 1))
return s.getsockname()[0]
finally:
s.close()
async def connect_and_serve(self):
server_ip = await self.discover_server()
self.my_ip = self.determine_my_ip(server_ip)
uri = "ws://%s:%d" % (server_ip, WS_PORT)
logging.info("Connecting to %s as %s ...", uri, self.my_ip)
async with websockets.connect(uri) as ws:
self.ws = ws
logging.info("Connected.")
# Drop what we knew: forces a fresh CameraUpsert for every camera,
# even ones already announced before the reconnect. The relay keeps
# no state (server.js just rebroadcasts), so nothing replays for us.
self.cameras = {}
await self.refresh_cameras()
watcher = asyncio.ensure_future(self.camera_watch_loop())
try:
async for raw in ws:
await self.handle_message(raw)
finally:
watcher.cancel()
# ---- outgoing -------------------------------------------------------------
async def send(self, msg_type, action, payload, request_id=None, target=None):
envelope = {
"Type": msg_type,
"Action": action,
"RequestId": request_id,
"Target": target,
"Payload": payload,
"Success": False,
"Error": None,
}
if self.ws is None:
logging.debug("Not connected — dropping %s/%s.", msg_type, action)
return
await self.ws.send(json.dumps(envelope))
async def camera_watch_loop(self):
"""Polls for USB plug/unplug every few seconds and syncs the mesh.
Simple polling (not pyudev/inotify) — no extra dependency, and a
couple of seconds of latency is irrelevant for a recording rig.
Also re-announces every camera every ~30s so freeSpace/
percentFreeSpace stay live in Unity's UI as recordings fill the disk."""
tick = 0
while True:
await asyncio.sleep(3)
tick += 1
try:
await self.refresh_cameras()
if self.preview_active:
await self._reconcile_streams()
if self.cameras and tick % 10 == 0:
await self.flush_cameras()
# Ecrit les reglages en attente (voir CamSettingsStore._save).
self.cam_settings.flush()
except asyncio.CancelledError:
# Python 3.6 (this Jetson): CancelledError derives from Exception,
# so the broad except below would swallow connect_and_serve's
# watcher.cancel() and leave this loop running after the socket
# it belongs to is gone — one zombie per reconnect, each logging
# a send failure every 30s.
raise
except Exception:
logging.exception("Error while polling for camera changes.")
async def refresh_cameras(self):
cams = [c for c in detect_cameras() if self._allowed(c["serial"])]
present = {c["serial"]: c for c in cams}
# The store refuses to swap an index with an absent camera, so it needs
# to know who is here.
self.cam_settings.present = set(present)
# Gone: one CameraRemove per camera. Iterate over a copy — we mutate.
for serial in list(self.cameras):
if serial not in present:
logging.info("Camera %s disconnected.", serial)
await self.send(MessageType.EVENT, ActionType.CAMERA_REMOVE,
{"Ip": serial}) # 'Ip' carries the Serial, as WSClient.cs:264 does
self.cameras.pop(serial, None)
proc = self.procs.pop(serial, None)
if proc is not None and proc.poll() is None:
logging.warning("Camera %s vanished mid-recording — killing "
"its recorder (pid %d).", serial, proc.pid)
proc.kill()
# The preview server keeps running for the cameras that remain;
# _reconcile_streams restarts it if this one comes back.
# New: one CameraUpsert per camera. The mesh has no "list" message —
# a list is N idempotent upserts sharing one computerHostname, exactly
# what CamerasBridge.FlushLocalCameras does on the Unity side.
for cam in cams:
if cam["serial"] in self.cameras:
continue
# Reglages restaures depuis le store : ce sont NOS cameras, donc c'est
# ici que vit leur configuration, pas cote Unity.
cam = dict(cam, camIndex=self.cam_settings.next_index(cam["serial"]))
cam.update(self.cam_settings.settings(cam["serial"]))
self.cameras[cam["serial"]] = cam
self._warned_no_camera = False
logging.info("Camera %s (%s, camIndex %d, port %s) connected — "
"announcing as %s", cam["serial"], cam["model"],
cam["camIndex"], cam["port"], self.my_ip)
await self.send_camera_upsert(cam)
if not cams and not self._warned_no_camera:
logging.warning("No Orbbec camera detected on this Jetson (sysfs found none).")
self._warned_no_camera = True
async def flush_cameras(self):
"""Re-announce every local camera. Idempotent (keyed by Serial), so it is
safe to call on reconnect, when a peer joins, or periodically."""
for cam in list(self.cameras.values()):
await self.send_camera_upsert(cam)
async def send_camera_upsert(self, cam):
free_gb, total_gb, pct = free_space(self.cfg["working_directory"])
payload = {
"Serial": cam["serial"],
# L'IP reste l'identite du poste dans le mesh (cle des widgets, de
# l'upsert du noeud Computer, et du retrait par ClientsBridge).
"computerHostname": self.my_ip,
# Nom de machine, pour l'affichage seulement : sans lui l'UI montrait
# "192.168.1.152 192.168.1.152" au lieu du nom du poste.
"computerName": socket.gethostname(),
"camType": CAM_TYPE_ORBBEC,
"camModel": cam["model"],
"freeSpace": free_gb,
"totalSpace": total_gb,
"percentFreeSpace": pct,
# Reglages tels que persistes : l'UI affiche donc ce qui sera
# reellement applique, et une valeur editee survit a un redemarrage.
"syncMode": self.cam_sync_wire(cam),
"camIndex": cam["camIndex"],
"exposureAuto": cam.get("exposureAuto", True),
"whiteBalanceAuto": cam.get("whiteBalanceAuto", True),
# Comme l'orientation : omis ici, il repartirait a son defaut (False)
# et recocherait la case que l'operateur vient de decocher.
"disabled": bool(cam.get("disabled", False)),
# Doit figurer ici, pas seulement dans le store : ce que la re-annonce
# omet arrive a 0 chez Unity et ecrase le reglage de l'operateur.
"orientationQuarters": cam.get("orientationQuarters", 0),
}
for field in MANUAL_CONTROLS:
payload[field] = cam.get(field, 0)
await self.send(MessageType.EVENT, ActionType.CAMERA_UPSERT, payload)
def cam_sync_wire(self, cam):
"""Mode de sync de la camera, en valeur de fil. config.json sert de defaut
tant que le mesh n'a rien impose."""
wire = cam.get("syncMode")
if wire in SYNC_MODE_NAME:
return wire
return SYNC_MODE_WIRE.get(self.cfg["sync_mode"], 0)
# ---- incoming ---------------------------------------------------------------
async def handle_message(self, raw):
raw = raw.strip()
# The server sends a plain-text "Welcome client!" greeting on connect —
# not JSON. Ignore anything that isn't a JSON object/array.
if not raw or raw[0] not in "{[":
logging.debug("Ignoring non-JSON message: %r", raw[:60])
return
try:
msg = json.loads(raw)
except ValueError:
logging.warning("Bad JSON from server: %r", raw[:200])
return
action = msg.get("Action")
payload = msg.get("Payload") or {}
if action == ActionType.RECORD:
await self.on_record(payload)
elif action == ActionType.PREVIEW:
await self.on_preview()
elif action == ActionType.STOP:
await self.on_stop()
elif action == ActionType.RECORD_SETTINGS:
self.on_record_settings(payload)
elif action == ActionType.CAMERA_UPSERT:
self.on_camera_upsert(payload)
elif action == ActionType.CLIENT_CONNECTED:
# A node just joined. It missed every upsert we sent before it
# connected (server.js rebroadcasts live, it replays nothing), so
# push our whole camera list at it. Same reason CamerasBridge.cs
# wires OnClientConnected -> FlushLocalCameras. Without this a
# newcomer would only see us at the next 30s periodic re-announce.
logging.info("Peer %s joined — re-announcing our %d camera(s).",
payload.get("Ip"), len(self.cameras))
await self.flush_cameras()
elif action == ActionType.GATHER_TAKES:
await self.on_gather_takes(payload)
# ClientDisconnected/TakeName/AutoGather/DeleteAfterUpload/GetClient etc.:
# nothing to do for a headless recorder node, ignored.
def on_record_settings(self, payload):
for k, v in payload.items():
if v is not None:
self.record_settings[k] = v
logging.info("RecordSettings updated: %s", self.record_settings)
def on_camera_upsert(self, payload):
"""Reglages recus pour UNE DE NOS cameras : on les applique et on les garde.
C'est le pendant de CameraPersistence cote Unity, qui ne sauvegarde
volontairement que les cameras locales. Sans ce stockage, tout reglage fait
dans l'UI etait recu puis jete, et notre re-annonce periodique (30 s)
remettait l'ancienne valeur : l'operateur voyait ses choix ne pas tenir."""
serial = payload.get("Serial")
cam = self.cameras.get(serial)
if cam is None:
return # camera d'un autre poste, pas la notre
changes = {}
idx = payload.get("camIndex")
if isinstance(idx, int) and idx >= 0:
if self.procs and idx != cam.get("camIndex"):
logging.warning("camIndex de %s non modifié : enregistrement en cours "
"(il détermine le nom du fichier).", serial)
else:
changes["camIndex"] = idx
wire = payload.get("syncMode")
if wire in SYNC_MODE_NAME:
changes["syncMode"] = wire
for field in ("exposureAuto", "whiteBalanceAuto", "disabled"):
value = payload.get(field)
if isinstance(value, bool):
changes[field] = value
# Tous les controles image sont memorises, sans tenter de deviner si la
# valeur vient d'un curseur deplace ou d'une remontee automatique du
# capteur : rien dans le payload ne permet de les distinguer. L'UI est la
# source de verite — si elle affiche contraste=1, c'est 1 qu'il faut
# stocker et appliquer, sinon l'interface mentirait sur l'etat reel.
# Le cout en ecritures disque est traite par l'amortissement du store, pas
# en refusant de stocker (une version precedente le faisait, et la
# re-annonce renvoyait alors 0 pour ces champs, ce qui effacait le reglage
# de l'operateur toutes les 30 s).
for field in MANUAL_CONTROLS:
value = payload.get(field)
if isinstance(value, int) and not isinstance(value, bool) and value >= 0:
changes[field] = value
# Orientation du montage. Ramenee dans 0-3 plutot que rejetee : un quart de
# tour de trop veut dire la meme chose qu'un tour complet en moins, et un
# reglage refuse en silence serait renvoye tel quel par la re-annonce.
quarters = payload.get("orientationQuarters")
if isinstance(quarters, int) and not isinstance(quarters, bool):
changes["orientationQuarters"] = quarters % ORIENTATION_QUARTERS
touched = self.cam_settings.update(serial, changes)
if not touched:
return
# Recharge depuis le store : un echange d'index a pu deplacer une AUTRE
# camera, dont l'UI doit etre corrigee elle aussi.
for c in self.cameras.values():
c.update(self.cam_settings.settings(c["serial"]))
asyncio.ensure_future(self.flush_cameras())
async def on_record(self, payload):
if self.procs:
logging.warning("Record received while already recording — ignoring.")
return
if not self.cameras:
logging.error("Record received but no camera was detected — ignoring.")
return
# The streamers hold the cameras exclusively, so a recorder would fail
# with uvc_open -6. Stop them and give the kernel a moment to release
# the USB devices before opening them again.
if self.stream_proc is not None:
logging.info("Record received while streaming — stopping the streams first.")
self.preview_active = False
await self.stop_streams()
await asyncio.sleep(3)
# The take name is computed once by whichever node triggered the
# recording and broadcast verbatim; every node must reuse it as-is
# (no local timestamping) so files land in the same take folder.
take = sanitize_take_name(payload.get("TakeName"))
self.current_take = take
out_dir = os.path.join(self.cfg["working_directory"], take, "cap", "mkv")
os.makedirs(out_dir, exist_ok=True)
# Deux master sur le meme bus est physiquement impossible : chacun envoie
# un reset de timestamp materiel + un pulse, et ils se dechassent du bus
# (segfault d'un cote, uvc_stream_open_ctrl failed de l'autre). On previent
# sans decider a la place de l'operateur.
masters = [c["serial"] for c in self._enabled_cams() if c.get("sync_mode") == "master"]
if len(masters) > 1:
logging.error("%d caméras réglées en 'master' sur cette Jetson (%s) — elles vont "
"se perturber mutuellement. Un seul master par bus, les autres en "
"'subordinate' avec le câble de sync branché.",
len(masters), ", ".join(masters))
# Record/Stop carry no Serial (WSClient.SendRecord only sends TakeName),
# so a Record means "every TICKED camera on this node" — one process each.
# They are launched in camIndex order for reproducible file numbering.
cams = self._enabled_cams()
if not cams:
logging.error("Record reçu mais toutes les caméras de ce nœud sont "
"décochées — rien à enregistrer.")
return
for cam in cams:
self._spawn_recorder(cam, take, out_dir)
def _spawn_recorder(self, cam, take, out_dir):
cam_index = cam["camIndex"]
out_file = os.path.join(out_dir, "%s_%d.mkv" % (take, cam_index))
log_file_path = os.path.join(out_dir, "%s_%d.log" % (take, cam_index))
rs = self.record_settings
args = [
self.cfg["recorder_bin"],
# Target this exact sensor. Without it the binary takes --device 0,
# so every process would grab the same camera.
"--serial", cam["serial"],
"--color-mode", COLOR_MODE_MAP.get(rs.get("ColorMode", 1), "1080p"),
"--depth-mode", DEPTH_MODE_MAP.get(rs.get("DepthMode", 1), "NFOV_UNBINNED"),
"--rate", str(rs.get("FpsMode", 30)),
]
if not rs.get("DurationBool", True) and rs.get("Duration"):
args += ["--record-length", str(rs["Duration"])]
# Keep concurrent recorders off each other's core; --pin-core defaults to
# --device, which we never pass, so it has to be set explicitly.
if len(self.cameras) > 1:
args += ["--pin-cpu", "--pin-core", str(cam_index)]
# Exposition et balance des blancs : le payload porte un booleen "auto"
# explicite, donc l'intention de l'operateur est connue sans ambiguite.
# Les valeurs partent telles quelles — le recorder interroge la plage
# reelle du capteur et clampe en journalisant (setColorInt), il n'y a donc
# pas d'echelle a convertir ici.
if not cam.get("exposureAuto", True) and cam.get("exposure"):
args += ["--exposure-control", str(cam["exposure"])]
if not cam.get("whiteBalanceAuto", True) and cam.get("whiteBalance"):
args += ["--whitebalance", str(cam["whiteBalance"])]
# Les cinq controles sans booleen "auto" : on applique ce que l'UI affiche.
# 0 reste le sentinelle "jamais transmis" (aucun de ces controles ne vaut
# 0 dans les plages Orbbec, sauf brightness dont 0 est aussi le defaut :
# l'omettre revient donc au meme).
if self.cfg.get("apply_image_controls", True):
for field, flag in (("brightness", "--brightness"), ("contrast", "--contrast"),
("saturation", "--saturation"), ("sharpness", "--sharpness"),
("gain", "--gain")):
if cam.get(field):
args += [flag, str(cam[field])]
# Mode par camera : celui envoye par le mesh (UI Unity) s'il est arrive,
# sinon la valeur de config.json.
sync_mode = SYNC_MODE_NAME.get(self.cam_sync_wire(cam), "standalone")
if sync_mode == "master":
args.append("--master")
elif sync_mode == "subordinate":
args.append("--subordinate")
if rs.get("SyncDelay"):
args += ["--sync-delay", str(rs["SyncDelay"])]
else:
args.append("--standalone")
args.append(out_file)
logging.info("Starting recording cam %d (%s) -> %s",
cam_index, cam["serial"], out_file)
logging.debug("argv: %s", args)
# IMPORTANT: never use stdout=subprocess.PIPE without something
# continuously reading it. The K4A wrapper's bundled OrbbecSDK v1.10.18
# logs verbosely to stdout; an unread pipe fills its ~64KB kernel
# buffer within seconds and the recorder blocks on write() forever —
# including inside its own SIGTERM cleanup path, so Stop would never
# take effect and the camera is left mid-stream until force-killed.
# A real file has no such backpressure.
log_file = open(log_file_path, "wb")
try:
proc = subprocess.Popen(
args, stdin=subprocess.DEVNULL, stdout=log_file, stderr=subprocess.STDOUT
)
finally:
log_file.close() # the child keeps its own dup'd fd open; we don't need ours
self.procs[cam["serial"]] = proc
asyncio.ensure_future(self._watch_process(cam["serial"], proc))
async def _watch_process(self, serial, proc):
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, proc.wait)
# Only clear our slot if it is still this process (a new take may have
# replaced it while we were waiting).
if self.procs.get(serial) is proc:
logging.info("Recorder for %s exited (code %s).", serial, proc.returncode)
del self.procs[serial]
# ---- rapatriement vers le partage SMB du serveur ------------------------
async def on_gather_takes(self, payload):
"""Uploade nos fichiers de prise vers le partage SMB annonce par le serveur.
Reproduit WSClient.HandleGatherTakes : meme regle de saut, meme
arborescence de destination, meme avancement pondere par les octets, meme
signal de fin. Le serveur enumere SES takes et nous envoie, pour chacun,
ce qu'il possede deja ; on ne transmet que le manquant.
Les identifiants viennent du payload et n'existent que le temps de la
demande : le partage est cree au demarrage du serveur Unity
(WindowsUserManager.SetupUserAndShare) et supprime a son arret. On ne les
conserve donc jamais."""
if self._gathering:
logging.warning("GatherTakes reçu alors qu'un rapatriement est déjà en cours — ignoré.")
return
if not payload or not payload.get("takes"):
logging.warning("GatherTakes sans takes — rien à faire.")
return
try:
from smb.SMBConnection import SMBConnection
except ImportError:
logging.error("GatherTakes impossible : pysmb absent (pip3 install --user pysmb).")
return
server_ip = payload.get("serverIp")
share = payload.get("shareName")
user = payload.get("username")
pwd = payload.get("password")
delete_after = bool(payload.get("deleteAfterUpload"))
if not server_ip or not share:
logging.error("GatherTakes : serverIp/shareName manquants — abandon.")
return
wd = self.cfg["working_directory"]
uploads = self._plan_uploads(payload["takes"], wd)
total_bytes = sum(u[2] for u in uploads)
logging.info("Gather : %d fichier(s) à envoyer (%.1f Mo) vers //%s/%s%s",
len(uploads), total_bytes / (1024.0 ** 2), server_ip, share,
" [suppression après upload]" if delete_after else "")
self._gathering = True
loop = asyncio.get_event_loop()
try:
await self.send_progress(0.0, "")
await loop.run_in_executor(
None, self._do_gather, SMBConnection, server_ip, share, user, pwd,
uploads, total_bytes, delete_after, loop)
await self.send_progress(100.0, "")
await self.send(MessageType.EVENT, ActionType.STATUT,
{"Client": self.my_ip, "File": "", "Statut": "done"})
logging.info("Gather terminé.")
except Exception:
logging.exception("Gather : échec.")
finally:
self._gathering = False
def _plan_uploads(self, takes, wd):
"""Diff local/serveur -> [(chemin_local, dest_rel, taille)].
dest_rel est en antislash ("<take>\\cap\\mkv\\x.mkv") : c'est la forme que
le serveur affiche dans la popup de progression."""
uploads = []
skipped = 0
# Une prise en cours d'écriture n'est pas finalisée : ses MKV n'ont pas
# d'index et deleteAfterUpload effacerait un fichier encore ouvert.
busy_take = self.current_take if self.procs else None
for entry in takes:
take = (entry or {}).get("take")
if not take:
continue
if take == busy_take:
logging.info("Gather : take '%s' en cours d'enregistrement — ignoré.", take)
continue
take_dir = os.path.join(wd, take)
cap_dir = os.path.join(take_dir, "cap")
if not os.path.isdir(cap_dir):
continue
server_map = {}
for f in (entry.get("serverFiles") or []):
if f and f.get("rel"):
server_map[norm_rel(f["rel"])] = f.get("size", 0)
for root, _dirs, files in os.walk(cap_dir):
for name in files:
local = os.path.join(root, name)
rel = os.path.relpath(local, take_dir)
size = os.path.getsize(local)
srv = server_map.get(norm_rel(rel))
# Present ET complet cote serveur -> on saute. Un fichier
# tronque (upload interrompu) est renvoye entierement.
if srv is not None and srv >= size:
skipped += 1
continue
uploads.append((local, os.path.join(take, rel).replace("/", "\\"), size))
if skipped:
logging.info("Gather : %d fichier(s) déjà présent(s) côté serveur — sautés.", skipped)
return uploads
def _do_gather(self, SMBConnection, server_ip, share, user, pwd,
uploads, total_bytes, delete_after, loop):
"""Partie bloquante, exécutée dans un thread : connexion et transferts."""
# remote_name : l'IP est acceptee par le serveur (le nom NetBIOS et
# '*SMBSERVER' aussi), et c'est la seule valeur que le payload nous donne.
conn = SMBConnection(user or "", pwd or "", socket.gethostname(), server_ip,
use_ntlm_v2=True, is_direct_tcp=True)
if not conn.connect(server_ip, 445, timeout=20):
raise RuntimeError("authentification refusée sur //%s/%s" % (server_ip, share))
try:
done_bytes = [0]
last_pct = [-1.0]
def report(nbytes):
done_bytes[0] += nbytes
pct = (done_bytes[0] * 100.0 / total_bytes) if total_bytes else 100.0
# Anti-flood : un envoi par point de pourcentage, comme le client Windows.
if pct - last_pct[0] >= 1.0:
last_pct[0] = pct
asyncio.run_coroutine_threadsafe(
self.send_progress(pct, current[0]), loop)
current = [""]
for local, dest_rel, size in uploads:
current[0] = dest_rel
smb_path = "/" + dest_rel.replace("\\", "/")
self._ensure_dirs(conn, share, smb_path)
try:
with ProgressReader(local, report) as fh:
conn.storeFile(share, smb_path, fh)
except Exception as e:
logging.error("Gather : échec de '%s' : %s", dest_rel, e)
continue
logging.info("Gather : envoyé %s (%.1f Mo)", dest_rel, size / (1024.0 ** 2))
if delete_after:
try:
os.remove(local)
except OSError as e:
logging.warning("Gather : suppression impossible '%s' : %s", local, e)
finally:
conn.close()
@staticmethod
def _ensure_dirs(conn, share, smb_path):
"""Cree l'arborescence du fichier. createDirectory echoue si le dossier
existe deja : c'est le cas nominal des le second fichier d'un take."""
parts = smb_path.strip("/").split("/")[:-1]
path = ""
for p in parts:
path += "/" + p
try:
conn.createDirectory(share, path)
except Exception:
pass
async def send_progress(self, percent, dest_rel):
try:
await self.send(MessageType.EVENT, ActionType.PROGRESS,
{"Client": self.my_ip, "File": dest_rel or "",
"Progress": round(percent, 2)})
except Exception:
logging.debug("Gather : envoi de progression impossible (socket fermée ?).")
# ---- Preview streaming --------------------------------------------------
#
# Unity's Play button in the Record context broadcasts Preview, and Stop
# broadcasts Stop. A camera can only be opened by one process, so streaming
# and recording are mutually exclusive: Record stops the streams first, and
# Preview refuses outright while a recording runs rather than half-starting
# and leaving some cameras dark.
async def on_preview(self):
if self.procs:
logging.warning("Preview received while recording — ignoring.")
return
if self.stream_proc is not None:
logging.info("Preview received but already streaming — nothing to do.")
return
if not self.cameras:
logging.error("Preview received but no camera was detected — ignoring.")
return
# Les cases, pas ce qui est branche : sans ce test on demarrait avec un
# --map vide, et un map vide veut dire "toutes les cameras".
if not self._enabled_cams():
logging.warning("Preview reçue mais toutes les caméras de ce nœud sont "
"décochées — rien à servir.")
return
# Take over any streamer we do not own: one started by hand, or left by a
# previous run of this bridge. They hold the cameras, so ours would fail
# to open them and the preview would stay dark.
self._reap_foreign_streamers()
self._spawn_streamer()
self.preview_active = True
async def _reconcile_streams(self):
"""Keeps the preview server covering every present camera.
Covers both directions of a hot unplug: the server dying takes the
preview with it, and a camera that comes back is not in the server that
is already running, since it opens its cameras once at start. Called
from the camera poll loop, so recovery follows a replug within a few
seconds.
A camera that merely goes away does not trigger a restart: the server
keeps serving the others, and interrupting them to react to a departure
would cost more than it gains."""
if self.stream_proc is not None and self.stream_proc.poll() is not None:
logging.warning("The preview server exited (rc=%d) — restarting it.",
self.stream_proc.returncode)
self.stream_proc = None
self.stream_serials = set()
self.stream_present = set()
present = {c["serial"] for c in self._enabled_cams()}
options = self._stream_options()
# Un branchement a chaud relance le serveur, comme avant : la camera n'a
# aucun pipeline dans le processus en cours, et seule une nouvelle enumeration
# la trouve. Un changement de CASE, non — il passe par le fichier de
# composition, plus bas, et le serveur ferme ou rouvre le capteur concerne en
# une seconde sans toucher aux autres. Relancer coutait une dizaine de
# secondes de noir pour TOUTES les cameras du noeud.
#
# Les deux se distinguent par ce qui etait sur le BUS au demarrage, pas par
# ce qui etait servi : une camera decochee puis recochee n'a jamais quitte le
# bus, donc elle ne compte pas comme une arrivee.
appeared = sorted(set(self.cameras) - self.stream_present)
# D'abord, avant toute comparaison : n'avoir rien a servir est un etat, pas un
# changement a reconcilier. Teste plus bas, le raccourci "rien n'a change"
# sortait avant et laissait le serveur tourner pour personne.
if not present:
# Pas d'interruption visible a craindre, il n'y a plus aucune carte.
if self.stream_proc is not None:
logging.info("Plus aucune caméra cochée — arrêt du serveur de preview.")
await self.stop_streams()
return
# Composition changee alors que le serveur tourne : on reecrit le fichier et
# c'est tout. Il fermera ou rouvrira la camera concernee en une seconde, les
# autres continuant d'emettre — ce que le redemarrage ne savait pas faire.
#
# Fait avant le raccourci ci-dessous, sinon "rien n'a change" sortirait le
# premier : une case cochee ne change ni les options ni ce qui est sur le bus.
align_changed = self._wanted_align() != self.stream_align
if self.stream_proc is not None and (present != self.stream_serials
or align_changed):
added = sorted(present - self.stream_serials)
gone = sorted(self.stream_serials - present)
before = self.stream_align
if self._write_stream_map(self._enabled_cams()):
logging.info("Composition mise à jour sans redémarrage%s%s%s.",
" (+%s)" % ", ".join(added) if added else "",
" (-%s)" % ", ".join(gone) if gone else "",
" (alignement %s -> %s)" % (before, self.stream_align)
if align_changed else "")
self.stream_serials = set(present)
covered = self.stream_proc is not None and not appeared
current = self.stream_proc is None or options == self.stream_options
if covered and current:
return # nothing to change
if self.stream_proc is not None:
if appeared:
logging.info("Camera(s) %s appeared while previewing — restarting "
"the preview server to include them.",
", ".join(appeared))
else:
logging.info("Capture settings changed while previewing — "
"restarting the preview server with %s.",
" ".join(options))
await self.stop_streams()
await asyncio.sleep(2) # let the USB devices be released
self._spawn_streamer()
@staticmethod
def _reap_foreign_streamers():
# The supervisors first: killed after their streamer, they would restart
# it. The [p]/[e] brackets keep each pattern from matching the pkill
# itself, and -f is required because Linux truncates the process name to
# 15 characters ("ob_stream_previ"), which -x would then never match.
#
# The first two belong to the retired H.264/WebRTC route, and are still
# reaped here: they hold the cameras just as firmly, and one left running
# by hand is exactly the case this exists for.
patterns = ["stream_on[e].sh", "ob_stream_preview_rts[p]",
"ithaca_rgbd_serve[r]"]
killed = False
for pattern in patterns:
if subprocess.call(["pkill", "-f", pattern]) == 0:
killed = True
if killed:
logging.info("Stopped streamers we did not start; the cameras need a "
"moment to be released.")
time.sleep(10)
def _stream_options(self):
"""The preview server's settings, as the command line it would be given.
Compared against what the running server was started with, since it reads
them once and a change has to become a restart to take effect."""
rs = self.record_settings
return [
"--color-mode", COLOR_MODE_MAP.get(rs.get("ColorMode", 0), "720p"),
"--depth-mode", DEPTH_MODE_MAP.get(rs.get("DepthMode", 0), "NFOV_2X2BINNED"),
"--rate", str(rs.get("FpsMode", 30)),
# Toujours raw ICI, et ce n'est pas le mode effectif : celui-ci est la
# deuxieme ligne du fichier de composition, que le serveur relit une fois
# par seconde (voir _wanted_align et _write_stream_map). Ainsi Compare
# s'allume et s'eteint sans relancer le serveur, donc sans couper les
# cameras du noeud une dizaine de secondes.
#
# Ce champ ne doit surtout PAS suivre le menu : il fait partie des options
# comparees pour decider d'un redemarrage, donc l'y mettre rendrait chaque
# changement d'alignement couteux a nouveau. Il reste le defaut sur pour un
# lancement a la main.
"--align", "raw",
]
def _wanted_align(self):
"""Le mode que le noeud doit produire.
Seuls deux nous concernent : "compare" quand l operateur veut la transformation
du SDK a cote des trames brutes, "raw" sinon. D2C et C2D sont calcules par le
viewer — c est la seule facon que le menu agisse en direct."""
return "compare" if self.record_settings.get("AlignMode", 0) == 3 else "raw"
def _stream_map_path(self):
# A cote du store de reglages : le Bridge ne connait pas le repertoire de
# config, mais le store, lui, porte son propre chemin.
return os.path.join(os.path.dirname(self.cam_settings.path), STREAM_MAP_FILE)
def _write_stream_map(self, cams):
"""Ecrit la composition que le serveur doit servir.
Ecriture atomique (fichier temporaire puis rename) : le serveur relit ce
fichier sur minuterie, et il ne doit jamais en lire une moitie. Il ignore
de son cote un fichier vide ou illisible, ce qui le protege du reste."""
mapping = ",".join("%s:%d" % (c["serial"], c["camIndex"]) for c in cams)
align = self._wanted_align()
path = self._stream_map_path()
try:
tmp = path + ".tmp"
with open(tmp, "w") as f:
# Ligne 1 la composition, ligne 2 le mode : un seul rename change les
# deux, ils ne peuvent donc jamais etre lus en desaccord.
f.write(mapping + "\n" + align + "\n")
os.replace(tmp, path)
except OSError as e:
logging.warning("Ecriture de %s impossible (%s) — la composition ne "
"pourra pas changer sans relancer le serveur.", path, e)
return False
self.stream_align = align
return True
def _enabled_cams(self):
"""Cameras whose checkbox is ticked, in camIndex order.
Une camera decochee dans la liste Recorder n'est pas ouverte du tout : ni
mappee dans le serveur de preview, ni confiee a un recorder. C'est tout
l'objet de la case — un capteur inutilise ne doit rien couter a ce noeud."""
return sorted((c for c in self.cameras.values() if not c.get("disabled")),
key=lambda c: c["camIndex"])
def _spawn_streamer(self):
"""Starts the one preview server that carries every ticked camera."""
cams = self._enabled_cams()
# The camera indices are the ones edited in the Recorder list, and they
# travel in each frame's header: that is how the viewer puts a stream in
# the right slot of the mosaic without asking anyone.
mapping = ",".join("%s:%d" % (c["serial"], c["camIndex"]) for c in cams)
# Un map vide fait servir TOUTES les cameras par le serveur : c'est son
# repli documente quand le bridge ne dit rien. On refuse donc plutot que
# de lancer l'inverse de ce qui est demande.
if not mapping:
logging.warning("Aucune caméra cochée — serveur de preview non lancé.")
return
# The settings come from _stream_options so that the line the server is
# given and the line it is compared against can never drift apart.
# Ecrit avant de lancer : le serveur lit ce fichier au demarrage aussi, et
# c'est lui qui fait autorite ensuite.
self._write_stream_map(cams)
options = self._stream_options()
args = [self.cfg.get("stream_server", DEFAULT_STREAM_SERVER),
"--map", mapping,
"--map-file", self._stream_map_path()] + options
# start_new_session gives it its own process group, so stop_streams can
# signal it and anything it spawns at once. Without it a survivor would
# keep the cameras and the next Record would fail with uvc_open -6.
proc = subprocess.Popen(
args, stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL,
stderr=subprocess.STDOUT, start_new_session=True)
self.stream_proc = proc
# Ce qui est REELLEMENT servi, donc les cameras cochees.
self.stream_serials = {c["serial"] for c in cams}
# Et ce qui etait sur le bus a cet instant : c'est cet ensemble que
# _reconcile_streams compare, pour ne relancer que sur un vrai branchement.
self.stream_present = set(self.cameras)
# Kept so a settings change made while previewing is noticed: the server
# reads its options once, at start, and cannot be told afterwards.
self.stream_options = options
logging.info("Preview server started (pid %d) for %d camera(s): %s",
proc.pid, len(cams), mapping)
async def stop_streams(self):
proc = self.stream_proc
if proc is None:
return
self.stream_proc = None
self.stream_serials = set()
self.stream_present = set()
self.stream_align = None
self.stream_options = None
self._signal_group(proc, signal.SIGTERM)
loop = asyncio.get_event_loop()
try:
await asyncio.wait_for(loop.run_in_executor(None, proc.wait), timeout=8)
except asyncio.TimeoutError:
# Closing the cameras goes through the SDK, which releases its EGL
# contexts on the way out and has been seen to block there. The
# cameras must be free for the next Record, so do not wait forever.
logging.warning("The preview server ignored SIGTERM — killing it.")
self._signal_group(proc, signal.SIGKILL)
logging.info("Streaming stopped.")
@staticmethod
def _signal_group(proc, sig):
try:
os.killpg(os.getpgid(proc.pid), sig)
except (ProcessLookupError, PermissionError):
pass # already gone
async def on_stop(self):
# Stop means "release the cameras", whichever pipeline holds them: the
# two cannot coexist, so a single Stop covers recording and streaming.
self.preview_active = False
await self.stop_streams()
if not self.procs:
return
running = list(self.procs.items())
# SIGTERM everything first, then wait: sequential terminate+wait would
# let the last camera run ~8s longer than the first.
for serial, proc in running:
logging.info("Stopping recording for %s (SIGTERM, pid %d)...", serial, proc.pid)
proc.terminate() # caught by the binary's posixSignalHandler for a clean MKV finalize
loop = asyncio.get_event_loop()
for serial, proc in running:
try:
await asyncio.wait_for(loop.run_in_executor(None, proc.wait), timeout=8)
except asyncio.TimeoutError:
logging.warning("Recorder for %s did not exit within 8s of "
"SIGTERM — killing it.", serial)
proc.kill()
self.procs.pop(serial, None)
def main():
default_cfg_path = os.path.join(os.path.dirname(os.path.abspath(__file__)), "config.json")
cfg_path = sys.argv[1] if len(sys.argv) > 1 else default_cfg_path
cfg = load_config(cfg_path)
logging.basicConfig(
level=getattr(logging, cfg.get("log_level", "INFO"), logging.INFO),
format="%(asctime)s [%(levelname)s] %(message)s",
)
logging.info("Config: %s", cfg_path)
cfg_dir = os.path.dirname(os.path.abspath(cfg_path))
cam_settings = CamSettingsStore(os.path.join(cfg_dir, CAM_SETTINGS_FILE),
os.path.join(cfg_dir, LEGACY_CAM_INDEX_FILE))
bridge = Bridge(cfg, cam_settings)
loop = asyncio.get_event_loop()
try:
loop.run_until_complete(bridge.run())
except KeyboardInterrupt:
logging.info("Interrupted, exiting.")
if __name__ == "__main__":
main()