Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 27 additions & 2 deletions src/youtube_extension/backend/services/real_video_processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,29 @@ def _write_cache_file(cache_path: Path, payload: dict[str, Any]) -> None:
os.unlink(tmp_name)
raise

@staticmethod
def _count_cached_files(cache_dir: Path) -> int:
"""Count published cache entries under ``cache_dir``.

Blocking: performs a full directory walk. Always run this off the event
loop. It is reached from the status endpoint, so an inline walk stalls
every concurrently-served request for as long as the scan takes — which
grows with the number of cached videos.

``glob`` already yields nothing for a missing directory, so no separate
existence check is needed (it would only add a redundant ``stat`` and
would not make the walk atomic). A genuine filesystem failure — e.g. a
permission error on a directory that does exist — is logged and
re-raised rather than being silently reported as an empty cache. The
count is accumulated lazily rather than materializing the whole listing,
since only the total is ever used.
"""
try:
return sum(1 for _ in cache_dir.glob("*_processed.json"))
except OSError:
logger.exception("Failed to count cached files in %s", cache_dir)
raise

Comment thread
coderabbitai[bot] marked this conversation as resolved.
async def _load_from_cache(self, video_id: str) -> Optional[dict[str, Any]]:
"""Load processed result from cache if available"""
if not self.enable_caching:
Expand Down Expand Up @@ -475,8 +498,10 @@ async def get_processing_status(self) -> dict[str, Any]:
try:
cost_dashboard = await cost_monitor.get_cost_dashboard()

# Count cached files
cached_files = len(list(self.cache_dir.glob("*_processed.json"))) if self.cache_dir.exists() else 0
# Count cached files (off the event loop; the walk is unbounded)
cached_files = await asyncio.to_thread(
self._count_cached_files, self.cache_dir
)

return {
'service_status': 'operational',
Expand Down
78 changes: 78 additions & 0 deletions tests/unit/test_real_processors.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,10 @@

from __future__ import annotations

import asyncio
import json
import sys
import threading
from pathlib import Path
from unittest.mock import AsyncMock, MagicMock, patch

Expand Down Expand Up @@ -1591,6 +1593,82 @@ async def test_error_in_status_returns_error_dict(self, tmp_path):
assert status["service_status"] == "error"
assert "error" in status

async def test_cache_scan_runs_off_the_event_loop(self, tmp_path):
"""The cache-directory scan must not run on the event loop thread.

``get_processing_status`` is served by a live HTTP endpoint, so scanning
the cache directory inline stalls every concurrently-served request for
the duration of the walk.
"""
proc = _make_video_processor(tmp_path)
(proc.cache_dir / "abc_processed.json").write_text("{}")

loop_thread = threading.get_ident()
scan_threads: list[int] = []
cache_dir = proc.cache_dir
real_glob = Path.glob

def recording_glob(self, pattern, *args, **kwargs):
# Only record the cache scan itself; a global patch would otherwise
# intercept unrelated Path.glob calls and validate the patch rather
# than production behavior.
if self == cache_dir and pattern == "*_processed.json":
scan_threads.append(threading.get_ident())
return real_glob(self, pattern, *args, **kwargs)

with patch("youtube_extension.backend.services.real_video_processor.cost_monitor") as cm:
cm.get_cost_dashboard = AsyncMock(return_value={})
with patch.object(Path, "glob", recording_glob):
status = await proc.get_processing_status()

# Guards against a vacuous pass: an unscanned directory would satisfy
# the membership assertion trivially.
assert scan_threads, "cache directory was never scanned"
assert loop_thread not in scan_threads
assert status["cache"]["cached_videos"] == 1

async def test_cache_scan_does_not_stall_the_event_loop(self, tmp_path):
"""The loop keeps scheduling coroutines while the scan is in flight.

The scan blocks until a coroutine running *on the loop* releases it. If
the scan were inline that coroutine could never be scheduled, so the
gather would exceed its timeout instead of completing.
"""
proc = _make_video_processor(tmp_path)
(proc.cache_dir / "abc_processed.json").write_text("{}")

scan_started = threading.Event()
may_finish = threading.Event()
cache_dir = proc.cache_dir
real_glob = Path.glob

def gated_glob(self, pattern, *args, **kwargs):
# Gate only the cache scan; a global patch would otherwise block on
# unrelated Path.glob calls and make the test assert the patch.
if self == cache_dir and pattern == "*_processed.json":
scan_started.set()
may_finish.wait(timeout=10)
return real_glob(self, pattern, *args, **kwargs)

async def release_once_scan_starts():
while not scan_started.is_set():
await asyncio.sleep(0.01)
may_finish.set()

with patch("youtube_extension.backend.services.real_video_processor.cost_monitor") as cm:
cm.get_cost_dashboard = AsyncMock(return_value={})
with patch.object(Path, "glob", gated_glob):
status, _ = await asyncio.wait_for(
asyncio.gather(
proc.get_processing_status(), release_once_scan_starts()
),
timeout=5,
)

assert scan_started.is_set(), "cache directory was never scanned"
assert may_finish.is_set()
assert status["cache"]["cached_videos"] == 1


class TestClose:
async def test_close_calls_youtube_service_close(self, tmp_path):
Expand Down
Loading