From 520c0613acb221d17c5fb403f9034c63aef74426 Mon Sep 17 00:00:00 2001 From: Rhainland Date: Tue, 29 Sep 2026 13:52:31 +0200 Subject: [PATCH 1/4] fix(downloads): keep every audio track and subtitle in prepared downloads Remux and transcode downloads kept only the first audio track and dropped all embedded subtitles. Prepared MP4s now carry every audio track and the plain-text subtitles as MP4 timed text; ASS/SSA and PGS tracks are offered as manifest sidecar files extracted from the source through the shared subtitle cache. A track recipe version and tracks_v1 queue states keep older workers and transcode nodes from producing the legacy layout. --- contracts/api/v2/openapi.json | 88 ++++++ docs/architecture/playback-protocol-v3.md | 9 +- docs/downloads-api.md | 33 ++- internal/api/router.go | 3 + internal/apiv2/download_delivery.go | 9 +- internal/downloadprepare/transport.go | 28 +- internal/downloads/artifact.go | 19 +- internal/downloads/artifact_repo.go | 47 ++-- internal/downloads/artifact_repo_test.go | 100 +++++++ internal/downloads/artifact_test.go | 37 +++ internal/downloads/artifacts.go | 18 +- internal/downloads/manifest.go | 81 +++++- internal/downloads/manifest_test.go | 113 ++++++++ internal/downloads/offline.go | 60 +++- internal/downloads/remote_preparer.go | 43 ++- internal/downloads/remote_preparer_test.go | 61 ++++ internal/downloads/repo.go | 4 +- internal/downloads/service.go | 5 + internal/lang/lang.go | 19 ++ internal/lang/lang_test.go | 8 + internal/playback/prepare_file.go | 13 +- internal/playback/prepare_tracks.go | 254 +++++++++++++++++ internal/playback/prepare_tracks_test.go | 262 ++++++++++++++++++ internal/playback/protocol_v3.go | 3 + internal/playback/transcode.go | 6 +- internal/transcodenode/server.go | 8 + internal/workmetrics/queues.go | 2 +- ...30_fence_track_recipe_artifact_workers.sql | 163 +++++++++++ web/src/api/v2/schema.ts | 1 + 29 files changed, 1427 insertions(+), 70 deletions(-) create mode 100644 internal/playback/prepare_tracks.go create mode 100644 internal/playback/prepare_tracks_test.go create mode 100644 migrations/sql/20260928222330_fence_track_recipe_artifact_workers.sql diff --git a/contracts/api/v2/openapi.json b/contracts/api/v2/openapi.json index f93ca43b08..633adfc0a6 100644 --- a/contracts/api/v2/openapi.json +++ b/contracts/api/v2/openapi.json @@ -31050,6 +31050,9 @@ }, "language": { "type": "string" + }, + "title": { + "type": "string" } }, "required": [ @@ -194093,6 +194096,9 @@ "input_path": { "type": "string" }, + "prepared_tracks": { + "$ref": "#/x-silo-worker-protocols/schemas/playback_PreparedTracks" + }, "software_video_decode": { "type": "boolean" }, @@ -194164,6 +194170,9 @@ "total_duration": { "format": "double", "type": "number" + }, + "track_recipe_version": { + "type": "string" } }, "required": [ @@ -194828,6 +194837,85 @@ ], "type": "object" }, + "playback_PreparedAudioTrack": { + "additionalProperties": false, + "properties": { + "codec": { + "type": "string" + }, + "default": { + "type": "boolean" + }, + "language": { + "type": "string" + }, + "source_channels": { + "format": "int64", + "type": "integer" + }, + "source_index": { + "format": "int64", + "type": "integer" + }, + "title": { + "type": "string" + } + }, + "required": [ + "source_index", + "codec" + ], + "type": "object" + }, + "playback_PreparedSubtitleTrack": { + "additionalProperties": false, + "properties": { + "default": { + "type": "boolean" + }, + "forced": { + "type": "boolean" + }, + "hearing_impaired": { + "type": "boolean" + }, + "language": { + "type": "string" + }, + "source_index": { + "format": "int64", + "type": "integer" + }, + "title": { + "type": "string" + } + }, + "required": [ + "source_index" + ], + "type": "object" + }, + "playback_PreparedTracks": { + "additionalProperties": false, + "properties": { + "audio": { + "items": { + "$ref": "#/x-silo-worker-protocols/schemas/playback_PreparedAudioTrack" + }, + "type": "array" + }, + "subtitles": { + "items": { + "$ref": "#/x-silo-worker-protocols/schemas/playback_PreparedSubtitleTrack" + }, + "type": "array" + } + }, + "required": [ + "audio" + ], + "type": "object" + }, "playback_RenderDeviceInfo": { "additionalProperties": false, "properties": { diff --git a/docs/architecture/playback-protocol-v3.md b/docs/architecture/playback-protocol-v3.md index 2fbba144e3..bd5d20c1de 100644 --- a/docs/architecture/playback-protocol-v3.md +++ b/docs/architecture/playback-protocol-v3.md @@ -1387,7 +1387,14 @@ session updates preserve codec, source/target channels, bitrate, and the transcode decision as one recipe. A failed Jellyfin audio switch restores the prior durable selection and executor facts so the same client report can retry. Prepared downloads persist the audio recipe version and use `audio_v2_*` queue -states that pre-v2 API workers cannot claim or publish as ready. +states that pre-v2 API workers cannot claim or publish as ready. The multi-track +prepared layout (every audio track, plain-text subtitles as MP4 timed text, +ASS/SSA and PGS as manifest sidecars) is a +separate `track_recipe_version` with `tracks_v1_*` queue states that outrank the +audio and tone-map families. Its per-track plan travels in the prepare request +and execution fingerprint; only transcode nodes advertising the +`prepared_tracks_v1` transport feature receive it, and an older node's legacy +receipt is rejected. They are advertised only if an eligible executor actually has the required capability. The ordinary FFmpeg feature probe is cached; the more expensive diff --git a/docs/downloads-api.md b/docs/downloads-api.md index 499b56576e..dfe7c3c00d 100644 --- a/docs/downloads-api.md +++ b/docs/downloads-api.md @@ -128,8 +128,23 @@ Yes, manifests include metadata needed to make the offline item feel native: - External and downloaded subtitle fetch URLs plus known subtitle file sizes. - Container, codecs, resolution, HDR, duration, selected audio track, and audio track inventory. For remux/transcode entries these describe the prepared - artifact the file endpoint actually delivers (single audio track, target - container/codecs), not the catalog source it was prepared from. + artifact the file endpoint actually delivers (target container/codecs), not + the catalog source it was prepared from. + +Prepared remux/transcode files keep every source audio track in source order, +so `audio_tracks[].index` and `selected_audio_track_index` address positions +in the delivered MP4. Transcodes encode each track to stereo AAC; remuxes copy +tracks that share the primary track's codec (or are AAC/MP3) and encode the +rest to stereo AAC. Embedded plain-text subtitles (SRT, WebVTT) are carried +inside the MP4 as timed text, with their language, title, and forced flag. +MP4 timed text would drop ASS/SSA styling, drawing commands, and overlapping +events, and MP4 cannot store bitmap subtitles, so each embedded ASS/SSA track is +listed in `subtitles[]` as an `ass` sidecar and each PGS track as a `sup` +sidecar, with the track's `title` when it has one; DVD and DVB bitmap subtitles are not carried. MP4 marks the first +embedded subtitle track as default, so clients choose subtitles from forced +flags and viewer preference rather than that flag. +Files prepared before this layout contain only the first audio track and no +subtitles, and their manifests keep describing them that way. - Stable provider identity and integrity metadata for local validation/rescan recovery. The client still needs to fetch artwork/subtitle bytes once while online and cache @@ -613,10 +628,14 @@ Other failures, such as `404 not_found`, won't succeed on a retry. GET /api/v2/downloads/{id}/subtitles/{ref} ``` -`ref` comes from `subtitles[].fetch_url` and encodes either `external:{index}` or -`downloaded:{id}`; `X-Silo-Device-Id` is required. Invalid refs return -`422 validation_failed`. Current content access is checked before asset delivery, -and downloaded-subtitle ownership must match the entry's media file. +`ref` comes from `subtitles[].fetch_url` and encodes `external:{index}`, +`embedded:{ordinal}`, or `downloaded:{id}`; `X-Silo-Device-Id` is required. +`embedded` refs name an embedded ASS/SSA or PGS track by subtitle ordinal and +return the complete track, as an ASS script or a `.sup` elementary stream, +extracted from the source file. +Invalid refs return `422 validation_failed`. Current content access is checked +before asset delivery, and downloaded-subtitle ownership must match the entry's +media file. ### 4.10 Direct download @@ -1598,7 +1617,7 @@ unavailable or ineligible proxy targets fall back to existing local delivery. `GET /api/v2/downloads/{id}/artwork/{kind}` and `GET /api/v2/downloads/{id}/subtitles/{ref}` require the device header. Artwork kinds are poster, backdrop and logo; subtitle references retain the -existing external:index and downloaded:id identity. Current content access is +external:index, embedded:ordinal, and downloaded:id identity. Current content access is checked before asset delivery, and downloaded subtitle ownership must match the entry's media file. These two asset routes preserve whole-object delivery and private caching; they do not advertise byte ranges. diff --git a/internal/api/router.go b/internal/api/router.go index fd16b47a2b..49032ac0bb 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -2005,6 +2005,9 @@ func newChiRouter(deps Dependencies) chi.Router { } downloadSvc.SetOfflineDeps(detailSvc, subtitleSource, nil) } + if streamHandler != nil { + downloadSvc.SetSubtitleCache(streamHandler.SubtitleCache) + } if deps.MarkerPopulation != nil { downloadSvc.SetMarkerPopulation(deps.MarkerPopulation) } diff --git a/internal/apiv2/download_delivery.go b/internal/apiv2/download_delivery.go index 3530954b3b..caf4aa9aa2 100644 --- a/internal/apiv2/download_delivery.go +++ b/internal/apiv2/download_delivery.go @@ -105,7 +105,7 @@ func (reg *Registry) serveDownloadDelivery(w http.ResponseWriter, r *http.Reques // written rather than before the first read, so that failure still gets // its problem response. Files keep ReadFrom for sendfile; ServeContent // writes their header before copying anyway. - asset := struct{ http.ResponseWriter }{writer} + asset := assetResponseWriter{writer} var err error switch kind { case "file": @@ -125,3 +125,10 @@ func (reg *Registry) serveDownloadDelivery(w http.ResponseWriter, r *http.Reques } writeProblem(w, r, downloadProblem(err)) } + +// assetResponseWriter hides ReadFrom from artwork and subtitle copies (see +// serveDownloadDelivery) while keeping Unwrap, so response controllers can +// still reach the connection for rolling write deadlines. +type assetResponseWriter struct{ http.ResponseWriter } + +func (w assetResponseWriter) Unwrap() http.ResponseWriter { return w.ResponseWriter } diff --git a/internal/downloadprepare/transport.go b/internal/downloadprepare/transport.go index 402077a5df..10aeb93bba 100644 --- a/internal/downloadprepare/transport.go +++ b/internal/downloadprepare/transport.go @@ -98,8 +98,13 @@ type Request struct { AudioRecipeVersion string `json:"audio_recipe_version,omitempty"` // SourceAudioChannels freezes the selected input stream's probed channel // count. Zero is the mixed-version-safe unknown value and never enables gain. - SourceAudioChannels int `json:"source_audio_channels,omitempty"` - TotalDuration float64 `json:"total_duration,omitempty"` + SourceAudioChannels int `json:"source_audio_channels,omitempty"` + // TrackRecipeVersion and PreparedTracks carry the multi-track stream + // layout. Older nodes ignore both, encode the legacy single-audio layout, + // and therefore cannot return the matching execution fingerprint. + TrackRecipeVersion string `json:"track_recipe_version,omitempty"` + PreparedTracks *playback.PreparedTracks `json:"prepared_tracks,omitempty"` + TotalDuration float64 `json:"total_duration,omitempty"` } // Result identifies a completed artifact without exposing the node's local @@ -194,11 +199,23 @@ func (r Request) StereoDownmixBoostRequested() bool { playback.IsAudioToAACStereoDownmixV3(r.SourceAudioChannels, r.TargetCodecAudio, r.TargetAudioChannels) } +// PreparedTracksRequested includes incomplete layouts so the node can reject a +// partial recipe instead of encoding the legacy single-audio layout. +func (r Request) PreparedTracksRequested() bool { + return r.TrackRecipeVersion != "" || r.PreparedTracks != nil +} + +// ValidPreparedTracks reports whether a requested layout is complete and uses +// the stream-layout version this build executes. +func (r Request) ValidPreparedTracks() bool { + return r.TrackRecipeVersion == playback.PreparedTracksRecipeVersion && r.PreparedTracks != nil +} + // ExecutionAttestationRequested reports whether accepting bytes requires a // receipt from a node that understood all newly transported recipe fields. // Explicit audio output settings affect bytes even when the v2 boost does not. func (r Request) ExecutionAttestationRequested() bool { - return r.ToneMapRequested() || r.AudioRecipeRequested() || + return r.ToneMapRequested() || r.AudioRecipeRequested() || r.PreparedTracksRequested() || r.TargetAudioChannels != 0 || r.TargetAudioBitrateKbps != 0 } @@ -252,6 +269,10 @@ func NewRequest(artifactID string, opts playback.TranscodeOpts) Request { AudioTrackIndex: opts.AudioTrackIndex, TotalDuration: opts.TotalDuration, } + if opts.PreparedTracks != nil { + request.TrackRecipeVersion = playback.PreparedTracksRecipeVersion + request.PreparedTracks = opts.PreparedTracks + } if playback.IsAudioToAACStereoDownmixV3(opts.SourceAudioChannels, request.TargetCodecAudio, request.TargetAudioChannels) { request.SourceAudioChannels = opts.SourceAudioChannels request.AudioRecipeVersion = playback.TransformationAudioToAACRecipeVersionV3 @@ -287,6 +308,7 @@ func (r Request) TranscodeOpts(ffmpegPath, hwAccel, hwDevice string, sink playba AudioTrackIndex: r.AudioTrackIndex, SourceAudioChannels: r.SourceAudioChannels, SubtitleTrackIndex: -1, + PreparedTracks: r.PreparedTracks, FFmpegPath: ffmpegPath, HWAccel: hwAccel, HWDevice: hwDevice, diff --git a/internal/downloads/artifact.go b/internal/downloads/artifact.go index cba931e952..a77d713b3f 100644 --- a/internal/downloads/artifact.go +++ b/internal/downloads/artifact.go @@ -22,11 +22,17 @@ const ( ArtifactAudioV2Queued = "audio_v2_queued" ArtifactAudioV2Running = "audio_v2_running" ArtifactAudioV2Ready = "audio_v2_ready" + ArtifactTracksQueued = "tracks_v1_queued" + ArtifactTracksRunning = "tracks_v1_running" + ArtifactTracksReady = "tracks_v1_ready" ArtifactReady = "ready" ArtifactFailed = "failed" ) -func queuedArtifactStatus(mode tonemap.Mode, audioRecipeVersion string) string { +func queuedArtifactStatus(mode tonemap.Mode, audioRecipeVersion, trackRecipeVersion string) string { + if trackRecipeVersion != "" { + return ArtifactTracksQueued + } if audioRecipeVersion != "" { return ArtifactAudioV2Queued } @@ -37,7 +43,8 @@ func queuedArtifactStatus(mode tonemap.Mode, audioRecipeVersion string) string { } func artifactReady(artifact *Artifact) bool { - return artifact != nil && (artifact.Status == ArtifactReady || artifact.Status == ArtifactToneMapReady || artifact.Status == ArtifactAudioV2Ready) + return artifact != nil && (artifact.Status == ArtifactReady || artifact.Status == ArtifactToneMapReady || + artifact.Status == ArtifactAudioV2Ready || artifact.Status == ArtifactTracksReady) } // ErrNoArtifactJob is returned by the queue when no claimable job exists. @@ -54,6 +61,7 @@ type Artifact struct { CodecVideo string CodecAudio string AudioRecipeVersion string + TrackRecipeVersion string // playback.PreparedTracksRecipeVersion; empty = legacy single-audio layout Resolution string AudioTrackIndex int TargetBitrateKbps int @@ -129,13 +137,14 @@ func paramsHashWithToneMapRevision(params paramsHashParams) string { } // artifactUsesExecutionFingerprint distinguishes source-sensitive recipes from -// legacy parameter-only artifacts. AudioRecipeVersion is also the durable -// queue discriminator that keeps a pre-v2 worker from claiming these bytes. +// legacy parameter-only artifacts. AudioRecipeVersion and TrackRecipeVersion +// are also the durable queue discriminators that keep an older worker from +// claiming bytes it would encode differently. func artifactUsesExecutionFingerprint(a *Artifact) bool { if a == nil { return false } - return a.ToneMapMode != "" || a.AudioRecipeVersion != "" + return a.ToneMapMode != "" || a.AudioRecipeVersion != "" || a.TrackRecipeVersion != "" } // effectiveArtifactDir resolves where prepared artifacts are written: the diff --git a/internal/downloads/artifact_repo.go b/internal/downloads/artifact_repo.go index c7f078194b..f967208940 100644 --- a/internal/downloads/artifact_repo.go +++ b/internal/downloads/artifact_repo.go @@ -12,7 +12,7 @@ import ( "github.com/Silo-Server/silo-server/internal/tonemap" ) -const artifactColumns = `id, media_file_id, format, params_hash, container, codec_video, codec_audio, audio_recipe_version, +const artifactColumns = `id, media_file_id, format, params_hash, container, codec_video, codec_audio, audio_recipe_version, track_recipe_version, resolution, audio_track_index, target_bitrate_kbps, tone_map_policy, tone_map_mode, tone_map_source_kind, tone_map_recipe_version, tone_map_preflight_required, tone_map_source_revision, tone_map_dv_config_present, tone_map_dv_bl_compat_id_present, tone_map_dv_bl_present, tone_map_dv_rpu_present, output_path, origin_node_id, origin_node_url, origin_node_group, origin_artifact_id, file_size, status, error_message, @@ -47,7 +47,7 @@ func scanArtifact(row pgx.Row) (*Artifact, error) { var a Artifact var leaseOwner *string if err := row.Scan( - &a.ID, &a.MediaFileID, &a.Format, &a.ParamsHash, &a.Container, &a.CodecVideo, &a.CodecAudio, &a.AudioRecipeVersion, + &a.ID, &a.MediaFileID, &a.Format, &a.ParamsHash, &a.Container, &a.CodecVideo, &a.CodecAudio, &a.AudioRecipeVersion, &a.TrackRecipeVersion, &a.Resolution, &a.AudioTrackIndex, &a.TargetBitrateKbps, &a.ToneMapPolicy, &a.ToneMapMode, &a.ToneMapSourceKind, &a.ToneMapRecipeVersion, &a.ToneMapPreflightRequired, &a.ToneMapSourceRevision, &a.ToneMapDVConfigPresent, &a.ToneMapDVBLCompatIDPresent, &a.ToneMapDVBLPresent, &a.ToneMapDVRPUPresent, &a.OutputPath, &a.OriginNodeID, &a.OriginNodeURL, &a.OriginNodeGroup, &a.OriginArtifactID, &a.FileSize, &a.Status, &a.ErrorMessage, @@ -67,15 +67,15 @@ func (r *ArtifactRepository) EnsureQueued(ctx context.Context, a *Artifact) (*Ar if a.ToneMapPolicy == "" { a.ToneMapPolicy = tonemap.PolicyNone } - status := queuedArtifactStatus(a.ToneMapMode, a.AudioRecipeVersion) + status := queuedArtifactStatus(a.ToneMapMode, a.AudioRecipeVersion, a.TrackRecipeVersion) tag, err := r.pool.Exec(ctx, `INSERT INTO download_artifacts - (id, media_file_id, format, params_hash, container, codec_video, codec_audio, audio_recipe_version, + (id, media_file_id, format, params_hash, container, codec_video, codec_audio, audio_recipe_version, track_recipe_version, resolution, audio_track_index, target_bitrate_kbps, tone_map_policy, tone_map_mode, tone_map_source_kind, tone_map_recipe_version, tone_map_preflight_required, tone_map_source_revision, tone_map_dv_config_present, tone_map_dv_bl_compat_id_present, tone_map_dv_bl_present, tone_map_dv_rpu_present, output_path, status, max_attempts) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25) ON CONFLICT (media_file_id, format, params_hash) DO NOTHING`, - a.ID, a.MediaFileID, a.Format, a.ParamsHash, a.Container, a.CodecVideo, a.CodecAudio, a.AudioRecipeVersion, + a.ID, a.MediaFileID, a.Format, a.ParamsHash, a.Container, a.CodecVideo, a.CodecAudio, a.AudioRecipeVersion, a.TrackRecipeVersion, a.Resolution, a.AudioTrackIndex, a.TargetBitrateKbps, a.ToneMapPolicy, a.ToneMapMode, a.ToneMapSourceKind, a.ToneMapRecipeVersion, a.ToneMapPreflightRequired, a.ToneMapSourceRevision, a.ToneMapDVConfigPresent, a.ToneMapDVBLCompatIDPresent, a.ToneMapDVBLPresent, a.ToneMapDVRPUPresent, a.OutputPath, status, a.MaxAttempts, ) @@ -129,6 +129,7 @@ func (r *ArtifactRepository) ClaimNext(ctx context.Context, owner string, lease a, err := scanArtifact(r.pool.QueryRow(ctx, `UPDATE download_artifacts SET status = CASE + WHEN status IN ('tracks_v1_queued', 'tracks_v1_running') THEN 'tracks_v1_running' WHEN status IN ('audio_v2_queued', 'audio_v2_running') THEN 'audio_v2_running' WHEN status IN ('tone_map_queued', 'tone_map_running') THEN 'tone_map_running' ELSE 'running' @@ -137,8 +138,8 @@ func (r *ArtifactRepository) ClaimNext(ctx context.Context, owner string, lease attempts = attempts + 1 WHERE id = ( SELECT id FROM download_artifacts - WHERE (status IN ('queued', 'tone_map_queued', 'audio_v2_queued') AND (next_retry_at IS NULL OR next_retry_at <= now())) - OR (status IN ('running', 'tone_map_running', 'audio_v2_running') AND lease_expires_at < now()) + WHERE (status IN ('queued', 'tone_map_queued', 'audio_v2_queued', 'tracks_v1_queued') AND (next_retry_at IS NULL OR next_retry_at <= now())) + OR (status IN ('running', 'tone_map_running', 'audio_v2_running', 'tracks_v1_running') AND lease_expires_at < now()) ORDER BY created_at LIMIT 1 FOR UPDATE SKIP LOCKED @@ -164,7 +165,7 @@ func (r *ArtifactRepository) Heartbeat(ctx context.Context, id, owner string, le } tag, err := r.pool.Exec(ctx, `UPDATE download_artifacts SET lease_expires_at = now() + make_interval(secs => $3) - WHERE id = $1 AND lease_owner = $2 AND status IN ('running', 'tone_map_running', 'audio_v2_running')`, + WHERE id = $1 AND lease_owner = $2 AND status IN ('running', 'tone_map_running', 'audio_v2_running', 'tracks_v1_running')`, id, owner, leaseSecs, ) if err != nil { @@ -183,6 +184,7 @@ func (r *ArtifactRepository) MarkReady(ctx context.Context, id, owner, outputPat tag, err := r.pool.Exec(ctx, `UPDATE download_artifacts SET status = CASE + WHEN track_recipe_version <> '' THEN 'tracks_v1_ready' WHEN audio_recipe_version <> '' THEN 'audio_v2_ready' WHEN tone_map_mode <> '' THEN 'tone_map_ready' ELSE 'ready' @@ -191,7 +193,7 @@ func (r *ArtifactRepository) MarkReady(ctx context.Context, id, owner, outputPat origin_node_group = $5, origin_artifact_id = $6, file_size = $7, error_message = '', completed_at = now(), last_used_at = now(), lease_owner = NULL, lease_expires_at = NULL, next_retry_at = NULL - WHERE id = $1 AND lease_owner = $8 AND status IN ('running', 'tone_map_running', 'audio_v2_running')`, + WHERE id = $1 AND lease_owner = $8 AND status IN ('running', 'tone_map_running', 'audio_v2_running', 'tracks_v1_running')`, id, outputPath, originNodeID, originNodeURL, originNodeGroup, originArtifactID, fileSize, owner, ) if err != nil { @@ -214,6 +216,7 @@ func (r *ArtifactRepository) MarkFailedOrRetry(ctx context.Context, id, owner, e err = r.pool.QueryRow(ctx, `UPDATE download_artifacts SET status = CASE WHEN attempts >= max_attempts THEN 'failed' + WHEN status = 'tracks_v1_running' THEN 'tracks_v1_queued' WHEN status = 'audio_v2_running' THEN 'audio_v2_queued' WHEN status = 'tone_map_running' THEN 'tone_map_queued' ELSE 'queued' END, @@ -221,7 +224,7 @@ func (r *ArtifactRepository) MarkFailedOrRetry(ctx context.Context, id, owner, e next_retry_at = CASE WHEN attempts >= max_attempts THEN NULL ELSE now() + make_interval(secs => $3) END, completed_at = CASE WHEN attempts >= max_attempts THEN now() ELSE NULL END, lease_owner = NULL, lease_expires_at = NULL - WHERE id = $1 AND lease_owner = $4 AND status IN ('running', 'tone_map_running', 'audio_v2_running') + WHERE id = $1 AND lease_owner = $4 AND status IN ('running', 'tone_map_running', 'audio_v2_running', 'tracks_v1_running') RETURNING status = 'failed'`, id, errMsg, backoffSecs, owner, ).Scan(&terminal) @@ -247,13 +250,14 @@ func (r *ArtifactRepository) ReclaimExpiredLeases(ctx context.Context) ([]reclai rows, err := r.pool.Query(ctx, `UPDATE download_artifacts SET status = CASE WHEN attempts >= max_attempts THEN 'failed' + WHEN status = 'tracks_v1_running' THEN 'tracks_v1_queued' WHEN status = 'audio_v2_running' THEN 'audio_v2_queued' WHEN status = 'tone_map_running' THEN 'tone_map_queued' ELSE 'queued' END, lease_owner = NULL, lease_expires_at = NULL, error_message = CASE WHEN attempts >= max_attempts THEN 'exceeded max attempts after lease expiry' ELSE error_message END, completed_at = CASE WHEN attempts >= max_attempts THEN now() ELSE completed_at END - WHERE status IN ('running', 'tone_map_running', 'audio_v2_running') AND lease_expires_at < now() + WHERE status IN ('running', 'tone_map_running', 'audio_v2_running', 'tracks_v1_running') AND lease_expires_at < now() RETURNING id, status = 'failed'`, ) if err != nil { @@ -279,6 +283,7 @@ func (r *ArtifactRepository) Requeue(ctx context.Context, id string) error { tag, err := r.pool.Exec(ctx, `UPDATE download_artifacts SET status = CASE + WHEN track_recipe_version <> '' THEN 'tracks_v1_queued' WHEN audio_recipe_version <> '' THEN 'audio_v2_queued' WHEN tone_map_mode <> '' THEN 'tone_map_queued' ELSE 'queued' @@ -319,7 +324,7 @@ const missingArtifactRetireGrace = 10 * time.Minute // unusedReadyArtifactPredicate selects ready rows that no active download // references and that nothing has used within the grace interval ($2, seconds). // Completed rows count as active because they remain re-downloadable. -const unusedReadyArtifactPredicate = `a.status IN ('ready', 'tone_map_ready', 'audio_v2_ready') +const unusedReadyArtifactPredicate = `a.status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready') AND a.last_used_at < now() - make_interval(secs => $2) AND NOT EXISTS (SELECT 1 FROM downloads d WHERE d.artifact_id = a.id AND d.status NOT IN ('cancelled', 'failed', 'revoked'))` @@ -347,13 +352,14 @@ func (r *ArtifactRepository) RecoverMissing(ctx context.Context, id string, grac tag, err = tx.Exec(ctx, `UPDATE download_artifacts SET status = CASE + WHEN track_recipe_version <> '' THEN 'tracks_v1_queued' WHEN audio_recipe_version <> '' THEN 'audio_v2_queued' WHEN tone_map_mode <> '' THEN 'tone_map_queued' ELSE 'queued' END, attempts = 0, error_message = '', next_retry_at = NULL, lease_owner = NULL, lease_expires_at = NULL, completed_at = NULL - WHERE id = $1 AND status IN ('ready', 'tone_map_ready', 'audio_v2_ready')`, + WHERE id = $1 AND status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready')`, id, ) if err != nil { @@ -446,6 +452,7 @@ func (r *ArtifactRepository) requeueRemote(ctx context.Context, artifact *Artifa } query := `UPDATE download_artifacts SET status = CASE + WHEN track_recipe_version <> '' THEN 'tracks_v1_queued' WHEN audio_recipe_version <> '' THEN 'audio_v2_queued' WHEN tone_map_mode <> '' THEN 'tone_map_queued' ELSE 'queued' @@ -453,7 +460,7 @@ func (r *ArtifactRepository) requeueRemote(ctx context.Context, artifact *Artifa attempts = 0, error_message = '', next_retry_at = NULL, lease_owner = NULL, lease_expires_at = NULL, completed_at = NULL, origin_node_id = 0, origin_node_url = '', origin_node_group = '', origin_artifact_id = '' - WHERE id = $1 AND status IN ('ready', 'tone_map_ready', 'audio_v2_ready') AND origin_node_id = $2 AND origin_artifact_id = $3` + WHERE id = $1 AND status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready') AND origin_node_id = $2 AND origin_artifact_id = $3` args := []any{artifact.ID, artifact.OriginNodeID, artifact.OriginArtifactID} if fenceURL { query += ` AND origin_node_url = $4` @@ -509,7 +516,7 @@ func (r *ArtifactRepository) TouchLastUsed(ctx context.Context, id string) error func (r *ArtifactRepository) TouchReady(ctx context.Context, id string) (bool, error) { tag, err := r.pool.Exec(ctx, `UPDATE download_artifacts SET last_used_at = now() - WHERE id = $1 AND status IN ('ready', 'tone_map_ready', 'audio_v2_ready')`, id) + WHERE id = $1 AND status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready')`, id) if err != nil { return false, fmt.Errorf("touching ready artifact: %w", err) } @@ -525,7 +532,7 @@ func (r *ArtifactRepository) RefreshRemoteLocator(ctx context.Context, artifact tag, err := r.pool.Exec(ctx, `UPDATE download_artifacts SET origin_node_url = $2, origin_node_group = $3 - WHERE id = $1 AND status IN ('ready', 'tone_map_ready', 'audio_v2_ready') + WHERE id = $1 AND status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready') AND origin_node_id = $4 AND origin_artifact_id = $5`, artifact.ID, artifact.OriginNodeURL, artifact.OriginNodeGroup, artifact.OriginNodeID, artifact.OriginArtifactID, @@ -539,7 +546,7 @@ func (r *ArtifactRepository) RefreshRemoteLocator(ctx context.Context, artifact // ListReady returns ready artifacts ordered by least-recently-used first. func (r *ArtifactRepository) ListReady(ctx context.Context) ([]*Artifact, error) { rows, err := r.pool.Query(ctx, - `SELECT `+artifactColumns+` FROM download_artifacts WHERE status IN ('ready', 'tone_map_ready', 'audio_v2_ready') ORDER BY last_used_at ASC`) + `SELECT `+artifactColumns+` FROM download_artifacts WHERE status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready') ORDER BY last_used_at ASC`) if err != nil { return nil, fmt.Errorf("listing ready artifacts: %w", err) } @@ -563,7 +570,7 @@ func scanArtifacts(rows pgx.Rows) ([]*Artifact, error) { func (r *ArtifactRepository) TotalReadyBytes(ctx context.Context) (int64, error) { var total int64 if err := r.pool.QueryRow(ctx, - `SELECT COALESCE(SUM(file_size), 0) FROM download_artifacts WHERE status IN ('ready', 'tone_map_ready', 'audio_v2_ready')`).Scan(&total); err != nil { + `SELECT COALESCE(SUM(file_size), 0) FROM download_artifacts WHERE status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready')`).Scan(&total); err != nil { return 0, fmt.Errorf("summing ready artifacts: %w", err) } return total, nil @@ -606,7 +613,7 @@ func (r *ArtifactRepository) ListFailedBefore(ctx context.Context, cutoff time.T func (r *ArtifactRepository) ListUnlinkedReadyBefore(ctx context.Context, cutoff time.Time) ([]*Artifact, error) { rows, err := r.pool.Query(ctx, `SELECT `+artifactColumns+` FROM download_artifacts a - WHERE a.status IN ('ready', 'tone_map_ready', 'audio_v2_ready') AND a.last_used_at < $1 + WHERE a.status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready') AND a.last_used_at < $1 AND NOT EXISTS (SELECT 1 FROM downloads d WHERE d.artifact_id = a.id)`, cutoff) if err != nil { return nil, fmt.Errorf("listing unlinked artifacts: %w", err) diff --git a/internal/downloads/artifact_repo_test.go b/internal/downloads/artifact_repo_test.go index cf8384d3d1..963a683646 100644 --- a/internal/downloads/artifact_repo_test.go +++ b/internal/downloads/artifact_repo_test.go @@ -365,6 +365,106 @@ func TestAudioV2ArtifactQueueRejectsMergeBaseWorkers(t *testing.T) { } } +// TestTrackRecipeArtifactQueueRejectsMergeBaseWorkers proves that a +// multi-track prepared file is fenced from workers that would encode the +// legacy single-audio layout, including rows that also carry the audio-v2 and +// tone-map recipes the merge-base worker already understands. +func TestTrackRecipeArtifactQueueRejectsMergeBaseWorkers(t *testing.T) { + repo, pool, fileID := newArtifactTestRepo(t) + ctx := context.Background() + var trackRecipeColumn *string + if err := pool.QueryRow(ctx, ` + SELECT column_name::text + FROM information_schema.columns + WHERE table_schema = 'public' AND table_name = 'download_artifacts' AND column_name = 'track_recipe_version' + `).Scan(&trackRecipeColumn); err != nil && !errors.Is(err, pgx.ErrNoRows) { + t.Fatalf("check track recipe worker fence: %v", err) + } + if trackRecipeColumn == nil { + t.Skip("migration 20260928222330_fence_track_recipe_artifact_workers has not been applied") + } + + a := newArtifact(t, fileID, "hash-track-recipe-worker-fence") + a.AudioRecipeVersion = playback.TransformationAudioToAACRecipeVersionV3 + a.TrackRecipeVersion = playback.PreparedTracksRecipeVersion + row, created, err := repo.EnsureQueued(ctx, a) + if err != nil || !created || row.Status != ArtifactTracksQueued || row.TrackRecipeVersion != playback.PreparedTracksRecipeVersion { + t.Fatalf("EnsureQueued = (%+v, created=%v, %v), want new tracks queued row", row, created, err) + } + + // Exact queue predicate from the merge-base worker. + var legacyClaimID string + err = pool.QueryRow(ctx, ` + UPDATE download_artifacts + SET status = CASE + WHEN status IN ('audio_v2_queued', 'audio_v2_running') THEN 'audio_v2_running' + WHEN status IN ('tone_map_queued', 'tone_map_running') THEN 'tone_map_running' + ELSE 'running' + END + WHERE id = ( + SELECT id FROM download_artifacts + WHERE (status IN ('queued', 'tone_map_queued', 'audio_v2_queued') AND (next_retry_at IS NULL OR next_retry_at <= now())) + OR (status IN ('running', 'tone_map_running', 'audio_v2_running') AND lease_expires_at < now()) + ORDER BY created_at LIMIT 1 FOR UPDATE SKIP LOCKED + ) + RETURNING id + `).Scan(&legacyClaimID) + if !errors.Is(err, pgx.ErrNoRows) { + t.Fatalf("merge-base worker claim = id %q, err %v; want no row", legacyClaimID, err) + } + + claim, err := repo.ClaimNext(ctx, "current-worker", time.Minute) + if err != nil || claim.ID != row.ID || claim.Status != ArtifactTracksRunning { + t.Fatalf("current ClaimNext = (%+v, %v), want tracks running", claim, err) + } + terminal, applied, err := repo.MarkFailedOrRetry(ctx, row.ID, "current-worker", "retry", time.Second) + if err != nil || !applied || terminal { + t.Fatalf("MarkFailedOrRetry = (%v, %v, %v), want retry", terminal, applied, err) + } + retried, err := repo.GetByID(ctx, row.ID) + if err != nil || retried.Status != ArtifactTracksQueued { + t.Fatalf("retried = (%+v, %v), want tracks queued", retried, err) + } + if _, err := pool.Exec(ctx, `UPDATE download_artifacts SET next_retry_at = now() - interval '1 second' WHERE id = $1`, row.ID); err != nil { + t.Fatal(err) + } + if claim, err = repo.ClaimNext(ctx, "current-worker", time.Minute); err != nil || claim.Status != ArtifactTracksRunning { + t.Fatalf("final ClaimNext = (%+v, %v), want tracks running", claim, err) + } + if applied, err := repo.MarkReady(ctx, row.ID, "current-worker", row.OutputPath, 0, "", "", "", 4242); err != nil || !applied { + t.Fatalf("MarkReady = (%v, %v), want applied", applied, err) + } + ready, err := repo.GetByID(ctx, row.ID) + if err != nil || ready.Status != ArtifactTracksReady || !artifactReady(ready) { + t.Fatalf("tracks ready row = (%+v, %v), want tracks ready", ready, err) + } + if total, err := repo.TotalReadyBytes(ctx); err != nil || total < 4242 { + t.Fatalf("TotalReadyBytes = (%d, %v), want the tracks artifact counted", total, err) + } + + var legacyReadyID string + err = pool.QueryRow(ctx, `SELECT id FROM download_artifacts WHERE id = $1 AND status IN ('ready', 'tone_map_ready', 'audio_v2_ready')`, row.ID).Scan(&legacyReadyID) + if !errors.Is(err, pgx.ErrNoRows) { + t.Fatalf("merge-base ready reader saw tracks artifact %q, err %v", legacyReadyID, err) + } + + // A merge-base API requeues with its audio-v2 status expression; the + // trigger must restore the tracks fence before another poll. + if _, err := pool.Exec(ctx, `UPDATE download_artifacts SET status = 'audio_v2_queued' WHERE id = $1`, row.ID); err != nil { + t.Fatal(err) + } + legacyRequeue, err := repo.GetByID(ctx, row.ID) + if err != nil || legacyRequeue.Status != ArtifactTracksQueued { + t.Fatalf("legacy requeue = (%+v, %v), want database-normalized tracks queued", legacyRequeue, err) + } + if err := repo.Requeue(ctx, row.ID); err != nil { + t.Fatal(err) + } + if requeued, err := repo.GetByID(ctx, row.ID); err != nil || requeued.Status != ArtifactTracksQueued { + t.Fatalf("Requeue = (%+v, %v), want tracks queued", requeued, err) + } +} + // TestArtifactRetryUntilTerminal verifies attempt counting and backoff: a job // retries behind its backoff gate until max_attempts, then goes terminal-failed. func TestArtifactRetryUntilTerminal(t *testing.T) { diff --git a/internal/downloads/artifact_test.go b/internal/downloads/artifact_test.go index 27c32e8750..800a36cd0e 100644 --- a/internal/downloads/artifact_test.go +++ b/internal/downloads/artifact_test.go @@ -206,6 +206,43 @@ func TestPreparedAudioBoostFreezesSelectedSourceChannelsInExecutionFingerprint(t } } +func TestPreparedTracksFreezeStreamLayoutInExecutionFingerprint(t *testing.T) { + manager := &ArtifactManager{} + file := &models.MediaFile{ + ID: 42, FilePath: "/media/movie.mkv", Duration: 3600, CodecAudio: "ac3", + AudioTracks: []models.AudioTrack{{Codec: "ac3", Channels: 6}, {Codec: "aac", Channels: 2, Language: "ja"}}, + SubtitleTracks: []models.SubtitleTrack{{Codec: "subrip", Language: "en"}}, + } + artifact := &Artifact{ + ID: "artifact-tracks", MediaFileID: file.ID, Format: "transcode", + Container: "mp4", CodecVideo: "h264", CodecAudio: "aac", AudioTrackIndex: -1, + TrackRecipeVersion: playback.PreparedTracksRecipeVersion, + } + opts := manager.buildOpts(file, artifact) + if opts.PreparedTracks == nil || len(opts.PreparedTracks.Audio) != 2 || len(opts.PreparedTracks.Subtitles) != 1 { + t.Fatalf("PreparedTracks = %+v, want every audio track and the text subtitle", opts.PreparedTracks) + } + if opts.SourceAudioChannels != 0 { + t.Fatalf("SourceAudioChannels = %d, want the per-track layout to own downmix facts", opts.SourceAudioChannels) + } + artifact.ParamsHash = downloadprepare.NewRequest(artifact.ID, opts).ExecutionFingerprint() + if !artifactUsesExecutionFingerprint(artifact) || !artifactExecutionFingerprintMatches(artifact, opts) { + t.Fatal("track recipe was not protected by its execution fingerprint") + } + // A rescan that changes the source track list must not reuse the frozen file. + rescanned := *file + rescanned.AudioTracks = append(rescanned.AudioTracks, models.AudioTrack{Codec: "aac", Channels: 2}) + if artifactExecutionFingerprintMatches(artifact, manager.buildOpts(&rescanned, artifact)) { + t.Fatal("changed source tracks reused the frozen multi-track artifact") + } + + legacy := *artifact + legacy.TrackRecipeVersion = "" + if opts := manager.buildOpts(file, &legacy); opts.PreparedTracks != nil { + t.Fatalf("legacy artifact gained a multi-track layout: %+v", opts.PreparedTracks) + } +} + func TestPreparedSourceAudioChannelsRequiresAACSurroundToDefaultStereo(t *testing.T) { file := &models.MediaFile{AudioTracks: []models.AudioTrack{{Channels: 2}, {Channels: 6}}} for _, test := range []struct { diff --git a/internal/downloads/artifacts.go b/internal/downloads/artifacts.go index 726b92f0c8..2c46920f4b 100644 --- a/internal/downloads/artifacts.go +++ b/internal/downloads/artifacts.go @@ -325,11 +325,14 @@ func (m *ArtifactManager) ensureResolved(ctx context.Context, file *models.Media OutputPath: artifactOutputPath(m.artifactDir(), file.ID, format, hash), MaxAttempts: artifactMaxAttempts, } + if playback.PreparedTracksAvailable(file) { + a.TrackRecipeVersion = playback.PreparedTracksRecipeVersion + } request := downloadprepare.NewRequest(a.ID, m.buildOpts(file, a)) if request.StereoDownmixBoostRequested() { a.AudioRecipeVersion = request.AudioRecipeVersion } - if a.ToneMapMode != "" || a.AudioRecipeVersion != "" { + if artifactUsesExecutionFingerprint(a) { a.ParamsHash = request.ExecutionFingerprint() a.OutputPath = artifactOutputPath(m.artifactDir(), file.ID, format, a.ParamsHash) } @@ -368,7 +371,7 @@ func (m *ArtifactManager) ensureResolved(ctx context.Context, file *models.Media case err != nil: return nil, err default: - row.Status = queuedArtifactStatus(row.ToneMapMode, row.AudioRecipeVersion) + row.Status = queuedArtifactStatus(row.ToneMapMode, row.AudioRecipeVersion, row.TrackRecipeVersion) } m.triggerDrain() return row, nil @@ -1089,6 +1092,14 @@ func (m *ArtifactManager) buildOpts(file *models.MediaFile, a *Artifact) playbac slog.Warn("download artifact tone-map source revision is invalid", "component", "downloads", "artifact_id", a.ID, "source_revision_length", len(a.ToneMapSourceRevision)) sourceRevision = tonemap.SourceRevision{MediaFileID: -1} } + // The multi-track layout carries each encoded track's channel count, so + // the single-track downmix recipe fields stay unset. + var preparedTracks *playback.PreparedTracks + sourceAudioChannels := preparedSourceAudioChannels(file, a.AudioTrackIndex, a.CodecAudio) + if a.TrackRecipeVersion != "" { + preparedTracks = playback.PlanPreparedTracks(file, a.CodecAudio, a.AudioTrackIndex) + sourceAudioChannels = 0 + } return playback.TranscodeOpts{ InputPath: file.FilePath, SourceVideoCodec: sourceVideoCodec, @@ -1110,8 +1121,9 @@ func (m *ArtifactManager) buildOpts(file *models.MediaFile, a *Artifact) playbac ToneMapDVBLPresent: a.ToneMapDVBLPresent, ToneMapDVRPUPresent: a.ToneMapDVRPUPresent, AudioTrackIndex: a.AudioTrackIndex, - SourceAudioChannels: preparedSourceAudioChannels(file, a.AudioTrackIndex, a.CodecAudio), + SourceAudioChannels: sourceAudioChannels, SubtitleTrackIndex: -1, + PreparedTracks: preparedTracks, FFmpegPath: cfg.Playback.FFmpegPath, HWAccel: cfg.Playback.HWAccel, HWDevice: cfg.Playback.HWDevice, diff --git a/internal/downloads/manifest.go b/internal/downloads/manifest.go index 6176d58cd9..eb654b47c7 100644 --- a/internal/downloads/manifest.go +++ b/internal/downloads/manifest.go @@ -11,6 +11,7 @@ import ( "github.com/Silo-Server/silo-server/internal/catalog" "github.com/Silo-Server/silo-server/internal/models" + "github.com/Silo-Server/silo-server/internal/playback" "github.com/Silo-Server/silo-server/internal/subtitles" ) @@ -55,6 +56,7 @@ type OfflineChapter struct { // authenticated proxy endpoint, never a presigned URL. type OfflineSubtitle struct { Language string `json:"language"` + Title string `json:"title,omitempty"` Format string `json:"format"` Forced bool `json:"forced"` HearingImpaired bool `json:"hearing_impaired"` @@ -282,19 +284,21 @@ func (b *ManifestBuilder) build(ctx context.Context, dl *Download, filter catalo // Remux/transcode entries deliver the prepared artifact, not the catalog // source: describe that file so the client picks the right decoder/tracks. + var prepared *Artifact if dl.Format != FormatOriginal && dl.ArtifactID != "" && b.artifact != nil { if a, err := b.artifact(ctx, dl.ArtifactID); err == nil && a != nil { - applyArtifactParams(m, a) + prepared = a + applyArtifactParams(m, a, file) } } - m.Subtitles = b.buildSubtitles(ctx, dl, file) + m.Subtitles = b.buildSubtitles(ctx, dl, file, prepared) return m, nil } // applyArtifactParams overwrites the source file's media parameters with the // prepared artifact's target parameters. "copy" targets keep the source value. -func applyArtifactParams(m *OfflineManifest, a *Artifact) { +func applyArtifactParams(m *OfflineManifest, a *Artifact, file *models.MediaFile) { if a.Container != "" { m.Container = a.Container } @@ -307,8 +311,12 @@ func applyArtifactParams(m *OfflineManifest, a *Artifact) { if a.Resolution != "" { m.Resolution = a.Resolution } - // A prepared file contains exactly one audio stream — the track the encode - // selected (playback.PrepareFile maps a single audio track). + if a.TrackRecipeVersion != "" { + applyPreparedAudioTracks(m, a, file) + return + } + // A legacy prepared file contains exactly one audio stream — the track the + // encode selected (playback.PrepareFile mapped a single audio track). if len(m.AudioTracks) > 0 { idx := a.AudioTrackIndex if idx < 0 || idx >= len(m.AudioTracks) { @@ -326,11 +334,47 @@ func applyArtifactParams(m *OfflineManifest, a *Artifact) { } } -// buildSubtitles enumerates external (sidecar) + downloaded (S3) subtitle assets -// for the download's media file (already loaded by build — no re-fetch). -// Embedded tracks live inside the downloaded video file and need no separate -// fetch. -func (b *ManifestBuilder) buildSubtitles(ctx context.Context, dl *Download, file *models.MediaFile) []OfflineSubtitle { +// applyPreparedAudioTracks describes a multi-track prepared file. It keeps +// every source audio track in source order, so output positions equal source +// positions and the viewer's catalog selection stays valid; encoded tracks +// report the AAC output layout. +func applyPreparedAudioTracks(m *OfflineManifest, a *Artifact, file *models.MediaFile) { + if file == nil { + return + } + plan := playback.PlanPreparedTracks(file, a.CodecAudio, a.AudioTrackIndex) + tracks := toOfflineAudioTracks(file.AudioTracks) + channels, bitrateKbps := playback.ResolveAACOutputV3(0, 0) + fileDefault := -1 + for i, track := range plan.Audio { + tracks[i].Default = track.Default + if track.Default { + fileDefault = i + } + if track.Codec == playback.PreparedAudioAAC { + tracks[i].Codec = playback.PreparedAudioAAC + tracks[i].Channels = channels + tracks[i].Layout = "stereo" + tracks[i].Bitrate = bitrateKbps + } + } + m.AudioTracks = tracks + if selected := m.SelectedAudioTrackIndex; selected != nil && *selected >= 0 && *selected < len(tracks) { + return + } + m.SelectedAudioTrackIndex = nil + if fileDefault >= 0 { + m.SelectedAudioTrackIndex = &fileDefault + } +} + +// buildSubtitles enumerates external (sidecar), embedded sidecar, and +// downloaded (S3) subtitle assets for the download's media file (already +// loaded by build — no re-fetch). Other embedded tracks live inside the +// downloaded video file and need no separate fetch. A multi-track prepared MP4 +// carries plain-text subtitles as timed text; ASS/SSA and PGS tracks are +// offered as .ass/.sup sidecars extracted from the source instead. +func (b *ManifestBuilder) buildSubtitles(ctx context.Context, dl *Download, file *models.MediaFile, prepared *Artifact) []OfflineSubtitle { out := []OfflineSubtitle{} if file != nil { @@ -341,6 +385,7 @@ func (b *ManifestBuilder) buildSubtitles(ctx context.Context, dl *Download, file } out = append(out, OfflineSubtitle{ Language: ext.Language, + Title: ext.Title, Format: ext.Format, Forced: ext.Forced, HearingImpaired: ext.HearingImpaired, @@ -349,6 +394,22 @@ func (b *ManifestBuilder) buildSubtitles(ctx context.Context, dl *Download, file FileSize: size, }) } + if prepared != nil && prepared.TrackRecipeVersion != "" { + for i, track := range file.SubtitleTracks { + format := playback.PreparedSubtitleSidecarFormat(track.Codec) + if track.External || format == "" { + continue + } + out = append(out, OfflineSubtitle{ + Language: track.Language, + Title: track.EmbeddedTitle, + Format: format, + Forced: track.Forced, + HearingImpaired: track.HearingImpaired, + FetchURL: subtitleProxyURL(dl.ID, fmt.Sprintf("%s:%d", subtitleRefEmbedded, i)), + }) + } + } } if b.subs != nil { diff --git a/internal/downloads/manifest_test.go b/internal/downloads/manifest_test.go index 06c4a3a378..a43a3f0d8b 100644 --- a/internal/downloads/manifest_test.go +++ b/internal/downloads/manifest_test.go @@ -5,10 +5,13 @@ import ( "context" "encoding/json" "errors" + "net/http" + "net/http/httptest" "testing" "github.com/Silo-Server/silo-server/internal/catalog" "github.com/Silo-Server/silo-server/internal/models" + "github.com/Silo-Server/silo-server/internal/playback" "github.com/Silo-Server/silo-server/internal/subtitles" ) @@ -160,6 +163,89 @@ func TestManifestBuilderAssembles(t *testing.T) { } } +func preparedManifestFixture(artifact *Artifact) (*ManifestBuilder, *Download) { + audio := []models.AudioTrack{ + {Codec: "truehd", Channels: 8, Language: "en", Layout: "7.1", Title: "TrueHD 7.1", Default: true}, + {Codec: "ac3", Channels: 2, Language: "ja", EmbeddedTitle: "Commentary", Title: "Commentary"}, + } + selected := 1 + detail := &catalog.ItemDetail{Type: "movie", Title: "The Movie", Versions: []catalog.FileVersion{{ + FileID: 99, Container: "mkv", CodecVideo: "hevc", CodecAudio: "truehd", Resolution: "2160p", + AudioTracks: audio, EffectiveAudioTrackIndex: &selected, + }}} + file := &models.MediaFile{ + ID: 99, CodecAudio: "truehd", AudioTracks: audio, + ExternalSubtitles: []models.ExternalSubtitle{{Path: "/media/sub.en.srt", Language: "en", Format: "srt"}}, + SubtitleTracks: []models.SubtitleTrack{ + {Codec: "subrip", Language: "en"}, + {Codec: "hdmv_pgs_subtitle", Language: "fr", Forced: true}, + {Codec: "ass", Language: "ja", EmbeddedTitle: "Signs & Songs", HearingImpaired: true}, + }, + } + lookup := func(context.Context, string) (*Artifact, error) { return artifact, nil } + b := NewManifestBuilder(fakeManifestSource{detail: detail}, nil, fakeFileResolver{file: file}, lookup) + return b, &Download{ID: "dl1", ContentID: "c1", MediaFileID: 99, Format: FormatTranscode, ArtifactID: artifact.ID} +} + +// TestManifestDescribesMultiTrackArtifact verifies a multi-track prepared file +// is described by output position, with encoded tracks reporting their AAC +// layout and PGS offered as a .sup sidecar the MP4 cannot store. +func TestManifestDescribesMultiTrackArtifact(t *testing.T) { + b, dl := preparedManifestFixture(&Artifact{ + ID: "a1", Container: "mp4", CodecVideo: "h264", CodecAudio: "aac", Resolution: "1080p", + AudioTrackIndex: -1, TrackRecipeVersion: playback.PreparedTracksRecipeVersion, + }) + m, err := b.Build(context.Background(), dl, catalog.AccessFilter{}) + if err != nil { + t.Fatalf("Build: %v", err) + } + // The viewer's catalog selection (the commentary track) survives because + // every source track is present at its source position. + if len(m.AudioTracks) != 2 || m.SelectedAudioTrackIndex == nil || *m.SelectedAudioTrackIndex != 1 { + t.Fatalf("audio tracks = %+v selected %v, want both tracks with output 1 selected", m.AudioTracks, m.SelectedAudioTrackIndex) + } + if !m.AudioTracks[0].Default { + t.Fatal("source default track lost its default flag") + } + for i, track := range m.AudioTracks { + if track.Index != i || track.Codec != "aac" || track.Channels != 2 || track.Layout != "stereo" { + t.Fatalf("audio track %d = %+v, want AAC stereo at output %d", i, track, i) + } + } + if m.AudioTracks[1].Language != "ja" || m.AudioTracks[1].Title != "Commentary" || m.AudioTracks[1].Default { + t.Fatalf("commentary track = %+v", m.AudioTracks[1]) + } + if len(m.Subtitles) != 3 { + t.Fatalf("subtitles = %+v, want external sidecar plus PGS and ASS sidecars", m.Subtitles) + } + pgs := m.Subtitles[1] + if pgs.FetchURL != "/api/v2/downloads/dl1/subtitles/embedded:1" || pgs.Format != "sup" || pgs.Language != "fr" || !pgs.Forced || pgs.External { + t.Fatalf("PGS sidecar = %+v", pgs) + } + ass := m.Subtitles[2] + if ass.FetchURL != "/api/v2/downloads/dl1/subtitles/embedded:2" || ass.Format != "ass" || ass.Language != "ja" || ass.Title != "Signs & Songs" || !ass.HearingImpaired { + t.Fatalf("ASS sidecar = %+v", ass) + } +} + +// TestManifestDescribesLegacySingleTrackArtifact keeps already-prepared files +// described as the one audio stream they contain, without PGS sidecars. +func TestManifestDescribesLegacySingleTrackArtifact(t *testing.T) { + b, dl := preparedManifestFixture(&Artifact{ + ID: "a1", Container: "mp4", CodecVideo: "h264", CodecAudio: "aac", Resolution: "1080p", AudioTrackIndex: -1, + }) + m, err := b.Build(context.Background(), dl, catalog.AccessFilter{}) + if err != nil { + t.Fatalf("Build: %v", err) + } + if len(m.AudioTracks) != 1 || m.AudioTracks[0].Codec != "aac" || *m.SelectedAudioTrackIndex != 0 { + t.Fatalf("legacy audio tracks = %+v", m.AudioTracks) + } + if len(m.Subtitles) != 1 || m.Subtitles[0].FetchURL != "/api/v2/downloads/dl1/subtitles/external:0" { + t.Fatalf("legacy subtitles = %+v, want only the external sidecar", m.Subtitles) + } +} + func TestParseSubtitleRef(t *testing.T) { cases := []struct { ref string @@ -170,6 +256,7 @@ func TestParseSubtitleRef(t *testing.T) { {"external:0", "external", 0, false}, {"external:12", "external", 12, false}, {"downloaded:7", "downloaded", 7, false}, + {"embedded:3", "embedded", 3, false}, {"bogus", "", 0, true}, {"external:x", "", 0, true}, {"weird:1", "", 0, true}, @@ -190,3 +277,29 @@ func TestParseSubtitleRef(t *testing.T) { }) } } + +func TestServeEmbeddedSubtitleOnlyServesSidecarTracks(t *testing.T) { + file := &models.MediaFile{ + ID: 99, FilePath: t.TempDir() + "/missing.mkv", + SubtitleTracks: []models.SubtitleTrack{ + {Codec: "subrip"}, + {Codec: "hdmv_pgs_subtitle"}, + }, + } + s := &Service{fileRepo: fakeFileResolver{file: file}} + dl := &Download{ID: "dl1", MediaFileID: 99} + req := httptest.NewRequest(http.MethodGet, "/", nil) + for _, ordinal := range []int{-1, 0, 2} { + if err := s.serveEmbeddedSubtitle(httptest.NewRecorder(), req, dl, ordinal); !errors.Is(err, ErrAssetNotFound) { + t.Fatalf("ordinal %d err = %v, want ErrAssetNotFound", ordinal, err) + } + } + // A failed PGS extract after the shared cache committed its 200 must abort + // the response rather than end it as a complete track. + defer func() { + if rec := recover(); rec != http.ErrAbortHandler { //nolint:errorlint // sentinel compared by identity, as net/http does + t.Fatalf("failed committed extract recovered %v, want http.ErrAbortHandler", rec) + } + }() + _ = s.serveEmbeddedSubtitle(httptest.NewRecorder(), req, dl, 1) +} diff --git a/internal/downloads/offline.go b/internal/downloads/offline.go index a9744f8b60..d63a76bf53 100644 --- a/internal/downloads/offline.go +++ b/internal/downloads/offline.go @@ -13,6 +13,8 @@ import ( "strings" "time" + "github.com/go-chi/chi/v5/middleware" + "github.com/Silo-Server/silo-server/internal/catalog" "github.com/Silo-Server/silo-server/internal/playback" "github.com/Silo-Server/silo-server/internal/subtitles" @@ -275,10 +277,15 @@ func (r *storeReader) Read(p []byte) (int, error) { return n, err } -// ServeSubtitle streams a subtitle asset (external sidecar or downloaded S3 file) -// for a managed entry, authorized on (user, profile, device) with a per-profile -// content-access re-check. ref encodes "external:{index}" or "downloaded:{id}". -func (s *Service) ServeSubtitle(ctx context.Context, w http.ResponseWriter, _ *http.Request, userID int, profileID, deviceID, downloadID, ref string, filter catalog.AccessFilter) error { +// subtitleRefEmbedded addresses an embedded ASS/SSA or PGS track by its +// subtitle ordinal (0:s:N); multi-track prepared MP4s deliver those as sidecars. +const subtitleRefEmbedded = "embedded" + +// ServeSubtitle streams a subtitle asset (external sidecar, embedded ASS or PGS +// track, or downloaded S3 file) for a managed entry, authorized on (user, profile, +// device) with a per-profile content-access re-check. ref encodes +// "external:{index}", "embedded:{ordinal}", or "downloaded:{id}". +func (s *Service) ServeSubtitle(ctx context.Context, w http.ResponseWriter, r *http.Request, userID int, profileID, deviceID, downloadID, ref string, filter catalog.AccessFilter) error { dl, err := s.authorizeManagedAsset(ctx, userID, profileID, deviceID, downloadID) if err != nil { return err @@ -309,6 +316,8 @@ func (s *Service) ServeSubtitle(ctx context.Context, w http.ResponseWriter, _ *h } writeSubtitle(w, ext.Format, data) return nil + case subtitleRefEmbedded: + return s.serveEmbeddedSubtitle(w, r.WithContext(ctx), dl, value) case "downloaded": if s.subtitleSource == nil { return ErrManifestUnavailable @@ -329,15 +338,52 @@ func (s *Service) ServeSubtitle(ctx context.Context, w http.ResponseWriter, _ *h } } -// parseSubtitleRef parses a subtitle reference of the form "external:{index}" -// or "downloaded:{id}" into its kind and integer value. +// serveEmbeddedSubtitle serves one complete embedded ASS/SSA script or PGS +// stream, the sidecar a multi-track prepared download advertises for subtitles +// its MP4 cannot carry faithfully. It shares the streaming subtitle cache, so a +// track is demuxed from the source at most once while its cache entry lives. +func (s *Service) serveEmbeddedSubtitle(w http.ResponseWriter, r *http.Request, dl *Download, ordinal int) error { + file, err := s.fileRepo.GetByID(r.Context(), dl.MediaFileID) + if err != nil { + return fmt.Errorf("loading media file: %w", err) + } + if file == nil || ordinal < 0 || ordinal >= len(file.SubtitleTracks) { + return ErrAssetNotFound + } + track := file.SubtitleTracks[ordinal] + if track.External || playback.PreparedSubtitleSidecarFormat(track.Codec) == "" { + return ErrAssetNotFound + } + ffmpegPath := "" + if s.artifacts != nil && s.artifacts.liveCfg != nil { + if cfg := s.artifacts.liveCfg(); cfg != nil { + ffmpegPath = cfg.Playback.FFmpegPath + } + } + response := middleware.NewWrapResponseWriter(w, r.ProtoMajor) + err = s.subtitleCache.ServeExtract(response, r, playback.StreamExtractOpts{ + InputPath: file.FilePath, + TrackIndex: ordinal, + SourceCodec: track.Codec, + FFmpegPath: playback.ResolveFFmpegPath(ffmpegPath), + }, playback.StreamExtractSubtitle) + if err == nil || response.Status() == 0 || r.Context().Err() != nil { + return err + } + playback.LogSubtitleStreamError(r.Context(), err, file.ID, ordinal) + // A clean EOF would let the client keep a truncated track; abort instead. + panic(http.ErrAbortHandler) +} + +// parseSubtitleRef parses a subtitle reference of the form "external:{index}", +// "embedded:{ordinal}", or "downloaded:{id}" into its kind and integer value. func parseSubtitleRef(ref string) (kind string, value int, err error) { k, v, ok := strings.Cut(ref, ":") if !ok { return "", 0, ErrInvalidSubtitleRef } switch k { - case "external", "downloaded": + case "external", subtitleRefEmbedded, "downloaded": n, perr := strconv.Atoi(v) if perr != nil { return "", 0, ErrInvalidSubtitleRef diff --git a/internal/downloads/remote_preparer.go b/internal/downloads/remote_preparer.go index 56400e3a13..d92719da99 100644 --- a/internal/downloads/remote_preparer.go +++ b/internal/downloads/remote_preparer.go @@ -6,6 +6,7 @@ import ( "fmt" "log/slog" "net/http" + "slices" "strings" "sync" "time" @@ -48,6 +49,7 @@ type NodeAwarePreparer struct { type remoteToneMapCapabilities struct { capabilities tonemap.Capabilities transformations []playback.TransformationV3 + transportFeatures []string err error expiresAt time.Time probeRequestTimeout time.Duration @@ -146,7 +148,7 @@ func (p *NodeAwarePreparer) PrepareFile(ctx context.Context, artifactID string, request := downloadprepare.NewRequest(artifactID, opts) var node *nodepool.Node var release func() - if request.ToneMapRequested() || request.StereoDownmixBoostRequested() { + if request.ToneMapRequested() || request.StereoDownmixBoostRequested() || request.PreparedTracksRequested() { selector, ok := p.planner.(eligibleTranscodeWorkPlanner) if ok { toneMapCapable := map[string]struct{}{} @@ -157,6 +159,10 @@ func (p *NodeAwarePreparer) PrepareFile(ctx context.Context, artifactID string, if request.StereoDownmixBoostRequested() { audioBoostCapable = p.audioBoostCapableNodeURLs(ctx) } + tracksCapable := map[string]struct{}{} + if request.PreparedTracksRequested() { + tracksCapable = p.preparedTracksCapableNodeURLs(ctx) + } node, release = selector.ReserveTranscodeWorkWith("download-prepare-"+artifactID, func(candidate *nodepool.Node) bool { if candidate == nil { return false @@ -172,6 +178,11 @@ func (p *NodeAwarePreparer) PrepareFile(ctx context.Context, artifactID string, return false } } + if request.PreparedTracksRequested() { + if _, supported := tracksCapable[nodeURL]; !supported { + return false + } + } return true }) } @@ -312,9 +323,26 @@ func (p *NodeAwarePreparer) capableToneMapNodeURLs(ctx context.Context, mode ton } // audioBoostCapableNodeURLs returns nodes advertising the exact audio_to_aac -// recipe version that consumes SourceAudioChannels. Capability fetches share -// the existing bounded cache and singleflight used by tone-map discovery. +// recipe version that consumes SourceAudioChannels. func (p *NodeAwarePreparer) audioBoostCapableNodeURLs(ctx context.Context) map[string]struct{} { + return p.nodeURLsSupporting(ctx, func(entry remoteToneMapCapabilities) bool { + return supportsAudioBoostTransformation(entry.transformations) + }) +} + +// preparedTracksCapableNodeURLs returns nodes that execute the multi-track +// prepared-download layout. An older node would encode the legacy layout and +// fail attestation only after spending the whole encode. +func (p *NodeAwarePreparer) preparedTracksCapableNodeURLs(ctx context.Context) map[string]struct{} { + return p.nodeURLsSupporting(ctx, func(entry remoteToneMapCapabilities) bool { + return slices.Contains(entry.transportFeatures, playback.TransportFeaturePreparedTracksV1) + }) +} + +// nodeURLsSupporting returns the normalized URLs of enabled nodes whose +// capability report satisfies supports. Capability fetches share the existing +// bounded cache and singleflight used by tone-map discovery. +func (p *NodeAwarePreparer) nodeURLsSupporting(ctx context.Context, supports func(remoteToneMapCapabilities) bool) map[string]struct{} { result := make(map[string]struct{}) enumerator, ok := p.planner.(transcodeNodeEnumerator) if !ok { @@ -327,7 +355,7 @@ func (p *NodeAwarePreparer) audioBoostCapableNodeURLs(ctx context.Context) map[s wg.Add(1) go func(i int, nodeURL string) { defer wg.Done() - supported[i], _ = p.audioBoostCapabilityForNode(ctx, nodeURL) + supported[i], _ = p.nodeCapabilitySupports(ctx, nodeURL, supports) }(i, nodeURL) } wg.Wait() @@ -339,10 +367,10 @@ func (p *NodeAwarePreparer) audioBoostCapableNodeURLs(ctx context.Context) map[s return result } -func (p *NodeAwarePreparer) audioBoostCapabilityForNode(ctx context.Context, nodeURL string) (bool, error) { +func (p *NodeAwarePreparer) nodeCapabilitySupports(ctx context.Context, nodeURL string, supports func(remoteToneMapCapabilities) bool) (bool, error) { nodeURL = nodepool.NormalizeNodeURL(nodeURL) if entry, ok := p.cachedRemoteCapabilitiesForNode(nodeURL, time.Now()); ok { - return supportsAudioBoostTransformation(entry.transformations), entry.err + return supports(entry), entry.err } if _, err := p.toneMapCapabilitiesForNode(ctx, nodeURL); err != nil { return false, err @@ -351,7 +379,7 @@ func (p *NodeAwarePreparer) audioBoostCapabilityForNode(ctx context.Context, nod if !ok { return false, errors.New("transcode node capability result was not cached") } - return supportsAudioBoostTransformation(entry.transformations), entry.err + return supports(entry), entry.err } func supportsAudioBoostTransformation(transformations []playback.TransformationV3) bool { @@ -476,6 +504,7 @@ func (p *NodeAwarePreparer) fetchToneMapCapabilitiesForNode(ctx context.Context, entry := remoteToneMapCapabilities{ capabilities: append(tonemap.Capabilities(nil), info.ToneMapCapabilities...), transformations: append([]playback.TransformationV3(nil), info.Transformations...), + transportFeatures: append([]string(nil), info.TransportFeatures...), expiresAt: time.Now().Add(remoteToneMapCapabilityTTL), probeRequestTimeout: normalizeRemoteToneMapProbeTimeout(info.ProbeRequestTimeoutMillis), } diff --git a/internal/downloads/remote_preparer_test.go b/internal/downloads/remote_preparer_test.go index 2a3d0127ae..1adac3449b 100644 --- a/internal/downloads/remote_preparer_test.go +++ b/internal/downloads/remote_preparer_test.go @@ -15,6 +15,7 @@ import ( "github.com/Silo-Server/silo-server/internal/config" "github.com/Silo-Server/silo-server/internal/downloadprepare" + "github.com/Silo-Server/silo-server/internal/models" "github.com/Silo-Server/silo-server/internal/nodepool" "github.com/Silo-Server/silo-server/internal/playback" "github.com/Silo-Server/silo-server/internal/tonemap" @@ -283,6 +284,66 @@ func TestNodeAwarePreparerRequiresAudioToAACV2ForSurroundDownmix(t *testing.T) { } } +func TestNodeAwarePreparerRequiresPreparedTracksNode(t *testing.T) { + capabilityNode := func(features ...string) *httptest.Server { + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + _ = json.NewEncoder(w).Encode(playback.HWAccelInfo{TransportFeatures: features}) + })) + } + legacy := capabilityNode() + defer legacy.Close() + current := capabilityNode(playback.TransportFeaturePreparedTracksV1) + defer current.Close() + + pool := nodepool.NewTranscodePool() + pool.SetNodes([]*nodepool.Node{ + {ID: 1, URL: legacy.URL, Enabled: true, Healthy: true}, + {ID: 2, URL: current.URL, Enabled: true, Healthy: true, ActiveJobs: 1}, + }) + local := &recordingEncodePreparer{} + remote := &recordingRemotePreparer{} + cfg := &config.Config{} + cfg.Auth.JWTSecret = "secret" + p := NewNodeAwarePreparer(local, nodepool.NewPlanner(nodepool.NewProxyPool(), pool), func() *config.Config { return cfg }) + p.remote = remote + file := &models.MediaFile{CodecAudio: "aac", AudioTracks: []models.AudioTrack{{Codec: "aac"}, {Codec: "ac3", Channels: 6}}} + opts := playback.TranscodeOpts{ + InputPath: "/media/movie.mkv", TargetCodecVideo: "h264", TargetCodecAudio: "aac", + PreparedTracks: playback.PlanPreparedTracks(file, "aac", -1), + } + prepared, err := p.PrepareFile(context.Background(), "artifact-tracks", opts, "/local/artifact.mp4") + if err != nil { + t.Fatal(err) + } + if local.calls != 0 || remote.nodeURL != current.URL || prepared.OriginNodeID != 2 { + t.Fatalf("local calls = %d, remote node = %q, want prepared-tracks node %q", local.calls, remote.nodeURL, current.URL) + } + if remote.request.TrackRecipeVersion != playback.PreparedTracksRecipeVersion || len(remote.request.PreparedTracks.Audio) != 2 { + t.Fatalf("remote request = %#v, want the multi-track layout", remote.request) + } +} + +func TestRemotePrepareResultRequiresPreparedTracksAttestation(t *testing.T) { + file := &models.MediaFile{AudioTracks: []models.AudioTrack{{Codec: "aac"}, {Codec: "aac"}}} + request := downloadprepare.NewRequest("artifact-tracks", playback.TranscodeOpts{ + TargetCodecAudio: "aac", PreparedTracks: playback.PlanPreparedTracks(file, "aac", -1), + }) + // A node without the layout ignores it and attests its legacy recipe. + legacyRequest := request + legacyRequest.TrackRecipeVersion, legacyRequest.PreparedTracks = "", nil + legacy := downloadprepare.Result{ArtifactID: request.ArtifactID, FileSize: 55, ExecutionFingerprint: legacyRequest.ExecutionFingerprint()} + if remotePrepareResultMatches(legacy, request.ArtifactID, request) { + t.Fatal("multi-track request accepted a legacy single-track receipt") + } + if remotePrepareResultMatches(downloadprepare.Result{ArtifactID: request.ArtifactID, FileSize: 55}, request.ArtifactID, request) { + t.Fatal("multi-track request accepted an unattested result") + } + attested := downloadprepare.Result{ArtifactID: request.ArtifactID, FileSize: 55, ExecutionFingerprint: request.ExecutionFingerprint()} + if !remotePrepareResultMatches(attested, request.ArtifactID, request) { + t.Fatal("multi-track request rejected its exact execution receipt") + } +} + func TestNodeAwarePreparerRejectsUnattestedOrMismatchedToneMapPrepareResult(t *testing.T) { revision := tonemap.SourceRevision{MediaFileID: 42, FileSize: 100, StreamSignature: "stream"} valid := downloadprepare.Result{ diff --git a/internal/downloads/repo.go b/internal/downloads/repo.go index 57aff4c356..414b431b83 100644 --- a/internal/downloads/repo.go +++ b/internal/downloads/repo.go @@ -787,7 +787,7 @@ func (r *Repository) ConfirmArtifactLink(ctx context.Context, d *Download) (*Dow return nil, fmt.Errorf("checking linked artifact status: %w", err) } switch artifactStatus { - case "queued", "tone_map_queued", "audio_v2_queued", "running", "tone_map_running", "audio_v2_running": + case "queued", "tone_map_queued", "audio_v2_queued", "tracks_v1_queued", "running", "tone_map_running", "audio_v2_running", "tracks_v1_running": if _, err := tx.Exec(ctx, `UPDATE downloads SET status = 'preparing', bytes_sent = 0, completed_at = NULL, error_message = '', updated_at = now() @@ -825,7 +825,7 @@ func (r *Repository) ReconcileLinkedDownloads(ctx context.Context) (ready []*Dow file_size = COALESCE((SELECT a.file_size FROM download_artifacts a WHERE a.id = downloads.artifact_id), file_size), updated_at = now() WHERE status = 'preparing' AND artifact_id IS NOT NULL - AND artifact_id IN (SELECT id FROM download_artifacts WHERE status IN ('ready', 'tone_map_ready', 'audio_v2_ready')) + AND artifact_id IN (SELECT id FROM download_artifacts WHERE status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready')) RETURNING `+downloadColumns, ) if err != nil { diff --git a/internal/downloads/service.go b/internal/downloads/service.go index de57f46aec..83eb064c7f 100644 --- a/internal/downloads/service.go +++ b/internal/downloads/service.go @@ -126,6 +126,7 @@ type Service struct { subtitleSource SubtitleSource artworkSource ManifestSource httpClient *http.Client + subtitleCache *playback.SubtitleCache // Prepare-to-file pipeline (Phase 3); nil until SetArtifactManager wires it. artifacts *ArtifactManager @@ -142,6 +143,10 @@ type Service struct { cfgLoadedAt time.Time } +// SetSubtitleCache shares the streaming subtitle cache with embedded subtitle +// sidecars. Nil disables caching. +func (s *Service) SetSubtitleCache(cache *playback.SubtitleCache) { s.subtitleCache = cache } + // SetOfflineDeps wires the offline-manifest dependencies (catalog detail for // manifest + artwork, subtitle assets, and an HTTP client for streaming // artwork bytes). When unset, the manifest/artwork/subtitle endpoints report diff --git a/internal/lang/lang.go b/internal/lang/lang.go index 06ca5c526d..193ad8b4fe 100644 --- a/internal/lang/lang.go +++ b/internal/lang/lang.go @@ -129,6 +129,25 @@ func CodeAliases(value string) []string { return aliases } +// ISO6392 returns the ISO 639-2/T code of value's primary language, the form +// container formats such as MP4 store per track. Undefined, private-use, and +// malformed values return "". +func ISO6392(value string) string { + primary := PrimaryLanguage(value) + if primary == "" { + return "" + } + tag, err := language.Parse(primary) + if err != nil { + return "" + } + base, _ := tag.Base() + if code := base.ISO3(); code != "und" { + return code + } + return "" +} + // PrimaryLanguage intentionally drops script and region for language matching. // It never infers a language from an undefined or private-use tag. func PrimaryLanguage(value string) string { diff --git a/internal/lang/lang_test.go b/internal/lang/lang_test.go index b2af0c2673..f0255b12ee 100644 --- a/internal/lang/lang_test.go +++ b/internal/lang/lang_test.go @@ -33,6 +33,14 @@ func TestPrimaryLanguage(t *testing.T) { } } +func TestISO6392(t *testing.T) { + for in, want := range map[string]string{"en": "eng", "eng": "eng", "pt-BR": "por", "zh-Hant": "zho", "fr": "fra", "fil": "fil", "": "", "und": "", "x-private": "", "unknown": ""} { + if got := ISO6392(in); got != want { + t.Errorf("ISO6392(%q) = %q, want %q", in, got, want) + } + } +} + func TestCanonical(t *testing.T) { cases := []struct { in, want string diff --git a/internal/playback/prepare_file.go b/internal/playback/prepare_file.go index ace99b37c6..8730e5902e 100644 --- a/internal/playback/prepare_file.go +++ b/internal/playback/prepare_file.go @@ -456,6 +456,9 @@ func buildPrepareFileArgs(opts TranscodeOpts, outputPath string) []string { opts = normalizeTranscodeOpts(opts) isVideoCopy := opts.TargetCodecVideo == "copy" isAudioCopy := opts.TargetCodecAudio == "copy" + if opts.PreparedTracks != nil { + isAudioCopy = opts.PreparedTracks.allAudioCopied() + } args := []string{"-nostdin", "-hide_banner", "-loglevel", "error"} @@ -469,7 +472,11 @@ func buildPrepareFileArgs(opts TranscodeOpts, outputPath string) []string { ) args = append(args, "-i", opts.InputPath) args = append(args, "-map_metadata", "-1", "-map_chapters", "-1") - args = appendStreamSelectionArgs(args, opts) + if opts.PreparedTracks != nil { + args = appendPreparedTrackArgs(args, opts) + } else { + args = appendStreamSelectionArgs(args, opts) + } if isVideoCopy { args = append(args, "-c:v", "copy") @@ -485,7 +492,9 @@ func buildPrepareFileArgs(opts TranscodeOpts, outputPath string) []string { if isVideoCopy && !isAudioCopy { args = append(args, "-threads", "1", "-filter_threads", "1", "-filter_complex_threads", "1") } - args = appendAudioArgs(args, opts) + if opts.PreparedTracks == nil { + args = appendAudioArgs(args, opts) + } if !isVideoCopy { args = appendVideoFilterArgs(args, opts) diff --git a/internal/playback/prepare_tracks.go b/internal/playback/prepare_tracks.go new file mode 100644 index 0000000000..5fa1299489 --- /dev/null +++ b/internal/playback/prepare_tracks.go @@ -0,0 +1,254 @@ +package playback + +import ( + "fmt" + "strconv" + "strings" + + "github.com/Silo-Server/silo-server/internal/lang" + "github.com/Silo-Server/silo-server/internal/models" +) + +// PreparedTracksRecipeVersion identifies the prepared-file stream layout that +// carries every source audio track and every text subtitle track. Artifacts +// without it are legacy single-audio, subtitle-free files. +const PreparedTracksRecipeVersion = "1" + +// Prepared audio track codecs. +const ( + PreparedAudioCopy = codecCopyV3 + PreparedAudioAAC = audioCodecAACV3 +) + +const ( + audioCodecMP3 = "mp3" + subtitleCodecTextV3 = "text" +) + +// PreparedAudioTrack is one audio stream of a prepared download, in output +// order. SourceIndex is the ffmpeg audio ordinal (0:a:N), which is also the +// position in models.MediaFile.AudioTracks. +type PreparedAudioTrack struct { + SourceIndex int `json:"source_index"` + // Codec is PreparedAudioCopy or PreparedAudioAAC. + Codec string `json:"codec"` + // SourceChannels selects the stereo downmix policy for an AAC encode. + SourceChannels int `json:"source_channels,omitempty"` + Language string `json:"language,omitempty"` + Title string `json:"title,omitempty"` + Default bool `json:"default,omitempty"` +} + +// PreparedSubtitleTrack is one embedded text subtitle converted to MP4 timed +// text (mov_text). SourceIndex is the ffmpeg subtitle ordinal (0:s:N), which is +// also the position in models.MediaFile.SubtitleTracks. +type PreparedSubtitleTrack struct { + SourceIndex int `json:"source_index"` + Language string `json:"language,omitempty"` + Title string `json:"title,omitempty"` + Default bool `json:"default,omitempty"` + Forced bool `json:"forced,omitempty"` + HearingImpaired bool `json:"hearing_impaired,omitempty"` +} + +// PreparedTracks is the frozen stream layout of a prepared download. It is part +// of the transported recipe, so every field affects the artifact's execution +// fingerprint. +type PreparedTracks struct { + Audio []PreparedAudioTrack `json:"audio"` + Subtitles []PreparedSubtitleTrack `json:"subtitles,omitempty"` +} + +// PreparedTracksAvailable reports whether file's probed audio inventory can +// drive the multi-track layout. Without one (a failed or legacy probe) the +// plan would map no audio, so such files keep the legacy layout, whose +// optional 0:a:0? mapping needs no inventory. +func PreparedTracksAvailable(file *models.MediaFile) bool { + return file != nil && len(file.AudioTracks) > 0 +} + +// PlanPreparedTracks derives the stream layout of a prepared download from the +// probed source. Every audio track is kept in source order. With a "copy" +// audio target a track is copied when it shares the primary track's codec +// (which client negotiation verified) or is universally decodable, otherwise +// it is encoded to AAC; an "aac" target encodes every track. audioTrackIndex +// marks the default track, falling back to the source default and then the +// first track. Embedded plain-text subtitles become mov_text. ASS/SSA keep +// their styling, typesetting, and overlapping events only as sidecars, and +// bitmap subtitles cannot be stored in MP4, so both are left to sidecar +// delivery (PreparedSubtitleSidecarFormat). +func PlanPreparedTracks(file *models.MediaFile, targetCodecAudio string, audioTrackIndex int) *PreparedTracks { + plan := &PreparedTracks{Audio: []PreparedAudioTrack{}} + if file == nil { + return plan + } + copyAudio := strings.EqualFold(strings.TrimSpace(targetCodecAudio), PreparedAudioCopy) + primaryCodec := normalizeCodecV3(file.CodecAudio) + if primaryCodec == "" && len(file.AudioTracks) > 0 { + primaryCodec = normalizeCodecV3(file.AudioTracks[0].Codec) + } + defaultIndex := preparedDefaultAudioIndex(file.AudioTracks, audioTrackIndex) + for i, track := range file.AudioTracks { + codec := PreparedAudioAAC + if copyAudio && preparedAudioCopyable(track.Codec, primaryCodec) { + codec = PreparedAudioCopy + } + sourceChannels := 0 + if codec == PreparedAudioAAC { + sourceChannels = track.Channels + } + plan.Audio = append(plan.Audio, PreparedAudioTrack{ + SourceIndex: i, + Codec: codec, + SourceChannels: sourceChannels, + Language: track.Language, + Title: track.EmbeddedTitle, + Default: i == defaultIndex, + }) + } + for i, track := range file.SubtitleTracks { + if track.External || !PreparedSubtitleEmbeddable(track.Codec) { + continue + } + plan.Subtitles = append(plan.Subtitles, PreparedSubtitleTrack{ + SourceIndex: i, + Language: track.Language, + Title: track.EmbeddedTitle, + Default: track.Default, + Forced: track.Forced, + HearingImpaired: track.HearingImpaired, + }) + } + return plan +} + +func (p *PreparedTracks) allAudioCopied() bool { + for _, track := range p.Audio { + if track.Codec != PreparedAudioCopy { + return false + } + } + return true +} + +func preparedDefaultAudioIndex(tracks []models.AudioTrack, requested int) int { + if requested >= 0 && requested < len(tracks) { + return requested + } + for i, track := range tracks { + if track.Default { + return i + } + } + return 0 +} + +// preparedAudioCopyable reports whether a non-primary audio track may be +// stream-copied into the prepared MP4 alongside a copied primary track. +func preparedAudioCopyable(codec, primaryCodec string) bool { + normalized := normalizeCodecV3(codec) + if normalized == "" { + return false + } + return normalized == primaryCodec || normalized == audioCodecAACV3 || normalized == audioCodecMP3 +} + +// preparedTextSubtitleCodecs lists embedded plain-text subtitle codecs that +// convert to MP4 timed text without losing content. +var preparedTextSubtitleCodecs = map[string]bool{ + subtitleCodecSubRip: true, + subtitleFormatSRT: true, + subtitleMuxerWebVTT: true, + subtitleCodecMovText: true, + subtitleCodecTextV3: true, +} + +// PreparedSubtitleEmbeddable reports whether an embedded subtitle codec is +// carried inside a prepared MP4 download as timed text. +func PreparedSubtitleEmbeddable(codec string) bool { + return preparedTextSubtitleCodecs[normalizeCodecV3(codec)] +} + +// PreparedSubtitleSidecarFormat returns the sidecar file format ("ass" or +// "sup") StreamExtractSubtitle produces for an embedded subtitle a prepared +// MP4 cannot carry faithfully, or "" when the track is embedded in the MP4 or +// not delivered. ASS/SSA would lose styling, drawing commands, and overlapping +// events as MP4 timed text; PGS has no MP4 representation. +func PreparedSubtitleSidecarFormat(codec string) string { + if _, format := streamExtractOutput(codec); format == subtitleFormatASS || format == subtitleFormatSUP { + return format + } + return "" +} + +// appendPreparedTrackArgs maps and encodes every stream of a prepared-file +// plan. The prepared output clears source metadata, so each stream's language, +// title (MP4 handler name), and disposition are written explicitly. +func appendPreparedTrackArgs(args []string, opts TranscodeOpts) []string { + plan := opts.PreparedTracks + args = append(args, "-map", "0:v:0") + for _, track := range plan.Audio { + args = append(args, "-map", fmt.Sprintf("0:a:%d", track.SourceIndex)) + } + for _, track := range plan.Subtitles { + args = append(args, "-map", fmt.Sprintf("0:s:%d", track.SourceIndex)) + } + args = append(args, "-dn") + + for i, track := range plan.Audio { + stream := strconv.Itoa(i) + if track.Codec == PreparedAudioCopy { + args = append(args, "-c:a:"+stream, "copy") + } else { + channels, bitrateKbps := ResolveAACOutputV3(0, 0) + filter := aacTimestampNormalizeFilterV3 + if IsAudioToAACStereoDownmixV3(track.SourceChannels, PreparedAudioAAC, 0) { + filter = stereoDownmixBoostFilterV3 + } + args = append(args, + "-c:a:"+stream, audioCodecAACV3, + "-b:a:"+stream, strconv.Itoa(bitrateKbps)+"k", + "-ac:a:"+stream, strconv.Itoa(channels), + "-filter:a:"+stream, filter, + ) + } + args = appendPreparedStreamMetadata(args, "a:"+stream, track.Language, track.Title) + args = append(args, "-disposition:a:"+stream, preparedDisposition(track.Default, false, false)) + } + for i, track := range plan.Subtitles { + stream := strconv.Itoa(i) + args = append(args, "-c:s:"+stream, subtitleCodecMovText) + args = appendPreparedStreamMetadata(args, "s:"+stream, track.Language, track.Title) + args = append(args, "-disposition:s:"+stream, preparedDisposition(track.Default, track.Forced, track.HearingImpaired)) + } + return args +} + +// appendPreparedStreamMetadata writes an MP4-compatible language (ISO 639-2) +// and title. The MP4 muxer stores a track title only as its handler name. +func appendPreparedStreamMetadata(args []string, stream, language, title string) []string { + if code := lang.ISO6392(language); code != "" { + args = append(args, "-metadata:s:"+stream, "language="+code) + } + if title = strings.TrimSpace(title); title != "" { + args = append(args, "-metadata:s:"+stream, "handler_name="+title) + } + return args +} + +func preparedDisposition(isDefault, forced, hearingImpaired bool) string { + var flags []string + if isDefault { + flags = append(flags, "default") + } + if forced { + flags = append(flags, "forced") + } + if hearingImpaired { + flags = append(flags, "hearing_impaired") + } + if len(flags) == 0 { + return "0" + } + return strings.Join(flags, "+") +} diff --git a/internal/playback/prepare_tracks_test.go b/internal/playback/prepare_tracks_test.go new file mode 100644 index 0000000000..c00ae7875e --- /dev/null +++ b/internal/playback/prepare_tracks_test.go @@ -0,0 +1,262 @@ +package playback + +import ( + "context" + "encoding/json" + "os" + "os/exec" + "path/filepath" + "reflect" + "strings" + "testing" + "time" + + "github.com/Silo-Server/silo-server/internal/models" +) + +func preparedTracksTestFile() *models.MediaFile { + return &models.MediaFile{ + CodecAudio: "eac3", + AudioTracks: []models.AudioTrack{ + {Codec: "eac3", Channels: 6, Language: "en", EmbeddedTitle: "Surround", Title: "Surround"}, + {Codec: "truehd", Channels: 8, Language: "en", Title: "TRUEHD", Default: true}, + {Codec: "aac", Channels: 2, Language: "ja", EmbeddedTitle: "Commentary", Title: "Commentary"}, + }, + SubtitleTracks: []models.SubtitleTrack{ + {Codec: "subrip", Language: "fr", EmbeddedTitle: "Forced", Forced: true}, + {Codec: "hdmv_pgs_subtitle", Language: "en"}, + {Codec: "ass", Language: "en", HearingImpaired: true}, + {Codec: "dvd_subtitle", Language: "de"}, + {Codec: "webvtt", Language: "es", HearingImpaired: true}, + }, + } +} + +func TestPlanPreparedTracksKeepsEveryAudioTrackAndTextSubtitle(t *testing.T) { + file := preparedTracksTestFile() + + remux := PlanPreparedTracks(file, "copy", -1) + wantRemuxAudio := []PreparedAudioTrack{ + {SourceIndex: 0, Codec: "copy", Language: "en", Title: "Surround"}, + {SourceIndex: 1, Codec: "aac", SourceChannels: 8, Language: "en", Default: true}, + {SourceIndex: 2, Codec: "copy", Language: "ja", Title: "Commentary"}, + } + if !reflect.DeepEqual(remux.Audio, wantRemuxAudio) { + t.Fatalf("remux audio plan = %+v, want %+v", remux.Audio, wantRemuxAudio) + } + wantSubtitles := []PreparedSubtitleTrack{ + {SourceIndex: 0, Language: "fr", Title: "Forced", Forced: true}, + {SourceIndex: 4, Language: "es", HearingImpaired: true}, + } + if !reflect.DeepEqual(remux.Subtitles, wantSubtitles) { + t.Fatalf("subtitle plan = %+v, want %+v", remux.Subtitles, wantSubtitles) + } + // wantRemuxAudio marks the source default track (1) as the default output. + + transcode := PlanPreparedTracks(file, "aac", 2) + for i, track := range transcode.Audio { + if track.Codec != "aac" { + t.Fatalf("transcode track %d codec = %q, want aac", i, track.Codec) + } + if track.SourceChannels != file.AudioTracks[i].Channels { + t.Fatalf("transcode track %d source channels = %d, want %d", i, track.SourceChannels, file.AudioTracks[i].Channels) + } + } + if !transcode.Audio[2].Default || transcode.Audio[1].Default { + t.Fatalf("explicit audio track index must be the only default output: %+v", transcode.Audio) + } +} + +func TestPreparedSubtitleSidecarFormat(t *testing.T) { + for codec, want := range map[string]string{ + "ass": "ass", "ssa": "ass", "hdmv_pgs_subtitle": "sup", "pgssub": "sup", + "subrip": "", "webvtt": "", "dvd_subtitle": "", "dvb_subtitle": "", + } { + if got := PreparedSubtitleSidecarFormat(codec); got != want { + t.Errorf("PreparedSubtitleSidecarFormat(%q) = %q, want %q", codec, got, want) + } + } +} + +func TestPreparedTracksAvailableRequiresAudioInventory(t *testing.T) { + for name, tc := range map[string]struct { + file *models.MediaFile + want bool + }{ + "probed audio": {file: &models.MediaFile{CodecAudio: "aac", AudioTracks: []models.AudioTrack{{Codec: "aac"}}}, want: true}, + "no probed audio": {file: &models.MediaFile{}, want: false}, + "codec without inventory": {file: &models.MediaFile{CodecAudio: "eac3"}, want: false}, + "missing file": {file: nil, want: false}, + } { + if got := PreparedTracksAvailable(tc.file); got != tc.want { + t.Errorf("%s: PreparedTracksAvailable = %v, want %v", name, got, tc.want) + } + } +} + +func TestPlanPreparedTracksWithoutAudio(t *testing.T) { + plan := PlanPreparedTracks(&models.MediaFile{}, "aac", -1) + if plan.Audio == nil || len(plan.Audio) != 0 { + t.Fatalf("empty source plan = %+v", plan) + } +} + +func TestBuildPrepareFileArgsMapsPreparedTracks(t *testing.T) { + plan := PlanPreparedTracks(preparedTracksTestFile(), "copy", -1) + args := buildPrepareFileArgs(TranscodeOpts{ + InputPath: "/media/in.mkv", TargetCodecVideo: "copy", TargetCodecAudio: "copy", + HWAccel: "none", AudioTrackIndex: -1, SubtitleTrackIndex: -1, PreparedTracks: plan, + }, "/artifacts/out.mp4") + joined := strings.Join(args, " ") + + for _, want := range []string{ + "-map 0:v:0 -map 0:a:0 -map 0:a:1 -map 0:a:2 -map 0:s:0 -map 0:s:4 -dn", + "-c:a:0 copy -metadata:s:a:0 language=eng -metadata:s:a:0 handler_name=Surround -disposition:a:0 0", + "-c:a:1 aac -b:a:1 192k -ac:a:1 2 -filter:a:1 " + stereoDownmixBoostFilterV3 + " -metadata:s:a:1 language=eng -disposition:a:1 default", + "-c:a:2 copy -metadata:s:a:2 language=jpn -metadata:s:a:2 handler_name=Commentary -disposition:a:2 0", + "-c:s:0 mov_text -metadata:s:s:0 language=fra -metadata:s:s:0 handler_name=Forced -disposition:s:0 forced", + "-c:s:1 mov_text -metadata:s:s:1 language=spa -disposition:s:1 hearing_impaired", + "-movflags +faststart -f mp4", + } { + if !strings.Contains(joined, want) { + t.Fatalf("prepared track args missing %q:\n%s", want, joined) + } + } + for _, forbidden := range []string{"-sn", "-c:a copy", "-af ", "0:s:1", "0:s:2", "0:s:3"} { + if strings.Contains(joined, forbidden) { + t.Fatalf("prepared track args contain %q:\n%s", forbidden, joined) + } + } + // A mixed copy/AAC plan still encodes audio, so remux keeps its single-thread cap. + if !strings.Contains(joined, "-threads 1") { + t.Fatalf("remux with an encoded track must cap threads:\n%s", joined) + } +} + +func TestBuildPrepareFileArgsWithoutPreparedTracksKeepsLegacyLayout(t *testing.T) { + joined := strings.Join(buildPrepareFileArgs(TranscodeOpts{ + InputPath: "/media/in.mkv", TargetCodecVideo: "copy", TargetCodecAudio: "copy", + HWAccel: "none", AudioTrackIndex: -1, SubtitleTrackIndex: -1, + }, "/artifacts/out.mp4"), " ") + if !strings.Contains(joined, "-map 0:v:0 -map 0:a:0? -sn -dn") || !strings.Contains(joined, "-c:a copy") { + t.Fatalf("legacy prepared layout changed:\n%s", joined) + } +} + +type preparedProbeStream struct { + CodecType string `json:"codec_type"` + CodecName string `json:"codec_name"` + Channels int `json:"channels"` + Tags map[string]string `json:"tags"` + Disposition map[string]int `json:"disposition"` +} + +func TestPrepareFileKeepsEveryAudioAndTextSubtitleTrack(t *testing.T) { + if testing.Short() { + t.Skip("real FFmpeg integration test") + } + ffmpeg, err := exec.LookPath("ffmpeg") + if err != nil { + t.Skip("ffmpeg is not installed") + } + ffprobe := ffprobePathFromFFmpeg(ffmpeg) + if _, err := exec.LookPath(ffprobe); err != nil { + t.Skip("ffprobe is not installed") + } + ctx, cancel := context.WithTimeout(t.Context(), 2*time.Minute) + defer cancel() + dir := t.TempDir() + srt := filepath.Join(dir, "cues.srt") + if err := os.WriteFile(srt, []byte("1\n00:00:00,500 --> 00:00:01,500\nHello\n"), 0o600); err != nil { + t.Fatal(err) + } + source := filepath.Join(dir, "source.mkv") + if output, err := exec.CommandContext(ctx, ffmpeg, "-v", "error", + "-f", "lavfi", "-i", "testsrc=d=2:s=320x240:r=24", + "-f", "lavfi", "-i", "sine=f=440:d=2", + "-f", "lavfi", "-i", "sine=f=880:d=2", + "-i", srt, "-i", srt, + "-map", "0", "-map", "1", "-map", "2", "-map", "3", "-map", "4", + "-c:v", "libx264", "-preset", "ultrafast", + "-c:a:0", "ac3", "-ac:a:0", "6", "-c:a:1", "aac", + "-c:s:0", "srt", "-c:s:1", "webvtt", + "-y", source, + ).CombinedOutput(); err != nil { + t.Fatalf("create source: %v\n%s", err, output) + } + file := &models.MediaFile{ + CodecAudio: "ac3", + AudioTracks: []models.AudioTrack{ + {Codec: "ac3", Channels: 6, Language: "en", EmbeddedTitle: "Main"}, + {Codec: "aac", Channels: 1, Language: "ja", EmbeddedTitle: "Commentary"}, + }, + SubtitleTracks: []models.SubtitleTrack{ + {Codec: "subrip", Language: "fr", Forced: true}, + {Codec: "webvtt", Language: "en"}, + }, + } + + for _, tc := range []struct { + name, video, audio string + wantAudioCodecs []string + }{ + {name: "transcode", video: "h264", audio: "aac", wantAudioCodecs: []string{"aac", "aac"}}, + {name: "remux", video: "copy", audio: "copy", wantAudioCodecs: []string{"ac3", "aac"}}, + } { + t.Run(tc.name, func(t *testing.T) { + output := filepath.Join(dir, tc.name+".mp4") + err := PrepareFile(ctx, TranscodeOpts{ + InputPath: source, SourceVideoCodec: "h264", TargetCodecVideo: tc.video, TargetCodecAudio: tc.audio, + HWAccel: "none", FFmpegPath: ffmpeg, AudioTrackIndex: -1, SubtitleTrackIndex: -1, + PreparedTracks: PlanPreparedTracks(file, tc.audio, -1), + }, output) + if err != nil { + t.Fatalf("PrepareFile: %v", err) + } + probe, err := exec.CommandContext(ctx, ffprobe, "-v", "error", "-show_streams", "-of", "json", output).Output() + if err != nil { + t.Fatalf("ffprobe: %v", err) + } + var parsed struct { + Streams []preparedProbeStream `json:"streams"` + } + if err := json.Unmarshal(probe, &parsed); err != nil { + t.Fatal(err) + } + var audio, subtitles []preparedProbeStream + for _, stream := range parsed.Streams { + switch stream.CodecType { + case "audio": + audio = append(audio, stream) + case "subtitle": + subtitles = append(subtitles, stream) + } + } + if len(audio) != 2 || len(subtitles) != 2 { + t.Fatalf("prepared file has %d audio and %d subtitle streams, want 2 and 2", len(audio), len(subtitles)) + } + for i, want := range tc.wantAudioCodecs { + if audio[i].CodecName != want { + t.Fatalf("audio %d codec = %q, want %q", i, audio[i].CodecName, want) + } + } + if tc.audio == "aac" && audio[0].Channels != 2 { + t.Fatalf("surround track was not downmixed to stereo: %d channels", audio[0].Channels) + } + if audio[0].Tags["language"] != "eng" || audio[1].Tags["language"] != "jpn" || audio[1].Tags["handler_name"] != "Commentary" { + t.Fatalf("audio metadata lost: %+v / %+v", audio[0].Tags, audio[1].Tags) + } + if audio[0].Disposition["default"] != 1 || audio[1].Disposition["default"] != 0 { + t.Fatalf("default audio disposition = %d/%d, want 1/0", audio[0].Disposition["default"], audio[1].Disposition["default"]) + } + for i, want := range []string{"fra", "eng"} { + if subtitles[i].CodecName != "mov_text" || subtitles[i].Tags["language"] != want { + t.Fatalf("subtitle %d = %s/%q, want mov_text/%q", i, subtitles[i].CodecName, subtitles[i].Tags["language"], want) + } + } + if subtitles[0].Disposition["forced"] != 1 { + t.Fatal("forced subtitle disposition was lost") + } + }) + } +} diff --git a/internal/playback/protocol_v3.go b/internal/playback/protocol_v3.go index 024ca6ef9d..7a1f6821ce 100644 --- a/internal/playback/protocol_v3.go +++ b/internal/playback/protocol_v3.go @@ -113,6 +113,9 @@ const ( // node that approves theme files as progressive AAC inputs. TransportFeatureThemeAudioEgressV1 = "theme_audio_egress_v1" TransportFeatureThemeAudioExecutionV1 = "theme_audio_execution_v1" + // A transcode node that executes the multi-track prepared-download layout + // (PreparedTracksRecipeVersion). + TransportFeaturePreparedTracksV1 = "prepared_tracks_v1" ) // Degradation warning codes reported by playback plans. diff --git a/internal/playback/transcode.go b/internal/playback/transcode.go index fa4117bf6c..fc39a022f0 100644 --- a/internal/playback/transcode.go +++ b/internal/playback/transcode.go @@ -121,7 +121,11 @@ type TranscodeOpts struct { // "hdmv_pgs_subtitle"). Bitmap codecs (PGS/DVD/DVB) select the overlay // filter_complex pipeline; text codecs use the libass subtitles filter. // Empty preserves the legacy text path for callers minted before the field. - SubtitleCodec string + SubtitleCodec string + // PreparedTracks selects the multi-track stream layout of a prepared + // download (PrepareFile only). Nil keeps the legacy single-audio, + // subtitle-free layout that older artifacts were encoded with. + PreparedTracks *PreparedTracks AudioTrackIndex int // -1 = default (first track), >= 0 = specific track // SourceAudioChannels is the selected source stream's channel count. Zero // means unknown and deliberately disables stereo downmix gain: boosting an diff --git a/internal/transcodenode/server.go b/internal/transcodenode/server.go index 6e10b48ec6..405ee38020 100644 --- a/internal/transcodenode/server.go +++ b/internal/transcodenode/server.go @@ -923,6 +923,10 @@ func (s *Server) handleDownloadPrepare(w http.ResponseWriter, r *http.Request) { http.Error(w, "invalid audio recipe", http.StatusBadRequest) return } + if req.PreparedTracksRequested() && !req.ValidPreparedTracks() { + http.Error(w, "invalid track recipe", http.StatusBadRequest) + return + } cfg := s.watcher.Config() if cfg == nil { @@ -1043,6 +1047,9 @@ func expectedDownloadPrepareResult(req downloadprepare.Request, fileSize int64) if req.AudioRecipeRequested() && !req.StereoDownmixBoostRequested() { return downloadprepare.Result{}, false } + if req.PreparedTracksRequested() && !req.ValidPreparedTracks() { + return downloadprepare.Result{}, false + } result.ExecutionFingerprint = req.ExecutionFingerprint() return result, result.ExecutionFingerprint != "" } @@ -1354,6 +1361,7 @@ func (s *Server) buildCapabilitySnapshotLocked(ctx context.Context) (playback.HW } } } + info.TransportFeatures = append(info.TransportFeatures, playback.TransportFeaturePreparedTracksV1) info.CapabilityHash = playback.ComputeCapabilityHash(info) return info, nil } diff --git a/internal/workmetrics/queues.go b/internal/workmetrics/queues.go index 9eb0544dfb..1840751333 100644 --- a/internal/workmetrics/queues.go +++ b/internal/workmetrics/queues.go @@ -16,7 +16,7 @@ var queueQueries = []struct{ name, sql string }{ {workloadScan, `SELECT status, count(*), min(requested_at) FROM scan_runs WHERE status IN ('accepted','running') GROUP BY status`}, {"admin", `SELECT status, count(*), min(requested_at) FROM admin_jobs WHERE status IN ('queued','running') GROUP BY status`}, {"history_import", `SELECT status, count(*), min(created_at) FROM history_import_runs WHERE status IN ('queued','running') GROUP BY status`}, - {workloadDownloads, `SELECT CASE WHEN status IN ('queued','tone_map_queued','audio_v2_queued') THEN 'queued' ELSE 'running' END, count(*), min(created_at) FROM download_artifacts WHERE status IN ('queued','running','tone_map_queued','tone_map_running','audio_v2_queued','audio_v2_running') GROUP BY 1`}, + {workloadDownloads, `SELECT CASE WHEN status IN ('queued','tone_map_queued','audio_v2_queued','tracks_v1_queued') THEN 'queued' ELSE 'running' END, count(*), min(created_at) FROM download_artifacts WHERE status IN ('queued','running','tone_map_queued','tone_map_running','audio_v2_queued','audio_v2_running','tracks_v1_queued','tracks_v1_running') GROUP BY 1`}, {"subtitles", `SELECT status, count(*), min(created_at) FROM subtitle_ai_jobs WHERE status IN ('pending','running') GROUP BY status`}, {workloadSearch, `SELECT 'queued', count(*), min(created_at) FROM catalog_search_index_events WHERE processed_at IS NULL`}, {workloadNotifications, `SELECT 'queued', count(*), min(created_at) FROM release_events WHERE processed_at IS NULL`}, diff --git a/migrations/sql/20260928222330_fence_track_recipe_artifact_workers.sql b/migrations/sql/20260928222330_fence_track_recipe_artifact_workers.sql new file mode 100644 index 0000000000..a7ce51acfd --- /dev/null +++ b/migrations/sql/20260928222330_fence_track_recipe_artifact_workers.sql @@ -0,0 +1,163 @@ +-- +goose NO TRANSACTION + +-- +goose Up +-- A prepared download that carries every audio track and text subtitle has a +-- stream layout an older API worker cannot reproduce. Give it a durable recipe +-- discriminator and a status family that the merge-base ClaimNext predicate +-- does not recognize. The track recipe outranks the audio and tone-map +-- families: its execution fingerprint already covers both. +DROP INDEX CONCURRENTLY IF EXISTS public.download_artifacts_lease_idx; +DROP INDEX CONCURRENTLY IF EXISTS public.download_artifacts_lru_idx; + +ALTER TABLE public.download_artifacts + ADD COLUMN track_recipe_version text NOT NULL DEFAULT '', + DROP CONSTRAINT download_artifacts_status_check, + ADD CONSTRAINT download_artifacts_status_check + CHECK (status IN ( + 'queued', 'running', 'ready', + 'tone_map_queued', 'tone_map_running', 'tone_map_ready', + 'audio_v2_queued', 'audio_v2_running', 'audio_v2_ready', + 'tracks_v1_queued', 'tracks_v1_running', 'tracks_v1_ready', + 'failed' + )) NOT VALID; + +-- +goose StatementBegin +CREATE OR REPLACE FUNCTION public.fence_tone_map_artifact_worker_status() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + IF NEW.track_recipe_version <> '' THEN + IF NEW.status IN ('queued', 'tone_map_queued', 'audio_v2_queued') THEN + NEW.status := 'tracks_v1_queued'; + ELSIF NEW.status IN ('running', 'tone_map_running', 'audio_v2_running') THEN + NEW.status := 'tracks_v1_running'; + ELSIF NEW.status IN ('ready', 'tone_map_ready', 'audio_v2_ready') THEN + NEW.status := 'tracks_v1_ready'; + END IF; + ELSIF NEW.audio_recipe_version <> '' THEN + IF NEW.status IN ('queued', 'tone_map_queued', 'tracks_v1_queued') THEN + NEW.status := 'audio_v2_queued'; + ELSIF NEW.status IN ('running', 'tone_map_running', 'tracks_v1_running') THEN + NEW.status := 'audio_v2_running'; + ELSIF NEW.status IN ('ready', 'tone_map_ready', 'tracks_v1_ready') THEN + NEW.status := 'audio_v2_ready'; + END IF; + ELSIF NEW.tone_map_mode <> '' THEN + IF NEW.status IN ('queued', 'audio_v2_queued', 'tracks_v1_queued') THEN + NEW.status := 'tone_map_queued'; + ELSIF NEW.status IN ('running', 'audio_v2_running', 'tracks_v1_running') THEN + NEW.status := 'tone_map_running'; + ELSIF NEW.status IN ('ready', 'audio_v2_ready', 'tracks_v1_ready') THEN + NEW.status := 'tone_map_ready'; + END IF; + ELSE + IF NEW.status IN ('tone_map_queued', 'audio_v2_queued', 'tracks_v1_queued') THEN + NEW.status := 'queued'; + ELSIF NEW.status IN ('tone_map_running', 'audio_v2_running', 'tracks_v1_running') THEN + NEW.status := 'running'; + ELSIF NEW.status IN ('tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready') THEN + NEW.status := 'ready'; + END IF; + END IF; + RETURN NEW; +END; +$$; +-- +goose StatementEnd + +DROP TRIGGER IF EXISTS download_artifacts_tone_map_worker_status ON public.download_artifacts; +CREATE TRIGGER download_artifacts_tone_map_worker_status +BEFORE INSERT OR UPDATE OF status, tone_map_mode, audio_recipe_version, track_recipe_version ON public.download_artifacts +FOR EACH ROW +EXECUTE FUNCTION public.fence_tone_map_artifact_worker_status(); + +ALTER TABLE public.download_artifacts + VALIDATE CONSTRAINT download_artifacts_status_check; + +CREATE INDEX CONCURRENTLY download_artifacts_lease_idx ON public.download_artifacts (lease_expires_at) + WHERE status IN ('running', 'tone_map_running', 'audio_v2_running', 'tracks_v1_running'); +CREATE INDEX CONCURRENTLY download_artifacts_lru_idx ON public.download_artifacts (last_used_at) + WHERE status IN ('ready', 'tone_map_ready', 'audio_v2_ready', 'tracks_v1_ready'); + +-- +goose Down +DROP INDEX CONCURRENTLY IF EXISTS public.download_artifacts_lease_idx; +DROP INDEX CONCURRENTLY IF EXISTS public.download_artifacts_lru_idx; + +DROP TRIGGER IF EXISTS download_artifacts_tone_map_worker_status ON public.download_artifacts; + +-- Older code serves multi-track files as ordinary artifacts; their first audio +-- stream stays valid for the single-track manifest it describes. +UPDATE public.download_artifacts +SET status = CASE + WHEN status = 'tracks_v1_queued' AND audio_recipe_version <> '' THEN 'audio_v2_queued' + WHEN status = 'tracks_v1_running' AND audio_recipe_version <> '' THEN 'audio_v2_running' + WHEN status = 'tracks_v1_ready' AND audio_recipe_version <> '' THEN 'audio_v2_ready' + WHEN status = 'tracks_v1_queued' AND tone_map_mode <> '' THEN 'tone_map_queued' + WHEN status = 'tracks_v1_running' AND tone_map_mode <> '' THEN 'tone_map_running' + WHEN status = 'tracks_v1_ready' AND tone_map_mode <> '' THEN 'tone_map_ready' + WHEN status = 'tracks_v1_queued' THEN 'queued' + WHEN status = 'tracks_v1_running' THEN 'running' + WHEN status = 'tracks_v1_ready' THEN 'ready' + ELSE status +END +WHERE status IN ('tracks_v1_queued', 'tracks_v1_running', 'tracks_v1_ready'); + +-- +goose StatementBegin +CREATE OR REPLACE FUNCTION public.fence_tone_map_artifact_worker_status() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + IF NEW.audio_recipe_version <> '' THEN + IF NEW.status IN ('queued', 'tone_map_queued') THEN + NEW.status := 'audio_v2_queued'; + ELSIF NEW.status IN ('running', 'tone_map_running') THEN + NEW.status := 'audio_v2_running'; + ELSIF NEW.status IN ('ready', 'tone_map_ready') THEN + NEW.status := 'audio_v2_ready'; + END IF; + ELSIF NEW.tone_map_mode <> '' THEN + IF NEW.status IN ('queued', 'audio_v2_queued') THEN + NEW.status := 'tone_map_queued'; + ELSIF NEW.status IN ('running', 'audio_v2_running') THEN + NEW.status := 'tone_map_running'; + ELSIF NEW.status IN ('ready', 'audio_v2_ready') THEN + NEW.status := 'tone_map_ready'; + END IF; + ELSE + IF NEW.status IN ('tone_map_queued', 'audio_v2_queued') THEN + NEW.status := 'queued'; + ELSIF NEW.status IN ('tone_map_running', 'audio_v2_running') THEN + NEW.status := 'running'; + ELSIF NEW.status IN ('tone_map_ready', 'audio_v2_ready') THEN + NEW.status := 'ready'; + END IF; + END IF; + RETURN NEW; +END; +$$; +-- +goose StatementEnd + +ALTER TABLE public.download_artifacts + DROP COLUMN track_recipe_version, + DROP CONSTRAINT download_artifacts_status_check, + ADD CONSTRAINT download_artifacts_status_check + CHECK (status IN ( + 'queued', 'running', 'ready', + 'tone_map_queued', 'tone_map_running', 'tone_map_ready', + 'audio_v2_queued', 'audio_v2_running', 'audio_v2_ready', + 'failed' + )) NOT VALID; + +CREATE TRIGGER download_artifacts_tone_map_worker_status +BEFORE INSERT OR UPDATE OF status, tone_map_mode, audio_recipe_version ON public.download_artifacts +FOR EACH ROW +EXECUTE FUNCTION public.fence_tone_map_artifact_worker_status(); + +ALTER TABLE public.download_artifacts + VALIDATE CONSTRAINT download_artifacts_status_check; + +CREATE INDEX CONCURRENTLY download_artifacts_lease_idx ON public.download_artifacts (lease_expires_at) + WHERE status IN ('running', 'tone_map_running', 'audio_v2_running'); +CREATE INDEX CONCURRENTLY download_artifacts_lru_idx ON public.download_artifacts (last_used_at) + WHERE status IN ('ready', 'tone_map_ready', 'audio_v2_ready'); diff --git a/web/src/api/v2/schema.ts b/web/src/api/v2/schema.ts index 9057100f11..0a2656e2d1 100644 --- a/web/src/api/v2/schema.ts +++ b/web/src/api/v2/schema.ts @@ -22166,6 +22166,7 @@ export interface components { format: string; hearing_impaired: boolean; language: string; + title?: string; }; OnboardingCapabilitiesOutputBody: { /** @description Whether the current principal may use the capability */ From 12caeb2ac288d43f1c8b5386dc1d5cef1a9471d3 Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:05:45 -0400 Subject: [PATCH 2/4] fix(downloads): freeze prepared audio tracks and copy only MP4-safe codecs A ready multi-track download's manifest now describes the audio tracks recorded when the file became ready, not the source's current probe, so replacing and rescanning the source cannot advertise tracks the MP4 lacks. A viewer selection that no longer names the same-language track falls back to the file's default. Remuxes copy an audio track only when MP4 can store its codec; a negotiated passthrough codec such as TrueHD, DTS, or PCM is encoded to AAC instead of failing the mux. The track-recipe migration's column additions are now idempotent so a rerun after a partial Up succeeds. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/downloads-api.md | 8 ++- internal/downloads/artifact.go | 3 +- internal/downloads/artifact_repo.go | 28 +++++++--- internal/downloads/artifact_repo_test.go | 36 ++++++++----- internal/downloads/artifacts.go | 4 +- internal/downloads/manifest.go | 51 ++++++++++++------- internal/downloads/manifest_test.go | 38 ++++++++++++++ internal/downloads/repo_managed_test.go | 2 +- internal/playback/prepare_tracks.go | 24 ++++++--- internal/playback/prepare_tracks_test.go | 19 +++++++ ...30_fence_track_recipe_artifact_workers.sql | 9 +++- 11 files changed, 172 insertions(+), 50 deletions(-) diff --git a/docs/downloads-api.md b/docs/downloads-api.md index dfe7c3c00d..df0d09d80a 100644 --- a/docs/downloads-api.md +++ b/docs/downloads-api.md @@ -134,8 +134,12 @@ Yes, manifests include metadata needed to make the offline item feel native: Prepared remux/transcode files keep every source audio track in source order, so `audio_tracks[].index` and `selected_audio_track_index` address positions in the delivered MP4. Transcodes encode each track to stereo AAC; remuxes copy -tracks that share the primary track's codec (or are AAC/MP3) and encode the -rest to stereo AAC. Embedded plain-text subtitles (SRT, WebVTT) are carried +tracks that share the primary track's codec (or are AAC/MP3) when MP4 can +store that codec (AAC, MP3, AC-3, E-AC-3, ALAC) and encode the rest to stereo +AAC. The server records the audio tracks when the file becomes ready, so a +later rescan of a replaced source does not change `audio_tracks`; a +`selected_audio_track_index` that no longer names the same-language track +falls back to the file's default track. Embedded plain-text subtitles (SRT, WebVTT) are carried inside the MP4 as timed text, with their language, title, and forced flag. MP4 timed text would drop ASS/SSA styling, drawing commands, and overlapping events, and MP4 cannot store bitmap subtitles, so each embedded ASS/SSA track is diff --git a/internal/downloads/artifact.go b/internal/downloads/artifact.go index a77d713b3f..929774b2db 100644 --- a/internal/downloads/artifact.go +++ b/internal/downloads/artifact.go @@ -61,7 +61,8 @@ type Artifact struct { CodecVideo string CodecAudio string AudioRecipeVersion string - TrackRecipeVersion string // playback.PreparedTracksRecipeVersion; empty = legacy single-audio layout + TrackRecipeVersion string // playback.PreparedTracksRecipeVersion; empty = legacy single-audio layout + PreparedAudioTracks []OfflineAudioTrack // multi-track audio inventory, frozen when the file became ready Resolution string AudioTrackIndex int TargetBitrateKbps int diff --git a/internal/downloads/artifact_repo.go b/internal/downloads/artifact_repo.go index f967208940..be2c3865bb 100644 --- a/internal/downloads/artifact_repo.go +++ b/internal/downloads/artifact_repo.go @@ -2,6 +2,7 @@ package downloads import ( "context" + "encoding/json" "errors" "fmt" "time" @@ -12,7 +13,7 @@ import ( "github.com/Silo-Server/silo-server/internal/tonemap" ) -const artifactColumns = `id, media_file_id, format, params_hash, container, codec_video, codec_audio, audio_recipe_version, track_recipe_version, +const artifactColumns = `id, media_file_id, format, params_hash, container, codec_video, codec_audio, audio_recipe_version, track_recipe_version, prepared_audio_tracks, resolution, audio_track_index, target_bitrate_kbps, tone_map_policy, tone_map_mode, tone_map_source_kind, tone_map_recipe_version, tone_map_preflight_required, tone_map_source_revision, tone_map_dv_config_present, tone_map_dv_bl_compat_id_present, tone_map_dv_bl_present, tone_map_dv_rpu_present, output_path, origin_node_id, origin_node_url, origin_node_group, origin_artifact_id, file_size, status, error_message, @@ -46,8 +47,9 @@ func NewArtifactRepository(pool *pgxpool.Pool) *ArtifactRepository { func scanArtifact(row pgx.Row) (*Artifact, error) { var a Artifact var leaseOwner *string + var preparedAudio []byte if err := row.Scan( - &a.ID, &a.MediaFileID, &a.Format, &a.ParamsHash, &a.Container, &a.CodecVideo, &a.CodecAudio, &a.AudioRecipeVersion, &a.TrackRecipeVersion, + &a.ID, &a.MediaFileID, &a.Format, &a.ParamsHash, &a.Container, &a.CodecVideo, &a.CodecAudio, &a.AudioRecipeVersion, &a.TrackRecipeVersion, &preparedAudio, &a.Resolution, &a.AudioTrackIndex, &a.TargetBitrateKbps, &a.ToneMapPolicy, &a.ToneMapMode, &a.ToneMapSourceKind, &a.ToneMapRecipeVersion, &a.ToneMapPreflightRequired, &a.ToneMapSourceRevision, &a.ToneMapDVConfigPresent, &a.ToneMapDVBLCompatIDPresent, &a.ToneMapDVBLPresent, &a.ToneMapDVRPUPresent, &a.OutputPath, &a.OriginNodeID, &a.OriginNodeURL, &a.OriginNodeGroup, &a.OriginArtifactID, &a.FileSize, &a.Status, &a.ErrorMessage, @@ -57,6 +59,11 @@ func scanArtifact(row pgx.Row) (*Artifact, error) { return nil, err } a.LeaseOwner = deref(leaseOwner) + if len(preparedAudio) > 0 { + if err := json.Unmarshal(preparedAudio, &a.PreparedAudioTracks); err != nil { + return nil, fmt.Errorf("decoding prepared audio tracks: %w", err) + } + } return &a, nil } @@ -174,13 +181,21 @@ func (r *ArtifactRepository) Heartbeat(ctx context.Context, id, owner string, le return tag.RowsAffected() > 0, nil } -// MarkReady transitions a job to ready, records its size/path, and clears the -// lease. The write is fenced on (lease_owner, status='running') so a worker that +// MarkReady transitions a job to ready, records its size/path and the audio +// inventory of a multi-track file (nil otherwise), and clears the lease. The write is fenced on (lease_owner, status='running') so a worker that // lost its lease — e.g. a slow encode whose lease expired and was reclaimed by // another node — cannot flip a row it no longer owns. Returns false when the // fence rejected the write (the lease was lost); the caller must then NOT flip // linked downloads, leaving that to the current owner. -func (r *ArtifactRepository) MarkReady(ctx context.Context, id, owner, outputPath string, originNodeID int, originNodeURL, originNodeGroup, originArtifactID string, fileSize int64) (bool, error) { +func (r *ArtifactRepository) MarkReady(ctx context.Context, id, owner, outputPath string, originNodeID int, originNodeURL, originNodeGroup, originArtifactID string, fileSize int64, preparedAudioTracks []OfflineAudioTrack) (bool, error) { + var preparedAudio []byte + if preparedAudioTracks != nil { + encoded, err := json.Marshal(preparedAudioTracks) + if err != nil { + return false, fmt.Errorf("encoding prepared audio tracks: %w", err) + } + preparedAudio = encoded + } tag, err := r.pool.Exec(ctx, `UPDATE download_artifacts SET status = CASE @@ -191,10 +206,11 @@ func (r *ArtifactRepository) MarkReady(ctx context.Context, id, owner, outputPat END, output_path = $2, origin_node_id = $3, origin_node_url = $4, origin_node_group = $5, origin_artifact_id = $6, file_size = $7, error_message = '', + prepared_audio_tracks = $9, completed_at = now(), last_used_at = now(), lease_owner = NULL, lease_expires_at = NULL, next_retry_at = NULL WHERE id = $1 AND lease_owner = $8 AND status IN ('running', 'tone_map_running', 'audio_v2_running', 'tracks_v1_running')`, - id, outputPath, originNodeID, originNodeURL, originNodeGroup, originArtifactID, fileSize, owner, + id, outputPath, originNodeID, originNodeURL, originNodeGroup, originArtifactID, fileSize, owner, preparedAudio, ) if err != nil { return false, fmt.Errorf("marking artifact ready: %w", err) diff --git a/internal/downloads/artifact_repo_test.go b/internal/downloads/artifact_repo_test.go index 963a683646..6ca8e9c419 100644 --- a/internal/downloads/artifact_repo_test.go +++ b/internal/downloads/artifact_repo_test.go @@ -6,6 +6,7 @@ import ( "fmt" "os" "path/filepath" + "reflect" "testing" "time" @@ -163,7 +164,7 @@ func TestArtifactQueueClaimAndLeaseRecovery(t *testing.T) { if err != nil || claim2.ID != row.ID || claim2.Attempts != 2 { t.Fatalf("reclaim ClaimNext = (%+v, %v), want attempts=2", claim2, err) } - if applied, err := repo.MarkReady(ctx, row.ID, "worker-2", claim2.OutputPath, 0, "", "", "", 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "worker-2", claim2.OutputPath, 0, "", "", "", 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v), want (true, nil)", applied, err) } done, err := repo.GetByKey(ctx, fileID, "transcode", "hash-recovery") @@ -252,7 +253,7 @@ func TestToneMapArtifactQueueRejectsLegacyWorkers(t *testing.T) { if err != nil || claim.Status != ArtifactToneMapRunning { t.Fatalf("final ClaimNext = (%+v, %v), want tone-map running", claim, err) } - if applied, err := repo.MarkReady(ctx, row.ID, "current-worker", row.OutputPath, 0, "", "", "", 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "current-worker", row.OutputPath, 0, "", "", "", 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v), want applied", applied, err) } ready, err := repo.GetByID(ctx, row.ID) @@ -340,7 +341,7 @@ func TestAudioV2ArtifactQueueRejectsMergeBaseWorkers(t *testing.T) { if err != nil || claim.Status != ArtifactAudioV2Running { t.Fatalf("final ClaimNext = (%+v, %v), want audio-v2 running", claim, err) } - if applied, err := repo.MarkReady(ctx, row.ID, "current-worker", row.OutputPath, 0, "", "", "", 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "current-worker", row.OutputPath, 0, "", "", "", 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v), want applied", applied, err) } ready, err := repo.GetByID(ctx, row.ID) @@ -431,13 +432,20 @@ func TestTrackRecipeArtifactQueueRejectsMergeBaseWorkers(t *testing.T) { if claim, err = repo.ClaimNext(ctx, "current-worker", time.Minute); err != nil || claim.Status != ArtifactTracksRunning { t.Fatalf("final ClaimNext = (%+v, %v), want tracks running", claim, err) } - if applied, err := repo.MarkReady(ctx, row.ID, "current-worker", row.OutputPath, 0, "", "", "", 4242); err != nil || !applied { + preparedAudio := []OfflineAudioTrack{ + {Index: 0, Language: "en", Codec: "eac3", Channels: 6, Default: true}, + {Index: 1, Language: "ja", Codec: "aac", Channels: 2, Layout: "stereo", Bitrate: 192}, + } + if applied, err := repo.MarkReady(ctx, row.ID, "current-worker", row.OutputPath, 0, "", "", "", 4242, preparedAudio); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v), want applied", applied, err) } ready, err := repo.GetByID(ctx, row.ID) if err != nil || ready.Status != ArtifactTracksReady || !artifactReady(ready) { t.Fatalf("tracks ready row = (%+v, %v), want tracks ready", ready, err) } + if !reflect.DeepEqual(ready.PreparedAudioTracks, preparedAudio) { + t.Fatalf("prepared audio tracks = %+v, want %+v", ready.PreparedAudioTracks, preparedAudio) + } if total, err := repo.TotalReadyBytes(ctx); err != nil || total < 4242 { t.Fatalf("TotalReadyBytes = (%d, %v), want the tracks artifact counted", total, err) } @@ -527,7 +535,7 @@ func TestArtifactMarkFencedByOwner(t *testing.T) { } // A non-owner cannot mark the job ready or failed. - if applied, err := repo.MarkReady(ctx, row.ID, "owner-2", "/tmp/x.mp4", 0, "", "", "", 10); err != nil || applied { + if applied, err := repo.MarkReady(ctx, row.ID, "owner-2", "/tmp/x.mp4", 0, "", "", "", 10, nil); err != nil || applied { t.Fatalf("MarkReady(non-owner) = (%v, %v), want (false, nil)", applied, err) } if _, applied, err := repo.MarkFailedOrRetry(ctx, row.ID, "owner-2", "boom", time.Second); err != nil || applied { @@ -541,7 +549,7 @@ func TestArtifactMarkFencedByOwner(t *testing.T) { } // The real owner succeeds. - if applied, err := repo.MarkReady(ctx, row.ID, "owner-1", "/tmp/x.mp4", 0, "", "", "", 10); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "owner-1", "/tmp/x.mp4", 0, "", "", "", 10, nil); err != nil || !applied { t.Fatalf("MarkReady(owner) = (%v, %v), want (true, nil)", applied, err) } } @@ -556,7 +564,7 @@ func TestArtifactRemoteLocatorRoundTripsAndRequeueClearsIt(t *testing.T) { if _, err := repo.ClaimNext(ctx, "worker", time.Minute); err != nil { t.Fatal(err) } - if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 17, "http://transcode", "host-a", "artifact-opaque", 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 17, "http://transcode", "host-a", "artifact-opaque", 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v)", applied, err) } ready, err := repo.GetByID(ctx, row.ID) @@ -600,7 +608,7 @@ func TestArtifactReadyPersistsRefreshedOriginLocator(t *testing.T) { if _, err := repo.ClaimNext(ctx, "worker", time.Minute); err != nil { t.Fatal(err) } - if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 17, "http://old-url", "old-group", "artifact-refresh", 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 17, "http://old-url", "old-group", "artifact-refresh", 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v)", applied, err) } manager := &ArtifactManager{ @@ -635,7 +643,7 @@ func TestArtifactReadyMapsRemovedOriginToInactiveWhenRequeueLosesFence(t *testin if _, err := repo.ClaimNext(ctx, "worker", time.Minute); err != nil { t.Fatal(err) } - if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 17, "http://removed", "host-a", "artifact-removed", 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 17, "http://removed", "host-a", "artifact-removed", 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v)", applied, err) } manager := &ArtifactManager{ @@ -663,7 +671,7 @@ func TestArtifactRecoveryContinuesAfterStaleLocatorRefresh(t *testing.T) { if _, err := repo.ClaimNext(ctx, "worker", time.Minute); err != nil { t.Fatal(err) } - if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 17, "http://old-url", "host-a", fmt.Sprintf("artifact-%d", i), 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 17, "http://old-url", "host-a", fmt.Sprintf("artifact-%d", i), 4242, nil); err != nil || !applied { t.Fatalf("MarkReady(%d) = (%v, %v)", i, applied, err) } artifact, err := repo.GetByID(ctx, row.ID) @@ -750,7 +758,7 @@ func TestArtifactRecoveryDeletesWrongSizedRemoteBeforeRequeue(t *testing.T) { if _, err := repo.ClaimNext(ctx, "worker", time.Minute); err != nil { t.Fatal(err) } - if applied, err := repo.MarkReady(ctx, row.ID, "worker", row.OutputPath, 17, "http://transcode", "host-a", "artifact-truncated", 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "worker", row.OutputPath, 17, "http://transcode", "host-a", "artifact-truncated", 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v)", applied, err) } @@ -823,7 +831,7 @@ func TestArtifactRemoteRequeueAtomicallyQueuesCleanup(t *testing.T) { if _, err := repo.ClaimNext(ctx, "worker", time.Minute); err != nil { t.Fatal(err) } - if applied, err := repo.MarkReady(ctx, row.ID, "worker", row.OutputPath, 23, "http://transcode-old", "host-a", "artifact-abandoned", 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "worker", row.OutputPath, 23, "http://transcode-old", "host-a", "artifact-abandoned", 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v)", applied, err) } ready, err := repo.GetByID(ctx, row.ID) @@ -914,7 +922,7 @@ func TestProxyMissingReportFencesCompleteRemoteLocator(t *testing.T) { if _, err := repo.ClaimNext(ctx, "worker", time.Minute); err != nil { t.Fatal(err) } - if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 23, "http://transcode-current", "host-a", "artifact-missing", 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "worker", "", 23, "http://transcode-current", "host-a", "artifact-missing", 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v)", applied, err) } manager := NewArtifactManager(repo, nil, nil, nil, "proxy-test", nil, nil) @@ -1125,7 +1133,7 @@ func readyArtifactForRecovery(t *testing.T, repo *ArtifactRepository, pool *pgxp if originArtifactID != "" { originNodeID, originNodeURL = 31, "http://transcode-recovery" } - if applied, err := repo.MarkReady(ctx, row.ID, "worker", row.OutputPath, originNodeID, originNodeURL, "", originArtifactID, 4242); err != nil || !applied { + if applied, err := repo.MarkReady(ctx, row.ID, "worker", row.OutputPath, originNodeID, originNodeURL, "", originArtifactID, 4242, nil); err != nil || !applied { t.Fatalf("MarkReady = (%v, %v)", applied, err) } if _, err := pool.Exec(ctx, `UPDATE download_artifacts SET last_used_at = now() - interval '1 hour' WHERE id = $1`, row.ID); err != nil { diff --git a/internal/downloads/artifacts.go b/internal/downloads/artifacts.go index 2c46920f4b..59b41e12d0 100644 --- a/internal/downloads/artifacts.go +++ b/internal/downloads/artifacts.go @@ -913,7 +913,9 @@ func (m *ArtifactManager) encodeOne(ctx context.Context, a *Artifact) { // remote-missing requeue can safely fall back to integrated preparation. outputPath = a.OutputPath } - applied, err := m.repo.MarkReady(ctx, a.ID, m.owner, outputPath, prepared.OriginNodeID, prepared.OriginNodeURL, prepared.OriginNodeGroup, prepared.OriginArtifactID, size) + // The fingerprint check above tied these bytes to file's current probe; + // freeze the audio inventory it describes with the ready transition. + applied, err := m.repo.MarkReady(ctx, a.ID, m.owner, outputPath, prepared.OriginNodeID, prepared.OriginNodeURL, prepared.OriginNodeGroup, prepared.OriginArtifactID, size, preparedAudioTracks(file, a)) if err != nil { slog.ErrorContext(ctx, "marking artifact ready failed", "component", "downloads", "artifact_id", a.ID, "error", err) m.cleanupRejectedPrepared(ctx, a.ID, prepared) diff --git a/internal/downloads/manifest.go b/internal/downloads/manifest.go index eb654b47c7..47459b7c93 100644 --- a/internal/downloads/manifest.go +++ b/internal/downloads/manifest.go @@ -334,23 +334,47 @@ func applyArtifactParams(m *OfflineManifest, a *Artifact, file *models.MediaFile } } -// applyPreparedAudioTracks describes a multi-track prepared file. It keeps -// every source audio track in source order, so output positions equal source -// positions and the viewer's catalog selection stays valid; encoded tracks -// report the AAC output layout. +// applyPreparedAudioTracks describes a multi-track prepared file from the +// audio inventory frozen when it became ready, falling back to the current +// source probe for an artifact that is not ready yet. Output positions equal +// the prepared source's positions, so the viewer's catalog selection stays +// valid unless the source has since changed at that position. func applyPreparedAudioTracks(m *OfflineManifest, a *Artifact, file *models.MediaFile) { - if file == nil { + tracks := a.PreparedAudioTracks + if len(tracks) == 0 { + tracks = preparedAudioTracks(file, a) + } + if len(tracks) == 0 { return } + m.AudioTracks = tracks + fileDefault := 0 + for i, track := range tracks { + if track.Default { + fileDefault = i + break + } + } + m.CodecAudio = tracks[fileDefault].Codec + if selected := m.SelectedAudioTrackIndex; selected != nil && *selected >= 0 && *selected < len(tracks) && + file != nil && *selected < len(file.AudioTracks) && file.AudioTracks[*selected].Language == tracks[*selected].Language { + return + } + m.SelectedAudioTrackIndex = &fileDefault +} + +// preparedAudioTracks describes the audio streams a multi-track prepared file +// built from file contains: every source track in source order, with encoded +// tracks reporting the AAC output layout. +func preparedAudioTracks(file *models.MediaFile, a *Artifact) []OfflineAudioTrack { + if file == nil || a == nil || a.TrackRecipeVersion == "" { + return nil + } plan := playback.PlanPreparedTracks(file, a.CodecAudio, a.AudioTrackIndex) tracks := toOfflineAudioTracks(file.AudioTracks) channels, bitrateKbps := playback.ResolveAACOutputV3(0, 0) - fileDefault := -1 for i, track := range plan.Audio { tracks[i].Default = track.Default - if track.Default { - fileDefault = i - } if track.Codec == playback.PreparedAudioAAC { tracks[i].Codec = playback.PreparedAudioAAC tracks[i].Channels = channels @@ -358,14 +382,7 @@ func applyPreparedAudioTracks(m *OfflineManifest, a *Artifact, file *models.Medi tracks[i].Bitrate = bitrateKbps } } - m.AudioTracks = tracks - if selected := m.SelectedAudioTrackIndex; selected != nil && *selected >= 0 && *selected < len(tracks) { - return - } - m.SelectedAudioTrackIndex = nil - if fileDefault >= 0 { - m.SelectedAudioTrackIndex = &fileDefault - } + return tracks } // buildSubtitles enumerates external (sidecar), embedded sidecar, and diff --git a/internal/downloads/manifest_test.go b/internal/downloads/manifest_test.go index a43a3f0d8b..b1da5f7e08 100644 --- a/internal/downloads/manifest_test.go +++ b/internal/downloads/manifest_test.go @@ -7,6 +7,7 @@ import ( "errors" "net/http" "net/http/httptest" + "reflect" "testing" "github.com/Silo-Server/silo-server/internal/catalog" @@ -228,6 +229,43 @@ func TestManifestDescribesMultiTrackArtifact(t *testing.T) { } } +// TestManifestDescribesFrozenAudioAfterSourceReprobe verifies a ready +// multi-track file keeps the audio inventory it was prepared with after the +// source is replaced at the same path and re-probed. +func TestManifestDescribesFrozenAudioAfterSourceReprobe(t *testing.T) { + prepared := []OfflineAudioTrack{ + {Index: 0, Language: "en", Codec: "aac", Channels: 2, Layout: "stereo", Bitrate: 192, Default: true}, + {Index: 1, Language: "ja", Codec: "aac", Channels: 2, Layout: "stereo", Bitrate: 192}, + } + b, dl := preparedManifestFixture(&Artifact{ + ID: "a1", Container: "mp4", CodecVideo: "h264", CodecAudio: "aac", Resolution: "1080p", + AudioTrackIndex: -1, TrackRecipeVersion: playback.PreparedTracksRecipeVersion, + PreparedAudioTracks: prepared, + }) + reprobed := []models.AudioTrack{ + {Codec: "ac3", Channels: 6, Language: "ja", Default: true}, + {Codec: "truehd", Channels: 8, Language: "en"}, + {Codec: "dts", Channels: 6, Language: "fr"}, + } + for _, selected := range []int{2, 0} { + b.detail.(fakeManifestSource).detail.Versions[0].AudioTracks = reprobed + b.detail.(fakeManifestSource).detail.Versions[0].EffectiveAudioTrackIndex = &selected + b.fileRepo.(fakeFileResolver).file.AudioTracks = reprobed + m, err := b.Build(context.Background(), dl, catalog.AccessFilter{}) + if err != nil { + t.Fatalf("Build: %v", err) + } + if !reflect.DeepEqual(m.AudioTracks, prepared) { + t.Fatalf("audio tracks = %+v, want the frozen inventory %+v", m.AudioTracks, prepared) + } + // Neither the out-of-range position nor the position now holding a + // different language describes a delivered track the viewer chose. + if m.SelectedAudioTrackIndex == nil || *m.SelectedAudioTrackIndex != 0 { + t.Fatalf("selection %d: selected = %v, want the delivered default 0", selected, m.SelectedAudioTrackIndex) + } + } +} + // TestManifestDescribesLegacySingleTrackArtifact keeps already-prepared files // described as the one audio stream they contain, without PGS sidecars. func TestManifestDescribesLegacySingleTrackArtifact(t *testing.T) { diff --git a/internal/downloads/repo_managed_test.go b/internal/downloads/repo_managed_test.go index e935bceb84..efb09bc39a 100644 --- a/internal/downloads/repo_managed_test.go +++ b/internal/downloads/repo_managed_test.go @@ -289,7 +289,7 @@ func TestReconcileLinkedDownloads(t *testing.T) { if _, err := arepo.ClaimNext(ctx, "w", time.Minute); err != nil { t.Fatalf("claim ready artifact: %v", err) } - if ok, err := arepo.MarkReady(ctx, readyArt.ID, "w", "/tmp/ready.mp4", 0, "", "", "", 4242); err != nil || !ok { + if ok, err := arepo.MarkReady(ctx, readyArt.ID, "w", "/tmp/ready.mp4", 0, "", "", "", 4242, nil); err != nil || !ok { t.Fatalf("MarkReady = (%v, %v)", ok, err) } diff --git a/internal/playback/prepare_tracks.go b/internal/playback/prepare_tracks.go index 5fa1299489..f5ab6bbd79 100644 --- a/internal/playback/prepare_tracks.go +++ b/internal/playback/prepare_tracks.go @@ -69,9 +69,10 @@ func PreparedTracksAvailable(file *models.MediaFile) bool { // PlanPreparedTracks derives the stream layout of a prepared download from the // probed source. Every audio track is kept in source order. With a "copy" -// audio target a track is copied when it shares the primary track's codec -// (which client negotiation verified) or is universally decodable, otherwise -// it is encoded to AAC; an "aac" target encodes every track. audioTrackIndex +// audio target a track is copied when the MP4 muxer accepts its codec and it +// either shares the primary track's codec (which client negotiation verified) +// or is universally decodable, otherwise it is encoded to AAC; an "aac" target +// encodes every track. audioTrackIndex // marks the default track, falling back to the source default and then the // first track. Embedded plain-text subtitles become mov_text. ASS/SSA keep // their styling, typesetting, and overlapping events only as sidecars, and @@ -143,11 +144,22 @@ func preparedDefaultAudioIndex(tracks []models.AudioTrack, requested int) int { return 0 } -// preparedAudioCopyable reports whether a non-primary audio track may be -// stream-copied into the prepared MP4 alongside a copied primary track. +// preparedMP4AudioCodecs lists audio codecs FFmpeg's MP4 muxer stores without +// experimental flags. Client negotiation checks decode support, not the +// container, so a passthrough codec such as TrueHD, DTS, or PCM is encoded. +var preparedMP4AudioCodecs = map[string]bool{ + audioCodecAACV3: true, + audioCodecMP3: true, + "ac3": true, + "eac3": true, + "alac": true, +} + +// preparedAudioCopyable reports whether an audio track may be stream-copied +// into the prepared MP4 under a "copy" audio target. func preparedAudioCopyable(codec, primaryCodec string) bool { normalized := normalizeCodecV3(codec) - if normalized == "" { + if !preparedMP4AudioCodecs[normalized] { return false } return normalized == primaryCodec || normalized == audioCodecAACV3 || normalized == audioCodecMP3 diff --git a/internal/playback/prepare_tracks_test.go b/internal/playback/prepare_tracks_test.go index c00ae7875e..e4ecdbdc85 100644 --- a/internal/playback/prepare_tracks_test.go +++ b/internal/playback/prepare_tracks_test.go @@ -67,6 +67,25 @@ func TestPlanPreparedTracksKeepsEveryAudioTrackAndTextSubtitle(t *testing.T) { } } +func TestPlanPreparedTracksEncodesCodecsMP4CannotCarry(t *testing.T) { + file := &models.MediaFile{ + CodecAudio: "truehd", + AudioTracks: []models.AudioTrack{ + {Codec: "truehd", Channels: 8, Language: "en"}, + {Codec: "truehd", Channels: 6, Language: "de"}, + {Codec: "pcm_s24le", Channels: 2, Language: "fr"}, + {Codec: "aac", Channels: 2, Language: "ja"}, + }, + } + plan := PlanPreparedTracks(file, "copy", -1) + want := []string{"aac", "aac", "aac", "copy"} + for i, track := range plan.Audio { + if track.Codec != want[i] { + t.Fatalf("track %d codec = %q, want %q (plan %+v)", i, track.Codec, want[i], plan.Audio) + } + } +} + func TestPreparedSubtitleSidecarFormat(t *testing.T) { for codec, want := range map[string]string{ "ass": "ass", "ssa": "ass", "hdmv_pgs_subtitle": "sup", "pgssub": "sup", diff --git a/migrations/sql/20260928222330_fence_track_recipe_artifact_workers.sql b/migrations/sql/20260928222330_fence_track_recipe_artifact_workers.sql index a7ce51acfd..a5ca861b38 100644 --- a/migrations/sql/20260928222330_fence_track_recipe_artifact_workers.sql +++ b/migrations/sql/20260928222330_fence_track_recipe_artifact_workers.sql @@ -6,11 +6,15 @@ -- discriminator and a status family that the merge-base ClaimNext predicate -- does not recognize. The track recipe outranks the audio and tone-map -- families: its execution fingerprint already covers both. +-- prepared_audio_tracks freezes the audio streams of the delivered file when +-- it becomes ready, so a later re-probe of the source cannot change how the +-- offline manifest describes bytes that were already prepared. DROP INDEX CONCURRENTLY IF EXISTS public.download_artifacts_lease_idx; DROP INDEX CONCURRENTLY IF EXISTS public.download_artifacts_lru_idx; ALTER TABLE public.download_artifacts - ADD COLUMN track_recipe_version text NOT NULL DEFAULT '', + ADD COLUMN IF NOT EXISTS track_recipe_version text NOT NULL DEFAULT '', + ADD COLUMN IF NOT EXISTS prepared_audio_tracks jsonb, DROP CONSTRAINT download_artifacts_status_check, ADD CONSTRAINT download_artifacts_status_check CHECK (status IN ( @@ -139,7 +143,8 @@ $$; -- +goose StatementEnd ALTER TABLE public.download_artifacts - DROP COLUMN track_recipe_version, + DROP COLUMN IF EXISTS track_recipe_version, + DROP COLUMN IF EXISTS prepared_audio_tracks, DROP CONSTRAINT download_artifacts_status_check, ADD CONSTRAINT download_artifacts_status_check CHECK (status IN ( From 1a869a78dfeb3b9be3b392e2c576a1b220ad2b44 Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:27:15 -0400 Subject: [PATCH 3/4] chore(api): refresh contract fixture after rebase --- contracts/api/v2/fixtures/get_system_info_ok.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/contracts/api/v2/fixtures/get_system_info_ok.json b/contracts/api/v2/fixtures/get_system_info_ok.json index 60e0b102ee..3a5fecd6c6 100644 --- a/contracts/api/v2/fixtures/get_system_info_ok.json +++ b/contracts/api/v2/fixtures/get_system_info_ok.json @@ -1,7 +1,7 @@ { "server_version": "unavailable", "api_major": 2, - "contract_digest": "1b83c215553d5b1bc2cc54037cd789ec0a1364fbb37e41d0dfdfd7cf135c9552", + "contract_digest": "a090db5a3b9114c65e340fd57ad7ca081b7452b01ab4b4d56d72f9ffc338b81b", "links": { "openapi": "/api/v2/openapi.json", "capabilities": "/api/v2/capabilities", From 9468ab83fe34d7e3050aef52f407d9fe460e6c73 Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:33:55 -0400 Subject: [PATCH 4/4] chore(downloads): name prepared track codec and layout constants --- internal/downloads/manifest.go | 4 +++- internal/playback/prepare_tracks.go | 6 ++++-- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/internal/downloads/manifest.go b/internal/downloads/manifest.go index 47459b7c93..ce3a0ee48c 100644 --- a/internal/downloads/manifest.go +++ b/internal/downloads/manifest.go @@ -18,6 +18,8 @@ import ( // manifestVersion is bumped whenever the OfflineManifest DTO shape changes. const manifestVersion = 2 +const preparedAudioLayoutStereo = "stereo" + // apiDownloadsPrefix is the namespace every offline asset reference is minted // under. Manifests are stored and handed to clients verbatim, so the reference // has to stay resolvable after the /api/v1 tombstone; the v2 projection only @@ -378,7 +380,7 @@ func preparedAudioTracks(file *models.MediaFile, a *Artifact) []OfflineAudioTrac if track.Codec == playback.PreparedAudioAAC { tracks[i].Codec = playback.PreparedAudioAAC tracks[i].Channels = channels - tracks[i].Layout = "stereo" + tracks[i].Layout = preparedAudioLayoutStereo tracks[i].Bitrate = bitrateKbps } } diff --git a/internal/playback/prepare_tracks.go b/internal/playback/prepare_tracks.go index f5ab6bbd79..11500b264f 100644 --- a/internal/playback/prepare_tracks.go +++ b/internal/playback/prepare_tracks.go @@ -22,6 +22,8 @@ const ( const ( audioCodecMP3 = "mp3" + audioCodecAC3 = "ac3" + audioCodecEAC3 = "eac3" subtitleCodecTextV3 = "text" ) @@ -150,8 +152,8 @@ func preparedDefaultAudioIndex(tracks []models.AudioTrack, requested int) int { var preparedMP4AudioCodecs = map[string]bool{ audioCodecAACV3: true, audioCodecMP3: true, - "ac3": true, - "eac3": true, + audioCodecAC3: true, + audioCodecEAC3: true, "alac": true, }