Files
silo-server/internal/transcodenode/server.go
T
854d07cf8f feat(playback): add protocol v3 planning and recovery (#398)
* docs(playback): plan protocol v3 server implementation

* docs(playback): incorporate protocol v3 review

* feat(playback): implement protocol v3 server

* fix(playback): persist empty route diagnostics

* feat(playback): harden protocol v3 HDR routing

* feat(playback): complete protocol v3 client contract

* fix(playback): harden protocol v3 recovery

* fix(playback): restore dovi_rpu strip filter for DV remuxes

The v3 work renamed the Dolby Vision strip recipe to a dovi_split=mode=bl
bitstream filter that does not exist in stock FFmpeg or jellyfin-ffmpeg;
the probe failed closed on every deployment, disabling the new validated
DV7-to-HDR10 route and regressing the previously working dovi_rpu=strip=1
remux path from main. Restore dovi_rpu across the probe, remux and HLS
copy arguments, and the recipe-card constant.

Also from review: validate the remux DV mode for every profile (garbage
modes on non-P7 sources silently no-opped), reject preserve mode for P7
outright (a base-layer-only remux cannot preserve dual-layer DV), tag
dvhe sample entries only for the explicit v3 preserve recipe so legacy
web/jellycompat remuxes keep their pre-v3 hev1 labeling, and honor the
token-frozen DV mode in the proxy remux path instead of legacy-auto.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(playback): correct v3 planner policy and contract validation

Review fixes to the v3 planner and wire contracts:

- Bar Profile 7 sources from the non-strip progressive remux route: a
  base-layer-only remux can never deliver native dual-layer DV, so the
  planner no longer emits plans claiming validated Dolby Vision while
  the executed remux drops the enhancement layer.
- Accept the device-quirks feature flag from either capability location,
  matching every other dual-location feature check.
- Treat legacy hdr_unknown rows as HDR10 for HDR10-capable clients with
  a degradation warning instead of leaving them unplayable under v3.
- Honor bandwidth_cap_kbps as a hard ceiling in every quality mode and
  wire the previously dead Metered signal into conservative auto rungs.
- Degrade to the validated source-quality route instead of a terminal
  when only an implicit quality reduction demanded an unsupported
  transcode; explicit user-selected rungs keep terminal behavior.
- Bound inner capability lists and strings; compare attempt keys exactly
  instead of case-folded; make ParseTrackIDV3 strict about canonical
  numerics; accept dvdsub/pgssub/dvbsub aliases and stop promising
  burn-in for unknown subtitle codecs; probe every h264 encoder rather
  than requiring libx264; normalize the file-level bitrate fallback.
- Evaluate subtitle renderability against the engine each candidate
  route executes on, not always media3_direct.
- Pin the with-quirks attempt-key preimage arity in the cross-language
  fixture so the Kotlin client stays in lockstep.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(playback): harden v3 control-plane reliability

Review fixes to the v3 session, store, and handler layer:

- Bound concurrent replans with a slot semaphore: each replan pins a
  pooled connection for its advisory lock while issuing further store
  queries from the same pool, so an unbounded recovery storm could turn
  every connection into a lock holder and deadlock the server.
- Make CompleteReplan a real compare-and-swap (base-revision predicate,
  ErrReplanSupersededV3) and map BeginReplan insert races to a replay
  instead of a raw unique violation.
- Fingerprint start requests (request_digest column): an attempt ID
  reused with different input is now a 409-style conflict rather than a
  silent replay, and both replay paths check session liveness so dead
  sessions surface as retryable terminals.
- Pre-delete expired attempt rows on SaveAttempt so a retry during the
  cleanup window cannot wedge on an unreachable conflict.
- Align the in-memory store's semantics with Postgres and add DB-backed
  planstore tests (SILO_TEST_DATABASE_URL), including a regression test
  inserting every route-event name against the real CHECK constraint.
- Session manager: v3 route-set updates own RemuxDVMode outright so a
  replan onto an SDR source clears a stale strip mode; replacement
  reservations survive unrelated legacy stream updates; replacement
  admission excludes the replaced session explicitly instead of
  decrementing totals it may no longer be part of; the admission CAS
  loop is bounded and decider errors are logged.
- Map transient store failures to 500s instead of terminal 404/403s;
  authorize route events via identity-only projections after the rate
  limiter; keep sanitized diagnostics deterministic.
- Merge the server-computed durable plan key into replan exclusions so
  unreproducible client history cannot re-select the failed route.
- Remap tracks only when the effective edition changes (a same-file
  replan no longer switches audio to a lookalike track) and remap
  ID-only subtitle selections on edition fallback.
- Cache the v3/shadow feature flags for five seconds instead of one
  settings SELECT per playback request; stop remote transports
  best-effort when the start call times out; carry dvm/tid claims and
  the transport-scoped job identity through the legacy audio-change
  re-mint; index playback_route_events(received_at) for the retention
  delete; run store maintenance for DB-less deployments too.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(transcode): reap idle node jobs and gate WebVTT conversion

- Add an idle reaper to the transcode node: a job untouched by manifest
  or segment requests for ten minutes is closed and unregistered. After
  a v3 replan retires a transport ID, a stale in-flight stream token
  could resurrect the old job via reconstruct and encode to end-of-file
  for nobody; jobs waiting on readiness count registration as access
  and are never reaped mid-wait, and reaping keeps the recipe so a
  still-valid token reconstructs on the next hit.
- Reject bitmap subtitle tracks (PGS) on the .vtt conversion path with
  415 before headers are written instead of spawning an ffmpeg command
  that always fails mid-response, and make the extract-format override
  fall back to source-driven mapping for bitmap codecs.
- Drain error bodies on non-202 node responses so the HTTP transport
  can reuse connections.
- Pin the transcode-dir cleanup separator-boundary semantics with a
  regression test (a session ID sharing another's prefix must not
  retain foreign directories).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(playback): close v3 planner policy gaps from review

- Clamp the final transcode bitrate to bandwidth_cap_kbps: the ladder has
  no rung below 480p/1500kbps, so lower caps were silently exceeded even
  though the cap is documented as a hard delivery ceiling.
- Treat video-only media as audio-compatible instead of forcing an AAC
  conversion (or an audio_conversion_unsupported terminal) onto a file
  with no audio stream. Tracks whose codec failed to probe keep the gate.
- Only promise a bitmap subtitle sidecar for embedded PGS with an engine
  that renders embedded bitmap: external/downloaded bitmap and embedded
  DVD/DVB published artifact URLs that always failed at fetch. They now
  fall through to burn-in or its terminal.
- Accept client_video_transformations_v1 from either client_features or
  the nested context when validating client-executor transformations,
  matching the planner's dual-source reads.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(playback): probe and execute DV remuxes with one ffmpeg binary

The v3 transformation registry probed the configured playback.ffmpeg_path
while progressive remux execution resolved the process-global discovery
path, so a deployment where only one binary carries dovi_rpu could plan a
server_dv7_to_hdr10 route and then fail it at stream time. Resolution now
goes through a shared ResolveFFmpegPath (configured path first, discovery
fallback — the same rule the transcode pipeline already used), the
dovi_rpu probe is cached per binary path, and the stream handler and proxy
worker pass their configured path into ServeRemuxWithDVMode.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(playback): harden v3 replan identity and control-plane limits

- Seed failure-replan track selections from the durable current plan
  before overlaying the request: after an alternate-version fallback the
  normalized request still carries requested-edition track IDs, so a
  replan omitting unchanged tracks was rejected as a track/file mismatch.
- Remap ID-only audio selections across edition changes (parse the ID to
  an index like the subtitle remap already does) instead of leaving a
  stale file-bound ID to fail validation.
- Release the node planner reservation when a prepared remote transport
  rolls back after the node accepted the job; repeated failed starts
  could otherwise pin max-job/bandwidth budgets for the full reservation
  age.
- Size the replan semaphore below the PostgreSQL pool via a store
  capacity advisor: with max_connections at or below the fixed bound,
  advisory-lock holders could starve the inner store queries they need
  to finish.
- Contain shadow-planner panics with a recover boundary; it runs on a
  bare goroutine where an escaped panic kills the process for what is
  telemetry-only work. Document why the memory store's session lock is
  deliberately a no-op.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(transcode): serialize node job teardown against reconstructs

- Look up and touch manifest/segment sessions in one critical section so
  the idle reaper cannot unregister a job between the lookup and its
  liveness refresh.
- Re-validate each reap candidate under the per-session lifecycle lock
  before closing it: Close removes the output directory, and without the
  lock it could race a token reconstruct and wipe the segments the fresh
  ffmpeg is writing.
- Take the lifecycle lock in handleStop so a stop racing a RequireReady
  start's readiness wait blocks until registration and tears the job
  down, instead of 404ing and orphaning the ffmpeg until the reaper.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
2026-07-14 11:51:27 -04:00

1001 lines
38 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package transcodenode
import (
"context"
"encoding/json"
"errors"
"log/slog"
"net/http"
"os"
"path/filepath"
"runtime"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/go-chi/chi/v5"
"golang.org/x/sync/singleflight"
"github.com/Silo-Server/silo-server/internal/chapterthumbs"
"github.com/Silo-Server/silo-server/internal/nodeconfig"
"github.com/Silo-Server/silo-server/internal/nodesessions"
"github.com/Silo-Server/silo-server/internal/playback"
"github.com/Silo-Server/silo-server/internal/streamtoken"
)
// TranscodeStartRequest is the JSON body for POST /transcode/start.
type TranscodeStartRequest struct {
SessionID string `json:"session_id"`
InputPath string `json:"input_path"`
SourceVideoCodec string `json:"source_video_codec"`
VideoBitstreamFilter string `json:"video_bitstream_filter,omitempty"`
SeekSeconds float64 `json:"seek_seconds"`
StartSegmentNumber int `json:"start_segment_number"`
TargetResolution string `json:"target_resolution"`
TargetCodecVideo string `json:"target_codec_video"`
TargetCodecAudio string `json:"target_codec_audio"`
TargetBitrateKbps int `json:"target_bitrate_kbps"`
SegmentDuration int `json:"segment_duration"`
HWAccel string `json:"hw_accel"`
AudioTrackIndex int `json:"audio_track_index"`
SubtitleTrackIndex int `json:"subtitle_track_index"`
SubtitleBurnIn bool `json:"subtitle_burn_in"`
SubtitleCodec string `json:"subtitle_codec,omitempty"`
TotalDuration float64 `json:"total_duration"`
RequireReady bool `json:"require_ready,omitempty"`
}
// TranscodeStartResponse is the JSON response for POST /transcode/start.
type TranscodeStartResponse struct {
SessionID string `json:"session_id"`
Status string `json:"status"`
HWAccel string `json:"hw_accel,omitempty"`
}
// HealthResponse is the JSON response for GET /api/v1/health.
type HealthResponse struct {
Status string `json:"status"`
ActiveJobs int32 `json:"active_jobs"`
}
// sessionIdleTTL is how long a job may go without a manifest or segment
// request before the idle reaper closes it. Reaping is safe because a client
// that comes back later re-presents its still-valid stream token and the job
// reconstructs seeked to the requested segment; without the reaper, a job
// whose audience vanished (e.g. a v3 replan retired its transport id and a
// stale in-flight token resurrected the old one) encodes to end-of-file for
// nobody.
const sessionIdleTTL = 10 * time.Minute
// sessionReapInterval is how often the idle reaper sweeps for stale jobs.
const sessionReapInterval = time.Minute
// Server is the HTTP handler for transcode mode.
type Server struct {
watcher *nodeconfig.Watcher
tracker *nodesessions.Tracker
ffmpegSink playback.FFmpegLogSink
sessions map[string]*playback.TranscodeSession
// lastAccess records, per registered session id, when a manifest or segment
// request last touched the job (registration counts as the first access).
// Guarded by mu alongside sessions; the idle reaper closes jobs whose entry
// is older than sessionIdleTTL.
lastAccess map[string]time.Time
reaperOnce sync.Once
mu sync.RWMutex
activeJobs atomic.Int32
// reconstructGroup single-flights node-side session reconstruction per session
// id so a post-restart wave of concurrent manifest/segment requests for the same
// lost session spawns exactly one ffmpeg, never racing duplicates into the shared
// output directory.
reconstructGroup singleflight.Group
// reconstructSem bounds how many sessions may be reconstructed (ffmpeg
// re-spawned) at once after a node restart, pacing the cold-start burst instead
// of stampeding the host. Lazily sized to NumCPU on first use.
reconstructSemOnce sync.Once
reconstructSem chan struct{}
// lifecycleMu guards lifecycleLocks, the per-session mutexes that serialize
// every path which spawns ffmpeg into a session's output dir (fresh start and
// reconstruct). reconstructGroup only single-flights reconstructs against each
// other; without this a reconstruct racing a fresh /transcode/start could run
// two ffmpeg writers against the same dir.
lifecycleMu sync.Mutex
lifecycleLocks map[string]*sessionLifecycleLock
// recipeStore is the control-plane recipe store consulted when a forwarded
// token carries no recipe (the jellycompat node hop). Nil disables that path.
recipeStore recipeStore
}
// sessionLifecycleLock is a refcounted per-session mutex; the refcount lets the
// node drop the map entry once no path holds or waits on it so the map stays
// bounded over the node's lifetime.
type sessionLifecycleLock struct {
mu sync.Mutex
refs int
}
// lockSessionLifecycle acquires the per-session lifecycle mutex and returns a
// release func. Held across "check existing → spawn → register" so a fresh start
// and a reconstruct never run concurrent ffmpeg writers for one session's dir.
func (s *Server) lockSessionLifecycle(sessionID string) func() {
s.lifecycleMu.Lock()
if s.lifecycleLocks == nil {
s.lifecycleLocks = make(map[string]*sessionLifecycleLock)
}
lk := s.lifecycleLocks[sessionID]
if lk == nil {
lk = &sessionLifecycleLock{}
s.lifecycleLocks[sessionID] = lk
}
lk.refs++
s.lifecycleMu.Unlock()
lk.mu.Lock()
return func() {
lk.mu.Unlock()
s.lifecycleMu.Lock()
lk.refs--
if lk.refs == 0 {
delete(s.lifecycleLocks, sessionID)
}
s.lifecycleMu.Unlock()
}
}
// restartSessionLocked re-spawns session under the per-session lifecycle lock so
// a segment-recovery restart can never race a fresh start, reconstruct, or
// another restart into the same output directory. It holds the lock only across
// the cancel→respawn transition inside Restart and releases it before the caller
// waits on segments. Under the lock it confirms session is still the live mapped
// session; a concurrent teardown or reconstruct that replaced it yields
// ErrSessionSuperseded rather than re-spawning the stale handle.
func (s *Server) restartSessionLocked(ctx context.Context, sessionID string, session *playback.TranscodeSession, seekSeconds float64, startSegment int) error {
unlock := s.lockSessionLifecycle(sessionID)
defer unlock()
s.mu.RLock()
live, ok := s.sessions[sessionID]
s.mu.RUnlock()
if !ok || live != session {
return playback.ErrSessionSuperseded
}
return session.Restart(ctx, seekSeconds, startSegment)
}
// NewServer creates a new transcode server.
func NewServer(watcher *nodeconfig.Watcher, tracker *nodesessions.Tracker) *Server {
s := &Server{
watcher: watcher,
tracker: tracker,
sessions: make(map[string]*playback.TranscodeSession),
lastAccess: make(map[string]time.Time),
}
if cfg := watcher.Config(); cfg != nil {
if cleaned, err := playback.CleanupOrphanedTranscodeDirs(cfg.Playback.TranscodeDir, nil, 0); err != nil {
slog.Warn("transcode node cleanup failed", "dir", cfg.Playback.TranscodeDir, "error", err)
} else if cleaned > 0 {
slog.Info("transcode node cleanup removed orphaned dirs", "dir", cfg.Playback.TranscodeDir, "count", cleaned)
}
}
return s
}
func (s *Server) SetFFmpegLogSink(sink playback.FFmpegLogSink) {
s.ffmpegSink = sink
}
// noteSessionAccessLocked records an access for a registered job. Callers must
// hold s.mu for writing. Lazily allocates so directly-constructed test servers
// work.
func (s *Server) noteSessionAccessLocked(sessionID string) {
if s.lastAccess == nil {
s.lastAccess = make(map[string]time.Time)
}
s.lastAccess[sessionID] = time.Now()
}
// touchSession refreshes a registered job's idle clock so the reaper spares
// it. Unknown ids are ignored — a reconstruct records its own first access
// when it registers the rebuilt job.
func (s *Server) touchSession(sessionID string) {
s.mu.Lock()
if _, ok := s.sessions[sessionID]; ok {
s.noteSessionAccessLocked(sessionID)
}
s.mu.Unlock()
}
// acquireSessionTouched returns the registered job for sessionID and, when
// found, refreshes its idle clock in the same critical section. Doing both
// under one lock closes the gap where the idle reaper could unregister the
// job between a read-lock lookup and a separate touch, leaving the request
// serving a session whose teardown is already removing its output dir.
func (s *Server) acquireSessionTouched(sessionID string) (*playback.TranscodeSession, bool) {
s.mu.Lock()
defer s.mu.Unlock()
session, ok := s.sessions[sessionID]
if ok {
s.noteSessionAccessLocked(sessionID)
}
return session, ok
}
// startIdleReaper launches the background sweep that closes jobs no client has
// touched for sessionIdleTTL. Called once when the node starts serving;
// subsequent calls are no-ops. The goroutine runs for the process lifetime,
// matching the node's own.
func (s *Server) startIdleReaper() {
s.reaperOnce.Do(func() {
go func() {
ticker := time.NewTicker(sessionReapInterval)
defer ticker.Stop()
for range ticker.C {
s.reapIdleSessions(sessionIdleTTL)
}
}()
})
}
// reapIdleSessions closes and unregisters every job whose last manifest or
// segment access is older than ttl. Registration counts as the first access,
// so a job still waiting on its manifest (the RequireReady flow) is never
// reaped mid-wait. Candidates are collected under the map lock, then each is
// re-validated and torn down under its per-session lifecycle lock: Close
// removes the output dir, and without that lock it could race a token
// reconstruct and wipe the segments a fresh ffmpeg is writing. The session's
// recipe is deliberately kept — an idle reap is not a client stop, and a
// still-valid token must be able to reconstruct on the next hit.
func (s *Server) reapIdleSessions(ttl time.Duration) {
cutoff := time.Now().Add(-ttl)
type idleJob struct {
id string
session *playback.TranscodeSession
}
var candidates []idleJob
s.mu.Lock()
for id, session := range s.sessions {
last, ok := s.lastAccess[id]
if !ok {
// Untracked registration (shouldn't happen): start its idle clock
// now rather than closing a job that may be actively serving.
s.noteSessionAccessLocked(id)
continue
}
if last.Before(cutoff) {
candidates = append(candidates, idleJob{id: id, session: session})
}
}
s.mu.Unlock()
for _, c := range candidates {
s.reapSession(c.id, c.session, cutoff)
}
}
// reapSession tears down one idle job under the per-session lifecycle lock so
// its Close can never overlap a start or reconstruct spawning a new ffmpeg
// into the same output dir. Ownership and idleness are re-checked under the
// lock; a job touched or replaced since the sweep scan is spared.
func (s *Server) reapSession(sessionID string, session *playback.TranscodeSession, cutoff time.Time) {
unlock := s.lockSessionLifecycle(sessionID)
defer unlock()
s.mu.Lock()
live, ok := s.sessions[sessionID]
if !ok || live != session {
s.mu.Unlock()
return
}
last, tracked := s.lastAccess[sessionID]
if tracked && !last.Before(cutoff) {
s.mu.Unlock()
return
}
delete(s.sessions, sessionID)
delete(s.lastAccess, sessionID)
s.mu.Unlock()
s.activeJobs.Add(-1)
if err := session.Close(); err != nil {
slog.Error("close idle transcode session", "component", "transcodenode", "error", err, "session", sessionID, "playback_session_id", sessionID)
}
if s.tracker != nil {
s.tracker.Remove(context.Background(), sessionID)
}
slog.Info("transcode node reaped idle session", "component", "transcodenode",
"session", sessionID, "playback_session_id", sessionID, "idle_ms", time.Since(last).Milliseconds())
}
// recipeStore reads a remote transcode's reconstruction recipe written by central
// at transcode start. The jellycompat node-hop token is identity-only by design —
// not because a Jellyfin client can't round-trip it, but because the recipe is
// mutated in place and the client can't be driven to refresh a stale token, so the
// authoritative recipe lives server-side (see internal/noderecipe). On a node
// restart the node fetches it here instead of 404ing. *noderecipe.Store implements it.
type recipeStore interface {
Get(ctx context.Context, sessionID string) (*playback.RecipeCard, bool)
// Delete drops a session's recipe so a buffered/retrying request after a node
// restart cannot reconstruct a brand-new ffmpeg for an already-stopped session.
// Called only on deliberate teardown; nil-safe and a missing key is a no-op.
Delete(ctx context.Context, sessionID string) error
}
// SetRecipeStore wires the control-plane recipe store so this node can rebuild a
// jellycompat transcode after its own restart. Optional; without it a recipe-less
// (jellycompat) token cannot reconstruct and the request 404s as before.
func (s *Server) SetRecipeStore(store recipeStore) {
s.recipeStore = store
}
// Handler returns the chi.Router with all transcode routes.
func (s *Server) Handler() http.Handler {
s.startIdleReaper()
r := chi.NewRouter()
r.Get("/api/v1/health", s.handleHealth)
r.Group(func(r chi.Router) {
r.Use(s.requireBearer)
r.Get("/hw-capabilities", s.handleHWCapabilities)
r.Post("/chapter-thumbnails/extract", s.handleChapterThumbnailExtract)
r.Post("/transcode/start", s.handleStart)
r.Delete("/transcode/{session_id}", s.handleStop)
r.Get("/transcode/{session_id}/master.m3u8", s.handleManifest)
r.Get("/transcode/{session_id}/segment/{name}", s.handleSegment)
r.Post("/admin/force-reload", s.handleForceReload)
r.Get("/status", s.handleStatus)
})
return r
}
func (s *Server) handleHealth(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(HealthResponse{
Status: "ok",
ActiveJobs: s.activeJobs.Load(),
})
}
func (s *Server) handleHWCapabilities(w http.ResponseWriter, r *http.Request) {
ffmpegPath := ""
if cfg := s.watcher.Config(); cfg != nil {
ffmpegPath = cfg.Playback.FFmpegPath
}
info := playback.DetectHWAccelWithFFmpeg(ffmpegPath)
info.Transformations = playback.ProbeTransformationRegistryV3(r.Context(), ffmpegPath).Advertised()
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(info)
}
func (s *Server) handleChapterThumbnailExtract(w http.ResponseWriter, r *http.Request) {
var req chapterthumbs.RemoteExtractRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeChapterThumbnailError(w, http.StatusBadRequest, "invalid_request", "invalid request body")
return
}
if strings.TrimSpace(req.InputPath) == "" {
writeChapterThumbnailError(w, http.StatusBadRequest, "invalid_request", "input_path is required")
return
}
cfg := s.watcher.Config()
frame, reason, err := chapterthumbs.ExtractFrame(r.Context(), chapterthumbs.FrameExtractOptions{
InputPath: req.InputPath,
SeekSeconds: req.SeekSeconds,
FFmpegPath: cfg.Playback.FFmpegPath,
HWAccel: cfg.Playback.HWAccel,
HWDevice: cfg.Playback.HWDevice,
ToneMap: req.ToneMap,
})
if err != nil {
writeChapterThumbnailError(w, http.StatusUnprocessableEntity, reason, err.Error())
return
}
w.Header().Set("Content-Type", "image/jpeg")
w.WriteHeader(http.StatusOK)
_, _ = w.Write(frame)
}
func writeChapterThumbnailError(w http.ResponseWriter, status int, reason string, message string) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(chapterthumbs.RemoteExtractErrorResponse{
Reason: reason,
Error: message,
})
}
// requireBearer is middleware that checks for Authorization: Bearer {secret}.
func (s *Server) requireBearer(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
cfg := s.watcher.Config()
auth := r.Header.Get("Authorization")
if !strings.HasPrefix(auth, "Bearer ") || strings.TrimPrefix(auth, "Bearer ") != cfg.Auth.JWTSecret {
http.Error(w, "unauthorized", http.StatusUnauthorized)
return
}
next.ServeHTTP(w, r)
})
}
func (s *Server) handleStart(w http.ResponseWriter, r *http.Request) {
var req TranscodeStartRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "invalid request body", http.StatusBadRequest)
return
}
if req.SessionID == "" || req.InputPath == "" {
http.Error(w, "session_id and input_path are required", http.StatusBadRequest)
return
}
cfg := s.watcher.Config()
outputDir := filepath.Join(cfg.Playback.TranscodeDir, req.SessionID)
opts := playback.TranscodeOpts{
InputPath: req.InputPath,
OutputDir: outputDir,
SessionID: req.SessionID,
SourceVideoCodec: req.SourceVideoCodec,
VideoBitstreamFilter: req.VideoBitstreamFilter,
SeekSeconds: req.SeekSeconds,
StartSegmentNumber: req.StartSegmentNumber,
TargetResolution: req.TargetResolution,
TargetCodecVideo: req.TargetCodecVideo,
TargetCodecAudio: req.TargetCodecAudio,
TargetBitrateKbps: req.TargetBitrateKbps,
SegmentDuration: req.SegmentDuration,
FFmpegPath: cfg.Playback.FFmpegPath,
HWAccel: req.HWAccel,
HWDevice: "",
AudioTrackIndex: req.AudioTrackIndex,
SubtitleTrackIndex: req.SubtitleTrackIndex,
SubtitleBurnIn: req.SubtitleBurnIn,
SubtitleCodec: req.SubtitleCodec,
TotalDuration: req.TotalDuration,
FastStart: true,
NodeType: "transcode",
ExecutionMode: "transcode_node",
FFmpegLogSink: s.ffmpegSink,
}
if opts.HWAccel == "" && cfg.Playback.HWAccel != "" {
opts.HWAccel = cfg.Playback.HWAccel
}
// Hold the per-session lifecycle lock across teardown → spawn → register so a
// concurrent reconstruct cannot run a second ffmpeg writer against this
// session's output dir while we replace it.
unlock := s.lockSessionLifecycle(req.SessionID)
// Defensively close any existing session for this ID so that a quality
// switch doesn't orphan the old ffmpeg process or leave stale segments.
s.mu.Lock()
if old, ok := s.sessions[req.SessionID]; ok {
delete(s.sessions, req.SessionID)
delete(s.lastAccess, req.SessionID)
s.mu.Unlock()
s.activeJobs.Add(-1)
_ = old.Close()
// Move the old segment directory aside and delete it in the
// background: removing a long session's segments can take seconds
// on slow disks, and the playback start that triggered this switch
// is blocked waiting for our 202.
staleDir := outputDir + ".stale-" + strconv.FormatInt(time.Now().UnixNano(), 10)
if err := os.Rename(outputDir, staleDir); err == nil {
go func() { _ = os.RemoveAll(staleDir) }()
} else {
os.RemoveAll(outputDir)
}
} else {
s.mu.Unlock()
}
session, err := playback.StartTranscode(context.WithoutCancel(r.Context()), opts)
if err != nil {
unlock()
slog.ErrorContext(r.Context(), "start transcode", "component", "transcodenode", "error", err, "session", req.SessionID, "playback_session_id", req.SessionID)
http.Error(w, "failed to start transcode", http.StatusInternalServerError)
return
}
if req.RequireReady {
if _, err := session.WaitForManifest(8 * time.Second); err != nil {
_ = session.Close()
unlock()
slog.ErrorContext(r.Context(), "transcode failed readiness check", "component", "transcodenode", "error", err, "session", req.SessionID, "playback_session_id", req.SessionID)
http.Error(w, "transcode did not become ready", http.StatusInternalServerError)
return
}
}
s.mu.Lock()
s.sessions[req.SessionID] = session
s.noteSessionAccessLocked(req.SessionID)
s.mu.Unlock()
unlock()
s.activeJobs.Add(1)
// Track session in Redis off the request path — the API server (and
// behind it the playback client) is blocked on this 202, and the
// tracking write is monitoring-only.
effectiveHWAccel := session.Opts().HWAccel
trackCtx := context.WithoutCancel(r.Context())
go s.tracker.Track(trackCtx, nodesessions.SessionInfo{
SessionID: req.SessionID,
NodeURL: s.tracker.NodeURL(),
NodeName: s.tracker.NodeName(),
Type: "transcode",
CodecVideo: req.TargetCodecVideo,
CodecAudio: req.TargetCodecAudio,
Resolution: req.TargetResolution,
HWAccel: effectiveHWAccel,
StartedAt: time.Now().UTC().Format(time.RFC3339),
})
w.WriteHeader(http.StatusAccepted)
json.NewEncoder(w).Encode(TranscodeStartResponse{
SessionID: req.SessionID,
Status: "started",
HWAccel: effectiveHWAccel,
})
}
// reconstructFromToken rebuilds a transcode session this node lost to its own
// restart. The proxy forwards the client's verified stream token in the
// X-Silo-Stream-Token header; the token carries the full byte-affecting recipe
// (the former Postgres "recipe card"), so the node can re-spawn ffmpeg seeked to
// the requested segment rather than 404ing — mirroring the integrated server's
// token-carried reconstruct. Returns nil when the request carries no usable
// transcode token, which the caller renders as a genuine not-found.
//
// requestedSegment is the segment the client is fetching, or negative on the
// manifest path. Reconstruction is single-flighted per session id so concurrent
// manifest and segment requests for the same lost session share one ffmpeg.
func (s *Server) reconstructFromToken(r *http.Request, sessionID string, requestedSegment int) *playback.TranscodeSession {
tokenStr := r.Header.Get("X-Silo-Stream-Token")
if tokenStr == "" {
return nil
}
cfg := s.watcher.Config()
claims, err := streamtoken.Verify(tokenStr, cfg.Auth.JWTSecret)
if err != nil {
slog.WarnContext(r.Context(), "transcode node reconstruct: invalid stream token", "component", "transcodenode", "error", err,
"session", sessionID, "playback_session_id", sessionID)
return nil
}
card := playback.RecipeCardFromClaims(claims)
// The token's recipe must be a transcode card for the session id in the URL: a
// mismatch is a forged or stale request, and direct/remux cards carry no encode
// parameters to rebuild. An empty PlayMethod is a transcode card (back-compat).
expectedTransportID := card.SessionID
if card.TranscodeTransportID != "" {
expectedTransportID = card.TranscodeTransportID
}
if expectedTransportID != sessionID || (card.PlayMethod != "" && card.PlayMethod != playback.PlayTranscode) {
return nil
}
// A native token carries the full byte-affecting recipe. The jellycompat node
// hop signs an identity-only token by design (see internal/noderecipe for why),
// so its card decodes with no encode parameters. For the jellycompat case the
// recipe is fetched from the control-plane recipe store below; without that
// store there is nothing to rebuild from, so 404.
tokenComplete := card.SegmentDuration > 0 && card.TargetCodecVideo != ""
if !tokenComplete && s.recipeStore == nil {
return nil
}
v, _, _ := s.reconstructGroup.Do(sessionID, func() (interface{}, error) {
// A concurrent reconstruct (or a fresh start) may already have registered the
// session; serve it rather than spawning a duplicate ffmpeg.
s.mu.RLock()
existing, ok := s.sessions[sessionID]
s.mu.RUnlock()
if ok {
return existing, nil
}
resolved := card
if !tokenComplete {
// Recipe-less (jellycompat) token: fetch the recipe central wrote to the
// control-plane store at transcode start. A miss / incomplete recipe is a
// genuine not-found (404), never a spawn from a bad recipe.
fetched, ok := s.recipeStore.Get(r.Context(), sessionID)
if !ok || fetched == nil || fetched.SessionID != sessionID ||
fetched.SegmentDuration <= 0 || fetched.TargetCodecVideo == "" {
return (*playback.TranscodeSession)(nil), nil
}
resolved = *fetched
}
return s.spawnReconstruct(r, sessionID, requestedSegment, resolved), nil
})
if session, _ := v.(*playback.TranscodeSession); session != nil {
return session
}
return nil
}
// spawnReconstruct re-spawns ffmpeg for a lost session from its recipe card and
// registers it in the live map. It is only ever called inside the per-session
// single-flight in reconstructFromToken, so it is the sole writer racing to
// register sessionID. Returns nil if the spawn fails or the slot wait is canceled.
func (s *Server) spawnReconstruct(r *http.Request, sessionID string, requestedSegment int, card playback.RecipeCard) *playback.TranscodeSession {
// Pace the cold-start burst so a node restart that loses many sessions does not
// launch every ffmpeg at once. A client that disconnects while waiting releases
// its slot rather than queueing dead work.
release, ok := s.acquireReconstructSlot(r.Context())
if !ok {
return nil
}
defer release()
// Serialize against a concurrent fresh /transcode/start for this session so the
// two never run ffmpeg writers against the same dir. Re-check under the lock and
// yield to any live session rather than spawning a duplicate.
unlock := s.lockSessionLifecycle(sessionID)
defer unlock()
s.mu.RLock()
existing, ok := s.sessions[sessionID]
s.mu.RUnlock()
if ok {
return existing
}
cfg := s.watcher.Config()
outputDir := filepath.Join(cfg.Playback.TranscodeDir, sessionID)
opts := card.TranscodeOpts(outputDir, cfg.Playback.FFmpegPath, s.ffmpegSink)
opts.SessionID = sessionID
// Re-resolve environment-specific encode knobs from this node's live config; the
// token deliberately omits HWAccel/HWDevice so an operator change applies on
// rebuild. Run as a transcode node, not integrated (card.TranscodeOpts defaults).
opts.HWAccel = cfg.Playback.HWAccel
opts.HWDevice = cfg.Playback.HWDevice
opts.NodeType = "transcode"
opts.ExecutionMode = "transcode_node"
// Resume near the segment the client is actually requesting. The card records
// the original start; if the client has played past it, spawning at the old
// position forces a wait-then-seek stall. A negative requestedSegment (manifest
// path) carries no segment context, so the card position stands.
//
// The fast seg×dur mapping is only valid for ENCODED transcodes, whose forced
// keyframes make every segment exactly SegmentDuration long. Copy-mode segments
// have variable durations, so seg×dur points at the wrong source time and causes
// multi-second A/V desync after a restart. For copy-mode cards leave the card's
// original start untouched and let the segment-recovery machinery seek forward
// once the manifest is rebuilt. This mirrors doReconstructTranscode in
// internal/playback/transcode_manager.go so both reconstruct paths stay consistent.
if requestedSegment > card.StartSegmentNumber && card.SegmentDuration > 0 &&
!strings.EqualFold(card.TargetCodecVideo, "copy") {
opts.StartSegmentNumber = requestedSegment
opts.SeekSeconds = float64(requestedSegment * card.SegmentDuration)
}
session, err := playback.StartTranscode(context.WithoutCancel(r.Context()), opts)
if err != nil {
slog.ErrorContext(r.Context(), "transcode node reconstruct start failed", "component", "transcodenode", "error", err,
"session", sessionID, "playback_session_id", sessionID)
return nil
}
// Yield to a winner registered by another path; close only the duplicate ffmpeg,
// never the shared output directory the winner is actively serving.
s.mu.Lock()
if existing, ok := s.sessions[sessionID]; ok {
s.mu.Unlock()
_ = session.CloseProcess()
return existing
}
s.sessions[sessionID] = session
s.noteSessionAccessLocked(sessionID)
s.mu.Unlock()
s.activeJobs.Add(1)
trackCtx := context.WithoutCancel(r.Context())
go s.tracker.Track(trackCtx, nodesessions.SessionInfo{
SessionID: sessionID,
NodeURL: s.tracker.NodeURL(),
NodeName: s.tracker.NodeName(),
Type: "transcode",
CodecVideo: card.TargetCodecVideo,
CodecAudio: card.TargetCodecAudio,
Resolution: card.TargetResolution,
HWAccel: session.Opts().HWAccel,
StartedAt: time.Now().UTC().Format(time.RFC3339),
AuthUserID: card.UserID,
ProfileID: card.ProfileID,
MediaFileID: card.MediaFileID,
})
slog.InfoContext(r.Context(), "transcode node session reconstructed from token", "component", "transcodenode",
"session", sessionID, "playback_session_id", sessionID,
"requested_segment", requestedSegment, "start_segment_number", opts.StartSegmentNumber)
return session
}
// acquireReconstructSlot blocks until a reconstruct slot is free or the request
// context is canceled, returning a release func and true on success. The semaphore
// is lazily sized to NumCPU so a node restart paces its ffmpeg cold starts.
func (s *Server) acquireReconstructSlot(ctx context.Context) (func(), bool) {
s.reconstructSemOnce.Do(func() {
n := runtime.NumCPU()
if n < 1 {
n = 4
}
s.reconstructSem = make(chan struct{}, n)
})
select {
case s.reconstructSem <- struct{}{}:
return func() { <-s.reconstructSem }, true
case <-ctx.Done():
return nil, false
}
}
func (s *Server) handleStop(w http.ResponseWriter, r *http.Request) {
sessionID := chi.URLParam(r, "session_id")
// Serialize against in-flight starts and reconstructs. A RequireReady
// start registers its job only after the readiness wait; a stop racing
// that wait (the API rolling back a timed-out start) would otherwise miss
// the map and 404, orphaning the just-spawned ffmpeg until the idle
// reaper finds it. Blocking here until the start registers turns that
// miss into a normal teardown.
unlock := s.lockSessionLifecycle(sessionID)
defer unlock()
s.mu.Lock()
session, ok := s.sessions[sessionID]
if !ok {
s.mu.Unlock()
http.Error(w, "session not found", http.StatusNotFound)
return
}
delete(s.sessions, sessionID)
delete(s.lastAccess, sessionID)
s.mu.Unlock()
s.activeJobs.Add(-1)
if err := session.Close(); err != nil {
slog.ErrorContext(r.Context(), "close transcode session", "component", "transcodenode", "error", err, "session", sessionID, "playback_session_id", sessionID)
}
cfg := s.watcher.Config()
outputDir := filepath.Join(cfg.Playback.TranscodeDir, sessionID)
os.RemoveAll(outputDir)
// Drop the recipe so a buffered/retrying request after a node restart cannot
// reconstruct a new ffmpeg for this now-stopped session. Best-effort: a stop
// must still succeed even if the recipe store is briefly unavailable.
if s.recipeStore != nil {
if err := s.recipeStore.Delete(r.Context(), sessionID); err != nil {
slog.WarnContext(r.Context(), "delete transcode recipe on stop", "component", "transcodenode", "error", err, "session", sessionID, "playback_session_id", sessionID)
}
}
s.tracker.Remove(r.Context(), sessionID)
w.WriteHeader(http.StatusNoContent)
}
func (s *Server) handleManifest(w http.ResponseWriter, r *http.Request) {
sessionID := chi.URLParam(r, "session_id")
// Lookup and liveness refresh happen atomically so the idle reaper can
// never unregister the job between them and tear down a session this
// request is about to serve from.
session, ok := s.acquireSessionTouched(sessionID)
if !ok {
// Lost the in-memory session (this node restarted): rebuild it from the
// stream token the proxy forwarded. The manifest path carries no segment
// context, so reconstruct at the recipe's original start position.
session = s.reconstructFromToken(r, sessionID, -1)
if session == nil {
http.Error(w, "session not found", http.StatusNotFound)
return
}
// A reconstruct that yielded to a concurrently registered winner has
// not recorded this hit; count it so the reaper sees the liveness.
s.touchSession(sessionID)
}
manifest, err := session.BuildPlaybackManifest("segment/", r.URL.RawQuery)
if err != nil {
slog.ErrorContext(r.Context(), "get manifest", "component", "transcodenode", "error", err, "session", sessionID, "playback_session_id", sessionID)
http.Error(w, "manifest not ready", http.StatusServiceUnavailable)
return
}
w.Header().Set("Content-Type", "application/vnd.apple.mpegurl")
w.Header().Set("Cache-Control", "no-store, max-age=0")
w.Header().Set("Pragma", "no-cache")
w.Write(manifest)
}
func (s *Server) handleSegment(w http.ResponseWriter, r *http.Request) {
sessionID := chi.URLParam(r, "session_id")
name := chi.URLParam(r, "name")
// Lookup and liveness refresh happen atomically so the idle reaper can
// never unregister the job between them and tear down a session this
// request is about to serve from.
session, ok := s.acquireSessionTouched(sessionID)
if !ok {
// Lost the in-memory session (this node restarted): rebuild it from the
// forwarded stream token, seeked to the segment the client is requesting so
// playback resumes near its position instead of restarting from the start.
requestedSegment := -1
if n, parseErr := playback.ParseSegmentNumber(name); parseErr == nil {
requestedSegment = n
}
session = s.reconstructFromToken(r, sessionID, requestedSegment)
if session == nil {
http.Error(w, "session not found", http.StatusNotFound)
return
}
// A reconstruct that yielded to a concurrently registered winner has
// not recorded this hit; count it so the reaper sees the liveness.
s.touchSession(sessionID)
}
segPath, err := session.GetSegment(name)
if err != nil && err == playback.ErrSegmentNotFound {
segNum, parseErr := playback.ParseSegmentNumber(name)
if parseErr == nil {
now := time.Now()
decision := session.SegmentRecoveryDecision(segNum, now)
lastProducedAgeMS := int64(-1)
if !decision.Progress.LastProducedAt.IsZero() {
lastProducedAgeMS = now.Sub(decision.Progress.LastProducedAt).Milliseconds()
}
slog.InfoContext(r.Context(), "transcode segment missing", "component", "transcodenode",
"segment", name,
"requested_segment", segNum,
"produced_head", decision.Progress.ProducedHead,
"last_requested_segment", decision.Progress.LastRequestedSegment,
"start_segment_number", decision.Progress.StartSegmentNumber,
"last_produced_age_ms", lastProducedAgeMS,
"wait_timeout_ms", decision.WaitTimeout.Milliseconds(),
"reason", decision.Reason,
"session", sessionID,
"playback_session_id", sessionID,
)
if decision.Wait {
slog.InfoContext(r.Context(), "transcode segment wait", "component", "transcodenode",
"segment", name,
"requested_segment", segNum,
"produced_head", decision.Progress.ProducedHead,
"last_requested_segment", decision.Progress.LastRequestedSegment,
"start_segment_number", decision.Progress.StartSegmentNumber,
"last_produced_age_ms", lastProducedAgeMS,
"wait_timeout_ms", decision.WaitTimeout.Milliseconds(),
"reason", decision.Reason,
"session", sessionID,
"playback_session_id", sessionID,
)
segPath, err = session.WaitForSegment(name, decision.WaitTimeout)
if err != nil && err == playback.ErrSegmentNotFound {
slog.InfoContext(r.Context(), "transcode segment wait timeout", "component", "transcodenode",
"segment", name,
"requested_segment", segNum,
"produced_head", decision.Progress.ProducedHead,
"last_requested_segment", decision.Progress.LastRequestedSegment,
"start_segment_number", decision.Progress.StartSegmentNumber,
"last_produced_age_ms", lastProducedAgeMS,
"wait_timeout_ms", decision.WaitTimeout.Milliseconds(),
"reason", decision.Reason,
"session", sessionID,
"playback_session_id", sessionID,
)
}
}
if err != nil && err == playback.ErrSegmentNotFound && decision.RestartOnTimeout {
seekSeconds, ok, seekErr := session.RestartSeekTarget(segNum)
if seekErr != nil && !errors.Is(seekErr, playback.ErrManifestNotReady) {
slog.ErrorContext(r.Context(), "resolve transcode node seek target", "component", "transcodenode", "error", seekErr, "segment", name, "session", sessionID, "playback_session_id", sessionID)
}
if ok {
slog.InfoContext(r.Context(), "transcode node seek restart", "component", "transcodenode",
"segment", name,
"requested_segment", segNum,
"produced_head", decision.Progress.ProducedHead,
"last_requested_segment", decision.Progress.LastRequestedSegment,
"start_segment_number", decision.Progress.StartSegmentNumber,
"last_produced_age_ms", lastProducedAgeMS,
"wait_timeout_ms", decision.WaitTimeout.Milliseconds(),
"reason", decision.Reason,
"seek_seconds", seekSeconds,
"session", sessionID,
"playback_session_id", sessionID,
)
if restartErr := s.restartSessionLocked(
context.WithoutCancel(r.Context()),
sessionID,
session,
seekSeconds,
segNum,
); restartErr == nil {
segPath, err = session.WaitForSegment(name, 30*time.Second)
}
}
if !ok && session.IsCopyVideo() {
err = playback.ErrSegmentNotFound
}
}
} else if session.IsRunning() {
// Non-numbered segment (e.g., init.mp4 for fMP4 HLS).
// Wait briefly — the init segment is written almost immediately.
segPath, err = session.WaitForSegment(name, 10*time.Second)
}
}
if err != nil {
http.Error(w, "segment not found", http.StatusNotFound)
return
}
w.Header().Set("Cache-Control", "no-store, max-age=0")
w.Header().Set("Pragma", "no-cache")
http.ServeFile(w, r, segPath)
}
func (s *Server) handleForceReload(w http.ResponseWriter, r *http.Request) {
if err := s.watcher.ForceReload(r.Context()); err != nil {
http.Error(w, "reload failed: "+err.Error(), http.StatusInternalServerError)
return
}
cfg := s.watcher.Config()
s.mu.Lock()
stopped := make([]string, 0, len(s.sessions))
for id, session := range s.sessions {
session.Close()
os.RemoveAll(filepath.Join(cfg.Playback.TranscodeDir, id))
delete(s.sessions, id)
delete(s.lastAccess, id)
stopped = append(stopped, id)
}
s.activeJobs.Store(0)
s.mu.Unlock()
// A force-reload tears every session down for good, so drop their recipes too:
// otherwise a buffered/retrying request could reconstruct a session this reload
// deliberately killed. Best-effort, done outside the map lock.
if s.recipeStore != nil {
for _, id := range stopped {
if err := s.recipeStore.Delete(r.Context(), id); err != nil {
slog.WarnContext(r.Context(), "delete transcode recipe on force reload", "component", "transcodenode", "error", err, "session", id, "playback_session_id", id)
}
}
}
s.tracker.Cleanup(r.Context())
slog.InfoContext(r.Context(), "transcode force reload completed", slog.String("component", "transcodenode"))
w.WriteHeader(http.StatusNoContent)
}
func (s *Server) handleStatus(w http.ResponseWriter, r *http.Request) {
s.mu.RLock()
sessionIDs := make([]string, 0, len(s.sessions))
for id := range s.sessions {
sessionIDs = append(sessionIDs, id)
}
s.mu.RUnlock()
w.Header().Set("Content-Type", "application/json")
type statusResponse struct {
Status string `json:"status"`
ActiveJobs int32 `json:"active_jobs"`
Sessions []string `json:"sessions"`
}
json.NewEncoder(w).Encode(statusResponse{
Status: "ok",
ActiveJobs: s.activeJobs.Load(),
Sessions: sessionIDs,
})
}