import 'dart:async'; import '../media/ids.dart'; import '../media/media_hub.dart'; import '../media/media_item.dart'; import '../media/media_kind.dart'; import '../media/media_library.dart'; import '../media/media_server_client.dart'; import '../exceptions/media_server_exceptions.dart'; import '../utils/app_logger.dart'; import '../utils/external_ids.dart'; import '../utils/global_key_utils.dart'; import '../utils/search_relevance.dart'; import 'local_playback_history.dart'; import 'multi_server_manager.dart'; typedef OnDeckAggregationResult = ({ List items, Set succeededServerIds, Set cancelledServerIds, }); typedef HubAggregationResult = ({List hubs, Set succeededServerIds, Set cancelledServerIds}); typedef LibraryAggregationResult = ({ List libraries, Set succeededServerIds, Set cancelledServerIds, }); typedef _FanOutResult = ({List items, Set succeededServerIds, Set cancelledServerIds}); /// Whether [error] is a client-side abort (client teardown mid-request) /// rather than a genuine server failure. Aggregation reports these servers /// in `cancelledServerIds` so callers can tell a *disrupted* pass — whose /// results say nothing about actual content — from a settled failure. bool _isCancellation(Object error) => error is MediaServerHttpException && error.isCancellation; /// Cross-server aggregation: fans calls out to every online client and /// merges the results. Single-server operations now go through the /// [MediaServerClient] interface directly (resolved via /// [ProviderExtensions.tryGetMediaClientForServer] etc.), so this service /// only owns the genuinely multi-server flows: home/discover hubs, on-deck, /// search, and the global library list. class DataAggregationService { final MultiServerManager _serverManager; DataAggregationService(this._serverManager); /// Online clients, optionally restricted to [serverIds] — delta refreshes /// fan out to newly-online servers only. Map _clientsFor(Set? serverIds) { final clients = _serverManager.onlineClients; if (serverIds == null) return clients; return { for (final entry in clients.entries) if (serverIds.contains(entry.key)) entry.key: entry.value, }; } /// Run [fetch] against every client in [clients] and concatenate the results /// in client order. A per-server failure is swallowed — logged with /// [failureMessage] and contributing nothing — so one bad server cannot sink /// the pass; that server is simply absent from `succeededServerIds`, and also /// lands in `cancelledServerIds` when the failure was a client-side abort. Future<_FanOutResult> _fanOut( Map clients, { required String Function(String serverId) failureMessage, required Future> Function(String serverId, MediaServerClient client) fetch, }) async { final cancelledServerIds = {}; final futures = clients.entries.map((entry) async { try { return (serverId: entry.key, items: await fetch(entry.key, entry.value)); } catch (e, stackTrace) { if (_isCancellation(e)) cancelledServerIds.add(entry.key); appLogger.e(failureMessage(entry.key), error: e, stackTrace: stackTrace); return (serverId: null, items: []); } }); final results = await Future.wait(futures); return ( items: [for (final result in results) ...result.items], succeededServerIds: { for (final result in results) if (result.serverId != null) result.serverId!, }, cancelledServerIds: cancelledServerIds, ); } /// Fetch libraries from all online clients regardless of backend, returning /// the merged neutral [MediaLibrary]s alongside the ids of the servers whose /// fetch actually succeeded. [serverIds] restricts the fan-out to those /// servers. /// /// A per-server `fetchLibraries()` failure is swallowed (that server simply /// contributes no libraries) so one unreachable server doesn't sink the whole /// list. [succeededServerIds] lets callers tell a *failed* fetch apart from a /// server that genuinely has no libraries — both contribute nothing, so /// conflating them would let a transient failure be cached as "loaded" and /// never retried. Servers whose fetch was aborted client-side land in /// `cancelledServerIds` — a disrupted pass, unlike a settled failure, must /// never be committed as authoritative. Future getMediaLibrariesFromAllServers({Set? serverIds}) async { final clients = _clientsFor(serverIds); if (clients.isEmpty) { appLogger.w('No online servers available for fetching libraries (neutral)'); return ( libraries: const [], succeededServerIds: const {}, cancelledServerIds: const {}, ); } final fetched = await _fanOut( clients, failureMessage: (serverId) => 'Failed neutral library fetch from $serverId', fetch: (_, client) => client.fetchLibraries(), ); return ( libraries: fetched.items, succeededServerIds: fetched.succeededServerIds, cancelledServerIds: fetched.cancelledServerIds, ); } /// Fetch "On Deck" (Continue Watching) from all servers and merge by recency. /// Items are tagged with server info by the underlying client. Returns /// neutral [MediaItem]s plus the ids of servers whose fetch succeeded. /// [serverIds] restricts the fan-out to those servers. Future getOnDeckFromAllServers({ int? limit, Set? hiddenLibraryKeys, Set? serverIds, }) async { final clients = _clientsFor(serverIds); if (clients.isEmpty) { appLogger.w('No online servers available for fetching on deck'); return (items: const [], succeededServerIds: const {}, cancelledServerIds: const {}); } final fetched = await _fanOut( clients, failureMessage: (serverId) => 'Failed on-deck fetch from $serverId', fetch: (_, client) => client.fetchContinueWatching(count: limit), ); final allOnDeck = fetched.items; // Filter out items from hidden libraries List filteredOnDeck = allOnDeck; if (hiddenLibraryKeys != null && hiddenLibraryKeys.isNotEmpty) { filteredOnDeck = allOnDeck.where((item) { if (item.libraryId == null || item.serverId == null) return true; final globalKey = buildGlobalKey(ServerId(item.serverId!), item.libraryId!); return !hiddenLibraryKeys.contains(globalKey); }).toList(); } // Sort by most recently viewed, falling back to addedAt for unwatched items. // Same key as JellyfinClient's continue-watching merge (MediaItem.recencySortKey) // so per-server and cross-server ordering can't drift apart. filteredOnDeck.sort((a, b) => b.recencySortKey.compareTo(a.recencySortKey)); filteredOnDeck = await _deduplicateContinueWatching(filteredOnDeck); // Apply limit if specified final items = limit != null && limit < filteredOnDeck.length ? filteredOnDeck.sublist(0, limit) : filteredOnDeck; appLogger.i('Fetched ${items.length} on deck items from all servers'); return ( items: items, succeededServerIds: fetched.succeededServerIds, cancelledServerIds: fetched.cancelledServerIds, ); } /// Merge an [existing] Continue Watching list with [fresh] rows from /// newly-online servers: same recency ordering and cross-server identity /// dedup as [getOnDeckFromAllServers], applied to the union. Future> mergeContinueWatching(List existing, List fresh, {int? limit}) async { final combined = [...existing, ...fresh]..sort((a, b) => b.recencySortKey.compareTo(a.recencySortKey)); final deduped = await _deduplicateContinueWatching(combined); return limit != null && limit < deduped.length ? deduped.sublist(0, limit) : deduped; } Future> _deduplicateContinueWatching(List items) async { if (items.length < 2) return items; final bucketCounts = {}; for (final item in items) { final bucket = _continueWatchingTitleBucket(item); if (bucket == null) continue; bucketCounts[bucket] = (bucketCounts[bucket] ?? 0) + 1; } final duplicateBuckets = { for (final entry in bucketCounts.entries) if (entry.value > 1) entry.key, }; if (duplicateBuckets.isEmpty) return items; final externalIdLoads = >{}; final identityKeysByIndex = >{}; final identityKeyLoads = >[]; for (var i = 0; i < items.length; i++) { if (!duplicateBuckets.contains(_continueWatchingTitleBucket(items[i]))) continue; final index = i; identityKeyLoads.add( _continueWatchingIdentityKeys(items[index], externalIdLoads).then((keys) => identityKeysByIndex[index] = keys), ); } await Future.wait(identityKeyLoads); // Group duplicates instead of greedily dropping them: the first item to // claim an identity key anchors the group and holds its shelf slot; // later items sharing a claimed key join as members without claiming // their own keys (same transitive semantics as the old drop). Each slot // then shows the member the user most recently played on this device — // servers sync watch state across guid-linked siblings, so their // lastViewedAt ties and can't tell the 4K copy from the 1080p one // (#1492). Without local history the anchor (recency order) stands. final keyToGroup = {}; final groups = >[]; final groupSlots = {}; final result = []; for (var i = 0; i < items.length; i++) { final item = items[i]; if (!duplicateBuckets.contains(_continueWatchingTitleBucket(item))) { result.add(item); continue; } final identityKeys = identityKeysByIndex[i] ?? const {}; if (identityKeys.isEmpty) { result.add(item); continue; } var joined = false; for (final key in identityKeys) { final groupIndex = keyToGroup[key]; if (groupIndex != null) { groups[groupIndex].add(item); joined = true; break; } } if (joined) continue; final groupIndex = groups.length; groups.add([item]); for (final key in identityKeys) { keyToGroup[key] = groupIndex; } groupSlots[result.length] = groupIndex; result.add(item); } if (groupSlots.isEmpty) return result; final lastPlayed = await LocalPlaybackHistory.snapshot(); for (final slot in groupSlots.entries) { final members = groups[slot.value]; if (members.length > 1) { result[slot.key] = _preferLocallyLastPlayed(members, lastPlayed); } } return result; } /// The duplicate-group member most recently played on this device (by item /// or series key), or the anchor — `members.first`, the group's most recent /// item by [MediaItem.recencySortKey] — when the local history has nothing /// newer to say. MediaItem _preferLocallyLastPlayed(List members, Map lastPlayed) { var winner = members.first; var winnerLastPlayedAt = 0; for (final member in members) { final itemTs = lastPlayed[member.globalKey] ?? 0; final seriesKey = member.seriesGlobalKey; final seriesTs = seriesKey != null ? (lastPlayed[seriesKey] ?? 0) : 0; final lastPlayedAt = itemTs > seriesTs ? itemTs : seriesTs; if (lastPlayedAt > winnerLastPlayedAt) { winner = member; winnerLastPlayedAt = lastPlayedAt; } } return winner; } String? _continueWatchingTitleBucket(MediaItem item) { final scope = _continueWatchingIdentityScope(item); if (scope == null) return null; final title = switch (item.kind) { MediaKind.episode || MediaKind.season => item.grandparentTitle ?? item.parentTitle ?? item.title, _ => item.title, }; final normalized = title?.trim().toLowerCase().replaceAll(RegExp(r'\s+'), ' '); if (normalized == null || normalized.isEmpty) return null; return '$scope:$normalized'; } Future> _continueWatchingIdentityKeys( MediaItem item, Map> externalIdLoads, ) async { final scope = _continueWatchingIdentityScope(item); if (scope == null) return const {}; final keys = {}; final serverId = item.serverId; final targetId = _continueWatchingIdentityTargetId(item); final client = serverId == null ? null : _serverManager.getClient(ServerId(serverId)); if (client != null && targetId != null && targetId.isNotEmpty) { try { final cacheKey = buildGlobalKey(ServerId(serverId!), targetId); final externalIds = await externalIdLoads.putIfAbsent(cacheKey, () => client.fetchExternalIds(targetId)); _addExternalIdentityKeys(keys, scope, externalIds); } catch (e, stackTrace) { appLogger.d( 'Failed to resolve Continue Watching identity for ${item.globalKey}', error: e, stackTrace: stackTrace, ); } } final stableGuid = _stableMediaGuid(item.guid); if (stableGuid != null) { final guidScope = item.kind == MediaKind.episode ? 'episode' : scope; keys.add('$guidScope:guid:$stableGuid'); } return keys; } String? _continueWatchingIdentityScope(MediaItem item) { return switch (item.kind) { MediaKind.episode || MediaKind.season || MediaKind.show => 'show', MediaKind.movie => 'movie', _ => null, }; } String? _continueWatchingIdentityTargetId(MediaItem item) { return switch (item.kind) { MediaKind.episode => item.grandparentId, MediaKind.season => item.grandparentId ?? item.parentId, MediaKind.show || MediaKind.movie => item.id, _ => null, }; } void _addExternalIdentityKeys(Set keys, String scope, ExternalIds externalIds) { final imdb = externalIds.imdb?.trim().toLowerCase(); if (imdb != null && imdb.isNotEmpty) keys.add('$scope:imdb:$imdb'); final tmdb = externalIds.tmdb; if (tmdb != null) keys.add('$scope:tmdb:$tmdb'); final tvdb = externalIds.tvdb; if (tvdb != null) keys.add('$scope:tvdb:$tvdb'); } String? _stableMediaGuid(String? guid) { final value = guid?.trim(); if (value == null || value.isEmpty) return null; if (!value.contains('://')) return null; if (value.contains('agents.none://')) return null; return value.toLowerCase(); } /// Fetch recommendation hubs from all servers as neutral [MediaHub]s. /// When useGlobalHubs is true (default), rich-hub backends use their true /// home page hubs (Plex's promoted/global hub endpoint). /// Backends without rich home hubs fall back to per-library hubs so one /// capped "Latest" response cannot hide whole library types. /// [serverIds] restricts the fan-out (including the library prefetch) to /// those servers. Returns the ids of servers whose hub fetch succeeded so /// callers do not cache transient per-server failures as loaded. Future getHubsFromAllServers({ int? limit, Set? hiddenLibraryKeys, bool useGlobalHubs = true, bool includePlaybackHubs = true, Set? serverIds, }) async { final clients = _clientsFor(serverIds); if (clients.isEmpty) { appLogger.w('No online servers available for fetching hubs'); return (hubs: const [], succeededServerIds: const {}, cancelledServerIds: const {}); } // Home layout needs the library list for every client: fallback backends // build all their rows from per-library hubs, and rich-hub backends // (Plex) need it to detect visible music libraries, whose hubs the // global-hub endpoint excludes. One `fetchLibraries` per server, served // from the per-backend API cache when warm. final libraries = useGlobalHubs ? _groupLibrariesByServer((await getMediaLibrariesFromAllServers(serverIds: serverIds)).libraries) : null; final fetched = await _fanOut( clients, failureMessage: (serverId) => 'Failed to fetch hubs from server $serverId', fetch: (serverId, client) async { final serverLibraries = libraries?[serverId]; final shouldUseGlobalHubs = useGlobalHubs && client.capabilities.richHubs; final hubItemLimit = limit ?? defaultHubPreviewLimit; final hubs = shouldUseGlobalHubs ? [ ...await client.fetchGlobalHubs(limit: hubItemLimit, includePlaybackHubs: includePlaybackHubs), // Plex's promoted/global hub endpoint never includes music // libraries — append their per-library hubs so music rows // reach home. No-op (zero extra calls) without a visible // music library. ...await _fetchLibraryHubsForClient( client, limit: hubItemLimit, hiddenLibraryKeys: hiddenLibraryKeys, includePlaybackHubs: includePlaybackHubs, libraries: serverLibraries ?? const [], kinds: const {MediaKind.artist}, ), ] : await _fetchLibraryHubsForClient( client, limit: hubItemLimit, hiddenLibraryKeys: hiddenLibraryKeys, includePlaybackHubs: includePlaybackHubs, libraries: useGlobalHubs ? serverLibraries : null, ); return _postProcessHubs(hubs, serverId: ServerId(serverId), hiddenLibraryKeys: hiddenLibraryKeys); }, ); final all = fetched.items; final hubs = limit != null && limit < all.length ? all.sublist(0, limit) : all; return (hubs: hubs, succeededServerIds: fetched.succeededServerIds, cancelledServerIds: fetched.cancelledServerIds); } /// Per-library hub fetch for a single client. Filters to visible libraries /// of [kinds] (movie/show/clip/artist by default — clip covers Jellyfin /// musicvideos/homevideos, #1476; artist brings music rows to home) and /// concatenates the results. The rich-hub music append passes /// `{MediaKind.artist}` to fetch only what the global endpoint misses. Future> _fetchLibraryHubsForClient( MediaServerClient client, { required int limit, Set? hiddenLibraryKeys, required bool includePlaybackHubs, List? libraries, Set kinds = const {MediaKind.movie, MediaKind.show, MediaKind.clip, MediaKind.artist}, }) async { final libs = libraries ?? await client.fetchLibraries(); final visible = libs.where((l) { if (!kinds.contains(l.kind)) return false; if (l.hidden) return false; if (hiddenLibraryKeys != null && hiddenLibraryKeys.contains(l.globalKey)) return false; return true; }).toList(); const concurrency = 3; final all = []; for (var start = 0; start < visible.length; start += concurrency) { final batch = visible.skip(start).take(concurrency); final results = await Future.wait( batch.map((l) async { try { return await client.fetchLibraryHubs( l.id, libraryName: l.title, limit: limit, includePlaybackHubs: includePlaybackHubs, libraryKind: l.kind, ); } catch (e, st) { appLogger.e('Failed to fetch library hubs for ${l.globalKey}', error: e, stackTrace: st); return []; } }), ); for (final list in results) { all.addAll(list); } } return all; } /// Filter hidden-library items and drop empty hubs. List _postProcessHubs(List hubs, {required ServerId serverId, Set? hiddenLibraryKeys}) { var filtered = hubs; if (hiddenLibraryKeys != null && hiddenLibraryKeys.isNotEmpty) { filtered = filtered .map((hub) { final filteredItems = hub.items.where((item) { final libraryId = item.libraryId; if (libraryId == null) return true; final globalKey = buildGlobalKey(ServerId(serverId), libraryId); return !hiddenLibraryKeys.contains(globalKey); }).toList(); if (filteredItems.isEmpty) return null; return hub.copyWith(items: filteredItems, size: filteredItems.length); }) .whereType() .toList(); } return filtered; } /// Search across all online servers (Plex + Jellyfin). Returns neutral /// [MediaItem]s. Future> searchAcrossServers(String query, {int? limit}) async { if (query.trim().isEmpty) { return []; } final clients = _serverManager.onlineClients; if (clients.isEmpty) return []; final resultLimit = limit ?? defaultMediaSearchLimit; final fetchLimit = resultLimit < defaultMediaSearchLimit ? defaultMediaSearchLimit : resultLimit; final fetched = await _fanOut( clients, failureMessage: (serverId) => 'Search failed on $serverId', fetch: (_, client) => client.searchItems(query, limit: fetchLimit), ); final result = rankMediaSearchResults(fetched.items, query, limit: resultLimit); appLogger.i('Found ${result.length} search results across all servers'); return result; } /// Reverse external-id lookup fanned out to every online server (see /// [MediaServerClient.findByExternalIds]). One request wave per tap on an /// Explore catalog item; per-server failures are logged and skipped. Future> findByExternalIdsAcrossServers( ExternalIds ids, { required MediaKind kind, String? title, int? year, }) async { if (!ids.hasAny) return []; final clients = _serverManager.onlineClients; if (clients.isEmpty) return []; final futures = clients.entries.map((entry) async { try { return await entry.value.findByExternalIds(ids, kind: kind, title: title, year: year); } catch (e, st) { appLogger.w('External-id lookup failed on ${entry.key}', error: e, stackTrace: st); return null; } }); return (await Future.wait(futures)).nonNulls.toList(); } /// Group libraries by server (internal aggregation helper). Map> _groupLibrariesByServer(List libraries) { final grouped = >{}; for (final library in libraries) { final serverId = library.serverId; if (serverId != null) { grouped.putIfAbsent(serverId, () => []).add(library); } } return grouped; } }