2026-05-22 20:26:11 -04:00
package transcodenode
import (
"context"
"encoding/json"
"errors"
"log/slog"
"net/http"
"os"
"path/filepath"
2026-07-03 02:23:14 +08:00
"runtime"
2026-06-10 17:18:18 -04:00
"strconv"
2026-05-22 20:26:11 -04:00
"strings"
"sync"
"sync/atomic"
"time"
"github.com/go-chi/chi/v5"
2026-07-03 02:23:14 +08:00
"golang.org/x/sync/singleflight"
2026-05-22 20:26:11 -04:00
"github.com/Silo-Server/silo-server/internal/chapterthumbs"
2026-08-08 15:23:16 -04:00
"github.com/Silo-Server/silo-server/internal/httpstream"
2026-05-22 20:26:11 -04:00
"github.com/Silo-Server/silo-server/internal/nodeconfig"
"github.com/Silo-Server/silo-server/internal/nodesessions"
"github.com/Silo-Server/silo-server/internal/playback"
2026-07-03 02:23:14 +08:00
"github.com/Silo-Server/silo-server/internal/streamtoken"
2026-08-11 14:30:55 -04:00
"github.com/Silo-Server/silo-server/internal/transcodeproxy"
2026-05-22 20:26:11 -04:00
)
// TranscodeStartRequest is the JSON body for POST /transcode/start.
type TranscodeStartRequest struct {
2026-07-20 19:25:46 +02:00
SessionID string `json:"session_id"`
InputPath string `json:"input_path"`
SourceVideoCodec string `json:"source_video_codec"`
2026-08-10 18:14:49 -04:00
SourceVideoProfile string `json:"source_video_profile,omitempty"`
SourceVideoBitDepth int `json:"source_video_bit_depth,omitempty"`
SoftwareVideoDecode bool `json:"software_video_decode,omitempty"`
2026-07-20 19:25:46 +02:00
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"`
2026-08-10 18:14:49 -04:00
TargetAudioBitrateKbps int `json:"target_audio_bitrate_kbps,omitempty"`
2026-07-20 19:25:46 +02:00
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"`
2026-05-22 20:26:11 -04:00
}
// TranscodeStartResponse is the JSON response for POST /transcode/start.
type TranscodeStartResponse struct {
SessionID string `json:"session_id"`
Status string `json:"status"`
2026-06-05 19:43:20 -07:00
HWAccel string `json:"hw_accel,omitempty"`
2026-05-22 20:26:11 -04:00
}
// HealthResponse is the JSON response for GET /api/v1/health.
type HealthResponse struct {
Status string `json:"status"`
ActiveJobs int32 `json:"active_jobs"`
}
2026-07-14 11:51:27 -04:00
// 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
2026-05-22 20:26:11 -04:00
// Server is the HTTP handler for transcode mode.
type Server struct {
watcher * nodeconfig . Watcher
tracker * nodesessions . Tracker
ffmpegSink playback . FFmpegLogSink
sessions map [ string ] * playback . TranscodeSession
2026-07-14 11:51:27 -04:00
// 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
2026-05-22 20:26:11 -04:00
mu sync . RWMutex
activeJobs atomic . Int32
2026-07-03 02:23:14 +08:00
// reconstructGroup single-flights node-side session reconstruction per session
// id so a post-restart wave of concurrent manifest/segment requests for the same
// lost session spawns exactly one ffmpeg, never racing duplicates into the shared
// output directory.
reconstructGroup singleflight . Group
// reconstructSem bounds how many sessions may be reconstructed (ffmpeg
// re-spawned) at once after a node restart, pacing the cold-start burst instead
// of stampeding the host. Lazily sized to NumCPU on first use.
reconstructSemOnce sync . Once
reconstructSem chan struct {}
// lifecycleMu guards lifecycleLocks, the per-session mutexes that serialize
// every path which spawns ffmpeg into a session's output dir (fresh start and
// reconstruct). reconstructGroup only single-flights reconstructs against each
// other; without this a reconstruct racing a fresh /transcode/start could run
// two ffmpeg writers against the same dir.
lifecycleMu sync . Mutex
lifecycleLocks map [ string ] * sessionLifecycleLock
// recipeStore is the control-plane recipe store consulted when a forwarded
// token carries no recipe (the jellycompat node hop). Nil disables that path.
recipeStore recipeStore
}
// sessionLifecycleLock is a refcounted per-session mutex; the refcount lets the
// node drop the map entry once no path holds or waits on it so the map stays
// bounded over the node's lifetime.
type sessionLifecycleLock struct {
mu sync . Mutex
refs int
}
// lockSessionLifecycle acquires the per-session lifecycle mutex and returns a
// release func. Held across "check existing → spawn → register" so a fresh start
// and a reconstruct never run concurrent ffmpeg writers for one session's dir.
func ( s * Server ) lockSessionLifecycle ( sessionID string ) func () {
s . lifecycleMu . Lock ()
if s . lifecycleLocks == nil {
s . lifecycleLocks = make ( map [ string ] * sessionLifecycleLock )
}
lk := s . lifecycleLocks [ sessionID ]
if lk == nil {
lk = & sessionLifecycleLock {}
s . lifecycleLocks [ sessionID ] = lk
}
lk . refs ++
s . lifecycleMu . Unlock ()
lk . mu . Lock ()
return func () {
lk . mu . Unlock ()
s . lifecycleMu . Lock ()
lk . refs --
if lk . refs == 0 {
delete ( s . lifecycleLocks , sessionID )
}
s . lifecycleMu . Unlock ()
}
}
// restartSessionLocked re-spawns session under the per-session lifecycle lock so
// a segment-recovery restart can never race a fresh start, reconstruct, or
// another restart into the same output directory. It holds the lock only across
// the cancel→respawn transition inside Restart and releases it before the caller
// waits on segments. Under the lock it confirms session is still the live mapped
// session; a concurrent teardown or reconstruct that replaced it yields
// ErrSessionSuperseded rather than re-spawning the stale handle.
func ( s * Server ) restartSessionLocked ( ctx context . Context , sessionID string , session * playback . TranscodeSession , seekSeconds float64 , startSegment int ) error {
unlock := s . lockSessionLifecycle ( sessionID )
defer unlock ()
s . mu . RLock ()
live , ok := s . sessions [ sessionID ]
s . mu . RUnlock ()
if ! ok || live != session {
return playback . ErrSessionSuperseded
}
return session . Restart ( ctx , seekSeconds , startSegment )
2026-05-22 20:26:11 -04:00
}
// NewServer creates a new transcode server.
func NewServer ( watcher * nodeconfig . Watcher , tracker * nodesessions . Tracker ) * Server {
s := & Server {
2026-07-14 11:51:27 -04:00
watcher : watcher ,
tracker : tracker ,
sessions : make ( map [ string ] * playback . TranscodeSession ),
lastAccess : make ( map [ string ] time . Time ),
2026-05-22 20:26:11 -04:00
}
2026-07-17 02:48:00 +08:00
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 := ""
if cfg := s . watcher . Config (); cfg != nil {
dir = cfg . Playback . TranscodeDir
}
playback . StartPeriodicOrphanCleanup ( ctx , "transcodenode" , dir , func () ( int , error ) {
// Re-read config each run so a hot-reloaded TranscodeDir is honored.
cfg := s . watcher . Config ()
if cfg == nil {
return 0 , nil
2026-05-22 20:26:11 -04:00
}
2026-07-17 02:48:00 +08:00
// 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 ( cfg . Playback . TranscodeDir , 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 ))
for id := range s . sessions {
active [ id ] = struct {}{}
2026-05-22 20:26:11 -04:00
}
2026-07-17 02:48:00 +08:00
return active
2026-05-22 20:26:11 -04:00
}
func ( s * Server ) SetFFmpegLogSink ( sink playback . FFmpegLogSink ) {
s . ffmpegSink = sink
}
2026-07-14 11:51:27 -04:00
// 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 ())
}
2026-07-03 02:23:14 +08:00
// 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
}
2026-05-22 20:26:11 -04:00
// Handler returns the chi.Router with all transcode routes.
func ( s * Server ) Handler () http . Handler {
2026-07-14 11:51:27 -04:00
s . startIdleReaper ()
2026-05-22 20:26:11 -04:00
r := chi . NewRouter ()
r . Get ( "/api/v1/health" , s . handleHealth )
r . Group ( func ( r chi . Router ) {
r . Use ( s . requireBearer )
r . Get ( "/hw-capabilities" , s . handleHWCapabilities )
r . Post ( "/chapter-thumbnails/extract" , s . handleChapterThumbnailExtract )
r . Post ( "/transcode/start" , s . handleStart )
r . Delete ( "/transcode/{session_id}" , s . handleStop )
r . Get ( "/transcode/{session_id}/master.m3u8" , s . handleManifest )
r . Get ( "/transcode/{session_id}/segment/{name}" , s . handleSegment )
2026-08-11 14:21:33 -04:00
r . Post ( "/transcode/{session_id}/segment/{name}/downloaded" , s . handleSegmentDownloaded )
2026-05-22 20:26:11 -04:00
r . Post ( "/admin/force-reload" , s . handleForceReload )
r . Get ( "/status" , s . handleStatus )
})
return r
}
func ( s * Server ) handleHealth ( w http . ResponseWriter , _ * http . Request ) {
w . Header (). Set ( "Content-Type" , "application/json" )
json . NewEncoder ( w ). Encode ( HealthResponse {
Status : "ok" ,
ActiveJobs : s . activeJobs . Load (),
})
}
2026-07-14 11:51:27 -04:00
func ( s * Server ) handleHWCapabilities ( w http . ResponseWriter , r * http . Request ) {
2026-06-08 02:27:56 +01:00
ffmpegPath := ""
if cfg := s . watcher . Config (); cfg != nil {
ffmpegPath = cfg . Playback . FFmpegPath
}
info := playback . DetectHWAccelWithFFmpeg ( ffmpegPath )
2026-07-14 11:51:27 -04:00
info . Transformations = playback . ProbeTransformationRegistryV3 ( r . Context (), ffmpegPath ). Advertised ()
2026-05-22 20:26:11 -04:00
w . Header (). Set ( "Content-Type" , "application/json" )
json . NewEncoder ( w ). Encode ( info )
}
func ( s * Server ) handleChapterThumbnailExtract ( w http . ResponseWriter , r * http . Request ) {
var req chapterthumbs . RemoteExtractRequest
if err := json . NewDecoder ( r . Body ). Decode ( & req ); err != nil {
writeChapterThumbnailError ( w , http . StatusBadRequest , "invalid_request" , "invalid request body" )
return
}
if strings . TrimSpace ( req . InputPath ) == "" {
writeChapterThumbnailError ( w , http . StatusBadRequest , "invalid_request" , "input_path is required" )
return
}
cfg := s . watcher . Config ()
frame , reason , err := chapterthumbs . ExtractFrame ( r . Context (), chapterthumbs . FrameExtractOptions {
2026-08-11 18:13:06 +02:00
InputPath : req . InputPath ,
SeekSeconds : req . SeekSeconds ,
FFmpegPath : cfg . Playback . FFmpegPath ,
HWAccel : cfg . Playback . HWAccel ,
HWDevice : cfg . Playback . HWDevice ,
ToneMap : req . ToneMap ,
AllowSoftwareToneMap : req . AllowSoftwareToneMap ,
2026-05-22 20:26:11 -04:00
})
if err != nil {
writeChapterThumbnailError ( w , http . StatusUnprocessableEntity , reason , err . Error ())
return
}
w . Header (). Set ( "Content-Type" , "image/jpeg" )
w . WriteHeader ( http . StatusOK )
_ , _ = w . Write ( frame )
}
func writeChapterThumbnailError ( w http . ResponseWriter , status int , reason string , message string ) {
w . Header (). Set ( "Content-Type" , "application/json" )
w . WriteHeader ( status )
_ = json . NewEncoder ( w ). Encode ( chapterthumbs . RemoteExtractErrorResponse {
Reason : reason ,
Error : message ,
})
}
// requireBearer is middleware that checks for Authorization: Bearer {secret}.
func ( s * Server ) requireBearer ( next http . Handler ) http . Handler {
return http . HandlerFunc ( func ( w http . ResponseWriter , r * http . Request ) {
cfg := s . watcher . Config ()
auth := r . Header . Get ( "Authorization" )
if ! strings . HasPrefix ( auth , "Bearer " ) || strings . TrimPrefix ( auth , "Bearer " ) != cfg . Auth . JWTSecret {
http . Error ( w , "unauthorized" , http . StatusUnauthorized )
return
}
next . ServeHTTP ( w , r )
})
}
func ( s * Server ) handleStart ( w http . ResponseWriter , r * http . Request ) {
var req TranscodeStartRequest
if err := json . NewDecoder ( r . Body ). Decode ( & req ); err != nil {
http . Error ( w , "invalid request body" , http . StatusBadRequest )
return
}
if req . SessionID == "" || req . InputPath == "" {
http . Error ( w , "session_id and input_path are required" , http . StatusBadRequest )
return
}
cfg := s . watcher . Config ()
outputDir := filepath . Join ( cfg . Playback . TranscodeDir , req . SessionID )
opts := playback . TranscodeOpts {
2026-08-08 15:23:16 -04:00
InputPath : req . InputPath ,
OutputDir : outputDir ,
SessionID : req . SessionID ,
SourceVideoCodec : req . SourceVideoCodec ,
2026-08-11 13:05:12 -04:00
SourceVideoProfile : req . SourceVideoProfile ,
SourceVideoBitDepth : req . SourceVideoBitDepth ,
SoftwareVideoDecode : req . SoftwareVideoDecode ,
2026-08-08 15:23:16 -04:00
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 ,
2026-08-11 13:05:12 -04:00
TargetAudioBitrateKbps : req . TargetAudioBitrateKbps ,
2026-08-08 15:23:16 -04:00
TargetBitrateKbps : req . TargetBitrateKbps ,
SegmentDuration : req . SegmentDuration ,
SegmentRetentionSeconds : cfg . Playback . SegmentRetentionSeconds ,
FFmpegPath : cfg . Playback . FFmpegPath ,
HWAccel : req . HWAccel ,
2026-08-04 17:33:15 +02:00
// 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 ,
2026-05-22 20:26:11 -04:00
}
if opts . HWAccel == "" && cfg . Playback . HWAccel != "" {
opts . HWAccel = cfg . Playback . HWAccel
}
2026-07-03 02:23:14 +08:00
// 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 )
2026-05-22 20:26:11 -04:00
// 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 )
2026-07-14 11:51:27 -04:00
delete ( s . lastAccess , req . SessionID )
2026-05-22 20:26:11 -04:00
s . mu . Unlock ()
s . activeJobs . Add ( - 1 )
_ = old . Close ()
2026-06-10 17:18:18 -04:00
// 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 )
}
2026-05-22 20:26:11 -04:00
} else {
s . mu . Unlock ()
}
session , err := playback . StartTranscode ( context . WithoutCancel ( r . Context ()), opts )
if err != nil {
2026-07-03 02:23:14 +08:00
unlock ()
2026-07-09 20:53:52 +08:00
slog . ErrorContext ( r . Context (), "start transcode" , "component" , "transcodenode" , "error" , err , "session" , req . SessionID , "playback_session_id" , req . SessionID )
2026-05-22 20:26:11 -04:00
http . Error ( w , "failed to start transcode" , http . StatusInternalServerError )
return
}
2026-07-14 11:51:27 -04:00
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
}
}
2026-05-22 20:26:11 -04:00
s . mu . Lock ()
s . sessions [ req . SessionID ] = session
2026-07-14 11:51:27 -04:00
s . noteSessionAccessLocked ( req . SessionID )
2026-05-22 20:26:11 -04:00
s . mu . Unlock ()
2026-07-03 02:23:14 +08:00
unlock ()
2026-05-22 20:26:11 -04:00
s . activeJobs . Add ( 1 )
2026-06-10 17:18:18 -04:00
// 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.
2026-06-05 19:43:20 -07:00
effectiveHWAccel := session . Opts (). HWAccel
2026-06-10 17:18:18 -04:00
trackCtx := context . WithoutCancel ( r . Context ())
go s . tracker . Track ( trackCtx , nodesessions . SessionInfo {
2026-05-22 20:26:11 -04:00
SessionID : req . SessionID ,
NodeURL : s . tracker . NodeURL (),
NodeName : s . tracker . NodeName (),
Type : "transcode" ,
CodecVideo : req . TargetCodecVideo ,
CodecAudio : req . TargetCodecAudio ,
Resolution : req . TargetResolution ,
2026-06-05 19:43:20 -07:00
HWAccel : effectiveHWAccel ,
2026-05-22 20:26:11 -04:00
StartedAt : time . Now (). UTC (). Format ( time . RFC3339 ),
})
w . WriteHeader ( http . StatusAccepted )
json . NewEncoder ( w ). Encode ( TranscodeStartResponse {
SessionID : req . SessionID ,
Status : "started" ,
2026-06-05 19:43:20 -07:00
HWAccel : effectiveHWAccel ,
2026-05-22 20:26:11 -04:00
})
}
2026-07-03 02:23:14 +08:00
// reconstructFromToken rebuilds a transcode session this node lost to its own
// restart. The proxy forwards the client's verified stream token in the
// X-Silo-Stream-Token header; the token carries the full byte-affecting recipe
// (the former Postgres "recipe card"), so the node can re-spawn ffmpeg seeked to
// the requested segment rather than 404ing — mirroring the integrated server's
// token-carried reconstruct. Returns nil when the request carries no usable
// transcode token, which the caller renders as a genuine not-found.
//
// requestedSegment is the segment the client is fetching, or negative on the
// manifest path. Reconstruction is single-flighted per session id so concurrent
// manifest and segment requests for the same lost session share one ffmpeg.
func ( s * Server ) reconstructFromToken ( r * http . Request , sessionID string , requestedSegment int ) * playback . TranscodeSession {
tokenStr := r . Header . Get ( "X-Silo-Stream-Token" )
if tokenStr == "" {
return nil
}
cfg := s . watcher . Config ()
claims , err := streamtoken . Verify ( tokenStr , cfg . Auth . JWTSecret )
if err != nil {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( r . Context (), "transcode node reconstruct: invalid stream token" , "component" , "transcodenode" , "error" , err ,
2026-07-03 02:23:14 +08:00
"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).
2026-07-14 11:51:27 -04:00
expectedTransportID := card . SessionID
if card . TranscodeTransportID != "" {
expectedTransportID = card . TranscodeTransportID
}
if expectedTransportID != sessionID || ( card . PlayMethod != "" && card . PlayMethod != playback . PlayTranscode ) {
2026-07-03 02:23:14 +08:00
return nil
}
// A native token carries the full byte-affecting recipe. The jellycompat node
// hop signs an identity-only token by design (see internal/noderecipe for why),
// so its card decodes with no encode parameters. For the jellycompat case the
// recipe is fetched from the control-plane recipe store below; without that
// store there is nothing to rebuild from, so 404.
tokenComplete := card . SegmentDuration > 0 && card . TargetCodecVideo != ""
if ! tokenComplete && s . recipeStore == nil {
return nil
}
v , _ , _ := s . reconstructGroup . Do ( sessionID , func () ( interface {}, error ) {
// A concurrent reconstruct (or a fresh start) may already have registered the
// session; serve it rather than spawning a duplicate ffmpeg.
s . mu . RLock ()
existing , ok := s . sessions [ sessionID ]
s . mu . RUnlock ()
if ok {
return existing , nil
}
resolved := card
if ! tokenComplete {
// Recipe-less (jellycompat) token: fetch the recipe central wrote to the
// control-plane store at transcode start. A miss / incomplete recipe is a
// genuine not-found (404), never a spawn from a bad recipe.
fetched , ok := s . recipeStore . Get ( r . Context (), sessionID )
if ! ok || fetched == nil || fetched . SessionID != sessionID ||
fetched . SegmentDuration <= 0 || fetched . TargetCodecVideo == "" {
return ( * playback . TranscodeSession )( nil ), nil
}
resolved = * fetched
}
return s . spawnReconstruct ( r , sessionID , requestedSegment , resolved ), nil
})
if session , _ := v .( * playback . TranscodeSession ); session != nil {
return session
}
return nil
}
// spawnReconstruct re-spawns ffmpeg for a lost session from its recipe card and
// registers it in the live map. It is only ever called inside the per-session
// single-flight in reconstructFromToken, so it is the sole writer racing to
// register sessionID. Returns nil if the spawn fails or the slot wait is canceled.
func ( s * Server ) spawnReconstruct ( r * http . Request , sessionID string , requestedSegment int , card playback . RecipeCard ) * playback . TranscodeSession {
// Pace the cold-start burst so a node restart that loses many sessions does not
// launch every ffmpeg at once. A client that disconnects while waiting releases
// its slot rather than queueing dead work.
release , ok := s . acquireReconstructSlot ( r . Context ())
if ! ok {
return nil
}
defer release ()
// Serialize against a concurrent fresh /transcode/start for this session so the
// two never run ffmpeg writers against the same dir. Re-check under the lock and
// yield to any live session rather than spawning a duplicate.
unlock := s . lockSessionLifecycle ( sessionID )
defer unlock ()
s . mu . RLock ()
existing , ok := s . sessions [ sessionID ]
s . mu . RUnlock ()
if ok {
return existing
}
cfg := s . watcher . Config ()
outputDir := filepath . Join ( cfg . Playback . TranscodeDir , sessionID )
opts := card . TranscodeOpts ( outputDir , cfg . Playback . FFmpegPath , s . ffmpegSink )
2026-07-14 11:51:27 -04:00
opts . SessionID = sessionID
2026-07-03 02:23:14 +08:00
// 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
2026-08-08 15:23:16 -04:00
opts . SegmentRetentionSeconds = cfg . Playback . SegmentRetentionSeconds
2026-07-03 02:23:14 +08:00
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 {
2026-07-09 20:53:52 +08:00
slog . ErrorContext ( r . Context (), "transcode node reconstruct start failed" , "component" , "transcodenode" , "error" , err ,
2026-07-03 02:23:14 +08:00
"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
2026-07-14 11:51:27 -04:00
s . noteSessionAccessLocked ( sessionID )
2026-07-03 02:23:14 +08:00
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 ,
})
2026-07-09 20:53:52 +08:00
slog . InfoContext ( r . Context (), "transcode node session reconstructed from token" , "component" , "transcodenode" ,
2026-07-03 02:23:14 +08:00
"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
}
}
2026-05-22 20:26:11 -04:00
func ( s * Server ) handleStop ( w http . ResponseWriter , r * http . Request ) {
sessionID := chi . URLParam ( r , "session_id" )
2026-07-14 11:51:27 -04:00
// 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 ()
2026-05-22 20:26:11 -04:00
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 )
2026-07-14 11:51:27 -04:00
delete ( s . lastAccess , sessionID )
2026-05-22 20:26:11 -04:00
s . mu . Unlock ()
s . activeJobs . Add ( - 1 )
if err := session . Close (); err != nil {
2026-07-09 20:53:52 +08:00
slog . ErrorContext ( r . Context (), "close transcode session" , "component" , "transcodenode" , "error" , err , "session" , sessionID , "playback_session_id" , sessionID )
2026-05-22 20:26:11 -04:00
}
cfg := s . watcher . Config ()
outputDir := filepath . Join ( cfg . Playback . TranscodeDir , sessionID )
os . RemoveAll ( outputDir )
2026-07-03 02:23:14 +08:00
// 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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( r . Context (), "delete transcode recipe on stop" , "component" , "transcodenode" , "error" , err , "session" , sessionID , "playback_session_id" , sessionID )
2026-07-03 02:23:14 +08:00
}
}
2026-05-22 20:26:11 -04:00
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" )
2026-07-14 11:51:27 -04:00
// 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 )
2026-05-22 20:26:11 -04:00
if ! ok {
2026-07-03 02:23:14 +08:00
// 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
}
2026-07-14 11:51:27 -04:00
// 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 )
2026-05-22 20:26:11 -04:00
}
2026-08-05 20:27:42 -04:00
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 )
}
2026-05-22 20:26:11 -04:00
if err != nil {
2026-07-09 20:53:52 +08:00
slog . ErrorContext ( r . Context (), "get manifest" , "component" , "transcodenode" , "error" , err , "session" , sessionID , "playback_session_id" , sessionID )
2026-05-22 20:26:11 -04:00
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" )
2026-07-14 11:51:27 -04:00
// 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 )
2026-05-22 20:26:11 -04:00
if ! ok {
2026-07-03 02:23:14 +08:00
// 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
}
2026-07-14 11:51:27 -04:00
// 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 )
2026-05-22 20:26:11 -04:00
}
2026-08-11 12:59:25 -04:00
segmentLease , err := session . OpenSegment ( name )
2026-05-22 20:26:11 -04:00
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 ()
}
2026-07-09 20:53:52 +08:00
slog . InfoContext ( r . Context (), "transcode segment missing" , "component" , "transcodenode" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-09 20:53:52 +08:00
slog . InfoContext ( r . Context (), "transcode segment wait" , "component" , "transcodenode" ,
2026-05-22 20:26:11 -04:00
"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 ,
)
2026-08-11 12:59:25 -04:00
segmentLease , err = session . WaitForOpenSegment ( name , decision . WaitTimeout )
2026-05-22 20:26:11 -04:00
if err != nil && err == playback . ErrSegmentNotFound {
2026-07-09 20:53:52 +08:00
slog . InfoContext ( r . Context (), "transcode segment wait timeout" , "component" , "transcodenode" ,
2026-05-22 20:26:11 -04:00
"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 ,
)
}
}
2026-06-07 04:30:24 +02:00
if err != nil && err == playback . ErrSegmentNotFound && decision . RestartOnTimeout {
2026-05-22 20:26:11 -04:00
seekSeconds , ok , seekErr := session . RestartSeekTarget ( segNum )
if seekErr != nil && ! errors . Is ( seekErr , playback . ErrManifestNotReady ) {
2026-07-09 20:53:52 +08:00
slog . ErrorContext ( r . Context (), "resolve transcode node seek target" , "component" , "transcodenode" , "error" , seekErr , "segment" , name , "session" , sessionID , "playback_session_id" , sessionID )
2026-05-22 20:26:11 -04:00
}
if ok {
2026-07-09 20:53:52 +08:00
slog . InfoContext ( r . Context (), "transcode node seek restart" , "component" , "transcodenode" ,
2026-05-22 20:26:11 -04:00
"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 ,
)
2026-07-03 02:23:14 +08:00
if restartErr := s . restartSessionLocked (
2026-05-22 20:26:11 -04:00
context . WithoutCancel ( r . Context ()),
2026-07-03 02:23:14 +08:00
sessionID ,
session ,
2026-05-22 20:26:11 -04:00
seekSeconds ,
segNum ,
); restartErr == nil {
2026-08-11 12:59:25 -04:00
segmentLease , err = session . WaitForOpenSegment ( name , 30 * time . Second )
2026-05-22 20:26:11 -04:00
}
}
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.
2026-08-11 12:59:25 -04:00
segmentLease , err = session . WaitForOpenSegment ( name , 10 * time . Second )
2026-05-22 20:26:11 -04:00
}
}
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" )
2026-08-11 14:30:55 -04:00
proxied := r . Header . Get ( transcodeproxy . RequestHeader ) == "1"
2026-08-11 14:21:33 -04:00
if proxied {
2026-08-11 14:30:55 -04:00
w . Header (). Set ( transcodeproxy . GenerationHeader , segmentLease . GenerationToken )
2026-08-11 14:21:33 -04:00
}
2026-08-11 12:59:25 -04:00
defer func () { _ = segmentLease . Close () }()
2026-08-08 15:23:16 -04:00
sw := httpstream . NewRollingDeadlineWriter ( w )
2026-08-11 12:59:25 -04:00
http . ServeContent ( sw , r , segmentLease . Info . Name (), segmentLease . Info . ModTime (), segmentLease . File )
2026-08-11 14:21:33 -04:00
if ! proxied && r . Method == http . MethodGet &&
2026-08-11 15:02:38 -04:00
sw . CompletedFullResponse ( segmentLease . Info . Size ()) {
2026-08-08 15:23:16 -04:00
if segmentNumber , parseErr := playback . ParseSegmentNumber ( name ); parseErr == nil {
2026-08-11 12:59:25 -04:00
session . ReportSegmentDownloadedForGeneration ( segmentNumber , segmentLease . Generation )
2026-08-08 15:23:16 -04:00
}
}
2026-05-22 20:26:11 -04:00
}
2026-08-11 14:21:33 -04:00
// handleSegmentDownloaded records completion only after the central API has
// delivered the full segment to its downstream playback client. The generation
// returned with the original segment response prevents a delayed acknowledgement
// from advancing a replacement FFmpeg timeline.
func ( s * Server ) handleSegmentDownloaded ( w http . ResponseWriter , r * http . Request ) {
sessionID := chi . URLParam ( r , "session_id" )
name := chi . URLParam ( r , "name" )
segmentNumber , err := playback . ParseSegmentNumber ( name )
if err != nil {
http . Error ( w , "invalid segment" , http . StatusBadRequest )
return
}
2026-08-11 14:30:55 -04:00
generationToken := r . Header . Get ( transcodeproxy . GenerationHeader )
2026-08-11 14:21:33 -04:00
if generationToken == "" {
http . Error ( w , "invalid segment generation" , http . StatusBadRequest )
return
}
session , ok := s . acquireSessionTouched ( sessionID )
if ! ok {
http . Error ( w , "session not found" , http . StatusNotFound )
return
}
session . ReportSegmentDownloadedForGenerationToken ( segmentNumber , generationToken )
w . WriteHeader ( http . StatusNoContent )
}
2026-05-22 20:26:11 -04:00
func ( s * Server ) handleForceReload ( w http . ResponseWriter , r * http . Request ) {
if err := s . watcher . ForceReload ( r . Context ()); err != nil {
http . Error ( w , "reload failed: " + err . Error (), http . StatusInternalServerError )
return
}
cfg := s . watcher . Config ()
s . mu . Lock ()
2026-07-03 02:23:14 +08:00
stopped := make ([] string , 0 , len ( s . sessions ))
2026-05-22 20:26:11 -04:00
for id , session := range s . sessions {
session . Close ()
os . RemoveAll ( filepath . Join ( cfg . Playback . TranscodeDir , id ))
delete ( s . sessions , id )
2026-07-14 11:51:27 -04:00
delete ( s . lastAccess , id )
2026-07-03 02:23:14 +08:00
stopped = append ( stopped , id )
2026-05-22 20:26:11 -04:00
}
s . activeJobs . Store ( 0 )
s . mu . Unlock ()
2026-07-03 02:23:14 +08:00
// 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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( r . Context (), "delete transcode recipe on force reload" , "component" , "transcodenode" , "error" , err , "session" , id , "playback_session_id" , id )
2026-07-03 02:23:14 +08:00
}
}
}
2026-05-22 20:26:11 -04:00
s . tracker . Cleanup ( r . Context ())
2026-07-09 20:53:52 +08:00
slog . InfoContext ( r . Context (), "transcode force reload completed" , slog . String ( "component" , "transcodenode" ))
2026-05-22 20:26:11 -04:00
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 ,
})
}