aboutsummaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-10-01 09:53:34 +0200
committerChristophe Besson <cbesson@gmail.com>2026-10-01 09:53:34 +0200
commit0673922e704f718037eeffd2724debcfa8ac0b4b (patch)
tree78fee9a5c04542f006c1757bc76c0af7f57a79db /packages
parent6426912946270bb02e7b94508008edf8f949d993 (diff)
downloadmeshbay-0673922e704f718037eeffd2724debcfa8ac0b4b.tar.gz
fix(node): ffprobe over a member's file is bounded, and stopped when it is
probe_video waits 30 s at most and kills ffprobe on a timeout or when its caller gives up — a cancelled wait left the process running. The seek probe kills what it timed out on. Stream, subtitle and enrichment requests no longer hang on a file that keeps ffprobe busy (F-18, timeouts; the protocol whitelist was dropped: ffmpeg already confines nested protocols of a local input, measured on 8.0 against HLS and concat inputs). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Diffstat (limited to 'packages')
-rw-r--r--packages/meshbay-node/src/meshbay_node/media_probe.py17
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py7
-rw-r--r--packages/meshbay-node/tests/test_ffprobe_is_bounded.py61
3 files changed, 84 insertions, 1 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/media_probe.py b/packages/meshbay-node/src/meshbay_node/media_probe.py
index a6267d9..e9090ea 100644
--- a/packages/meshbay-node/src/meshbay_node/media_probe.py
+++ b/packages/meshbay-node/src/meshbay_node/media_probe.py
@@ -10,6 +10,12 @@ import asyncio
import json
from dataclasses import dataclass, field
+# How long ffprobe may take over one file's headers. The file is a member's
+# upload as often as the operator's own: one that keeps ffprobe busy must not
+# keep the stream request, the subtitle request or the enrichment slot that
+# asked for it — the same bound the index-time enrichment already put around it.
+FFPROBE_TIMEOUT_SECS = 30
+
_H264_PROFILES = {"Baseline": "42", "Main": "4d", "High": "64", "High 10": "6e"}
# Source video codecs whose MSE codec string is real but which no mainstream
@@ -158,7 +164,16 @@ async def probe_video(path: str) -> VideoProbe:
"-of", "json", path,
stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE,
)
- stdout, _ = await proc.communicate()
+ try:
+ stdout, _ = await asyncio.wait_for(proc.communicate(), FFPROBE_TIMEOUT_SECS)
+ except (TimeoutError, asyncio.CancelledError) as e:
+ # Killed, not abandoned: a cancelled wait leaves the process running,
+ # and a caller's own timeout (enrich.py) cancels exactly this wait.
+ proc.kill()
+ await proc.wait()
+ if isinstance(e, asyncio.CancelledError):
+ raise
+ raise RuntimeError(f"ffprobe timed out after {FFPROBE_TIMEOUT_SECS}s") from None
info = json.loads(stdout)
duration = float(info.get("format", {}).get("duration", 0))
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py
index 7ecb0a5..22c1690 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py
@@ -165,6 +165,7 @@ async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> floa
fd, tmp_name = tempfile.mkstemp(suffix=".mp4")
os.close(fd)
tmp_path = Path(tmp_name)
+ proc = probe = None
try:
proc = await asyncio.create_subprocess_exec(
platform.ffmpeg_cmd(), "-hide_banner", "-loglevel", "error", "-y",
@@ -189,6 +190,12 @@ async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> floa
log.warning("stream: seek probe failed at %.1fs: %r", t, e)
return None
finally:
+ # A timed-out wait leaves its process running; it is stopped here, not
+ # left to finish a seek nobody is waiting for.
+ for p in (proc, probe):
+ if p is not None and p.returncode is None:
+ p.kill()
+ await p.wait()
await _discard_scratch(tmp_path)
text = stdout.decode(errors="replace").strip().rstrip(",")
try:
diff --git a/packages/meshbay-node/tests/test_ffprobe_is_bounded.py b/packages/meshbay-node/tests/test_ffprobe_is_bounded.py
new file mode 100644
index 0000000..c96a19a
--- /dev/null
+++ b/packages/meshbay-node/tests/test_ffprobe_is_bounded.py
@@ -0,0 +1,61 @@
+"""
+ffprobe over a member's file is bounded, and a bounded wait stops the process.
+
+`probe_video` runs before every stream and subtitle request and in the
+enrichment pool. A file that keeps ffprobe busy must not hold any of them, and
+giving up on the wait is not enough: an asyncio subprocess whose wait was
+cancelled keeps running. These run a stand-in ffprobe that never answers and
+check both — the call returns, and the process is gone.
+"""
+
+import asyncio
+import os
+import sys
+import time
+
+import pytest
+from meshbay_node import media_probe, platform
+
+pytestmark = pytest.mark.skipif(sys.platform == "win32", reason="a POSIX shell stand-in")
+
+
+@pytest.fixture
+def hanging_ffprobe(tmp_path, monkeypatch):
+ pid_file = tmp_path / "pid"
+ tool = tmp_path / "ffprobe"
+ tool.write_text(f"#!/bin/sh\necho $$ > {pid_file}\nexec sleep 600\n")
+ tool.chmod(0o755)
+ monkeypatch.setattr(platform, "_ffprobe_path", str(tool))
+ return pid_file
+
+
+def _gone(pid: int) -> bool:
+ try:
+ os.kill(pid, 0)
+ except ProcessLookupError:
+ return True
+ # A zombie still answers kill(0); its state says it has exited.
+ try:
+ with open(f"/proc/{pid}/stat") as f:
+ return f.read().split()[2] == "Z"
+ except OSError:
+ return True
+
+
+async def test_a_probe_that_never_answers_times_out_and_is_stopped(hanging_ffprobe, monkeypatch):
+ monkeypatch.setattr(media_probe, "FFPROBE_TIMEOUT_SECS", 0.5)
+ started = time.monotonic()
+ with pytest.raises(RuntimeError, match="timed out"):
+ await media_probe.probe_video("/nonexistent/file.mkv")
+ assert time.monotonic() - started < 5
+ assert _gone(int(hanging_ffprobe.read_text()))
+
+
+async def test_a_caller_giving_up_stops_it_too(hanging_ffprobe):
+ with pytest.raises(TimeoutError):
+ await asyncio.wait_for(media_probe.probe_video("/nonexistent/file.mkv"), 0.5)
+ for _ in range(50):
+ if hanging_ffprobe.exists():
+ break
+ await asyncio.sleep(0.05)
+ assert _gone(int(hanging_ffprobe.read_text()))