diff --git a/src/main/java/com/filmexa/stream/modules/streaming/config/StreamProperties.java b/src/main/java/com/filmexa/stream/modules/streaming/config/StreamProperties.java index 25accba..7d1962d 100644 --- a/src/main/java/com/filmexa/stream/modules/streaming/config/StreamProperties.java +++ b/src/main/java/com/filmexa/stream/modules/streaming/config/StreamProperties.java @@ -20,9 +20,13 @@ public class StreamProperties { private String ffprobePath = "ffprobe"; - private String preset = "veryfast"; + private String preset = "ultrafast"; - private int maxConcurrentTranscodes = 4; + /** One active encoder per simultaneous viewer on the expected two-core host. */ + private int maxConcurrentTranscodes = 2; + + /** Prevent one ffmpeg process from taking every core and starving another viewer. */ + private int encoderThreads = 1; private int transcodeTimeoutSeconds = 120; diff --git a/src/main/java/com/filmexa/stream/modules/streaming/ffmpeg/Ffmpeg.java b/src/main/java/com/filmexa/stream/modules/streaming/ffmpeg/Ffmpeg.java index 0db59dc..b06fd5a 100644 --- a/src/main/java/com/filmexa/stream/modules/streaming/ffmpeg/Ffmpeg.java +++ b/src/main/java/com/filmexa/stream/modules/streaming/ffmpeg/Ffmpeg.java @@ -152,6 +152,7 @@ public byte[] encodeSegment(MediaInfo info, Resolution resolution, int segmentIn "-vf", "scale=-2:" + resolution.getHeight(), "-c:v", "libx264", "-preset", properties.getPreset(), + "-threads", String.valueOf(Math.max(1, properties.getEncoderThreads())), "-profile:v", "high", "-level", "4.1", "-pix_fmt", "yuv420p", @@ -209,8 +210,10 @@ private byte[] runCollecting(List command, int segmentIndex) { errThread.join(5000); if (process.exitValue() != 0) { - throw new IllegalStateException("ffmpeg failed on segment " + segmentIndex - + ": " + errors.toString().trim()); + log.debug("ffmpeg could not read segment {} yet: {}", segmentIndex, + errors.toString().trim()); + throw new StreamNotReadyException( + "Segment " + segmentIndex + " is not readable yet", 5); } // Asked to seek past the end of the data actually on disk, ffmpeg exits 0 and diff --git a/src/main/java/com/filmexa/stream/modules/streaming/service/StreamService.java b/src/main/java/com/filmexa/stream/modules/streaming/service/StreamService.java index 48160f2..c97b606 100644 --- a/src/main/java/com/filmexa/stream/modules/streaming/service/StreamService.java +++ b/src/main/java/com/filmexa/stream/modules/streaming/service/StreamService.java @@ -32,6 +32,6 @@ public interface StreamService { * @param audioTrackIndex which of the file's audio streams to encode - the original * language, resolved once when the session is prepared */ - record Segment(MediaInfo info, Resolution resolution, int index, int audioTrackIndex) { + record Segment(Long movieId, MediaInfo info, Resolution resolution, int index, int audioTrackIndex) { } } diff --git a/src/main/java/com/filmexa/stream/modules/streaming/serviceImpl/StreamServiceImpl.java b/src/main/java/com/filmexa/stream/modules/streaming/serviceImpl/StreamServiceImpl.java index 12f6acf..c79dca1 100644 --- a/src/main/java/com/filmexa/stream/modules/streaming/serviceImpl/StreamServiceImpl.java +++ b/src/main/java/com/filmexa/stream/modules/streaming/serviceImpl/StreamServiceImpl.java @@ -16,6 +16,8 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -92,6 +94,10 @@ public class StreamServiceImpl implements StreamService { private final Map probes = new ConcurrentHashMap<>(); + /** One in-flight encoder per segment, shared by every browser/session. */ + private final Map> segmentEncodes = + new ConcurrentHashMap<>(); + /** TMDB's original_language per movie. One lookup per movie, then served from memory. */ private final Map originalLanguages = new ConcurrentHashMap<>(); @@ -102,9 +108,6 @@ public StreamSessionDto createSession(Long movieId, String imdbId, User viewer) MediaInfo info; try { info = mediaInfo(movieId); - // READY has to mean "the player will actually get its first segment", not just - // "the download flipped a flag" - so check the real gate the segments use. - requireDownloaded(movieId, info, 0); } catch (StreamNotReadyException | NotFoundException notReadyYet) { log.debug("Movie {} not playable yet: {}", movieId, notReadyYet.getMessage()); return preparing(movieId, download); @@ -227,14 +230,114 @@ public Segment prepareSegment(Long movieId, int height, int segmentIndex) { throw new NotFoundException("Segment " + segmentIndex + " is outside movie " + movieId); } - requireDownloaded(movieId, info, segmentIndex); - return new Segment(info, resolution, segmentIndex, originalAudioIndex(movieId, info)); + // For the first segment ffmpeg is the authoritative readiness check. Mapping a + // timestamp to torrent pieces is only an estimate (especially for VBR MP4), and + // can reject a segment that the decoder can already read. A not-yet-readable + // source becomes a retryable 503 in Ffmpeg rather than holding the session in + // PREPARING indefinitely. Later segments still use the piece gate for fast seeks. + if (segmentIndex > 0) { + requireDownloaded(movieId, info, segmentIndex); + } + return new Segment(movieId, info, resolution, segmentIndex, originalAudioIndex(movieId, info)); } @Override public byte[] segmentBytes(Segment segment) { - return ffmpeg.encodeSegment(segment.info(), segment.resolution(), segment.index(), - segment.audioTrackIndex()); + Path cached = segmentCacheFile(segment); + byte[] existing = readCachedSegment(cached); + if (existing != null) { + return existing; + } + + SegmentCacheKey key = new SegmentCacheKey(segment.movieId(), + segment.resolution().getHeight(), segment.index(), segment.audioTrackIndex()); + CompletableFuture mine = new CompletableFuture<>(); + CompletableFuture running = segmentEncodes.putIfAbsent(key, mine); + + if (running != null) { + try { + return running.join(); + } catch (CompletionException failure) { + throw rethrow(failure.getCause()); + } + } + + try { + // Check again after winning the single-flight slot: a previous encoder may + // have completed between the first disk read and putIfAbsent. + existing = readCachedSegment(cached); + if (existing != null) { + mine.complete(existing); + return existing; + } + + byte[] encoded = ffmpeg.encodeSegment(segment.info(), segment.resolution(), + segment.index(), segment.audioTrackIndex()); + writeCachedSegment(cached, encoded); + mine.complete(encoded); + return encoded; + } catch (RuntimeException failure) { + mine.completeExceptionally(failure); + throw failure; + } finally { + segmentEncodes.remove(key, mine); + } + } + + private Path segmentCacheFile(Segment segment) { + return Paths.get(videoStoragePath, segment.movieId().toString(), "hls", + segment.resolution().getHeight() + "p-a" + segment.audioTrackIndex(), + "seg-" + segment.index() + ".ts").toAbsolutePath(); + } + + private byte[] readCachedSegment(Path cached) { + if (!Files.isRegularFile(cached)) { + return null; + } + try { + byte[] bytes = Files.readAllBytes(cached); + return bytes.length == 0 ? null : bytes; + } catch (IOException e) { + log.warn("Could not read cached HLS segment {}: {}", cached, e.getMessage()); + return null; + } + } + + private void writeCachedSegment(Path destination, byte[] bytes) { + Path temporary = null; + try { + Files.createDirectories(destination.getParent()); + temporary = Files.createTempFile(destination.getParent(), "seg-", ".ts.part"); + Files.write(temporary, bytes); + try { + Files.move(temporary, destination, StandardCopyOption.REPLACE_EXISTING, + StandardCopyOption.ATOMIC_MOVE); + } catch (java.nio.file.AtomicMoveNotSupportedException unsupported) { + Files.move(temporary, destination, StandardCopyOption.REPLACE_EXISTING); + } + temporary = null; + } catch (IOException e) { + // Caching is an optimisation. The encoded response is still valid, so a + // read-only/full cache directory must not break playback. + log.warn("Could not cache HLS segment {}: {}", destination, e.getMessage()); + } finally { + if (temporary != null) { + try { + Files.deleteIfExists(temporary); + } catch (IOException ignored) { + log.debug("Could not remove temporary HLS segment {}", temporary); + } + } + } + } + + private RuntimeException rethrow(Throwable failure) { + return failure instanceof RuntimeException runtime + ? runtime + : new IllegalStateException("Concurrent segment encoding failed", failure); + } + + private record SegmentCacheKey(Long movieId, int height, int segmentIndex, int audioTrackIndex) { } @Override diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties index 5c7b468..650c1eb 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -85,8 +85,9 @@ piratebay.base-url=${PIRATEBAY_BASE_URL} app.stream.segment-seconds=${STREAM_SEGMENT_SECONDS:6} app.stream.ffmpeg-path=${FFMPEG_PATH:ffmpeg} app.stream.ffprobe-path=${FFPROBE_PATH:ffprobe} -app.stream.preset=${STREAM_PRESET:veryfast} -app.stream.max-concurrent-transcodes=${STREAM_MAX_TRANSCODES:4} +app.stream.preset=${STREAM_PRESET:ultrafast} +app.stream.max-concurrent-transcodes=${STREAM_MAX_TRANSCODES:2} +app.stream.encoder-threads=${STREAM_ENCODER_THREADS:1} app.stream.transcode-timeout-seconds=${STREAM_TRANSCODE_TIMEOUT:120} app.stream.probe-timeout-seconds=${STREAM_PROBE_TIMEOUT:30} app.stream.token-secret=${STREAM_TOKEN_SECRET:${SECURITY_JWT_SECRET_KEY}} diff --git a/src/test/resources/application.properties b/src/test/resources/application.properties index f50fd93..c24cc01 100644 --- a/src/test/resources/application.properties +++ b/src/test/resources/application.properties @@ -65,8 +65,9 @@ piratebay.base-url=${PIRATEBAY_BASE_URL:https://piratebay.org} app.stream.segment-seconds=${STREAM_SEGMENT_SECONDS:6} app.stream.ffmpeg-path=${FFMPEG_PATH:ffmpeg} app.stream.ffprobe-path=${FFPROBE_PATH:ffprobe} -app.stream.preset=${STREAM_PRESET:veryfast} -app.stream.max-concurrent-transcodes=${STREAM_MAX_TRANSCODES:4} +app.stream.preset=${STREAM_PRESET:ultrafast} +app.stream.max-concurrent-transcodes=${STREAM_MAX_TRANSCODES:2} +app.stream.encoder-threads=${STREAM_ENCODER_THREADS:1} app.stream.transcode-timeout-seconds=${STREAM_TRANSCODE_TIMEOUT:120} app.stream.probe-timeout-seconds=${STREAM_PROBE_TIMEOUT:30} app.stream.token-secret=${STREAM_TOKEN_SECRET:${SECURITY_JWT_SECRET_KEY}}