aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/indexer/enrich.py
blob: ae3f3dc243eac495a48b382f767dacf718566d3f (plain) (blame)
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)