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