#cloud-config # Nodo stream Hetzner: MediaMTX + ffmpeg + relay-agent (YouTube sul nodo). # HCLOUD_USER_DATA_FILE=.../infra/stream-node/cloud-init.yaml package_update: true packages: - docker.io - ffmpeg - curl - apparmor - python3 - fonts-dejavu-core write_files: - path: /opt/stream-node/mediamtx.yml permissions: "0644" content: | logLevel: info authInternalUsers: - user: any pass: ips: [] permissions: - action: api - action: metrics - action: publish - action: read - action: playback api: yes apiAddress: :9997 rtmp: yes rtmpAddress: :1935 hls: yes hlsAddress: :8888 paths: all_others: - path: /opt/stream-node/relay-agent.service permissions: "0644" content: | [Unit] Description=Match Live TV YouTube relay agent After=network-online.target docker.service Wants=network-online.target [Service] Type=simple EnvironmentFile=-/opt/stream-node/agent.env ExecStart=/usr/bin/python3 /opt/stream-node/relay-agent.py Restart=always RestartSec=2 [Install] WantedBy=multi-user.target - path: /opt/stream-node/agent.env permissions: "0600" content: | STREAM_NODE_AGENT_SECRET=mediamtx_webhook_dev_secret STREAM_NODE_AGENT_LISTEN=0.0.0.0:9100 STREAM_NODE_LOCAL_HLS_URL=http://127.0.0.1:8888 STREAM_NODE_RELAY_LOG_DIR=/var/log STREAM_NODE_SLATES_DIR=/slates/custom STREAM_NODE_RECORDINGS_DIR=/recordings - path: /opt/stream-node/relay-agent.py permissions: "0755" content: | #!/usr/bin/env python3 """Avvia/ferma ffmpeg sul nodo stream (loopback MediaMTX → YouTube). Il control plane chiama questo agent all'avvio della diretta; ffmpeg resta sul CPX. """ from __future__ import annotations import json import os import shutil import signal import subprocess import tarfile import tempfile import threading from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from urllib.parse import unquote, urlparse SECRET = os.environ.get("STREAM_NODE_AGENT_SECRET", "") LISTEN = os.environ.get("STREAM_NODE_AGENT_LISTEN", "0.0.0.0:9100") HLS_BASE = os.environ.get("STREAM_NODE_LOCAL_HLS_URL", "http://127.0.0.1:8888").rstrip("/") LOG_DIR = os.environ.get("STREAM_NODE_RELAY_LOG_DIR", "/var/log") SLATES_DIR = os.environ.get("STREAM_NODE_SLATES_DIR", "/slates/custom") RECORDINGS_DIR = os.environ.get("STREAM_NODE_RECORDINGS_DIR", "/recordings") _lock = threading.Lock() # session_id -> {"pid": int, "path": str, "proc": Popen} _relays: dict[str, dict] = {} def _authorized(handler: BaseHTTPRequestHandler) -> bool: if not SECRET: return True hdr = handler.headers.get("Authorization", "") return hdr == f"Bearer {SECRET}" def _ffmpeg_cmd(path: str, rtmps: str) -> list[str]: return [ "ffmpeg", "-nostdin", "-hide_banner", "-loglevel", "warning", "-fflags", "+genpts+discardcorrupt", "-reconnect", "1", "-reconnect_streamed", "1", "-reconnect_on_network_error", "1", "-reconnect_delay_max", "2", "-rw_timeout", "15000000", "-live_start_index", "-1", "-i", f"{HLS_BASE}/{path}/index.m3u8", "-c:v", "copy", "-c:a", "copy", "-bsf:a", "aac_adtstoasc", "-f", "flv", rtmps, ] def _alive(pid: int) -> bool: try: os.kill(pid, 0) return True except OSError: return False def _stop_session(session_id: str) -> None: info = _relays.pop(session_id, None) if not info: return proc = info.get("proc") pid = int(info.get("pid") or 0) if proc and proc.poll() is None: try: os.killpg(proc.pid, signal.SIGTERM) except OSError: pass try: proc.wait(timeout=2) except Exception: try: os.killpg(proc.pid, signal.SIGKILL) except OSError: pass elif pid: try: os.kill(pid, signal.SIGTERM) except OSError: pass def _recordings_root(path_name): name = unquote(path_name or "").strip().lstrip("/") if not name or ".." in name.split("/"): return None base = os.path.realpath(RECORDINGS_DIR) root = os.path.realpath(os.path.join(base, name)) if root != base and not root.startswith(base + os.sep): return None return root def _has_recording_files(root: str) -> bool: if not os.path.isdir(root): return False for dirpath, _dirnames, filenames in os.walk(root): for name in filenames: if name.lower().endswith((".mp4", ".fmp4", ".m4s", ".ts")): return True return False def _start_session(session_id: str, path: str, rtmps: str) -> int: existing = _relays.get(session_id) if existing: proc = existing.get("proc") pid = int(existing.get("pid") or 0) if proc is not None and proc.poll() is None and existing.get("path") == path: return pid _stop_session(session_id) os.makedirs(LOG_DIR, exist_ok=True) safe = path.replace("/", "_") log_path = os.path.join(LOG_DIR, f"youtube-relay-{safe}.log") log = open(log_path, "ab", buffering=0) proc = subprocess.Popen( _ffmpeg_cmd(path, rtmps), stdout=log, stderr=log, start_new_session=True, ) _relays[session_id] = {"pid": proc.pid, "path": path, "proc": proc} return proc.pid class Handler(BaseHTTPRequestHandler): def log_message(self, fmt: str, *args) -> None: print(f"[relay-agent] {self.address_string()} {fmt % args}") def _json(self, code: int, payload: dict) -> None: body = json.dumps(payload).encode("utf-8") self.send_response(code) self.send_header("Content-Type", "application/json") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) def do_GET(self) -> None: parsed = urlparse(self.path) if parsed.path in ("/health", "/"): self._json(200, {"ok": True, "relays": len(_relays)}) return if not _authorized(self): self._json(401, {"error": "unauthorized"}) return if parsed.path.startswith("/relays/"): session_id = parsed.path.split("/relays/", 1)[1].strip("/") with _lock: info = _relays.get(session_id) running = bool(info and _alive(int(info["pid"]))) if info and not running: _relays.pop(session_id, None) if not running: self._json(404, {"running": False, "session_id": session_id}) return self._json(200, {"running": True, "session_id": session_id, "pid": info["pid"], "path": info["path"]}) return if parsed.path.startswith("/recordings/"): name = parsed.path.split("/recordings/", 1)[1].strip("/") root = _recordings_root(name) if not root or not _has_recording_files(root): self._json(404, {"error": "not found"}) return tmp = tempfile.NamedTemporaryFile(prefix="mltv-rec-", suffix=".tar.gz", delete=False) tmp.close() try: with tarfile.open(tmp.name, "w:gz") as tar: tar.add(root, arcname=".") size = os.path.getsize(tmp.name) self.send_response(200) self.send_header("Content-Type", "application/gzip") self.send_header("Content-Length", str(size)) self.end_headers() with open(tmp.name, "rb") as fh: shutil.copyfileobj(fh, self.wfile) finally: try: os.unlink(tmp.name) except OSError: pass return self._json(404, {"error": "not found"}) def do_POST(self) -> None: if not _authorized(self): self._json(401, {"error": "unauthorized"}) return if urlparse(self.path).path != "/relays": self._json(404, {"error": "not found"}) return length = int(self.headers.get("Content-Length") or 0) try: data = json.loads(self.rfile.read(length) or b"{}") except json.JSONDecodeError: self._json(400, {"error": "invalid json"}) return session_id = str(data.get("session_id") or "").strip() path = str(data.get("path") or "").strip() rtmps = str(data.get("rtmps") or "").strip() if not session_id or not path or not rtmps.startswith("rtmps://"): self._json(400, {"error": "session_id, path, rtmps required"}) return with _lock: pid = _start_session(session_id, path, rtmps) self._json(201, {"pid": pid, "session_id": session_id, "path": path}) def do_DELETE(self) -> None: if not _authorized(self): self._json(401, {"error": "unauthorized"}) return parsed = urlparse(self.path) if parsed.path.startswith("/relays/"): session_id = parsed.path.split("/relays/", 1)[1].strip("/") with _lock: _stop_session(session_id) self._json(200, {"stopped": True, "session_id": session_id}) return if parsed.path.startswith("/recordings/"): name = parsed.path.split("/recordings/", 1)[1].strip("/") root = _recordings_root(name) if not root or not os.path.isdir(root): self._json(404, {"error": "not found"}) return shutil.rmtree(root, ignore_errors=True) self._json(200, {"deleted": True, "path": name}) return self._json(404, {"error": "not found"}) def do_PUT(self) -> None: if not _authorized(self): self._json(401, {"error": "unauthorized"}) return parsed = urlparse(self.path) if not parsed.path.startswith("/slates/"): self._json(404, {"error": "not found"}) return name = parsed.path.split("/slates/", 1)[1].strip("/") if not name or "/" in name or ".." in name: self._json(400, {"error": "invalid slate name"}) return length = int(self.headers.get("Content-Length") or 0) if length <= 0: self._json(400, {"error": "empty body"}) return data = self.rfile.read(length) os.makedirs(SLATES_DIR, exist_ok=True) dest = os.path.join(SLATES_DIR, name) with open(dest, "wb") as out: out.write(data) self._json(201, {"path": dest, "bytes": len(data), "name": name}) def main() -> None: host, port_s = LISTEN.rsplit(":", 1) httpd = ThreadingHTTPServer((host, int(port_s)), Handler) print(f"[relay-agent] listen {LISTEN} hls={HLS_BASE}") httpd.serve_forever() if __name__ == "__main__": main() - path: /opt/stream-node/bootstrap.sh permissions: "0755" content: | #!/bin/bash set -euo pipefail mkdir -p /recordings /slates/custom /var/log systemctl enable --now docker for i in $(seq 1 60); do if docker info >/dev/null 2>&1; then break; fi sleep 2 done docker pull bluenviron/mediamtx:latest docker rm -f mediamtx 2>/dev/null || true FONT=/usr/share/fonts/truetype/dejavu/DejaVuSans-Bold.ttf # Copertina alwaysAvailable (NON nera): resta in HLS/YouTube se l'app cade. ffmpeg -y -f lavfi -i color=c=0x0a0a0e:s=1280x720:r=30:d=30 \ -f lavfi -i anullsrc=r=48000:cl=mono \ -vf "drawtext=fontfile=${FONT}:text='Match Live TV':fontsize=64:fontcolor=white:x=(w-text_w)/2:y=(h-text_h)/2-48,drawtext=fontfile=${FONT}:text='Trasmissione in pausa':fontsize=36:fontcolor=white:x=(w-text_w)/2:y=(h-text_h)/2+36" \ -c:v libx264 -pix_fmt yuv420p -profile:v baseline -level 3.1 \ -x264-params "keyint=30:min-keyint=30:scenecut=0:bframes=0" \ -g 30 -keyint_min 30 -preset fast \ -c:a aac -b:a 128k -ac 1 -ar 48000 -shortest /slates/offline.mp4 docker run -d --name mediamtx --restart unless-stopped --network host \ --security-opt apparmor=unconfined \ -v /opt/stream-node/mediamtx.yml:/mediamtx.yml:ro \ -v /recordings:/recordings \ -v /slates:/slates:ro \ bluenviron/mediamtx:latest /mediamtx.yml if [[ -f /opt/stream-node/relay-agent.py ]]; then cp /opt/stream-node/relay-agent.service /etc/systemd/system/mltv-relay-agent.service systemctl daemon-reload systemctl enable --now mltv-relay-agent.service fi for i in $(seq 1 30); do if curl -fsS http://127.0.0.1:9997/v3/paths/list >/dev/null; then echo "mediamtx ready" | tee /opt/stream-node/ready exit 0 fi sleep 2 done echo "mediamtx NOT ready" | tee /opt/stream-node/ready docker logs mediamtx || true exit 1 runcmd: - [bash, /opt/stream-node/bootstrap.sh]