chore: snapshot livekit-cameras after serial pin + speechbrain lengths fix

This commit is contained in:
Hermes
2026-08-28 05:15:00 +00:00
commit dc46cbef04
14 changed files with 1369 additions and 0 deletions
+32
View File
@@ -0,0 +1,32 @@
# LiveKit server connection
LIVEKIT_URL=ws://localhost:7880
LIVEKIT_API_KEY=APIKey_xxx
LIVEKIT_API_SECRET=APIsecret_xxx
# Room / naming
LIVEKIT_ROOM=cameras
PARTICIPANT_PREFIX=cam
# Streams
# "all" = every discovered camera, "1,3,5" or "0-19" = subset
CAMERAS=all
VIDEO_WIDTH=1280
VIDEO_HEIGHT=720
VIDEO_FPS=30
VIDEO_MIN_FPS=5
# Audio
AUDIO_SAMPLE_RATE=16000
AUDIO_CHANNELS=1
# Audio cleanup (speechbrain)
# mode: auto (default: enhanced if the model is cached, otherwise passthrough)
# force (always enhance), off (always passthrough)
ENHANCE_MODE=auto
# speechbrain model repo (must be 16 kHz; metricGAN-plus runs ~20-40ms/frame on CPU)
ENHANCE_MODEL=speechbrain/metricgan-plus-voicebank
ENHANCE_CHUNK_S=0.2
# Misc
PUBLISH_TIMEOUT_S=60
LOG_LEVEL=INFO
+8
View File
@@ -0,0 +1,8 @@
.venv/
__pycache__/
*.pyc
.env
*.wav
.hermes/
_health_discover.py
_test_enhance.py
+79
View File
@@ -0,0 +1,79 @@
# livekit-cameras
Publishes 20 Logitech C920 webcam streams (video + cleaned audio) to a
LiveKit room.
Each camera becomes one LiveKit participant (`cam-01` .. `cam-20`) that
publishes:
- a camera video track (1280x720 @ 30 fps by default, MJPG from V4L2)
- a microphone audio track (16 kHz mono, cleaned with speechbrain)
## Layout
```
run.py orchestrator: discovery + per-camera worker processes
run_worker.py one camera: ffmpeg capture -> LiveKit track publish
discovery.py serial-based camera <-> ALSA card discovery
audio_cleanup.py speechbrain SpectralMaskEnhancement (16 kHz) + fallback
tokens.py LiveKit JWT generation
config.py .env / environment configuration
prefetch_model.py one-time speechbrain model download
generate_udev_rules.py (re)generate udev pinning + slots.json from live state
slots.json stable slot map: cam-NN -> USB serial, /dev/camNN, ALSA card
```
## Setup
```bash
uv venv --python 3.12
uv pip install -p .venv/bin/python -r pyproject.toml # or: uv sync
cp .env.example .env # then fill in LIVEKIT_URL / API_KEY / API_SECRET
.venv/bin/python prefetch_model.py # download speechbrain model once
```
## Stable camera identity (cam-01 is always cam-01)
Each C920 has a unique built-in USB serial — that is the camera's UUID.
udev rules (`/etc/udev/rules.d/99-c920-pin.rules`, generated from the live
state) bind every slot to a serial, so identity survives reboots, hub
reordering, and any connect sequence:
- video: `/dev/cam01` .. `/dev/cam20` (primary UVC node; `.meta` = metadata)
- audio: fixed ALSA card per slot, `hw:2,0` .. `hw:21,0` (SOUND_CARD_INDEX)
`discovery.py` assigns identities from `slots.json` (serial lookup, with a
port-based fallback for cameras not yet in the map). To re-map slots (e.g.
after physically swapping cameras), regenerate:
```bash
.venv/bin/python generate_udev_rules.py --write
sudo udevadm control --reload-rules
sudo udevadm trigger --subsystem-match=video4linux --action=add
sudo udevadm trigger --subsystem-match=sound --action=add
```
## Run
```bash
.venv/bin/python run.py # all 20 cameras
.venv/bin/python run.py --cameras 1,5 # subset
.venv/bin/python run.py --once # status report after timeout
```
## Audio cleanup
The mic audio is cleaned per chunk (default 200 ms) by
`speechbrain/metricgan-plus-voicebank` (16 kHz, CPU). One worker process per
camera loads the model once and reuses it; on this 28-core box all 20 stay
real-time. `ENHANCE_MODE=off` (or a failed model load with `auto`) falls back
to a lightweight high-pass + soft-clip DSP path so streams still work.
## Notes
- Video uses MJPG directly from the UVC interface (`ffmpeg -f v4l2
-input_format mjpeg`) — cheap, and converts to BGRA for LiveKit.
- Audio: ALSA mics on these C920s deliver 16 kHz **stereo** only; the worker
downmixes to mono for LiveKit.
- Camera <-> card pairing is resolved by matching USB port paths under /sys,
so each mic belongs to its own camera.
- Workers auto-restart (up to 10 times) if a camera drops.
+142
View File
@@ -0,0 +1,142 @@
"""Audio cleanup for camera mics, using speechbrain.
Pipeline per audio frame (16 kHz mono int16, ~100 ms chunks from the mic):
int16 -> float32 [-1, 1] -> speechbrain enhancement -> int16
Two implementations, selected by availability and the ENHANCE_MODE config:
* "model" - speechbrain SpectralMaskEnhancement (metricGAN-plus). Real
ML-based noise suppression. Runs on CPU; ~20-40 ms per 100 ms
chunk on a modern core, so latency stays well under a frame
budget if workers are spread across cores.
* "gain" - lightweight DSP fallback (soft-clip + high-pass) used when the
model can't be loaded (ENHANCE_MODE=off / auto + no model).
The enhancer is created once per worker process and reused across frames.
"""
from __future__ import annotations
import logging
from typing import Optional
import numpy as np
log = logging.getLogger("audio.cleanup")
SAMPLE_RATE = 16000
class AudioCleaner:
"""Stateful audio cleaner. Call .process(int16 mono) -> int16 mono."""
def __init__(self, mode: str = "auto",
model_source: str = "speechbrain/metricgan-plus-voicebank",
sample_rate: int = SAMPLE_RATE):
self.mode = mode
self.sample_rate = sample_rate
self._enhancer = None
self._torch = None
self.backend = "off"
if mode == "force":
self._load_model(model_source)
if self.backend == "off":
raise RuntimeError(
f"ENHANCE_MODE=force but failed to load model "
f"'{model_source}'")
elif mode == "auto":
self._load_model(model_source)
# mode == "off": passthrough + basic gain shaping
def _load_model(self, source: str) -> None:
try:
import torch
from speechbrain.inference.enhancement import (
SpectralMaskEnhancement,
)
except Exception as e: # noqa: BLE001
log.info("speechbrain unavailable (%s); using fallback", e)
self._enhancer = None
return
try:
log.info("loading speechbrain enhancement model %s ...", source)
self._torch = torch
model = SpectralMaskEnhancement.from_hparams(
source=source,
savedir="/root/.cache/speechbrain-enhancement",
)
model.mods.enhance_model.eval()
self._enhancer = model
self.backend = "model"
log.info("speechbrain enhancement ready (16 kHz in/out)")
except Exception as e: # noqa: BLE001
log.warning("failed to load enhancement model %s: %s "
"(using fallback)", source, e)
self._enhancer = None
# -- DSP fallback ------------------------------------------------------
@staticmethod
def _simple_clean(x: np.ndarray) -> np.ndarray:
"""Lightweight cleanup: high-pass + soft knee + gentle gain.
Good enough to remove DC offset/low rumble and taming clipping
without any ML model.
"""
# first-order high-pass at ~80 Hz
if len(x) > 1:
alpha = 0.998
y = np.empty_like(x)
y[0] = x[0]
for i in range(1, len(x)):
y[i] = alpha * (y[i - 1] + x[i] - x[i - 1])
x = y
# soft-clip to avoid harsh distortion
x = np.tanh(x * 0.9) / np.tanh(0.9)
# gentle makeup gain
x = x * 1.1
return np.clip(x, -1.0, 1.0)
# -- public API ---------------------------------------------------------
def process(self, pcm: np.ndarray) -> np.ndarray:
"""pcm: int16 mono numpy array. Returns int16 mono, same length.
If the model changes the length, we trim/pad back to the original
length so downstream chunking stays aligned.
"""
if self._enhancer is None:
if len(pcm) == 0:
return pcm
x = pcm.astype(np.float32) / 32768.0
y = self._simple_clean(x)
return (y * 32767.0).astype(np.float32).astype(np.int16)
x = pcm.astype(np.float32) / 32768.0
torch = self._torch
try:
w = torch.from_numpy(x).unsqueeze(0).to(self._enhancer.device)
# speechbrain >= 1.1 (e.g. metricGAN-plus) requires `lengths`
# on enhance_batch. The value is RELATIVE (multiplied by the
# number of frames internally), so 1.0 = "the whole sequence".
# Older versions accept the call without it (TypeError fallback).
lengths = torch.tensor([1.0]).to(self._enhancer.device)
try:
with torch.no_grad():
out = self._enhancer.enhance_batch(w, lengths=lengths)
except TypeError:
with torch.no_grad():
out = self._enhancer.enhance_batch(w)
out = out.squeeze(0).cpu().numpy()
except Exception as e: # noqa: BLE001
log.warning("enhance_batch failed (%s); using fallback", e)
y = self._simple_clean(x)
return (y * 32767.0).astype(np.float32).astype(np.int16)
# align length
n = len(x)
if len(out) >= n:
out = out[:n]
else:
out = np.pad(out, (0, n - len(out)))
return (out * 32767.0).astype(np.float32).astype(np.int16)
+91
View File
@@ -0,0 +1,91 @@
"""Load configuration from environment / .env file."""
import os
from dataclasses import dataclass, field
from pathlib import Path
from typing import Optional
def _load_dotenv(path: Path) -> None:
if not path.is_file():
return
for line in path.read_text().splitlines():
line = line.strip()
if not line or line.startswith("#") or "=" not in line:
continue
key, _, value = line.partition("=")
os.environ.setdefault(key.strip(), value.strip())
@dataclass
class Config:
# connection
url: str = ""
api_key: str = ""
api_secret: str = ""
room: str = "cameras"
participant_prefix: str = "cam"
# video
width: int = 1280
height: int = 720
fps: int = 30
min_fps: int = 5
# audio
audio_rate: int = 16000
audio_channels: int = 1
# speechbrain cleanup
enhance_mode: str = "auto" # auto | force | off
enhance_model: str = "speechbrain/metricgan-plus-voicebank"
enhance_chunk_s: float = 0.2
# misc
publish_timeout_s: float = 60.0
log_level: str = "INFO"
cameras: list[str] = field(default_factory=lambda: ["all"])
def _parse_cameras(spec: str) -> list[str]:
"""Expand 'all' / '1,3,5' / '0-19' into a list of indices ('' = all)."""
spec = (spec or "all").strip().lower()
if spec in ("all", ""):
return ["all"]
out = []
for part in spec.split(","):
part = part.strip()
if not part:
continue
if "-" in part:
lo, hi = part.split("-", 1)
out.extend(str(i) for i in range(int(lo), int(hi) + 1))
else:
out.append(part)
return out
def load_config(env_path: Optional[Path] = None) -> Config:
_load_dotenv(Path(env_path or Path(__file__).parent / ".env"))
cfg = Config(
url=os.environ.get("LIVEKIT_URL", ""),
api_key=os.environ.get("LIVEKIT_API_KEY", ""),
api_secret=os.environ.get("LIVEKIT_API_SECRET", ""),
room=os.environ.get("LIVEKIT_ROOM", "cameras"),
participant_prefix=os.environ.get("PARTICIPANT_PREFIX", "cam"),
width=int(os.environ.get("VIDEO_WIDTH", "1280")),
height=int(os.environ.get("VIDEO_HEIGHT", "720")),
fps=int(os.environ.get("VIDEO_FPS", "30")),
min_fps=int(os.environ.get("VIDEO_MIN_FPS", "5")),
audio_rate=int(os.environ.get("AUDIO_SAMPLE_RATE", "16000")),
audio_channels=int(os.environ.get("AUDIO_CHANNELS", "1")),
enhance_mode=os.environ.get("ENHANCE_MODE", "auto").lower(),
enhance_model=os.environ.get(
"ENHANCE_MODEL", "speechbrain/metricgan-plus-voicebank"
),
enhance_chunk_s=float(os.environ.get("ENHANCE_CHUNK_S", "0.2")),
publish_timeout_s=float(os.environ.get("PUBLISH_TIMEOUT_S", "60")),
log_level=os.environ.get("LOG_LEVEL", "INFO").upper(),
cameras=_parse_cameras(os.environ.get("CAMERAS", "all")),
)
if cfg.enhance_mode not in ("auto", "force", "off"):
raise ValueError(f"ENHANCE_MODE must be auto|force|off, got {cfg.enhance_mode}")
return cfg
+244
View File
@@ -0,0 +1,244 @@
"""Discover C920 cameras: pair each /dev/video* node with its ALSA card.
Identity is stable across reboots and connect order: each camera is bound
to its USB *serial number* by udev rules (see generate_udev_rules.py and
slots.json in this directory), which pin
- video: /dev/camNN symlinks (primary UVC node per serial)
- audio: fixed ALSA card index per serial (SOUND_CARD_INDEX)
discover_cameras() therefore assigns cam-01..cam-20 by serial, never by
enumeration order. If slots.json is missing (first run), it falls back to
USB-port-based pairing and enumeration order.
"""
import json
import logging
import re
import subprocess
from dataclasses import dataclass, field
from pathlib import Path
log = logging.getLogger("cameras.discovery")
_SLOTS_FILE = Path(__file__).parent / "slots.json"
@dataclass
class Camera:
index: int
identity: str
video_device: str
audio_card: int | None
usb_port: str = ""
serial: str = ""
extra_video_devices: list[str] = field(default_factory=list)
def _run(cmd: list[str]) -> str:
return subprocess.run(cmd, capture_output=True, text=True, check=False).stdout
def _strip_interface(port: str) -> str:
"""Drop the trailing ':config.interface' so the video interface (x:1.0)
and audio interface (x:1.2) of the same C920 compare equal."""
if ":" in port:
port = port.rsplit(":", 1)[0]
return port
def _serial_of(sysfs_path: str) -> str:
"""Walk up from a sysfs path (interface, function, or card device)
until we reach the USB device that carries the 'serial' attribute."""
p = Path(sysfs_path)
for _ in range(12):
cand = p / "serial"
if cand.is_file():
return cand.read_text().strip()
if str(p).startswith("/sys/devices"):
p = p.parent
else:
return ""
return ""
def _load_slots() -> dict[str, dict]:
"""slots.json: identity -> {serial, video, card}. Empty if absent."""
if _SLOTS_FILE.exists():
try:
return json.loads(_SLOTS_FILE.read_text())
except (OSError, ValueError) as e:
log.warning("could not read %s (%s); using port-based fallback",
_SLOTS_FILE, e)
return {}
def _usb_port_of_video_node(video_node: str) -> str:
"""Resolve the USB port (e.g. '0c:00.0-1.1.2.1') for a video device node."""
link = Path(f"/sys/class/video4linux/{Path(video_node).name}")
if not link.exists():
return ""
dev = link / "device"
try:
real = dev.resolve()
except OSError:
return ""
# The USB port is the segment after the last '/' of the USB part, e.g.
# /sys/devices/.../usb9/9-1/9-1.1/9-1.1.2/9-1.1.2.1/...
parts = real.parts
for i, p in enumerate(parts):
if p.startswith("usb"):
# take everything after the usbX part as the port path
port_parts = [q for q in parts[i + 1:] if re.match(r"^\d+-", q)]
if port_parts:
return "/".join(port_parts)
return ""
def _usb_port_of_sound_card(card: int) -> str:
dev = Path(f"/sys/class/sound/card{card}/device")
if not dev.exists():
return ""
try:
real = dev.resolve()
except OSError:
return ""
parts = real.parts
for i, p in enumerate(parts):
if p.startswith("usb"):
port_parts = [q for q in parts[i + 1:] if re.match(r"^\d+-", q)]
if port_parts:
return "/".join(port_parts)
return ""
def _sound_card_names() -> dict[int, str]:
cards: dict[int, str] = {}
out = _run(["cat", "/proc/asound/cards"])
# format: " 2 [C920 ]: USB-Audio - HD Pro Webcam C920"
# (the card name is left-justified in a fixed-width bracket field)
for line in out.splitlines():
m = re.match(r"\s*(\d+)\s+\[([A-Za-z0-9_]+)\s*\]:\s*(.*)", line)
if m:
cards[int(m.group(1))] = m.group(3).strip()
return cards
def _serial_of_video_node(video_node: str) -> str:
link = Path(f"/sys/class/video4linux/{Path(video_node).name}/device")
try:
real = link.resolve()
except OSError:
return ""
return _serial_of(str(real))
def _serial_of_sound_card(card: int) -> str:
dev = Path(f"/sys/class/sound/card{card}/device")
try:
real = dev.resolve()
except OSError:
return ""
return _serial_of(str(real))
def discover_cameras() -> list[Camera]:
"""Return cameras with stable serial-based identities (cam-01..cam-20)."""
slots = _load_slots()
slots_by_serial = {s["serial"]: (ident, s)
for ident, s in slots.items() if s.get("serial")}
v4l_out = _run(["v4l2-ctl", "--list-devices"])
if not v4l_out:
raise RuntimeError("v4l2-ctl not available; is v4l-utils installed?")
# group video nodes by device name from v4l2-ctl output
# v4l2-ctl output: "NAME:\n\t/dev/videoX\n\t/dev/videoY\n\t/dev/mediaZ"
current: str | None = None
groups: dict[str, list[str]] = {}
for raw in v4l_out.splitlines():
line = raw.strip()
if not line:
continue
if line.endswith(":") and "\t" not in raw:
current = line[:-1]
groups.setdefault(current, [])
elif current and line.startswith("/dev/video"):
groups[current].append(line)
# sound cards keyed by USB serial (authoritative) with port fallback
card_names = _sound_card_names()
card_by_serial: dict[str, int] = {}
card_by_port: dict[str, int] = {}
for card, name in card_names.items():
if "C920" not in name and "920" not in name:
continue
serial = _serial_of_sound_card(card)
if serial:
card_by_serial[serial] = card
else:
port = _usb_port_of_sound_card(card)
if port:
card_by_port[_strip_interface(port)] = card
cameras: list[Camera] = []
for name, nodes in groups.items():
if not nodes:
continue
nodes.sort(key=lambda n: int(Path(n).name.replace("video", "")))
primary = nodes[0]
serial = _serial_of_video_node(primary)
port = _usb_port_of_video_node(primary)
identity = None
video_device = primary
audio_card = None
slot = slots_by_serial.get(serial)
if slot is not None:
identity, slot_info = slot
# prefer the pinned udev symlink if it exists
pinned = slot_info.get("video", "")
if pinned and Path(pinned).exists():
video_device = pinned
if slot_info.get("card") is not None:
audio_card = slot_info["card"]
if serial in card_by_serial:
audio_card = card_by_serial[serial] # live serial match wins
else:
log.warning("camera %s (%s): serial %s not in %s; "
"port-based fallback", name, primary, serial,
_SLOTS_FILE.name)
if serial:
audio_card = card_by_serial.get(serial)
if audio_card is None:
audio_card = card_by_port.get(_strip_interface(port))
if audio_card is None:
log.warning("camera %s (%s): no matching sound card "
"(port=%s)", name, primary, port)
identity = None # assigned below by enumeration order
cameras.append(Camera(
index=len(cameras),
identity=identity or f"cam-{len(cameras) + 1:02d}",
video_device=video_device,
audio_card=audio_card,
usb_port=port,
serial=serial,
extra_video_devices=nodes[1:],
))
# stable order: known slots by slot number, unknowns by video node number
def _video_num(cam: Camera) -> int:
try:
# resolve a pinned /dev/camNN symlink to its real /dev/videoN node
real = Path(cam.video_device).resolve().name
return int(real.replace("video", ""))
except (ValueError, OSError):
return 10_000
def _slot_num(cam: Camera) -> int:
m = re.match(r"cam-(\d+)$", cam.identity)
return int(m.group(1)) if m else 1000
cameras.sort(key=lambda c: (_slot_num(c), _video_num(c)))
for i, cam in enumerate(cameras):
cam.index = i
return cameras
+141
View File
@@ -0,0 +1,141 @@
#!/usr/bin/env python3
"""Generate udev rules that pin each C920 webcam to a fixed slot.
Slot N (cam-01..cam-20) is bound to the USB *serial number* of the camera
that currently occupies it, giving stable:
- video: /dev/camNN (primary capture node), /dev/camNN.meta (UVC metadata)
- audio: fixed ALSA card index (hw:INDEX,0)
Run: python3 generate_udev_rules.py # prints to stdout
python3 generate_udev_rules.py --write # also writes /etc/udev/rules.d/99-c920-pin.rules
"""
import re
import subprocess
import sys
HEADER = """\
# Stable slots for 20x Logitech HD Pro Webcam C920 (046d:08e5).
# Slot = physical camera, identified by USB serial. Survives reboots,
# re-enumeration, and hub reordering. Generated by generate_udev_rules.py.
#
# /dev/camNN -> primary video capture node (UVC stream 0)
# /dev/camNN.meta -> UVC metadata node (stream 1)
# hw:N,0 -> mic of cam-NN (N from SOUND_CARD_INDEX below)
"""
def sh(cmd):
return subprocess.run(cmd, capture_output=True, text=True).stdout
def _serial_of(sysfs_path):
"""Walk up from a sysfs path (interface, function, or card device)
until we reach the USB device that carries the 'serial' attribute."""
import pathlib
p = pathlib.Path(sysfs_path)
for _ in range(12):
cand = p / "serial"
if cand.is_file():
return cand.read_text().strip()
if str(p).startswith("/sys/devices"):
p = p.parent
else:
return ""
return ""
def current_state():
"""Map serial -> (primary_video_node, sound_card_index)."""
import pathlib
# sound cards: index -> serial
cards = {}
for line in sh(["cat", "/proc/asound/cards"]).splitlines():
m = re.match(r"\s*(\d+)\s+\[([A-Za-z0-9_]+)\s*\]:\s*(.*)", line)
if not m:
continue
idx = int(m.group(1))
if "C920" not in m.group(3) and "920" not in m.group(2):
continue
dev = f"/sys/class/sound/card{idx}/device"
real = sh(["readlink", "-f", dev]).strip()
cards[idx] = _serial_of(real)
# video: primary node (index 0) per USB serial, via sysfs
primary_by_serial = {}
for node in sorted(pathlib.Path("/sys/class/video4linux").glob("video*"),
key=lambda p: int(p.name.replace("video", ""))):
idx_attr = (node / "index").read_text().strip()
if idx_attr != "0":
continue
real = sh(["readlink", "-f", str(node / "device")]).strip()
s = _serial_of(real)
primary_by_serial[s] = f"/dev/{node.name}"
return primary_by_serial, cards
def main():
primary_by_serial, cards = current_state()
# assign slots: cam-01..cam-20 in current video-node order
ordered = []
for serial, node in primary_by_serial.items():
n = int(node.replace("/dev/video", ""))
ordered.append((n, serial, node))
ordered.sort()
if len(ordered) != 20:
print(f"warning: expected 20 cameras, got {len(ordered)}", file=sys.stderr)
lines = [HEADER]
for i, (vnum, serial, node) in enumerate(ordered, start=1):
card = next((c for c, s in cards.items() if s == serial), None)
slot = f"cam{i:02d}"
lines.append(f"# {slot}: serial {serial} (was {node}, card {card})")
lines.append(
f'SUBSYSTEM=="video4linux", KERNEL=="video*", '
f'ENV{{ID_VENDOR_ID}}=="046d", ENV{{ID_MODEL_ID}}=="08e5", '
f'ENV{{ID_USB_SERIAL_SHORT}}=="{serial}", '
f'ATTR{{index}}=="0", SYMLINK+="{slot}"')
lines.append(
f'SUBSYSTEM=="video4linux", KERNEL=="video*", '
f'ENV{{ID_VENDOR_ID}}=="046d", ENV{{ID_MODEL_ID}}=="08e5", '
f'ENV{{ID_USB_SERIAL_SHORT}}=="{serial}", '
f'ATTR{{index}}=="1", SYMLINK+="{slot}.meta"')
if card is not None:
lines.append(
f'SUBSYSTEM=="sound", KERNEL=="card*", '
f'ENV{{ID_VENDOR_ID}}=="046d", ENV{{ID_MODEL_ID}}=="08e5", '
f'ENV{{ID_USB_SERIAL_SHORT}}=="{serial}", '
f'ENV{{SOUND_CARD_INDEX}}="{card}"')
else:
print(f"warning: no sound card for serial {serial}", file=sys.stderr)
lines.append("")
text = "\n".join(lines)
print(text)
if "--write" in sys.argv:
path = "/etc/udev/rules.d/99-c920-pin.rules"
with open(path, "w") as f:
f.write(text)
print(f"# wrote {path}", file=sys.stderr)
# durable slot map for the app (discovery.py reads this)
import json
import os
slots = {}
for i, (vnum, serial, node) in enumerate(ordered, start=1):
card = next((c for c, s in cards.items() if s == serial), None)
slots[f"cam-{i:02d}"] = {
"serial": serial,
"video": f"/dev/cam{i:02d}",
"card": card,
}
base = os.path.dirname(os.path.abspath(__file__))
slots_path = os.path.join(base, "slots.json")
with open(slots_path, "w") as f:
json.dump(slots, f, indent=2)
print(f"# wrote {slots_path}", file=sys.stderr)
# summary table to stderr
for i, (vnum, serial, node) in enumerate(ordered, start=1):
card = next((c for c, s in cards.items() if s == serial), "?")
print(f"# cam-{i:02d} serial={serial} video={node} card={card}",
file=sys.stderr)
if __name__ == "__main__":
main()
+49
View File
@@ -0,0 +1,49 @@
"""One-time download of the speechbrain enhancement model.
Each worker loads the model from the local cache, so run this once before
starting all 20 workers:
.venv/bin/python prefetch_model.py
"""
import logging
import sys
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger("prefetch")
def main() -> int:
from config import load_config
cfg = load_config()
import torch
from speechbrain.inference.enhancement import SpectralMaskEnhancement
log.info("loading %s (first run downloads weights) ...", cfg.enhance_model)
model = SpectralMaskEnhancement.from_hparams(
source=cfg.enhance_model,
savedir="/root/.cache/speechbrain-enhancement",
)
model.mods.enhance_model.eval()
# quick self-test: 0.5 s of 16 kHz sine + noise
import numpy as np
n = int(cfg.audio_rate * 0.5)
t = np.arange(n) / cfg.audio_rate
x = (0.3 * np.sin(2 * np.pi * 440 * t) + 0.05 * np.random.randn(n))
x = x / np.abs(x).max()
with torch.no_grad():
w = torch.from_numpy(x.astype(np.float32)).unsqueeze(0).to(model.device)
# lengths are relative (1.0 = full sequence); older speechbrain
# builds don't accept the kwarg at all -> TypeError fallback.
try:
out = model.enhance_batch(w, lengths=torch.tensor([1.0]))
except TypeError:
out = model.enhance_batch(w)
log.info("self-test OK: in %d samples -> out %d samples",
n, out.numel())
return 0
if __name__ == "__main__":
sys.exit(main())
+13
View File
@@ -0,0 +1,13 @@
[project]
name = "livekit-cameras"
version = "0.1.0"
description = "Publish 20 C920 webcam streams (video + speechbrain-cleaned audio) to a LiveKit room"
requires-python = ">=3.10"
dependencies = [
"livekit>=1.1",
"livekit-api>=1.2",
"speechbrain",
"numpy",
"opencv-python-headless",
"scipy",
]
+180
View File
@@ -0,0 +1,180 @@
"""Orchestrator: discover cameras and run one publisher worker each.
Usage:
python run.py # use .env (or environment variables)
python run.py --cameras 1,5 # publish only cameras 1 and 5
python run.py --once # run, print a status report, exit (for CI)
"""
from __future__ import annotations
import argparse
import asyncio
import contextlib
import json
import logging
import os
import signal
import subprocess
import sys
import time
from pathlib import Path
from config import load_config, _load_dotenv
from discovery import discover_cameras
from tokens import make_token
log = logging.getLogger("cameras.main")
def worker_args_for(cfg, cam) -> list[str]:
return [
"--camera-json", json.dumps(cam.__dict__),
"--room", cfg.room,
"--url", cfg.url,
"--token", make_token(cfg.url, cfg.api_key, cfg.api_secret,
cfg.room, cam.identity),
"--width", str(cfg.width),
"--height", str(cfg.height),
"--fps", str(cfg.fps),
"--audio-rate", str(cfg.audio_rate),
"--enhance-mode", cfg.enhance_mode,
"--enhance-model", cfg.enhance_model,
"--enhance-chunk-s", str(cfg.enhance_chunk_s),
]
class Supervisor:
def __init__(self, cfg, cameras):
self.cfg = cfg
self.cameras = cameras
self.procs: dict[str, subprocess.Popen] = {}
self.last_exit: dict[str, int | None] = {}
self.restarts: dict[str, int] = {}
def spawn(self, cam) -> None:
args = worker_args_for(self.cfg, cam)
env = dict(os.environ)
p = subprocess.Popen(
[sys.executable, str(Path(__file__).parent / "run_worker.py"),
*args],
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
text=True, bufsize=1, env=env,
)
self.procs[cam.identity] = p
self.restarts[cam.identity] = self.restarts.get(cam.identity, 0) + 1
def pump(self) -> None:
"""Read one line from each running worker; print them."""
for ident, p in list(self.procs.items()):
line = p.stdout.readline() if p.stdout else ""
if line:
sys.stdout.write(f"[{ident}] {line}")
sys.stdout.flush()
if p.poll() is not None and not line:
pass # EOF handled by poll below
def reap(self) -> None:
for ident, p in list(self.procs.items()):
rc = p.poll()
if rc is None:
continue
self.last_exit[ident] = rc
if rc != 0:
log.warning("%s exited rc=%s (restart %d/%d)", ident, rc,
self.restarts.get(ident, 0), 10)
if self.restarts.get(ident, 0) < 10:
time.sleep(1.0)
# find the camera object again
cam = next(c for c in self.cameras
if c.identity == ident)
self.spawn(cam)
else:
log.error("%s gave up after 10 restarts", ident)
else:
log.info("%s exited cleanly", ident)
def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--cameras", default=None,
help="subset spec, e.g. 'all', '1,5', '0-19'")
ap.add_argument("--once", action="store_true",
help="spawn workers, report status after "
"PUBLISH_TIMEOUT_S, then exit")
args = ap.parse_args()
cfg = load_config()
logging.basicConfig(
level=getattr(logging, cfg.log_level, logging.INFO),
format="%(asctime)s [main] %(levelname)s %(message)s")
if args.cameras:
from config import _parse_cameras
cfg.cameras = _parse_cameras(args.cameras)
if not cfg.url or not cfg.api_key or not cfg.api_secret:
log.error("LIVEKIT_URL / LIVEKIT_API_KEY / LIVEKIT_API_SECRET "
"not set (create .env from .env.example)")
return 2
cameras = discover_cameras()
log.info("discovered %d cameras:", len(cameras))
for c in cameras:
log.info(" %s video=%s audio=hw:%s usb=%s",
c.identity, c.video_device,
c.audio_card if c.audio_card is not None else "-",
c.usb_port or "-")
if cfg.cameras != ["all"]:
wanted = set(int(i) - 1 for i in cfg.cameras)
cameras = [c for c in cameras if c.index in wanted]
if not cameras:
log.error("no cameras match --cameras %s", args.cameras)
return 2
sup = Supervisor(cfg, cameras)
for cam in cameras:
sup.spawn(cam)
log.info("spawned %d workers", len(cameras))
if args.once:
deadline = time.monotonic() + cfg.publish_timeout_s
while time.monotonic() < deadline:
sup.pump()
sup.reap()
time.sleep(0.5)
# final report
up = [i for i, p in sup.procs.items() if p.poll() is None]
log.info("STATUS: %d/%d workers alive", len(up), len(sup.procs))
for ident in up:
pass
# stop everything
for p in sup.procs.values():
p.send_signal(signal.SIGTERM)
for p in sup.procs.values():
with contextlib.suppress(Exception):
p.wait(timeout=10)
return 0 if len(up) == len(sup.procs) else 1
try:
while True:
sup.pump()
sup.reap()
time.sleep(0.2)
except KeyboardInterrupt:
log.info("shutting down ...")
for p in sup.procs.values():
if p.poll() is None:
p.send_signal(signal.SIGTERM)
for p in sup.procs.values():
with contextlib.suppress(Exception):
p.wait(timeout=10)
for p in sup.procs.values():
if p.poll() is None:
p.kill()
return 0
if __name__ == "__main__":
sys.exit(main())
+225
View File
@@ -0,0 +1,225 @@
"""One camera -> one LiveKit participant.
Run as a subprocess per camera: python -m cameras.run_worker --index 0 ...
"""
from __future__ import annotations
import argparse
import asyncio
import contextlib
import json
import logging
import os
import signal
import subprocess
import sys
import time
import numpy as np
log = logging.getLogger("cameras.worker")
def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--camera-json", required=True,
help="JSON of one discovered camera dict")
ap.add_argument("--room", required=True)
ap.add_argument("--url", required=True)
ap.add_argument("--token", required=True)
ap.add_argument("--width", type=int, default=1280)
ap.add_argument("--height", type=int, default=720)
ap.add_argument("--fps", type=int, default=30)
ap.add_argument("--audio-rate", type=int, default=16000)
ap.add_argument("--enhance-mode", default="auto")
ap.add_argument("--enhance-model",
default="speechbrain/metricgan-plus-voicebank")
ap.add_argument("--enhance-chunk-s", type=float, default=0.2)
args = ap.parse_args()
cam = json.loads(args.camera_json)
logging.basicConfig(
level=logging.INFO,
format=f"%(asctime)s [{cam['identity']}] %(levelname)s %(message)s",
)
from config import Config # noqa: F401 (ensures .env is loaded)
import livekit.rtc as rtc
from audio_cleanup import AudioCleaner
identity = cam["identity"]
video_dev = cam["video_device"]
audio_card = cam.get("audio_card")
async def run() -> int:
cleaner = AudioCleaner(
mode=args.enhance_mode,
model_source=args.enhance_model,
sample_rate=args.audio_rate,
)
log.info("audio backend: %s", cleaner.backend)
room = rtc.Room()
await room.connect(args.url, args.token)
log.info("connected to room %r as %r", room.name, identity)
# ---------------- video ----------------
video_source = rtc.VideoSource(args.width, args.height)
video_track = rtc.LocalVideoTrack.create_video_track(
f"{identity}-video", video_source)
vproc = subprocess.Popen(
["ffmpeg", "-hide_banner", "-loglevel", "error",
"-f", "v4l2",
"-input_format", "mjpeg",
"-video_size", f"{args.width}x{args.height}",
"-framerate", str(args.fps),
"-i", video_dev,
"-f", "rawvideo", "-pix_fmt", "bgra",
"pipe:1"],
stdout=subprocess.PIPE,
)
frame_bytes = args.width * args.height * 4
video_frames = 0
async def video_loop():
nonlocal video_frames
loop = asyncio.get_running_loop()
last = time.monotonic()
while vproc.poll() is None:
data = await loop.run_in_executor(
None, vproc.stdout.read, frame_bytes)
if not data or len(data) < frame_bytes:
if not data:
break
# partial read: accumulate (v4l2/ffmpeg shouldn't do this)
keep = b""
while len(keep) < frame_bytes:
more = await loop.run_in_executor(
None, vproc.stdout.read,
frame_bytes - len(keep))
if not more:
return
keep += more
data = keep
now = time.monotonic()
ts_us = int((now - last) * 1_000_000)
try:
video_source.capture_frame(
rtc.VideoFrame(
args.width, args.height,
rtc.VideoBufferType.BGRA, data),
timestamp_us=ts_us,
)
video_frames += 1
except Exception: # noqa: BLE001
log.exception("video capture_frame failed")
# ---------------- audio ----------------
aproc: subprocess.Popen | None = None
atrack: rtc.LocalAudioTrack | None = None
if audio_card is not None:
audio_source = rtc.AudioSource(
sample_rate=args.audio_rate, num_channels=1)
atrack = rtc.LocalAudioTrack.create_audio_track(
f"{identity}-audio", audio_source)
chunk_s = args.enhance_chunk_s
chunk = int(chunk_s * args.audio_rate) # samples per channel
async def audio_loop():
assert aproc is not None and atrack is not None
loop = asyncio.get_running_loop()
# mic delivers stereo s16le: bytes per chunk =
# chunk * 2 ch * 2 bytes
need = chunk * 2 * 2
buf = b""
while aproc.poll() is None:
data = await loop.run_in_executor(
None, aproc.stdout.read, 4096)
if not data:
break
buf += data
while len(buf) >= need:
piece, buf = buf[:need], buf[need:]
pcm_stereo = np.frombuffer(piece, dtype=np.int16)
pcm_mono = pcm_stereo.reshape(-1, 2).mean(
axis=1).astype(np.int16)
# align to full chunks
usable = (len(pcm_mono) // chunk) * chunk
if usable:
cleaned = await loop.run_in_executor(
None, cleaner.process,
pcm_mono[:usable])
await audio_source.capture_frame(
rtc.AudioFrame(
cleaned.tobytes(),
args.audio_rate,
1,
len(cleaned)))
aproc = subprocess.Popen(
["ffmpeg", "-hide_banner", "-loglevel", "error",
"-f", "alsa", "-i", f"hw:{audio_card},0",
"-ar", str(args.audio_rate), "-ac", "2",
"-f", "s16le", "pipe:1"],
stdout=subprocess.PIPE,
)
# ---------------- publish ----------------
from livekit.rtc import TrackSource, TrackPublishOptions
pub_opts = TrackPublishOptions(
source=TrackSource.SOURCE_CAMERA, simulcast=False)
await room.local_participant.publish_track(video_track, pub_opts)
if atrack is not None:
await room.local_participant.publish_track(
atrack,
TrackPublishOptions(
source=TrackSource.SOURCE_MICROPHONE,
dtx=False))
log.info("published video+audio (audio backend: %s)",
cleaner.backend)
else:
log.info("published video only (no audio card matched)")
log.info("status: UP identity=%s video=%s audio=%s",
identity, video_dev,
f"hw:{audio_card},0" if audio_card is not None else "none")
# ---------------- run ----------------
tasks = [asyncio.create_task(video_loop(), name="video")]
if aproc is not None:
tasks.append(asyncio.create_task(audio_loop(), name="audio"))
stop = asyncio.Event()
loop = asyncio.get_running_loop()
for sig in (signal.SIGTERM, signal.SIGINT):
with contextlib.suppress(NotImplementedError):
loop.add_signal_handler(sig, stop.set)
await stop.wait()
for t in tasks:
t.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
for p in (vproc, aproc):
if p is not None and p.poll() is None:
with contextlib.suppress(Exception):
p.terminate()
with contextlib.suppress(Exception):
p.wait(timeout=3)
log.info("stopped; delivered %d video frames", video_frames)
with contextlib.suppress(Exception):
await room.disconnect()
return 0
try:
return asyncio.run(run())
except Exception: # noqa: BLE001
log.exception("worker fatal")
return 1
if __name__ == "__main__":
sys.exit(main())
+102
View File
@@ -0,0 +1,102 @@
{
"cam-01": {
"serial": "503A123F",
"video": "/dev/cam01",
"card": 2
},
"cam-02": {
"serial": "4123BF6F",
"video": "/dev/cam02",
"card": 3
},
"cam-03": {
"serial": "84D270BF",
"video": "/dev/cam03",
"card": 4
},
"cam-04": {
"serial": "F8F2BF6F",
"video": "/dev/cam04",
"card": 5
},
"cam-05": {
"serial": "B638DFEF",
"video": "/dev/cam05",
"card": 6
},
"cam-06": {
"serial": "7972BF6F",
"video": "/dev/cam06",
"card": 7
},
"cam-07": {
"serial": "6AA1BFAF",
"video": "/dev/cam07",
"card": 8
},
"cam-08": {
"serial": "FDA67EEF",
"video": "/dev/cam08",
"card": 9
},
"cam-09": {
"serial": "65B2BF2F",
"video": "/dev/cam09",
"card": 10
},
"cam-10": {
"serial": "FAA65E2F",
"video": "/dev/cam10",
"card": 11
},
"cam-11": {
"serial": "33C920FF",
"video": "/dev/cam11",
"card": 12
},
"cam-12": {
"serial": "5292BF6F",
"video": "/dev/cam12",
"card": 13
},
"cam-13": {
"serial": "357A8E2F",
"video": "/dev/cam13",
"card": 14
},
"cam-14": {
"serial": "2EB1BF2F",
"video": "/dev/cam14",
"card": 15
},
"cam-15": {
"serial": "B478307F",
"video": "/dev/cam15",
"card": 16
},
"cam-16": {
"serial": "9BF2BF6F",
"video": "/dev/cam16",
"card": 17
},
"cam-17": {
"serial": "0642BF6F",
"video": "/dev/cam17",
"card": 18
},
"cam-18": {
"serial": "8CE2BF6F",
"video": "/dev/cam18",
"card": 19
},
"cam-19": {
"serial": "D41E8E6F",
"video": "/dev/cam19",
"card": 20
},
"cam-20": {
"serial": "9B91BF6F",
"video": "/dev/cam20",
"card": 21
}
}
+43
View File
@@ -0,0 +1,43 @@
"""Smoke test: one camera worker against the local LiveKit server, ~20 s."""
import json
import signal
import subprocess
import sys
import time
from pathlib import Path
HERE = Path(__file__).parent
sys.path.insert(0, str(HERE))
from discovery import discover_cameras
from tokens import make_token
cam = discover_cameras()[0]
tok = make_token("ws://localhost:7880", "devkey", "devsecret", "smoke", "smoke-test")
p = subprocess.Popen(
[str(HERE / ".venv/bin/python"), str(HERE / "run_worker.py"),
"--camera-json", json.dumps({
"index": 0, "identity": "smoke-test",
"video_device": cam.video_device, "audio_card": cam.audio_card,
"usb_port": cam.usb_port}),
"--room", "smoke", "--url", "ws://localhost:7880", "--token", tok,
"--width", "640", "--height", "480", "--fps", "30",
"--audio-rate", "16000", "--enhance-mode", "off"],
stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1,
)
t0 = time.time()
lines = []
while time.time() - t0 < 20:
line = p.stdout.readline()
if line:
lines.append(line.rstrip())
print(line.rstrip(), flush=True)
p.send_signal(signal.SIGTERM)
try:
rc = p.wait(timeout=15)
except subprocess.TimeoutExpired:
p.kill()
rc = "killed"
print(f"--- smoke test done: rc={rc}")
print("published markers:", [l for l in lines if "published" in l or "status" in l or "UP" in l])
+20
View File
@@ -0,0 +1,20 @@
"""Generate LiveKit access tokens for the camera participants."""
from livekit.api import AccessToken, VideoGrants
def make_token(url: str, api_key: str, api_secret: str,
room: str, identity: str) -> str:
token = (
AccessToken(api_key, api_secret)
.with_identity(identity)
.with_name(identity)
.with_grants(VideoGrants(
room_join=True,
room=room,
can_publish=True,
can_subscribe=True,
can_publish_data=True,
))
)
return token.to_jwt()