Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -209,8 +210,10 @@ private byte[] runCollecting(List<String> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -92,6 +94,10 @@ public class StreamServiceImpl implements StreamService {

private final Map<Long, MediaInfo> probes = new ConcurrentHashMap<>();

/** One in-flight encoder per segment, shared by every browser/session. */
private final Map<SegmentCacheKey, CompletableFuture<byte[]>> segmentEncodes =
new ConcurrentHashMap<>();

/** TMDB's original_language per movie. One lookup per movie, then served from memory. */
private final Map<Long, String> originalLanguages = new ConcurrentHashMap<>();

Expand All @@ -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);
Expand Down Expand Up @@ -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<byte[]> mine = new CompletableFuture<>();
CompletableFuture<byte[]> 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
Expand Down
5 changes: 3 additions & 2 deletions src/main/resources/application.properties
Original file line number Diff line number Diff line change
Expand Up @@ -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}}
Expand Down
5 changes: 3 additions & 2 deletions src/test/resources/application.properties
Original file line number Diff line number Diff line change
Expand Up @@ -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}}
Expand Down
Loading