* docs(playback): add v3 neutral-contract finalization plan Supersedes the wire-contract sections of the 2026-07-12 v3 plan: server-owned attempt keys, delivery-keyed negotiation without Media3 engine names, tiered capability evidence, neutral device/output context, track/quality replan operations, audio-only planning, and coordinated no-back-compat rollout across server, Android, Apple, and web. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * feat(playback): make v3 attempt keys server-owned and replace engines with deliveries Contract core of the platform-neutral v3 finalization (plan sections 3.1 and 3.2), breaking on purpose — v3 is dark and all clients move together: - Every PlanV3 now carries plan_attempt_key, an opaque server-computed token clients store and echo in attempted_plan_keys; ReplanRequestV3 gains bounded local_mutations that the replan handler folds into the failed plan's key. Clients never hash anything. - KotlinName() is deleted from DeliveryV3, StreamProtocolV3 and SubtitleModeV3; the attempt-key canonical string now uses lowercase wire tokens, and PlanRecipeVersionV3 bumps to v3.3 so no key or plan ID computed under the old canonicalization can collide. - EngineV3 leaves the wire: ClientPlaybackContextV3.Engines (media3_*) becomes Deliveries keyed original_http|progressive|hls, with EngineCapabilityV3 renamed DeliveryCapabilityV3. PlanV3.Engine is removed; the planner, subtitle policy and quirk registry re-key on delivery class, and the media3_only feature token is deleted. - Validated-claim strings drop the prefix: media3_h264_decode -> h264_decode, media3_audio_decode -> audio_decode. - Golden fixtures in testdata/protocol_v3 are regenerated by Go and are now the cross-repo source of truth. Part of the playback protocol v3 neutral-contract train (steps 2-3 of docs/superpowers/plans/2026-07-30-playback-protocol-v3-neutral-contract.md). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * feat(playback): add v3 evidence tiers and neutral device/output context Implement plan sections 3.3 and 3.4 of the v3 neutral-contract pass: - ClientCodecCapabilitiesV3 gains required video_evidence and audio_evidence closed enums (exact | platform_attested | declared). Planner strictness follows the tier: exact keeps the strict decode-entry validation, platform_attested validates codec/resolution/bit-depth/ frame-rate but skips profile/level matching, declared grants copy routes from the flat codec lists. Only exact audio evidence earns passthrough claims. The detailed_decode_capabilities feature token is deleted (subsumed by video_evidence=exact), and evidence-blocked direct routes carry the new evidence_insufficient_for_direct reason/warning. - DeviceContextV3 is now platform/os_version/manufacturer/model plus a bounded platform_details map (<=16 entries, <=128 chars); the Android Build dump fields are gone. Fire TV quirks keep matching on manufacturer/model (brand fallback removed with the field). - output_route_generation (int64, dual-location) becomes an optional opaque output_context_id string on the output context; the dual-location consistency validation is deleted. Attempt keys, plan invalidation, route events, and the planstore column follow (new Goose migration). - Feature advertisement collapses to the top-level client_features list only; ClientPlaybackContextV3.Features is deleted and ReplanRequestV3 gains an optional client_features refresh. - PlanRecipeVersionV3 bumped v3.3 -> v3.4; fixtures re-keyed. Part of the playback protocol v3 neutral-contract finalization plan (docs/superpowers/plans/2026-07-30-playback-protocol-v3-neutral-contract.md). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * feat(playback): add v3 intent replans, quality menu, and audio-only routes Protocol v3 could only replan after a failure, so changing the audio track or the quality still required the legacy audio PATCH and the client-recipe transcode start — the two endpoints v3 is meant to replace. Clients also had to own a resolution ladder to render a quality menu, and a source with no video track was terminaled by the video/HDR gates, keeping audiobooks on the legacy path. Add track_change and quality_change replan operations. They carry no failure classification and route through the existing replan transaction, so they inherit its idempotency, capacity reservation, and staged-successor commit for free. Because nothing failed, the previous route stays eligible: neither the attempted-key history nor the failed-plan exclusion applies to them. Publish the server ladder on the plan as available_qualities so the quality menu is server-owned; the rungs come from the same resolutionLabelV3 and ladderBitrateKbpsV3 helpers the planner itself uses, not a parallel table. Plan audio-only sources through their own reduced route family: original_http when the client decodes the codec, otherwise a progressive AAC conversion. The plan advertises audio/mp4 for that remux and the transport now serves the same value, because a declared-tier client probes the advertised MIME with isTypeSupported before attaching a source buffer, and "video/mp4" on a stream with no video track is exactly the mismatch that makes the probe lie. Name the protocol's string vocabulary (dynamic ranges, transformations, executors, validated claims, terminal reasons) as constants while touching these lines, so the wire values have one definition. Part of #135 * docs(playback): publish the v3 protocol contract and fix subtitle ordinals Protocol v3 exists only as Go code today, so the Android and Apple ports have no authority to implement against other than reading this repository. Publish the contract as a normative document, machine-checkable schemas, and generated golden fixtures, and fix the one place where the server's own wire output disagreed with the ordinal space it publishes. - docs/architecture/playback-protocol-v3.md is self-contained enough for a third-party client: endpoints and status codes, evidence tiers and their bound-matching rules, delivery classes, the timeline model, replan semantics, registries, track identity, plan identity, quality, and transformations. - docs/design/schemas/playback-v3/ carries JSON Schemas for the five wire shapes plus valid and invalid fixtures, following the client-diagnostics layout. internal/playback/contract validates every fixture against its schema, so a schema that drifts from the Go types fails the Go suite. - cmd/playbackfixtures generates internal/playback/testdata/protocol_v3 from the production planner. `make playback-fixtures` writes them and `make verify-playback-fixtures` (wired into CI) fails when they are stale. These files are what the client ports consume, so drift would otherwise surface as a playback bug on three platforms at once. The subtitle fix: combined ordinals are one dense space over externals, then embedded tracks, then downloaded ones, but the legacy URL builder skipped burn-in-only tracks while assigning indices, so every track after a DVD/DVB track was numbered one too low and resolved to its neighbour. Ordinal assignment now lives in playback.BuildSubtitleInventoryV3 and both the plan inventory and the legacy `subtitle_urls` shape project from it; the legacy shape still filters burn-in-only entries but keeps each track's real index. Part of #135 * feat(web): migrate the players to the neutral playback v3 contract The web player was the last client still speaking the legacy start protocol: it picked its own file version from a codec probe, posted an ffmpeg recipe to start a transcode, PATCHed an endpoint to change audio tracks, and derived its own quality ladder. None of that survives a server-owned plan, and none of it produced telemetry the apps could be compared against. Video player: starts with a v3 request that advertises `declared` evidence from `isTypeSupported` probes and the three delivery classes, then consumes the returned plan for its URL, timeline, tracks and warnings. Quality and track changes become replans (`quality_change`, `track_change`), the quality menu renders `available_qualities` instead of computing rungs, and playback failures emit `route-events` so web failures land in the same diagnostics as Android and Apple. The duration comes from `source.duration_seconds` rather than the playback engine, and the "how was this delivered" overlay reads the plan's delivery and server transformations instead of comparing codec strings. Audiobook player: starts against the audio-only planner path with a single `original` rung, and takes its seek anchor from `timeline.player_start_seconds` so the progressive-remux route (which anchors the stream and restarts the player clock at zero) does not seek twice. Server side, `disable_progress_persistence` left the wire, so the rule it encoded is now derived. Resume state is keyed on the item, but every part of a multipart presentation shares that key while carrying its own file-local clock — persisting part 4's position would store "12 minutes in" as the book's resume point. `PresentationPartTotal > 1` expresses that directly and generalizes to multipart movies and split episodes, and a client can no longer forget to ask or lie about it. `useTranscodeQuality` and the legacy response types are deleted, and `WEBTEST_KNOWN_FAILURES` loses the audiobook entry along with its fix. Part of #135 * feat(playback)!: make v3 the only playback protocol Protocol v3 shipped behind a flag, alongside the legacy start path it was designed to replace. Running both meant every planner change had to be made twice, in two shapes that disagree about who decides the route: the legacy body carried a decision the client had already made, while v3 asks the server to make it. This deletes the legacy half. Removed: - `handleStartPlaybackLegacy` and its request/response bodies. The `POST /playback/start` route stays, but the protocol-version dispatch envelope is now a strict v3 decode — a body that does not declare `protocol_version: 3` gets `426 client_upgrade_required` so an outdated app can render a clear "update required" state instead of misreading a plan. Deliberately not a `400`: the request may be well-formed for the protocol it was written against. - `POST /playback/transcode/start`, superseded by the `quality_change` replan operation, and `PATCH /playback/{session_id}/audio`, superseded by `track_change`. Both mutated a session without re-planning. - The shadow planner and both rollout settings rows. With v3 the only protocol, `playback.protocol_v3_enabled` would mean "no playback at all"; `playback.protocol_v3_shadow_enabled` gated a comparison against a path that no longer exists. `409 protocol_disabled` on route-events goes with them, and capability `enabled` is now constant `true` (the field stays — clients feature-detect against it). - Version-selection helpers in `internal/playback/resolver.go` that only legacy start reached. `Resolve`/`ClientCapabilities`/`PlayDecision` stay: downloads consumes them. `internal/jellycompat` has its own resolution surface and is untouched. Behaviour the legacy handlers owned and v3 now owns explicitly: series version and audio-track preferences are persisted on start and on a `track_change` replan (not on failure recovery, whose forced route is not a user choice); an omitted `start_position` resolves to the profile's saved resume point; and an omitted audio track resolves through the series preference, the profile audio language, then the library override. Both are settled before planning, because the plan's timeline is cut at the start position. Spec §2.2 documents this as "omission is a request, not a default". The encode-target clamp that lived in the deleted transcode handler is already enforced in the planner, twice — `availableQualitiesV3` omits rungs at or above the source height, and the encode path clamps `targetHeight` to it. Unchanged: progress, stop, HLS manifest and segment delivery, the realtime control socket, stream tokens and restart reconstruction, watch together, downloads, jellycompat. Every removal is recorded in the pre-lock removals table in docs/architecture/v1-scope.md. Part of #135 * fix(scanner): stop recording embedded cover art as a video track ffprobe reports embedded cover art as a video stream carrying disposition.attached_pic. convertProbeData appended every "video" stream to VideoTracks without consulting isMainVideoStream, the predicate that already existed for duration decisions, so the picture was persisted as a playable track. That misreports the file twice: - An audio file with a cover picks up a video track, so it no longer satisfies MediaFile.IsAudioOnly and the v3 planner routes an audiobook through the video path instead of planAudioOnlyV3. - When the picture is ordered ahead of the real stream, the flat codec_video/resolution/hdr columns describe the poster: a 954x720 h264 episode was stored as mjpeg 480x480. Filter attached_pic streams out of the track loop. The guard is the disposition flag, not the codec name, so a genuine MJPEG video is still probed as video — the library has one. Already-probed rows self-heal on the next playback: NeedsCriticalProbeRepair already reprobes tracks missing color_range, which covers 21 of the 23 affected rows, and applyProbeData overwrites VideoTracks wholesale. The remaining two need a rescan; nothing persisted records attached_pic, and keying repair off still-image codec names would reprobe the genuine MJPEG file on every playback forever. Part of the playback v3 neutral-contract work: it is what lets Android drop AUDIOBOOK_COVER_ART_CODECS, which fabricated decode support the client cannot honestly claim under video_evidence: "exact". Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix(playback): publish subtitle URLs even when playback starts with subtitles off The v3 plan's subtitle inventory is the authoritative track list a client builds its subtitle menu from, but the handler only rewrote it with session-scoped URLs when a track was actually selected. A start or replan that resolved to `subtitle.mode: "off"` therefore returned the planner's URL-less inventory, so a client whose picker reads the inventory had a menu it could not fetch anything from. The Cast path hits this every time: it starts with subtitles off and needs the receiver's text tracks up front. attachSubtitleArtifactV3 now scopes and publishes the inventory unconditionally and gates only the artifact stamping on the selection. Spec §8 records that the `url` on a sidecar entry does not depend on the current selection. Part of the v3 neutral-contract finalization. * chore(playback): reconcile neutral v3 with main * fix(playback): preserve subtitle intent across replans * fix(playback): retain subtitle inventory on adapted routes * fix(playback): software-decode High10 AVC for QSV * fix(playback): scale High10 frames before QSV upload * fix(playback): preserve empty subtitle inventories * fix(playback): freeze terminal attempt contract * chore(playback): name fixture contract tokens * fix(playback): close v3 conformance review gaps * chore(playback): name conformance category * fix(playback): complete v3 conformance contract * fix(playback): keep schema fixtures generated * fix(playback): emit schema-valid conformance arrays * fix(playback): omit empty replan failures * fix(web): omit empty replan failures * fix(playback): close neutral v3 contract gaps * fix(playback): harden v3 replan, transcode, and quality-ladder edge cases Review remediation for the neutral v3 cutover, server side: - A failed replan no longer overwrites the durable StartResponse with a terminal or advances the replan request ID; an idempotent start replay of a still-healthy session returns the original plan. - SoftwareVideoDecode is now derived inside the transcode layer from source facts (codec/profile/bit depth) carried on TranscodeOpts, so jellycompat, downloads, recipe-card reconstruction, and transcode nodes get the High10 software-decode fix, not just the v3 handler. video_to_h264 recipe version bumps to 2 so mixed-version node pools that would silently drop the flag fail validation instead. - Local transport startup shares the 30s ManifestStartupTimeout; a timeout with the process still running stays retryable and is no longer persisted as a durable terminal against the attempt. - Sparse replan bodies (failure_recovery et al) no longer reset a user-selected quality preference to auto; the empty-value guard now covers every operation. - availableQualitiesV3 publishes no fixed rungs when the source height is unknown, keeping the no-upscaling ladder contract. - The proxy remux path serves audio-only fMP4 as audio/mp4 via a new additive AudioOnly token claim, matching the integrated path. - Plain text subtitle sidecars accept any requested extension again (served as VTT), restoring the permissive v1 behavior; ASS and bitmap handling is unchanged. - The 4K-disallowed terminal message discloses when a lower-resolution alternate exists but was pinned away by quality "original". Part of #135. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix(web): keep playback alive through failed replans and honest audio claims Review remediation for the neutral v3 cutover, web player: - A failed or refused replan no longer unmounts the player: the fatal error screen is reserved for loads with no adopted plan, and replan failures surface through the existing non-fatal replanError path. - changeQuality rolls its optimistic preference back when the replan is refused or errors, so a failed switch is not silently applied by the next unrelated replan and the menu shows the real active rung. - The capability probe now tests mp3/vorbis codecs and mp3/flac/ogg containers (MediaSource with a canPlayType fallback), restoring direct play for mp3 audiobooks instead of per-part AAC re-encodes. - Reanchor seeks issued while a replan is in flight coalesce and run when it settles instead of being silently dropped with the scrubber pinned to a phantom position. - Subtitle refresh/translation replans use the resume anchor while the media element has no metadata, so a subtitle_ready broadcast during startup no longer restarts a resumed stream at 0:00. - An exhausted failure-recovery chain sets a visible error instead of returning silently. Part of #135. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix(playback): accept video-only and VP9 probe metadata Treat audio and video probe completeness independently so legitimate video-only assets converge without repeated ffprobe repair. Allow unknown codec profile/level metadata to fall through to server adaptation while preserving exact direct-decode constraints. Fixes #574 * fix(playback): address protocol v3 review findings * fix(playback): harden lease and probe repair decisions * fix(playback): close remaining v3 review gaps * fix(playback): recover failed transcode starts * fix(playback): address remaining review-bot findings on v3 replan and audio planning Server: - The deferred replan lease release is bounded by a 3s timeout so a saturated pool or DB outage cannot wedge a handler goroutine that holds the per-session store lock on an uncancellable context. - planAudioOnlyV3 honors the request bandwidth cap: an over-cap source skips the original_http direct route and converts to AAC with the same bandwidth_cap_applied warning and decision reason the video ladder uses. Unknown source bitrate never triggers the cap. - A copy-audio progressive plan rejected only by a per-delivery audio_decode_codecs subset retries as an AAC conversion instead of returning adaptation_unavailable, and the AAC recipe respects the delivery's max_channels. Web: - failure_recovery replans issued while another replan is in flight queue (superseding a pending seek reanchor) instead of being silently dropped with the fatal overlay already suppressed. - A terminal response to a fresh non-preserving start clears the previous plan and stops its session, so episode navigation cannot keep rendering the prior item under the new title. - A refused recovery replan for a transport-dead plan surfaces the error and re-arms the plan failure key, so transient recovery failures no longer strand an endless spinner; the audiobook player gets the same guard reset. - A track-less subtitle_translation_completed hands off to the refreshed persisted track once the inventory settles, clearing the live overlay, instead of pinning the synthetic live track forever. Part of #135. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix(playback): reuse HLS transport for sidecar replans * fix(playback): stabilize copy HLS remount timeline * fix(playback): address v3 review findings * fix(playback): satisfy player contract types --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
1060 lines
41 KiB
Go
1060 lines
41 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/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
|
||
|
||
// Server is the HTTP handler for transcode mode.
|
||
type Server struct {
|
||
watcher *nodeconfig.Watcher
|
||
tracker *nodesessions.Tracker
|
||
ffmpegSink playback.FFmpegLogSink
|
||
sessions map[string]*playback.TranscodeSession
|
||
// lastAccess records, per registered session id, when a manifest or segment
|
||
// request last touched the job (registration counts as the first access).
|
||
// Guarded by mu alongside sessions; the idle reaper closes jobs whose entry
|
||
// is older than sessionIdleTTL.
|
||
lastAccess map[string]time.Time
|
||
reaperOnce sync.Once
|
||
mu sync.RWMutex
|
||
activeJobs atomic.Int32
|
||
|
||
// reconstructGroup single-flights node-side session reconstruction per session
|
||
// id so a post-restart wave of concurrent manifest/segment requests for the same
|
||
// lost session spawns exactly one ffmpeg, never racing duplicates into the shared
|
||
// output directory.
|
||
reconstructGroup singleflight.Group
|
||
// reconstructSem bounds how many sessions may be reconstructed (ffmpeg
|
||
// re-spawned) at once after a node restart, pacing the cold-start burst instead
|
||
// of stampeding the host. Lazily sized to NumCPU on first use.
|
||
reconstructSemOnce sync.Once
|
||
reconstructSem chan struct{}
|
||
|
||
// lifecycleMu guards lifecycleLocks, the per-session mutexes that serialize
|
||
// every path which spawns ffmpeg into a session's output dir (fresh start and
|
||
// reconstruct). reconstructGroup only single-flights reconstructs against each
|
||
// other; without this a reconstruct racing a fresh /transcode/start could run
|
||
// two ffmpeg writers against the same dir.
|
||
lifecycleMu sync.Mutex
|
||
lifecycleLocks map[string]*sessionLifecycleLock
|
||
|
||
// recipeStore is the control-plane recipe store consulted when a forwarded
|
||
// token carries no recipe (the jellycompat node hop). Nil disables that path.
|
||
recipeStore recipeStore
|
||
}
|
||
|
||
// sessionLifecycleLock is a refcounted per-session mutex; the refcount lets the
|
||
// node drop the map entry once no path holds or waits on it so the map stays
|
||
// bounded over the node's lifetime.
|
||
type sessionLifecycleLock struct {
|
||
mu sync.Mutex
|
||
refs int
|
||
}
|
||
|
||
// lockSessionLifecycle acquires the per-session lifecycle mutex and returns a
|
||
// release func. Held across "check existing → spawn → register" so a fresh start
|
||
// and a reconstruct never run concurrent ffmpeg writers for one session's dir.
|
||
func (s *Server) lockSessionLifecycle(sessionID string) func() {
|
||
s.lifecycleMu.Lock()
|
||
if s.lifecycleLocks == nil {
|
||
s.lifecycleLocks = make(map[string]*sessionLifecycleLock)
|
||
}
|
||
lk := s.lifecycleLocks[sessionID]
|
||
if lk == nil {
|
||
lk = &sessionLifecycleLock{}
|
||
s.lifecycleLocks[sessionID] = lk
|
||
}
|
||
lk.refs++
|
||
s.lifecycleMu.Unlock()
|
||
|
||
lk.mu.Lock()
|
||
return func() {
|
||
lk.mu.Unlock()
|
||
s.lifecycleMu.Lock()
|
||
lk.refs--
|
||
if lk.refs == 0 {
|
||
delete(s.lifecycleLocks, sessionID)
|
||
}
|
||
s.lifecycleMu.Unlock()
|
||
}
|
||
}
|
||
|
||
// restartSessionLocked re-spawns session under the per-session lifecycle lock so
|
||
// a segment-recovery restart can never race a fresh start, reconstruct, or
|
||
// another restart into the same output directory. It holds the lock only across
|
||
// the cancel→respawn transition inside Restart and releases it before the caller
|
||
// waits on segments. Under the lock it confirms session is still the live mapped
|
||
// session; a concurrent teardown or reconstruct that replaced it yields
|
||
// ErrSessionSuperseded rather than re-spawning the stale handle.
|
||
func (s *Server) restartSessionLocked(ctx context.Context, sessionID string, session *playback.TranscodeSession, seekSeconds float64, startSegment int) error {
|
||
unlock := s.lockSessionLifecycle(sessionID)
|
||
defer unlock()
|
||
s.mu.RLock()
|
||
live, ok := s.sessions[sessionID]
|
||
s.mu.RUnlock()
|
||
if !ok || live != session {
|
||
return playback.ErrSessionSuperseded
|
||
}
|
||
return session.Restart(ctx, seekSeconds, startSegment)
|
||
}
|
||
|
||
// NewServer creates a new transcode server.
|
||
func NewServer(watcher *nodeconfig.Watcher, tracker *nodesessions.Tracker) *Server {
|
||
s := &Server{
|
||
watcher: watcher,
|
||
tracker: tracker,
|
||
sessions: make(map[string]*playback.TranscodeSession),
|
||
lastAccess: make(map[string]time.Time),
|
||
}
|
||
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
|
||
}
|
||
// 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{}{}
|
||
}
|
||
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
|
||
}
|
||
|
||
// Handler returns the chi.Router with all transcode routes.
|
||
func (s *Server) Handler() http.Handler {
|
||
s.startIdleReaper()
|
||
r := chi.NewRouter()
|
||
r.Get("/api/v1/health", s.handleHealth)
|
||
|
||
r.Group(func(r chi.Router) {
|
||
r.Use(s.requireBearer)
|
||
r.Get("/hw-capabilities", s.handleHWCapabilities)
|
||
r.Post("/chapter-thumbnails/extract", s.handleChapterThumbnailExtract)
|
||
r.Post("/transcode/start", s.handleStart)
|
||
r.Delete("/transcode/{session_id}", s.handleStop)
|
||
r.Get("/transcode/{session_id}/master.m3u8", s.handleManifest)
|
||
r.Get("/transcode/{session_id}/segment/{name}", s.handleSegment)
|
||
r.Post("/admin/force-reload", s.handleForceReload)
|
||
r.Get("/status", s.handleStatus)
|
||
})
|
||
return r
|
||
}
|
||
|
||
func (s *Server) handleHealth(w http.ResponseWriter, _ *http.Request) {
|
||
w.Header().Set("Content-Type", "application/json")
|
||
json.NewEncoder(w).Encode(HealthResponse{
|
||
Status: "ok",
|
||
ActiveJobs: s.activeJobs.Load(),
|
||
})
|
||
}
|
||
|
||
func (s *Server) handleHWCapabilities(w http.ResponseWriter, r *http.Request) {
|
||
ffmpegPath := ""
|
||
if cfg := s.watcher.Config(); cfg != nil {
|
||
ffmpegPath = cfg.Playback.FFmpegPath
|
||
}
|
||
info := playback.DetectHWAccelWithFFmpeg(ffmpegPath)
|
||
info.Transformations = playback.ProbeTransformationRegistryV3(r.Context(), ffmpegPath).Advertised()
|
||
w.Header().Set("Content-Type", "application/json")
|
||
json.NewEncoder(w).Encode(info)
|
||
}
|
||
|
||
func (s *Server) handleChapterThumbnailExtract(w http.ResponseWriter, r *http.Request) {
|
||
var req chapterthumbs.RemoteExtractRequest
|
||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||
writeChapterThumbnailError(w, http.StatusBadRequest, "invalid_request", "invalid request body")
|
||
return
|
||
}
|
||
if strings.TrimSpace(req.InputPath) == "" {
|
||
writeChapterThumbnailError(w, http.StatusBadRequest, "invalid_request", "input_path is required")
|
||
return
|
||
}
|
||
|
||
cfg := s.watcher.Config()
|
||
frame, reason, err := chapterthumbs.ExtractFrame(r.Context(), chapterthumbs.FrameExtractOptions{
|
||
InputPath: req.InputPath,
|
||
SeekSeconds: req.SeekSeconds,
|
||
FFmpegPath: cfg.Playback.FFmpegPath,
|
||
HWAccel: cfg.Playback.HWAccel,
|
||
HWDevice: cfg.Playback.HWDevice,
|
||
ToneMap: req.ToneMap,
|
||
})
|
||
if err != nil {
|
||
writeChapterThumbnailError(w, http.StatusUnprocessableEntity, reason, err.Error())
|
||
return
|
||
}
|
||
|
||
w.Header().Set("Content-Type", "image/jpeg")
|
||
w.WriteHeader(http.StatusOK)
|
||
_, _ = w.Write(frame)
|
||
}
|
||
|
||
func writeChapterThumbnailError(w http.ResponseWriter, status int, reason string, message string) {
|
||
w.Header().Set("Content-Type", "application/json")
|
||
w.WriteHeader(status)
|
||
_ = json.NewEncoder(w).Encode(chapterthumbs.RemoteExtractErrorResponse{
|
||
Reason: reason,
|
||
Error: message,
|
||
})
|
||
}
|
||
|
||
// requireBearer is middleware that checks for Authorization: Bearer {secret}.
|
||
func (s *Server) requireBearer(next http.Handler) http.Handler {
|
||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||
cfg := s.watcher.Config()
|
||
auth := r.Header.Get("Authorization")
|
||
if !strings.HasPrefix(auth, "Bearer ") || strings.TrimPrefix(auth, "Bearer ") != cfg.Auth.JWTSecret {
|
||
http.Error(w, "unauthorized", http.StatusUnauthorized)
|
||
return
|
||
}
|
||
next.ServeHTTP(w, r)
|
||
})
|
||
}
|
||
|
||
func (s *Server) handleStart(w http.ResponseWriter, r *http.Request) {
|
||
var req TranscodeStartRequest
|
||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||
http.Error(w, "invalid request body", http.StatusBadRequest)
|
||
return
|
||
}
|
||
|
||
if req.SessionID == "" || req.InputPath == "" {
|
||
http.Error(w, "session_id and input_path are required", http.StatusBadRequest)
|
||
return
|
||
}
|
||
|
||
cfg := s.watcher.Config()
|
||
outputDir := filepath.Join(cfg.Playback.TranscodeDir, req.SessionID)
|
||
|
||
opts := playback.TranscodeOpts{
|
||
InputPath: req.InputPath,
|
||
OutputDir: outputDir,
|
||
SessionID: req.SessionID,
|
||
SourceVideoCodec: req.SourceVideoCodec,
|
||
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,
|
||
})
|
||
}
|
||
|
||
// reconstructFromToken rebuilds a transcode session this node lost to its own
|
||
// restart. The proxy forwards the client's verified stream token in the
|
||
// X-Silo-Stream-Token header; the token carries the full byte-affecting recipe
|
||
// (the former Postgres "recipe card"), so the node can re-spawn ffmpeg seeked to
|
||
// the requested segment rather than 404ing — mirroring the integrated server's
|
||
// token-carried reconstruct. Returns nil when the request carries no usable
|
||
// transcode token, which the caller renders as a genuine not-found.
|
||
//
|
||
// requestedSegment is the segment the client is fetching, or negative on the
|
||
// manifest path. Reconstruction is single-flighted per session id so concurrent
|
||
// manifest and segment requests for the same lost session share one ffmpeg.
|
||
func (s *Server) reconstructFromToken(r *http.Request, sessionID string, requestedSegment int) *playback.TranscodeSession {
|
||
tokenStr := r.Header.Get("X-Silo-Stream-Token")
|
||
if tokenStr == "" {
|
||
return nil
|
||
}
|
||
cfg := s.watcher.Config()
|
||
claims, err := streamtoken.Verify(tokenStr, cfg.Auth.JWTSecret)
|
||
if err != nil {
|
||
slog.WarnContext(r.Context(), "transcode node reconstruct: invalid stream token", "component", "transcodenode", "error", err,
|
||
"session", sessionID, "playback_session_id", sessionID)
|
||
return nil
|
||
}
|
||
card := playback.RecipeCardFromClaims(claims)
|
||
// The token's recipe must be a transcode card for the session id in the URL: a
|
||
// mismatch is a forged or stale request, and direct/remux cards carry no encode
|
||
// parameters to rebuild. An empty PlayMethod is a transcode card (back-compat).
|
||
expectedTransportID := card.SessionID
|
||
if card.TranscodeTransportID != "" {
|
||
expectedTransportID = card.TranscodeTransportID
|
||
}
|
||
if expectedTransportID != sessionID || (card.PlayMethod != "" && card.PlayMethod != playback.PlayTranscode) {
|
||
return nil
|
||
}
|
||
// A native token carries the full byte-affecting recipe. The jellycompat node
|
||
// hop signs an identity-only token by design (see internal/noderecipe for why),
|
||
// so its card decodes with no encode parameters. For the jellycompat case the
|
||
// recipe is fetched from the control-plane recipe store below; without that
|
||
// store there is nothing to rebuild from, so 404.
|
||
tokenComplete := card.SegmentDuration > 0 && card.TargetCodecVideo != ""
|
||
if !tokenComplete && s.recipeStore == nil {
|
||
return nil
|
||
}
|
||
|
||
v, _, _ := s.reconstructGroup.Do(sessionID, func() (interface{}, error) {
|
||
// A concurrent reconstruct (or a fresh start) may already have registered the
|
||
// session; serve it rather than spawning a duplicate ffmpeg.
|
||
s.mu.RLock()
|
||
existing, ok := s.sessions[sessionID]
|
||
s.mu.RUnlock()
|
||
if ok {
|
||
return existing, nil
|
||
}
|
||
resolved := card
|
||
if !tokenComplete {
|
||
// Recipe-less (jellycompat) token: fetch the recipe central wrote to the
|
||
// control-plane store at transcode start. A miss / incomplete recipe is a
|
||
// genuine not-found (404), never a spawn from a bad recipe.
|
||
fetched, ok := s.recipeStore.Get(r.Context(), sessionID)
|
||
if !ok || fetched == nil || fetched.SessionID != sessionID ||
|
||
fetched.SegmentDuration <= 0 || fetched.TargetCodecVideo == "" {
|
||
return (*playback.TranscodeSession)(nil), nil
|
||
}
|
||
resolved = *fetched
|
||
}
|
||
return s.spawnReconstruct(r, sessionID, requestedSegment, resolved), nil
|
||
})
|
||
if session, _ := v.(*playback.TranscodeSession); session != nil {
|
||
return session
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// spawnReconstruct re-spawns ffmpeg for a lost session from its recipe card and
|
||
// registers it in the live map. It is only ever called inside the per-session
|
||
// single-flight in reconstructFromToken, so it is the sole writer racing to
|
||
// register sessionID. Returns nil if the spawn fails or the slot wait is canceled.
|
||
func (s *Server) spawnReconstruct(r *http.Request, sessionID string, requestedSegment int, card playback.RecipeCard) *playback.TranscodeSession {
|
||
// Pace the cold-start burst so a node restart that loses many sessions does not
|
||
// launch every ffmpeg at once. A client that disconnects while waiting releases
|
||
// its slot rather than queueing dead work.
|
||
release, ok := s.acquireReconstructSlot(r.Context())
|
||
if !ok {
|
||
return nil
|
||
}
|
||
defer release()
|
||
|
||
// Serialize against a concurrent fresh /transcode/start for this session so the
|
||
// two never run ffmpeg writers against the same dir. Re-check under the lock and
|
||
// yield to any live session rather than spawning a duplicate.
|
||
unlock := s.lockSessionLifecycle(sessionID)
|
||
defer unlock()
|
||
s.mu.RLock()
|
||
existing, ok := s.sessions[sessionID]
|
||
s.mu.RUnlock()
|
||
if ok {
|
||
return existing
|
||
}
|
||
|
||
cfg := s.watcher.Config()
|
||
outputDir := filepath.Join(cfg.Playback.TranscodeDir, sessionID)
|
||
opts := card.TranscodeOpts(outputDir, cfg.Playback.FFmpegPath, s.ffmpegSink)
|
||
opts.SessionID = sessionID
|
||
// Re-resolve environment-specific encode knobs from this node's live config; the
|
||
// token deliberately omits HWAccel/HWDevice so an operator change applies on
|
||
// rebuild. Run as a transcode node, not integrated (card.TranscodeOpts defaults).
|
||
opts.HWAccel = cfg.Playback.HWAccel
|
||
opts.HWDevice = cfg.Playback.HWDevice
|
||
opts.NodeType = "transcode"
|
||
opts.ExecutionMode = "transcode_node"
|
||
|
||
// Resume near the segment the client is actually requesting. The card records
|
||
// the original start; if the client has played past it, spawning at the old
|
||
// position forces a wait-then-seek stall. A negative requestedSegment (manifest
|
||
// path) carries no segment context, so the card position stands.
|
||
//
|
||
// The fast seg×dur mapping is only valid for ENCODED transcodes, whose forced
|
||
// keyframes make every segment exactly SegmentDuration long. Copy-mode segments
|
||
// have variable durations, so seg×dur points at the wrong source time and causes
|
||
// multi-second A/V desync after a restart. For copy-mode cards leave the card's
|
||
// original start untouched and let the segment-recovery machinery seek forward
|
||
// once the manifest is rebuilt. This mirrors doReconstructTranscode in
|
||
// internal/playback/transcode_manager.go so both reconstruct paths stay consistent.
|
||
if requestedSegment > card.StartSegmentNumber && card.SegmentDuration > 0 &&
|
||
!strings.EqualFold(card.TargetCodecVideo, "copy") {
|
||
opts.StartSegmentNumber = requestedSegment
|
||
opts.SeekSeconds = float64(requestedSegment * card.SegmentDuration)
|
||
}
|
||
|
||
session, err := playback.StartTranscode(context.WithoutCancel(r.Context()), opts)
|
||
if err != nil {
|
||
slog.ErrorContext(r.Context(), "transcode node reconstruct start failed", "component", "transcodenode", "error", err,
|
||
"session", sessionID, "playback_session_id", sessionID)
|
||
return nil
|
||
}
|
||
|
||
// Yield to a winner registered by another path; close only the duplicate ffmpeg,
|
||
// never the shared output directory the winner is actively serving.
|
||
s.mu.Lock()
|
||
if existing, ok := s.sessions[sessionID]; ok {
|
||
s.mu.Unlock()
|
||
_ = session.CloseProcess()
|
||
return existing
|
||
}
|
||
s.sessions[sessionID] = session
|
||
s.noteSessionAccessLocked(sessionID)
|
||
s.mu.Unlock()
|
||
s.activeJobs.Add(1)
|
||
|
||
trackCtx := context.WithoutCancel(r.Context())
|
||
go s.tracker.Track(trackCtx, nodesessions.SessionInfo{
|
||
SessionID: sessionID,
|
||
NodeURL: s.tracker.NodeURL(),
|
||
NodeName: s.tracker.NodeName(),
|
||
Type: "transcode",
|
||
CodecVideo: card.TargetCodecVideo,
|
||
CodecAudio: card.TargetCodecAudio,
|
||
Resolution: card.TargetResolution,
|
||
HWAccel: session.Opts().HWAccel,
|
||
StartedAt: time.Now().UTC().Format(time.RFC3339),
|
||
AuthUserID: card.UserID,
|
||
ProfileID: card.ProfileID,
|
||
MediaFileID: card.MediaFileID,
|
||
})
|
||
|
||
slog.InfoContext(r.Context(), "transcode node session reconstructed from token", "component", "transcodenode",
|
||
"session", sessionID, "playback_session_id", sessionID,
|
||
"requested_segment", requestedSegment, "start_segment_number", opts.StartSegmentNumber)
|
||
return session
|
||
}
|
||
|
||
// acquireReconstructSlot blocks until a reconstruct slot is free or the request
|
||
// context is canceled, returning a release func and true on success. The semaphore
|
||
// is lazily sized to NumCPU so a node restart paces its ffmpeg cold starts.
|
||
func (s *Server) acquireReconstructSlot(ctx context.Context) (func(), bool) {
|
||
s.reconstructSemOnce.Do(func() {
|
||
n := runtime.NumCPU()
|
||
if n < 1 {
|
||
n = 4
|
||
}
|
||
s.reconstructSem = make(chan struct{}, n)
|
||
})
|
||
select {
|
||
case s.reconstructSem <- struct{}{}:
|
||
return func() { <-s.reconstructSem }, true
|
||
case <-ctx.Done():
|
||
return nil, false
|
||
}
|
||
}
|
||
|
||
func (s *Server) handleStop(w http.ResponseWriter, r *http.Request) {
|
||
sessionID := chi.URLParam(r, "session_id")
|
||
|
||
// Serialize against in-flight starts and reconstructs. A RequireReady
|
||
// start registers its job only after the readiness wait; a stop racing
|
||
// that wait (the API rolling back a timed-out start) would otherwise miss
|
||
// the map and 404, orphaning the just-spawned ffmpeg until the idle
|
||
// reaper finds it. Blocking here until the start registers turns that
|
||
// miss into a normal teardown.
|
||
unlock := s.lockSessionLifecycle(sessionID)
|
||
defer unlock()
|
||
|
||
s.mu.Lock()
|
||
session, ok := s.sessions[sessionID]
|
||
if !ok {
|
||
s.mu.Unlock()
|
||
http.Error(w, "session not found", http.StatusNotFound)
|
||
return
|
||
}
|
||
delete(s.sessions, sessionID)
|
||
delete(s.lastAccess, sessionID)
|
||
s.mu.Unlock()
|
||
s.activeJobs.Add(-1)
|
||
|
||
if err := session.Close(); err != nil {
|
||
slog.ErrorContext(r.Context(), "close transcode session", "component", "transcodenode", "error", err, "session", sessionID, "playback_session_id", sessionID)
|
||
}
|
||
|
||
cfg := s.watcher.Config()
|
||
outputDir := filepath.Join(cfg.Playback.TranscodeDir, sessionID)
|
||
os.RemoveAll(outputDir)
|
||
|
||
// Drop the recipe so a buffered/retrying request after a node restart cannot
|
||
// reconstruct a new ffmpeg for this now-stopped session. Best-effort: a stop
|
||
// must still succeed even if the recipe store is briefly unavailable.
|
||
if s.recipeStore != nil {
|
||
if err := s.recipeStore.Delete(r.Context(), sessionID); err != nil {
|
||
slog.WarnContext(r.Context(), "delete transcode recipe on stop", "component", "transcodenode", "error", err, "session", sessionID, "playback_session_id", sessionID)
|
||
}
|
||
}
|
||
|
||
s.tracker.Remove(r.Context(), sessionID)
|
||
|
||
w.WriteHeader(http.StatusNoContent)
|
||
}
|
||
|
||
func (s *Server) handleManifest(w http.ResponseWriter, r *http.Request) {
|
||
sessionID := chi.URLParam(r, "session_id")
|
||
|
||
// Lookup and liveness refresh happen atomically so the idle reaper can
|
||
// never unregister the job between them and tear down a session this
|
||
// request is about to serve from.
|
||
session, ok := s.acquireSessionTouched(sessionID)
|
||
if !ok {
|
||
// Lost the in-memory session (this node restarted): rebuild it from the
|
||
// stream token the proxy forwarded. The manifest path carries no segment
|
||
// context, so reconstruct at the recipe's original start position.
|
||
session = s.reconstructFromToken(r, sessionID, -1)
|
||
if session == nil {
|
||
http.Error(w, "session not found", http.StatusNotFound)
|
||
return
|
||
}
|
||
// A reconstruct that yielded to a concurrently registered winner has
|
||
// not recorded this hit; count it so the reaper sees the liveness.
|
||
s.touchSession(sessionID)
|
||
}
|
||
|
||
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) {
|
||
if err := s.watcher.ForceReload(r.Context()); err != nil {
|
||
http.Error(w, "reload failed: "+err.Error(), http.StatusInternalServerError)
|
||
return
|
||
}
|
||
cfg := s.watcher.Config()
|
||
s.mu.Lock()
|
||
stopped := make([]string, 0, len(s.sessions))
|
||
for id, session := range s.sessions {
|
||
session.Close()
|
||
os.RemoveAll(filepath.Join(cfg.Playback.TranscodeDir, id))
|
||
delete(s.sessions, id)
|
||
delete(s.lastAccess, id)
|
||
stopped = append(stopped, id)
|
||
}
|
||
s.activeJobs.Store(0)
|
||
s.mu.Unlock()
|
||
|
||
// A force-reload tears every session down for good, so drop their recipes too:
|
||
// otherwise a buffered/retrying request could reconstruct a session this reload
|
||
// deliberately killed. Best-effort, done outside the map lock.
|
||
if s.recipeStore != nil {
|
||
for _, id := range stopped {
|
||
if err := s.recipeStore.Delete(r.Context(), id); err != nil {
|
||
slog.WarnContext(r.Context(), "delete transcode recipe on force reload", "component", "transcodenode", "error", err, "session", id, "playback_session_id", id)
|
||
}
|
||
}
|
||
}
|
||
|
||
s.tracker.Cleanup(r.Context())
|
||
|
||
slog.InfoContext(r.Context(), "transcode force reload completed", slog.String("component", "transcodenode"))
|
||
w.WriteHeader(http.StatusNoContent)
|
||
}
|
||
|
||
func (s *Server) handleStatus(w http.ResponseWriter, r *http.Request) {
|
||
s.mu.RLock()
|
||
sessionIDs := make([]string, 0, len(s.sessions))
|
||
for id := range s.sessions {
|
||
sessionIDs = append(sessionIDs, id)
|
||
}
|
||
s.mu.RUnlock()
|
||
|
||
w.Header().Set("Content-Type", "application/json")
|
||
type statusResponse struct {
|
||
Status string `json:"status"`
|
||
ActiveJobs int32 `json:"active_jobs"`
|
||
Sessions []string `json:"sessions"`
|
||
}
|
||
json.NewEncoder(w).Encode(statusResponse{
|
||
Status: "ok",
|
||
ActiveJobs: s.activeJobs.Load(),
|
||
Sessions: sessionIDs,
|
||
})
|
||
}
|