summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/indexer/enrich.py
blob: 784b2a34bcfcd345cf9f37b7722c7a82114ecfcc (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
"""
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)