diff --git a/lib/screens/video_player/parts/episode_navigation.dart b/lib/screens/video_player/parts/episode_navigation.dart index 4b3fea76..8a60069c 100644 --- a/lib/screens/video_player/parts/episode_navigation.dart +++ b/lib/screens/video_player/parts/episode_navigation.dart @@ -632,11 +632,22 @@ extension _VideoPlayerEpisodeNavigationMethods on VideoPlayerScreenState { } if (!isCurrentReload()) return _MediaReloadOutcome.superseded; - // Overlap the old item's stop report with the resolve round-trip; it - // is awaited again right before the open below. + // Local resume lookup can overlap the old stop, but source resolution + // below must not: Plex can use that stop to terminate any new + // transcode sharing this playback session identifier. final stoppedProgressFuture = _sendStoppedProgressOnce(); + var openResumePosition = await _resolveOpenResumePosition( + metadata: metadata, + isOffline: _offlineLibraryMode, + offlineWatchService: offlineWatchService, + requested: resumePosition, + ); + if (!isCurrentReload()) return _MediaReloadOutcome.superseded; + final playbackResolver = PlaybackSourceResolver(serverManager: serverManager, database: database); + await stoppedProgressFuture; + if (!isCurrentReload()) return _MediaReloadOutcome.superseded; final playbackContext = await playbackResolver.resolve( PlaybackInitializationOptions( metadata: metadata, @@ -649,6 +660,7 @@ extension _VideoPlayerEpisodeNavigationMethods on VideoPlayerScreenState { preferredSubtitleTrack: initializationSubtitleTrack, sessionIdentifier: _playbackSessionIdentifier, transcodeSessionId: _playbackTranscodeSessionId, + transcodeOffset: openResumePosition, ), offlineLibraryMode: _offlineLibraryMode, ); @@ -661,7 +673,17 @@ extension _VideoPlayerEpisodeNavigationMethods on VideoPlayerScreenState { if (result.videoUrl == null) { throw PlaybackException('No video URL available'); } - + if (result.isOffline && !_offlineLibraryMode) { + // The pre-resolve lookup assumed an online source; a download won + // instead, so consult locally tracked progress after all. + openResumePosition = await _resolveOpenResumePosition( + metadata: metadata, + isOffline: true, + offlineWatchService: offlineWatchService, + requested: resumePosition, + ); + if (!isCurrentReload()) return _MediaReloadOutcome.superseded; + } var subtitleSelection = await _resolveSubtitleSelectionForOpen( metadata: metadata, result: result, @@ -688,14 +710,6 @@ extension _VideoPlayerEpisodeNavigationMethods on VideoPlayerScreenState { showErrorSnackBar(context, t.videoControls.transcodeUnavailableFallback); } - final openResumePosition = await _resolveOpenResumePosition( - metadata: metadata, - isOffline: _offlineLibraryMode || result.isOffline, - offlineWatchService: offlineWatchService, - requested: resumePosition, - ); - if (!isCurrentReload()) return _MediaReloadOutcome.superseded; - final displayCriteria = result.mediaInfo?.displayCriteria; final settingsService = await SettingsService.getInstance(); if (!isCurrentReload()) return _MediaReloadOutcome.superseded; @@ -731,7 +745,6 @@ extension _VideoPlayerEpisodeNavigationMethods on VideoPlayerScreenState { resumePosition: openResumePosition, durationMs: metadata.durationMs, ); - await stoppedProgressFuture; _progressTracker?.stopTracking(); _progressTracker?.dispose(); _progressTracker = null; @@ -753,6 +766,12 @@ extension _VideoPlayerEpisodeNavigationMethods on VideoPlayerScreenState { externalSubtitles: subtitleSelection.sidecarsAtOpen, ); var effectiveExternalSubtitlePlan = externalSubtitlePlan; + await _awaitTranscodeReadiness( + client: mediaClient, + isTranscoding: result.isTranscoding, + videoUrl: result.videoUrl!, + ); + if (!isCurrentReload()) return _MediaReloadOutcome.superseded; final openResult = await _openMediaOnPlayer( player: currentPlayer, settingsService: settingsService, diff --git a/lib/screens/video_player/parts/playback_open.dart b/lib/screens/video_player/parts/playback_open.dart index 59377ab8..d00afbdf 100644 --- a/lib/screens/video_player/parts/playback_open.dart +++ b/lib/screens/video_player/parts/playback_open.dart @@ -602,6 +602,24 @@ extension _VideoPlayerOpenMethods on VideoPlayerScreenState { await player.setProperty('stream-buffer-size', '${ringBytes ?? mpvDefaultStreamBufferBytes}'); } + /// Best-effort wait for an offset transcode session's segment at the + /// resume point, run immediately before the player opens the URL so the + /// wait hides behind the other pre-open work and the guarantee is fresh + /// when the player attaches. A not-ready session still opens — mpv + /// classifies whatever the server actually returns — and no-offset URLs + /// return immediately. Starting a new probe aborts the previous one so a + /// superseded open never leaves it polling out its window. + Future _awaitTranscodeReadiness({ + required MediaServerClient? client, + required bool isTranscoding, + required String videoUrl, + }) async { + if (!isTranscoding || client is! PlexClient) return; + _transcodeReadinessAbort?.abort(); + final abort = _transcodeReadinessAbort = AbortController(); + await client.waitForTranscodeReady(videoUrl, abort: abort); + } + /// Open [videoUrl] on [player]: stream tuning → open → native subtitle style. /// /// [shouldContinue] is re-checked between the awaits so stale generations diff --git a/lib/screens/video_player/parts/playback_start.dart b/lib/screens/video_player/parts/playback_start.dart index 1ef39e91..2c558f3f 100644 --- a/lib/screens/video_player/parts/playback_start.dart +++ b/lib/screens/video_player/parts/playback_start.dart @@ -259,6 +259,12 @@ extension _VideoPlayerPlaybackStartMethods on VideoPlayerScreenState { resumePosition: resumePosition, durationMs: _currentMetadata.durationMs, ); + await _awaitTranscodeReadiness( + client: playbackContext.reportingClient, + isTranscoding: result.isTranscoding, + videoUrl: result.videoUrl!, + ); + if (!attempt.isCurrent) return; final openResult = await _openMediaOnPlayer( player: currentPlayer, settingsService: settingsService, diff --git a/lib/screens/video_player_screen.dart b/lib/screens/video_player_screen.dart index e6bf5f5d..34a12ab9 100644 --- a/lib/screens/video_player_screen.dart +++ b/lib/screens/video_player_screen.dart @@ -75,6 +75,7 @@ import '../providers/shader_provider.dart'; import '../providers/user_profile_provider.dart'; import '../utils/app_logger.dart'; import '../utils/dialogs.dart'; +import '../utils/media_server_http_client.dart' show AbortController; import '../utils/log_redaction_manager.dart'; import '../utils/live_tv_player_navigation.dart'; import '../utils/player_utils.dart'; @@ -431,6 +432,10 @@ class VideoPlayerScreenState extends State with WidgetsBindin _PlaybackTransitionLease? _playbackTransitionLease; Completer? _playbackTransitionIdleCompleter; bool _playbackIntentShouldPlay = true; + + /// In-flight transcode readiness probe, aborted by the next probe or by + /// dispose so a superseded open never leaves it polling out its window. + AbortController? _transcodeReadinessAbort; int _pendingSubtitleCycleCount = 0; bool _subtitleCycleDrainActive = false; @@ -1268,6 +1273,13 @@ class VideoPlayerScreenState extends State with WidgetsBindin preferredSubtitleTrack: _preferredSubtitleTrack, sessionIdentifier: _playbackSessionIdentifier, transcodeSessionId: _playbackTranscodeSessionId, + // The initial resume position is the server view offset (the + // online open resolves the same value later), so a resumed + // transcode starts producing at the resume point instead of + // seeking a stream that begins at zero. + transcodeOffset: _currentMetadata.viewOffsetMs != null + ? Duration(milliseconds: _currentMetadata.viewOffsetMs!) + : null, ), offlineLibraryMode: false, ); @@ -1688,6 +1700,7 @@ class VideoPlayerScreenState extends State with WidgetsBindin @override void dispose() { unawaited(AndroidExitDiagnostics.markUiState(AndroidUiState.mainScreen)); + _transcodeReadinessAbort?.abort(); _playerInitializationGeneration++; _frameRate.dispose(); WidgetsBinding.instance.removeObserver(this); diff --git a/lib/services/playback_initialization_types.dart b/lib/services/playback_initialization_types.dart index f2151193..9ad3da05 100644 --- a/lib/services/playback_initialization_types.dart +++ b/lib/services/playback_initialization_types.dart @@ -61,6 +61,13 @@ class PlaybackInitializationOptions { /// for Plex transcode. final String? transcodeSessionId; + /// Absolute VOD position at which a new Plex transcode must begin. Sent + /// with both Plex's decision and HLS start request so the server and the + /// player agree on the first available segment. Only the Plex client + /// consumes this today; Jellyfin's StartTimeTicks equivalent is + /// intentionally unwired. + final Duration? transcodeOffset; + const PlaybackInitializationOptions({ required this.metadata, required this.selectedMediaIndex, @@ -73,6 +80,7 @@ class PlaybackInitializationOptions { this.preferredSubtitleTrack, this.sessionIdentifier, this.transcodeSessionId, + this.transcodeOffset, }); } diff --git a/lib/services/plex_client.dart b/lib/services/plex_client.dart index a30c6d91..118b89b5 100644 --- a/lib/services/plex_client.dart +++ b/lib/services/plex_client.dart @@ -1,6 +1,7 @@ import 'dart:async'; import '../utils/isolate_helper.dart'; import '../utils/json_utils.dart'; +import 'package:clock/clock.dart'; import 'package:flutter/foundation.dart'; import 'package:http/http.dart' as http; import 'package:uuid/uuid.dart'; @@ -2617,6 +2618,7 @@ class PlexClient required String sessionIdentifier, required String transcodeSessionId, int? audioStreamId, + Duration? offset, }) async { try { final allParams = _buildTranscodeParams( @@ -2627,6 +2629,7 @@ class PlexClient sessionIdentifier: sessionIdentifier, transcodeSessionId: transcodeSessionId, audioStreamId: audioStreamId, + offset: offset, ); return await _runTranscodeDecision( startEndpoint: _plexVideoHlsStartEndpoint, @@ -2639,6 +2642,205 @@ class PlexClient } } + /// Absolute media position a transcode start URL was requested at, or null + /// when the URL is not an offset HLS start request. Matched on decoded + /// path segments so a percent-encoded spelling of the same URL cannot + /// silently switch the readiness probe off. + static Duration? transcodeStreamOffsetFromUrl(String videoUrl) { + final uri = Uri.tryParse(videoUrl); + if (uri == null || !'/${uri.pathSegments.join('/')}'.endsWith('/video/:/transcode/universal/start.m3u8')) { + return null; + } + final offsetSeconds = double.tryParse(uri.queryParameters['offset'] ?? ''); + if (offsetSeconds == null || offsetSeconds <= 0) return null; + return Duration(microseconds: (offsetSeconds * Duration.microsecondsPerSecond).round()); + } + + /// Picks the playlist entry the readiness probe should touch: the segment + /// whose duration window contains [offset]. + /// + /// Plex media playlists always cover the full title from segment zero, so + /// probing the first entry would steer the transcoder back to the start — + /// requesting a segment is how a client seeks a Plex HLS session. A master + /// playlist (no `#EXTINF` durations) descends into its first variant. A + /// media playlist whose durations never cross [offset] returns null: it + /// cannot say where the offset lives, and a probe aimed at the wrong + /// segment would seek the session, so the caller skips probing instead. + @visibleForTesting + static String? selectReadinessProbeTarget(String body, Duration offset) { + String? firstEntry; + var sawSegmentDurations = false; + var cumulative = Duration.zero; + var pending = Duration.zero; + for (final raw in body.split(RegExp(r'\r?\n'))) { + final line = raw.trim(); + if (line.isEmpty) continue; + if (line.startsWith('#')) { + if (line.startsWith('#EXTINF:')) { + sawSegmentDurations = true; + final seconds = double.tryParse(line.substring('#EXTINF:'.length).split(',').first); + if (seconds != null) pending = Duration(microseconds: (seconds * Duration.microsecondsPerSecond).round()); + } + continue; + } + firstEntry ??= line; + cumulative += pending; + pending = Duration.zero; + if (sawSegmentDurations && cumulative > offset) return line; + } + return sawSegmentDurations ? null : (firstEntry ?? ''); + } + + /// Waits for a just-started Plex offset HLS session to serve the segment at + /// the requested offset before a native player opens its playlist. Plex can + /// return a manifest before the segment is ready; mpv treats that 404 as an + /// HLS error and races through the rest of the manifest. + /// + /// Best-effort by design: the probe never fails an open, it only stops + /// waiting, and callers ignore the returned bool — it exists for tests. The + /// player then sees whatever the server is actually doing and the existing + /// log-stream classification applies unchanged. To that end a 500 stops the + /// wait immediately — a persistent 500 must keep failing fast so the + /// server-limit dialog appears promptly — whether it arrives as a response + /// or inside a decode exception, and a cancellation ([abort] fired or the + /// owning client closing) stops it too rather than sleeping out the window. + /// URLs without an offset return immediately: probing a no-offset playlist + /// would touch segment zero, and requesting a segment is how a client seeks + /// a Plex HLS session. + /// + /// Other non-2xx responses are the expected not-ready signal. `_http.get` + /// does not throw on the status, though its body decode can throw carrying + /// one — both paths share [handOffStatus] so they cannot drift. Every + /// not-ready round waits [pollInterval], doubling up to 4x after three + /// consecutive failed round-trips so a stalled transcode is not hammered; + /// the accepted trade is that a session whose segments 404 for real + /// reaches the player, and its media-unreadable dialog, one probe window + /// later than an unprobed open would. The probe carries this retry budget + /// itself, so its requests bypass endpoint failover, and each request has a + /// hard timeout (5s, shrinking as the overall deadline approaches) so a + /// single hung request cannot consume the entire window. + Future waitForTranscodeReady( + String videoUrl, { + Duration timeout = const Duration(seconds: 15), + Duration pollInterval = const Duration(milliseconds: 500), + AbortController? abort, + }) async { + final startUri = Uri.tryParse(videoUrl); + final probeOffset = transcodeStreamOffsetFromUrl(videoUrl); + if (startUri == null || probeOffset == null) return true; + + // One rule for terminal statuses, applied to responses and to + // status-bearing exceptions alike. + bool handOffStatus(int? statusCode) { + if (statusCode != 500) return false; + // Hand off without classifying: mpv opens the URL, hits the same 500, + // and the log-stream path raises the server-limit dialog. + appLogger.i('Plex transcode readiness probe handing off on HTTP 500'); + return true; + } + + final deadline = clock.now().add(timeout); + var candidate = startUri; + var playlistDepth = 0; + var consecutiveFailures = 0; + int? lastStatus; + while (true) { + final remaining = deadline.difference(clock.now()); + if (remaining <= Duration.zero) break; + if (abort?.isAborted ?? false) return false; + try { + final requestTimeout = remaining < const Duration(seconds: 5) ? remaining : const Duration(seconds: 5); + final isPlaylist = candidate.path.toLowerCase().endsWith('.m3u8'); + // The default Accept is application/json (PlexConfig.headers); the + // probe mirrors the player's request shape instead. Segments go + // through getStatus so a server that ignores Range never routes a + // full media segment through text decoding. + final int statusCode; + var body = ''; + Uri? effectiveUri; + if (isPlaylist) { + final response = await _http.get( + candidate.toString(), + headers: const {'Accept': '*/*'}, + timeout: requestTimeout, + abort: abort, + allowEndpointFailover: false, + ); + statusCode = response.statusCode; + body = response.data?.toString() ?? ''; + effectiveUri = response.effectiveUri; + } else { + final response = await _http.getStatus( + candidate.toString(), + headers: const {'Range': 'bytes=0-0', 'Accept': '*/*'}, + timeout: requestTimeout, + abort: abort, + ); + statusCode = response.statusCode; + } + lastStatus = statusCode; + if (statusCode >= 200 && statusCode < 300) { + consecutiveFailures = 0; + if (body.trimLeft().startsWith('#EXTM3U')) { + final child = selectReadinessProbeTarget(body, probeOffset); + if (child == null) { + // The playlist has segments but its durations never reach the + // offset — a playlist shape this client has never observed + // against a real PMS. It cannot say where the offset lives, + // and a probe aimed at the wrong segment would seek the + // session, so skip probing and let the player negotiate. + return true; + } + if (child.isNotEmpty) { + candidate = (effectiveUri ?? candidate).resolve(child); + playlistDepth++; + if (playlistDepth > 4) { + appLogger.w('Plex transcode readiness exceeded the HLS playlist depth limit'); + return false; + } + // Descending into a child playlist is progress, not a poll. + continue; + } + // A manifest with no media entries yet: not ready, poll again. + } else if (!isPlaylist) { + // The segment at the offset answered: the session is ready. + return true; + } + } else if (handOffStatus(statusCode)) { + return false; + } else { + consecutiveFailures++; + } + } on MediaServerHttpException catch (e) { + if (e.isCancellation) { + // Cancellation is not a not-ready signal, so stop instead of + // sleeping out the window. + return false; + } + lastStatus = e.statusCode ?? lastStatus; + if (handOffStatus(e.statusCode)) return false; + // Transport failure — same treatment as a not-ready response. + consecutiveFailures++; + appLogger.d('Plex transcode readiness probe transport failure', error: e); + } catch (e) { + consecutiveFailures++; + appLogger.d('Plex transcode readiness probe transport failure', error: e); + } + var delay = pollInterval; + if (consecutiveFailures > 3) { + delay = pollInterval * (1 << (consecutiveFailures - 3).clamp(0, 2)); + } + final timeLeft = deadline.difference(clock.now()); + if (timeLeft <= Duration.zero) break; + await Future.delayed(delay < timeLeft ? delay : timeLeft); + } + appLogger.w( + 'Plex transcode did not become ready within ${timeout.inMilliseconds}ms ' + '(playlistDepth=$playlistDepth, lastStatus=${lastStatus ?? 'none'}, consecutiveFailures=$consecutiveFailures)', + ); + return false; + } + /// Build a music transcode stream URL (decision + start path). /// /// Mirrors [buildTranscodeStartPath] for audio tracks: the same @@ -2736,6 +2938,7 @@ class PlexClient required String sessionIdentifier, required String transcodeSessionId, int? audioStreamId, + Duration? offset, }) { final isOriginal = preset.isOriginal; final clientProfileExtra = _buildPlexHlsClientProfileExtra( @@ -2760,6 +2963,7 @@ class PlexClient 'directStreamAudio': '0', 'mediaBufferSize': '102400', 'session': transcodeSessionId, + if (offset != null && offset > Duration.zero) 'offset': (offset.inMilliseconds / 1000).toStringAsFixed(6), 'subtitles': 'none', if (audioStreamId != null) 'audioStreamID': audioStreamId.toString(), 'Accept-Language': 'en', @@ -2789,6 +2993,7 @@ class PlexClient required String sessionIdentifier, required String transcodeSessionId, int? audioStreamId, + Duration? offset, }) { return _buildTranscodeParams( ratingKey: ratingKey, @@ -2798,6 +3003,7 @@ class PlexClient sessionIdentifier: sessionIdentifier, transcodeSessionId: transcodeSessionId, audioStreamId: audioStreamId, + offset: offset, ); } @@ -3330,6 +3536,7 @@ class PlexClient sessionIdentifier: options.sessionIdentifier!, transcodeSessionId: options.transcodeSessionId!, audioStreamId: resolvedAudioId, + offset: options.transcodeOffset, ); if (result.outcome == TranscodeDecisionOutcome.transcodeOk && result.startPath != null) { diff --git a/lib/utils/media_server_http_client.dart b/lib/utils/media_server_http_client.dart index 5daf4001..cc9663d5 100644 --- a/lib/utils/media_server_http_client.dart +++ b/lib/utils/media_server_http_client.dart @@ -177,6 +177,45 @@ class MediaServerHttpClient { ); } + /// Issue a GET and return only status and headers, draining the body + /// unread — the shape for probes that ask "does this answer?" rather than + /// "what does it say?". + /// + /// Unlike [getBytes] the status code is surfaced instead of only logged. + /// Unlike [get] nothing is ever decoded, so a body that fails decoding + /// cannot convert a status into an exception, and — because + /// [FailoverHttpClient] overrides [get] alone — this method structurally + /// never enters the endpoint-failover cascade. Non-2xx is returned, not + /// thrown, matching [get]. + Future getStatus( + String url, { + Map? headers, + Duration? timeout, + AbortController? abort, + }) { + return _perform( + 'GET', + url, + headers: headers, + timeout: timeout, + abort: abort, + consume: (streamed, scope) async { + final effectiveUri = switch (streamed) { + http.BaseResponseWithUrl(:final url) => url, + _ => scope.uri, + }; + await scope.receive(streamed.stream.drain()); + scope.logResponse(streamed.statusCode); + return MediaServerResponse( + statusCode: streamed.statusCode, + headers: streamed.headers, + requestUri: scope.uri, + effectiveUri: effectiveUri, + ); + }, + ); + } + /// Stream-download a URL directly into a file. Future downloadFile( String url, diff --git a/pubspec.lock b/pubspec.lock index dec11cb7..a203af13 100644 --- a/pubspec.lock +++ b/pubspec.lock @@ -215,7 +215,7 @@ packages: source: hosted version: "0.4.2" clock: - dependency: transitive + dependency: "direct main" description: name: clock sha256: fddb70d9b5277016c77a80201021d40a2247104d9f4aa7bab7157b7e3f05b84b diff --git a/pubspec.yaml b/pubspec.yaml index 93dafca8..2ccb0f17 100644 --- a/pubspec.yaml +++ b/pubspec.yaml @@ -10,6 +10,7 @@ environment: dependencies: flutter: sdk: flutter + clock: ^1.1.2 intl: ^0.20.2 json_annotation: ^4.12.0 shared_preferences: ^2.5.4 diff --git a/test/services/plex_playback_data_request_test.dart b/test/services/plex_playback_data_request_test.dart index 84058090..3511c3f7 100644 --- a/test/services/plex_playback_data_request_test.dart +++ b/test/services/plex_playback_data_request_test.dart @@ -1,4 +1,5 @@ import 'dart:convert'; +import 'package:fake_async/fake_async.dart'; import 'package:plezy/media/ids.dart'; import 'package:drift/native.dart'; @@ -6,6 +7,7 @@ import 'package:flutter_test/flutter_test.dart'; import 'package:http/http.dart' as http; import 'package:plezy/database/app_database.dart'; import 'package:plezy/exceptions/media_server_exceptions.dart'; +import 'package:plezy/utils/media_server_http_client.dart' show AbortController; import 'package:plezy/media/media_backend.dart'; import 'package:plezy/media/media_kind.dart'; @@ -319,6 +321,7 @@ void main() { qualityPreset: TranscodeQualityPreset.p720_3mbps, sessionIdentifier: 'session-id', transcodeSessionId: 'transcode-id', + transcodeOffset: const Duration(minutes: 8, seconds: 15), ), ); @@ -326,10 +329,14 @@ void main() { (request) => request.url.path == '/video/:/transcode/universal/decision', ); expect(decisionRequest.url.queryParameters['subtitles'], 'none'); + expect(decisionRequest.url.queryParameters['offset'], '495.000000'); expect(decisionRequest.url.queryParameters.containsKey('subtitleStreamID'), isFalse); expect(decisionRequest.url.queryParameters.containsKey('advancedSubtitles'), isFalse); + expect(decisionRequest.url.queryParameters['X-Plex-Incomplete-Segments'], '1'); expect(result.isTranscoding, isTrue); expect(result.videoUrl, contains('/video/:/transcode/universal/start.m3u8?')); + expect(Uri.parse(result.videoUrl!).queryParameters['offset'], '495.000000'); + expect(PlexClient.transcodeStreamOffsetFromUrl(result.videoUrl!), const Duration(minutes: 8, seconds: 15)); expect(result.subtitleSidecars.map((sidecar) => sidecar.sourceStreamId), [401, 402]); expect(result.subtitleSidecars.every((sidecar) => sidecar.preload), isTrue); expect(result.subtitleSidecars.first.track.isContainer, isTrue); @@ -661,10 +668,373 @@ void main() { expect(startPath, startsWith('/video/:/transcode/universal/start.m3u8?')); expect(startPath, contains('protocol=hls')); + // Without a requested offset the start path must stay offset-free so the + // session transcodes from the beginning exactly as before. expect(startPath, isNot(contains('offset='))); expect(startPath, isNot(contains('X-Plex-Token'))); }); + test('transcode start path carries the requested decision offset', () { + final client = makeClient((_) async => http.Response('not used', 500)); + addTearDown(client.close); + + final params = client.buildTranscodeParamsForTesting( + ratingKey: '42', + mediaIndex: 0, + preset: TranscodeQualityPreset.p720_3mbps, + sessionIdentifier: 'session-id', + transcodeSessionId: 'transcode-id', + offset: const Duration(minutes: 8, seconds: 15), + ); + + final startPath = client.buildTranscodeStartPathFromParamsForTesting(params); + + expect(startPath, startsWith('/video/:/transcode/universal/start.m3u8?')); + expect(Uri.parse(startPath).queryParameters['offset'], '495.000000'); + expect(startPath, isNot(contains('X-Plex-Token'))); + }); + + test('transcode stream offset is parsed from an offset start URL only', () { + expect( + PlexClient.transcodeStreamOffsetFromUrl( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=495.000000', + ), + const Duration(minutes: 8, seconds: 15), + ); + expect( + PlexClient.transcodeStreamOffsetFromUrl( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?session=abc', + ), + isNull, + ); + expect( + PlexClient.transcodeStreamOffsetFromUrl( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=0.000000', + ), + isNull, + ); + expect(PlexClient.transcodeStreamOffsetFromUrl('https://plex.example.com/library/parts/1/file.mkv'), isNull); + // Percent-encoded spelling of the same path must not silently switch the + // probe off. + expect( + PlexClient.transcodeStreamOffsetFromUrl( + 'https://plex.example.com/video/%3A/transcode/universal/start.m3u8?offset=495.000000', + ), + const Duration(minutes: 8, seconds: 15), + ); + }); + + test('transcode readiness probes the segment at the offset, not the first one', () async { + final requests = []; + final rangeHeaders = []; + final acceptHeaders = []; + var segmentAttempts = 0; + // Plex media playlists cover the full title from segment zero; the probe + // must touch the offset's segment because requesting a segment is how a + // client seeks a Plex HLS session — probing 00000.ts would relocate the + // transcoder back to the start. + final mediaPlaylist = StringBuffer('#EXTM3U\n#EXT-X-TARGETDURATION:1\n'); + for (var i = 0; i < 600; i++) { + mediaPlaylist.write('#EXTINF:1.000000,\n${i.toString().padLeft(5, '0')}.ts\n'); + } + final client = makeClient((request) async { + requests.add(request.url); + rangeHeaders.add(request.headers['Range']); + acceptHeaders.add(request.headers['Accept']); + return switch (request.url.path) { + '/video/:/transcode/universal/start.m3u8' => http.Response( + '#EXTM3U\n#EXT-X-STREAM-INF:BANDWIDTH=1500000\nindex.m3u8?session=transcode-id\n', + 200, + ), + '/video/:/transcode/universal/index.m3u8' => http.Response(mediaPlaylist.toString(), 200), + '/video/:/transcode/universal/00495.ts' => + ++segmentAttempts == 1 + ? http.Response('not ready', 404) + : http.Response('segment bytes', 206, headers: {'content-type': 'video/mp2t'}), + _ => http.Response('unexpected request', 500), + }; + }); + addTearDown(client.close); + + final ready = await client.waitForTranscodeReady( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=495.500000', + timeout: const Duration(seconds: 1), + pollInterval: Duration.zero, + ); + + expect(ready, isTrue); + expect(segmentAttempts, 2); + expect(requests.map((uri) => uri.path), [ + '/video/:/transcode/universal/start.m3u8', + '/video/:/transcode/universal/index.m3u8', + '/video/:/transcode/universal/00495.ts', + '/video/:/transcode/universal/00495.ts', + ]); + expect(requests[1].queryParameters['session'], 'transcode-id'); + expect(rangeHeaders, [null, null, 'bytes=0-0', 'bytes=0-0']); + // The client's default Accept is application/json; the probe must mirror + // the player's request shape instead. + expect(acceptHeaders, everyElement('*/*')); + }); + + test('readiness probe target selection walks segment durations to the offset', () { + const mediaPlaylist = + '#EXTM3U\n' + '#EXT-X-TARGETDURATION:2\n' + '#EXTINF:2.000000,\n00000.ts\n' + '#EXTINF:2.000000,\n00001.ts\n' + '#EXTINF:2.000000,\n00002.ts\n'; + expect(PlexClient.selectReadinessProbeTarget(mediaPlaylist, Duration.zero), '00000.ts'); + expect(PlexClient.selectReadinessProbeTarget(mediaPlaylist, const Duration(seconds: 3)), '00001.ts'); + // Durations never cross the offset: the playlist cannot say where the + // offset lives, and probing a guessed segment would seek the session. + expect(PlexClient.selectReadinessProbeTarget(mediaPlaylist, const Duration(seconds: 30)), isNull); + // A master playlist has no segment durations: descend into the first variant. + const masterPlaylist = '#EXTM3U\n#EXT-X-STREAM-INF:BANDWIDTH=1500000\nindex.m3u8\nfallback.m3u8\n'; + expect(PlexClient.selectReadinessProbeTarget(masterPlaylist, const Duration(minutes: 10)), 'index.m3u8'); + }); + + test('transcode readiness skips probing when the playlist cannot locate the offset', () { + fakeAsync((async) { + final requestPaths = []; + // The offset-relative playlist shape this client has never observed + // against a real PMS: segments named from the resume point, durations + // summing to a minute. Probing its last segment would ask the session + // to produce a minute past the resume point. + final mediaPlaylist = StringBuffer('#EXTM3U\n#EXT-X-TARGETDURATION:1\n'); + for (var i = 8062; i <= 8121; i++) { + mediaPlaylist.write('#EXTINF:1.000000,\n${i.toString().padLeft(5, '0')}.ts\n'); + } + final client = makeClient((request) async { + requestPaths.add(request.url.path); + if (request.url.path.endsWith('/start.m3u8')) { + return http.Response(mediaPlaylist.toString(), 200); + } + return http.Response('segment bytes', 206, headers: {'content-type': 'video/mp2t'}); + }); + addTearDown(client.close); + + bool? ready; + client + .waitForTranscodeReady( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=8062.000000', + timeout: const Duration(seconds: 15), + pollInterval: const Duration(milliseconds: 500), + ) + .then((value) => ready = value); + async.flushMicrotasks(); + + expect(ready, isTrue); + expect(requestPaths.where((path) => path.endsWith('.ts')), isEmpty); + }); + }); + + test('transcode readiness hands off on a 500 that arrives inside a decode exception', () { + fakeAsync((async) { + var requestCount = 0; + // A proxy or gateway in front of PMS answering 500 with a JSON + // content-type and a non-JSON body: the client's body decode throws a + // status-bearing exception instead of returning the response, and the + // exception path must apply the same hand-off rule as the response + // path. + final client = makeClient((request) async { + requestCount++; + return http.Response('gateway error', 500, headers: {'content-type': 'application/json'}); + }); + addTearDown(client.close); + + bool? ready; + client + .waitForTranscodeReady( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=495.000000', + timeout: const Duration(seconds: 15), + pollInterval: const Duration(milliseconds: 500), + ) + .then((value) => ready = value); + // No elapse: like the response-path 500, the exception-path 500 must + // complete the probe without a single poll wait. + async.flushMicrotasks(); + + expect(ready, isFalse); + expect(requestCount, 1); + }); + }); + + test('transcode readiness returns immediately when the probe is already aborted', () async { + var requestCount = 0; + final client = makeClient((request) async { + requestCount++; + return http.Response('#EXTM3U\n#EXTINF:1.0,\nmedia-00000.ts\n', 200); + }); + addTearDown(client.close); + + final abort = AbortController()..abort(); + final ready = await client.waitForTranscodeReady( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=0.500000', + abort: abort, + ); + + expect(ready, isFalse); + expect(requestCount, 0); + }); + + test('transcode readiness polls a bounded number of times while the segment stays unavailable', () { + fakeAsync((async) { + var requestCount = 0; + final client = makeClient((request) async { + requestCount++; + if (request.url.path.endsWith('/start.m3u8')) { + return http.Response('#EXTM3U\n#EXTINF:1.0,\nmedia-00000.ts\n', 200); + } + return http.Response('not ready', 404); + }); + addTearDown(client.close); + + bool? ready; + client + .waitForTranscodeReady( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=0.500000', + timeout: const Duration(milliseconds: 600), + pollInterval: const Duration(milliseconds: 50), + ) + .then((value) => ready = value); + async.elapse(const Duration(milliseconds: 700)); + + expect(ready, isFalse); + // The not-ready signal is a non-2xx *response*, not an exception: every + // failed probe must still wait out the poll interval, and after three + // consecutive failures the interval doubles to a 4x cap. Under a fake + // clock the cadence is exact: the playlist hop, then segment polls at + // 0/50/100/150/250/450ms. Drafts of this probe that skipped the delay + // on non-2xx or discarded the backoff measured 3,651 and 12 requests + // respectively in this same window; neither shape ever shipped. + expect(requestCount, 7); + }); + }); + + test('transcode readiness requests bypass endpoint failover', () { + fakeAsync((async) { + var exhaustedSignals = 0; + var requestCount = 0; + final client = testPlexClient( + serverId: ServerId('server-id'), + prioritizedEndpoints: const ['https://plex.example.com'], + onAllEndpointsExhausted: () => exhaustedSignals++, + // 503 on the playlist itself: the segment leg goes through getRaw, + // which never enters the failover path, so the playlist request is + // the one that could drive the cascade. + handler: (request) async { + requestCount++; + return http.Response('busy', 503); + }, + ); + addTearDown(client.close); + + bool? ready; + client + .waitForTranscodeReady( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=495.000000', + timeout: const Duration(milliseconds: 200), + pollInterval: const Duration(milliseconds: 50), + ) + .then((value) => ready = value); + async.elapse(const Duration(milliseconds: 300)); + + expect(ready, isFalse); + expect(requestCount, greaterThanOrEqualTo(3)); + // The probe carries its own retry budget, so it must not drive the + // endpoint-failover cascade: its URL is absolute (the retry re-hits the + // same host), each switch rewrites config.baseUrl underneath the + // videoUrl already handed to the open path, and on a single-endpoint + // server every failed poll fires the all-endpoints-exhausted signal — + // the manager's cue to flip server status and reconnect. + expect(exhaustedSignals, 0); + expect(client.config.baseUrl, 'https://plex.example.com'); + }); + }); + + test('transcode readiness stops on cancellation instead of sleeping out the window', () { + fakeAsync((async) { + var requestCount = 0; + late final PlexClient client; + client = testPlexClient( + serverId: ServerId('server-id'), + handler: (request) async { + requestCount++; + if (request.url.path.endsWith('/start.m3u8')) { + return http.Response('#EXTM3U\n#EXTINF:1.0,\nmedia-00000.ts\n', 200); + } + // The owner closes mid-probe; the next request must surface as a + // cancellation, not as one more not-ready round. + client.close(); + return http.Response('not ready', 404); + }, + ); + + bool? ready; + client + .waitForTranscodeReady( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=0.500000', + timeout: const Duration(milliseconds: 600), + pollInterval: const Duration(milliseconds: 50), + ) + .then((value) => ready = value); + // One poll interval is all it may consume after the close; a probe + // that counts cancellation as not-ready sleeps out the full 600ms. + async.elapse(const Duration(milliseconds: 100)); + + expect(ready, isFalse); + expect(requestCount, 2); + }); + }); + + test('transcode readiness hands off immediately on HTTP 500', () { + fakeAsync((async) { + var requestCount = 0; + final client = makeClient((request) async { + requestCount++; + if (request.url.path.endsWith('/start.m3u8')) { + return http.Response('#EXTM3U\n#EXTINF:1.0,\nmedia-00000.ts\n', 200); + } + return http.Response('limit rejected', 500); + }); + addTearDown(client.close); + + bool? ready; + client + .waitForTranscodeReady( + 'https://plex.example.com/video/:/transcode/universal/start.m3u8?offset=0.500000', + timeout: const Duration(seconds: 15), + pollInterval: const Duration(milliseconds: 500), + ) + .then((value) => ready = value); + // No elapse: a 500 must complete the probe without a single poll wait, + // so the player opens promptly and the server-limit dialog path runs. + async.flushMicrotasks(); + + expect(ready, isFalse); + expect(requestCount, 2); + }); + }); + + test('transcode readiness returns immediately for URLs without an offset', () async { + var requestCount = 0; + final client = makeClient((request) async { + requestCount++; + return http.Response('should not be called', 500); + }); + addTearDown(client.close); + + // Probing a no-offset playlist would touch segment zero, and requesting + // a segment is how a client seeks a Plex HLS session. + expect( + await client.waitForTranscodeReady('https://plex.example.com/video/:/transcode/universal/start.m3u8?session=x'), + isTrue, + ); + expect(await client.waitForTranscodeReady('https://plex.example.com/library/parts/1/file.mkv'), isTrue); + expect(requestCount, 0); + }); + test('transcode params preserve resolved media and part indices', () { final client = makeClient((_) async => http.Response('not used', 500)); addTearDown(client.close);