diff --git a/lib/services/plex_client.dart b/lib/services/plex_client.dart index 435d05ba..29770b2f 100644 --- a/lib/services/plex_client.dart +++ b/lib/services/plex_client.dart @@ -44,6 +44,7 @@ import '../models/plex/plex_video_playback_data.dart'; import '../models/transcode_quality_preset.dart'; import '../utils/endpoint_failover_interceptor.dart'; import '../utils/app_logger.dart'; +import '../utils/media_server_retry.dart'; import '../utils/media_server_timeouts.dart'; import '../utils/log_redaction_manager.dart'; import '../utils/plex_cache_parser.dart'; @@ -1068,13 +1069,20 @@ class PlexClient with MediaServerCacheMixin, _PlexLiveTvClientMethods implements /// Uses /hubs?identifier=home.continue,home.ondeck which respects the /// server's OnDeckWindow preference (unlike /library/onDeck). Future> _getContinueWatching({int count = 20}) async { - final response = await _getWithFailover( - '/hubs', - queryParameters: {'identifier': 'home.continue,home.ondeck', 'count': count, 'includeGuids': 1}, + final response = await retryTransientMediaServerCall( + operation: 'Plex continue watching hubs', + attemptTimeouts: MediaServerTimeouts.homeHubAttemptTimeouts, + call: (timeout, abort) => _getWithFailover( + '/hubs', + queryParameters: {'identifier': 'home.continue,home.ondeck', 'count': count, 'includeGuids': 1}, + timeout: timeout, + abort: abort, + ), ); final sid = serverId; final sname = serverName; - final hubs = await tryIsolateRun(() => _processHubResponse(response.data as Map, sid, sname)); + final data = response.data as Map; + final hubs = await tryIsolateRun(() => _processHubResponse(data, sid, sname)); // Deduplicate across home.continue and home.ondeck hubs. // Like plex-web, episodes from the same show (same grandparentRatingKey) // are deduplicated, preferring the in-progress item (has viewOffset). @@ -1521,7 +1529,8 @@ class PlexClient with MediaServerCacheMixin, _PlexLiveTvClientMethods implements ); final sid = serverId; final sname = serverName; - return await tryIsolateRun(() => _processHubResponse(response.data as Map, sid, sname)); + final data = response.data as Map; + return await tryIsolateRun(() => _processHubResponse(data, sid, sname)); } catch (e) { appLogger.e('Failed to get library hubs: $e'); } @@ -1533,10 +1542,20 @@ class PlexClient with MediaServerCacheMixin, _PlexLiveTvClientMethods implements /// This matches the official Plex client's home page layout. Future> _getGlobalHubs({int limit = 10}) async { try { - final response = await _getWithFailover('/hubs', queryParameters: {'count': limit, 'includeGuids': 1}); + final response = await retryTransientMediaServerCall( + operation: 'Plex global hubs', + attemptTimeouts: MediaServerTimeouts.homeHubAttemptTimeouts, + call: (timeout, abort) => _getWithFailover( + '/hubs', + queryParameters: {'count': limit, 'includeGuids': 1}, + timeout: timeout, + abort: abort, + ), + ); final sid = serverId; final sname = serverName; - return await tryIsolateRun(() => _processHubResponse(response.data as Map, sid, sname)); + final data = response.data as Map; + return await tryIsolateRun(() => _processHubResponse(data, sid, sname)); } catch (e) { appLogger.e('Failed to get global hubs: $e'); } @@ -1549,9 +1568,10 @@ class PlexClient with MediaServerCacheMixin, _PlexLiveTvClientMethods implements final response = await _getWithFailover('/hubs/metadata/$ratingKey/related', queryParameters: {'count': count}); final sid = serverId; final sname = serverName; + final data = response.data as Map; return await tryIsolateRun( () => _processHubResponse( - response.data as Map, + data, sid, sname, filter: (item) { diff --git a/lib/utils/isolate_helper.dart b/lib/utils/isolate_helper.dart index 2d13e42e..2cb33f5c 100644 --- a/lib/utils/isolate_helper.dart +++ b/lib/utils/isolate_helper.dart @@ -10,5 +10,8 @@ Future tryIsolateRun(R Function() computation) async { return await Isolate.run(computation); } on StateError { return computation(); + } on ArgumentError catch (e) { + if (!e.toString().contains('Illegal argument in isolate message')) rethrow; + return computation(); } } diff --git a/lib/utils/media_server_retry.dart b/lib/utils/media_server_retry.dart new file mode 100644 index 00000000..45e1b4eb --- /dev/null +++ b/lib/utils/media_server_retry.dart @@ -0,0 +1,47 @@ +import '../exceptions/media_server_exceptions.dart'; +import 'app_logger.dart'; +import 'media_server_http_client.dart'; + +typedef MediaServerRetryCall = Future Function(Duration timeout, AbortController abort); + +/// Retries media-server calls only when the failure is transient transport +/// noise. Callers pass per-attempt timeouts so cold-start surfaces can use a +/// bounded retry budget without changing global HTTP defaults. +Future retryTransientMediaServerCall({ + required String operation, + required List attemptTimeouts, + required MediaServerRetryCall call, +}) async { + if (attemptTimeouts.isEmpty) { + throw ArgumentError.value(attemptTimeouts, 'attemptTimeouts', 'must contain at least one timeout'); + } + + for (var attempt = 0; attempt < attemptTimeouts.length; attempt++) { + final timeout = attemptTimeouts[attempt]; + final abort = AbortController(); + try { + return await call(timeout, abort); + } on MediaServerHttpException catch (e, st) { + abort.abort(); + final isLastAttempt = attempt == attemptTimeouts.length - 1; + if (!e.isTransient || isLastAttempt) { + Error.throwWithStackTrace(e, st); + } + + appLogger.w( + 'Retrying $operation after transient media-server failure', + error: { + 'attempt': attempt + 1, + 'maxAttempts': attemptTimeouts.length, + 'nextTimeoutMs': attemptTimeouts[attempt + 1].inMilliseconds, + 'type': e.type.name, + }, + ); + } catch (e, st) { + abort.abort(); + Error.throwWithStackTrace(e, st); + } + } + + throw StateError('unreachable retry state'); +} diff --git a/lib/utils/media_server_timeouts.dart b/lib/utils/media_server_timeouts.dart index a8b2eaee..21069e4f 100644 --- a/lib/utils/media_server_timeouts.dart +++ b/lib/utils/media_server_timeouts.dart @@ -10,6 +10,10 @@ class MediaServerTimeouts { /// HTTP receive timeout for streaming/large responses from a media server. static const receive = Duration(seconds: 120); + /// Retry budget for home `/hubs` startup calls. These endpoints can be slow + /// while Plex wakes idle disks, but should not block forever. + static const homeHubAttemptTimeouts = [Duration(seconds: 10), Duration(seconds: 5), Duration(milliseconds: 2500)]; + // ── Plex server discovery / endpoint racing ──────────────────── /// Timeout for probing a cached/preferred endpoint before falling back to /// the full candidate race (used in [PlexServer.findBestWorkingConnection]). diff --git a/test/services/plex_home_retry_test.dart b/test/services/plex_home_retry_test.dart new file mode 100644 index 00000000..cc2f5e5d --- /dev/null +++ b/test/services/plex_home_retry_test.dart @@ -0,0 +1,90 @@ +import 'dart:async'; +import 'dart:convert'; + +import 'package:drift/native.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:http/http.dart' as http; +import 'package:plezy/database/app_database.dart'; +import 'package:plezy/models/plex/plex_config.dart'; +import 'package:plezy/services/plex_api_cache.dart'; +import 'package:plezy/services/plex_client.dart'; + +typedef _RequestHandler = Future Function(http.BaseRequest request); + +class _SequenceClient extends http.BaseClient { + _SequenceClient(this._handlers); + + final List<_RequestHandler> _handlers; + final requests = []; + + @override + Future send(http.BaseRequest request) { + requests.add(request); + if (_handlers.isEmpty) { + throw StateError('Unexpected request: ${request.url}'); + } + return _handlers.removeAt(0)(request); + } +} + +void main() { + group('PlexClient home hub retries', () { + test('fetchGlobalHubs retries a transient first failure', () async { + final db = AppDatabase.forTesting(NativeDatabase.memory()); + PlexApiCache.initialize(db); + addTearDown(db.close); + + final httpClient = _SequenceClient([ + (_) async => throw TimeoutException('cold Plex start'), + (_) async => _jsonResponse(_globalHubsPayload()), + ]); + final client = PlexClient.forTesting( + config: PlexConfig( + baseUrl: 'http://server:32400', + token: 'token', + clientIdentifier: 'client-id', + product: 'Plezy', + version: 'test', + ), + serverId: 'server-id', + serverName: 'Server', + httpClient: httpClient, + ); + addTearDown(client.close); + + final hubs = await client.fetchGlobalHubs(limit: 12); + + expect(hubs, hasLength(1)); + expect(hubs.single.title, 'Recently Added Movies'); + expect(hubs.single.items.single.title, 'Movie A'); + expect(httpClient.requests, hasLength(2)); + expect(httpClient.requests.map((r) => r.url.path), everyElement('/hubs')); + expect(httpClient.requests.map((r) => r.url.queryParameters['count']), everyElement('12')); + }); + }); +} + +Future _jsonResponse(Map body) async { + return http.StreamedResponse( + Stream.value(utf8.encode(jsonEncode(body))), + 200, + headers: {'content-type': 'application/json'}, + ); +} + +Map _globalHubsPayload() => { + 'MediaContainer': { + 'Hub': [ + { + 'key': '/hubs/movie.recentlyAdded', + 'title': 'Recently Added Movies', + 'type': 'movie', + 'hubIdentifier': 'movie.recentlyAdded.1', + 'size': 1, + 'Metadata': [ + {'ratingKey': '1', 'type': 'movie', 'title': 'Movie A'}, + ], + }, + ], + }, +}; diff --git a/test/utils/media_server_retry_test.dart b/test/utils/media_server_retry_test.dart new file mode 100644 index 00000000..ffdeb402 --- /dev/null +++ b/test/utils/media_server_retry_test.dart @@ -0,0 +1,73 @@ +import 'package:flutter_test/flutter_test.dart'; +import 'package:plezy/exceptions/media_server_exceptions.dart'; +import 'package:plezy/utils/media_server_http_client.dart'; +import 'package:plezy/utils/media_server_retry.dart'; + +void main() { + group('retryTransientMediaServerCall', () { + test('retries transient failures in timeout order', () async { + const timeouts = [Duration(seconds: 10), Duration(seconds: 5), Duration(milliseconds: 2500)]; + final seenTimeouts = []; + final aborts = []; + + final result = await retryTransientMediaServerCall( + operation: 'test operation', + attemptTimeouts: timeouts, + call: (timeout, abort) async { + seenTimeouts.add(timeout); + aborts.add(abort); + if (seenTimeouts.length < 3) { + throw MediaServerHttpException(type: MediaServerHttpErrorType.connectionTimeout, message: 'timed out'); + } + return 'ok'; + }, + ); + + expect(result, 'ok'); + expect(seenTimeouts, timeouts); + expect(aborts[0].isAborted, isTrue); + expect(aborts[1].isAborted, isTrue); + expect(aborts[2].isAborted, isFalse); + }); + + test('does not retry non-transient failures', () async { + var attempts = 0; + + await expectLater( + retryTransientMediaServerCall( + operation: 'test operation', + attemptTimeouts: const [Duration(seconds: 10), Duration(seconds: 5)], + call: (_, _) async { + attempts++; + throw MediaServerHttpException( + type: MediaServerHttpErrorType.unknown, + statusCode: 404, + message: 'HTTP 404', + ); + }, + ), + throwsA(isA().having((e) => e.statusCode, 'statusCode', 404)), + ); + + expect(attempts, 1); + }); + + test('rethrows final transient failure after exhausting attempts', () async { + var attempts = 0; + + await expectLater( + retryTransientMediaServerCall( + operation: 'test operation', + attemptTimeouts: const [Duration(seconds: 10), Duration(seconds: 5), Duration(milliseconds: 2500)], + call: (_, _) async { + attempts++; + throw MediaServerHttpException(type: MediaServerHttpErrorType.receiveTimeout, message: 'receive timed out'); + }, + ), + throwsA(isA().having((e) => e.type, 'type', MediaServerHttpErrorType.receiveTimeout)), + ); + + expect(attempts, 3); + }); + }); +}