Files
QuickandGitHub 47e45f4d67 feat(downloads): relay node-local artifacts through proxies (#608)
* feat(downloads): relay node-local artifacts through proxies

* fix(downloads): harden remote artifact relay

* fix(downloads): close distributed artifact review gaps

* fix(downloads): bound relays and refresh origins

* fix(downloads): recover stalled remote artifacts

* fix(transcode): report session cleanup failures

* fix(downloads): harden artifact cleanup recovery

* fix(downloads): close remote artifact lifecycle races

* fix(downloads): honor configured artifact storage

* fix(downloads): recover proxy-observed artifact misses

* test(downloads): isolate remote cleanup assertions

* fix(downloads): bound remote artifact recovery

* fix(downloads): allow concurrent artifact relays

* fix(downloads): harden distributed artifact recovery
2026-08-12 15:14:08 -04:00

1354 lines
52 KiB
Go
Raw Permalink 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/config"
"github.com/Silo-Server/silo-server/internal/downloadprepare"
"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"`
SourceVideoProfile string `json:"source_video_profile,omitempty"`
SourceVideoBitDepth int `json:"source_video_bit_depth,omitempty"`
SoftwareVideoDecode bool `json:"software_video_decode,omitempty"`
VideoBitstreamFilter string `json:"video_bitstream_filter,omitempty"`
SeekSeconds float64 `json:"seek_seconds"`
StreamOriginSeconds float64 `json:"stream_origin_seconds,omitempty"`
CopySeekAnchorResolved bool `json:"copy_seek_anchor_resolved,omitempty"`
StartSegmentNumber int `json:"start_segment_number"`
TargetResolution string `json:"target_resolution"`
TargetCodecVideo string `json:"target_codec_video"`
TargetCodecAudio string `json:"target_codec_audio"`
TargetAudioChannels int `json:"target_audio_channels,omitempty"`
TargetAudioBitrateKbps int `json:"target_audio_bitrate_kbps,omitempty"`
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
// Session tracking is monitoring-only and must never make a healthy node's
// control-plane response depend on Redis latency.
const sessionTrackingOperationTimeout = 2 * time.Second
type sessionTracker interface {
Track(context.Context, nodesessions.SessionInfo)
Remove(context.Context, string)
Cleanup(context.Context)
NodeURL() string
NodeName() string
}
// Server is the HTTP handler for transcode mode.
type Server struct {
watcher *nodeconfig.Watcher
tracker sessionTracker
ffmpegSink playback.FFmpegLogSink
inputPaths InputPathAuthorizer
transcodeDir string
artifactRoot string
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
// reloadMu keeps force-reload teardown atomic with session creation and
// reconstruction. It is always acquired before lifecycleMu or mu.
reloadMu 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 locks 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. Artifact readers use the shared side
// so concurrent relays remain independent while prepare/delete stay exclusive.
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 lock; 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.RWMutex
refs int
}
func (s *Server) retainSessionLifecycleLock(sessionID string) *sessionLifecycleLock {
s.lifecycleMu.Lock()
defer s.lifecycleMu.Unlock()
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++
return lk
}
func (s *Server) releaseSessionLifecycleLock(sessionID string, lk *sessionLifecycleLock) {
s.lifecycleMu.Lock()
defer s.lifecycleMu.Unlock()
lk.refs--
if lk.refs == 0 {
delete(s.lifecycleLocks, sessionID)
}
}
// 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() {
lk := s.retainSessionLifecycleLock(sessionID)
lk.mu.Lock()
return func() {
lk.mu.Unlock()
s.releaseSessionLifecycleLock(sessionID, lk)
}
}
// lockSessionLifecycleRead holds the shared side of a lifecycle lock. Artifact
// relays can therefore proceed concurrently, while preparation and deletion
// remain exclusive for the full transfer.
func (s *Server) lockSessionLifecycleRead(sessionID string) func() {
lk := s.retainSessionLifecycleLock(sessionID)
lk.mu.RLock()
return func() {
lk.mu.RUnlock()
s.releaseSessionLifecycleLock(sessionID, lk)
}
}
// 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 {
var trackerImpl sessionTracker
if tracker != nil {
trackerImpl = tracker
}
transcodeDir := config.DefaultTranscodeDir
artifactDir := ""
if watcher != nil {
if cfg := watcher.Config(); cfg != nil {
if strings.TrimSpace(cfg.Playback.TranscodeDir) != "" {
transcodeDir = cfg.Playback.TranscodeDir
}
artifactDir = cfg.Download.ArtifactDir
}
}
artifactRoot := filepath.Join(transcodeDir, downloadprepare.ArtifactDirectoryName)
if strings.TrimSpace(artifactDir) != "" {
artifactRoot = config.EffectiveDownloadArtifactDir(artifactDir, transcodeDir)
}
s := &Server{
watcher: watcher,
tracker: trackerImpl,
transcodeDir: transcodeDir,
artifactRoot: artifactRoot,
sessions: make(map[string]*playback.TranscodeSession),
lastAccess: make(map[string]time.Time),
}
return s
}
// StartOrphanSweeper runs the age-guarded orphan-transcode sweep immediately and
// then hourly until ctx is cancelled. It never blocks (a slow network-filesystem
// delete runs in its own goroutine), so it is safe to call before the node binds
// its listener. This is the node's only filesystem-level reclaimer of dirs left
// behind by a session that was dropped without its output dir being removed — the
// idle reaper only deletes dirs it still tracks in s.sessions, so without this
// periodic pass such orphans would linger until the next process restart. The
// MaxTokenTTL age guard keeps a delete from racing a token-carried reconstruct
// writing into TranscodeDir/<sessionID>: a dir younger than the max token
// lifetime may still be reused, while older dirs are never reconstructable.
func (s *Server) StartOrphanSweeper(ctx context.Context) {
dir := s.transcodeDir
playback.StartPeriodicOrphanCleanup(ctx, "transcodenode", dir, func() (int, error) {
// Spare the live registered jobs by id, not by age alone: now that the
// sweep runs during live traffic, a long-lived session that re-serves
// already-written segments stops advancing its dir mtime, so the age
// guard could misclassify it as orphaned. The live set is authoritative
// (in-flight reconstructs are covered by their fresh writes + age guard).
return playback.CleanupOrphanedTranscodeDirs(dir, s.activeSessionIDs(), playback.MaxTokenTTL)
}, playback.OrphanCleanupInterval)
}
// activeSessionIDs snapshots the ids of currently registered jobs so the orphan
// sweep spares their output dirs regardless of directory mtime, mirroring the
// central TranscodeManager's live-set snapshot.
func (s *Server) activeSessionIDs() map[string]struct{} {
s.mu.RLock()
defer s.mu.RUnlock()
active := make(map[string]struct{}, len(s.sessions)+1)
for id := range s.sessions {
active[id] = struct{}{}
}
// Node-local prepared downloads may live under the existing persistent
// transcode volume. They have their own lifecycle and must never be mistaken
// for an orphaned HLS session directory. Protect the top-level container for
// both the default and an explicitly nested artifact directory.
if s.transcodeDir != "" && s.artifactRoot != "" {
if rel, err := filepath.Rel(s.transcodeDir, s.artifactRoot); err == nil && rel != "." && rel != ".." && !strings.HasPrefix(rel, ".."+string(filepath.Separator)) {
if first, _, ok := strings.Cut(rel, string(filepath.Separator)); ok {
active[first] = struct{}{}
} else {
active[rel] = struct{}{}
}
}
}
return active
}
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
}
// SetInputPathAuthorizer wires the library-root authority used by every node
// endpoint that accepts an FFmpeg input path.
func (s *Server) SetInputPathAuthorizer(authorizer InputPathAuthorizer) {
s.inputPaths = authorizer
}
// 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("/downloads/prepare", s.handleDownloadPrepare)
r.Head("/downloads/artifacts/{artifact_id}", s.handleDownloadArtifact)
r.Get("/downloads/artifacts/{artifact_id}", s.handleDownloadArtifact)
r.Delete("/downloads/artifacts/{artifact_id}", s.handleDeleteDownloadArtifact)
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) handleDownloadPrepare(w http.ResponseWriter, r *http.Request) {
var req downloadprepare.Request
if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 64<<10)).Decode(&req); err != nil {
http.Error(w, "invalid request body", http.StatusBadRequest)
return
}
if !downloadprepare.ValidArtifactID(req.ArtifactID) || strings.TrimSpace(req.InputPath) == "" {
http.Error(w, "a valid artifact_id and input_path are required", http.StatusBadRequest)
return
}
cfg := s.watcher.Config()
if cfg == nil {
http.Error(w, "node not configured", http.StatusServiceUnavailable)
return
}
if !s.requireApprovedInputPath(w, r, req.InputPath) {
return
}
artifactRoot := s.artifactRoot
if err := os.MkdirAll(artifactRoot, 0o755); err != nil {
http.Error(w, "artifact directory unavailable", http.StatusInternalServerError)
return
}
outputPath := filepath.Join(artifactRoot, req.ArtifactID+".mp4")
unlock := s.lockSessionLifecycle("download-artifact-" + req.ArtifactID)
defer unlock()
if stat, err := os.Stat(outputPath); err == nil && stat.Mode().IsRegular() && stat.Size() > 0 {
writeDownloadPrepareResult(w, req.ArtifactID, stat.Size())
return
}
jobCtx := r.Context()
s.activeJobs.Add(1)
defer s.activeJobs.Add(-1)
if s.tracker != nil {
finishTracking := s.trackDownloadPrepare(jobCtx, nodesessions.SessionInfo{
SessionID: "download-" + req.ArtifactID,
NodeURL: s.tracker.NodeURL(),
NodeName: s.tracker.NodeName(),
Type: "download_prepare",
CodecVideo: req.TargetCodecVideo,
CodecAudio: req.TargetCodecAudio,
Resolution: req.TargetResolution,
StartedAt: time.Now().UTC().Format(time.RFC3339),
})
defer finishTracking()
}
opts := req.TranscodeOpts(cfg.Playback.FFmpegPath, cfg.Playback.HWAccel, cfg.Playback.HWDevice, s.ffmpegSink)
if err := playback.PrepareFile(jobCtx, opts, outputPath); err != nil {
if jobCtx.Err() == nil {
slog.ErrorContext(jobCtx, "prepare download artifact", "component", "transcodenode", "artifact_id", req.ArtifactID, "error", err)
}
http.Error(w, "failed to prepare download artifact", http.StatusInternalServerError)
return
}
stat, err := os.Stat(outputPath)
if err != nil || !stat.Mode().IsRegular() {
http.Error(w, "prepared download artifact unavailable", http.StatusInternalServerError)
return
}
writeDownloadPrepareResult(w, req.ArtifactID, stat.Size())
}
func writeDownloadPrepareResult(w http.ResponseWriter, artifactID string, fileSize int64) {
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(downloadprepare.Result{ArtifactID: artifactID, FileSize: fileSize})
}
func (s *Server) sessionOutputDir(sessionID string) string {
return filepath.Join(s.transcodeDir, sessionID)
}
func (s *Server) handleDownloadArtifact(w http.ResponseWriter, r *http.Request) {
artifactID := chi.URLParam(r, "artifact_id")
if !downloadprepare.ValidArtifactID(artifactID) {
http.NotFound(w, r)
return
}
cfg := s.watcher.Config()
if cfg == nil {
http.Error(w, "node not configured", http.StatusServiceUnavailable)
return
}
// Share the lock with other readers, but serialize with preparation and
// deletion. A recovery HEAD must not report a definitive 404 while
// PrepareFile is still publishing this artifact.
unlock := s.lockSessionLifecycleRead("download-artifact-" + artifactID)
defer unlock()
path := filepath.Join(s.artifactRoot, artifactID+".mp4")
f, err := os.Open(path)
if err != nil {
if os.IsNotExist(err) {
http.NotFound(w, r)
return
}
http.Error(w, "artifact unavailable", http.StatusInternalServerError)
return
}
defer func() { _ = f.Close() }()
stat, err := f.Stat()
if err != nil || !stat.Mode().IsRegular() {
http.Error(w, "artifact unavailable", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Disposition", `attachment; filename="`+artifactID+`.mp4"`)
w.Header().Set("Content-Type", playback.MimeFromExtension(path))
w.Header().Set("ETag", `"`+artifactID+`-`+strconv.FormatInt(stat.Size(), 10)+`"`)
http.ServeContent(w, r, stat.Name(), stat.ModTime(), f)
}
func (s *Server) handleDeleteDownloadArtifact(w http.ResponseWriter, r *http.Request) {
artifactID := chi.URLParam(r, "artifact_id")
if !downloadprepare.ValidArtifactID(artifactID) {
http.NotFound(w, r)
return
}
cfg := s.watcher.Config()
if cfg == nil {
http.Error(w, "node not configured", http.StatusServiceUnavailable)
return
}
unlock := s.lockSessionLifecycle("download-artifact-" + artifactID)
defer unlock()
path := filepath.Join(s.artifactRoot, artifactID+".mp4")
for _, candidate := range []string{path, path + ".part"} {
if err := os.Remove(candidate); err != nil && !os.IsNotExist(err) {
http.Error(w, "failed to remove artifact", http.StatusInternalServerError)
return
}
}
w.WriteHeader(http.StatusNoContent)
}
// trackDownloadPrepare runs the monitoring lifecycle off the request path.
// One goroutine owns both operations so Remove can never overtake Track, even
// when the encode completes before Redis responds. finish only signals that
// the job ended; it never waits for either bounded Redis operation.
func (s *Server) trackDownloadPrepare(ctx context.Context, info nodesessions.SessionInfo) func() {
finished := make(chan struct{})
var finishOnce sync.Once
baseCtx := context.WithoutCancel(ctx)
go func() {
trackCtx, cancelTrack := context.WithTimeout(baseCtx, sessionTrackingOperationTimeout)
s.tracker.Track(trackCtx, info)
cancelTrack()
<-finished
removeCtx, cancelRemove := context.WithTimeout(baseCtx, sessionTrackingOperationTimeout)
s.tracker.Remove(removeCtx, info.SessionID)
cancelRemove()
}()
return func() {
finishOnce.Do(func() { close(finished) })
}
}
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
}
if !s.requireApprovedInputPath(w, r, req.InputPath) {
return
}
cfg := s.watcher.Config()
if cfg == nil {
writeChapterThumbnailError(w, http.StatusServiceUnavailable, "node_unavailable", "node not configured")
return
}
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,
AllowSoftwareToneMap: req.AllowSoftwareToneMap,
})
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()
if cfg == nil {
http.Error(w, "node not configured", http.StatusServiceUnavailable)
return
}
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
}
if !s.requireApprovedInputPath(w, r, req.InputPath) {
return
}
s.reloadMu.RLock()
defer s.reloadMu.RUnlock()
cfg := s.watcher.Config()
if cfg == nil {
http.Error(w, "node not configured", http.StatusServiceUnavailable)
return
}
outputDir := s.sessionOutputDir(req.SessionID)
opts := playback.TranscodeOpts{
InputPath: req.InputPath,
OutputDir: outputDir,
SessionID: req.SessionID,
SourceVideoCodec: req.SourceVideoCodec,
SourceVideoProfile: req.SourceVideoProfile,
SourceVideoBitDepth: req.SourceVideoBitDepth,
SoftwareVideoDecode: req.SoftwareVideoDecode,
VideoBitstreamFilter: req.VideoBitstreamFilter,
SeekSeconds: req.SeekSeconds,
StreamOriginSeconds: req.StreamOriginSeconds,
CopySeekAnchorResolved: req.CopySeekAnchorResolved,
StartSegmentNumber: req.StartSegmentNumber,
TargetResolution: req.TargetResolution,
TargetCodecVideo: req.TargetCodecVideo,
TargetCodecAudio: req.TargetCodecAudio,
TargetAudioChannels: req.TargetAudioChannels,
TargetAudioBitrateKbps: req.TargetAudioBitrateKbps,
TargetBitrateKbps: req.TargetBitrateKbps,
SegmentDuration: req.SegmentDuration,
FFmpegPath: cfg.Playback.FFmpegPath,
HWAccel: req.HWAccel,
// This node's configured device (or device list — StartTranscode
// resolves it to one GPU), matching what reconstruction uses so fresh
// and reconstructed sessions balance identically.
HWDevice: cfg.Playback.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,
})
}
func (s *Server) requireApprovedInputPath(w http.ResponseWriter, r *http.Request, path string) bool {
if s.inputPaths == nil {
http.Error(w, "input path authority unavailable", http.StatusServiceUnavailable)
return false
}
allowed, err := s.inputPaths.Allowed(r.Context(), path)
if err != nil {
slog.ErrorContext(r.Context(), "authorize transcode input path", "component", "transcodenode", "error", err)
http.Error(w, "input path authority unavailable", http.StatusServiceUnavailable)
return false
}
if !allowed {
http.Error(w, "input_path must be an approved absolute media file", http.StatusBadRequest)
return false
}
return true
}
// 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()
if cfg == nil {
return nil
}
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 {
if s.inputPaths == nil {
slog.ErrorContext(r.Context(), "transcode node reconstruct input authority unavailable", "component", "transcodenode", "session", sessionID)
return nil
}
allowed, err := s.inputPaths.Allowed(r.Context(), card.InputPath)
if err != nil || !allowed {
slog.WarnContext(r.Context(), "transcode node reconstruct input rejected", "component", "transcodenode", "session", sessionID, "error", err)
return nil
}
// 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()
s.reloadMu.RLock()
defer s.reloadMu.RUnlock()
// 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()
if cfg == nil {
return nil
}
outputDir := s.sessionOutputDir(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)
}
if err := os.RemoveAll(s.sessionOutputDir(sessionID)); err != nil {
slog.WarnContext(r.Context(), "remove transcode session directory", "component", "transcodenode", "session", sessionID, "error", err)
}
// 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)
}
var manifest []byte
var err error
if r.URL.Query().Get(playback.SourceTimelineQueryParam) == "1" {
manifest, err = session.BuildSourceAlignedPlaybackManifest("segment/", r.URL.RawQuery)
} else {
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) {
s.reloadMu.Lock()
defer s.reloadMu.Unlock()
if err := s.watcher.ForceReload(r.Context()); err != nil {
http.Error(w, "reload failed: "+err.Error(), http.StatusInternalServerError)
return
}
s.mu.Lock()
stopped := make([]string, 0, len(s.sessions))
for id, session := range s.sessions {
session.Close()
if err := os.RemoveAll(s.sessionOutputDir(id)); err != nil {
slog.WarnContext(r.Context(), "remove transcode session directory during reload", "component", "transcodenode", "session", id, "error", err)
}
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,
})
}