* fix(player): keep text subtitles in sync across copy-mode restarts and sparse cue windows - Rebase already-loaded cues in place when streamOriginSeconds changes (copy-mode session restart) instead of leaving them offset by the delta. - Stop inferring end-of-input from where a window's cues stop; only the known media duration marks EOF, so a dialogue gap no longer silently ends prefetching for the rest of playback. - Anchor the first window fetch to the intended start position (resume target or pending seek) while the element still reports currentTime=0, and reset coverage on forward seeks past the fetched window. * fix(playback): align encoded transcode start with the declared segment boundary A mid-segment start (resume, seek restart, audio switch) spawned ffmpeg at the raw seek position while labeling its first segment with the grid number, whose synthetic-manifest start is up to one segment earlier. hls.js aligns the first fragment's content to that declared position, shifting the whole session's timeline late by seek mod segment_duration (0-2s): subtitles trail dialogue by a constant per-session offset and progress/resume positions drift by the same amount. Snap the ffmpeg start position down to the segment boundary for encoded sessions so declared and produced timelines match exactly; the player still seeks to the precise requested position. Copy-mode sessions serve ffmpeg's real manifest and keep the raw seek. * fix(api): forward http.Flusher through response-writer middleware wrappers Streamed subtitle extracts (and any progressive response) flush per chunk via an http.Flusher assertion, but none of the status-capturing middleware wrappers implemented Flush, so the assertion failed and cues sat in Go's response buffer until ffmpeg finished. On large remuxes where a 600s window takes 20s+ to demux, captions appeared only when the whole window completed instead of within the first seconds. Give every wrapper a Flush() (satisfies plain assertions, including chi's Compress) and Unwrap() (satisfies http.ResponseController). The jellycompat image-proxy tag rewriter flushes only in passthrough mode since it buffers JSON bodies for rewriting. Regression test asserts the API chain forwards Flush end to end. * feat(player): PGS subtitles honor size, position, and background settings Port the tvOS/iOS bitmap-cue styling to the web player. libpgs now decodes in worker mode but draws to a hidden source canvas on the main thread; a compositor detects cue regions from the frame's alpha channel and re-places them on a visible overlay canvas per the shared subtitle appearance settings: size scale (with the 0.85 authored-size compensation), vertical position preset (dialogue-band cues only — floating signs keep authored placement, matching the Apple implementation), and the background box. Font family, text color, and outline are baked into the source pixels and remain inapplicable. * fix(player): anchor PGS position presets to the text overlay's reference frame The initial port used silo-apple's 30/1080 bottom margin and video-relative lower-third/top anchors; the web text overlay anchors to a 16:9 reference frame with 7%/18% offsets that extends into the letterbox for wide content. Use the same anchors so PGS dialogue lands exactly where SRT text does. * feat(player): size PGS cues to match the text subtitle line height Replace the authored-size ladder (0.85 × font-size ratio) with per-cue text-line matching: the region detector reports the tallest text line inside each cue, and the compositor scales the cue so one line of bitmap text renders at the same pixel height as the SRT overlay's font at the current preset. Authored size differences between discs no longer leak through; upscaling is capped at 2.5× to keep small bitmaps from going blurry. * feat(subtitles): opt-in windowed PGS extraction to cut mid-file load latency PGS extracts always demuxed the source from byte 0, so starting a large remux mid-file meant minutes before the first bitmap cue. The web player now opts in to windowed extraction (?windowed=1&position=&duration=) and re-points libpgs at a fresh window on seeks and near coverage end; ffmpeg input-side -ss with -copyts keeps absolute source timestamps. Without the explicit opt-in the endpoint behaves byte-identically, so Apple/Android and other single-fetch consumers are unaffected. ASS remains unconditionally non-windowed (its header only exists at offset 0). * feat(playback): cache extracted PGS subtitle tracks Every selection of an embedded PGS track re-ran a full ffmpeg extract that demuxes the entire source file from byte 0 — minutes for a large remux — and responses were Cache-Control: no-store, so repeat selections, re-watches, and multiple viewers all paid full price. Add a disk cache for full-track .sup extracts under <transcode_dir>/subtitle-cache, created lazily: - Keyed by source path hash + subtitle stream ordinal + source mtime+size (encoded in the filename), so a replaced source file implicitly invalidates its entries; the source is stat'ed on every lookup. - Cache miss: ffmpeg stdout is teed to the response (first viewer still streams progressively, first-byte latency unchanged) and into a temp file that is fsynced and atomically renamed into the cache on clean ffmpeg exit. Any error — ffmpeg failure, client disconnect, tee write failure, or the source changing mid-extract — discards the temp file, so a partial entry is never served. - Cache hit: served via http.ServeContent (Range support, Content-Length, Last-Modified from the source mtime) with a revalidatable Cache-Control instead of no-store. - Concurrent requests for the same in-flight track run their own uncached extract (mutex + in-flight key set) rather than blocking on another client's connection. - Scan-on-commit LRU eviction under a 2 GiB cap (recency tracked by bumping entry mtime on hit; atime is unreliable under relatime), plus sweep of crash-orphaned .part temp files. - Windowed PGS requests (?windowed=) bypass the cache in both directions: their output covers only a slice of the track. Both the integrated API handler and the standalone proxy subtitle path share the same playback.SubtitleCache.ServeSUPExtract helper. VTT (already windowed and fast) and ASS (small) stay uncached. No API surface change. AI-use disclosure: implemented with Claude Code. * fix(playback): check Close error returns in subtitle cache paths Silence errcheck on the cache-hit defer and the test's simulated disk-full Close. AI-use disclosure: implemented with Claude Code. * feat(playback): warm PGS cache in background and window from cached track Windowed PGS requests bypassed the cache entirely, so every window fetch re-demuxed the multi-GB original file. Now a windowed miss kicks off a detached background warm (full-track extract into the cache, at most 2 concurrent server-wide, coalesced with client-driven fills), and once the entry exists windowed extracts read the 15-80MB cached .sup instead — seeks and re-enables become near-instant after the first load. Verified empirically that ffmpeg preserves absolute PTS when windowing a sup input. * feat(player): hold playback while PGS subtitle cues load When a PGS track is enabled (or a seek lands outside the fetched window), extraction takes seconds and dialogue could play unsubtitled. The player now pauses until the renderer's parsed data covers the playhead — tracked via libpgs' parsed-timestamp watermark, the exact predicate it renders by — showing a "Loading subtitles…" indicator after 500ms. User play/pause always wins over the hold, a 20s safety timeout prevents stranding playback, and background prefetch never pauses. If future libpgs versions reshape the observed internals the hook degrades to the old play-through behavior. * perf(player): shrink uncached PGS window to 600s Draining a windowed extract from the source reads the full interleaved container across the window (~1GB per 100s of remux on measured hardware); a 3600s window cost ~12GB of reads per fetch while cold. Once the server cache is warm a window costs milliseconds regardless of size, so smaller windows only add trivially cheap re-fetches. * feat(subtitles): burn in PGS/bitmap subtitles for the web player The web player rendered PGS client-side via libpgs, which required extracting the .sup track — a cold ffmpeg demux that took seconds even windowed, since c:s copy still reads the whole interleaved container up to the playhead. Every other server (Plex, Jellyfin default, Emby) burns image subtitles into the video instead, and that is the only path with no per-seek extraction cost. Selecting a bitmap subtitle (PGS/DVD/DVB) now restarts the transcode with subtitle_burn_in at the current aligned position, reusing the same restart machinery as an audio/quality switch so the segment-boundary timeline alignment holds. The server composites the decoded subtitle onto the video with an overlay filter_complex graph (libass's subtitles= filter is text-only); overlay runs at native resolution before any target scaling, and hardware pipelines round-trip through CPU like the text path. Burn-in forces a video encode, so copy-video recipes are upgraded to h264 both client- and server-side. Text subtitles keep the instant, styled, client-side path. The .sup streaming endpoints, cache, and windowing are retained for the Apple client, which renders PGS natively. The now-dead web PGS stack (usePGSSubtitles, pgsPlacement, libpgs dep) is removed. Tradeoff: bitmap subtitles no longer honor web appearance settings (baked into the video) and toggling one restarts the transcode (~1-2s buffering), matching Plex behavior. * fix(player): rebuild text subtitle track when turning off PGS burn-in Selecting an SRT track that turned off bitmap burn-in restarted the transcode, and the client TextTrack built in that same moment was orphaned when the <video> element reloaded, so the subtitles never rendered (and a seek could not recover the dead track). Rebuild the text track once the new stream settles, gated on the burn-in-off transition so quality/audio switches and copy-mode seek restarts keep their subtitles without a needless re-extract. * fix(player): render web subtitles behind the control HUD The text subtitle overlay sat at z-20, above the controls layer (z-10), so cues painted over the bottom HUD and cluttered the control bar. Drop it to z-[5] — above the video, below the controls — so the HUD paints over the cues while it is visible. When controls are hidden the whole controls layer is opacity-0, so cues remain fully visible. * feat(player): lift web subtitles above the control bar while it's visible Rather than hiding bottom-anchored cues behind the HUD, raise them just above the control bar (measured height + a small gap) whenever the bar is visible in the foreground player, then settle them back when it hides. The bar is a roughly fixed pixel height while the cue offset scales with the player, so the bar is measured via ResizeObserver rather than hardcoded. Top-anchored cues never collide with the bottom HUD, so they stay put. z-[5] is retained as a safety so any residual overlap tucks behind the bar. * fix(player): coalesce same-tick transcode restarts into one dispatch Starting playback with a persisted bitmap subtitle fired transcode/start twice within milliseconds: the auto-start effect dispatched before subtitle auto-selection restored the burn-in, whose effect then forced a second start. The first request was already on the wire (no abort signal was passed to fetch), so the server spawned an ffmpeg only to kill it for the second start — visible in production as an ffmpeg exit error ~1ms after every such session start, and slowing time to first frame. Defer the network dispatch by one macrotask so back-to-back restart calls in a tick collapse into a single request carrying the final parameters; state updates stay synchronous. Pass the abort signal into playerFetch so a superseded in-flight request is actually cancelled, and drop any deferred dispatch on unmount so a stray transcode/start cannot land after the session's exit DELETE. * fix(catalog): resolve effective subtitle defaults for movie item details Movie pre-play subtitle selectors were missing the effective defaults (including per-item overrides saved from a previous play) that episodes and watch payloads already resolve. Extract applyToItemDetail/ applyToWatchDetail helpers and apply defaults for movies in buildMediaItemDetail. The SubtitlesPopover now also eagerly loads downloaded subtitles when the saved preference points at one so the closed trigger's Auto summary reflects the override. * feat(player): scale subtitle font size with the rendered video Replace fixed rem font sizes with px values defined at a 720px 16:9 reference height, scaled proportionally with the actually-rendered video (object-fit: contain) so subtitles keep the same relative size as the window grows or shrinks, with a 12px legibility floor. Rename useSubtitlePositionStyle to useSubtitleLayout, returning both the position style and the font scale, and add unit tests for the appearance helpers. * fix(player): satisfy strict index checks in transcode quality test * feat(player): let the pre-play Auto option clear the saved subtitle override A manual in-player subtitle selection persists as an 'always' override for that movie/series, but nothing in the UI could undo it — auto selection stayed pinned to the chosen track forever. Choosing 'Auto' in the pre-play subtitles popover now also deletes the stored preference (movie content ID / episode series ID) and invalidates item details so profile-level auto selection applies again. * feat(player): persist pre-play subtitle selections as the item override Choosing a track (or Off) in the pre-play subtitles popover only lived in component state: it applied to that playback session but vanished on returning to the detail page. Persist it through PUT /subtitle-prefs — the same 'always'/'off' override a manual in-player selection saves — keyed by movie content ID or episode series ID, and invalidate item details so the effective defaults reflect it immediately. * feat(ui): show the saved subtitle override and richer pre-play pill summaries A stored per-item override displayed as 'Auto: <language>', hiding both that an override exists and which track it is. The pre-play subtitle pill now shows the resolved track directly (name with (SDH)/(Forced) markers plus format, skipping markers the name already carries), the matching list row gets the checkmark instead of Auto, and the Auto row reads 'Reset to profile defaults'. Subtitle, audio, edition, and version pill summaries also truncate much later (max-w-44/sm:max-w-64). * fix(player): recover text subtitles from stream reloads and failed window fetches Three failure modes could silently freeze or stop web text subtitles: - A stream restart (seek-triggered transcode restart, quality/audio switch) reloads the <video> element and can orphan the programmatic TextTrack — cuechange stops firing and the last cue freezes on screen. Only the PGS-burn-in-off transition rebuilt the track. Now every settled stream URL change bumps the generation, and the rebuild carries loaded cues (converted back to source time) and window coverage over so it costs no refetch. - The sliding-window fetcher committed windowEnd before the fetch ran, so a failed or hung window counted as covered and was never retried — subtitles silently stopped for up to 10 minutes. Coverage now commits only after the window streams in fully; failures leave the range uncovered and retry after a 5s backoff. - A hung extraction (one fetch in flight at a time, no deadline) blocked every future window for the session. Reads now arm a 30s stall timer that aborts a response which stops delivering chunks; slow-but- progressing streams keep resetting it. Diagnosed from a session where ffmpeg took 69s to stream one subtitle window and a transcode restart landed mid-fetch, freezing the active cue. * fix(playback): keep subtitle selections stable across file changes * fix(web): clarify subtitle labels and positioning * fix(player): hide HUD when pointer leaves * fix(web): tidy subtitle track badges * fix(playback): remap audio tracks across file versions * fix(web): tidy audio track labels * docs(playback): clarify bitmap subtitle appearance * fix(playback): preserve selection state on restart * fix(http): preserve response state across flushes * fix(web): preserve pending quality for burn-in * fix(playback): preserve subtitle inventory identity --------- Co-authored-by: Quick <31828688+Quick104@users.noreply.github.com>
2073 lines
65 KiB
Go
2073 lines
65 KiB
Go
package playback
|
||
|
||
import (
|
||
"bytes"
|
||
"context"
|
||
"errors"
|
||
"fmt"
|
||
"io"
|
||
"log"
|
||
"log/slog"
|
||
"math"
|
||
"mime"
|
||
"os"
|
||
"os/exec"
|
||
"path/filepath"
|
||
"strconv"
|
||
"strings"
|
||
"sync"
|
||
"syscall"
|
||
"time"
|
||
)
|
||
|
||
func init() {
|
||
// Register .m4s so http.ServeFile sets the correct Content-Type for
|
||
// fMP4 HLS segments. Go's default MIME database does not include it.
|
||
_ = mime.AddExtensionType(".m4s", "video/mp4")
|
||
}
|
||
|
||
// TranscodeOpts holds configuration for an HLS transcode session.
|
||
type TranscodeOpts struct {
|
||
InputPath string
|
||
OutputDir string // e.g., /tmp/silo-transcode/{session_id}/
|
||
SessionID string
|
||
SourceVideoCodec string
|
||
SeekSeconds float64
|
||
TargetResolution string // e.g., 1080p, 720p
|
||
TargetCodecVideo string // e.g., h264 (or hevc if allowed)
|
||
TargetCodecAudio string // e.g., aac
|
||
SegmentDuration int // seconds, default 6
|
||
StartSegmentNumber int // -hls_segment_start_number, default 0
|
||
FFmpegPath string // optional explicit ffmpeg binary path
|
||
HWAccel string // auto, qsv, vaapi, nvenc, none
|
||
HWDevice string // e.g., /dev/dri/renderD128 (default if empty)
|
||
SubtitleTrackIndex int // -1 = no subtitles
|
||
SubtitleBurnIn bool
|
||
// SubtitleCodec is the probed codec of the burn-in track (e.g. "subrip",
|
||
// "hdmv_pgs_subtitle"). Bitmap codecs (PGS/DVD/DVB) select the overlay
|
||
// filter_complex pipeline; text codecs use the libass subtitles filter.
|
||
// Empty preserves the legacy text path for callers minted before the field.
|
||
SubtitleCodec string
|
||
AudioTrackIndex int // -1 = default (first track), >= 0 = specific track
|
||
TargetBitrateKbps int // max video bitrate in kbps; 0 = CRF-only (no cap)
|
||
TotalDuration float64 // total media duration in seconds (for VOD manifest)
|
||
FastStart bool // use superfast preset for faster first-segment production
|
||
NodeType string
|
||
ExecutionMode string
|
||
FFmpegLogSink FFmpegLogSink
|
||
}
|
||
|
||
// TranscodeSession manages a running ffmpeg HLS transcode process.
|
||
type TranscodeSession struct {
|
||
cmd *exec.Cmd
|
||
cancel context.CancelFunc
|
||
opts TranscodeOpts
|
||
outputDir string
|
||
running bool
|
||
restarting bool
|
||
waitErr error
|
||
stderr *boundedTailBuffer
|
||
mu sync.Mutex
|
||
done chan struct{} // closed when the monitor goroutine finishes
|
||
stdinPipe io.WriteCloser
|
||
lastRequestedSegment int
|
||
throttler *TranscodeThrottler
|
||
stderrLinesLogged int
|
||
stderrBytesLogged int
|
||
stderrDroppedLines int
|
||
stderrCapLogged bool
|
||
restartCount int
|
||
stderrLineIndex int
|
||
stderrWriter *ffmpegStderrWriter
|
||
restartHook func(context.Context)
|
||
}
|
||
|
||
// SetRestartHook registers a callback fired after every successful Restart.
|
||
// The owning handler uses it to re-arm the transcode throttler and the exit
|
||
// monitor; firing it from Restart itself keeps every restart caller (web
|
||
// segment recovery, audio switch, jellycompat seek) consistent.
|
||
func (s *TranscodeSession) SetRestartHook(fn func(context.Context)) {
|
||
s.mu.Lock()
|
||
s.restartHook = fn
|
||
s.mu.Unlock()
|
||
}
|
||
|
||
// SegmentProgress describes the media ffmpeg has actually produced on disk.
|
||
type SegmentProgress struct {
|
||
ProducedHead int
|
||
ProducedCount int
|
||
LastProducedAt time.Time
|
||
ManifestModTime time.Time
|
||
HasManifest bool
|
||
Running bool
|
||
Restarting bool
|
||
StartSegmentNumber int
|
||
SegmentDuration int
|
||
LastRequestedSegment int
|
||
}
|
||
|
||
// SegmentRecoveryDecision tells the segment handler whether to briefly wait
|
||
// for ffmpeg or seek-restart immediately.
|
||
type SegmentRecoveryDecision struct {
|
||
Wait bool
|
||
WaitTimeout time.Duration
|
||
RestartOnTimeout bool
|
||
Reason string
|
||
Progress SegmentProgress
|
||
}
|
||
|
||
// defaultSegmentDuration is the segment length when not specified. Short
|
||
// segments (2s) allow the player to start quickly while still maintaining
|
||
// efficient HTTP delivery. This matches the approach used by Plex.
|
||
const defaultSegmentDuration = 2
|
||
|
||
// DefaultSegmentDuration is the exported segment length used when a transcode
|
||
// request does not specify one. Callers minting a reconstruct recipe must embed
|
||
// a concrete (>0) value so the token passes the node's completeness gate and the
|
||
// embedded length matches what the node actually produces.
|
||
const DefaultSegmentDuration = defaultSegmentDuration
|
||
|
||
const maxPersistedFFmpegLines = 2000
|
||
const maxPersistedFFmpegBytes = 256 * 1024
|
||
const maxPersistedFFmpegChars = 2000
|
||
|
||
const (
|
||
maxSequentialMissingSegments = 2
|
||
activeSegmentWait = 12 * time.Second
|
||
segmentWaitGrace = 1500 * time.Millisecond
|
||
maxSegmentWait = 6 * time.Second
|
||
minSegmentWait = 3 * time.Second
|
||
minStaleProducedWindow = 5 * time.Second
|
||
)
|
||
|
||
// StartTranscode launches an ffmpeg process that produces HLS segments.
|
||
func StartTranscode(ctx context.Context, opts TranscodeOpts) (*TranscodeSession, error) {
|
||
if opts.SegmentDuration <= 0 {
|
||
opts.SegmentDuration = defaultSegmentDuration
|
||
}
|
||
opts.HWAccel = resolveEffectiveTranscodeHWAccel(opts)
|
||
|
||
// Ensure output directory exists.
|
||
if err := os.MkdirAll(opts.OutputDir, 0o755); err != nil {
|
||
return nil, fmt.Errorf("create output dir: %w", err)
|
||
}
|
||
|
||
ctx, cancel := context.WithCancel(ctx)
|
||
s := &TranscodeSession{
|
||
cancel: cancel,
|
||
opts: opts,
|
||
outputDir: opts.OutputDir,
|
||
running: true,
|
||
done: make(chan struct{}),
|
||
stderr: newBoundedTailBuffer(stderrTailMaxBytes),
|
||
lastRequestedSegment: opts.StartSegmentNumber,
|
||
}
|
||
|
||
args := buildFFmpegArgs(opts)
|
||
bin := opts.FFmpegPath
|
||
if bin == "" {
|
||
bin = ffmpegBinary()
|
||
}
|
||
|
||
log.Printf("playback: ffmpeg cmd: %s %s", bin, strings.Join(args, " "))
|
||
s.logFFmpegEvent(ctx, "ffmpeg process starting", "")
|
||
|
||
cmd := exec.CommandContext(ctx, bin, args...)
|
||
stdinPipe, err := cmd.StdinPipe()
|
||
if err != nil {
|
||
cancel()
|
||
return nil, fmt.Errorf("create stdin pipe: %w", err)
|
||
}
|
||
cmd.Dir = opts.OutputDir
|
||
cmd.Stderr = s.newStderrWriter(ctx)
|
||
cmd.WaitDelay = 3 * time.Second
|
||
|
||
if err := cmd.Start(); err != nil {
|
||
cancel()
|
||
s.logFFmpegEvent(ctx, "ffmpeg process exit error", err.Error())
|
||
return nil, fmt.Errorf("start ffmpeg: %w", err)
|
||
}
|
||
s.cmd = cmd
|
||
s.stdinPipe = stdinPipe
|
||
s.logFFmpegEvent(ctx, "ffmpeg process started", "")
|
||
|
||
// Monitor ffmpeg in background.
|
||
go func() {
|
||
waitErr := cmd.Wait()
|
||
s.flushStderr(ctx)
|
||
s.mu.Lock()
|
||
s.running = false
|
||
s.waitErr = waitErr
|
||
s.mu.Unlock()
|
||
s.logWaitResult(ctx, waitErr)
|
||
close(s.done)
|
||
}()
|
||
|
||
return s, nil
|
||
}
|
||
|
||
// IsMPEG2VideoCodec reports whether a probed video codec name identifies
|
||
// MPEG-2 video. It accepts common FFmpeg aliases because codec strings can
|
||
// come from scan metadata, direct probes, or client capability lists.
|
||
func IsMPEG2VideoCodec(codec string) bool {
|
||
normalized := strings.NewReplacer(
|
||
" ", "",
|
||
"-", "",
|
||
"_", "",
|
||
".", "",
|
||
).Replace(strings.ToLower(strings.TrimSpace(codec)))
|
||
switch normalized {
|
||
case "mpeg2video", "mpeg2", "mp2v":
|
||
return true
|
||
default:
|
||
return false
|
||
}
|
||
}
|
||
|
||
// IsMPEG4Part2VideoCodec reports whether a codec name identifies MPEG-4 Part 2
|
||
// video, commonly found in older XviD/DivX AVI files.
|
||
func IsMPEG4Part2VideoCodec(codec string) bool {
|
||
normalized := strings.NewReplacer(
|
||
" ", "",
|
||
"-", "",
|
||
"_", "",
|
||
".", "",
|
||
).Replace(strings.ToLower(strings.TrimSpace(codec)))
|
||
switch normalized {
|
||
case "mpeg4", "mp4v", "xvid", "divx", "dx50":
|
||
return true
|
||
default:
|
||
return false
|
||
}
|
||
}
|
||
|
||
// buildFFmpegArgs constructs the full ffmpeg argument list from TranscodeOpts.
|
||
func buildFFmpegArgs(opts TranscodeOpts) []string {
|
||
// Resolve "auto" into a concrete accel method once so all downstream
|
||
// helpers (appendHWAccelArgs, appendVideoArgs, etc.) see the real value.
|
||
opts.HWAccel = resolveEffectiveTranscodeHWAccel(opts)
|
||
|
||
isVideoCopy := opts.TargetCodecVideo == "copy"
|
||
isAudioCopy := opts.TargetCodecAudio == "copy"
|
||
|
||
args := []string{
|
||
"-hide_banner",
|
||
"-loglevel", "error",
|
||
}
|
||
|
||
// Hardware acceleration — skip when copying video (no encoding needed).
|
||
if !isVideoCopy {
|
||
args = appendHWAccelArgs(args, opts)
|
||
}
|
||
|
||
// Limit input probing to speed up startup, especially on network storage.
|
||
// -fflags +genpts generates PTS for files with missing timestamps;
|
||
// +fastseek enables faster input seeking (matches Jellyfin).
|
||
args = append(args,
|
||
"-fflags", "+genpts+fastseek",
|
||
"-analyzeduration", "3000000", // 3 seconds (default 5s)
|
||
"-probesize", "5000000", // 5 MB (default 5MB, explicit for clarity)
|
||
)
|
||
|
||
// Seek before input for fast seeking.
|
||
if opts.SeekSeconds > 0 {
|
||
args = append(args, "-ss", fmt.Sprintf("%.3f", opts.SeekSeconds))
|
||
// When video is copied but audio is transcoded, accurate_seek causes
|
||
// A/V desync: video must start at a keyframe but audio is trimmed to
|
||
// the exact seek point. Disabling it keeps both streams aligned.
|
||
if isVideoCopy && !isAudioCopy {
|
||
args = append(args, "-noaccurate_seek")
|
||
}
|
||
}
|
||
|
||
// Input file.
|
||
args = append(args, "-i", opts.InputPath)
|
||
args = append(args, "-map_metadata", "-1")
|
||
args = append(args, "-map_chapters", "-1")
|
||
args = appendStreamSelectionArgs(args, opts)
|
||
args = appendTimestampNormalizationArgs(args, opts)
|
||
|
||
// Video codec and encoding settings.
|
||
if isVideoCopy {
|
||
args = append(args, "-c:v", "copy")
|
||
} else {
|
||
args = appendVideoArgs(args, opts)
|
||
}
|
||
|
||
// Copy-video sessions only do audio work on the filter/encode side.
|
||
// ffmpeg's default thread selection spawns one filter thread per CPU
|
||
// (observed 14 idle `af#0:1` threads for a 5.1→2.0 downmix), so pin
|
||
// audio filter + encode to a single thread.
|
||
if isVideoCopy && !isAudioCopy {
|
||
args = append(args, "-threads", "1", "-filter_threads", "1", "-filter_complex_threads", "1")
|
||
}
|
||
|
||
// Audio codec.
|
||
args = appendAudioArgs(args, opts)
|
||
|
||
// Subtitle burn-in and resolution scaling — only when encoding video.
|
||
if !isVideoCopy {
|
||
args = appendVideoFilterArgs(args, opts)
|
||
args = appendSegmentBoundaryArgs(args, opts)
|
||
}
|
||
|
||
// HLS output options.
|
||
// Codec-copy sessions usually use fMP4 segments — no transmuxing needed in
|
||
// hls.js, which avoids Safari MSE compatibility issues with certain codecs
|
||
// in TS. MPEG-2 video is the exception: Apple consumes it as compatibility
|
||
// HLS, so package it in MPEG-TS while still copying the video stream.
|
||
// Actual transcoding uses MPEG-TS segments to avoid the hls.js endOfStream()
|
||
// race with fMP4 (hls.js #6337).
|
||
var segmentPattern string
|
||
segmentType := "mpegts"
|
||
copyVideoUsesFMP4 := isVideoCopy && !IsMPEG2VideoCodec(opts.SourceVideoCodec)
|
||
if copyVideoUsesFMP4 {
|
||
segmentType = "fmp4"
|
||
segmentPattern = filepath.Join(opts.OutputDir, "seg_%05d.m4s")
|
||
} else {
|
||
segmentPattern = filepath.Join(opts.OutputDir, "seg_%05d.ts")
|
||
}
|
||
manifestPath := filepath.Join(opts.OutputDir, "stream.m3u8")
|
||
|
||
args = append(args,
|
||
"-max_muxing_queue_size", "2048",
|
||
"-max_delay", "5000000",
|
||
"-f", "hls",
|
||
"-hls_time", fmt.Sprintf("%d", opts.SegmentDuration),
|
||
"-hls_list_size", "0",
|
||
"-hls_segment_type", segmentType,
|
||
// Write segments to temp files first so the player never fetches a
|
||
// partially-written segment during a quality switch.
|
||
"-hls_flags", "independent_segments+temp_file",
|
||
"-hls_segment_filename", segmentPattern,
|
||
)
|
||
// fMP4 segments need movflags=+frag_discont so each fragment writes
|
||
// audio DTS/PTS including the initial delay into MOOF→TRAF→TFDT.
|
||
// Without this, some browsers (notably Chromium on macOS) can experience
|
||
// A/V sync issues during copy-mode HLS playback. Matches Jellyfin's
|
||
// proven fMP4 HLS pipeline.
|
||
if copyVideoUsesFMP4 {
|
||
args = append(args, "-hls_segment_options", "movflags=+frag_discont")
|
||
}
|
||
if opts.StartSegmentNumber > 0 {
|
||
args = append(args, "-start_number", fmt.Sprintf("%d", opts.StartSegmentNumber))
|
||
}
|
||
args = append(args, manifestPath)
|
||
|
||
return args
|
||
}
|
||
|
||
func resolveEffectiveTranscodeHWAccel(opts TranscodeOpts) string {
|
||
hwAccel := ResolveHWAccelWithFFmpeg(opts.HWAccel, opts.FFmpegPath)
|
||
if hwAccel == "" {
|
||
return ""
|
||
}
|
||
if strings.EqualFold(opts.TargetCodecVideo, "copy") {
|
||
return "none"
|
||
}
|
||
if IsMPEG4Part2VideoCodec(opts.SourceVideoCodec) {
|
||
return "none"
|
||
}
|
||
return hwAccel
|
||
}
|
||
|
||
// bitmapBurnInActive reports whether this transcode composites a bitmap
|
||
// subtitle track (PGS/VOBSUB/DVB) into the video via the overlay
|
||
// filter_complex pipeline. Bitmap burn-in requires a video encode: copy-video
|
||
// sessions never activate it (the API layer forces an encoding recipe before
|
||
// starting a burn-in transcode, this is a defensive backstop).
|
||
func bitmapBurnInActive(opts TranscodeOpts) bool {
|
||
return opts.SubtitleBurnIn &&
|
||
opts.SubtitleTrackIndex >= 0 &&
|
||
NeedsBurnIn(opts.SubtitleCodec) &&
|
||
!strings.EqualFold(opts.TargetCodecVideo, "copy")
|
||
}
|
||
|
||
// appendStreamSelectionArgs limits output to primary video/audio streams.
|
||
// When bitmap burn-in is active the video output comes from the overlay
|
||
// filter_complex graph's labeled pad instead of the raw input stream —
|
||
// mapping both would emit two video streams into the HLS mux.
|
||
func appendStreamSelectionArgs(args []string, opts TranscodeOpts) []string {
|
||
audioTrackIndex := opts.AudioTrackIndex
|
||
if bitmapBurnInActive(opts) {
|
||
args = append(args, "-map", "[vout]")
|
||
} else {
|
||
args = append(args, "-map", "0:v:0")
|
||
}
|
||
if audioTrackIndex >= 0 {
|
||
args = append(args, "-map", fmt.Sprintf("0:a:%d?", audioTrackIndex))
|
||
} else {
|
||
args = append(args, "-map", "0:a:0?")
|
||
}
|
||
args = append(args, "-sn")
|
||
args = append(args, "-dn")
|
||
return args
|
||
}
|
||
|
||
// appendTimestampNormalizationArgs selects timestamp handling based on the
|
||
// playback mode. Copy-video full-file starts use zero-based timestamps so
|
||
// fMP4 fragments always have sane local durations. Copy-video resumes
|
||
// preserve source timestamps so each fragment's TFDT matches its playlist
|
||
// position (segment K sits at playlist-time K*segDur); zero-basing here
|
||
// makes seg_K carry TFDT=0, and strict players (Jellyfin Android TV /
|
||
// ExoPlayer) read EXT-X-START, jump to seg_K expecting media at K*segDur,
|
||
// see TFDT=0, treat the gap as a discontinuity, reload init.mp4, and
|
||
// eventually abort — the symptom that crashes ATV on a second resume.
|
||
// Encoded transcodes keep the source-timestamp policy unconditionally.
|
||
func appendTimestampNormalizationArgs(args []string, opts TranscodeOpts) []string {
|
||
if strings.EqualFold(opts.TargetCodecVideo, "copy") {
|
||
if opts.SeekSeconds > 0 {
|
||
return append(args,
|
||
"-copyts",
|
||
"-avoid_negative_ts", "disabled",
|
||
)
|
||
}
|
||
return append(args,
|
||
"-avoid_negative_ts", "make_zero",
|
||
)
|
||
}
|
||
return append(args,
|
||
"-copyts",
|
||
"-avoid_negative_ts", "disabled",
|
||
)
|
||
}
|
||
|
||
// appendSegmentBoundaryArgs forces keyframes on segment boundaries so each HLS
|
||
// fragment starts cleanly and can be appended independently by the player.
|
||
//
|
||
// With -copyts, the output timestamp t starts at the seek position rather than
|
||
// 0. Subtracting SeekSeconds prevents a "catch-up storm" where n_forced races
|
||
// from 0 to seek_position/segment_duration, making every frame an I-frame and
|
||
// grinding encoding to a halt for large seeks.
|
||
func appendSegmentBoundaryArgs(args []string, opts TranscodeOpts) []string {
|
||
args = append(args, "-sc_threshold", "0")
|
||
if opts.SeekSeconds > 0 {
|
||
args = append(args, "-force_key_frames",
|
||
fmt.Sprintf("expr:gte(t-%.3f,n_forced*%d)", opts.SeekSeconds, opts.SegmentDuration))
|
||
} else {
|
||
args = append(args, "-force_key_frames",
|
||
fmt.Sprintf("expr:gte(t,n_forced*%d)", opts.SegmentDuration))
|
||
}
|
||
|
||
// Hardware encoders (QSV, VAAPI, NVENC) may not reliably honor
|
||
// force_key_frames expressions. Set explicit GOP size so segment
|
||
// boundaries always start with an IDR frame. We assume 30 fps as a
|
||
// safe ceiling — the GOP will be at most segmentDuration * 30 frames.
|
||
// Matches Jellyfin's approach for hardware encoders.
|
||
if opts.HWAccel == "qsv" || opts.HWAccel == "vaapi" || opts.HWAccel == "nvenc" {
|
||
gopSize := fmt.Sprintf("%d", opts.SegmentDuration*30)
|
||
args = append(args, "-g", gopSize, "-keyint_min", gopSize)
|
||
}
|
||
|
||
return args
|
||
}
|
||
|
||
// appendHWAccelArgs adds hardware acceleration flags based on the HWAccel setting.
|
||
// The caller must resolve "auto" via ResolveHWAccel before calling this.
|
||
func appendHWAccelArgs(args []string, opts TranscodeOpts) []string {
|
||
switch opts.HWAccel {
|
||
case "qsv":
|
||
hwDevice := PickRenderDevice(opts.HWDevice)
|
||
if hwDevice == "" {
|
||
slog.Warn("no GPU render device found, QSV transcode may fail")
|
||
hwDevice = "/dev/dri/renderD128" // last-resort fallback
|
||
}
|
||
// VAAPI→QSV hardware pipeline: derive QSV from VAAPI device.
|
||
args = append(args,
|
||
"-init_hw_device", fmt.Sprintf("vaapi=va:%s,driver=iHD,kernel_driver=i915,vendor_id=0x8086", hwDevice),
|
||
"-init_hw_device", "qsv=qs@va",
|
||
"-filter_hw_device", "va",
|
||
"-hwaccel", "vaapi",
|
||
"-hwaccel_output_format", "vaapi",
|
||
"-noautorotate",
|
||
)
|
||
case "vaapi":
|
||
vaapiDevice := PickRenderDevice(opts.HWDevice)
|
||
if vaapiDevice == "" {
|
||
vaapiDevice = "/dev/dri/renderD128" // last-resort fallback
|
||
}
|
||
args = append(args,
|
||
"-init_hw_device", fmt.Sprintf("vaapi=hw:%s", vaapiDevice),
|
||
"-filter_hw_device", "hw",
|
||
"-hwaccel", "vaapi",
|
||
"-hwaccel_output_format", "vaapi",
|
||
)
|
||
case "nvenc":
|
||
args = append(args,
|
||
"-hwaccel", "cuda",
|
||
"-hwaccel_output_format", "cuda",
|
||
"-noautorotate",
|
||
)
|
||
if hwDevice := strings.TrimSpace(opts.HWDevice); hwDevice != "" {
|
||
args = append(args, "-hwaccel_device", hwDevice)
|
||
}
|
||
}
|
||
return args
|
||
}
|
||
|
||
// videoPreset returns an encoder-compatible preset. CPU encoders use a faster
|
||
// fast-start preset for initial playback, while QSV stays on the fastest
|
||
// preset family it supports.
|
||
func videoPreset(opts TranscodeOpts, hwAccel string) string {
|
||
if hwAccel == "qsv" {
|
||
return "veryfast"
|
||
}
|
||
if opts.FastStart {
|
||
return "superfast"
|
||
}
|
||
return "veryfast"
|
||
}
|
||
|
||
// appendVideoArgs adds video codec arguments.
|
||
func appendVideoArgs(args []string, opts TranscodeOpts) []string {
|
||
codec := opts.TargetCodecVideo
|
||
if codec == "" {
|
||
codec = "h264"
|
||
}
|
||
|
||
if codec == "copy" {
|
||
return append(args, "-c:v", "copy")
|
||
}
|
||
|
||
preset := videoPreset(opts, opts.HWAccel)
|
||
hasBitrateCap := opts.TargetBitrateKbps > 0
|
||
|
||
switch {
|
||
case opts.HWAccel == "qsv" && codec == "h264":
|
||
if hasBitrateCap {
|
||
// VBR mode with bitrate cap instead of global_quality.
|
||
args = append(args, "-c:v", "h264_qsv", "-preset", preset,
|
||
"-b:v", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-maxrate", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-bufsize", fmt.Sprintf("%dk", opts.TargetBitrateKbps*2))
|
||
} else {
|
||
args = append(args, "-c:v", "h264_qsv", "-preset", preset, "-global_quality", "23")
|
||
}
|
||
case opts.HWAccel == "qsv" && codec == "hevc":
|
||
if hasBitrateCap {
|
||
args = append(args, "-c:v", "hevc_qsv", "-preset", preset,
|
||
"-b:v", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-maxrate", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-bufsize", fmt.Sprintf("%dk", opts.TargetBitrateKbps*2))
|
||
} else {
|
||
args = append(args, "-c:v", "hevc_qsv", "-preset", preset, "-global_quality", "28")
|
||
}
|
||
case opts.HWAccel == "vaapi" && codec == "h264":
|
||
args = append(args, "-c:v", "h264_vaapi", "-qp", "23")
|
||
if hasBitrateCap {
|
||
args = append(args,
|
||
"-maxrate", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-bufsize", fmt.Sprintf("%dk", opts.TargetBitrateKbps*2))
|
||
}
|
||
case opts.HWAccel == "vaapi" && codec == "hevc":
|
||
args = append(args, "-c:v", "hevc_vaapi", "-qp", "28")
|
||
if hasBitrateCap {
|
||
args = append(args,
|
||
"-maxrate", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-bufsize", fmt.Sprintf("%dk", opts.TargetBitrateKbps*2))
|
||
}
|
||
case opts.HWAccel == "nvenc" && codec == "h264":
|
||
args = append(args, "-c:v", "h264_nvenc", "-rc:v", "vbr")
|
||
if hasBitrateCap {
|
||
args = append(args,
|
||
"-b:v", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-maxrate", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-bufsize", fmt.Sprintf("%dk", opts.TargetBitrateKbps*2))
|
||
} else {
|
||
args = append(args, "-cq:v", "23", "-b:v", "0")
|
||
}
|
||
case opts.HWAccel == "nvenc" && codec == "hevc":
|
||
args = append(args, "-c:v", "hevc_nvenc", "-rc:v", "vbr")
|
||
if hasBitrateCap {
|
||
args = append(args,
|
||
"-b:v", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-maxrate", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-bufsize", fmt.Sprintf("%dk", opts.TargetBitrateKbps*2))
|
||
} else {
|
||
args = append(args, "-cq:v", "28", "-b:v", "0")
|
||
}
|
||
default:
|
||
// CPU fallback — match Jellyfin's proven browser-compatible settings.
|
||
// Force yuv420p to ensure 8-bit output (10-bit sources produce High 10
|
||
// Profile which browsers cannot decode via MSE).
|
||
if codec == "hevc" {
|
||
args = append(args, "-c:v", "libx265", "-preset", preset, "-crf", "28", "-pix_fmt", "yuv420p")
|
||
} else {
|
||
args = append(args, "-c:v", "libx264", "-preset", preset, "-crf", "23",
|
||
"-pix_fmt", "yuv420p", "-profile:v", "high", "-level", "4.1")
|
||
}
|
||
if hasBitrateCap {
|
||
args = append(args,
|
||
"-maxrate", fmt.Sprintf("%dk", opts.TargetBitrateKbps),
|
||
"-bufsize", fmt.Sprintf("%dk", opts.TargetBitrateKbps*2))
|
||
}
|
||
}
|
||
|
||
return args
|
||
}
|
||
|
||
// appendVideoFilterArgs appends the -vf selection for an encoding (non-copy)
|
||
// video stream: the subtitle burn-in chain (which includes scaling and hw
|
||
// download/upload for QSV/VAAPI) or the hwaccel-appropriate standalone scale
|
||
// filter. The ONE home of this decision — the HLS builder and the single-file
|
||
// prepare builder must always produce identical filter chains (a fix landing
|
||
// in only one of them silently ships wrong cached artifacts).
|
||
func appendVideoFilterArgs(args []string, opts TranscodeOpts) []string {
|
||
switch {
|
||
case bitmapBurnInActive(opts):
|
||
return appendBitmapSubtitleBurnInArgs(args, opts)
|
||
case opts.SubtitleBurnIn && opts.SubtitleTrackIndex >= 0:
|
||
return appendSubtitleBurnInArgs(args, opts)
|
||
case opts.HWAccel == "qsv":
|
||
return append(args, "-vf", qsvScaleFilter(opts.TargetResolution))
|
||
case opts.HWAccel == "vaapi":
|
||
return append(args, "-vf", vaapiScaleFilter(opts.TargetResolution))
|
||
case opts.HWAccel == "nvenc":
|
||
return append(args, "-vf", nvencScaleFilter(opts.TargetResolution))
|
||
case opts.TargetResolution != "":
|
||
if scale := resolutionToScale(opts.TargetResolution); scale != "" {
|
||
return append(args, "-vf", scale)
|
||
}
|
||
}
|
||
return args
|
||
}
|
||
|
||
// appendAudioArgs adds audio codec arguments. Supports "copy" for passthrough,
|
||
// plus opus / aac / eac3 / ac3 as re-encode targets. EAC3 and AC3 are useful
|
||
// when we must transcode video but want to preserve surround channels for an
|
||
// HDMI receiver — both are legal in HLS fMP4 (not MPEG-TS; ensure the HLS
|
||
// packager is fMP4 when emitting these).
|
||
func appendAudioArgs(args []string, opts TranscodeOpts) []string {
|
||
codec := opts.TargetCodecAudio
|
||
if codec == "" {
|
||
codec = "aac"
|
||
}
|
||
|
||
switch codec {
|
||
case "copy":
|
||
args = append(args, "-c:a", "copy")
|
||
case "opus":
|
||
args = append(args, "-c:a", "libopus", "-b:a", "192k", "-ac", "2")
|
||
case "eac3":
|
||
// Typical Dolby Digital Plus 5.1 bitrate; let the source dictate channel
|
||
// count so we preserve surround when possible.
|
||
args = append(args, "-c:a", "eac3", "-b:a", "384k")
|
||
case "ac3":
|
||
// Legacy Dolby Digital; universal AVR support.
|
||
args = append(args, "-c:a", "ac3", "-b:a", "448k")
|
||
default:
|
||
args = append(args, "-c:a", "aac", "-b:a", "192k", "-ac", "2")
|
||
}
|
||
|
||
return args
|
||
}
|
||
|
||
// appendBitmapSubtitleBurnInArgs adds burn-in arguments for BITMAP subtitle
|
||
// codecs (PGS/VOBSUB/DVB). libass's subtitles= filter cannot render bitmap
|
||
// tracks, so the decoded subtitle stream is composited onto the video with
|
||
// overlay in a -filter_complex graph (the "Plex route"). The graph's output
|
||
// pad [vout] replaces the raw video stream in stream mapping (see
|
||
// appendStreamSelectionArgs), so -vf must never be emitted alongside this.
|
||
//
|
||
// Overlay runs at the source's native resolution FIRST and any target scaling
|
||
// happens after, so bitmap subtitle geometry is never distorted by a
|
||
// pre-scaling mismatch between the video and subtitle planes.
|
||
// eof_action=pass keeps the video flowing untouched once the subtitle stream
|
||
// ends instead of freezing the last overlay frame on screen.
|
||
//
|
||
// Hardware pipelines mirror appendSubtitleBurnInArgs: frames are downloaded
|
||
// to CPU memory for the overlay, then re-uploaded for the hardware encoder.
|
||
func appendBitmapSubtitleBurnInArgs(args []string, opts TranscodeOpts) []string {
|
||
// [0:s:N] indexes subtitle streams only, matching the si=N semantics of
|
||
// the text path — SubtitleTrackIndex is the embedded subtitle ordinal.
|
||
cpuFilters := fmt.Sprintf("[0:s:%d]overlay=eof_action=pass", opts.SubtitleTrackIndex)
|
||
if scale := resolutionToScale(opts.TargetResolution); scale != "" {
|
||
cpuFilters += "," + scale
|
||
}
|
||
|
||
var graph string
|
||
switch opts.HWAccel {
|
||
case "qsv":
|
||
// VAAPI→QSV pipeline: download decoded frames to CPU, overlay, convert
|
||
// to nv12, upload back to VAAPI, then map to QSV for the encoder.
|
||
graph = "[0:v:0]hwdownload,format=yuv420p[vmain];[vmain]" + cpuFilters +
|
||
",format=nv12,hwupload,hwmap=derive_device=qsv,format=qsv[vout]"
|
||
case "vaapi":
|
||
graph = "[0:v:0]hwdownload,format=yuv420p[vmain];[vmain]" + cpuFilters +
|
||
",format=nv12,hwupload[vout]"
|
||
case "nvenc":
|
||
graph = "[0:v:0]hwdownload,format=yuv420p[vmain];[vmain]" + cpuFilters +
|
||
",format=nv12,hwupload_cuda[vout]"
|
||
default:
|
||
// CPU encoding: overlay directly on decoded frames.
|
||
graph = "[0:v:0]" + cpuFilters + "[vout]"
|
||
}
|
||
|
||
return append(args, "-filter_complex", graph)
|
||
}
|
||
|
||
// appendSubtitleBurnInArgs adds subtitle burn-in filter arguments for TEXT
|
||
// subtitle codecs (SRT/ASS/…) via the libass-based subtitles= filter; bitmap
|
||
// codecs take the overlay path in appendBitmapSubtitleBurnInArgs.
|
||
// For CPU encoding, the filter chain is: [scale,]subtitles.
|
||
// For QSV/VAAPI, frames must be downloaded from hardware, processed on CPU,
|
||
// then re-uploaded: hwdownload → format=yuv420p → [scale,] subtitles → hwupload → hwmap.
|
||
func appendSubtitleBurnInArgs(args []string, opts TranscodeOpts) []string {
|
||
scale := resolutionToScale(opts.TargetResolution)
|
||
subFilter := fmt.Sprintf("subtitles='%s':si=%d",
|
||
escapeFilterPath(opts.InputPath), opts.SubtitleTrackIndex)
|
||
|
||
// Build the CPU filter portion: scale (if any) then subtitle overlay.
|
||
// Scale must come before subtitles so text is rendered at target resolution.
|
||
var cpuFilters string
|
||
if scale != "" {
|
||
cpuFilters = scale + "," + subFilter
|
||
} else {
|
||
cpuFilters = subFilter
|
||
}
|
||
|
||
switch opts.HWAccel {
|
||
case "qsv":
|
||
// VAAPI→QSV pipeline: download from VAAPI surface to CPU, apply subtitle
|
||
// and scale filters, convert to nv12 (required by hwupload for VAAPI
|
||
// surfaces), upload back to VAAPI, then map to QSV for the encoder.
|
||
vf := "hwdownload,format=yuv420p," + cpuFilters + ",format=nv12,hwupload,hwmap=derive_device=qsv,format=qsv"
|
||
args = append(args, "-vf", vf)
|
||
case "vaapi":
|
||
// VAAPI-only: download, apply CPU filters, convert to nv12, upload back.
|
||
vf := "hwdownload,format=yuv420p," + cpuFilters + ",format=nv12,hwupload"
|
||
args = append(args, "-vf", vf)
|
||
case "nvenc":
|
||
// NVENC/CUDA: download to CPU for subtitle rendering, then upload back.
|
||
vf := "hwdownload,format=yuv420p," + cpuFilters + ",format=nv12,hwupload_cuda"
|
||
args = append(args, "-vf", vf)
|
||
default:
|
||
// CPU encoding: filters run directly on decoded frames.
|
||
args = append(args, "-vf", cpuFilters)
|
||
}
|
||
|
||
return args
|
||
}
|
||
|
||
// resolutionToScale returns an ffmpeg scale filter string for the target resolution.
|
||
func resolutionToScale(res string) string {
|
||
switch res {
|
||
case "2160p":
|
||
return "scale=-2:2160"
|
||
case "1080p":
|
||
return "scale=-2:1080"
|
||
case "720p":
|
||
return "scale=-2:720"
|
||
case "480p":
|
||
return "scale=-2:480"
|
||
case "420p":
|
||
return "scale=-2:420"
|
||
case "328p":
|
||
return "scale=-2:328"
|
||
default:
|
||
return ""
|
||
}
|
||
}
|
||
|
||
// qsvScaleFilter returns the VAAPI→QSV filter chain with optional resolution scaling.
|
||
func qsvScaleFilter(res string) string {
|
||
switch res {
|
||
case "2160p":
|
||
return "scale_vaapi=w=-2:h=2160:format=nv12,hwmap=derive_device=qsv,format=qsv"
|
||
case "1080p":
|
||
return "scale_vaapi=w=-2:h=1080:format=nv12,hwmap=derive_device=qsv,format=qsv"
|
||
case "720p":
|
||
return "scale_vaapi=w=-2:h=720:format=nv12,hwmap=derive_device=qsv,format=qsv"
|
||
case "480p":
|
||
return "scale_vaapi=w=-2:h=480:format=nv12,hwmap=derive_device=qsv,format=qsv"
|
||
case "420p":
|
||
return "scale_vaapi=w=-2:h=420:format=nv12,hwmap=derive_device=qsv,format=qsv"
|
||
case "328p":
|
||
return "scale_vaapi=w=-2:h=328:format=nv12,hwmap=derive_device=qsv,format=qsv"
|
||
default:
|
||
return "scale_vaapi=format=nv12,hwmap=derive_device=qsv,format=qsv"
|
||
}
|
||
}
|
||
|
||
// vaapiScaleFilter keeps VAAPI frames in hardware and converts them to a
|
||
// browser-compatible encoder format. Using the CPU scale filter on VAAPI frames
|
||
// causes FFmpeg auto_scale format-negotiation failures.
|
||
func vaapiScaleFilter(res string) string {
|
||
switch res {
|
||
case "2160p":
|
||
return "scale_vaapi=w=-2:h=2160:format=nv12"
|
||
case "1080p":
|
||
return "scale_vaapi=w=-2:h=1080:format=nv12"
|
||
case "720p":
|
||
return "scale_vaapi=w=-2:h=720:format=nv12"
|
||
case "480p":
|
||
return "scale_vaapi=w=-2:h=480:format=nv12"
|
||
case "420p":
|
||
return "scale_vaapi=w=-2:h=420:format=nv12"
|
||
case "328p":
|
||
return "scale_vaapi=w=-2:h=328:format=nv12"
|
||
default:
|
||
return "scale_vaapi=format=nv12"
|
||
}
|
||
}
|
||
|
||
func nvencScaleFilter(res string) string {
|
||
switch res {
|
||
case "2160p":
|
||
return "scale_cuda=w=-2:h=2160:format=nv12"
|
||
case "1080p":
|
||
return "scale_cuda=w=-2:h=1080:format=nv12"
|
||
case "720p":
|
||
return "scale_cuda=w=-2:h=720:format=nv12"
|
||
case "480p":
|
||
return "scale_cuda=w=-2:h=480:format=nv12"
|
||
case "420p":
|
||
return "scale_cuda=w=-2:h=420:format=nv12"
|
||
case "328p":
|
||
return "scale_cuda=w=-2:h=328:format=nv12"
|
||
default:
|
||
return "scale_cuda=format=nv12"
|
||
}
|
||
}
|
||
|
||
// filterPathReplacer escapes special characters in file paths for ffmpeg filter syntax.
|
||
var filterPathReplacer = strings.NewReplacer(
|
||
"'", "'\\''",
|
||
"[", "\\[",
|
||
"]", "\\]",
|
||
";", "\\;",
|
||
",", "\\,",
|
||
)
|
||
|
||
// escapeFilterPath escapes special characters in file paths for ffmpeg filter syntax.
|
||
func escapeFilterPath(path string) string {
|
||
return filterPathReplacer.Replace(path)
|
||
}
|
||
|
||
// minManifestSegments is the standard startup lead for actively encoded HLS.
|
||
// True video transcodes benefit from a larger cushion so playback does not
|
||
// outrun the encoder immediately after the first frame appears.
|
||
const minManifestSegments = 3
|
||
|
||
// minCopyManifestSegments is the startup lead for codec-copy sessions.
|
||
// Copying video while transcoding only audio can produce startup files far
|
||
// faster than real-time encoding, so waiting for 3 full segments adds
|
||
// unnecessary latency at playback start.
|
||
const minCopyManifestSegments = 2
|
||
|
||
func startupSegmentRequirement(opts TranscodeOpts) int {
|
||
if strings.EqualFold(opts.TargetCodecVideo, "copy") {
|
||
return minCopyManifestSegments
|
||
}
|
||
return minManifestSegments
|
||
}
|
||
|
||
// GetManifest returns the HLS m3u8 manifest content.
|
||
// It returns ErrManifestNotReady if the manifest does not yet contain enough
|
||
// segments for reliable HLS playback (see minManifestSegments).
|
||
func (s *TranscodeSession) GetManifest() ([]byte, error) {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
|
||
manifestPath := filepath.Join(s.outputDir, "stream.m3u8")
|
||
data, err := os.ReadFile(manifestPath)
|
||
if err != nil {
|
||
if os.IsNotExist(err) {
|
||
if !s.running {
|
||
if s.restarting {
|
||
return nil, ErrManifestNotReady
|
||
}
|
||
if s.waitErr != nil {
|
||
stderr := truncateStderr(s.stderr.String())
|
||
if stderr != "" {
|
||
return nil, fmt.Errorf("%w: %v (stderr: %s)", ErrTranscodeFailed, s.waitErr, stderr)
|
||
}
|
||
return nil, fmt.Errorf("%w: %v", ErrTranscodeFailed, s.waitErr)
|
||
}
|
||
return nil, ErrTranscodeFailed
|
||
}
|
||
return nil, ErrManifestNotReady
|
||
}
|
||
return nil, fmt.Errorf("read manifest: %w", err)
|
||
}
|
||
|
||
requiredSegments := startupSegmentRequirement(s.opts)
|
||
|
||
// Wait until enough startup media exists before serving. Counting #EXTINF
|
||
// lines alone is not enough for FFmpeg's live-written manifest because the
|
||
// playlist can reference copy-mode segments before the files are fully
|
||
// flushed to disk, especially on resumed sessions with a non-zero media
|
||
// sequence. Requiring the referenced startup files prevents the browser from
|
||
// stalling on its very first segment fetch.
|
||
if s.running && !startupFilesReady(data, s.outputDir, requiredSegments) {
|
||
return nil, ErrManifestNotReady
|
||
}
|
||
if strings.EqualFold(s.opts.TargetCodecVideo, "copy") {
|
||
if err := validateCopyPlaybackManifest(data); err != nil {
|
||
return nil, fmt.Errorf("invalid copy playback manifest: %w", err)
|
||
}
|
||
}
|
||
|
||
return data, nil
|
||
}
|
||
|
||
// WaitForManifest polls until the manifest is ready for playback or the timeout
|
||
// expires. It keeps the initial request open long enough for FFmpeg to write
|
||
// the first safe playback window instead of forcing the client to race a 503.
|
||
func (s *TranscodeSession) WaitForManifest(timeout time.Duration) ([]byte, error) {
|
||
deadline := time.After(timeout)
|
||
for {
|
||
manifest, err := s.GetManifest()
|
||
if err == nil {
|
||
return manifest, nil
|
||
}
|
||
if err != nil && err != ErrManifestNotReady {
|
||
return nil, err
|
||
}
|
||
|
||
select {
|
||
case <-deadline:
|
||
return nil, s.manifestTimeoutError(timeout)
|
||
case <-time.After(100 * time.Millisecond):
|
||
}
|
||
}
|
||
}
|
||
|
||
// BuildPlaybackManifest returns the manifest we should expose to clients.
|
||
//
|
||
// Copy-video sessions always expose FFmpeg's real manifest so the playlist
|
||
// timing matches the variable-length fragments FFmpeg actually writes and the
|
||
// seekable window reflects what FFmpeg has produced so far. Encoded transcodes
|
||
// still use the synthetic full VOD manifest when duration is known because
|
||
// forced keyframes make that timeline stable and seek-anywhere friendly.
|
||
func (s *TranscodeSession) BuildPlaybackManifest(segPrefix, rawQuery string) ([]byte, error) {
|
||
opts := s.Opts()
|
||
if strings.EqualFold(opts.TargetCodecVideo, "copy") || opts.TotalDuration <= 0 {
|
||
// Copy-video or unknown-duration sessions must use FFmpeg's real manifest.
|
||
manifest, err := s.WaitForManifest(30 * time.Second)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return RewriteManifestPaths(manifest, segPrefix, rawQuery)
|
||
}
|
||
|
||
return s.GenerateFullManifest(segPrefix, rawQuery), nil
|
||
}
|
||
|
||
func firstNonEmptyManifestLine(manifest []byte) []byte {
|
||
for line := range bytes.SplitSeq(manifest, []byte("\n")) {
|
||
trimmed := bytes.TrimSpace(line)
|
||
if len(trimmed) > 0 {
|
||
return trimmed
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func validateManifestHeader(manifest []byte) error {
|
||
if len(bytes.TrimSpace(manifest)) == 0 {
|
||
return fmt.Errorf("manifest is empty")
|
||
}
|
||
if line := firstNonEmptyManifestLine(manifest); !bytes.Equal(line, []byte("#EXTM3U")) {
|
||
return fmt.Errorf("manifest missing #EXTM3U header")
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func parseTargetDuration(manifest []byte) (int, error) {
|
||
for line := range bytes.SplitSeq(manifest, []byte("\n")) {
|
||
trimmed := bytes.TrimSpace(line)
|
||
if bytes.HasPrefix(trimmed, []byte("#EXT-X-TARGETDURATION:")) {
|
||
value := strings.TrimSpace(strings.TrimPrefix(string(trimmed), "#EXT-X-TARGETDURATION:"))
|
||
targetDuration, err := strconv.Atoi(value)
|
||
if err != nil {
|
||
return 0, fmt.Errorf("parse target duration %q: %w", value, err)
|
||
}
|
||
return targetDuration, nil
|
||
}
|
||
}
|
||
return 0, fmt.Errorf("manifest missing #EXT-X-TARGETDURATION")
|
||
}
|
||
|
||
func validateCopyPlaybackManifest(manifest []byte) error {
|
||
if err := validateManifestHeader(manifest); err != nil {
|
||
return err
|
||
}
|
||
|
||
targetDuration, err := parseTargetDuration(manifest)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if targetDuration <= 0 {
|
||
return fmt.Errorf("manifest target duration must be positive, got %d", targetDuration)
|
||
}
|
||
|
||
timeline, err := parseManifestTimeline(manifest)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if len(timeline.entries) == 0 {
|
||
return fmt.Errorf("manifest contains no playable media segments")
|
||
}
|
||
for _, entry := range timeline.entries {
|
||
if entry.duration <= 0 {
|
||
return fmt.Errorf("segment %d has non-positive duration %.6f", entry.number, entry.duration)
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func extractMapURI(line []byte) string {
|
||
const marker = `URI="`
|
||
text := string(line)
|
||
start := strings.Index(text, marker)
|
||
if start < 0 {
|
||
return ""
|
||
}
|
||
start += len(marker)
|
||
end := strings.Index(text[start:], `"`)
|
||
if end < 0 {
|
||
return ""
|
||
}
|
||
return text[start : start+end]
|
||
}
|
||
|
||
func manifestURIToFilename(uri string) string {
|
||
base, _, _ := strings.Cut(uri, "?")
|
||
return filepath.Base(base)
|
||
}
|
||
|
||
func manifestStartupFiles(manifest []byte, maxSegments int) ([]string, int) {
|
||
files := make([]string, 0, maxSegments+1)
|
||
segmentCount := 0
|
||
|
||
for line := range bytes.SplitSeq(manifest, []byte("\n")) {
|
||
trimmed := bytes.TrimSpace(line)
|
||
if len(trimmed) == 0 {
|
||
continue
|
||
}
|
||
if bytes.HasPrefix(trimmed, []byte("#EXT-X-MAP:")) {
|
||
if uri := extractMapURI(trimmed); uri != "" {
|
||
files = append(files, manifestURIToFilename(uri))
|
||
}
|
||
continue
|
||
}
|
||
if trimmed[0] == '#' {
|
||
continue
|
||
}
|
||
|
||
files = append(files, manifestURIToFilename(string(trimmed)))
|
||
segmentCount++
|
||
if segmentCount >= maxSegments {
|
||
break
|
||
}
|
||
}
|
||
|
||
return files, segmentCount
|
||
}
|
||
|
||
func startupFilesReady(manifest []byte, outputDir string, requiredSegments int) bool {
|
||
files, segmentCount := manifestStartupFiles(manifest, requiredSegments)
|
||
if segmentCount < requiredSegments {
|
||
return false
|
||
}
|
||
|
||
for _, name := range files {
|
||
info, err := os.Stat(filepath.Join(outputDir, name))
|
||
if err != nil || info.Size() <= 0 {
|
||
return false
|
||
}
|
||
}
|
||
|
||
return true
|
||
}
|
||
|
||
type manifestSegmentEntry struct {
|
||
number int
|
||
duration float64
|
||
}
|
||
|
||
type manifestTimeline struct {
|
||
mediaSequence int
|
||
entries []manifestSegmentEntry
|
||
}
|
||
|
||
func parseManifestTimeline(manifest []byte) (manifestTimeline, error) {
|
||
if err := validateManifestHeader(manifest); err != nil {
|
||
return manifestTimeline{}, err
|
||
}
|
||
|
||
timeline := manifestTimeline{}
|
||
currentNumber := 0
|
||
var pendingDuration float64
|
||
var haveDuration bool
|
||
|
||
for line := range bytes.SplitSeq(manifest, []byte("\n")) {
|
||
trimmed := bytes.TrimSpace(line)
|
||
if len(trimmed) == 0 {
|
||
continue
|
||
}
|
||
|
||
if bytes.HasPrefix(trimmed, []byte("#EXT-X-MEDIA-SEQUENCE:")) {
|
||
value := strings.TrimSpace(strings.TrimPrefix(string(trimmed), "#EXT-X-MEDIA-SEQUENCE:"))
|
||
sequence, err := strconv.Atoi(value)
|
||
if err != nil {
|
||
return manifestTimeline{}, fmt.Errorf("parse media sequence %q: %w", value, err)
|
||
}
|
||
timeline.mediaSequence = sequence
|
||
currentNumber = sequence
|
||
continue
|
||
}
|
||
|
||
if bytes.HasPrefix(trimmed, []byte("#EXTINF:")) {
|
||
value := strings.TrimSpace(strings.TrimSuffix(strings.TrimPrefix(string(trimmed), "#EXTINF:"), ","))
|
||
duration, err := strconv.ParseFloat(value, 64)
|
||
if err != nil {
|
||
return manifestTimeline{}, fmt.Errorf("parse segment duration %q: %w", value, err)
|
||
}
|
||
pendingDuration = duration
|
||
haveDuration = true
|
||
continue
|
||
}
|
||
|
||
if trimmed[0] == '#' {
|
||
continue
|
||
}
|
||
|
||
if !haveDuration {
|
||
continue
|
||
}
|
||
|
||
segmentNumber := currentNumber
|
||
if parsed, err := ParseSegmentNumber(filepath.Base(string(trimmed))); err == nil {
|
||
segmentNumber = parsed
|
||
}
|
||
|
||
timeline.entries = append(timeline.entries, manifestSegmentEntry{
|
||
number: segmentNumber,
|
||
duration: pendingDuration,
|
||
})
|
||
currentNumber = segmentNumber + 1
|
||
haveDuration = false
|
||
}
|
||
|
||
return timeline, nil
|
||
}
|
||
|
||
func hlsSegmentExtension(opts TranscodeOpts) string {
|
||
if strings.EqualFold(opts.TargetCodecVideo, "copy") && !IsMPEG2VideoCodec(opts.SourceVideoCodec) {
|
||
return ".m4s"
|
||
}
|
||
return ".ts"
|
||
}
|
||
|
||
func segmentFilename(segNum int, opts TranscodeOpts) string {
|
||
return fmt.Sprintf("seg_%05d%s", segNum, hlsSegmentExtension(opts))
|
||
}
|
||
|
||
func segmentWaitTimeout(segmentDuration int) time.Duration {
|
||
if segmentDuration <= 0 {
|
||
segmentDuration = defaultSegmentDuration
|
||
}
|
||
timeout := time.Duration(segmentDuration)*time.Second + segmentWaitGrace
|
||
if timeout < minSegmentWait {
|
||
timeout = minSegmentWait
|
||
}
|
||
if timeout > maxSegmentWait {
|
||
timeout = maxSegmentWait
|
||
}
|
||
return timeout
|
||
}
|
||
|
||
func staleProducedWindow(segmentDuration int) time.Duration {
|
||
if segmentDuration <= 0 {
|
||
segmentDuration = defaultSegmentDuration
|
||
}
|
||
window := 2*time.Duration(segmentDuration)*time.Second + segmentWaitGrace
|
||
if window < minStaleProducedWindow {
|
||
window = minStaleProducedWindow
|
||
}
|
||
return window
|
||
}
|
||
|
||
// SegmentProgress reports the highest manifest-referenced segment that exists
|
||
// on disk with data. This is the produced media source of truth.
|
||
func (s *TranscodeSession) SegmentProgress(time.Time) SegmentProgress {
|
||
s.mu.Lock()
|
||
opts := s.opts
|
||
progress := SegmentProgress{
|
||
ProducedHead: opts.StartSegmentNumber - 1,
|
||
Running: s.running,
|
||
Restarting: s.restarting,
|
||
StartSegmentNumber: opts.StartSegmentNumber,
|
||
SegmentDuration: opts.SegmentDuration,
|
||
LastRequestedSegment: s.lastRequestedSegment,
|
||
}
|
||
s.mu.Unlock()
|
||
|
||
if progress.SegmentDuration <= 0 {
|
||
progress.SegmentDuration = defaultSegmentDuration
|
||
}
|
||
|
||
manifestPath := filepath.Join(s.outputDir, "stream.m3u8")
|
||
manifestInfo, statErr := os.Stat(manifestPath)
|
||
if statErr != nil {
|
||
return progress
|
||
}
|
||
progress.HasManifest = true
|
||
progress.ManifestModTime = manifestInfo.ModTime()
|
||
|
||
manifest, err := os.ReadFile(manifestPath)
|
||
if err != nil {
|
||
return progress
|
||
}
|
||
timeline, err := parseManifestTimeline(manifest)
|
||
if err != nil {
|
||
return progress
|
||
}
|
||
|
||
for _, entry := range timeline.entries {
|
||
segmentPath := filepath.Join(s.outputDir, segmentFilename(entry.number, opts))
|
||
info, err := os.Stat(segmentPath)
|
||
if err != nil || info.Size() <= 0 {
|
||
continue
|
||
}
|
||
progress.ProducedCount++
|
||
if entry.number > progress.ProducedHead {
|
||
progress.ProducedHead = entry.number
|
||
}
|
||
if info.ModTime().After(progress.LastProducedAt) {
|
||
progress.LastProducedAt = info.ModTime()
|
||
}
|
||
}
|
||
|
||
return progress
|
||
}
|
||
|
||
// SegmentRecoveryDecision determines whether a missing segment should briefly
|
||
// wait for ffmpeg or immediately use the seek-restart path.
|
||
func (s *TranscodeSession) SegmentRecoveryDecision(segNum int, now time.Time) SegmentRecoveryDecision {
|
||
progress := s.SegmentProgress(now)
|
||
decision := SegmentRecoveryDecision{
|
||
WaitTimeout: segmentWaitTimeout(progress.SegmentDuration),
|
||
RestartOnTimeout: true,
|
||
Progress: progress,
|
||
}
|
||
|
||
switch {
|
||
// Restarting must be checked before Running: the restart window runs
|
||
// with running=false, and a concurrent segment request must wait out
|
||
// the in-flight restart rather than trigger another one. Dueling
|
||
// restarts keep preempting the segment the player is blocked on,
|
||
// which surfaces as a seek/intro-skip freeze (issue #243).
|
||
case progress.Restarting:
|
||
decision.Wait = true
|
||
decision.WaitTimeout = activeSegmentWait
|
||
decision.RestartOnTimeout = false
|
||
decision.Reason = "transcode_restarting"
|
||
case !progress.Running:
|
||
decision.Reason = "transcode_not_running"
|
||
case segNum < progress.StartSegmentNumber:
|
||
decision.Reason = "before_start_segment"
|
||
case segNum <= progress.ProducedHead:
|
||
decision.Reason = "segment_missing_behind_produced_head"
|
||
case !progress.HasManifest:
|
||
if segNum <= progress.StartSegmentNumber+1 {
|
||
decision.Wait = true
|
||
decision.WaitTimeout = activeSegmentWait
|
||
decision.RestartOnTimeout = false
|
||
decision.Reason = "startup_manifest_not_ready"
|
||
} else {
|
||
decision.Reason = "startup_request_beyond_window"
|
||
}
|
||
case segNum > progress.ProducedHead+maxSequentialMissingSegments:
|
||
decision.Reason = "request_beyond_produced_window"
|
||
case progress.ProducedHead >= progress.StartSegmentNumber && now.Sub(progress.LastProducedAt) > staleProducedWindow(progress.SegmentDuration):
|
||
decision.Reason = "produced_output_stale"
|
||
default:
|
||
decision.Wait = true
|
||
decision.WaitTimeout = activeSegmentWait
|
||
decision.RestartOnTimeout = false
|
||
decision.Reason = "near_produced_head"
|
||
}
|
||
|
||
return decision
|
||
}
|
||
|
||
// GenerateFullManifest builds a complete VOD-style HLS manifest that lists
|
||
// every segment for the full media duration, matching Jellyfin's approach.
|
||
// The player can seek to any position immediately; the backend produces
|
||
// segments on demand when they are requested via HandleGetTranscodeSegment.
|
||
//
|
||
// segPrefix is prepended to each segment filename (e.g. "segment/") and
|
||
// rawQuery is appended as a query string (e.g. auth tokens).
|
||
func (s *TranscodeSession) GenerateFullManifest(segPrefix, rawQuery string) []byte {
|
||
opts := s.Opts()
|
||
totalDur := opts.TotalDuration
|
||
segDur := opts.SegmentDuration
|
||
if segDur <= 0 {
|
||
segDur = defaultSegmentDuration
|
||
}
|
||
if totalDur <= 0 {
|
||
totalDur = float64(segDur) // fallback: single segment
|
||
}
|
||
|
||
segCount := int(math.Ceil(totalDur / float64(segDur)))
|
||
if segCount < 1 {
|
||
segCount = 1
|
||
}
|
||
|
||
var suffix string
|
||
if rawQuery != "" {
|
||
suffix = "?" + rawQuery
|
||
}
|
||
|
||
segExt := hlsSegmentExtension(opts)
|
||
hlsVersion := 3
|
||
if segExt == ".m4s" {
|
||
hlsVersion = 7
|
||
}
|
||
|
||
var buf bytes.Buffer
|
||
buf.WriteString("#EXTM3U\n")
|
||
buf.WriteString(fmt.Sprintf("#EXT-X-VERSION:%d\n", hlsVersion))
|
||
buf.WriteString(fmt.Sprintf("#EXT-X-TARGETDURATION:%d\n", segDur))
|
||
buf.WriteString("#EXT-X-MEDIA-SEQUENCE:0\n")
|
||
buf.WriteString("#EXT-X-PLAYLIST-TYPE:VOD\n")
|
||
|
||
if segExt == ".m4s" {
|
||
buf.WriteString(fmt.Sprintf("#EXT-X-MAP:URI=\"%sinit.mp4%s\"\n", segPrefix, suffix))
|
||
}
|
||
|
||
for i := range segCount {
|
||
dur := float64(segDur)
|
||
if i == segCount-1 {
|
||
// Last segment covers the remainder.
|
||
dur = totalDur - float64(i)*float64(segDur)
|
||
if dur <= 0 {
|
||
dur = float64(segDur)
|
||
}
|
||
}
|
||
buf.WriteString(fmt.Sprintf("#EXTINF:%.6f,\n", dur))
|
||
buf.WriteString(fmt.Sprintf("%sseg_%05d%s%s\n", segPrefix, i, segExt, suffix))
|
||
}
|
||
|
||
buf.WriteString("#EXT-X-ENDLIST\n")
|
||
return buf.Bytes()
|
||
}
|
||
|
||
// GetSegment returns the file path of a named segment if it exists.
|
||
func (s *TranscodeSession) GetSegment(name string) (string, error) {
|
||
// Sanitize the name to prevent directory traversal.
|
||
clean := filepath.Base(name)
|
||
segPath := filepath.Join(s.outputDir, clean)
|
||
|
||
info, err := os.Stat(segPath)
|
||
if err != nil {
|
||
if os.IsNotExist(err) {
|
||
return "", ErrSegmentNotFound
|
||
}
|
||
return "", fmt.Errorf("stat segment: %w", err)
|
||
}
|
||
if info.Size() <= 0 {
|
||
// ffmpeg can create init.mp4 before it has written any bytes. Treat
|
||
// zero-byte files as not-ready so callers fall back to WaitForSegment.
|
||
return "", ErrSegmentNotFound
|
||
}
|
||
return segPath, nil
|
||
}
|
||
|
||
// Close terminates the ffmpeg process and removes the temporary output directory.
|
||
func (s *TranscodeSession) Close() error {
|
||
return s.shutdown(true)
|
||
}
|
||
|
||
// CloseProcess terminates the ffmpeg process but leaves the output directory in
|
||
// place. It is used when another session owns the same output directory (e.g. a
|
||
// concurrent reconstruct race loser): tearing down this duplicate must not wipe
|
||
// the segments and init.mp4 the winning session is actively serving.
|
||
func (s *TranscodeSession) CloseProcess() error {
|
||
return s.shutdown(false)
|
||
}
|
||
|
||
// shutdown kills the ffmpeg process and, when removeOutput is true, removes the
|
||
// temporary output directory.
|
||
func (s *TranscodeSession) shutdown(removeOutput bool) error {
|
||
s.StopThrottler()
|
||
// Cancel the context to kill the process (no mutex needed for cancel).
|
||
if s.cancel != nil {
|
||
s.cancel()
|
||
}
|
||
|
||
// Wait for the monitor goroutine to finish reaping the process.
|
||
// This avoids a deadlock: the goroutine needs s.mu to mark running=false,
|
||
// so we must not hold s.mu while waiting.
|
||
// done is nil when no process was started (e.g. test-only sessions).
|
||
if s.done != nil {
|
||
<-s.done
|
||
}
|
||
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
|
||
s.running = false
|
||
|
||
// Clean up temporary directory.
|
||
if removeOutput && s.outputDir != "" {
|
||
if err := os.RemoveAll(s.outputDir); err != nil {
|
||
return fmt.Errorf("remove output dir: %w", err)
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// Done returns a channel that closes when the current ffmpeg process exits.
|
||
func (s *TranscodeSession) Done() <-chan struct{} {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
return s.done
|
||
}
|
||
|
||
// IsRunning reports whether the ffmpeg process is still running.
|
||
func (s *TranscodeSession) IsRunning() bool {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
return s.running
|
||
}
|
||
|
||
// WaitError returns the error from the last ffmpeg process exit, or nil if
|
||
// the process exited cleanly. A nil return means all output was written
|
||
// successfully and segments should remain servable.
|
||
func (s *TranscodeSession) WaitError() error {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
return s.waitErr
|
||
}
|
||
|
||
// Opts returns the TranscodeOpts used to create this session (for testing).
|
||
func (s *TranscodeSession) Opts() TranscodeOpts {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
return s.opts
|
||
}
|
||
|
||
// SetAudioTrackIndex updates the audio track index in the session's opts.
|
||
// Must be called before Restart() to take effect on the new ffmpeg process.
|
||
func (s *TranscodeSession) SetAudioTrackIndex(index int) {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
s.opts.AudioTrackIndex = index
|
||
}
|
||
|
||
// cleanStaleSegments removes segment files at or after startSegment and the
|
||
// old manifest so a restarted copy-mode FFmpeg process writes clean output.
|
||
// The init.mp4 is preserved — its codec configuration is derived from the
|
||
// source file and is identical across restarts.
|
||
func (s *TranscodeSession) cleanStaleSegments(startSegment int) {
|
||
entries, err := os.ReadDir(s.outputDir)
|
||
if err != nil {
|
||
return
|
||
}
|
||
for _, entry := range entries {
|
||
name := entry.Name()
|
||
if name == "stream.m3u8" {
|
||
os.Remove(filepath.Join(s.outputDir, name))
|
||
continue
|
||
}
|
||
if name == "init.mp4" {
|
||
continue
|
||
}
|
||
segNum, parseErr := ParseSegmentNumber(name)
|
||
if parseErr != nil {
|
||
continue
|
||
}
|
||
if segNum >= startSegment {
|
||
os.Remove(filepath.Join(s.outputDir, name))
|
||
}
|
||
}
|
||
}
|
||
|
||
// Restart kills the current ffmpeg process and starts a new one seeking to
|
||
// the given position. startSegment sets -hls_segment_start_number so that
|
||
// output filenames align with the expected segment numbering. Existing
|
||
// segment files are preserved so backward seeks can reuse them; for
|
||
// copy-mode sessions, stale segments at or after the restart point are
|
||
// cleaned to prevent serving data from the wrong timeline position.
|
||
func (s *TranscodeSession) Restart(ctx context.Context, seekSeconds float64, startSegment int) error {
|
||
s.mu.Lock()
|
||
// Single-flight: a second caller arriving while a restart is in
|
||
// progress must not kill the process the first restart just started.
|
||
// It returns immediately and the caller falls through to
|
||
// WaitForSegment, which polls through the in-flight restart.
|
||
if s.restarting {
|
||
s.mu.Unlock()
|
||
return nil
|
||
}
|
||
s.restarting = true
|
||
cancelCurrent := s.cancel
|
||
done := s.done
|
||
s.mu.Unlock()
|
||
s.StopThrottler()
|
||
|
||
// Kill current process without removing output directory.
|
||
if cancelCurrent != nil {
|
||
cancelCurrent()
|
||
}
|
||
if done != nil {
|
||
<-done
|
||
}
|
||
|
||
s.mu.Lock()
|
||
s.running = false
|
||
s.waitErr = nil
|
||
if s.stderr != nil {
|
||
s.stderr.Reset()
|
||
}
|
||
s.restartCount++
|
||
opts := s.opts
|
||
s.mu.Unlock()
|
||
|
||
// Copy-mode restarts must clean stale segments so ffmpeg writes fresh
|
||
// output. Encoded transcodes keep old segments for backward seek reuse.
|
||
if strings.EqualFold(opts.TargetCodecVideo, "copy") {
|
||
s.cleanStaleSegments(startSegment)
|
||
}
|
||
|
||
opts.SeekSeconds = seekSeconds
|
||
opts.StartSegmentNumber = startSegment
|
||
opts.FastStart = false // seek-restarts use veryfast for better quality
|
||
|
||
args := buildFFmpegArgs(opts)
|
||
bin := opts.FFmpegPath
|
||
if bin == "" {
|
||
bin = ffmpegBinary()
|
||
}
|
||
|
||
log.Printf("playback: ffmpeg restart cmd: %s %s", bin, strings.Join(args, " "))
|
||
s.logFFmpegEvent(ctx, "ffmpeg process restart", "")
|
||
|
||
ctx, cancel := context.WithCancel(ctx)
|
||
cmd := exec.CommandContext(ctx, bin, args...)
|
||
stdinPipe, err := cmd.StdinPipe()
|
||
if err != nil {
|
||
cancel()
|
||
s.mu.Lock()
|
||
s.restarting = false
|
||
s.waitErr = err
|
||
s.mu.Unlock()
|
||
return fmt.Errorf("create stdin pipe: %w", err)
|
||
}
|
||
cmd.Dir = opts.OutputDir
|
||
cmd.Stderr = s.newStderrWriter(ctx)
|
||
cmd.WaitDelay = 3 * time.Second
|
||
|
||
if err := cmd.Start(); err != nil {
|
||
cancel()
|
||
s.mu.Lock()
|
||
s.restarting = false
|
||
s.waitErr = err
|
||
s.mu.Unlock()
|
||
s.logFFmpegEvent(ctx, "ffmpeg process exit error", err.Error())
|
||
return fmt.Errorf("restart ffmpeg: %w", err)
|
||
}
|
||
|
||
s.mu.Lock()
|
||
if s.stdinPipe != nil {
|
||
s.stdinPipe.Close()
|
||
}
|
||
s.cmd = cmd
|
||
s.cancel = cancel
|
||
s.opts = opts
|
||
s.running = true
|
||
s.restarting = false
|
||
s.stdinPipe = stdinPipe
|
||
s.lastRequestedSegment = startSegment
|
||
s.done = make(chan struct{})
|
||
hook := s.restartHook
|
||
s.mu.Unlock()
|
||
|
||
go func() {
|
||
waitErr := cmd.Wait()
|
||
s.flushStderr(ctx)
|
||
s.mu.Lock()
|
||
s.running = false
|
||
s.waitErr = waitErr
|
||
s.mu.Unlock()
|
||
s.logWaitResult(ctx, waitErr)
|
||
close(s.done)
|
||
}()
|
||
|
||
if hook != nil {
|
||
hook(ctx)
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// WaitForSegment polls until the named segment file exists on disk or the
|
||
// timeout expires. Returns the segment file path on success.
|
||
//
|
||
// Segments are served as soon as they appear on disk. The -hls_flags temp_file
|
||
// flag ensures ffmpeg writes to a .tmp file and atomically renames on completion,
|
||
// so a successful stat means the segment is fully written.
|
||
func (s *TranscodeSession) WaitForSegment(name string, timeout time.Duration) (string, error) {
|
||
clean := filepath.Base(name)
|
||
segPath := filepath.Join(s.outputDir, clean)
|
||
|
||
deadline := time.After(timeout)
|
||
for {
|
||
info, statErr := os.Stat(segPath)
|
||
segReady := statErr == nil && info.Size() > 0
|
||
if segReady {
|
||
return segPath, nil
|
||
}
|
||
|
||
s.mu.Lock()
|
||
running := s.running
|
||
restarting := s.restarting
|
||
waitErr := s.waitErr
|
||
s.mu.Unlock()
|
||
|
||
if restarting {
|
||
select {
|
||
case <-deadline:
|
||
return "", ErrSegmentNotFound
|
||
case <-time.After(100 * time.Millisecond):
|
||
continue
|
||
}
|
||
}
|
||
|
||
if !running && waitErr != nil {
|
||
return "", fmt.Errorf("%w: %v", ErrTranscodeFailed, waitErr)
|
||
}
|
||
// If ffmpeg finished cleanly but the segment doesn't exist,
|
||
// it won't appear later — fail fast.
|
||
if !running {
|
||
return "", ErrSegmentNotFound
|
||
}
|
||
|
||
select {
|
||
case <-deadline:
|
||
return "", ErrSegmentNotFound
|
||
case <-time.After(100 * time.Millisecond):
|
||
}
|
||
}
|
||
}
|
||
|
||
// RewriteManifestPaths prefixes relative segment references in an HLS manifest
|
||
// with segPrefix (e.g. "segment/") and optionally appends rawQuery as a query
|
||
// string. This ensures the HLS player's segment requests match server routes
|
||
// and preserve any auth or cache-busting parameters from the manifest URL.
|
||
func RewriteManifestPaths(manifest []byte, segPrefix, rawQuery string) ([]byte, error) {
|
||
if err := validateManifestHeader(manifest); err != nil {
|
||
return nil, fmt.Errorf("invalid manifest: %w", err)
|
||
}
|
||
|
||
var suffix string
|
||
if rawQuery != "" {
|
||
suffix = "?" + rawQuery
|
||
}
|
||
|
||
lines := bytes.Split(manifest, []byte("\n"))
|
||
for i, line := range lines {
|
||
trimmed := bytes.TrimSpace(line)
|
||
if len(trimmed) == 0 {
|
||
continue
|
||
}
|
||
|
||
// Rewrite #EXT-X-MAP:URI="filename"
|
||
if bytes.HasPrefix(trimmed, []byte("#EXT-X-MAP:")) {
|
||
rewritten, err := rewriteMapURI(trimmed, segPrefix, suffix)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
lines[i] = rewritten
|
||
continue
|
||
}
|
||
|
||
// Skip other tags/comments.
|
||
if trimmed[0] == '#' {
|
||
continue
|
||
}
|
||
|
||
// Segment filename line.
|
||
lines[i] = []byte(segPrefix + string(trimmed) + suffix)
|
||
}
|
||
return bytes.Join(lines, []byte("\n")), nil
|
||
}
|
||
|
||
// rewriteMapURI rewrites the URI value inside an #EXT-X-MAP tag.
|
||
func rewriteMapURI(line []byte, segPrefix, suffix string) ([]byte, error) {
|
||
uriStart := bytes.Index(line, []byte(`URI="`))
|
||
if uriStart < 0 {
|
||
return nil, fmt.Errorf("invalid #EXT-X-MAP line: missing URI attribute")
|
||
}
|
||
uriStart += 5 // skip past URI="
|
||
uriEnd := bytes.IndexByte(line[uriStart:], '"')
|
||
if uriEnd < 0 {
|
||
return nil, fmt.Errorf("invalid #EXT-X-MAP line: unterminated URI attribute")
|
||
}
|
||
uriEnd += uriStart
|
||
|
||
oldURI := string(line[uriStart:uriEnd])
|
||
newURI := segPrefix + oldURI + suffix
|
||
|
||
result := make([]byte, 0, len(line)+len(newURI)-len(oldURI))
|
||
result = append(result, line[:uriStart]...)
|
||
result = append(result, []byte(newURI)...)
|
||
result = append(result, line[uriEnd:]...)
|
||
return result, nil
|
||
}
|
||
|
||
// AppendManifestQueryParam appends a single "key=value" query parameter to every
|
||
// segment and #EXT-X-MAP init URI in an HLS media playlist, preserving any query
|
||
// the URI already carries (using "?" or "&" as appropriate). The value is
|
||
// appended verbatim — callers must supply a URL-safe value (a signed stream token
|
||
// is base64url and safe). A manifest without a valid #EXTM3U header is returned
|
||
// unchanged so a non-manifest body is never corrupted.
|
||
//
|
||
// It exists for the API proxy boundary: a transcode node builds its manifest from
|
||
// the forwarded request query, which deliberately omits the signed stream token
|
||
// ("st") so the token never reaches the node URL or its logs. That leaves the
|
||
// node-built segment URIs token-less, so a later segment fetch after a node/API
|
||
// restart cannot reconstruct the session. Re-injecting the client-facing token
|
||
// here keeps reconstruction working without ever exposing it to the node.
|
||
func AppendManifestQueryParam(manifest []byte, key, value string) []byte {
|
||
if key == "" || validateManifestHeader(manifest) != nil {
|
||
return manifest
|
||
}
|
||
param := key + "=" + value
|
||
lines := bytes.Split(manifest, []byte("\n"))
|
||
for i, line := range lines {
|
||
trimmed := bytes.TrimSpace(line)
|
||
if len(trimmed) == 0 {
|
||
continue
|
||
}
|
||
if bytes.HasPrefix(trimmed, []byte("#EXT-X-MAP:")) {
|
||
lines[i] = appendMapURIQueryParam(trimmed, param)
|
||
continue
|
||
}
|
||
if trimmed[0] == '#' {
|
||
continue
|
||
}
|
||
lines[i] = appendURIQueryParam(trimmed, param)
|
||
}
|
||
return bytes.Join(lines, []byte("\n"))
|
||
}
|
||
|
||
// appendURIQueryParam appends a "key=value" param to a bare URI line.
|
||
func appendURIQueryParam(uri []byte, param string) []byte {
|
||
sep := []byte("?")
|
||
if bytes.IndexByte(uri, '?') >= 0 {
|
||
sep = []byte("&")
|
||
}
|
||
out := make([]byte, 0, len(uri)+len(sep)+len(param))
|
||
out = append(out, uri...)
|
||
out = append(out, sep...)
|
||
out = append(out, param...)
|
||
return out
|
||
}
|
||
|
||
// appendMapURIQueryParam appends a "key=value" param to the URI inside an
|
||
// #EXT-X-MAP tag. A line without a well-formed URI attribute is returned
|
||
// unchanged.
|
||
func appendMapURIQueryParam(line []byte, param string) []byte {
|
||
uriStart := bytes.Index(line, []byte(`URI="`))
|
||
if uriStart < 0 {
|
||
return line
|
||
}
|
||
uriStart += 5 // skip past URI="
|
||
uriEnd := bytes.IndexByte(line[uriStart:], '"')
|
||
if uriEnd < 0 {
|
||
return line
|
||
}
|
||
uriEnd += uriStart
|
||
|
||
sep := []byte("?")
|
||
if bytes.IndexByte(line[uriStart:uriEnd], '?') >= 0 {
|
||
sep = []byte("&")
|
||
}
|
||
result := make([]byte, 0, len(line)+len(sep)+len(param))
|
||
result = append(result, line[:uriEnd]...)
|
||
result = append(result, sep...)
|
||
result = append(result, param...)
|
||
result = append(result, line[uriEnd:]...)
|
||
return result
|
||
}
|
||
|
||
func (s *TranscodeSession) manifestTimeoutError(timeout time.Duration) error {
|
||
s.mu.Lock()
|
||
running := s.running
|
||
waitErr := s.waitErr
|
||
stderr := ""
|
||
if s.stderr != nil {
|
||
stderr = truncateStderr(s.stderr.String())
|
||
}
|
||
s.mu.Unlock()
|
||
|
||
switch {
|
||
case waitErr != nil && stderr != "":
|
||
return fmt.Errorf("%w after %s: ffmpeg exited: %v (stderr: %s)", ErrManifestNotReady, timeout, waitErr, stderr)
|
||
case waitErr != nil:
|
||
return fmt.Errorf("%w after %s: ffmpeg exited: %v", ErrManifestNotReady, timeout, waitErr)
|
||
case running:
|
||
return fmt.Errorf("%w after %s: ffmpeg still running", ErrManifestNotReady, timeout)
|
||
default:
|
||
return fmt.Errorf("%w after %s: ffmpeg is no longer running", ErrManifestNotReady, timeout)
|
||
}
|
||
}
|
||
|
||
// IsCopyVideo reports whether this session is repackaging video without
|
||
// re-encoding. Copy-mode manifests must reflect FFmpeg's real fragment timing.
|
||
func (s *TranscodeSession) IsCopyVideo() bool {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
return strings.EqualFold(s.opts.TargetCodecVideo, "copy")
|
||
}
|
||
|
||
// SegmentStartTime reports the source-timeline start time of the requested
|
||
// segment using the current on-disk manifest. The bool return is false when
|
||
// the segment is not present in the manifest yet.
|
||
func (s *TranscodeSession) SegmentStartTime(segNum int) (float64, bool, error) {
|
||
s.mu.Lock()
|
||
manifestPath := filepath.Join(s.outputDir, "stream.m3u8")
|
||
baseSeekSeconds := s.opts.SeekSeconds
|
||
s.mu.Unlock()
|
||
|
||
manifest, err := os.ReadFile(manifestPath)
|
||
if err != nil {
|
||
if os.IsNotExist(err) {
|
||
return 0, false, ErrManifestNotReady
|
||
}
|
||
return 0, false, fmt.Errorf("read manifest: %w", err)
|
||
}
|
||
|
||
timeline, err := parseManifestTimeline(manifest)
|
||
if err != nil {
|
||
return 0, false, fmt.Errorf("parse manifest timeline: %w", err)
|
||
}
|
||
|
||
if len(timeline.entries) == 0 {
|
||
return 0, false, ErrManifestNotReady
|
||
}
|
||
|
||
currentTime := baseSeekSeconds
|
||
for _, entry := range timeline.entries {
|
||
if entry.number == segNum {
|
||
return currentTime, true, nil
|
||
}
|
||
currentTime += entry.duration
|
||
}
|
||
|
||
return 0, false, nil
|
||
}
|
||
|
||
// RestartSeekTarget resolves the source-timeline time to restart FFmpeg for
|
||
// the requested segment. Copy-mode sessions prefer the current manifest's real
|
||
// timing when available; encoded sessions use fixed-duration seek math
|
||
// matching the synthetic VOD manifest.
|
||
func (s *TranscodeSession) RestartSeekTarget(segNum int) (float64, bool, error) {
|
||
if strings.EqualFold(s.Opts().TargetCodecVideo, "copy") {
|
||
seekSeconds, ok, err := s.SegmentStartTime(segNum)
|
||
switch {
|
||
case err == nil && ok:
|
||
return seekSeconds, true, nil
|
||
case err != nil && !errors.Is(err, ErrManifestNotReady):
|
||
return 0, false, err
|
||
}
|
||
// Copy-mode fragments have variable durations, so the encoded
|
||
// `seg×dur` math would seek FFmpeg to the wrong source time and
|
||
// desync A/V after restart. When the manifest can't resolve this
|
||
// segment yet (ok=false with no error, including ErrManifestNotReady
|
||
// in a freshly reconstructed window), report the seek target as
|
||
// unresolved (0, false, nil) rather than guessing. The caller treats
|
||
// this as a retryable miss so the session keeps producing manifest
|
||
// until real timing is available.
|
||
return 0, false, nil
|
||
}
|
||
|
||
segDuration := defaultSegmentDuration
|
||
if opts := s.Opts(); opts.SegmentDuration > 0 {
|
||
segDuration = opts.SegmentDuration
|
||
}
|
||
return float64(segNum * segDuration), true, nil
|
||
}
|
||
|
||
// ReportSegmentDownloaded records that the client has downloaded the given
|
||
// segment number. Only updates if segNum exceeds the current high-water mark.
|
||
func (s *TranscodeSession) ReportSegmentDownloaded(segNum int) {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
if segNum > s.lastRequestedSegment {
|
||
s.lastRequestedSegment = segNum
|
||
}
|
||
}
|
||
|
||
// LastRequestedSegment returns the highest segment number downloaded by the client.
|
||
func (s *TranscodeSession) LastRequestedSegment() int {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
return s.lastRequestedSegment
|
||
}
|
||
|
||
// StartThrottler creates and starts a throttler for this session.
|
||
// No-op if thresholdSeconds <= 0 or stdinPipe is nil.
|
||
func (s *TranscodeSession) StartThrottler(thresholdSeconds int) {
|
||
s.mu.Lock()
|
||
if s.stdinPipe == nil || thresholdSeconds <= 0 {
|
||
s.mu.Unlock()
|
||
return
|
||
}
|
||
t := NewTranscodeThrottler(s, s.stdinPipe, thresholdSeconds, s.opts.SegmentDuration)
|
||
s.throttler = t
|
||
s.mu.Unlock()
|
||
t.Start()
|
||
}
|
||
|
||
// StopThrottler stops the throttler if one is running.
|
||
func (s *TranscodeSession) StopThrottler() {
|
||
s.mu.Lock()
|
||
t := s.throttler
|
||
s.throttler = nil
|
||
s.mu.Unlock()
|
||
if t != nil {
|
||
t.Stop()
|
||
}
|
||
}
|
||
|
||
type ffmpegStderrWriter struct {
|
||
session *TranscodeSession
|
||
ctx context.Context
|
||
partial []byte
|
||
}
|
||
|
||
func (w *ffmpegStderrWriter) Write(p []byte) (int, error) {
|
||
if len(p) == 0 {
|
||
return 0, nil
|
||
}
|
||
w.partial = append(w.partial, p...)
|
||
for {
|
||
idx := bytes.IndexByte(w.partial, '\n')
|
||
if idx < 0 {
|
||
break
|
||
}
|
||
line := strings.TrimRight(string(w.partial[:idx]), "\r")
|
||
w.session.logFFmpegLine(w.ctx, line)
|
||
w.partial = append([]byte(nil), w.partial[idx+1:]...)
|
||
}
|
||
return len(p), nil
|
||
}
|
||
|
||
func (w *ffmpegStderrWriter) Flush() {
|
||
if len(w.partial) == 0 {
|
||
return
|
||
}
|
||
w.session.logFFmpegLine(w.ctx, strings.TrimRight(string(w.partial), "\r"))
|
||
w.partial = nil
|
||
}
|
||
|
||
func (s *TranscodeSession) newStderrWriter(ctx context.Context) io.Writer {
|
||
lineWriter := &ffmpegStderrWriter{session: s, ctx: ctx}
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
if s.stderr == nil {
|
||
s.stderr = newBoundedTailBuffer(stderrTailMaxBytes)
|
||
}
|
||
s.stderrWriter = lineWriter
|
||
return io.MultiWriter(s.stderr, lineWriter)
|
||
}
|
||
|
||
func (s *TranscodeSession) flushStderr(ctx context.Context) {
|
||
s.mu.Lock()
|
||
writer := s.stderrWriter
|
||
s.stderrWriter = nil
|
||
s.mu.Unlock()
|
||
if writer != nil {
|
||
writer.Flush()
|
||
}
|
||
}
|
||
|
||
func (s *TranscodeSession) logFFmpegLine(ctx context.Context, line string) {
|
||
if s == nil || s.opts.FFmpegLogSink == nil {
|
||
return
|
||
}
|
||
|
||
line = strings.ToValidUTF8(line, "\uFFFD")
|
||
if strings.TrimSpace(line) == "" {
|
||
return
|
||
}
|
||
trimmed, truncated := truncateUTF8String(line, maxPersistedFFmpegChars)
|
||
if truncated {
|
||
trimmed += "...[truncated]"
|
||
}
|
||
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
|
||
if s.stderrLinesLogged >= maxPersistedFFmpegLines || s.stderrBytesLogged+len(trimmed) > maxPersistedFFmpegBytes {
|
||
s.stderrDroppedLines++
|
||
if !s.stderrCapLogged {
|
||
s.stderrCapLogged = true
|
||
attrs := s.ffmpegAttrsLocked()
|
||
attrs.DroppedLines = s.stderrDroppedLines
|
||
s.opts.FFmpegLogSink.WriteEvent(ctx, s.opts.SessionID, attrs, "ffmpeg stderr logging capped")
|
||
}
|
||
return
|
||
}
|
||
|
||
s.stderrLinesLogged++
|
||
s.stderrBytesLogged += len(trimmed)
|
||
s.stderrLineIndex++
|
||
attrs := s.ffmpegAttrsLocked()
|
||
attrs.LineIndex = s.stderrLineIndex
|
||
s.opts.FFmpegLogSink.WriteLine(ctx, s.opts.SessionID, attrs, trimmed)
|
||
}
|
||
|
||
func (s *TranscodeSession) logFFmpegEvent(ctx context.Context, message, exitError string) {
|
||
if s == nil || s.opts.FFmpegLogSink == nil {
|
||
return
|
||
}
|
||
s.mu.Lock()
|
||
attrs := s.ffmpegAttrsLocked()
|
||
attrs.ExitError = exitError
|
||
s.mu.Unlock()
|
||
s.opts.FFmpegLogSink.WriteEvent(ctx, s.opts.SessionID, attrs, message)
|
||
}
|
||
|
||
func (s *TranscodeSession) logWaitResult(ctx context.Context, waitErr error) {
|
||
if waitErr == nil {
|
||
s.logFFmpegEvent(ctx, "ffmpeg process exited", "")
|
||
return
|
||
}
|
||
s.logFFmpegEvent(ctx, "ffmpeg process exit error", formatWaitError(waitErr))
|
||
}
|
||
|
||
func (s *TranscodeSession) ffmpegAttrsLocked() FFmpegLogAttrs {
|
||
return FFmpegLogAttrs{
|
||
NodeType: s.opts.NodeType,
|
||
ExecutionMode: s.opts.ExecutionMode,
|
||
InputPath: s.opts.InputPath,
|
||
OutputDir: s.opts.OutputDir,
|
||
TargetResolution: s.opts.TargetResolution,
|
||
TargetVideoCodec: s.opts.TargetCodecVideo,
|
||
TargetAudioCodec: s.opts.TargetCodecAudio,
|
||
HWAccel: s.opts.HWAccel,
|
||
SeekSeconds: s.opts.SeekSeconds,
|
||
StartSegmentNumber: s.opts.StartSegmentNumber,
|
||
RestartCount: s.restartCount,
|
||
DroppedLines: s.stderrDroppedLines,
|
||
}
|
||
}
|
||
|
||
func formatWaitError(err error) string {
|
||
if err == nil {
|
||
return ""
|
||
}
|
||
if exitErr, ok := err.(*exec.ExitError); ok {
|
||
if status, ok := exitErr.Sys().(syscall.WaitStatus); ok {
|
||
return fmt.Sprintf("exit_code=%d: %v", status.ExitStatus(), err)
|
||
}
|
||
}
|
||
return err.Error()
|
||
}
|