#!/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 ("_.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_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 " / 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 ("\\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()