import 'dart:async'; import '../i18n/strings.g.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 '../utils/app_logger.dart'; import '../utils/global_key_utils.dart'; import 'multi_server_manager.dart'; /// 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); /// Fetch libraries from all online clients regardless of backend, returning /// neutral [MediaLibrary]s. Future> getMediaLibrariesFromAllServers() async { final clients = _serverManager.onlineClients; if (clients.isEmpty) { appLogger.w('No online servers available for fetching libraries (neutral)'); return []; } final futures = clients.entries.map((entry) async { try { return await entry.value.fetchLibraries(); } catch (e, stackTrace) { appLogger.e('Failed neutral library fetch from ${entry.key}', error: e, stackTrace: stackTrace); return []; } }); final results = await Future.wait(futures); return [for (final list in results) ...list]; } /// 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. Future> getOnDeckFromAllServers({int? limit, Set? hiddenLibraryKeys}) async { final clients = _serverManager.onlineClients; if (clients.isEmpty) { appLogger.w('No online servers available for fetching on deck'); return []; } final futures = clients.entries.map((entry) async { final client = entry.value; try { return await client.fetchContinueWatching(); } catch (e, st) { appLogger.e('Failed on-deck fetch from ${entry.key}', error: e, stackTrace: st); return []; } }); final allOnDeck = (await Future.wait(futures)).expand((l) => l).toList(); // 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(item.serverId!, item.libraryId!); return !hiddenLibraryKeys.contains(globalKey); }).toList(); } // Sort by most recently viewed, falling back to addedAt for unwatched items filteredOnDeck.sort((a, b) { final aTime = a.lastViewedAt ?? a.addedAt ?? 0; final bTime = b.lastViewedAt ?? b.addedAt ?? 0; return bTime.compareTo(aTime); // Descending (most recent first) }); // Apply limit if specified final result = limit != null && limit < filteredOnDeck.length ? filteredOnDeck.sublist(0, limit) : filteredOnDeck; appLogger.i('Fetched ${result.length} on deck items from all servers'); return result; } /// Fetch recommendation hubs from all servers as neutral [MediaHub]s. /// When useGlobalHubs is true (default), uses the global /hubs endpoint /// to get the true home page hubs like "Recently Added Movies", "Recently /// Added TV"; when false, uses per-library hubs from /// /hubs/sections/{sectionId}. Future> getHubsFromAllServers({ int? limit, Set? hiddenLibraryKeys, bool useGlobalHubs = true, }) async { final clients = _serverManager.onlineClients; if (clients.isEmpty) { appLogger.w('No online servers available for fetching hubs'); return []; } // For global hubs, pre-fetch libraries to split "Recently Added" hubs // by library and resolve human-readable names. final libraries = useGlobalHubs ? _groupLibrariesByServer(await getMediaLibrariesFromAllServers()) : null; final futures = clients.entries.map((entry) async { final serverId = entry.key; final client = entry.value; try { final hubs = useGlobalHubs ? await client.fetchGlobalHubs(limit: limit ?? 10) : await _fetchLibraryHubsForClient(client, limit: limit ?? 10, hiddenLibraryKeys: hiddenLibraryKeys); return _postProcessHubs( hubs, serverId: serverId, hiddenLibraryKeys: hiddenLibraryKeys, libraries: libraries?[serverId], splitRecentlyAdded: useGlobalHubs, ); } catch (e, stackTrace) { appLogger.e('Failed to fetch hubs from server $serverId', error: e, stackTrace: stackTrace); return []; } }); final results = await Future.wait(futures); final all = []; for (final list in results) { all.addAll(list); } return limit != null && limit < all.length ? all.sublist(0, limit) : all; } /// Per-library hub fetch for a single client. Filters to visible /// movie/show libraries (Plex hides music libraries from this surface) and /// concatenates the results. Future> _fetchLibraryHubsForClient( MediaServerClient client, { required int limit, Set? hiddenLibraryKeys, }) async { final libs = await client.fetchLibraries(); final visible = libs.where((l) { if (l.kind != MediaKind.movie && l.kind != MediaKind.show) return false; if (l.hidden) return false; if (hiddenLibraryKeys != null && hiddenLibraryKeys.contains(l.globalKey)) return false; return true; }); final futures = visible.map((l) => client.fetchLibraryHubs(l.id, libraryName: l.title, limit: limit)); final results = await Future.wait(futures); return [for (final list in results) ...list]; } /// Filter hidden-library items, optionally split multi-library "Recently /// Added" hubs by section, and drop empty hubs. List _postProcessHubs( List hubs, { required String serverId, Set? hiddenLibraryKeys, List? libraries, required bool splitRecentlyAdded, }) { 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, libraryId); return !hiddenLibraryKeys.contains(globalKey); }).toList(); if (filteredItems.isEmpty) return null; return hub.copyWith(items: filteredItems, size: filteredItems.length); }) .whereType() .toList(); } if (splitRecentlyAdded) { filtered = _splitRecentlyAddedHubs(filtered, libraries); } 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 futures = clients.entries.map((entry) async { final client = entry.value; try { return await client.searchItems(query, limit: limit ?? 30); } catch (e, st) { appLogger.e('Search failed on ${entry.key}', error: e, stackTrace: st); return []; } }); final allResults = (await Future.wait(futures)).expand((l) => l).toList(); final result = limit != null && limit < allResults.length ? allResults.sublist(0, limit) : allResults; appLogger.i('Found ${result.length} search results across all servers'); return result; } /// 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; } /// Split "Recently Added" hubs that contain items from multiple libraries /// into separate per-library hubs, matching the official Plex client behavior. List _splitRecentlyAddedHubs(List hubs, List? libraries) { final result = []; for (final hub in hubs) { final hubId = hub.identifier?.toLowerCase() ?? ''; if (!hubId.contains('.recent')) { result.add(hub); continue; } // Group items by libraryId final groups = >{}; final ungrouped = []; for (final item in hub.items) { final libraryId = item.libraryId; if (libraryId == null) { ungrouped.add(item); } else { groups.putIfAbsent(libraryId, () => []).add(item); } } // Single library (or no groupable items) — keep hub unchanged if (groups.length <= 1) { result.add(hub); continue; } // Multiple libraries — create one hub per library for (final entry in groups.entries) { final items = entry.value; final libraryName = _resolveLibraryName(items.first, libraries); final title = libraryName != null ? t.discover.recentlyAddedIn(library: libraryName) : hub.title; result.add( hub.copyWith( title: title, identifier: '${hub.identifier}_${entry.key}', size: items.length, items: items, libraryId: entry.key, ), ); } // Keep ungrouped items in a hub with the original title if (ungrouped.isNotEmpty) { result.add(hub.copyWith(size: ungrouped.length, items: ungrouped)); } } return result; } /// Resolve a library name from an item's [libraryTitle] or by looking up /// the library in the supplied list. String? _resolveLibraryName(MediaItem item, List? libraries) { if (item.libraryTitle != null && item.libraryTitle!.isNotEmpty) { return item.libraryTitle; } if (libraries != null && item.libraryId != null) { for (final lib in libraries) { if (lib.id == item.libraryId) { return lib.title; } } } return null; } }