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
2 changes: 2 additions & 0 deletions backend/app/schemas/domain/addon_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,8 @@ class MediaServerSyncEventMeta(BaseModel):
file_count: int = 0
nfo_count: int = 0
image_count: int = 0
episode_number: int | None = None
episode_numbers: list[int] = Field(default_factory=list)
trigger: str = ""
error: str = ""

Expand Down
1 change: 1 addition & 0 deletions backend/app/schemas/domain/media_server_sync.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ class MediaServerSyncItemResult(BaseModel):
class MediaServerSyncTargetFile(BaseModel):
destination_path: str
episode_number: int | None = None
episode_numbers: list[int] = Field(default_factory=list)


class MediaServerChangeType(str, Enum):
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -164,10 +164,10 @@ def _season_metadata_transfer_results(

@staticmethod
def _unique_transfer_results(results: list[MediaServerSyncTargetFile]) -> list[MediaServerSyncTargetFile]:
seen: set[tuple[str, int | None]] = set()
seen: set[tuple[str, int | None, tuple[int, ...]]] = set()
unique: list[MediaServerSyncTargetFile] = []
for item in results:
key = (item.destination_path, item.episode_number)
key = (item.destination_path, item.episode_number, tuple(item.episode_numbers))
if key in seen:
continue
seen.add(key)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,14 +79,15 @@ async def sync_one_season(
error=str(exc),
)
raise
emit_media_server_sync_events(
EventType.MEDIA_SERVER_SYNC_COMPLETED,
media,
needs.anchor_file or "",
needs.transfer_results,
media_server.id,
trigger="scheduler",
)
if self._should_emit_scheduler_completed_event(needs.missing_flags):
emit_media_server_sync_events(
EventType.MEDIA_SERVER_SYNC_COMPLETED,
media,
needs.anchor_file or "",
needs.transfer_results,
media_server.id,
trigger="scheduler",
)
if sync_cfg.write_nfo:
await media_server_sync_artifacts.mark_nfo_artifacts(
files,
Expand Down Expand Up @@ -127,5 +128,9 @@ async def apply_updates(
change_type=MediaServerChangeType.UPDATED,
)

@staticmethod
def _should_emit_scheduler_completed_event(missing_flags: list[str]) -> bool:
return set(missing_flags) != {"stale"}


media_server_sync_season_runner = MediaServerSyncSeasonRunner()
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ async def handle_import_completed(self, event: Event) -> None:
MediaServerSyncTargetFile(
destination_path=item.destination_path,
episode_number=item.episode_number,
episode_numbers=item.episode_numbers,
)
for item in context.imported_files
]
Expand Down Expand Up @@ -335,6 +336,7 @@ def _rerun_target_files(self, library_files: list[LibraryFile]) -> list[MediaSer
MediaServerSyncTargetFile(
destination_path=str(build_library_file_path(library_file.path, library_file.file_name)),
episode_number=self._episode_number_for_library_file(library_file),
episode_numbers=self._episode_numbers_for_library_file(library_file),
)
for library_file in library_files
if library_service.is_primary_file(library_file) and library_service.file_exists(library_file)
Expand All @@ -357,6 +359,12 @@ def _episode_number_for_library_file(library_file: LibraryFile) -> int | None:
episodes = list(attrs.episodes or []) if attrs else []
return int(episodes[0]) if len(episodes) == 1 else None

@staticmethod
def _episode_numbers_for_library_file(library_file: LibraryFile) -> list[int]:
attrs = library_file.resource_attributes
episodes = list(attrs.episodes or []) if attrs else []
return sorted({int(episode) for episode in episodes if int(episode) > 0})

async def _run_server_once(
self,
media_server: JellyfinConfig,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ def build_input(self, media: MediaFullInfo, layout: LibraryMediaLayout) -> Media
MediaServerSyncTargetFile(
destination_path=target.destination_path,
episode_number=target.episode_number,
episode_numbers=[target.episode_number] if target.episode_number else [],
)
for target in decision.target_files
],
Expand Down
16 changes: 16 additions & 0 deletions backend/app/services/audit/workflow_event_emitters.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ def emit_media_server_sync_events(
target_paths = _sync_event_paths(anchor_file, transfer_results)
nfo_count, image_count = _sync_artifact_counts(media, anchor_file, transfer_results)
for path in target_paths:
episode_numbers = _sync_event_episode_numbers(path, transfer_results) if media.media_type.value == "tv" else []
event_service.emit_media(
MediaEventCreate(
type=event_type,
Expand All @@ -95,6 +96,8 @@ def emit_media_server_sync_events(
file_count=len(target_paths),
nfo_count=nfo_count,
image_count=image_count,
episode_number=episode_numbers[0] if len(episode_numbers) == 1 else None,
episode_numbers=episode_numbers,
Comment thread
120318 marked this conversation as resolved.
trigger=trigger,
error=error,
),
Expand All @@ -108,6 +111,17 @@ def _sync_event_paths(anchor_file: str, transfer_results: list[MediaServerSyncTa
return sorted(set(paths))


def _sync_event_episode_numbers(path: str, transfer_results: list[MediaServerSyncTargetFile]) -> list[int]:
episodes: set[int] = set()
for item in transfer_results:
if item.destination_path != path:
continue
for episode in [item.episode_number, *item.episode_numbers]:
if episode and int(episode) > 0:
episodes.add(int(episode))
return sorted(episodes)


def _sync_artifact_counts(
media: MediaFullInfo,
anchor_file: str,
Expand Down Expand Up @@ -164,4 +178,6 @@ def _expected_nfo_paths(
paths.add(season_dir / "season.nfo")
if target.episode_number:
paths.add(target_path.with_suffix(".nfo"))
if target.episode_numbers:
paths.add(target_path.with_suffix(".nfo"))
return paths
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ class TelegramEventMeta(BaseModel):
selected_episodes: list[int] = Field(default_factory=list)
imported_files: list[TelegramImportedFileMeta] = Field(default_factory=list)
episode_number: int | None = None
episode_numbers: list[int] = Field(default_factory=list)

@field_validator(
"resource_title",
Expand Down Expand Up @@ -126,6 +127,8 @@ def _event_episodes(event: Event, meta: TelegramEventMeta) -> list[int]:
return _download_episodes(meta)
if event.type == EventType.MEDIA_IMPORT_COMPLETED:
return _imported_episodes(meta)
if event.type in {EventType.MEDIA_SERVER_SYNC_COMPLETED, EventType.MEDIA_SERVER_SYNC_FAILED}:
return _episode_range([meta.episode_number, *meta.episode_numbers])
if event.type in {EventType.DANMU_GENERATE_COMPLETED, EventType.DANMU_GENERATE_FAILED}:
return _episode_range([meta.episode_number])
return []
Expand Down
201 changes: 199 additions & 2 deletions backend/tests/test_media_server_sync_boundaries.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,16 +3,19 @@

import pytest

from app.schemas.config import MediaServerSyncConfig
from app.schemas.config import JellyfinConfig, MediaServerSyncConfig
from app.schemas.domain.library import LibraryFile, LibraryFileArtifactStatus, LibraryMediaLayout, LibraryMediaLayoutEntry
from app.schemas.domain.event import EventType
from app.schemas.domain.media import EpisodeInfo, MediaFullInfo, SeasonDetails
from app.schemas.domain.media_context import MediaCapabilities
from app.schemas.domain.media_server_sync import MediaServerSyncState
from app.schemas.domain.media_server_sync import MediaServerSyncDetectNeeds, MediaServerSyncState
from app.schemas.domain.media_server_sync import MediaServerSyncTargetFile
from app.schemas.domain.media_types import MediaType
from app.schemas.media_id import MediaID
from app.services.application.workflows.media_server_sync.artifacts import media_server_sync_artifacts
from app.services.application.workflows.media_server_sync.needs import media_server_sync_needs
from app.services.application.workflows.media_server_sync.season_runner import MediaServerSyncSeasonRunner
from app.services.audit.workflow_event_emitters import emit_media_server_sync_events
from app.services.domain.library.sidecar_files import library_sidecar_files
from app.services.integration.tmdb.images import to_tmdb_image_url

Expand Down Expand Up @@ -727,3 +730,197 @@ async def test_media_server_sync_needs_uses_profile_refresh_tiers_for_stale(tmp_
assert not cold_needs.should_run
assert hot_needs.should_run
assert hot_needs.missing_flags == ["stale"]


@pytest.mark.asyncio
async def test_media_server_sync_scheduler_suppresses_completed_event_for_pure_stale(monkeypatch, tmp_path: Path):
media = _tv()
video = tmp_path / "Show" / "Season 01" / "Show.S01E01.mkv"
video.parent.mkdir(parents=True)
video.write_text("video")
files = [
LibraryFile(
id="file-1",
task_id="task-1",
directory_id="dir-1",
media_id=media.media_id,
path=str(video.parent),
file_name=video.name,
created_at=1.0,
)
]
target = MediaServerSyncTargetFile(destination_path=str(video), episode_number=None)
needs = MediaServerSyncDetectNeeds(
should_run=True,
missing_flags=["stale"],
transfer_results=[target],
anchor_file=str(video),
media_root_dir=str(tmp_path / "Show"),
)
runner = MediaServerSyncSeasonRunner()
events = []
applied = []

async def fake_info(media_id, season_number=None):
return media

async def fake_layout(media_id, library_files):
return LibraryMediaLayout(media_id=media_id, media_type=media.media_type, primary_anchor_file=str(video))

async def fake_detect(*args, **kwargs):
return needs

async def fake_apply_updates(*args, **kwargs):
applied.append(args)

monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.media_service.info", fake_info)
monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.library_service.get_media_layout_for_files", fake_layout)
monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.media_server_sync_needs.detect", fake_detect)
monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.media_server_sync_state.record_success", lambda *args, **kwargs: None)
monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.emit_media_server_sync_events", lambda *args, **kwargs: events.append((args, kwargs)))
monkeypatch.setattr(runner, "apply_updates", fake_apply_updates)

result = await runner.sync_one_season(
JellyfinConfig(id="jellyfin-1", name="Jellyfin", url="http://jellyfin", api_key="token"),
media.media_id,
1,
files,
_state(media.media_id, last_success_at=time.time() - 8 * 86400),
time.time(),
MediaServerSyncConfig(write_nfo=False, download_images=False),
)

assert result
assert applied
assert events == []


@pytest.mark.asyncio
async def test_media_server_sync_scheduler_emits_completed_event_for_missing_artifacts(monkeypatch, tmp_path: Path):
media = _tv()
video = tmp_path / "Show" / "Season 01" / "Show.S01E01.mkv"
video.parent.mkdir(parents=True)
video.write_text("video")
files = [
LibraryFile(
id="file-1",
task_id="task-1",
directory_id="dir-1",
media_id=media.media_id,
path=str(video.parent),
file_name=video.name,
created_at=1.0,
)
]
target = MediaServerSyncTargetFile(destination_path=str(video), episode_number=1)
needs = MediaServerSyncDetectNeeds(
should_run=True,
missing_flags=["episode_nfo_missing"],
transfer_results=[target],
anchor_file=str(video),
media_root_dir=str(tmp_path / "Show"),
)
runner = MediaServerSyncSeasonRunner()
events = []

async def fake_info(media_id, season_number=None):
return media

async def fake_layout(media_id, library_files):
return LibraryMediaLayout(media_id=media_id, media_type=media.media_type, primary_anchor_file=str(video))

async def fake_detect(*args, **kwargs):
return needs

async def fake_apply_updates(*args, **kwargs):
return None

monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.media_service.info", fake_info)
monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.library_service.get_media_layout_for_files", fake_layout)
monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.media_server_sync_needs.detect", fake_detect)
monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.media_server_sync_state.record_success", lambda *args, **kwargs: None)
monkeypatch.setattr("app.services.application.workflows.media_server_sync.season_runner.emit_media_server_sync_events", lambda *args, **kwargs: events.append((args, kwargs)))
monkeypatch.setattr(runner, "apply_updates", fake_apply_updates)

result = await runner.sync_one_season(
JellyfinConfig(id="jellyfin-1", name="Jellyfin", url="http://jellyfin", api_key="token"),
media.media_id,
1,
files,
_state(media.media_id, last_success_at=time.time()),
time.time(),
MediaServerSyncConfig(write_nfo=False, download_images=False),
)

assert result
assert len(events) == 1
assert events[0][1]["trigger"] == "scheduler"


def test_media_server_sync_event_meta_includes_target_episodes(monkeypatch, tmp_path: Path):
media = _tv()
video = tmp_path / "Show" / "Season 01" / "Show.S01E03-E04.mkv"
video.parent.mkdir(parents=True)
video.write_text("video")
emitted = []

def fake_emit_media(event, meta=None):
emitted.append((event, meta))

monkeypatch.setattr("app.services.audit.workflow_event_emitters.event_service.emit_media", fake_emit_media)

emit_media_server_sync_events(
EventType.MEDIA_SERVER_SYNC_COMPLETED,
media,
str(video),
[
MediaServerSyncTargetFile(
destination_path=str(video),
episode_numbers=[3, 4],
)
],
"jellyfin-1",
trigger="import",
)

assert len(emitted) == 1
event, meta = emitted[0]
assert event.type == EventType.MEDIA_SERVER_SYNC_COMPLETED
assert meta.file_path == str(video)
assert meta.episode_number is None
assert meta.episode_numbers == [3, 4]
assert meta.trigger == "import"


def test_media_server_sync_event_meta_ignores_movie_episode_attributes(monkeypatch, tmp_path: Path):
media = _movie()
video = tmp_path / "Movie" / "Movie.2026.mkv"
video.parent.mkdir(parents=True)
video.write_text("video")
emitted = []

def fake_emit_media(event, meta=None):
emitted.append((event, meta))

monkeypatch.setattr("app.services.audit.workflow_event_emitters.event_service.emit_media", fake_emit_media)

emit_media_server_sync_events(
EventType.MEDIA_SERVER_SYNC_COMPLETED,
media,
str(video),
[
MediaServerSyncTargetFile(
destination_path=str(video),
episode_number=1,
episode_numbers=[1],
)
],
"jellyfin-1",
trigger="manual",
)

assert len(emitted) == 1
_event, meta = emitted[0]
assert meta.file_path == str(video)
assert meta.episode_number is None
assert meta.episode_numbers == []
Loading
Loading