From c167c836e61b6f39e6560a079521918c537aaeed Mon Sep 17 00:00:00 2001 From: AngePierre Date: Thu, 10 Sep 2026 15:57:15 +0200 Subject: [PATCH] Portable Ithaca capture node: one install for every Jetson MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- .gitattributes | 12 + .gitignore | 16 + README.md | 104 +++ bridge/bridge.py | 1431 ++++++++++++++++++++++++++++++ bridge/config.template.json | 9 + install.sh | 101 +++ scripts/build_sdk.sh | 82 ++ scripts/build_server.sh | 29 + scripts/detect_platform.sh | 57 ++ scripts/install_deps.sh | 32 + server/CMakeLists.txt | 93 ++ server/ithaca_rgbd_server.cpp | 1058 ++++++++++++++++++++++ systemd/ithaca-bridge.service.in | 21 + systemd/usbfs-memory.service | 12 + udev/99-obsensor-libusb.rules | 111 +++ 15 files changed, 3168 insertions(+) create mode 100644 .gitattributes create mode 100644 .gitignore create mode 100644 README.md create mode 100644 bridge/bridge.py create mode 100644 bridge/config.template.json create mode 100755 install.sh create mode 100755 scripts/build_sdk.sh create mode 100755 scripts/build_server.sh create mode 100755 scripts/detect_platform.sh create mode 100755 scripts/install_deps.sh create mode 100644 server/CMakeLists.txt create mode 100644 server/ithaca_rgbd_server.cpp create mode 100644 systemd/ithaca-bridge.service.in create mode 100644 systemd/usbfs-memory.service create mode 100644 udev/99-obsensor-libusb.rules diff --git a/.gitattributes b/.gitattributes new file mode 100644 index 0000000..f9ff551 --- /dev/null +++ b/.gitattributes @@ -0,0 +1,12 @@ +# This tree runs on Linux (Jetson). Everything checks out with LF, whatever the +# committer's OS — a CRLF shebang (`#!/usr/bin/env bash\r`) is a broken script. +* text=auto eol=lf +*.py text eol=lf +*.sh text eol=lf +*.cpp text eol=lf +*.txt text eol=lf +*.json text eol=lf +*.md text eol=lf +*.in text eol=lf +*.rules text eol=lf +*.service text eol=lf diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..2f6a21e --- /dev/null +++ b/.gitignore @@ -0,0 +1,16 @@ +# Generated for the local user by install.sh — never committed, it holds this +# machine's paths. +bridge/config.json + +# Build outputs. +server/build/ + +# The runtime state the bridge writes beside itself. +bridge/cam_settings.json +bridge/cam_index.json +bridge/stream_map.txt +bridge/*.log + +# Python. +__pycache__/ +*.pyc diff --git a/README.md b/README.md new file mode 100644 index 0000000..ef3223c --- /dev/null +++ b/README.md @@ -0,0 +1,104 @@ +# Ithaca capture node + +The software a Jetson runs to join an Ithaca RGBD capture rig: a Python bridge that +puts the node on the WebSocket mesh and answers Preview / Record, and a C++ preview +server that streams each camera's colour (the sensor's own JPEG) and depth (its own +16-bit millimetres, LZ4-compressed) straight to the Unity viewer. + +One checkout, one command, on any Jetson: + +```sh +git clone ~/ithaca-node +cd ~/ithaca-node +./install.sh +``` + +## Why an install script and not a disk image + +The rig currently spans two Jetsons with **no JetPack in common**: + +| Node | SoC | JetPack | Ubuntu / CUDA | +|------|-----|---------|---------------| +| Jetson Nano (2019) | tegra210 | **4.x only** (EOL) | 18.04 / 10.2 | +| Jetson Orin Nano | tegra234 | **5 or 6** | 22.04 / 12.x | + +tegra210 never got a newer L4T and the Orin refuses the old one, so no single OS +image covers both, and no compiled binary is portable between them — CUDA, glibc +and the C++ ABI all differ by generation. + +What *is* portable is the **source and the procedure**. `install.sh` detects the +platform and builds the C++ locally against a locally-built OrbbecSDK, so the same +command produces a correct binary on each machine. Adding a future Orin NX or a Thor +needs nothing new here — the script compiles against whatever toolchain it finds. + +The bridge is pure Python and already runs identically everywhere; keep it that way +by never giving it a dependency that has to be compiled. + +## Layout + +``` +install.sh one entry point, idempotent +bridge/ + bridge.py the node client — shipped verbatim, tested in place + config.template.json @PLACEHOLDERS@ filled per user by install.sh +server/ + ithaca_rgbd_server.cpp the preview server + CMakeLists.txt links a prebuilt libOrbbecSDK, no GStreamer +systemd/ + ithaca-bridge.service.in templated with the user and repo path + usbfs-memory.service raises usbfs_memory_mb for the camera bandwidth +udev/ + 99-obsensor-libusb.rules reference copy (the SDK's own installer wins if present) +scripts/ + detect_platform.sh L4T / JetPack / SoC, as KEY=value + install_deps.sh apt + websockets + build_sdk.sh reuse or build OrbbecSDK v2 for this platform + build_server.sh cmake + make the preview server +``` + +## What install.sh does + +1. Detects L4T / JetPack / SoC (`scripts/detect_platform.sh`). +2. Installs build and run dependencies (no GStreamer — the direct route encodes no + video). +3. Locates a working OrbbecSDK build, or clones and builds one. A node that already + has a good build (the Nano) is left untouched. +4. Builds `ithaca_rgbd_server` against it. +5. Generates `bridge/config.json` for the current user — this is where all the + machine-specific paths live, so `bridge.py` itself stays unmodified. +6. Installs the udev rules and the `usbfs-memory` service, and adds the user to the + `video` group. +7. Generates and starts the `ithaca-bridge` systemd service. + +Re-run it after a `git pull`: it rebuilds, regenerates the unit and config, and +restarts the service. It keeps an existing `config.json` (delete it to regenerate) +and an existing SDK build. + +## The SDK is the one thing to watch on a new JetPack + +`build_sdk.sh` reuses an existing build when it finds one. On a machine with none it +clones the official OrbbecSDK v2 and builds it — this is the step most likely to +need attention on a JetPack it has not been tried on. If it fails, build the SDK by +hand once and re-run with its location: + +```sh +OB_SDK_ROOT=/path/to/OrbbecSDK_v2 ./install.sh +``` + +Override the source with `OB_SDK_REPO` / `OB_SDK_REF` to pin a tag known to build on +your JetPack. + +## Operating the node + +```sh +systemctl status ithaca-bridge # is it up +journalctl -u ithaca-bridge -f # live log +``` + +Runtime state the bridge writes beside itself (git-ignored): + +- `config.json` — this machine's config (generated). +- `cam_settings.json` — per-camera settings, kept across restarts. +- `stream_map.txt` — the live composition: line 1 the cameras to stream, line 2 the + alignment (`raw` or `compare`). The preview server re-reads it once a second, so a + camera can be turned on or off, and Compare switched, without restarting. diff --git a/bridge/bridge.py b/bridge/bridge.py new file mode 100644 index 0000000..ee4ec19 --- /dev/null +++ b/bridge/bridge.py @@ -0,0 +1,1431 @@ +#!/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() diff --git a/bridge/config.template.json b/bridge/config.template.json new file mode 100644 index 0000000..d9ac947 --- /dev/null +++ b/bridge/config.template.json @@ -0,0 +1,9 @@ +{ + "server_ip": null, + "working_directory": "@HOME@/IthacaRecordings", + "recorder_bin": "@SDK_BUILD@/bin/ob_stream_depth_color_k4arec_pipeline", + "stream_server": "@REPO@/server/build/ithaca_rgbd_server", + "sync_mode": "standalone", + "camera_serial": null, + "log_level": "INFO" +} diff --git a/install.sh b/install.sh new file mode 100755 index 0000000..90d107a --- /dev/null +++ b/install.sh @@ -0,0 +1,101 @@ +#!/usr/bin/env bash +# One command, the same on every Jetson: installs the Ithaca capture node from this +# checkout. Detects the platform, installs dependencies, builds the C++ against a +# locally-built OrbbecSDK, generates a config and a service for the CURRENT user, +# and starts the bridge. +# +# Idempotent: safe to re-run after a git pull. It rebuilds the server, rewrites the +# config and unit, and restarts the service; it does not touch a working SDK build. +# +# Why a script and not a prebuilt image: the two Jetsons here have no common +# JetPack — tegra210 caps at 4, tegra234 needs 5+ — so nothing binary is portable +# between them. The source and this procedure are; the binaries are made locally. +set -euo pipefail + +REPO=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +USER_NAME=$(id -un) +HOME_DIR="$HOME" + +echo "==============================================================" +echo " Ithaca node install" +eval "$("$REPO/scripts/detect_platform.sh")" +echo " modele : $MODEL" +echo " L4T : $L4T (JetPack $JETPACK, $SOC, $ARCH)" +echo " user : $USER_NAME home: $HOME_DIR" +echo " repo : $REPO" +echo "==============================================================" + +if [[ "$ARCH" != "aarch64" ]]; then + echo "!! Cette install vise les Jetson (aarch64). Arch detectee: $ARCH." >&2 + echo "!! Continue quand meme dans 3 s (Ctrl-C pour arreter)…" >&2 + sleep 3 +fi + +# ---- 1. dependencies ------------------------------------------------------ +"$REPO/scripts/install_deps.sh" + +# ---- 2. the SDK, then the server ------------------------------------------ +# build_sdk.sh reuses an existing build (leaving a working node untouched) or +# builds one; build_server.sh compiles ithaca_rgbd_server against it. +echo ">> Localisation / compilation de l'OrbbecSDK…" +eval "$("$REPO/scripts/build_sdk.sh")" +echo " OB_SDK_ROOT = $OB_SDK_ROOT" +echo " OB_SDK_BUILD = $OB_SDK_BUILD" + +echo ">> Compilation du serveur de preview…" +SERVER_BIN=$("$REPO/scripts/build_server.sh") +echo " $SERVER_BIN" + +# ---- 3. config, generated for this user ----------------------------------- +# The tested bridge.py ships unchanged; portability lives here, in the config it +# reads. @PLACEHOLDERS@ become real paths for the current user. +CONFIG="$REPO/bridge/config.json" +if [[ -f "$CONFIG" ]]; then + echo ">> config.json existe deja — conserve (supprime-le pour regenerer)." +else + sed -e "s#@HOME@#$HOME_DIR#g" \ + -e "s#@REPO@#$REPO#g" \ + -e "s#@SDK_BUILD@#$OB_SDK_BUILD#g" \ + "$REPO/bridge/config.template.json" > "$CONFIG" + echo ">> config.json genere: $CONFIG" +fi +mkdir -p "$HOME_DIR/IthacaRecordings" + +# ---- 4. udev + usbfs (need root, affect the whole machine) ---------------- +echo ">> Regles udev (acces camera sans root) + usbfs_memory…" +# Prefer the SDK's own installer when it is there — it is the maintained source of +# the rule set; fall back to the copy shipped here. +if [[ -x "$OB_SDK_ROOT/scripts/env_setup/install_udev_rules.sh" ]]; then + sudo bash "$OB_SDK_ROOT/scripts/env_setup/install_udev_rules.sh" || true +else + sudo install -m 0644 "$REPO/udev/99-obsensor-libusb.rules" \ + /etc/udev/rules.d/99-obsensor-libusb.rules +fi +sudo udevadm control --reload-rules && sudo udevadm trigger || true + +sudo install -m 0644 "$REPO/systemd/usbfs-memory.service" \ + /etc/systemd/system/usbfs-memory.service + +# So the camera is reachable without being root right now, not only after a +# reboot: the video group is what the udev rules grant to. +sudo usermod -aG video "$USER_NAME" || true + +# ---- 5. the bridge service, generated for this user ----------------------- +UNIT=/etc/systemd/system/ithaca-bridge.service +sudo sed -e "s#@REPO@#$REPO#g" -e "s#@USER@#$USER_NAME#g" \ + "$REPO/systemd/ithaca-bridge.service.in" | sudo tee "$UNIT" >/dev/null + +sudo systemctl daemon-reload +sudo systemctl enable --now usbfs-memory.service +sudo systemctl enable ithaca-bridge.service +sudo systemctl restart ithaca-bridge.service + +echo "==============================================================" +echo " Installe. Etat du bridge :" +systemctl is-active ithaca-bridge.service || true +echo +echo " Journal en direct : journalctl -u ithaca-bridge -f" +echo " Si la camera n'est vue qu'apres reconnexion USB, c'est le" +echo " groupe 'video' qui vient d'etre ajoute — deconnecte/reconnecte" +echo " la session, ou rebranche la camera." +echo "==============================================================" diff --git a/scripts/build_sdk.sh b/scripts/build_sdk.sh new file mode 100755 index 0000000..5681220 --- /dev/null +++ b/scripts/build_sdk.sh @@ -0,0 +1,82 @@ +#!/usr/bin/env bash +# Makes sure an OrbbecSDK v2 build exists for THIS platform, and prints where it +# is as OB_SDK_ROOT=... / OB_SDK_BUILD=... / OB_SDK_GENERATED=... for build_server.sh. +# +# Never ships a prebuilt binary between Jetsons: CUDA and glibc differ by +# generation, so the SDK is built from source on the machine that will run it. +# The heavy artifact stays out of git. +# +# Order of preference: +# 1. An existing build pointed to by $OB_SDK_ROOT — reused as-is. This is what +# keeps a working node (the Nano already has one) untouched. +# 2. A build at one of the known locations, likewise reused. +# 3. A fresh clone + build. On a new JetPack this is the step to watch; if it +# fails, build the SDK by hand once and re-run install.sh with OB_SDK_ROOT set. +set -euo pipefail + +here=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +eval "$("$here/detect_platform.sh")" + +# The official v2 source. Pinned loosely: a caller can override both to match a +# tag known to build on their JetPack. +SDK_REPO="${OB_SDK_REPO:-https://github.com/orbbec/OrbbecSDK_v2.git}" +SDK_REF="${OB_SDK_REF:-main}" + +find_build() { + local root="$1" + [[ -f "$root/include/libobsensor/ObSensor.hpp" ]] || return 1 + local b + b=$(ls -d "$root"/build/linux_* 2>/dev/null | head -1) + [[ -z "$b" ]] && b="$root/build" + [[ -f "$b/lib/libOrbbecSDK.so" ]] || return 1 + echo "$b" +} + +emit() { + local root="$1" build="$2" + echo "OB_SDK_ROOT=$root" + echo "OB_SDK_BUILD=$build" + echo "OB_SDK_GENERATED=$root/build/src/generated" +} + +# 1 & 2: reuse whatever is already built. +for candidate in "${OB_SDK_ROOT:-}" \ + "$HOME/OrbbecSDK_v2" \ + "$HOME/OrbbecSDK_v2_jetson"; do + [[ -z "$candidate" ]] && continue + if b=$(find_build "$candidate" 2>/dev/null); then + echo ">> SDK deja bati, reutilise: $candidate" >&2 + emit "$candidate" "$b" + exit 0 + fi +done + +# 3: clone + build from source. +root="${OB_SDK_ROOT:-$HOME/OrbbecSDK_v2}" +echo ">> Aucun SDK bati — clone et compilation (JetPack ${JETPACK}, ${SOC})" >&2 +echo ">> ${SDK_REPO} @ ${SDK_REF} -> ${root}" >&2 + +if [[ ! -d "$root/.git" ]]; then + git clone --depth 1 --branch "$SDK_REF" "$SDK_REPO" "$root" >&2 +fi + +build="$root/build/linux_$(uname -m)" +mkdir -p "$build" +( + cd "$build" + # -DOB_BUILD_EXAMPLES=ON gives the k4a recorder sample used as recorder_bin. + cmake ../.. -DCMAKE_BUILD_TYPE=Release -DOB_BUILD_EXAMPLES=ON >&2 + make -j"$(nproc)" >&2 +) + +if ! b=$(find_build "$root" 2>/dev/null); then + cat >&2 < ./install.sh +EOF + exit 1 +fi +echo ">> SDK bati: $b" >&2 +emit "$root" "$b" diff --git a/scripts/build_server.sh b/scripts/build_server.sh new file mode 100755 index 0000000..e8ecf64 --- /dev/null +++ b/scripts/build_server.sh @@ -0,0 +1,29 @@ +#!/usr/bin/env bash +# Builds ithaca_rgbd_server against the SDK build_sdk.sh located. Prints nothing +# but the binary path on success. +set -euo pipefail + +here=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +repo=$(cd "$here/.." && pwd) + +# SDK locations come from build_sdk.sh (or the environment, when a caller already +# knows them). +if [[ -z "${OB_SDK_ROOT:-}" || -z "${OB_SDK_BUILD:-}" ]]; then + eval "$("$here/build_sdk.sh")" +fi + +build="$repo/server/build" +mkdir -p "$build" +( + cd "$build" + cmake "$repo/server" \ + -DCMAKE_BUILD_TYPE=Release \ + -DOB_SDK_ROOT="$OB_SDK_ROOT" \ + -DOB_SDK_BUILD="$OB_SDK_BUILD" \ + -DOB_SDK_GENERATED="${OB_SDK_GENERATED:-$OB_SDK_ROOT/build/src/generated}" >&2 + make -j"$(nproc)" >&2 +) + +bin="$build/ithaca_rgbd_server" +[[ -x "$bin" ]] || { echo "build failed: $bin absent" >&2; exit 1; } +echo "$bin" diff --git a/scripts/detect_platform.sh b/scripts/detect_platform.sh new file mode 100755 index 0000000..26735a9 --- /dev/null +++ b/scripts/detect_platform.sh @@ -0,0 +1,57 @@ +#!/usr/bin/env bash +# Prints the facts every other script branches on, one KEY=value per line, so a +# caller can `eval "$(detect_platform.sh)"`. +# +# The point of the whole repo is that these differ between Jetsons and nothing is +# assumed: the original Nano is tegra210 on L4T R32 (JetPack 4, Ubuntu 18.04, +# CUDA 10.2), an Orin is tegra234 on R35/R36 (JetPack 5/6). One install procedure, +# values read here, binaries built locally to match. +set -euo pipefail + +# L4T release, e.g. "32.7.1" or "36.4.3". Two sources; the file is the reliable +# one on every generation, the package is the cross-check. +l4t="" +if [[ -f /etc/nv_tegra_release ]]; then + # "# R32 (release), REVISION: 7.1, ..." -> 32.7.1 + rel=$(sed -n 's/^# R\([0-9]*\).*REVISION: \([0-9.]*\).*/\1.\2/p' /etc/nv_tegra_release) + l4t="$rel" +fi +if [[ -z "$l4t" ]] && command -v dpkg-query >/dev/null 2>&1; then + l4t=$(dpkg-query --showformat='${Version}' -W nvidia-l4t-core 2>/dev/null \ + | sed 's/-.*//') || true +fi + +major="${l4t%%.*}" + +# SoC family, from the device tree — the ground truth for what can run here. +soc="unknown" +if [[ -r /proc/device-tree/compatible ]]; then + comp=$(tr '\0' '\n' < /proc/device-tree/compatible) + case "$comp" in + *tegra210*) soc="tegra210" ;; # Nano, TX1 + *tegra186*) soc="tegra186" ;; # TX2 + *tegra194*) soc="tegra194" ;; # Xavier + *tegra234*) soc="tegra234" ;; # Orin + *tegra264*) soc="tegra264" ;; # Thor + esac +fi + +model="unknown" +[[ -r /proc/device-tree/model ]] && model=$(tr -d '\0' < /proc/device-tree/model) + +# JetPack generation follows from the L4T major, and it is what decides toolchain +# and dependency names. +case "$major" in + 32) jetpack="4" ;; + 35) jetpack="5" ;; + 36) jetpack="6" ;; + 38) jetpack="7" ;; + *) jetpack="unknown" ;; +esac + +echo "L4T=${l4t:-unknown}" +echo "L4T_MAJOR=${major:-unknown}" +echo "JETPACK=${jetpack}" +echo "SOC=${soc}" +echo "ARCH=$(uname -m)" +echo "MODEL=${model}" diff --git a/scripts/install_deps.sh b/scripts/install_deps.sh new file mode 100755 index 0000000..b71093c --- /dev/null +++ b/scripts/install_deps.sh @@ -0,0 +1,32 @@ +#!/usr/bin/env bash +# Installs the build- and run-time packages the node needs. Idempotent: apt is a +# no-op for what is already there. +# +# The set is the same across Jetson generations — only the distro release behind +# it changes (18.04 -> 22.04), and apt resolves the right versions itself. No +# GStreamer: the direct preview route encodes no video. +set -euo pipefail + +here=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +eval "$("$here/detect_platform.sh")" + +echo ">> Dependances pour JetPack ${JETPACK} (L4T ${L4T}, ${SOC})" + +# python3 is present on every JetPack image; websockets is the only bridge dep and +# is pure Python, so it works identically everywhere. Prefer apt's package, fall +# back to pip when the distro is too old to carry it. +sudo apt-get update -qq +sudo apt-get install -y --no-install-recommends \ + build-essential cmake pkg-config \ + libusb-1.0-0-dev libudev-dev liblz4-dev libjpeg-turbo8-dev \ + python3 python3-pip + +if ! python3 -c 'import websockets' 2>/dev/null; then + echo ">> websockets absent — installation via pip" + # 8.1 is what the Nano runs and what bridge.py was written against; anything + # >=8 keeps the same asyncio API surface the bridge uses. + python3 -m pip install --user "websockets>=8,<13" || \ + sudo apt-get install -y python3-websockets +fi + +echo ">> Dependances OK" diff --git a/server/CMakeLists.txt b/server/CMakeLists.txt new file mode 100644 index 0000000..bed9669 --- /dev/null +++ b/server/CMakeLists.txt @@ -0,0 +1,93 @@ +# Copyright (c) French Touch Factory. All Rights Reserved. +# +# Builds the one preview server, ithaca_rgbd_server, against an already-built +# OrbbecSDK v2. Deliberately outside the SDK tree: reconfiguring the SDK on an +# older Jetson can break the recorder build that already works, so we link the +# library it produced instead of building inside it. +# +# Portable across Jetson generations by construction: every SDK path comes from a +# cache variable that build_server.sh fills from OB_SDK_ROOT, and nothing here is +# tied to one L4T, one home directory, or one arm64 build folder name. +# +# The retired GStreamer/RTSP target is gone — the direct route encodes no video, +# so GStreamer is not a dependency any more. + +cmake_minimum_required(VERSION 3.10) +project(ithaca_node_server CXX) + +# Where the SDK source (headers) and its build output (libs, generated headers, +# extensions) live. build_server.sh passes these; the defaults match the layout +# build_sdk.sh produces. +set(OB_SDK_ROOT "$ENV{HOME}/OrbbecSDK_v2" CACHE PATH + "OrbbecSDK v2 source root (headers)") +set(OB_SDK_BUILD "" CACHE PATH + "OrbbecSDK v2 build output (lib/, lib/extensions/)") +set(OB_SDK_GENERATED "" CACHE PATH + "OrbbecSDK generated headers, where Export.h lands") + +# The build folder is named per platform (linux_arm64, linux_aarch64, ...), so it +# is discovered rather than assumed when the caller did not pass it. +if(NOT OB_SDK_BUILD) + file(GLOB _ob_build_candidates "${OB_SDK_ROOT}/build/linux_*") + if(_ob_build_candidates) + list(GET _ob_build_candidates 0 OB_SDK_BUILD) + else() + set(OB_SDK_BUILD "${OB_SDK_ROOT}/build") + endif() +endif() +if(NOT OB_SDK_GENERATED) + set(OB_SDK_GENERATED "${OB_SDK_ROOT}/build/src/generated") +endif() + +if(NOT EXISTS "${OB_SDK_ROOT}/include/libobsensor/ObSensor.hpp") + message(FATAL_ERROR "OrbbecSDK headers not found under ${OB_SDK_ROOT}/include " + "— set OB_SDK_ROOT or run scripts/build_sdk.sh first.") +endif() +if(NOT EXISTS "${OB_SDK_BUILD}/lib/libOrbbecSDK.so") + message(FATAL_ERROR "libOrbbecSDK.so not found in ${OB_SDK_BUILD}/lib " + "— build the SDK first (scripts/build_sdk.sh).") +endif() +if(NOT EXISTS "${OB_SDK_GENERATED}/Export.h") + message(FATAL_ERROR "Export.h not found in ${OB_SDK_GENERATED} " + "— the SDK is not fully built yet.") +endif() + +# Directory-scoped rather than target_link_directories(): that one needs CMake +# >= 3.13 and the oldest supported Jetson (Nano, JetPack 4) ships 3.10.2. Must +# precede add_executable(). +link_directories("${OB_SDK_BUILD}/lib") + +add_executable(ithaca_rgbd_server ithaca_rgbd_server.cpp) +set_property(TARGET ithaca_rgbd_server PROPERTY CXX_STANDARD 14) +target_compile_options(ithaca_rgbd_server PRIVATE -Wall -Wextra -O2) + +# Where the depth engine and filter plugins are loaded from at runtime. +target_compile_definitions(ithaca_rgbd_server PRIVATE + OB_EXTENSIONS_DIR="${OB_SDK_BUILD}/lib/extensions") + +# SYSTEM: the SDK's public headers are not warning-clean and this builds -Wall +# -Wextra. ObTypes.h includes Export.h, which the SDK generates into its build +# tree rather than shipping in include/, hence both directories. +target_include_directories(ithaca_rgbd_server SYSTEM PRIVATE + "${OB_SDK_ROOT}/include" + "${OB_SDK_GENERATED}") + +# libjpeg by absolute path, not -ljpeg: the SDK ships its own libjpeg.a in the +# directory link_directories() adds above, built without its SIMD objects — the +# linker would pick it first and fail on jsimd_* symbols. Only C2D alignment +# needs an encoder at all. +find_library(SYSTEM_JPEG NAMES jpeg PATHS + /usr/lib/aarch64-linux-gnu /usr/lib/arm-linux-gnueabihf NO_DEFAULT_PATH) +if(NOT SYSTEM_JPEG) + message(FATAL_ERROR "system libjpeg not found — install libjpeg-turbo8-dev.") +endif() + +target_link_libraries(ithaca_rgbd_server PRIVATE OrbbecSDK lz4 ${SYSTEM_JPEG} pthread) + +# The SDK lives outside any standard search path; without an rpath the binary +# only runs from a shell that already exports LD_LIBRARY_PATH. +set_target_properties(ithaca_rgbd_server PROPERTIES + BUILD_RPATH "${OB_SDK_BUILD}/lib" + INSTALL_RPATH "${OB_SDK_BUILD}/lib") + +install(TARGETS ithaca_rgbd_server RUNTIME DESTINATION bin) diff --git a/server/ithaca_rgbd_server.cpp b/server/ithaca_rgbd_server.cpp new file mode 100644 index 0000000..0cb61a2 --- /dev/null +++ b/server/ithaca_rgbd_server.cpp @@ -0,0 +1,1058 @@ +// Copyright (c) French Touch Factory. All Rights Reserved. +// +// Ithaca RGBD preview server — Jetson. +// +// Serves every camera on this node over one TCP connection: colour as the +// sensor's own JPEG, forwarded untouched, and depth as the sensor's own 16-bit +// millimetres, losslessly compressed. +// +// Why not the H.264/WebRTC path it replaces: +// +// Depth is data, not a picture. Squeezing it into 8 bits for a video codec put +// a floor of 19.6 mm on precision before compression even began, and H.264's +// smoothing then bled object edges into the background — a geometric error, +// not a cosmetic one. Sent as-is it is exact, and LZ4 over separated byte +// planes costs 3.35 ms a frame for a 2.95x ratio (measured, 640x576). +// +// Colour is already JPEG when it leaves the sensor. Decoding it to re-encode +// as H.264 spent NVDEC and NVENC, added a generation of loss, and made the two +// cameras contend for one decoder — which is where their differing latency +// came from. Forwarded as-is it costs nothing. +// +// Colour and depth of one frameset leave in the same message, under the same +// lock. They cannot drift apart. Over WebRTC they were two sessions with two +// jitter buffers and no shared clock, and MediaMTX's WHEP serves only one video +// track per session, so putting them in one was not available either. +// +// And nothing here buffers: no rate-smoothing window, no group of pictures, no +// jitter buffer. When the link falls behind, whole frames are dropped instead, +// which is safe because every frame stands alone. +// +// The cost of all this is bandwidth — roughly 55 Mbit/s per camera binned — so +// it is a local-network path. WebRTC remains the answer for anything leaving the +// site. +// +// Usage: +// ithaca_rgbd_server --map :[,:...] +// +// --map is both the numbering AND the selection: a camera named there is +// served under that index, and a camera absent from a NON-EMPTY map is not +// streamed. An empty (or missing) map serves every camera, numbered by +// enumeration order — what a run by hand relies on. +// +// --map-file names a file holding that same text, re-read once a second, so the +// composition can change WITHOUT restarting: a camera dropped from it has its +// streams stopped, one added has them started, and the others keep running +// throughout. That is how the operator's per-camera checkbox reaches the sensor +// while a preview is live — restarting instead took every camera of the node +// down for about ten seconds. +// +// A SECOND line of that file, if present, is the alignment: "raw" or "compare". +// Those two share their stream configuration, so the node can be told to start +// or stop producing the SDK's own transform at any moment. The other modes are +// fixed for the life of the process, needing a different colour format or a +// different target, and are computed by the viewer anyway. +// +// Stopping the streams is what saves the USB bandwidth, the depth engine and the +// heat. The device handle stays open, so turning the camera back on is immediate; +// Record does not care, since it stops this server before opening its own +// recorders. +// [--depth-mode M] [--color-mode M] [--rate N] +// [--align raw|d2c|d2c-fov|c2d|compare] [--port P] +#include +// Not pulled in by ObSensor.hpp, and this is where the ray table lives. +#include +#include + +#include // jpeglib.h needs FILE declared before it +#include + +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +namespace { + +std::atomic g_running{true}; +void onSignal(int) { g_running = false; } + +uint64_t nowUs() { + return std::chrono::duration_cast( + std::chrono::steady_clock::now().time_since_epoch()).count(); +} + +// Interleaved 16-bit samples compress poorly: the high bytes are smooth and +// repetitive, the low bytes are mostly sensor noise, and mixed together they +// spoil each other's matches. Separated, each compresses on its own terms — +// measured 2.95x against 2.27x, and faster with it. +void splitPlanes(const uint8_t *src, size_t bytes, std::vector &out) { + const size_t n = bytes / 2; + out.resize(bytes); + for(size_t i = 0; i < n; i++) { + out[i] = src[2 * i + 1]; + out[n + i] = src[2 * i]; + } +} + +// Re-encodes a warped colour frame, which only C2D needs: in every other mode +// the colour that arrives is already a JPEG and is forwarded as it is. +// +// One of these per camera, kept alive across frames — jpeg_create_compress +// allocates its tables and its Huffman state, and building them thirty times a +// second per camera would cost more than the encode. +class JpegEncoder { +public: + JpegEncoder() { + info_.err = jpeg_std_error(&err_); + jpeg_create_compress(&info_); + } + + ~JpegEncoder() { + jpeg_destroy_compress(&info_); + if(buffer_ != nullptr) free(buffer_); + } + + JpegEncoder(const JpegEncoder &) = delete; + JpegEncoder &operator=(const JpegEncoder &) = delete; + + /// Encodes packed 24-bit RGB. Returns false and leaves size at 0 on failure. + /// The returned pointer belongs to this encoder and is valid until the next + /// call, which is enough: the frame is sent before the next one is encoded. + bool encode(const uint8_t *rgb, int width, int height, int quality, + const uint8_t *&data, size_t &size) { + data = nullptr; + size = 0; + + // jpeg_mem_dest grows this itself and hands back the grown pointer, so it + // is kept between calls and settles at the largest frame seen. + jpeg_mem_dest(&info_, &buffer_, &capacity_); + + info_.image_width = static_cast(width); + info_.image_height = static_cast(height); + info_.input_components = 3; + info_.in_color_space = JCS_RGB; + jpeg_set_defaults(&info_); + jpeg_set_quality(&info_, quality, TRUE); + + jpeg_start_compress(&info_, TRUE); + const int stride = width * 3; + while(info_.next_scanline < info_.image_height) { + JSAMPROW row = const_cast(rgb + info_.next_scanline * stride); + jpeg_write_scanlines(&info_, &row, 1); + } + jpeg_finish_compress(&info_); + + if(buffer_ == nullptr || capacity_ == 0) return false; + data = buffer_; + size = capacity_; + return true; + } + +private: + jpeg_compress_struct info_{}; + jpeg_error_mgr err_{}; + uint8_t *buffer_ = nullptr; + unsigned long capacity_ = 0; +}; + +bool sendAll(int fd, const void *data, size_t size) { + const uint8_t *p = static_cast(data); + while(size > 0) { + const ssize_t n = ::send(fd, p, size, MSG_NOSIGNAL); + if(n <= 0) return false; + p += n; + size -= static_cast(n); + } + return true; +} + +// Bytes the kernel has still to put on the wire. On a link that cannot keep up +// this is what grows, and with it the delay: TCP queues frames happily until the +// picture is seconds behind. A preview must drop instead. +bool backedUp(int fd, int limitBytes) { + int pending = 0; + if(::ioctl(fd, TIOCOUTQ, &pending) != 0) return false; + return pending > limitBytes; +} + +// A camera's calibration, as it goes on the wire: enough for the receiver to put +// a depth pixel on the colour image itself, and nothing more. +// +// What it holds, and why exactly this: +// +// ray table two floats per depth pixel, the direction that pixel looks in. +// The depth camera's distortion is already solved into it. It has +// to be, because that model cannot be inverted at the far end — +// k1 = 17.4 on this sensor, strong enough that an iterative +// inverse was 10 pixels out at the edges while looking perfect in +// the middle, and the SDK itself refuses to evaluate the corners. +// +// colour intrinsics and the rigid transform. No distortion coefficients: the +// SDK's own projection applies none, and reproducing it exactly — +// 0.000 pixel over 123 points across the whole image at three +// distances — means a plain pinhole. Applying the stored colour +// coefficients, or inverting them, was 9 pixels out. +// +// The receiver then does, per depth pixel: P = ray * millimetres, P' = R P + T, +// and u = P'.x / P'.z * fx + cx. That is a texture read and a few multiplies — +// which is the whole point of sending this instead of transforming here. +#pragma pack(push, 1) +struct WireCalibrationHeader { + int32_t depthWidth, depthHeight; + int32_t colourWidth, colourHeight; + float colourFx, colourFy, colourCx, colourCy; + float rot[9]; // depth -> colour, row major + float trans[3]; // depth -> colour, in millimetres + int32_t rayTableBytes; // LZ4 payload that follows this header + int32_t rayTablePlainBytes; // what it expands to: depth pixels * 2 floats +}; +#pragma pack(pop) + +// The viewers currently connected. A frameset is written to each in turn under +// one lock, so no two cameras can interleave their messages on any of them. +struct Viewers { + std::mutex gate; + std::vector fds; + + void add(int fd) { + std::lock_guard guard(gate); + fds.push_back(fd); + } + + std::vector snapshot() { + std::lock_guard guard(gate); + return fds; + } + + void drop(int fd) { + std::lock_guard guard(gate); + for(size_t i = 0; i < fds.size(); i++) { + if(fds[i] == fd) { + fds.erase(fds.begin() + static_cast(i)); + ::close(fd); + return; + } + } + } + + void closeAll() { + std::lock_guard guard(gate); + for(int fd : fds) ::close(fd); + fds.clear(); + } + + size_t count() { + std::lock_guard guard(gate); + return fds.size(); + } +}; + +#pragma pack(push, 1) +struct Header { + uint32_t magic; // 'ITHD' + uint16_t camera; // camIndex, as the Recorder panel numbers it + uint16_t kind; // 0 = depth (LZ4, split planes), 1 = colour (JPEG), + // 2 = calibration (once per viewer), 3 = the SDK's own + // depth-to-colour result (LZ4, split planes) + uint16_t width; + uint16_t height; + uint32_t payload; // bytes that follow + uint32_t ageUs; // age of the frame when it was sent + uint32_t sequence; // shared by the colour and depth of one frameset +}; +#pragma pack(pop) + +const uint32_t kMagic = 0x44485449; // 'ITHD' little-endian + +// Only C2D re-encodes, and this is the one place a loss is introduced anywhere +// in this server. High enough that the artefacts stay below what the sensor's +// own JPEG already has, since a second generation is what is being paid for. +const int kJpegQuality = 90; + +struct Mode { int width, height; }; + +bool depthMode(const std::string &name, Mode &out) { + if(name == "NFOV_2X2BINNED") { out = {320, 288}; return true; } + if(name == "NFOV_UNBINNED") { out = {640, 576}; return true; } + if(name == "WFOV_2X2BINNED") { out = {512, 512}; return true; } + if(name == "WFOV_UNBINNED") { out = {1024, 1024}; return true; } + return false; +} + +// Which grid the two images are put on before they leave. +// +// Raw each sensor's own view. The two sit a couple of centimetres apart, +// so a given object falls on different pixels in the two images and +// nothing downstream can pair them without the calibration. Costs +// nothing, and the colour is the sensor's own JPEG untouched. +// +// D2C the depth resampled into the colour camera's view, at the colour +// resolution: pixel (x, y) is the same point of the world in both. +// This is k4a_transformation_depth_image_to_color_camera — a mesh warp +// rather than a per-pixel reprojection, which is what keeps holes out +// of the result and resolves occlusions with a z-buffer. The colour +// still leaves untouched; only the depth is transformed, and it grows +// to the colour resolution (4 MB a frame at 1080p, before compression). +// +// D2CFov the same view and the same aspect, kept at the depth sensor's own +// resolution. K4A has no such variant; this one is the v2 SDK's own +// (Align::setMatchTargetResolution). Texture coordinates still line up, +// which is what a shader or a point cloud actually samples with; only +// per-pixel identity is given up, and with it most of the bandwidth. +// +// C2D the colour resampled into the depth camera's view, at the depth +// resolution — k4a_transformation_color_image_to_depth_camera. This is +// the expensive one and unavoidably so: the sensor puts MJPG on the +// USB, so the frame has to be decoded to be warped, and then encoded +// again to be sent. Raw and D2C both avoid touching the colour at all. +// Compare the raw depth as in Raw, plus the SDK's own depth-to-colour result +// alongside it, so a viewer can overlay its own transform on the +// reference and see where the two place their points differently. +// Costs what the transform costs — 22 frames a second instead of 29, +// and a second depth map on the wire. For diagnosis, not capture. +enum class Align { Raw, D2C, D2CFov, C2D, Compare }; + +bool alignMode(const std::string &name, Align &out) { + if(name == "raw") { out = Align::Raw; return true; } + if(name == "d2c") { out = Align::D2C; return true; } + if(name == "d2c-fov") { out = Align::D2CFov; return true; } + if(name == "c2d") { out = Align::C2D; return true; } + if(name == "compare") { out = Align::Compare; return true; } + return false; +} + +// For messages about the mode in force NOW, which is not necessarily the one named +// on the command line. +const char *alignText(Align a) { + switch(a) { + case Align::Raw: return "raw"; + case Align::D2C: return "d2c"; + case Align::D2CFov: return "d2c-fov"; + case Align::C2D: return "c2d"; + case Align::Compare: return "compare"; + } + return "?"; +} + +bool colorMode(const std::string &name, Mode &out) { + if(name == "720p") { out = {1280, 720}; return true; } + if(name == "1080p") { out = {1920, 1080}; return true; } + if(name == "1440p") { out = {2560, 1440}; return true; } + if(name == "1536p") { out = {2048, 1536}; return true; } + if(name == "2160p") { out = {3840, 2160}; return true; } + return false; +} + +// "serial:index,serial:index" — the bridge owns these numbers, and they name the +// slot each stream occupies in the mosaic. +std::map parseMap(const std::string &text) { + std::map out; + size_t start = 0; + while(start < text.size()) { + const size_t comma = text.find(',', start); + const std::string item = text.substr(start, comma - start); + const size_t colon = item.find(':'); + if(colon != std::string::npos) { + out[item.substr(0, colon)] = + static_cast(std::atoi(item.substr(colon + 1).c_str())); + } + if(comma == std::string::npos) break; + start = comma + 1; + } + return out; +} + +// What the node should be doing right now: which cameras, and in which mode. +// +// Read on a timer rather than pushed, so nothing has to be added to the viewer +// protocol and a crashed writer cannot leave the server waiting on a command that +// never comes. Both live in ONE file so a single atomic rename changes them +// together — they can never be read half-updated against each other. +struct Composition { + std::map indices; + Align align = Align::Raw; + bool hasAlign = false; +}; + +Composition readComposition(const std::string &path) { + Composition out; + std::ifstream in(path); + if(!in) return out; + + std::string line; + std::getline(in, line); + out.indices = parseMap(line); + + // Second line optional: a writer that does not know about it simply leaves the + // mode as the command line set it. + if(std::getline(in, line)) { + while(!line.empty() && (line.back() == '\r' || line.back() == ' ')) line.pop_back(); + Align parsed{}; + if(!line.empty() && alignMode(line, parsed)) { + out.align = parsed; + out.hasAlign = true; + } + } + return out; +} + +// One camera the server can turn on and off in place. The frame callback is kept +// here because start() needs it again on every restart of this one pipeline. +struct Cam { + std::string serial; + uint16_t camIndex = 0; + std::shared_ptr pipeline; + std::shared_ptr config; + std::function)> callback; + bool running = false; +}; + +} // namespace + +int main(int argc, char **argv) { + std::string mapText, mapFile, depthName = "NFOV_2X2BINNED", colorName = "720p", + alignName = "none"; + int port = 5020, rate = 30, backlogLimit = 1 << 20; + + for(int i = 1; i < argc; i++) { + const std::string arg = argv[i]; + if(arg == "--map" && i + 1 < argc) mapText = argv[++i]; + else if(arg == "--map-file" && i + 1 < argc) mapFile = argv[++i]; + else if(arg == "--depth-mode" && i + 1 < argc) depthName = argv[++i]; + else if(arg == "--color-mode" && i + 1 < argc) colorName = argv[++i]; + else if(arg == "--align" && i + 1 < argc) alignName = argv[++i]; + else if(arg == "--rate" && i + 1 < argc) rate = std::atoi(argv[++i]); + else if(arg == "--port" && i + 1 < argc) port = std::atoi(argv[++i]); + else if(arg == "--backlog" && i + 1 < argc) backlogLimit = std::atoi(argv[++i]); + else { + std::cerr << "unknown option: " << arg << std::endl; + return 2; + } + } + + Mode depth{}, colour{}; + if(!depthMode(depthName, depth)) { + std::cerr << "unknown depth mode: " << depthName << std::endl; + return 2; + } + if(!colorMode(colorName, colour)) { + std::cerr << "unknown colour mode: " << colorName << std::endl; + return 2; + } + Align align{}; + if(!alignMode(alignName, align)) { + std::cerr << "unknown alignment: " << alignName + << " (expected none, color or color-fov)" << std::endl; + return 2; + } + std::map indices = parseMap(mapText); + // The file wins when it is there and readable: it is the live truth, and --map + // is then only the starting point the bridge passed on the command line. + if(!mapFile.empty()) { + const Composition start = readComposition(mapFile); + if(!start.indices.empty()) indices = start.indices; + if(start.hasAlign) align = start.align; + } + + // Only these two can be switched under way: they share the stream + // configuration, so nothing has to be reopened. Fixed otherwise. + const bool liveSwitch = (align == Align::Raw || align == Align::Compare); + std::atomic liveAlign{static_cast(align)}; + + std::signal(SIGINT, onSignal); + std::signal(SIGTERM, onSignal); + + // ---- listening socket --------------------------------------------------- + const int listenFd = ::socket(AF_INET, SOCK_STREAM, 0); + int yes = 1; + ::setsockopt(listenFd, SOL_SOCKET, SO_REUSEADDR, &yes, sizeof(yes)); + + sockaddr_in addr{}; + addr.sin_family = AF_INET; + addr.sin_addr.s_addr = INADDR_ANY; + addr.sin_port = htons(static_cast(port)); + if(::bind(listenFd, reinterpret_cast(&addr), sizeof(addr)) != 0) { + std::cerr << "cannot bind port " << port << std::endl; + return 1; + } + ::listen(listenFd, 2); + + Viewers viewers; + std::thread accepter([&]() { + while(g_running) { + const int fd = ::accept(listenFd, nullptr, nullptr); + if(fd < 0) break; + // Send at once rather than let Nagle group small writes: every + // millisecond held here is a millisecond of preview delay. + ::setsockopt(fd, IPPROTO_TCP, TCP_NODELAY, &yes, sizeof(yes)); + viewers.add(fd); + std::cout << "Viewer connected (" << viewers.count() << " watching)." + << std::endl; + } + }); + + // ---- cameras ------------------------------------------------------------ + ob::Context::setExtensionsDirectory(OB_EXTENSIONS_DIR); + ob::Context ctx; + auto list = ctx.queryDeviceList(); + if(list->getCount() == 0) { + std::cerr << "no camera found" << std::endl; + // The accepter is still joinable, and letting a joinable thread reach its + // destructor calls std::terminate — an abort on the ordinary path where + // nothing is plugged in. + ::shutdown(listenFd, SHUT_RDWR); + ::close(listenFd); + accepter.detach(); + return 1; + } + + std::mutex sendLock; + std::vector cams; + std::vector>> dropCounters; + + for(uint32_t i = 0; i < list->getCount(); i++) { + auto dev = list->getDevice(i); + const std::string serial = list->getSerialNumber(i); + + // Fall back to enumeration order only if the bridge did not say: its + // numbering is what the operator sees and edits. + uint16_t camIndex = static_cast(i); + auto known = indices.find(serial); + if(known != indices.end()) camIndex = known->second; + + // A non-empty map is the list to STREAM. The pipeline is still built for + // the others, so that a camera ticked back on starts at once instead of + // paying for a device open; what an unticked camera must not do is stream, + // which is where the USB bandwidth, the depth engine and the heat go. + const bool wanted = indices.empty() || known != indices.end(); + + auto pipeline = std::make_shared(dev); + auto config = std::make_shared(); + std::shared_ptr colourProfile, depthProfile; + try { + depthProfile = pipeline->getStreamProfileList(OB_SENSOR_DEPTH) + ->getVideoStreamProfile(depth.width, depth.height, + OB_FORMAT_Y16, rate); + config->enableStream(depthProfile); + // MJPG everywhere except C2D: the sensor puts JPEG on the USB and + // that is exactly what we want to forward. C2D has to warp the + // colour, so it needs pixels, and asking the SDK for RGB is what + // makes it decode — there is no way to warp a JPEG. + const OBFormat colourFormat = (align == Align::C2D) ? OB_FORMAT_RGB + : OB_FORMAT_MJPG; + colourProfile = pipeline->getStreamProfileList(OB_SENSOR_COLOR) + ->getVideoStreamProfile(colour.width, colour.height, + colourFormat, rate); + config->enableStream(colourProfile); + } + catch(const ob::Error &e) { + std::cerr << "camera " << serial << ": " << e.what() << std::endl; + return 1; + } + + // One filter per camera, not one shared: it holds the calibration of the + // pair it is aligning, and a filter is stateful. + // Built for Raw too when the mode can change: a server started in Raw must + // be able to produce the reference a moment later, and the filter is only + // ever RUN when the live mode asks for it (see the callback). + const Align filterFor = liveSwitch ? Align::Compare : align; + + std::shared_ptr aligner; + if(filterFor != Align::Raw) { + try { + aligner = std::make_shared( + filterFor == Align::C2D ? OB_STREAM_DEPTH : OB_STREAM_COLOR); + aligner->setMatchTargetResolution(filterFor != Align::D2CFov); + // Told the target profile explicitly rather than left to infer it + // from the first frame of the other stream: the intrinsics and + // extrinsics are what the transform needs, and this way the very + // first frame is aligned like all the others. + aligner->setAlignToStreamProfile( + filterFor == Align::C2D ? depthProfile : colourProfile); + } + catch(const ob::Error &e) { + std::cerr << "camera " << serial << ": cannot set up the " + << alignText(filterFor) << " transform: " << e.what() + << std::endl; + return 1; + } + } + + // Read once, here: it depends on the profiles just enabled, and the + // intrinsics are resolution-dependent — asking before configuring would + // return numbers for the wrong image size. + // + // Built as the bytes that go on the wire, header then compressed table, + // so the sending path has nothing left to assemble. + auto calibration = std::make_shared>(); + try { + OBCalibrationParam param = pipeline->getCalibrationParam(config); + + const size_t rays = static_cast(depth.width) * depth.height * 2; + std::vector table(rays); + // Sized from the profile rather than queried: passing a null buffer + // to ask for the size throws. Two floats per depth pixel is what the + // table is by definition. The SDK reports the count back in floats. + uint32_t reported = static_cast(rays); + OBXYTables tables{}; + if(!ob::CoordinateTransformHelper::transformationInitXYTables( + param, OB_SENSOR_DEPTH, table.data(), &reported, &tables)) { + throw std::runtime_error("the SDK would not build the ray table"); + } + + // The two tables come back as separate arrays; interleaved they are + // one texture on the far end, read in a single fetch. + std::vector interleaved(rays); + const size_t pixels = static_cast(tables.width) * tables.height; + for(size_t i = 0; i < pixels; i++) { + interleaved[2 * i] = tables.xTable[i]; + interleaved[2 * i + 1] = tables.yTable[i]; + } + + const int plainBytes = static_cast(interleaved.size() * sizeof(float)); + std::vector packed( + static_cast(LZ4_compressBound(plainBytes))); + const int packedBytes = LZ4_compress_default( + reinterpret_cast(interleaved.data()), + reinterpret_cast(packed.data()), + plainBytes, static_cast(packed.size())); + if(packedBytes <= 0) throw std::runtime_error("could not compress the ray table"); + + WireCalibrationHeader head{}; + head.depthWidth = tables.width; + head.depthHeight = tables.height; + head.colourWidth = param.intrinsics[OB_SENSOR_COLOR].width; + head.colourHeight = param.intrinsics[OB_SENSOR_COLOR].height; + head.colourFx = param.intrinsics[OB_SENSOR_COLOR].fx; + head.colourFy = param.intrinsics[OB_SENSOR_COLOR].fy; + head.colourCx = param.intrinsics[OB_SENSOR_COLOR].cx; + head.colourCy = param.intrinsics[OB_SENSOR_COLOR].cy; + const OBExtrinsic &d2c = param.extrinsics[OB_SENSOR_DEPTH][OB_SENSOR_COLOR]; + std::memcpy(head.rot, d2c.rot, sizeof(head.rot)); + std::memcpy(head.trans, d2c.trans, sizeof(head.trans)); + head.rayTableBytes = packedBytes; + head.rayTablePlainBytes = plainBytes; + + calibration->resize(sizeof(head) + static_cast(packedBytes)); + std::memcpy(calibration->data(), &head, sizeof(head)); + std::memcpy(calibration->data() + sizeof(head), packed.data(), + static_cast(packedBytes)); + + std::cout << " camera " << camIndex << " calibration: ray table " + << tables.width << "x" << tables.height << ", " + << (plainBytes / 1024) << " Ko -> " << (packedBytes / 1024) + << " Ko" << std::endl; + } + catch(const std::exception &e) { + // Not fatal: without it a receiver simply cannot align anything + // itself, which is the state everything was in until now. + calibration->clear(); + std::cerr << "camera " << serial << ": no calibration available: " + << e.what() << std::endl; + } + + // Which sockets already have it. Erased when a viewer goes, so a + // reconnection on the same descriptor number is sent it again. + auto calibrationSent = std::make_shared>(); + + // Per-camera buffers: the callbacks run on their own threads and would + // otherwise overwrite each other between the split and the compress. + auto planes = std::make_shared>(); + auto comp = std::make_shared>(); + // A second pair, for the reference map in Compare: the first is still + // holding the frame being sent. + auto refPlanes = std::make_shared>(); + auto refComp = std::make_shared>(); + auto counter = std::make_shared(0); + auto dropped = std::make_shared>(0); + + auto warned = std::make_shared(false); + // The size the depth comes out at is only known once a frame has been + // through the transform, so it is announced from there rather than guessed. + auto announced = std::make_shared(false); + // Built only where it is used: in Raw and D2C no colour byte is ever + // re-encoded, so there is nothing for it to do. + auto encoder = (align == Align::C2D) ? std::make_shared() : nullptr; + // Held rather than handed straight to start(): the same callback is given + // again every time this one camera is turned back on. + std::function)> callback = + [&, planes, comp, counter, dropped, camIndex, aligner, + warned, announced, encoder, calibration, + calibrationSent, refPlanes, refComp]( + std::shared_ptr fs) { + if(!g_running || fs == nullptr) return; + + // Read once for this whole frameset: it can change between framesets, + // and a frame decided half one way and half the other would be worse + // than either answer. + const Align mode = static_cast(liveAlign.load()); + + // Who gets this frameset: everyone connected whose link is keeping + // up. Decided once, here, and never revisited while the messages are + // being written — a viewer skipped halfway through would be left + // reading a payload against the wrong header. + std::vector targets; + for(int fd : viewers.snapshot()) { + if(backedUp(fd, backlogLimit)) { // behind: drop, do not queue + dropped->fetch_add(1); + continue; + } + targets.push_back(fd); + } + if(targets.empty()) return; // nobody watching, or all behind + + // The calibration goes to a viewer once, before any frame it could + // be needed for. Tracked per camera and per socket, so a new viewer + // gets it without the others being sent it again. + if(!calibration->empty()) { + Header kh{}; + kh.magic = kMagic; + kh.camera = camIndex; + kh.kind = 2; + kh.payload = static_cast(calibration->size()); + + std::lock_guard guard(sendLock); + for(int fd : targets) { + if(calibrationSent->count(fd) != 0) continue; + if(sendAll(fd, &kh, sizeof(kh)) + && sendAll(fd, calibration->data(), calibration->size())) { + calibrationSent->insert(fd); + } + } + } + + const uint64_t captured = nowUs(); + + // Which frameset each half is taken from is the whole difference + // between the modes. In Raw there is no transform at all. In D2C only + // the depth is transformed, and the colour is still the sensor's own + // JPEG, forwarded untouched. In C2D it is the colour that moves, so + // both halves come from the result — and that colour is pixels now, + // which is why it has to be encoded again below. + // + // process() runs here, on the thread the camera delivers on. Handing it + // to the thread inside the filter instead was measured and changed + // nothing: 22 frames a second in D2C either way. What limits it is the + // transform, not this thread, and the filter also discards silently + // when its queue is full, which cost the age its meaning. + std::shared_ptr source = fs; + std::shared_ptr aligned; + if(aligner != nullptr && mode != Align::Raw) { + try { + auto result = aligner->process(fs); + if(result == nullptr) return; + aligned = result->as(); + source = aligned; + } + catch(const ob::Error &e) { + if(!*warned) { // once, not thirty times a second + *warned = true; + std::cerr << "camera " << camIndex << ": the " + << alignText(mode) << " transform failed: " + << e.what() << std::endl; + } + return; + } + } + + // In Compare the frames sent are the camera's own; the transform's + // result travels beside them rather than replacing them. + if(mode == Align::Compare) source = fs; + + auto depthFrame = source->getFrame(OB_FRAME_DEPTH); + // Depth from the result, colour from whichever frameset actually holds + // it: in D2C the transform returns the depth alone about a third of the + // time, and taking the colour from there cost one camera almost all of + // its colour frames. Only C2D has a transformed colour to take. + auto colourFrame = (mode == Align::C2D ? source : fs)->getFrame(OB_FRAME_COLOR); + if(depthFrame == nullptr) return; + + const uint32_t stamp = (*counter)++; + auto df = depthFrame->as(); + const size_t bytes = df->getDataSize(); + + if(!*announced) { + *announced = true; + std::cout << " camera " << camIndex << " depth " << df->getWidth() + << "x" << df->getHeight() << ", " << (bytes / 1024) + << " Ko per frame before compression" << std::endl; + } + + // The protocol says millimetres, and the receiver reads the 16-bit + // values as such. The transform is not supposed to change that, but a + // silent change of unit would look like a room the wrong size rather + // than like an error, so it is worth one line to notice. + if(!*warned && df->getValueScale() != 1.0f) { + *warned = true; + std::cerr << "camera " << camIndex << ": depth value scale is " + << df->getValueScale() << ", not 1 — the receiver reads " + << "millimetres and will be wrong by that factor." + << std::endl; + } + + splitPlanes(static_cast(df->getData()), bytes, *planes); + comp->resize(static_cast(LZ4_compressBound(static_cast(bytes)))); + const int packed = LZ4_compress_default( + reinterpret_cast(planes->data()), + reinterpret_cast(comp->data()), + static_cast(bytes), static_cast(comp->size())); + if(packed <= 0) return; + + Header hd{}; + hd.magic = kMagic; + hd.camera = camIndex; + hd.kind = 0; + hd.width = static_cast(df->getWidth()); + hd.height = static_cast(df->getHeight()); + hd.payload = static_cast(packed); + hd.ageUs = static_cast(nowUs() - captured); + hd.sequence = stamp; + + // The colour to send, and its size. In every mode but C2D these are + // the sensor's own JPEG bytes; the receiver cannot tell the difference + // and does not need to, which is what keeps the protocol the same + // across all three modes. + const uint8_t *colourData = nullptr; + size_t colourSize = 0; + int colourW = 0, colourH = 0; + + if(colourFrame != nullptr) { + auto cf = colourFrame->as(); + colourW = static_cast(cf->getWidth()); + colourH = static_cast(cf->getHeight()); + + if(align == Align::C2D) { + if(cf->getFormat() != OB_FORMAT_RGB) { + if(!*warned) { + *warned = true; + std::cerr << "camera " << camIndex << ": the transform " + << "returned colour in an unexpected format; " + << "sending depth only." << std::endl; + } + } + else if(!encoder->encode(static_cast(cf->getData()), + colourW, colourH, kJpegQuality, + colourData, colourSize)) { + colourData = nullptr; + } + } + else { + colourData = static_cast(cf->getData()); + colourSize = cf->getDataSize(); + } + } + + // The SDK's own result, in Compare only. Same compression as the + // depth: split planes then LZ4, so the receiver has one routine. + const uint8_t *refData = nullptr; + size_t refSize = 0; + int refW = 0, refH = 0; + if(mode == Align::Compare && aligned != nullptr) { + auto refFrame = aligned->getFrame(OB_FRAME_DEPTH); + if(refFrame != nullptr) { + auto rf = refFrame->as(); + const size_t refBytes = rf->getDataSize(); + refW = static_cast(rf->getWidth()); + refH = static_cast(rf->getHeight()); + splitPlanes(static_cast(rf->getData()), + refBytes, *refPlanes); + refComp->resize(static_cast( + LZ4_compressBound(static_cast(refBytes)))); + const int packedRef = LZ4_compress_default( + reinterpret_cast(refPlanes->data()), + reinterpret_cast(refComp->data()), + static_cast(refBytes), static_cast(refComp->size())); + if(packedRef > 0) { + refData = refComp->data(); + refSize = static_cast(packedRef); + } + } + } + + Header ch = hd; + ch.kind = 1; + ch.width = static_cast(colourW); + ch.height = static_cast(colourH); + ch.payload = static_cast(colourSize); + + // One lock for the whole frameset: two cameras must never interleave + // their messages, since a reader takes a header then exactly its + // payload. Sending depth and colour together under it is also what + // makes the pair synchronised by construction rather than by + // negotiation. + std::vector lost; + { + std::lock_guard guard(sendLock); + for(int fd : targets) { + hd.ageUs = static_cast(nowUs() - captured); + bool ok = sendAll(fd, &hd, sizeof(hd)) + && sendAll(fd, comp->data(), static_cast(packed)); + if(ok && colourData != nullptr) { + ch.ageUs = static_cast(nowUs() - captured); + ok = sendAll(fd, &ch, sizeof(ch)) + && sendAll(fd, colourData, colourSize); + } + if(ok && refData != nullptr) { + // Under the same lock as the pair, so all three belong to + // one frameset on the wire and cannot be interleaved. + Header rh = hd; + rh.kind = 3; + rh.width = static_cast(refW); + rh.height = static_cast(refH); + rh.payload = static_cast(refSize); + rh.ageUs = static_cast(nowUs() - captured); + ok = sendAll(fd, &rh, sizeof(rh)) + && sendAll(fd, refData, refSize); + } + if(!ok) lost.push_back(fd); + } + } + + for(int fd : lost) { + calibrationSent->erase(fd); + viewers.drop(fd); + std::cout << "Viewer gone (" << viewers.count() << " watching)." + << std::endl; + } + }; + + Cam cam; + cam.serial = serial; + cam.camIndex = camIndex; + cam.pipeline = pipeline; + cam.config = config; + cam.callback = callback; + if(wanted) { + pipeline->start(config, callback); + cam.running = true; + } + cams.push_back(cam); + dropCounters.push_back(dropped); + std::cout << " camera " << camIndex << " : " << serial + << (wanted ? "" : " (built, not streaming)") << std::endl; + } + + size_t streaming = 0; + for(const auto &c : cams) if(c.running) streaming++; + + // With a composition file there is nothing wrong with starting empty — it can + // be filled a moment later. Without one, an empty selection is a mistake worth + // failing on rather than holding the port and sending nothing. + if(streaming == 0 && mapFile.empty()) { + std::cerr << "nothing to serve: every camera is excluded by --map" + << std::endl; + // Same teardown as the no-camera path above, and for the same reason. + ::shutdown(listenFd, SHUT_RDWR); + ::close(listenFd); + accepter.detach(); + return 1; + } + + std::cout << "Serving " << streaming << " camera(s) on port " << port + << " — colour " << colour.width << "x" << colour.height + << " JPEG, depth " << depth.width << "x" << depth.height + << " 16-bit lossless, " << rate << " fps, alignment " << alignName + << std::endl; + + uint32_t reported = 0; + int ticks = 0; + while(g_running) { + std::this_thread::sleep_for(std::chrono::seconds(1)); + ticks++; + + // Bring the running set in line with the file. A camera leaving it is + // stopped and the others carry on: that is the whole reason this exists + // rather than a restart. + if(!mapFile.empty()) { + const Composition now = readComposition(mapFile); + const std::map &wantedNow = now.indices; + + // The mode first, and independently of the camera list: it costs + // nothing but a flag, and the next frameset picks it up. + if(now.hasAlign && liveSwitch) { + const int wantMode = static_cast(now.align); + if(wantMode != liveAlign.load() + && (now.align == Align::Raw || now.align == Align::Compare)) { + liveAlign.store(wantMode); + std::cout << "alignment now " << alignText(now.align) + << " (no restart)." << std::endl; + } + } + + // An unreadable or empty file is far more likely to be a half-written + // one than a real request to stop everything, and acting on it would be + // a spectacular way to lose a take. So it is ignored until it reads. + if(!wantedNow.empty()) { + for(auto &cam : cams) { + const bool want = wantedNow.find(cam.serial) != wantedNow.end(); + if(want == cam.running) continue; + try { + if(want) { + cam.pipeline->start(cam.config, cam.callback); + cam.running = true; + std::cout << "camera " << cam.camIndex << " (" << cam.serial + << ") streaming again." << std::endl; + } + else { + cam.pipeline->stop(); + cam.running = false; + std::cout << "camera " << cam.camIndex << " (" << cam.serial + << ") stopped, it left the composition." + << std::endl; + } + } + catch(const ob::Error &e) { + // Leaves the flag alone so the next pass tries again, and + // says which camera: one sensor refusing to stop must not + // take the others down with it. + std::cerr << "camera " << cam.camIndex << ": " + << (want ? "start" : "stop") << " failed: " + << e.what() << std::endl; + } + } + } + } + + if(ticks % 5 != 0) continue; + uint32_t total = 0; + for(auto &d : dropCounters) total += d->load(); + if(total != reported) { + std::cout << " dropped " << (total - reported) + << " frame(s): the link is behind, or the transform is" + << std::endl; + reported = total; + } + } + + // Stopping the pipelines can block in the SDK while it releases its EGL + // contexts, so shut the socket first: a viewer sees the end immediately + // instead of waiting on a teardown it has no interest in. + ::shutdown(listenFd, SHUT_RDWR); + ::close(listenFd); + accepter.detach(); + viewers.closeAll(); + + for(auto &cam : cams) if(cam.running) cam.pipeline->stop(); + std::cout << "Stopped." << std::endl; + return 0; +} diff --git a/systemd/ithaca-bridge.service.in b/systemd/ithaca-bridge.service.in new file mode 100644 index 0000000..c8c2e97 --- /dev/null +++ b/systemd/ithaca-bridge.service.in @@ -0,0 +1,21 @@ +[Unit] +Description=Ithaca WebSocket bridge (Jetson node) +Documentation=file:@REPO@/bridge/bridge.py +After=network.target usbfs-memory.service +Wants=network.target + +[Service] +Type=simple +User=@USER@ +WorkingDirectory=@REPO@/bridge +ExecStart=/usr/bin/python3 -u @REPO@/bridge/bridge.py @REPO@/bridge/config.json +Restart=always +RestartSec=5 +KillSignal=SIGINT +TimeoutStopSec=20 +StandardOutput=journal +StandardError=journal +SyslogIdentifier=ithaca-bridge + +[Install] +WantedBy=multi-user.target diff --git a/systemd/usbfs-memory.service b/systemd/usbfs-memory.service new file mode 100644 index 0000000..419b37d --- /dev/null +++ b/systemd/usbfs-memory.service @@ -0,0 +1,12 @@ +[Unit] +Description=Raise usbfs_memory_mb for Orbbec depth cameras +After=sysinit.target +Before=multi-user.target + +[Service] +Type=oneshot +RemainAfterExit=yes +ExecStart=/bin/sh -c 'echo 1000 > /sys/module/usbcore/parameters/usbfs_memory_mb' + +[Install] +WantedBy=multi-user.target diff --git a/udev/99-obsensor-libusb.rules b/udev/99-obsensor-libusb.rules new file mode 100644 index 0000000..598b95f --- /dev/null +++ b/udev/99-obsensor-libusb.rules @@ -0,0 +1,111 @@ +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0501", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Bootloader Device" + +# UVC Modules +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0635", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Femto" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0638", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Femto-w" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0668", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Femto-live" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0636", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Astra+" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0637", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Astra+s" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0536", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Astra+_rgb" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0537", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Astra+s_rgb" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0669", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Femto-mega" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="066b", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Femto Bolt" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0660", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Astra 2" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0670", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 2" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0671", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 2 XL" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0673", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 2 L" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0675", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 2 VL" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0800", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 335" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0801", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 330" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0802", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini dm330" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0803", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 336" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0804", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 335L" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0805", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 330L" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0806", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini dm330L" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0807", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 336L" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="080b", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 335Lg" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="080d", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 336Lg" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="080e", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 335Le" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0810", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 336Le" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0674", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Orbbec Gemini 2 I" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="0701", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Dabai DCL" +SUBSYSTEMS=="usb", ATTRS{idVendor}=="2bc5", ATTRS{idProduct}=="069d", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="Astra Pro2" + +# OpenNI Modules +SUBSYSTEM=="usb", ATTR{idProduct}=="0401", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra" +SUBSYSTEM=="usb", ATTR{idProduct}=="0402", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra_s" +SUBSYSTEM=="usb", ATTR{idProduct}=="0403", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra_pro" +SUBSYSTEM=="usb", ATTR{idProduct}=="0404", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra_mini" +SUBSYSTEM=="usb", ATTR{idProduct}=="0407", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra_mini_s" +SUBSYSTEM=="usb", ATTR{idProduct}=="0601", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra_NH_GLST" +SUBSYSTEM=="usb", ATTR{idProduct}=="060b", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="deeyea" +SUBSYSTEM=="usb", ATTR{idProduct}=="050b", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="deeyea_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="060e", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="petrel" +SUBSYSTEM=="usb", ATTR{idProduct}=="050e", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="petrel_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="060f", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astro_pro_plus" +SUBSYSTEM=="usb", ATTR{idProduct}=="050f", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astro_pro_plus_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0610", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="bus_cl" +SUBSYSTEM=="usb", ATTR{idProduct}=="0603", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="atlas" +SUBSYSTEM=="usb", ATTR{idProduct}=="0510", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="atlas_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0614", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="gemini" +SUBSYSTEM=="usb", ATTR{idProduct}=="0511", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="gemini_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0616", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="u1s" +SUBSYSTEM=="usb", ATTR{idProduct}=="0516", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="u1s_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0617", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="projector" +SUBSYSTEM=="usb", ATTR{idProduct}=="0517", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="projector_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0618", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="butterfly" +SUBSYSTEM=="usb", ATTR{idProduct}=="0518", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="butterfly_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="061b", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="pavo" +SUBSYSTEM=="usb", ATTR{idProduct}=="051b", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="pavo_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="062b", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="petrel_pro" +SUBSYSTEM=="usb", ATTR{idProduct}=="052b", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="petrel_pro_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="062c", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="petrel_plus" +SUBSYSTEM=="usb", ATTR{idProduct}=="052c", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="petrel_plus_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="062d", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="pictor" +SUBSYSTEM=="usb", ATTR{idProduct}=="0632", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra+" +SUBSYSTEM=="usb", ATTR{idProduct}=="0532", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra+rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0633", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra+s" +SUBSYSTEM=="usb", ATTR{idProduct}=="0533", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra+s_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0634", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_petal_b" +SUBSYSTEM=="usb", ATTR{idProduct}=="0534", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_petal_b_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0635", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="femto" +SUBSYSTEM=="usb", ATTR{idProduct}=="0636", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="jarvis" +SUBSYSTEM=="usb", ATTR{idProduct}=="0536", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="jarvis_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0637", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="jarvis+s" +SUBSYSTEM=="usb", ATTR{idProduct}=="0537", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="jarvis+s_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0638", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="femto_w" +SUBSYSTEM=="usb", ATTR{idProduct}=="0639", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="argus" +SUBSYSTEM=="usb", ATTR{idProduct}=="0539", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="argus_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="063a", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="u2" +SUBSYSTEM=="usb", ATTR{idProduct}=="0650", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astradepth" +SUBSYSTEM=="usb", ATTR{idProduct}=="0651", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astradepth" +SUBSYSTEM=="usb", ATTR{idProduct}=="0654", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="qi_long_zhu" +SUBSYSTEM=="usb", ATTR{idProduct}=="0554", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="qi_long_zhu_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0655", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_plus" +SUBSYSTEM=="usb", ATTR{idProduct}=="0656", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_mini" +SUBSYSTEM=="usb", ATTR{idProduct}=="0657", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dc1" +SUBSYSTEM=="usb", ATTR{idProduct}=="0557", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dc1_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0658", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_d1" +SUBSYSTEM=="usb", ATTR{idProduct}=="0657", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dc1" +SUBSYSTEM=="usb", ATTR{idProduct}=="0557", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dc1_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="0659", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dcw" +SUBSYSTEM=="usb", ATTR{idProduct}=="0559", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dcw_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="065a", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dw" +SUBSYSTEM=="usb", ATTR{idProduct}=="065b", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra_mini_pro" +SUBSYSTEM=="usb", ATTR{idProduct}=="065c", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="gemini_e" +SUBSYSTEM=="usb", ATTR{idProduct}=="055c", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="gemini_e_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="065d", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="gemini_e_lite" +SUBSYSTEM=="usb", ATTR{idProduct}=="065e", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="astra_mini_s_pro" +SUBSYSTEM=="usb", ATTR{idProduct}=="0698", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="j1" +SUBSYSTEM=="usb", ATTR{idProduct}=="069c", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="TB2201" +SUBSYSTEM=="usb", ATTR{idProduct}=="06a0", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dcw2" +SUBSYSTEM=="usb", ATTR{idProduct}=="0561", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dcw2_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="069f", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_dw2" +SUBSYSTEM=="usb", ATTR{idProduct}=="069a", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_max" +SUBSYSTEM=="usb", ATTR{idProduct}=="069e", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_max_pro" +SUBSYSTEM=="usb", ATTR{idProduct}=="0560", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_max_pro_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="06aa", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_gemini_uw" +SUBSYSTEM=="usb", ATTR{idProduct}=="05aa", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="dabai_gemini_uw_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="06a6", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="gemini_ew" +SUBSYSTEM=="usb", ATTR{idProduct}=="05a6", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="gemini_ew_rgb" +SUBSYSTEM=="usb", ATTR{idProduct}=="06a7", ATTR{idVendor}=="2bc5", MODE:="0666", OWNER:="root", GROUP:="video", SYMLINK+="gemini_ew_lite"