* 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
1354 lines
52 KiB
Go
1354 lines
52 KiB
Go
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,
|
||
})
|
||
}
|