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
|
"""
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 = {}
duration: float | None = None
try:
_codec, duration, _has_audio, width, height = 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)
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
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)
|