From d533301880a7b57ef3bd0f82c014acc3871c1964 Mon Sep 17 00:00:00 2001 From: Breadway Date: Fri, 24 Jul 2026 23:57:58 +0800 Subject: [PATCH] Add native GPU-accelerated AV1 transcode capability Post-grab async background swap, upgrade_locked flag to prevent the upgrade loop from re-inflating a locally-transcoded file, and a transcode-library CLI backfill for the existing catalog. Bitrate model calibrated against Silicon Valley's real HEVC bitrate, scaled by resolution and AV1's encoding efficiency. HDR/2160p and anime excluded from this first pass. Jellyfin session polling throttles batch encoding back to 1 stream during active playback. --- breadarr-shared/src/config.rs | 110 ++++++++ breadarr-shared/src/dto.rs | 1 + breadarrd/src/api/mod.rs | 1 + breadarrd/src/api/routes/health.rs | 1 + breadarrd/src/db.rs | 35 ++- breadarrd/src/importer/mod.rs | 189 ++++++++++++- breadarrd/src/jellyfin.rs | 29 ++ breadarrd/src/main.rs | 91 +++++- breadarrd/src/scheduler.rs | 48 +++- breadarrd/src/transcode/mod.rs | 434 +++++++++++++++++++++++++++++ 10 files changed, 924 insertions(+), 15 deletions(-) create mode 100644 breadarrd/src/transcode/mod.rs diff --git a/breadarr-shared/src/config.rs b/breadarr-shared/src/config.rs index a82cf4d..684d6c7 100644 --- a/breadarr-shared/src/config.rs +++ b/breadarr-shared/src/config.rs @@ -23,6 +23,8 @@ pub struct Config { pub sources: SourcesConfig, #[serde(default)] pub notifications: NotificationsConfig, + #[serde(default)] + pub transcode: TranscodeConfig, } /// Where the TUI's "add show" flow places new series by default. Sonarr/ @@ -277,6 +279,114 @@ pub struct JellyfinConfig { pub api_key: String, } +/// GPU-accelerated AV1 transcode: re-encodes freshly-grabbed and existing +/// library files down to a space-reasonable size instead of keeping +/// whatever the source release happened to be (REMUX, huge season packs, +/// etc). `enabled` defaults false — this needs a manual calibration pass +/// against real content on the target GPU before it's safe to run +/// unattended against a whole library. +#[derive(Debug, Clone, Deserialize)] +pub struct TranscodeConfig { + #[serde(default)] + pub enabled: bool, + /// Tight on purpose — the actual pace is bottlenecked by encode time + /// (minutes per file), not this interval; a short poll just means a + /// freshly-completed job's slot gets refilled promptly instead of + /// sitting idle for the rest of a longer interval. + #[serde(default = "default_transcode_poll_interval_secs")] + pub poll_interval_secs: u64, + #[serde(default = "default_vaapi_device")] + pub vaapi_device: String, + #[serde(default = "default_parallelism_min")] + pub parallelism_min: usize, + /// Ramp-up ceiling for concurrent encode streams during a backfill — + /// tune this against how many simultaneous `av1_vaapi` sessions the + /// target GPU can actually sustain before per-stream throughput starts + /// dropping, not just picked arbitrarily. + #[serde(default = "default_parallelism_max")] + pub parallelism_max: usize, + /// The "looks fine, no complaints" calibration reference: a real + /// bitrate (Mbps, in kbps here) from content already in the library at + /// `reference_height` that the user is happy with. New AV1 encodes are + /// targeted relative to this, scaled by resolution and AV1's encoding + /// efficiency, rather than picking a bitrate out of thin air. + #[serde(default = "default_reference_bitrate_kbps")] + pub reference_bitrate_kbps: u32, + #[serde(default = "default_reference_height")] + pub reference_height: u32, + /// AV1 reaches equivalent perceived quality to HEVC at a meaningfully + /// lower bitrate — this factor is applied on top of the resolution + /// scaling so the AV1 target isn't just a like-for-like copy of the + /// HEVC/H264 reference bitrate. Conservative (not maximally aggressive) + /// on purpose: erring toward "still clearly smaller" over "as small as + /// AV1 could theoretically go" leaves margin against visible artifacts. + #[serde(default = "default_av1_efficiency_factor")] + pub av1_efficiency_factor: f32, + /// HDR10/Dolby Vision metadata preservation through the GPU encoder + /// hasn't been verified yet — excluded from both the backfill and the + /// post-grab path until that's specifically checked on a few samples. + #[serde(default = "default_exclude_hdr")] + pub exclude_hdr: bool, + /// Excludes the 2160p tier from the first pass for the same reason as + /// `exclude_hdr` (most current 4K content in this library is HDR + /// anyway) — revisit once HDR handling is confirmed safe. + #[serde(default = "default_exclude_min_height")] + pub exclude_min_height: u32, +} + +impl Default for TranscodeConfig { + fn default() -> Self { + Self { + enabled: false, + poll_interval_secs: default_transcode_poll_interval_secs(), + vaapi_device: default_vaapi_device(), + parallelism_min: default_parallelism_min(), + parallelism_max: default_parallelism_max(), + reference_bitrate_kbps: default_reference_bitrate_kbps(), + reference_height: default_reference_height(), + av1_efficiency_factor: default_av1_efficiency_factor(), + exclude_hdr: default_exclude_hdr(), + exclude_min_height: default_exclude_min_height(), + } + } +} + +fn default_transcode_poll_interval_secs() -> u64 { + 60 +} + +fn default_vaapi_device() -> String { + "/dev/dri/renderD128".to_string() +} + +fn default_parallelism_min() -> usize { + 1 +} + +fn default_parallelism_max() -> usize { + 4 +} + +fn default_reference_bitrate_kbps() -> u32 { + 5320 +} + +fn default_reference_height() -> u32 { + 1080 +} + +fn default_av1_efficiency_factor() -> f32 { + 0.7 +} + +fn default_exclude_hdr() -> bool { + true +} + +fn default_exclude_min_height() -> u32 { + 2000 +} + /// TVDB v4 API key, exchanged for a short-lived JWT at request time. #[derive(Debug, Clone, Default, Deserialize)] pub struct TvdbConfig { diff --git a/breadarr-shared/src/dto.rs b/breadarr-shared/src/dto.rs index 6e103d3..070734e 100644 --- a/breadarr-shared/src/dto.rs +++ b/breadarr-shared/src/dto.rs @@ -143,6 +143,7 @@ pub struct HealthDetail { pub last_import_cycle: Option, pub last_search_cycle: Option, pub last_upgrade_cycle: Option, + pub last_transcode_cycle: Option, pub search_halted: bool, } diff --git a/breadarrd/src/api/mod.rs b/breadarrd/src/api/mod.rs index 2c6c309..6fa0ebe 100644 --- a/breadarrd/src/api/mod.rs +++ b/breadarrd/src/api/mod.rs @@ -78,6 +78,7 @@ pub struct CycleStatus { pub last_import: Option, pub last_search: Option, pub last_upgrade: Option, + pub last_transcode: Option, /// Set once the search-driven loop's consecutive-failure backoff hits /// its ceiling — still ticking at max backoff underneath (self-healing /// if the source recovers), but worth a loud, easy-to-spot signal that diff --git a/breadarrd/src/api/routes/health.rs b/breadarrd/src/api/routes/health.rs index 4333efd..e69aa37 100644 --- a/breadarrd/src/api/routes/health.rs +++ b/breadarrd/src/api/routes/health.rs @@ -20,6 +20,7 @@ pub async fn health(State(state): State) -> Json { last_import_cycle: status.last_import.as_ref().map(to_info), last_search_cycle: status.last_search.as_ref().map(to_info), last_upgrade_cycle: status.last_upgrade.as_ref().map(to_info), + last_transcode_cycle: status.last_transcode.as_ref().map(to_info), search_halted: status.search_halted, }) } diff --git a/breadarrd/src/db.rs b/breadarrd/src/db.rs index 88ff19a..d356d88 100644 --- a/breadarrd/src/db.rs +++ b/breadarrd/src/db.rs @@ -353,7 +353,26 @@ pub fn init(conn: &Connection) -> anyhow::Result<()> { fetched_at TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_torrent_fetch_hash ON torrent_fetch(torrent_hash); - CREATE INDEX IF NOT EXISTS idx_torrent_fetch_fetched_at ON torrent_fetch(fetched_at);", + CREATE INDEX IF NOT EXISTS idx_torrent_fetch_fetched_at ON torrent_fetch(fetched_at); + + -- One row per file queued for AV1 transcoding, whether from the + -- post-import async hook or the `transcode-library` backfill sweep. + -- Persisted (not an in-memory queue) so a `pending`/`running` row + -- left over from a daemon crash mid-encode just gets picked up + -- again on the next tick instead of silently vanishing. + CREATE TABLE IF NOT EXISTS transcode_job ( + id INTEGER PRIMARY KEY, + episode_file_id INTEGER NOT NULL REFERENCES episode_file(id) ON DELETE CASCADE, + status TEXT NOT NULL DEFAULT 'pending' + CHECK (status IN ('pending','running','done','failed','skipped')), + original_codec TEXT, + original_bytes INTEGER, + new_bytes INTEGER, + error TEXT, + queued_at TEXT NOT NULL, + finished_at TEXT + ); + CREATE INDEX IF NOT EXISTS idx_transcode_job_status ON transcode_job(status);", )?; // Progress watermark for stalled-download detection (added after the @@ -411,6 +430,18 @@ pub fn init(conn: &Connection) -> anyhow::Result<()> { // less-common stream metadata). See `ffprobe::MediaProbe::raw_json`. add_column_if_missing(conn, "media_file_probe", "raw_ffprobe_json", "TEXT")?; + // Set once a file has been through a successful local AV1 transcode — + // stops the upgrade cycle from treating it as still needing a bigger + // HEVC/H264 release, since `best_existing_score` otherwise only ever + // sees the stored `release.score` from original-grab time, which a + // local re-encode never touches. + add_column_if_missing( + conn, + "episode_file", + "upgrade_locked", + "INTEGER NOT NULL DEFAULT 0", + )?; + Ok(()) } @@ -473,7 +504,7 @@ mod tests { |row| row.get(0), ) .unwrap(); - assert_eq!(table_count, 17); + assert_eq!(table_count, 18); } #[test] diff --git a/breadarrd/src/importer/mod.rs b/breadarrd/src/importer/mod.rs index d8e807a..1aa86df 100644 --- a/breadarrd/src/importer/mod.rs +++ b/breadarrd/src/importer/mod.rs @@ -1325,6 +1325,7 @@ async fn wait_for_relocation( false } +#[allow(clippy::too_many_arguments)] pub async fn run_import_cycle( conn: &Connection, qbit: &QbitClient, @@ -1332,6 +1333,7 @@ pub async fn run_import_cycle( category: &str, container_downloads_path: &str, host_downloads_path: &str, + transcode_enabled: bool, ) -> Result { let pending = fetch_pending_grabs(conn)?; if pending.is_empty() { @@ -1359,6 +1361,7 @@ pub async fn run_import_cycle( &torrents, container_downloads_path, host_downloads_path, + transcode_enabled, )?; if stats.imported > 0 { @@ -1384,6 +1387,7 @@ fn process_pending_grabs( torrents: &[crate::qbit::TorrentInfo], container_downloads_path: &str, host_downloads_path: &str, + transcode_enabled: bool, ) -> Result { let mut stats = ImportStats::default(); @@ -1471,7 +1475,7 @@ fn process_pending_grabs( continue; } - match import_one(conn, grab, &content_path) { + match import_one(conn, grab, &content_path, transcode_enabled) { Ok(ImportOutcome::Imported { remuxed, quality_flagged, @@ -1522,7 +1526,12 @@ enum ImportOutcome { SkippedAlreadyHaveBetter, } -fn import_one(conn: &Connection, grab: &PendingGrab, content_path: &Path) -> Result { +fn import_one( + conn: &Connection, + grab: &PendingGrab, + content_path: &Path, + transcode_enabled: bool, +) -> Result { let source_path = locate_video_file(content_path)?; let ext = source_path .extension() @@ -1763,6 +1772,37 @@ fn import_one(conn: &Connection, grab: &PendingGrab, content_path: &Path) -> Res )?; } + // Queue this freshly-imported file for the async AV1 transcode swap + // (see `transcode::run_cycle`) — the file is available in the library + // immediately, exactly as today; the transcode happens later in the + // background. Never blocks or fails the import itself: enqueue errors + // are logged and swallowed, same treatment as probing failures above. + if transcode_enabled { + match crate::transcode::is_anime(conn, grab.media_item_id()) { + Ok(true) => {} // Anime excluded until its own encode tuning exists. + Ok(false) => { + let original_codec: Option = conn + .query_row( + "SELECT video_codec FROM media_file_probe WHERE episode_file_id = ?1", + params![episode_file_id], + |row| row.get(0), + ) + .optional()? + .flatten(); + if original_codec.as_deref() != Some("av1") { + if let Err(e) = + crate::transcode::enqueue(conn, episode_file_id, original_codec.as_deref(), size_bytes as i64) + { + tracing::warn!(episode_file_id, error = %e, "failed to enqueue transcode job"); + } + } + } + Err(e) => { + tracing::warn!(episode_file_id, error = %e, "failed to check anime status for transcode enqueue"); + } + } + } + Ok(ImportOutcome::Imported { remuxed, quality_flagged, @@ -2969,7 +3009,7 @@ mod tests { // absent — simulating torrents qBit no longer knows about. ]; - let stats = process_pending_grabs(&conn, &pending, &torrents, "", "").unwrap(); + let stats = process_pending_grabs(&conn, &pending, &torrents, "", "", false).unwrap(); assert_eq!(stats.imported, 1); assert_eq!(stats.skipped_incomplete, 1); @@ -3031,7 +3071,7 @@ mod tests { content_path: "/tmp/somewhere-mid-move".to_string(), }]; - let stats = process_pending_grabs(&conn, &pending, &torrents, "", "").unwrap(); + let stats = process_pending_grabs(&conn, &pending, &torrents, "", "", false).unwrap(); assert_eq!(stats.imported, 0); assert_eq!(stats.skipped_incomplete, 1); assert_eq!(stats.errors, 0); @@ -3088,7 +3128,7 @@ mod tests { for i in 1..MAX_IMPORT_ERRORS { let pending = fetch_pending_grabs(&conn).unwrap(); - let stats = process_pending_grabs(&conn, &pending, &torrents, "", "").unwrap(); + let stats = process_pending_grabs(&conn, &pending, &torrents, "", "", false).unwrap(); assert_eq!(stats.errors, 1, "iteration {i}"); assert_eq!(stats.failed, 0, "iteration {i}"); let status: String = conn @@ -3099,7 +3139,7 @@ mod tests { // The Nth failure crosses the threshold and gives up. let pending = fetch_pending_grabs(&conn).unwrap(); - let stats = process_pending_grabs(&conn, &pending, &torrents, "", "").unwrap(); + let stats = process_pending_grabs(&conn, &pending, &torrents, "", "", false).unwrap(); assert_eq!(stats.failed, 1); assert_eq!(stats.errors, 0); let status: String = conn @@ -3153,7 +3193,7 @@ mod tests { root_folder: dest_root.to_string_lossy().to_string(), }; - let outcome = import_one(&conn, &grab, &content).unwrap(); + let outcome = import_one(&conn, &grab, &content, false).unwrap(); let ImportOutcome::Imported { remuxed, .. } = outcome else { panic!("expected a real import, got a skip"); }; @@ -3193,6 +3233,133 @@ mod tests { std::fs::remove_dir_all(&dir).unwrap(); } + #[test] + fn import_one_enqueues_a_transcode_job_when_transcode_is_enabled() { + let conn = Connection::open_in_memory().unwrap(); + crate::db::init(&conn).unwrap(); + conn.execute( + "INSERT INTO media_item (id, kind, title, year, monitored, quality_profile_id, root_folder) + VALUES (1, 'movie', 'Some Movie', 2016, 1, 1, '/tmp')", + [], + ) + .unwrap(); + conn.execute( + "INSERT INTO source (id, name, kind, base_url) VALUES (1, 'tpb', 'scrape', 'http://x')", + [], + ) + .unwrap(); + conn.execute( + "INSERT INTO release (id, media_item_id, episode_id, raw_title, source_id, guid, status, torrent_hash, grabbed_at) + VALUES (1, 1, NULL, 'Some Movie 2016 1080p', 1, 'guid-1', 'grabbed', 'deadbeef', datetime('now'))", + [], + ) + .unwrap(); + + let dir = std::env::temp_dir().join(format!( + "breadarr-transcode-enqueue-{}", + std::process::id() + )); + std::fs::create_dir_all(&dir).unwrap(); + let dest_root = dir.join("library"); + std::fs::create_dir_all(&dest_root).unwrap(); + let content = dir.join("Some.Movie.2016.1080p.mp4"); + std::fs::write(&content, b"fake movie data").unwrap(); + + let grab = PendingGrab::Movie { + release_id: 1, + media_item_id: 1, + torrent_hash: "deadbeef".to_string(), + title: "Some Movie".to_string(), + year: Some(2016), + root_folder: dest_root.to_string_lossy().to_string(), + }; + + import_one(&conn, &grab, &content, true).unwrap(); + + let episode_file_id: i64 = conn + .query_row( + "SELECT id FROM episode_file WHERE media_item_id = 1", + [], + |r| r.get(0), + ) + .unwrap(); + let job_count: i64 = conn + .query_row( + "SELECT count(*) FROM transcode_job WHERE episode_file_id = ?1 AND status = 'pending'", + params![episode_file_id], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(job_count, 1); + + std::fs::remove_dir_all(&dir).unwrap(); + } + + #[test] + fn import_one_does_not_enqueue_a_transcode_job_for_anime() { + let conn = Connection::open_in_memory().unwrap(); + crate::db::init(&conn).unwrap(); + conn.execute( + "INSERT INTO media_item (id, kind, title, tvdb_id, monitored, quality_profile_id, root_folder) + VALUES (1, 'series', 'Some Anime', 999, 1, 1, '/tmp')", + [], + ) + .unwrap(); + conn.execute( + "INSERT INTO anime_mapping (anidb_id, tvdb_id) VALUES (1, 999)", + [], + ) + .unwrap(); + conn.execute( + "INSERT INTO episode (id, media_item_id, season_number, episode_number, monitored, has_file) + VALUES (1, 1, 1, 1, 1, 0)", + [], + ) + .unwrap(); + conn.execute( + "INSERT INTO source (id, name, kind, base_url) VALUES (1, 'nyaa', 'rss', 'http://x')", + [], + ) + .unwrap(); + conn.execute( + "INSERT INTO release (id, media_item_id, episode_id, raw_title, source_id, guid, status, torrent_hash, grabbed_at) + VALUES (1, 1, 1, 'Some Anime S01E01', 1, 'guid-1', 'grabbed', 'deadbeef', datetime('now'))", + [], + ) + .unwrap(); + + let dir = std::env::temp_dir().join(format!( + "breadarr-transcode-anime-skip-{}", + std::process::id() + )); + std::fs::create_dir_all(&dir).unwrap(); + let dest_root = dir.join("library"); + std::fs::create_dir_all(&dest_root).unwrap(); + let content = dir.join("Some.Anime.S01E01.mp4"); + std::fs::write(&content, b"fake anime data").unwrap(); + + let grab = PendingGrab::Episode { + release_id: 1, + media_item_id: 1, + episode_id: 1, + torrent_hash: "deadbeef".to_string(), + series_title: "Some Anime".to_string(), + season_number: 1, + episode_number: 1, + episode_title: None, + root_folder: dest_root.to_string_lossy().to_string(), + }; + + import_one(&conn, &grab, &content, true).unwrap(); + + let job_count: i64 = conn + .query_row("SELECT count(*) FROM transcode_job", [], |r| r.get(0)) + .unwrap(); + assert_eq!(job_count, 0); + + std::fs::remove_dir_all(&dir).unwrap(); + } + #[test] fn import_one_places_a_single_episode_grab_under_its_season_subfolder() { // Regression: a real single-episode grab (Mushoku Tensei S03E02/E03, @@ -3248,7 +3415,7 @@ mod tests { root_folder: dest_root.to_string_lossy().to_string(), }; - let outcome = import_one(&conn, &grab, &content).unwrap(); + let outcome = import_one(&conn, &grab, &content, false).unwrap(); assert!(matches!(outcome, ImportOutcome::Imported { .. })); let dest = dest_root.join("Season 03").join("Some Show - S03E02.mp4"); @@ -3312,7 +3479,7 @@ mod tests { root_folder: dest_root.to_string_lossy().to_string(), }; - let outcome = import_one(&conn, &grab, &content).unwrap(); + let outcome = import_one(&conn, &grab, &content, false).unwrap(); let ImportOutcome::Imported { quality_flagged, .. } = outcome @@ -3522,7 +3689,7 @@ mod tests { root_folder: dest_root.to_string_lossy().to_string(), }; - let outcome = import_one(&conn, &grab, &content).unwrap(); + let outcome = import_one(&conn, &grab, &content, false).unwrap(); assert!(matches!(outcome, ImportOutcome::SkippedAlreadyHaveBetter)); // The existing better file must be untouched, not overwritten. @@ -3600,7 +3767,7 @@ mod tests { root_folder: dest_root.to_string_lossy().to_string(), }; - let outcome = import_one(&conn, &grab, &content).unwrap(); + let outcome = import_one(&conn, &grab, &content, false).unwrap(); assert!(matches!(outcome, ImportOutcome::Imported { .. })); // The new file's content landed at the shared deterministic path. diff --git a/breadarrd/src/jellyfin.rs b/breadarrd/src/jellyfin.rs index eacb6a0..e70e194 100644 --- a/breadarrd/src/jellyfin.rs +++ b/breadarrd/src/jellyfin.rs @@ -38,4 +38,33 @@ impl JellyfinClient { } Ok(()) } + + /// Counts sessions Jellyfin is actively transcoding for right now (as + /// opposed to direct-play/direct-stream, which cost the GPU nothing) — + /// used to throttle the AV1 batch-transcode worker back so it doesn't + /// contend with a real viewer for the same encode/decode engines. + /// `TranscodingInfo` is only present on a session object while that + /// session is actually transcoding. + pub async fn active_transcode_sessions(&self) -> Result { + let resp = self + .client + .get(format!("{}/Sessions", self.base_url)) + .header("X-Emby-Token", &self.api_key) + .send() + .await + .context("jellyfin sessions request failed")?; + + let status = resp.status(); + if !status.is_success() { + let body = resp.text().await.unwrap_or_default(); + bail!("jellyfin sessions request failed: status={status} body={body:?}"); + } + + let sessions: Vec = + resp.json().await.context("failed to parse jellyfin sessions response")?; + Ok(sessions + .iter() + .filter(|s| !s["TranscodingInfo"].is_null()) + .count()) + } } diff --git a/breadarrd/src/main.rs b/breadarrd/src/main.rs index b730a91..bd66aeb 100644 --- a/breadarrd/src/main.rs +++ b/breadarrd/src/main.rs @@ -11,6 +11,7 @@ mod qbit; mod scheduler; mod scoring; mod sources; +mod transcode; use std::env; @@ -114,6 +115,9 @@ async fn main() -> Result<()> { Some("probe-library") => { return probe_library_cmd(&config).await; } + Some("transcode-library") => { + return transcode_library_cmd(&config).await; + } Some("verify-library") => { return verify_library_cmd(&config).await; } @@ -349,6 +353,9 @@ async fn background_loop( let mut upgrade_ticker = tokio::time::interval(std::time::Duration::from_secs( config.sources.upgrade_poll_interval_secs, )); + let mut transcode_ticker = tokio::time::interval(std::time::Duration::from_secs( + config.transcode.poll_interval_secs, + )); // Disk state doesn't change on its own — hourly is plenty to catch a // file deleted/moved by hand without adding meaningful load (one query // per tracked episode file, all local). Deliberately does *not* fire at @@ -366,6 +373,7 @@ async fn background_loop( search_ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); upgrade_ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); reconcile_ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + transcode_ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); // Cycle-level backoff on top of the search loop's own per-mirror // cooldowns: a whole cycle failing (source exhausted, or an outright @@ -407,7 +415,7 @@ async fn background_loop( _ = import_ticker.tick() => { let result = { let conn = conn.lock().await; - importer::run_import_cycle(&conn, &qbit, jellyfin.as_ref(), &config.qbit.category, &config.qbit.container_downloads_path, &config.qbit.host_downloads_path).await + importer::run_import_cycle(&conn, &qbit, jellyfin.as_ref(), &config.qbit.category, &config.qbit.container_downloads_path, &config.qbit.host_downloads_path, config.transcode.enabled).await }; if let (Ok(stats), Some(n)) = (&result, ¬ifier) { if stats.failed > 0 { @@ -536,6 +544,23 @@ async fn background_loop( } } } + _ = transcode_ticker.tick() => { + if !config.transcode.enabled { + continue; + } + let result = transcode::run_cycle(conn.clone(), config.transcode.clone(), jellyfin.as_ref()).await; + let record = match &result { + Ok(stats) => { + info!(?stats, "transcode cycle complete"); + api::CycleRecord { at: chrono::Utc::now(), ok: true, detail: format!("{stats:?}") } + } + Err(e) => { + error!(error = %e, "transcode cycle failed"); + api::CycleRecord { at: chrono::Utc::now(), ok: false, detail: e.to_string() } + } + }; + cycle_status.lock().expect("cycle_status poisoned").last_transcode = Some(record); + } _ = reconcile_ticker.tick() => { let result = { let conn = conn.lock().await; @@ -849,6 +874,7 @@ async fn debug_import_cycle(config: &Config) -> Result<()> { &config.qbit.category, &config.qbit.container_downloads_path, &config.qbit.host_downloads_path, + config.transcode.enabled, ) .await?; println!("{stats:?}"); @@ -1175,6 +1201,69 @@ async fn probe_library_cmd(config: &Config) -> Result<()> { Ok(()) } +/// One-time backfill: enqueues every existing-library file eligible for +/// AV1 transcoding (see `transcode::find_backlog_candidates` for the exact +/// eligibility rules — not already AV1, not anime, not HDR/2160p+ per the +/// first pass's scope) as a `transcode_job` row, then drives the same +/// worker loop the daemon's steady-state ticker uses +/// (`transcode::run_cycle`) until nothing is left pending. Shares that one +/// code path deliberately — there is exactly one place that actually runs +/// an encode, whether triggered by a backlog sweep or a fresh grab. +async fn transcode_library_cmd(config: &Config) -> Result<()> { + if !config.transcode.enabled { + bail!("transcode.enabled is false in config — enable it before running a backfill"); + } + if let Some(parent) = config.db_path().parent() { + std::fs::create_dir_all(parent)?; + } + let conn = Connection::open(config.db_path())?; + db::init(&conn)?; + + let candidates = transcode::find_backlog_candidates(&conn, &config.transcode)?; + println!( + "found {} backlog candidate(s), highest-bitrate first", + candidates.len() + ); + for candidate in &candidates { + transcode::enqueue( + &conn, + candidate.episode_file_id, + candidate.video_codec.as_deref(), + candidate.size_bytes, + )?; + } + + let jellyfin = if config.jellyfin.base_url.is_empty() { + None + } else { + Some(JellyfinClient::new( + config.jellyfin.base_url.clone(), + config.jellyfin.api_key.clone(), + )) + }; + + let conn = std::sync::Arc::new(tokio::sync::Mutex::new(conn)); + let mut total_succeeded = 0usize; + let mut total_failed = 0usize; + let mut total_bytes_saved: i64 = 0; + loop { + let stats = + transcode::run_cycle(conn.clone(), config.transcode.clone(), jellyfin.as_ref()).await?; + if stats.attempted == 0 { + break; + } + total_succeeded += stats.succeeded; + total_failed += stats.failed; + total_bytes_saved += stats.bytes_saved; + println!("batch: {stats:?}"); + } + println!( + "done: {total_succeeded} succeeded, {total_failed} failed, {:.1} GB saved total", + total_bytes_saved as f64 / 1_073_741_824.0 + ); + Ok(()) +} + /// Runs `importer::verify_library` — the expensive full-decode corruption /// check (`ffmpeg -xerror`, actually decoding every frame) against every /// header-probed-ok file that hasn't been decode-verified yet. Unlike diff --git a/breadarrd/src/scheduler.rs b/breadarrd/src/scheduler.rs index 2af2e01..b6bba16 100644 --- a/breadarrd/src/scheduler.rs +++ b/breadarrd/src/scheduler.rs @@ -222,7 +222,22 @@ fn movie_eligible_for_upgrade(conn: &Connection, media_item_id: i64) -> Result 0 { + return Ok(false); + } + // Same reasoning as the `upgrade_locked` check in + // `enumerate_upgrade_targets`: a locally AV1-transcoded file is a + // deliberate shrink, not something the upgrade loop should try to + // replace with the next bigger HEVC/x264 release it finds. + let upgrade_locked: i64 = conn + .query_row( + "SELECT upgrade_locked FROM episode_file + WHERE media_item_id = ?1 AND episode_id IS NULL", + params![media_item_id], + |row| row.get(0), + ) + .unwrap_or(0); + Ok(upgrade_locked == 0) } fn is_anime(conn: &Connection, tvdb_id: i64) -> Result { @@ -1170,6 +1185,7 @@ fn enumerate_upgrade_targets( (SELECT tvdb_id FROM anime_mapping WHERE tvdb_id IS NOT NULL)) AND NOT EXISTS (SELECT 1 FROM release r WHERE r.episode_id = e.id AND r.status IN ('grabbed','downloading')) + AND (ef.upgrade_locked IS NULL OR ef.upgrade_locked = 0) AND (us.last_checked_at IS NULL OR {})", UPGRADE_DUE_CLAUSE .replace("last_checked_at", "us.last_checked_at") @@ -2317,6 +2333,22 @@ mod tests { assert_eq!(targets[1].episode_id, Some(1)); } + #[test] + fn enumerate_upgrade_targets_excludes_an_upgrade_locked_episode() { + let conn = seeded_upgrade_conn_with_one_flagged_episode(); + // Episode 2 is the probe-flagged one that would otherwise sort + // first — lock it (as a completed local AV1 transcode would) and + // confirm it drops out entirely rather than just losing priority. + conn.execute( + "UPDATE episode_file SET upgrade_locked = 1 WHERE episode_id = 2", + [], + ) + .unwrap(); + let targets = enumerate_upgrade_targets(&conn, 10, 5.0).unwrap(); + assert_eq!(targets.len(), 1); + assert_eq!(targets[0].episode_id, Some(1)); + } + #[test] fn record_grab_writes_both_a_release_row_and_a_torrent_fetch_audit_row() { let conn = seeded_conn(); @@ -2586,6 +2618,20 @@ mod tests { assert!(movie_eligible_for_upgrade(&conn, 1).unwrap()); } + #[test] + fn movie_eligible_for_upgrade_is_false_once_upgrade_locked() { + let conn = seeded_movie_conn(); + conn.execute( + "INSERT INTO episode_file (media_item_id, episode_id, path, size_bytes, upgrade_locked) + VALUES (1, NULL, '/tmp/movie.mkv', 100, 1)", + [], + ) + .unwrap(); + // A locally AV1-transcoded file is a deliberate shrink, not + // something the upgrade loop should try to replace. + assert!(!movie_eligible_for_upgrade(&conn, 1).unwrap()); + } + #[test] fn movie_eligible_for_upgrade_is_false_when_unmonitored() { let conn = seeded_movie_conn(); diff --git a/breadarrd/src/transcode/mod.rs b/breadarrd/src/transcode/mod.rs new file mode 100644 index 0000000..44414ae --- /dev/null +++ b/breadarrd/src/transcode/mod.rs @@ -0,0 +1,434 @@ +use std::path::{Path, PathBuf}; +use std::process::Command; +use std::sync::Arc; + +use anyhow::{bail, Context, Result}; +use breadarr_shared::config::TranscodeConfig; +use rusqlite::{params, Connection}; +use tokio::sync::Mutex; + +use crate::importer::{self, ffprobe}; +use crate::jellyfin::JellyfinClient; + +/// Computes the target AV1 bitrate for a given resolution, scaled from the +/// configured reference (a real HEVC/H264 bitrate at `reference_height` the +/// user is already happy with) by pixel-count ratio, then discounted by +/// `av1_efficiency_factor` — AV1 reaches equivalent perceived quality to +/// HEVC/H264 at a meaningfully lower bitrate, so a like-for-like copy of the +/// reference bitrate would leave savings on the table. +pub fn target_bitrate_kbps(width: i64, height: i64, cfg: &TranscodeConfig) -> u32 { + let reference_width = cfg.reference_height as f64 * 16.0 / 9.0; + let reference_pixels = reference_width * cfg.reference_height as f64; + let pixel_ratio = ((width * height) as f64 / reference_pixels).max(0.1); + let bitrate = cfg.reference_bitrate_kbps as f64 * pixel_ratio * cfg.av1_efficiency_factor as f64; + bitrate.round() as u32 +} + +/// Runs the actual GPU encode via `av1_vaapi`, decoding through VAAPI too +/// (`-hwaccel_output_format vaapi`) so the whole pipeline stays on-GPU +/// rather than round-tripping frames through the CPU. Video-only re-encode +/// — every audio/subtitle/data stream is copied verbatim (`-c:a copy -c:s +/// copy -c:d copy`), and 10-bit sources stay 10-bit (AV1 handles this +/// natively; the VAAPI driver preserves the surface format through the +/// pipeline without any extra flags needed). +/// +/// Blocking and slow by design (a real GPU encode, potentially minutes per +/// file) — callers must run this inside `tokio::task::spawn_blocking`, never +/// directly on an async task, and never while holding the shared DB mutex. +fn run_ffmpeg_encode(input: &Path, output: &Path, bitrate_kbps: u32, vaapi_device: &str) -> Result<()> { + let maxrate = bitrate_kbps * 3 / 2; + let bufsize = bitrate_kbps * 2; + + let result = Command::new("ffmpeg") + .arg("-y") + .args(["-hwaccel", "vaapi"]) + .args(["-hwaccel_device", vaapi_device]) + .args(["-hwaccel_output_format", "vaapi"]) + .arg("-i") + .arg(input) + .args(["-map", "0"]) + .args(["-c:v", "av1_vaapi"]) + .args(["-b:v", &format!("{bitrate_kbps}k")]) + .args(["-maxrate", &format!("{maxrate}k")]) + .args(["-bufsize", &format!("{bufsize}k")]) + .args(["-c:a", "copy"]) + .args(["-c:s", "copy"]) + .args(["-c:d", "copy"]) + .arg(output) + .output() + .context("failed to run ffmpeg av1_vaapi encode")?; + + if !result.status.success() { + let _ = std::fs::remove_file(output); + bail!( + "ffmpeg av1_vaapi encode failed: {}", + String::from_utf8_lossy(&result.stderr) + ); + } + Ok(()) +} + +/// The blocking half of a transcode: encode to a temp file alongside the +/// original (same filesystem, required for the atomic rename-based swap +/// later — never a separate staging drive), then verify the result is +/// actually good *before* anything touches the original. Leaves the temp +/// file on disk on success (the caller finalizes the swap under the DB +/// lock); cleans it up itself on any failure, so the original is never at +/// risk regardless of what goes wrong here. +fn encode_and_verify( + input: PathBuf, + cfg: TranscodeConfig, + width: i64, + height: i64, + original_duration: Option, +) -> Result { + let tmp_path = input.with_extension("av1.mkv"); + let bitrate = target_bitrate_kbps(width, height, &cfg); + + run_ffmpeg_encode(&input, &tmp_path, bitrate, &cfg.vaapi_device)?; + + match ffprobe::verify_decodable(&tmp_path) { + Ok(ffprobe::DecodeCheck::Ok) => {} + Ok(ffprobe::DecodeCheck::Corrupt(detail)) => { + let _ = std::fs::remove_file(&tmp_path); + bail!("transcoded output failed decode verification: {detail}"); + } + Err(e) => { + let _ = std::fs::remove_file(&tmp_path); + return Err(e.context("failed to run decode verification on transcoded output")); + } + } + + if let Some(original_secs) = original_duration { + match ffprobe::probe(&tmp_path) { + Ok(new_probe) => { + if let Some(new_secs) = new_probe.duration_secs { + if (new_secs - original_secs).abs() > 1.0 { + let _ = std::fs::remove_file(&tmp_path); + bail!( + "duration mismatch after transcode: original={original_secs}s new={new_secs}s" + ); + } + } + } + Err(e) => { + let _ = std::fs::remove_file(&tmp_path); + return Err(e.context("failed to probe transcoded output for duration check")); + } + } + } + + Ok(tmp_path) +} + +/// Whether `media_item_id` is anime — same two membership checks +/// (`anime_mapping` for TV, `anime_tmdb_movie` for movies) already used +/// elsewhere for scoring/upgrade exclusions. Anime needs its own encode +/// tuning (thin lines, flat color, grain) not designed yet, so it's +/// excluded from both the backfill and the post-grab path for now. +pub fn is_anime(conn: &Connection, media_item_id: i64) -> Result { + let result: bool = conn.query_row( + "SELECT EXISTS( + SELECT 1 FROM media_item m + WHERE m.id = ?1 + AND ( + (m.tvdb_id IS NOT NULL AND m.tvdb_id IN + (SELECT tvdb_id FROM anime_mapping WHERE tvdb_id IS NOT NULL)) + OR + (m.tmdb_id IS NOT NULL AND m.tmdb_id IN (SELECT tmdb_id FROM anime_tmdb_movie)) + ) + )", + params![media_item_id], + |row| row.get(0), + )?; + Ok(result) +} + +/// Queues one file for transcoding. `original_codec`/`original_bytes` are +/// just recorded for the eventual report — not used for any decision. +pub fn enqueue( + conn: &Connection, + episode_file_id: i64, + original_codec: Option<&str>, + original_bytes: i64, +) -> Result<()> { + conn.execute( + "INSERT INTO transcode_job (episode_file_id, status, original_codec, original_bytes, queued_at) + VALUES (?1, 'pending', ?2, ?3, datetime('now'))", + params![episode_file_id, original_codec, original_bytes], + )?; + Ok(()) +} + +/// Every existing-library file eligible for the `transcode-library` backfill: +/// not already AV1, not anime, not HDR/2160p+ (per `cfg.exclude_hdr` / +/// `cfg.exclude_min_height` — the first pass is scoped to SDR 1080p/720p), +/// not already queued or done. Deliberately re-derives the anime exclusion +/// inline (rather than calling `is_anime` per row) so it's one query instead +/// of N+1 against a ~1400-file backlog. +pub fn find_backlog_candidates(conn: &Connection, cfg: &TranscodeConfig) -> Result> { + let mut stmt = conn.prepare( + "SELECT ef.id, ef.path, ef.size_bytes, p.video_codec, p.width, p.height + FROM episode_file ef + JOIN media_file_probe p ON p.episode_file_id = ef.id + LEFT JOIN episode e ON e.id = ef.episode_id + JOIN media_item m ON m.id = COALESCE(e.media_item_id, ef.media_item_id) + WHERE p.video_codec IS NOT NULL AND p.video_codec != 'av1' + AND (ef.upgrade_locked IS NULL OR ef.upgrade_locked = 0) + AND (?1 = 0 OR p.hdr = 0) + AND (p.height IS NULL OR p.height < ?2) + AND NOT ( + (m.tvdb_id IS NOT NULL AND m.tvdb_id IN + (SELECT tvdb_id FROM anime_mapping WHERE tvdb_id IS NOT NULL)) + OR + (m.tmdb_id IS NOT NULL AND m.tmdb_id IN (SELECT tmdb_id FROM anime_tmdb_movie)) + ) + AND ef.id NOT IN ( + SELECT episode_file_id FROM transcode_job WHERE status IN ('pending','running','done') + ) + ORDER BY (ef.size_bytes * 8.0 / NULLIF(p.duration_secs, 0)) DESC", + )?; + let rows = stmt + .query_map(params![cfg.exclude_hdr as i64, cfg.exclude_min_height], |row| { + Ok(BacklogCandidate { + episode_file_id: row.get(0)?, + path: row.get(1)?, + size_bytes: row.get(2)?, + video_codec: row.get(3)?, + }) + })? + .collect::>>()?; + Ok(rows) +} + +pub struct BacklogCandidate { + pub episode_file_id: i64, + pub path: String, + pub size_bytes: i64, + pub video_codec: Option, +} + +/// A `transcode_job` row claimed for processing this cycle, with the +/// probe/path data `encode_and_verify` needs already attached — claimed +/// under the DB lock, then the actual encode runs entirely outside it. +struct ClaimedJob { + job_id: i64, + episode_file_id: i64, + path: PathBuf, + width: i64, + height: i64, + duration_secs: Option, +} + +/// Claims up to `limit` `pending` jobs (marking them `running` so a crash +/// mid-cycle doesn't leave them silently re-claimable forever without at +/// least having been attempted once) and returns everything the encode step +/// needs. Only claims jobs whose file already has probe data — a job +/// enqueued moments after import but before `ensure_probed` has run yet +/// simply isn't claimed this tick, and picks up naturally on the next one. +async fn claim_pending_jobs(conn: &Arc>, limit: usize) -> Result> { + let conn = conn.lock().await; + // Highest current bitrate first — the worst offenders (REMUX, huge + // season packs) free the most space per file transcoded, so they're + // worth reaching before smaller, already-reasonable files. + let mut stmt = conn.prepare( + "SELECT j.id, j.episode_file_id, ef.path, p.width, p.height, p.duration_secs + FROM transcode_job j + JOIN episode_file ef ON ef.id = j.episode_file_id + JOIN media_file_probe p ON p.episode_file_id = ef.id + WHERE j.status = 'pending' AND p.width IS NOT NULL AND p.height IS NOT NULL + ORDER BY (ef.size_bytes * 8.0 / NULLIF(p.duration_secs, 0)) DESC + LIMIT ?1", + )?; + let claimed = stmt + .query_map(params![limit as i64], |row| { + Ok(ClaimedJob { + job_id: row.get(0)?, + episode_file_id: row.get(1)?, + path: PathBuf::from(row.get::<_, String>(2)?), + width: row.get(3)?, + height: row.get(4)?, + duration_secs: row.get(5)?, + }) + })? + .collect::>>()?; + + for job in &claimed { + conn.execute( + "UPDATE transcode_job SET status = 'running' WHERE id = ?1", + params![job.job_id], + )?; + } + Ok(claimed) +} + +/// Finalizes one job under the DB lock: on success, atomically swaps the +/// verified temp file over the original, updates `episode_file` (new size + +/// `upgrade_locked = 1`, the flag that keeps the upgrade cycle from ever +/// trying to replace a file breadarr itself just intentionally shrank), and +/// forces a fresh `ensure_probed` so `media_file_probe` reflects the real +/// AV1 ground truth — same shape as `remux_one_backlog_file`. On failure, +/// the original is left completely untouched; the job is marked `failed` +/// with the error recorded, no automatic retry. +async fn finalize_job(conn: &Arc>, job: ClaimedJob, encode_result: Result) -> Result { + let conn = conn.lock().await; + match encode_result { + Ok(tmp_path) => { + let original_bytes = std::fs::metadata(&job.path)?.len(); + std::fs::rename(&tmp_path, &job.path)?; + let new_bytes = std::fs::metadata(&job.path)?.len(); + + conn.execute( + "UPDATE episode_file SET size_bytes = ?1, upgrade_locked = 1 WHERE id = ?2", + params![new_bytes as i64, job.episode_file_id], + )?; + conn.execute( + "DELETE FROM media_file_probe WHERE episode_file_id = ?1", + params![job.episode_file_id], + )?; + importer::ensure_probed(&conn, job.episode_file_id, &job.path)?; + conn.execute( + "UPDATE transcode_job SET status = 'done', new_bytes = ?1, finished_at = datetime('now') WHERE id = ?2", + params![new_bytes as i64, job.job_id], + )?; + + Ok(TranscodeOutcome { original_bytes, new_bytes }) + } + Err(e) => { + conn.execute( + "UPDATE transcode_job SET status = 'failed', error = ?1, finished_at = datetime('now') WHERE id = ?2", + params![e.to_string(), job.job_id], + )?; + Err(e) + } + } +} + +pub struct TranscodeOutcome { + pub original_bytes: u64, + pub new_bytes: u64, +} + +#[derive(Debug, Default)] +pub struct TranscodeCycleStats { + pub attempted: usize, + pub succeeded: usize, + pub failed: usize, + pub bytes_saved: i64, +} + +/// Active parallelism for this cycle: full `parallelism_max` when nobody's +/// actively watching a transcoded Jellyfin stream, dropped to +/// `parallelism_min` (1) the moment anyone is — the batch job and a real +/// viewer are contending for the same GPU encode/decode engines, and a +/// stutter during someone's actual show loses every time. +async fn effective_parallelism(jellyfin: Option<&JellyfinClient>, cfg: &TranscodeConfig) -> usize { + let Some(client) = jellyfin else { + return cfg.parallelism_max; + }; + match client.active_transcode_sessions().await { + Ok(0) => cfg.parallelism_max, + Ok(_) => cfg.parallelism_min, + Err(e) => { + tracing::warn!(error = %e, "failed to poll jellyfin sessions; assuming worst case"); + cfg.parallelism_min + } + } +} + +/// One transcode-worker pass: claims up to the current (Jellyfin-aware) +/// parallelism worth of pending jobs, encodes them concurrently entirely +/// outside the DB lock (each encode can take minutes — holding the shared +/// mutex for that long would stall every other cycle: import, search, +/// review-queue actions, everything), then finalizes each result under a +/// brief lock. Shared by both the steady-state daemon ticker and the +/// `transcode-library` backfill CLI, so there's exactly one code path that +/// actually runs an encode. +pub async fn run_cycle( + conn: Arc>, + cfg: TranscodeConfig, + jellyfin: Option<&JellyfinClient>, +) -> Result { + let parallelism = effective_parallelism(jellyfin, &cfg).await; + let claimed = claim_pending_jobs(&conn, parallelism).await?; + + let mut stats = TranscodeCycleStats::default(); + if claimed.is_empty() { + return Ok(stats); + } + + let mut handles = Vec::with_capacity(claimed.len()); + for job in claimed { + let cfg = cfg.clone(); + let path = job.path.clone(); + let (width, height, duration) = (job.width, job.height, job.duration_secs); + let encode_handle = + tokio::task::spawn_blocking(move || encode_and_verify(path, cfg, width, height, duration)); + handles.push((job, encode_handle)); + } + + for (job, encode_handle) in handles { + stats.attempted += 1; + let encode_result = match encode_handle.await { + Ok(result) => result, + Err(join_err) => Err(anyhow::anyhow!("encode task panicked: {join_err}")), + }; + match finalize_job(&conn, job, encode_result).await { + Ok(outcome) => { + stats.succeeded += 1; + stats.bytes_saved += outcome.original_bytes as i64 - outcome.new_bytes as i64; + } + Err(e) => { + stats.failed += 1; + tracing::warn!(error = %e, "transcode job failed"); + } + } + } + + Ok(stats) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn cfg() -> TranscodeConfig { + TranscodeConfig { + enabled: true, + poll_interval_secs: 60, + vaapi_device: "/dev/dri/renderD128".to_string(), + parallelism_min: 1, + parallelism_max: 4, + reference_bitrate_kbps: 5320, + reference_height: 1080, + av1_efficiency_factor: 0.7, + exclude_hdr: true, + exclude_min_height: 2000, + } + } + + #[test] + fn target_bitrate_matches_reference_at_reference_resolution() { + let bitrate = target_bitrate_kbps(1920, 1080, &cfg()); + // pixel_ratio == 1.0 at the reference resolution, so this should be + // exactly reference_bitrate_kbps * av1_efficiency_factor. + assert_eq!(bitrate, (5320.0_f64 * 0.7).round() as u32); + } + + #[test] + fn target_bitrate_scales_down_for_720p() { + let bitrate_1080 = target_bitrate_kbps(1920, 1080, &cfg()); + let bitrate_720 = target_bitrate_kbps(1280, 720, &cfg()); + assert!(bitrate_720 < bitrate_1080); + // Roughly proportional to pixel count (~0.44x), not some flat cut. + let ratio = bitrate_720 as f64 / bitrate_1080 as f64; + assert!((0.4..0.5).contains(&ratio), "ratio was {ratio}"); + } + + #[test] + fn target_bitrate_scales_up_for_1440p() { + let bitrate_1080 = target_bitrate_kbps(1920, 1080, &cfg()); + let bitrate_1440 = target_bitrate_kbps(2560, 1440, &cfg()); + assert!(bitrate_1440 > bitrate_1080); + } +}