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", 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..df0d09d80a 100644 --- a/docs/downloads-api.md +++ b/docs/downloads-api.md @@ -128,8 +128,27 @@ 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) 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 +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 +632,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 +1621,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..929774b2db 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,8 @@ type Artifact struct { CodecVideo string CodecAudio string AudioRecipeVersion string + 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 @@ -129,13 +138,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..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, +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.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 } @@ -67,15 +74,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 +136,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 +145,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 +172,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 { @@ -173,26 +181,36 @@ 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 + 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' 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')`, - id, outputPath, originNodeID, originNodeURL, originNodeGroup, originArtifactID, fileSize, owner, + 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, preparedAudio, ) if err != nil { return false, fmt.Errorf("marking artifact ready: %w", err) @@ -214,6 +232,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 +240,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 +266,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 +299,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 +340,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 +368,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 +468,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 +476,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 +532,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 +548,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 +562,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 +586,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 +629,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..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) @@ -365,6 +366,113 @@ 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) + } + 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) + } + + 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) { @@ -427,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 { @@ -441,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) } } @@ -456,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) @@ -500,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{ @@ -535,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{ @@ -563,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) @@ -650,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) } @@ -723,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) @@ -814,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) @@ -1025,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/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..59b41e12d0 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 @@ -910,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) @@ -1089,6 +1094,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 +1123,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..ce3a0ee48c 100644 --- a/internal/downloads/manifest.go +++ b/internal/downloads/manifest.go @@ -11,12 +11,15 @@ 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" ) // 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 @@ -55,6 +58,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 +286,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 +313,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 +336,64 @@ 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 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) { + 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) + for i, track := range plan.Audio { + tracks[i].Default = track.Default + if track.Codec == playback.PreparedAudioAAC { + tracks[i].Codec = playback.PreparedAudioAAC + tracks[i].Channels = channels + tracks[i].Layout = preparedAudioLayoutStereo + tracks[i].Bitrate = bitrateKbps + } + } + return tracks +} + +// 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 +404,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 +413,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..b1da5f7e08 100644 --- a/internal/downloads/manifest_test.go +++ b/internal/downloads/manifest_test.go @@ -5,10 +5,14 @@ import ( "context" "encoding/json" "errors" + "net/http" + "net/http/httptest" + "reflect" "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 +164,126 @@ 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) + } +} + +// 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) { + 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 +294,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 +315,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/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/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..11500b264f --- /dev/null +++ b/internal/playback/prepare_tracks.go @@ -0,0 +1,268 @@ +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" + audioCodecAC3 = "ac3" + audioCodecEAC3 = "eac3" + 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 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 +// 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 +} + +// 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, + audioCodecAC3: true, + audioCodecEAC3: 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 !preparedMP4AudioCodecs[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..e4ecdbdc85 --- /dev/null +++ b/internal/playback/prepare_tracks_test.go @@ -0,0 +1,281 @@ +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 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", + "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..a5ca861b38 --- /dev/null +++ b/migrations/sql/20260928222330_fence_track_recipe_artifact_workers.sql @@ -0,0 +1,168 @@ +-- +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. +-- 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 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 ( + '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 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 ( + '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 */