From edd0618e641b8b372451c524daac2103464da221 Mon Sep 17 00:00:00 2001 From: edde746 <86283021+edde746@users.noreply.github.com> Date: Sun, 14 Jun 2026 09:28:48 +0200 Subject: [PATCH] fix: harden playback and discover state --- .../edde746/plezy/exoplayer/ExoPlayerCore.kt | 8 +- .../edde746/plezy/libass/media/AssHandler.kt | 20 ++- lib/main.dart | 4 +- lib/providers/discover_provider.dart | 134 +++++++++++------- lib/screens/video_player/parts/live_tv.dart | 35 +++-- lib/services/data_aggregation_service.dart | 54 ++++--- lib/utils/live_tv_player_navigation.dart | 14 +- test/providers/discover_provider_test.dart | 57 +++++++- .../data_aggregation_bridge_test.dart | 31 ++-- 9 files changed, 252 insertions(+), 105 deletions(-) diff --git a/android/app/src/main/kotlin/com/edde746/plezy/exoplayer/ExoPlayerCore.kt b/android/app/src/main/kotlin/com/edde746/plezy/exoplayer/ExoPlayerCore.kt index 671ddf0f..84e65095 100644 --- a/android/app/src/main/kotlin/com/edde746/plezy/exoplayer/ExoPlayerCore.kt +++ b/android/app/src/main/kotlin/com/edde746/plezy/exoplayer/ExoPlayerCore.kt @@ -112,8 +112,8 @@ class ExoPlayerCore(private val activity: Activity) : Player.Listener { private var assGlCrashHandlerInstalled = false - private var cronetEngine: CronetEngine? = null - private var cronetUnavailable = false + @Volatile private var cronetEngine: CronetEngine? = null + @Volatile private var cronetUnavailable = false private fun getCronetEngine(context: Context): CronetEngine? { cronetEngine?.let { return it } if (cronetUnavailable) return null @@ -2836,6 +2836,10 @@ class ExoPlayerCore(private val activity: Activity) : Player.Listener { } else { positionMs.coerceAtLeast(0L) } + // A user seek is authoritative. Do not let a pending start-position restore + // from open/reload recovery re-seek back over an early seek near zero once + // ExoPlayer reports STATE_READY. + pendingStartPositionMs = 0L player.seekTo(clampedPositionMs) lastPosition = clampedPositionMs delegate?.onPropertyChange("time-pos", clampedPositionMs / 1000.0) diff --git a/android/libass/src/main/java/com/edde746/plezy/libass/media/AssHandler.kt b/android/libass/src/main/java/com/edde746/plezy/libass/media/AssHandler.kt index e271a877..6ed0c365 100644 --- a/android/libass/src/main/java/com/edde746/plezy/libass/media/AssHandler.kt +++ b/android/libass/src/main/java/com/edde746/plezy/libass/media/AssHandler.kt @@ -115,12 +115,25 @@ class AssHandler( override fun onMediaItemTransition(mediaItem: MediaItem?, reason: Int) { super.onMediaItemTransition(mediaItem, reason) Log.i("AssHandler", "onMediaItemTransition: item = $mediaItem, reason = $reason") + resetMediaState(releaseNative = true) + } + + private fun resetMediaState(releaseNative: Boolean) { + val oldRender = render + val oldTracks = availableTracks.values.toList() + render = null track = null + format = null availableTracks.clear() pendingFonts.clear() videoSize = Size.ZERO renderCallback?.invoke(null) + + if (releaseNative) { + oldRender?.release() + oldTracks.forEach { it.release() } + } } /** @@ -359,12 +372,7 @@ class AssHandler( videoFrameCallback = null player?.clearVideoFrameMetadataListener(videoFrameMetadataListener) player = null - render?.release() - render = null - availableTracks.values.forEach { it.release() } - availableTracks.clear() - track = null - pendingFonts.clear() + resetMediaState(releaseNative = true) if (assDelegate.isInitialized()) { ass.release() } diff --git a/lib/main.dart b/lib/main.dart index c5f82e50..4591be92 100644 --- a/lib/main.dart +++ b/lib/main.dart @@ -95,13 +95,13 @@ const String _sentryDist = String.fromEnvironment('SENTRY_DIST'); bool _zeroOffsetPointerGuardInstalled = false; void _installZeroOffsetPointerGuard() { - if (_zeroOffsetPointerGuardInstalled) return; + if (_zeroOffsetPointerGuardInstalled || !Platform.isIOS) return; GestureBinding.instance.pointerRouter.addGlobalRoute(_absorbZeroOffsetPointerEvent); _zeroOffsetPointerGuardInstalled = true; } void _absorbZeroOffsetPointerEvent(PointerEvent event) { - if (event.position == Offset.zero) { + if (event is PointerDownEvent && event.position == Offset.zero) { GestureBinding.instance.cancelPointer(event.pointer); } } diff --git a/lib/providers/discover_provider.dart b/lib/providers/discover_provider.dart index ce75759d..6af75ad9 100644 --- a/lib/providers/discover_provider.dart +++ b/lib/providers/discover_provider.dart @@ -9,6 +9,7 @@ import '../media/media_server_client.dart'; import '../mixins/disposable_change_notifier_mixin.dart'; import '../mixins/event_aware.dart'; import '../services/settings_service.dart'; +import '../services/data_aggregation_service.dart'; import '../services/system_shelf_service.dart'; import '../utils/app_logger.dart'; import '../utils/global_key_utils.dart'; @@ -77,8 +78,16 @@ class DiscoverProvider extends ChangeNotifier with DisposableChangeNotifierMixin Set _lastSeenHiddenKeys = {}; List _lastSeenLibraryOrderKeys = const []; - /// Online servers that contributed to the last successful [load] pass. - Set _loadedOnlineServerIds = {}; + /// Online servers whose Continue Watching fetch succeeded in the current + /// on-deck list. Tracked separately from hubs so a transient failure in one + /// surface does not cache the other as loaded forever or force unnecessary + /// refetches. + Set _loadedOnDeckServerIds = {}; + + /// Online servers whose home-hub fetch succeeded in the current hub list. + Set _loadedHubServerIds = {}; + + Set get _fullyLoadedServerIds => _loadedOnDeckServerIds.intersection(_loadedHubServerIds); Future? _inFlightLoad; bool _hasPendingLoad = false; @@ -117,12 +126,16 @@ class DiscoverProvider extends ChangeNotifier with DisposableChangeNotifierMixin /// and merged in; already-loaded servers are not refetched. Future syncToOnlineServers(Set onlineServerIds) { if (onlineServerIds.isEmpty || isProfileBinding()) return Future.value(); - if (_onDeckState == DiscoverLoadState.loaded && _loadedOnlineServerIds.containsAll(onlineServerIds)) { + if ( + _onDeckState == DiscoverLoadState.loaded && + _hubsState == DiscoverLoadState.loaded && + _fullyLoadedServerIds.containsAll(onlineServerIds) + ) { return Future.value(); } // Nothing (or a failed pass) to merge into yet — run the full load. if (_onDeckState != DiscoverLoadState.loaded || _hubsState != DiscoverLoadState.loaded) return load(); - _pendingDeltaServerIds.addAll(onlineServerIds.difference(_loadedOnlineServerIds)); + _pendingDeltaServerIds.addAll(onlineServerIds.difference(_fullyLoadedServerIds)); return _ensureLoadLoop(); } @@ -187,26 +200,25 @@ class DiscoverProvider extends ChangeNotifier with DisposableChangeNotifierMixin includePlaybackHubs: false, ); - final fetchedFromServerIds = Set.of(_multiServer.onlineServerIds); - final fetchedOnDeck = await onDeckFuture; if (isDisposed) return; - _applyOnDeck(fetchedOnDeck); + _applyOnDeck(fetchedOnDeck.items); _onDeckState = DiscoverLoadState.loaded; - _loadedOnlineServerIds = fetchedFromServerIds; + _loadedOnDeckServerIds = fetchedOnDeck.succeededServerIds; _loadGeneration++; safeNotifyListeners(); unawaited(_syncSystemShelf(_onDeck)); - final allHubs = await hubsFuture; + final fetchedHubs = await hubsFuture; if (isDisposed) return; - final filteredHubs = _filterDiscoverHubs(allHubs); + final filteredHubs = _filterDiscoverHubs(fetchedHubs.hubs); sortMediaHubsByLibraryOrder(filteredHubs, _libraries.libraries); appLogger.d('DiscoverProvider: ${_onDeck.length} on-deck items, ${filteredHubs.length} hubs'); _hubs = filteredHubs; _hubsState = DiscoverLoadState.loaded; + _loadedHubServerIds = fetchedHubs.succeededServerIds; safeNotifyListeners(); } catch (e) { appLogger.e('Failed to load discover content', error: e); @@ -224,9 +236,11 @@ class DiscoverProvider extends ChangeNotifier with DisposableChangeNotifierMixin /// status emission retries them. Future _loadDeltaOnce(Set serverIds) async { // A full pass may have covered these ids while they sat in the queue. - final ids = serverIds.difference(_loadedOnlineServerIds); - if (ids.isEmpty) return; - appLogger.d('DiscoverProvider: merging content from newly-online servers $ids'); + final ids = serverIds.difference(_fullyLoadedServerIds); + final onDeckIds = ids.difference(_loadedOnDeckServerIds); + final hubIds = ids.difference(_loadedHubServerIds); + if (onDeckIds.isEmpty && hubIds.isEmpty) return; + appLogger.d('DiscoverProvider: merging content from newly-online servers $ids (onDeck=$onDeckIds, hubs=$hubIds)'); try { await _hiddenLibraries.ensureInitialized(); @@ -236,43 +250,53 @@ class DiscoverProvider extends ChangeNotifier with DisposableChangeNotifierMixin final useGlobalHubs = settings.read(SettingsService.useGlobalHubs); final aggregation = _multiServer.aggregationService; - final onDeckFuture = aggregation.getOnDeckFromAllServers( - limit: _continueWatchingProbeLimit, - hiddenLibraryKeys: _hiddenLibraries.hiddenLibraryKeys, - serverIds: ids, - ); - final hubsFuture = aggregation.getHubsFromAllServers( - hiddenLibraryKeys: _hiddenLibraries.hiddenLibraryKeys, - useGlobalHubs: useGlobalHubs, - includePlaybackHubs: false, - serverIds: ids, - ); + final Future onDeckFuture = onDeckIds.isEmpty + ? Future.value() + : aggregation.getOnDeckFromAllServers( + limit: _continueWatchingProbeLimit, + hiddenLibraryKeys: _hiddenLibraries.hiddenLibraryKeys, + serverIds: onDeckIds, + ); + final Future hubsFuture = hubIds.isEmpty + ? Future.value() + : aggregation.getHubsFromAllServers( + hiddenLibraryKeys: _hiddenLibraries.hiddenLibraryKeys, + useGlobalHubs: useGlobalHubs, + includePlaybackHubs: false, + serverIds: hubIds, + ); final freshOnDeck = await onDeckFuture; final freshHubs = await hubsFuture; if (isDisposed) return; - final hadMore = _hasMoreContinueWatching; - final mergedOnDeck = await aggregation.mergeContinueWatching( - _onDeck, - freshOnDeck, - limit: _continueWatchingProbeLimit, - ); - if (isDisposed) return; - _applyOnDeck(mergedOnDeck); - // The stored list is already trimmed, so the merge can't see old items - // past the cap — a previously-true "more" affordance stays true. - if (hadMore) _hasMoreContinueWatching = true; - // No _loadGeneration bump: a delta behaves like the background Continue - // Watching refresh (the hero clamps instead of resetting). + if (freshOnDeck != null) { + final hadMore = _hasMoreContinueWatching; + final mergedOnDeck = await aggregation.mergeContinueWatching( + _onDeck, + freshOnDeck.items, + limit: _continueWatchingProbeLimit, + ); + if (isDisposed) return; + _applyOnDeck(mergedOnDeck); + // The stored list is already trimmed, so the merge can't see old items + // past the cap — a previously-true "more" affordance stays true. + if (hadMore) _hasMoreContinueWatching = true; + _loadedOnDeckServerIds = {..._loadedOnDeckServerIds, ...freshOnDeck.succeededServerIds}; + // No _loadGeneration bump: a delta behaves like the background Continue + // Watching refresh (the hero clamps instead of resetting). + } - final mergedHubs = [ - ..._hubs.where((hub) => !ids.contains(hub.serverId)), - ..._filterDiscoverHubs(freshHubs), - ]; - sortMediaHubsByLibraryOrder(mergedHubs, _libraries.libraries); - _hubs = mergedHubs; - _loadedOnlineServerIds = {..._loadedOnlineServerIds, ...ids}; + if (freshHubs != null) { + final succeededHubIds = freshHubs.succeededServerIds; + final mergedHubs = [ + ..._hubs.where((hub) => hub.serverId == null || !succeededHubIds.contains(hub.serverId)), + ..._filterDiscoverHubs(freshHubs.hubs), + ]; + sortMediaHubsByLibraryOrder(mergedHubs, _libraries.libraries); + _hubs = mergedHubs; + _loadedHubServerIds = {..._loadedHubServerIds, ...succeededHubIds}; + } appLogger.d('DiscoverProvider: ${_onDeck.length} on-deck items, ${_hubs.length} hubs after merging $ids'); safeNotifyListeners(); @@ -308,7 +332,8 @@ class DiscoverProvider extends ChangeNotifier with DisposableChangeNotifierMixin hiddenLibraryKeys: _hiddenLibraries.hiddenLibraryKeys, ); if (isDisposed) return; - _applyOnDeck(fetched); + _applyOnDeck(fetched.items); + _loadedOnDeckServerIds = fetched.succeededServerIds; safeNotifyListeners(); unawaited(_syncSystemShelf(_onDeck)); } catch (e) { @@ -321,9 +346,10 @@ class DiscoverProvider extends ChangeNotifier with DisposableChangeNotifierMixin if (!_multiServer.hasConnectedServers) return const []; await _hiddenLibraries.ensureInitialized(); if (isDisposed) return const []; - return _multiServer.aggregationService.getOnDeckFromAllServers( + final fetched = await _multiServer.aggregationService.getOnDeckFromAllServers( hiddenLibraryKeys: _hiddenLibraries.hiddenLibraryKeys, ); + return fetched.items; } /// Refetch a single item (post-edit refresh from a hub row) and swap it @@ -464,9 +490,13 @@ class DiscoverProvider extends ChangeNotifier with DisposableChangeNotifierMixin try { final settings = await SettingsService.getInstance(); + final syncableOnDeck = onDeck.where((item) { + final serverId = item.serverId; + return serverId != null && _multiServer.getClientForServer(ServerId(serverId)) != null; + }).toList(growable: false); await SystemShelfService().syncFromContinueWatching( - onDeck, - _clientWithFallback, + syncableOnDeck, + _clientForShelfItem, hideSpoilers: settings.read(SettingsService.hideSpoilers), ); } catch (e) { @@ -478,14 +508,10 @@ class DiscoverProvider extends ChangeNotifier with DisposableChangeNotifierMixin } } - MediaServerClient _clientWithFallback(ServerId serverId) { + MediaServerClient _clientForShelfItem(ServerId serverId) { final direct = _multiServer.getClientForServer(serverId); if (direct != null) return direct; - for (final id in _multiServer.onlineServerIds) { - final fallback = _multiServer.getClientForServer(ServerId(id)); - if (fallback != null) return fallback; - } - throw Exception('No client available for $serverId'); + throw Exception('No owning client available for $serverId'); } @override diff --git a/lib/screens/video_player/parts/live_tv.dart b/lib/screens/video_player/parts/live_tv.dart index 6c438599..8bf7a62f 100644 --- a/lib/screens/video_player/parts/live_tv.dart +++ b/lib/screens/video_player/parts/live_tv.dart @@ -248,19 +248,17 @@ extension _VideoPlayerLiveTvMethods on VideoPlayerScreenState { _playbackTransition = _PlaybackTransition.switchingChannel; _liveSeek.cancel(); - // Stop old session heartbeats and notify server - _stopLiveTimelineUpdates(); - await _sendLiveTimeline('stopped'); - + final previousSession = _live.session; final channel = channels[newIndex]; appLogger.d('Switching to channel: ${channel.displayName} (${channel.key})'); - if (!mounted) return; - _setPlayerState(() => _hasFirstFrame.value = false); - LiveTvPlaybackSession? session; + var replacementOpenStarted = false; try { - // Channel switch IS a fresh start: same resolution path as launch. + // Channel switch IS a fresh start: same resolution path as launch. Keep + // the old session alive until the replacement stream is actually open so + // a failed zap does not tell the server to reclaim the still-playing + // tuner/transcode session. session = await _startLiveSession(channel); if (session == null) return; if (!mounted || player != currentPlayer) { @@ -275,7 +273,25 @@ extension _VideoPlayerLiveTvMethods on VideoPlayerScreenState { } await _setLiveStreamOptions(currentPlayer); + if (!mounted || player != currentPlayer) { + _abandonLiveSession(session); + return; + } + + _setPlayerState(() => _hasFirstFrame.value = false); + replacementOpenStarted = true; await currentPlayer.open(Media(streamUrl, headers: const {'Accept-Language': 'en'}), play: true, isLive: true); + if (!mounted || player != currentPlayer) { + _abandonLiveSession(session); + return; + } + + // The new stream is now the active local playback. Stop the old heartbeat + // and send its terminal timeline before adopting the replacement session. + _stopLiveTimelineUpdates(); + if (previousSession != null) { + await _sendLiveTimeline('stopped'); + } _live.adoptSession(session); _live.fallbackLevel = 0; @@ -294,6 +310,9 @@ extension _VideoPlayerLiveTvMethods on VideoPlayerScreenState { // would otherwise hold its server-side tuner until the backend times out. final orphan = session; if (orphan != null && _live.session != orphan) _abandonLiveSession(orphan); + if (replacementOpenStarted && mounted && _live.session == previousSession) { + _setPlayerState(() => _hasFirstFrame.value = true); + } appLogger.e('Failed to switch channel', error: e); if (mounted) showErrorSnackBar(context, e.toString()); } finally { diff --git a/lib/services/data_aggregation_service.dart b/lib/services/data_aggregation_service.dart index 01a68df4..74adc1e8 100644 --- a/lib/services/data_aggregation_service.dart +++ b/lib/services/data_aggregation_service.dart @@ -12,6 +12,9 @@ import '../utils/global_key_utils.dart'; import '../utils/search_relevance.dart'; import 'multi_server_manager.dart'; +typedef OnDeckAggregationResult = ({List items, Set succeededServerIds}); +typedef HubAggregationResult = ({List hubs, Set succeededServerIds}); + /// 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 @@ -70,8 +73,9 @@ class DataAggregationService { /// 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. [serverIds] restricts the fan-out to those servers. - Future> getOnDeckFromAllServers({ + /// 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, @@ -79,18 +83,25 @@ class DataAggregationService { final clients = _clientsFor(serverIds); if (clients.isEmpty) { appLogger.w('No online servers available for fetching on deck'); - return []; + return (items: const [], succeededServerIds: const {}); } + final futures = clients.entries.map((entry) async { final client = entry.value; try { - return await client.fetchContinueWatching(count: limit); + final items = await client.fetchContinueWatching(count: limit); + return (serverId: entry.key, items: items); } catch (e, st) { appLogger.e('Failed on-deck fetch from ${entry.key}', error: e, stackTrace: st); - return []; + return (serverId: null, items: []); } }); - final allOnDeck = (await Future.wait(futures)).expand((l) => l).toList(); + final results = await Future.wait(futures); + final succeededServerIds = { + for (final result in results) + if (result.serverId != null) result.serverId!, + }; + final allOnDeck = results.expand((result) => result.items).toList(); // Filter out items from hidden libraries List filteredOnDeck = allOnDeck; @@ -110,11 +121,11 @@ class DataAggregationService { filteredOnDeck = await _deduplicateContinueWatching(filteredOnDeck); // Apply limit if specified - final result = limit != null && limit < filteredOnDeck.length ? filteredOnDeck.sublist(0, limit) : filteredOnDeck; + final items = limit != null && limit < filteredOnDeck.length ? filteredOnDeck.sublist(0, limit) : filteredOnDeck; - appLogger.i('Fetched ${result.length} on deck items from all servers'); + appLogger.i('Fetched ${items.length} on deck items from all servers'); - return result; + return (items: items, succeededServerIds: succeededServerIds); } /// Merge an [existing] Continue Watching list with [fresh] rows from @@ -266,8 +277,9 @@ class DataAggregationService { /// 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. - Future> getHubsFromAllServers({ + /// 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, @@ -277,7 +289,7 @@ class DataAggregationService { final clients = _clientsFor(serverIds); if (clients.isEmpty) { appLogger.w('No online servers available for fetching hubs'); - return []; + return (hubs: const [], succeededServerIds: const {}); } // Only fallback clients need a library prefetch when home layout is on; @@ -303,19 +315,27 @@ class DataAggregationService { includePlaybackHubs: includePlaybackHubs, libraries: useGlobalHubs ? serverLibraries : null, ); - return _postProcessHubs(hubs, serverId: ServerId(serverId), hiddenLibraryKeys: hiddenLibraryKeys); + return ( + serverId: serverId, + hubs: _postProcessHubs(hubs, serverId: ServerId(serverId), hiddenLibraryKeys: hiddenLibraryKeys), + ); } catch (e, stackTrace) { appLogger.e('Failed to fetch hubs from server $serverId', error: e, stackTrace: stackTrace); - return []; + return (serverId: null, hubs: []); } }); final results = await Future.wait(futures); + final succeededServerIds = { + for (final result in results) + if (result.serverId != null) result.serverId!, + }; final all = []; - for (final list in results) { - all.addAll(list); + for (final result in results) { + all.addAll(result.hubs); } - return limit != null && limit < all.length ? all.sublist(0, limit) : all; + final hubs = limit != null && limit < all.length ? all.sublist(0, limit) : all; + return (hubs: hubs, succeededServerIds: succeededServerIds); } /// Per-library hub fetch for a single client. Filters to visible diff --git a/lib/utils/live_tv_player_navigation.dart b/lib/utils/live_tv_player_navigation.dart index 0c387767..b7ab6d88 100644 --- a/lib/utils/live_tv_player_navigation.dart +++ b/lib/utils/live_tv_player_navigation.dart @@ -54,14 +54,24 @@ Future navigateToLiveTv( raw: {'key': channel.key}, ); + final normalizedChannels = List.of(channels); + var currentChannelIndex = normalizedChannels.indexWhere( + (ch) => liveTvChannelScopeKey(ch) == liveTvChannelScopeKey(channel), + ); + if (currentChannelIndex < 0) { + normalizedChannels.insert(0, channel); + currentChannelIndex = 0; + appLogger.w('Live TV launch channel was not present in navigation list; prepending ${channel.key}'); + } + final route = PageRouteBuilder( settings: const RouteSettings(name: kVideoPlayerRouteName), pageBuilder: (context, animation, secondaryAnimation) => VideoPlayerScreen( metadata: placeholder, live: LiveTvSessionArgs( channel: channel, - channels: channels, - currentChannelIndex: channels.indexWhere((ch) => liveTvChannelScopeKey(ch) == liveTvChannelScopeKey(channel)), + channels: normalizedChannels, + currentChannelIndex: currentChannelIndex, ), ), transitionDuration: Duration.zero, diff --git a/test/providers/discover_provider_test.dart b/test/providers/discover_provider_test.dart index 23138749..27343e29 100644 --- a/test/providers/discover_provider_test.dart +++ b/test/providers/discover_provider_test.dart @@ -50,11 +50,13 @@ class _FakeAggregationService extends DataAggregationService { int hubCalls = 0; Set? lastOnDeckServerIds; Set? lastHubsServerIds; + Set? onDeckSucceededServerIds; + Set? hubSucceededServerIds; List Function() onDeckResult = () => const []; List Function() hubsResult = () => const []; @override - Future> getOnDeckFromAllServers({ + Future getOnDeckFromAllServers({ int? limit, Set? hiddenLibraryKeys, Set? serverIds, @@ -62,11 +64,14 @@ class _FakeAggregationService extends DataAggregationService { onDeckCalls++; lastOnDeckServerIds = serverIds; final items = onDeckResult(); - return limit != null && items.length > limit ? items.sublist(0, limit) : items; + return ( + items: limit != null && items.length > limit ? items.sublist(0, limit) : items, + succeededServerIds: onDeckSucceededServerIds ?? serverIds ?? const {'server_1'}, + ); } @override - Future> getHubsFromAllServers({ + Future getHubsFromAllServers({ int? limit, Set? hiddenLibraryKeys, bool useGlobalHubs = true, @@ -75,7 +80,10 @@ class _FakeAggregationService extends DataAggregationService { }) async { hubCalls++; lastHubsServerIds = serverIds; - return hubsResult(); + return ( + hubs: hubsResult(), + succeededServerIds: hubSucceededServerIds ?? serverIds ?? const {'server_1'}, + ); } } @@ -341,6 +349,47 @@ void main() { expect(aggregation.onDeckCalls, callsAfterDelta); }); + test('full load partial hub failure retries hubs without refetching continue watching', () async { + aggregation.onDeckResult = () => [_item('a')]; + aggregation.hubsResult = () => const []; + aggregation.hubSucceededServerIds = const {}; + await provider.load(); + final onDeckCallsBefore = aggregation.onDeckCalls; + final hubCallsBefore = aggregation.hubCalls; + + aggregation.hubsResult = () => [_hub('hub-1')]; + aggregation.hubSucceededServerIds = const {'server_1'}; + await provider.syncToOnlineServers({'server_1'}); + + expect(aggregation.onDeckCalls, onDeckCallsBefore); + expect(aggregation.hubCalls, hubCallsBefore + 1); + expect(aggregation.lastHubsServerIds, {'server_1'}); + expect(provider.hubs.map((h) => h.id), ['hub-1']); + }); + + test('delta partial hub failure retries only the missing surface', () async { + aggregation.onDeckResult = () => [_item('a')]; + aggregation.hubsResult = () => [_hub('hub-1')]; + await provider.load(); + + aggregation.onDeckResult = () => [_item('b', serverId: 'server_2')]; + aggregation.hubsResult = () => const []; + aggregation.hubSucceededServerIds = const {}; + await provider.syncToOnlineServers({'server_1', 'server_2'}); + final onDeckCallsAfterPartial = aggregation.onDeckCalls; + final hubCallsAfterPartial = aggregation.hubCalls; + + aggregation.hubsResult = () => [_hub('hub-2', serverId: 'server_2')]; + aggregation.hubSucceededServerIds = const {'server_2'}; + await provider.syncToOnlineServers({'server_1', 'server_2'}); + + expect(aggregation.onDeckCalls, onDeckCallsAfterPartial); + expect(aggregation.hubCalls, hubCallsAfterPartial + 1); + expect(aggregation.lastHubsServerIds, {'server_2'}); + expect(provider.onDeck.map((i) => i.id), containsAll(['a', 'b'])); + expect(provider.hubs.map((h) => h.id), containsAll(['hub-1', 'hub-2'])); + }); + test('delta failure keeps the loaded state and retries on the next emission', () async { aggregation.onDeckResult = () => [_item('a')]; aggregation.hubsResult = () => [_hub('hub-1')]; diff --git a/test/services/data_aggregation_bridge_test.dart b/test/services/data_aggregation_bridge_test.dart index 0ef4809a..3b692af4 100644 --- a/test/services/data_aggregation_bridge_test.dart +++ b/test/services/data_aggregation_bridge_test.dart @@ -60,7 +60,9 @@ void main() { test('searchAcrossServers and getOnDeckFromAllServers return empty when no clients', () async { expect(await service.searchAcrossServers('hello'), isEmpty); - expect(await service.getOnDeckFromAllServers(), isEmpty); + final onDeck = await service.getOnDeckFromAllServers(); + expect(onDeck.items, isEmpty); + expect(onDeck.succeededServerIds, isEmpty); }); test('searchAcrossServers overfetches and ranks before trimming across backends', () async { @@ -161,9 +163,10 @@ void main() { addTearDown(client.close); manager.debugRegisterClientForTesting(client); - final items = await service.getOnDeckFromAllServers(limit: 21); + final result = await service.getOnDeckFromAllServers(limit: 21); - expect(items.map((item) => item.id), ['movie-1']); + expect(result.items.map((item) => item.id), ['movie-1']); + expect(result.succeededServerIds, {'plex-1'}); expect(captured.single.path, '/hubs'); expect(captured.single.queryParameters['count'], '21'); }); @@ -239,9 +242,10 @@ void main() { addTearDown(client.close); manager.debugRegisterClientForTesting(client); - final items = await service.getOnDeckFromAllServers(limit: 10); + final result = await service.getOnDeckFromAllServers(limit: 10); - expect(items.map((item) => item.id), ['new-episode']); + expect(result.items.map((item) => item.id), ['new-episode']); + expect(result.succeededServerIds, {'plex-1'}); }); test('getOnDeckFromAllServers keeps duplicate titles without stable ids', () async { @@ -307,9 +311,10 @@ void main() { addTearDown(client.close); manager.debugRegisterClientForTesting(client); - final items = await service.getOnDeckFromAllServers(limit: 10); + final result = await service.getOnDeckFromAllServers(limit: 10); - expect(items.map((item) => item.id), ['new-unmatched', 'old-unmatched']); + expect(result.items.map((item) => item.id), ['new-unmatched', 'old-unmatched']); + expect(result.succeededServerIds, {'plex-1'}); }); test('per-library hubs skip playback rows and fetch in bounded batches', () async { @@ -352,8 +357,10 @@ void main() { addTearDown(client.close); manager.debugRegisterJellyfinClientForTesting(client); - final hubs = await service.getHubsFromAllServers(useGlobalHubs: false, includePlaybackHubs: false); + final result = await service.getHubsFromAllServers(useGlobalHubs: false, includePlaybackHubs: false); + final hubs = result.hubs; + expect(result.succeededServerIds, {'srv-1'}); expect(hubs.map((h) => h.identifier), [ 'library.lib-1.recent', 'library.lib-2.recent', @@ -410,8 +417,10 @@ void main() { addTearDown(client.close); manager.debugRegisterJellyfinClientForTesting(client); - final hubs = await service.getHubsFromAllServers(useGlobalHubs: true, includePlaybackHubs: false); + final result = await service.getHubsFromAllServers(useGlobalHubs: true, includePlaybackHubs: false); + final hubs = result.hubs; + expect(result.succeededServerIds, {'srv-1'}); expect(hubs.map((h) => h.identifier), ['library.movies.recent', 'library.shows.recent']); expect(hubs.map((h) => h.items.single.id), ['movie-1', 'show-1']); expect(captured.where((uri) => uri.path == '/Users/user-1/Views'), hasLength(1)); @@ -473,8 +482,10 @@ void main() { addTearDown(client.close); manager.debugRegisterClientForTesting(client); - final hubs = await service.getHubsFromAllServers(useGlobalHubs: true, includePlaybackHubs: false); + final result = await service.getHubsFromAllServers(useGlobalHubs: true, includePlaybackHubs: false); + final hubs = result.hubs; + expect(result.succeededServerIds, {'plex-1'}); expect(hubs, hasLength(1)); expect(hubs.single.title, 'Recently Added TV'); expect(hubs.single.identifier, 'home.television.recent');