1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
|
"""
Index-time enrichment for the Videos group app: technical probe (ffprobe),
filename parsing (title_parse), and thumbnail generation (ffmpeg) for a
newly-added video IndexEntry.
Runs through its own small bounded worker pool — separate from the streaming
transcode pool (docs/mediacenter.md §5.2, mirroring webrtc_server.py's
`_transcode_semaphore`) — so indexing a large library never blocks on this,
and enrichment never competes with an active viewer for CPU. The scan itself
already put the entry in the index with hash/size/type only; this fills in
the rest asynchronously and hands the result back via a callback.
"""
import asyncio
import logging
from pathlib import Path
from typing import Awaitable, Callable
import blake3
from meshbay_common.protocol import IndexEntry
from meshbay_node.indexer import title_parse
from meshbay_node.indexer.indexer import MEDIA_EXTENSIONS
from meshbay_node.media_cache import MediaCache
from meshbay_node.media_probe import probe_video
log = logging.getLogger(__name__)
DEFAULT_MAX_CONCURRENT = 2
PROBE_TIMEOUT_SECS = 30
THUMB_TIMEOUT_SECS = 30
THUMB_WIDTH = 320
# "Show/SeasonFolder/episode.mkv" is the expected shape, with a little slack
# for an extra wrapper folder — not an attempt to find the exact group root.
MAX_ANCESTOR_DEPTH = 4
# Bounds the "borrow a title from a sibling episode filename" scan (§3.4) so
# a folder with thousands of files costs a fixed, small amount of work.
MAX_SIBLINGS_CHECKED = 20
def _season_from_ancestors(file_path: Path) -> int | None:
folder = file_path.parent
for _ in range(MAX_ANCESTOR_DEPTH):
if folder is None or folder == folder.parent:
break
season = title_parse.season_from_folder_name(folder.name)
if season is not None:
return season
folder = folder.parent
return None
def _title_from_siblings(file_path: Path) -> str | None:
"""
§3.4: an episode filename with no show name in it borrows the title from
a representative sibling in the same folder, never from the folder name
alone (an acronym-named show folder is a real, observed case).
"""
try:
names = sorted(p.name for p in file_path.parent.iterdir() if p.is_file())
except OSError:
return None
checked = 0
for name in names:
if name == file_path.name:
continue
if Path(name).suffix.lower() not in MEDIA_EXTENSIONS["video"]:
continue
checked += 1
if checked > MAX_SIBLINGS_CHECKED:
break
parsed = title_parse.parse_episode_filename(name)
if parsed.display_title:
return parsed.display_title
return None
async def _make_thumbnail(file_path: Path, duration: float | None) -> bytes | None:
"""One ffmpeg frame grab at ~10% of duration (or 5s if unknown), scaled down."""
seek = max(0.0, (duration or 50.0) * 0.1)
proc = await asyncio.create_subprocess_exec(
"ffmpeg", "-v", "error", "-ss", str(seek), "-i", str(file_path),
"-frames:v", "1", "-vf", f"scale={THUMB_WIDTH}:-1",
"-f", "image2", "-c:v", "mjpeg", "pipe:1",
stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE,
)
try:
stdout, _ = await asyncio.wait_for(proc.communicate(), THUMB_TIMEOUT_SECS)
except asyncio.TimeoutError:
proc.kill()
await proc.wait()
return None
return stdout or None
class Enricher:
"""Owns the node's bounded index-time enrichment pool."""
def __init__(self, media_cache: MediaCache, max_concurrent: int = DEFAULT_MAX_CONCURRENT):
self._media_cache = media_cache
self._sem = asyncio.Semaphore(max_concurrent)
self._tasks: set[asyncio.Task] = set()
def spawn(
self, entry: IndexEntry, file_path: Path,
on_done: Callable[[str, dict], Awaitable[None]],
) -> asyncio.Task:
"""
Fire-and-forget one file's enrichment. `on_done(file_id, fields)` is
awaited with the index fields to merge in once ready — never blocks
the caller (a scan or watchdog event). The reference this method
returns is what keeps the task alive; callers should hold it the
same way `WebRTCPeerSession._spawn` holds streaming tasks.
"""
task = asyncio.ensure_future(self._run(entry, file_path, on_done))
self._tasks.add(task)
def _cleanup(t: asyncio.Task) -> None:
self._tasks.discard(t)
if not t.cancelled() and t.exception():
log.error("Enrichment failed for %s: %s", entry.id[:12], t.exception(),
exc_info=t.exception())
task.add_done_callback(_cleanup)
return task
async def _run(
self, entry: IndexEntry, file_path: Path,
on_done: Callable[[str, dict], Awaitable[None]],
) -> None:
async with self._sem:
fields: dict = {}
# entry.id is the file's own content hash — a probe/thumbnail
# cache hit here means this exact content was already handled
# (this run, an earlier one, even a previous daemon process).
# Neither survives a restart on its own (the in-memory
# GroupIndex entry is rebuilt from scratch every time), but
# media_cache.db does — nothing was checking it before spawning
# ffprobe/ffmpeg again on every file, every restart.
#
# Deliberately *not* cached this way: display_title/season/
# episode. Those come from guessit against entry.name, which is
# exactly what a rename needs re-derived —
# _reenrich_renamed_video_entries exists for precisely that —
# and reusing a stale parse under a new name would silently
# defeat it. ffprobe's own output has no such concern: the same
# bytes probe the same regardless of what the file is called.
cached_meta = await self._media_cache.get_video_meta(entry.id)
duration: float | None = None
if cached_meta is not None:
fields["duration"] = cached_meta["duration"]
fields["width"] = cached_meta["width"]
fields["height"] = cached_meta["height"]
duration = cached_meta["duration"]
else:
try:
_codec, duration, _has_audio, width, height, _raw = await asyncio.wait_for(
probe_video(str(file_path)), timeout=PROBE_TIMEOUT_SECS)
fields["duration"] = int(duration) if duration else None
fields["width"] = width
fields["height"] = height
except Exception as e:
log.warning("Probe failed for %s: %s", file_path, e)
await self._media_cache.put_video_meta(
entry.id, fields.get("duration"), fields.get("width"), fields.get("height"))
ep = title_parse.parse_episode_filename(entry.name)
if ep.episode is not None:
title = ep.display_title or await asyncio.to_thread(
_title_from_siblings, file_path)
season = ep.season
if season is None:
season = await asyncio.to_thread(_season_from_ancestors, file_path)
fields["display_title"] = title or title_parse.naive_title(entry.name)
fields["season"] = season
fields["episode"] = ep.episode
else:
mv = title_parse.parse_movie_filename(entry.name)
fields["display_title"] = mv.display_title or mv.naive_title
cached_thumb_hash = await self._media_cache.get_thumb_hash_by_file_id(entry.id)
if cached_thumb_hash:
fields["thumb_hash"] = cached_thumb_hash
else:
try:
thumb = await _make_thumbnail(file_path, duration)
except Exception as e:
log.warning("Thumbnail generation failed for %s: %s", file_path, e)
thumb = None
if thumb:
thumb_hash = blake3.blake3(thumb).hexdigest()
await self._media_cache.put_thumb(thumb_hash, entry.id, thumb)
fields["thumb_hash"] = thumb_hash
await on_done(entry.id, fields)
|