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.
This commit is contained in:
Miguel Palhas
2026-08-25 02:40:09 +01:00
parent 4d84625332
commit 3a7d80c295
+240 -3
View File
@@ -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<String>,
providers_enabled: BTreeSet<String>,
translation_engine: Option<String>,
/// 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<String, i64>,
/// Engine name -> daily character allowance (§15, #197). Same absent ==
/// unlimited rule as `provider_daily_budgets`.
translator_daily_budgets: BTreeMap<String, i64>,
}
impl Settings {
fn provider_allowance(&self, provider: &str) -> Option<i64> {
self.provider_daily_budgets.get(provider).copied()
}
fn translator_allowance(&self, engine: &str) -> Option<i64> {
self.translator_daily_budgets.get(engine).copied()
}
}
async fn load_settings(database: &Db) -> Result<Settings, SubtitleError> {
let row = sqlx::query!(
r#"SELECT wanted_languages AS "wanted_languages!: String",
providers_enabled AS "providers_enabled!: String",
translation_engine AS "translation_engine: 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<Settings, SubtitleError> {
.map_err(|error| SubtitleError::Settings(error.to_string()))?;
let providers_enabled: BTreeSet<String> = serde_json::from_str(&row.providers_enabled)
.map_err(|error| SubtitleError::Settings(error.to_string()))?;
let provider_daily_budgets: BTreeMap<String, i64> =
serde_json::from_str(&row.provider_daily_budgets)
.map_err(|error| SubtitleError::Settings(error.to_string()))?;
let translator_daily_budgets: BTreeMap<String, i64> =
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<Candidate>,
languages: &[Language],
settings: &Settings,
) -> Result<(Option<Candidate>, 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::<usize>(),
)
.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::SubtitleAttempt> {
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": [