From 3a7d80c2958ba8d2143f23335f310746de374496 Mon Sep 17 00:00:00 2001 From: Miguel Palhas Date: Tue, 25 Aug 2026 02:40:09 +0100 Subject: [PATCH] feat(arr): spend the subtitle budget in the reconcile loop Providers are charged one unit per download, claimed atomically right before the call; a provider at its cap is skipped in favour of the next-ranked candidate rather than failing the whole gap. Translators are charged the source character count before translating. Either cap is a queue state, same as a provider's own 429. --- crates/arr-daemon/src/subtitles.rs | 247 ++++++++++++++++++++++++++++- 1 file changed, 242 insertions(+), 5 deletions(-) diff --git a/crates/arr-daemon/src/subtitles.rs b/crates/arr-daemon/src/subtitles.rs index 72d88be..f014e27 100644 --- a/crates/arr-daemon/src/subtitles.rs +++ b/crates/arr-daemon/src/subtitles.rs @@ -20,12 +20,13 @@ //! machine translation too — the loop stops. There is no re-search and no //! replacement; that is a manual action through the API. -use std::collections::{BTreeSet, HashSet}; +use std::collections::{BTreeMap, BTreeSet, HashSet}; use std::path::{Path, PathBuf}; use std::sync::Arc; use arr_core::subs::{rank, SubtitleTarget, SubtitleVerdict}; use arr_core::{layout, Language}; +use arr_db::subtitle_budget::{self as budget, BudgetKind}; use arr_db::subtitles::{self as db, PendingSubtitle, SubtitleOrigin, SubtitleState}; use arr_db::Db; use arr_probe::Extractor; @@ -60,13 +61,32 @@ struct Settings { wanted: Vec, providers_enabled: BTreeSet, translation_engine: Option, + /// Provider id -> daily download allowance (§15, #197). A name absent + /// from the map has no cap configured, which is unlimited (§10), not + /// zero. + provider_daily_budgets: BTreeMap, + /// Engine name -> daily character allowance (§15, #197). Same absent == + /// unlimited rule as `provider_daily_budgets`. + translator_daily_budgets: BTreeMap, +} + +impl Settings { + fn provider_allowance(&self, provider: &str) -> Option { + self.provider_daily_budgets.get(provider).copied() + } + + fn translator_allowance(&self, engine: &str) -> Option { + self.translator_daily_budgets.get(engine).copied() + } } async fn load_settings(database: &Db) -> Result { let row = sqlx::query!( - r#"SELECT wanted_languages AS "wanted_languages!: String", - providers_enabled AS "providers_enabled!: String", - translation_engine AS "translation_engine: String" + r#"SELECT wanted_languages AS "wanted_languages!: String", + providers_enabled AS "providers_enabled!: String", + translation_engine AS "translation_engine: String", + provider_daily_budgets AS "provider_daily_budgets!: String", + translator_daily_budgets AS "translator_daily_budgets!: String" FROM subtitle_settings WHERE id = 1"# ) .fetch_one(database.pool()) @@ -75,10 +95,18 @@ async fn load_settings(database: &Db) -> Result { .map_err(|error| SubtitleError::Settings(error.to_string()))?; let providers_enabled: BTreeSet = serde_json::from_str(&row.providers_enabled) .map_err(|error| SubtitleError::Settings(error.to_string()))?; + let provider_daily_budgets: BTreeMap = + serde_json::from_str(&row.provider_daily_budgets) + .map_err(|error| SubtitleError::Settings(error.to_string()))?; + let translator_daily_budgets: BTreeMap = + serde_json::from_str(&row.translator_daily_budgets) + .map_err(|error| SubtitleError::Settings(error.to_string()))?; Ok(Settings { wanted, providers_enabled, translation_engine: row.translation_engine, + provider_daily_budgets, + translator_daily_budgets, }) } @@ -456,6 +484,42 @@ impl Worker { self.in_flight.lock().await.remove(&key); } + /// Rank `candidates` and claim the winner's provider download budget + /// (§15, #197), retrying without a provider that is at its cap until one + /// is affordable or none remain. Returns the claimed winner, if any, and + /// whether a provider was skipped for budget along the way — the caller + /// needs that even on a miss, since it changes whether "nothing found" + /// means no provider has it or just that arr would not pay for it today. + async fn claim_within_budget( + &self, + database: &Db, + target: &Target, + hash: Option<&str>, + mut candidates: Vec, + languages: &[Language], + settings: &Settings, + ) -> Result<(Option, bool), SubtitleError> { + let mut budget_capped = false; + while let Some(winner) = choose(target, hash, &candidates, languages) { + let provider_name = winner.provider.to_string(); + let allowance = settings.provider_allowance(&provider_name); + if budget::try_spend( + database.pool(), + BudgetKind::Provider, + &provider_name, + 1, + allowance, + ) + .await? + { + return Ok((Some(winner), budget_capped)); + } + budget_capped = true; + candidates.retain(|candidate| candidate.provider != winner.provider); + } + Ok((None, budget_capped)) + } + async fn close( &self, database: &Db, @@ -506,10 +570,36 @@ impl Worker { } let hash = moviehash(&target.path, target.size).await; - if let Some(winner) = choose(&target, hash.as_deref(), &offered, &languages) { + let (winner, budget_capped) = self + .claim_within_budget( + database, + &target, + hash.as_deref(), + offered, + &languages, + settings, + ) + .await?; + if let Some(winner) = winner { return self.fetch(database, &target, wanted_tag, winner).await; } + if budget_capped { + let reason = "provider daily budget exhausted".to_owned(); + db::record_attempt( + database.pool(), + media_file_id, + wanted_tag, + SubtitleState::Capped, + Some(&reason), + ) + .await?; + return Ok(Closed::Recorded { + state: SubtitleState::Capped, + reason, + }); + } + // A provider error is not an empty answer (§15): the candidate may // exist where nothing could be asked, so translating now would // satisfy the language forever with a worse subtitle. Back off and @@ -713,6 +803,32 @@ impl Worker { } }; + // Translators bill per character, not per call (§15), and the text + // going out is known before any of it is sent — no need to ask the + // backend afterwards (#197). + let characters = i64::try_from( + cues.iter() + .map(|cue| cue.text.chars().count()) + .sum::(), + ) + .unwrap_or(i64::MAX); + let allowance = settings.translator_allowance(&engine); + if !budget::try_spend( + database.pool(), + BudgetKind::Translator, + &engine, + characters, + allowance, + ) + .await? + { + return record( + SubtitleState::Capped, + "translator daily budget exhausted".to_owned(), + ) + .await; + } + let broken_name = format!("translator {engine}"); let translated = match arr_subs::translate::translate(backend.as_ref(), &cues, &source_language, wanted) @@ -1262,6 +1378,28 @@ mod tests { .unwrap(); } + async fn configure_provider_budget(&self, provider: &str, allowance: i64) { + sqlx::query( + "UPDATE subtitle_settings SET provider_daily_budgets = json_set(provider_daily_budgets, ?, ?)", + ) + .bind(format!("$.{provider}")) + .bind(allowance) + .execute(self.database.pool()) + .await + .unwrap(); + } + + async fn configure_translator_budget(&self, engine: &str, allowance: i64) { + sqlx::query( + "UPDATE subtitle_settings SET translator_daily_budgets = json_set(translator_daily_budgets, ?, ?)", + ) + .bind(format!("$.{engine}")) + .bind(allowance) + .execute(self.database.pool()) + .await + .unwrap(); + } + async fn attempt(&self, language: &str) -> Option { db::attempts_for(self.database.pool(), self.media_file_id) .await @@ -1471,6 +1609,60 @@ mod tests { assert_eq!(searches.load(Ordering::SeqCst), 1); } + #[tokio::test] + async fn a_provider_at_its_daily_budget_is_a_queue_state_not_a_download() { + let fixture = Fixture::new(&no_tracks()).await; + fixture.configure(r#"["pt-PT"]"#, None).await; + fixture.configure_provider_budget("opensubtitles", 0).await; + let provider = StubProvider::new( + "opensubtitles", + vec![candidate( + "opensubtitles", + "42", + Language::PortuguesePortugal, + )], + ); + let action = action(vec![Arc::new(provider)], Vec::new()); + + action.run(&fixture.database).await.unwrap(); + + let attempt = fixture.attempt("pt-PT").await.unwrap(); + assert_eq!(attempt.state, SubtitleState::Capped); + assert!(attempt + .last_failure + .unwrap() + .contains("provider daily budget")); + assert!(fixture.files().await.is_empty(), "nothing was downloaded"); + } + + #[tokio::test] + async fn a_cheaper_candidate_is_downloaded_when_the_ranked_winner_is_over_budget() { + let fixture = Fixture::new(&no_tracks()).await; + fixture.configure(r#"["pt-PT"]"#, None).await; + fixture.configure_provider_budget("opensubtitles", 0).await; + let winner = Candidate { + hash_match: true, + ..candidate("opensubtitles", "hash", Language::PortuguesePortugal) + }; + let opensubtitles = StubProvider::new("opensubtitles", vec![winner]); + let podnapisi = StubProvider::new( + "podnapisi", + vec![candidate("podnapisi", "7", Language::PortuguesePortugal)], + ); + let action = action( + vec![Arc::new(opensubtitles), Arc::new(podnapisi)], + Vec::new(), + ); + + action.run(&fixture.database).await.unwrap(); + + let files = fixture.files().await; + assert_eq!(files.len(), 1); + assert_eq!(files[0].provider.as_deref(), Some("podnapisi")); + let attempt = fixture.attempt("pt-PT").await.unwrap(); + assert_eq!(attempt.state, SubtitleState::Satisfied); + } + #[tokio::test] async fn no_candidate_translates_immediately_from_a_fetched_subtitle() { let fixture = Fixture::new(&no_tracks()).await; @@ -1552,6 +1744,51 @@ mod tests { ); } + #[tokio::test] + async fn a_translator_at_its_daily_character_budget_is_a_queue_state() { + let fixture = Fixture::new(&no_tracks()).await; + fixture.configure(r#"["pt-PT"]"#, Some("openai")).await; + fixture.configure_translator_budget("openai", 0).await; + let source = fixture + .directory + .path() + .join("Movie (2024) - [1080p].en.srt"); + tokio::fs::write(&source, SRT).await.unwrap(); + db::record_file( + fixture.database.pool(), + &arr_db::NewSubtitleFile::fetched( + fixture.media_file_id, + "en", + "opensubtitles", + "1", + &source.to_string_lossy(), + ), + ) + .await + .unwrap(); + let action = action( + vec![Arc::new(StubProvider::new("opensubtitles", Vec::new()))], + vec![Arc::new(StubBackend)], + ); + + action.run(&fixture.database).await.unwrap(); + + let attempt = fixture.attempt("pt-PT").await.unwrap(); + assert_eq!(attempt.state, SubtitleState::Capped); + assert!(attempt + .last_failure + .unwrap() + .contains("translator daily budget")); + assert!( + fixture + .files() + .await + .iter() + .all(|file| file.origin != SubtitleOrigin::Translated), + "nothing was translated" + ); + } + #[tokio::test] async fn a_text_embedded_track_is_extracted_as_the_translation_source() { let fixture = Fixture::new(&serde_json::json!({ "sub_tracks": [