Targeted-search backoff and release-date gating (#90)
This commit was merged in pull request #90.
This commit is contained in:
+68
@@ -0,0 +1,68 @@
|
|||||||
|
{
|
||||||
|
"db_name": "SQLite",
|
||||||
|
"query": "\n SELECT m.id AS \"id!: i64\",\n m.tmdb_id AS \"tmdb_id!: i64\",\n m.title AS \"title!: String\",\n m.year,\n m.original_language,\n m.search_attempts AS \"search_attempts!: i64\",\n m.last_searched_at,\n m.digital_release,\n m.metadata_refreshed_at\n FROM movies m\n WHERE m.wanted = 1\n AND m.blocked = 0\n AND NOT EXISTS (\n SELECT 1 FROM media_files f\n WHERE f.owner_kind = 'movie' AND f.owner_id = m.id\n )\n AND NOT EXISTS (\n SELECT 1 FROM grabs g\n WHERE g.target_kind = 'movie' AND g.target_id = m.id\n AND g.state IN ('sent', 'downloaded', 'imported')\n )\n ORDER BY m.last_searched_at IS NOT NULL, m.last_searched_at, m.id\n LIMIT ?\n ",
|
||||||
|
"describe": {
|
||||||
|
"columns": [
|
||||||
|
{
|
||||||
|
"name": "id!: i64",
|
||||||
|
"ordinal": 0,
|
||||||
|
"type_info": "Integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "tmdb_id!: i64",
|
||||||
|
"ordinal": 1,
|
||||||
|
"type_info": "Integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "title!: String",
|
||||||
|
"ordinal": 2,
|
||||||
|
"type_info": "Text"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "year",
|
||||||
|
"ordinal": 3,
|
||||||
|
"type_info": "Integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "original_language",
|
||||||
|
"ordinal": 4,
|
||||||
|
"type_info": "Text"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "search_attempts!: i64",
|
||||||
|
"ordinal": 5,
|
||||||
|
"type_info": "Integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "last_searched_at",
|
||||||
|
"ordinal": 6,
|
||||||
|
"type_info": "Text"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "digital_release",
|
||||||
|
"ordinal": 7,
|
||||||
|
"type_info": "Text"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "metadata_refreshed_at",
|
||||||
|
"ordinal": 8,
|
||||||
|
"type_info": "Text"
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"parameters": {
|
||||||
|
"Right": 1
|
||||||
|
},
|
||||||
|
"nullable": [
|
||||||
|
true,
|
||||||
|
false,
|
||||||
|
false,
|
||||||
|
true,
|
||||||
|
true,
|
||||||
|
false,
|
||||||
|
true,
|
||||||
|
true,
|
||||||
|
true
|
||||||
|
]
|
||||||
|
},
|
||||||
|
"hash": "0db845eb00a34dce18b6a26efc83b99efeddf05aae9a32765680a97b5db73bb0"
|
||||||
|
}
|
||||||
-38
@@ -1,38 +0,0 @@
|
|||||||
{
|
|
||||||
"db_name": "SQLite",
|
|
||||||
"query": "\n SELECT m.id AS \"id!: i64\",\n m.title AS \"title!: String\",\n m.year,\n m.original_language\n FROM movies m\n WHERE m.wanted = 1\n AND m.blocked = 0\n AND NOT EXISTS (\n SELECT 1 FROM media_files f\n WHERE f.owner_kind = 'movie' AND f.owner_id = m.id\n )\n AND NOT EXISTS (\n SELECT 1 FROM grabs g\n WHERE g.target_kind = 'movie' AND g.target_id = m.id\n AND g.state IN ('sent', 'downloaded', 'imported')\n )\n ORDER BY m.last_searched_at IS NOT NULL, m.last_searched_at, m.id\n LIMIT ?\n ",
|
|
||||||
"describe": {
|
|
||||||
"columns": [
|
|
||||||
{
|
|
||||||
"name": "id!: i64",
|
|
||||||
"ordinal": 0,
|
|
||||||
"type_info": "Integer"
|
|
||||||
},
|
|
||||||
{
|
|
||||||
"name": "title!: String",
|
|
||||||
"ordinal": 1,
|
|
||||||
"type_info": "Text"
|
|
||||||
},
|
|
||||||
{
|
|
||||||
"name": "year",
|
|
||||||
"ordinal": 2,
|
|
||||||
"type_info": "Integer"
|
|
||||||
},
|
|
||||||
{
|
|
||||||
"name": "original_language",
|
|
||||||
"ordinal": 3,
|
|
||||||
"type_info": "Text"
|
|
||||||
}
|
|
||||||
],
|
|
||||||
"parameters": {
|
|
||||||
"Right": 1
|
|
||||||
},
|
|
||||||
"nullable": [
|
|
||||||
true,
|
|
||||||
false,
|
|
||||||
true,
|
|
||||||
true
|
|
||||||
]
|
|
||||||
},
|
|
||||||
"hash": "2c329cab6e494ef1a765a7de05e1351358d4ed583a36fe5d0edc38cc12603e44"
|
|
||||||
}
|
|
||||||
+12
@@ -0,0 +1,12 @@
|
|||||||
|
{
|
||||||
|
"db_name": "SQLite",
|
||||||
|
"query": "UPDATE movies\n SET title = ?, year = ?, original_language = ?, digital_release = ?,\n metadata_refreshed_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now'),\n search_attempts = 0, last_searched_at = NULL,\n updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')\n WHERE id = ? AND (\n title IS NOT ? OR year IS NOT ? OR original_language IS NOT ?\n OR digital_release IS NOT ?\n )",
|
||||||
|
"describe": {
|
||||||
|
"columns": [],
|
||||||
|
"parameters": {
|
||||||
|
"Right": 9
|
||||||
|
},
|
||||||
|
"nullable": []
|
||||||
|
},
|
||||||
|
"hash": "9da028029e93a01bd4be5a3b065556e60ba014a782d297aa01ba4b80dc94b485"
|
||||||
|
}
|
||||||
+12
@@ -0,0 +1,12 @@
|
|||||||
|
{
|
||||||
|
"db_name": "SQLite",
|
||||||
|
"query": "UPDATE movies SET metadata_refreshed_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')\n WHERE id = ?",
|
||||||
|
"describe": {
|
||||||
|
"columns": [],
|
||||||
|
"parameters": {
|
||||||
|
"Right": 1
|
||||||
|
},
|
||||||
|
"nullable": []
|
||||||
|
},
|
||||||
|
"hash": "de2153eefd6ec03ceeaf640d599271834cb06c94eda7a9956c5bd7399593c0ae"
|
||||||
|
}
|
||||||
+439
-13
@@ -15,6 +15,7 @@
|
|||||||
|
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::path::PathBuf;
|
use std::path::PathBuf;
|
||||||
|
use std::sync::Arc;
|
||||||
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
||||||
|
|
||||||
use arr_core::policy::{evaluate, Candidate};
|
use arr_core::policy::{evaluate, Candidate};
|
||||||
@@ -22,6 +23,7 @@ use arr_core::{score::score, Language, Policy, TitleOverrides, Verdict};
|
|||||||
use arr_db::{blacklist, Blacklist, Db, MoviePolicy};
|
use arr_db::{blacklist, Blacklist, Db, MoviePolicy};
|
||||||
use arr_dl::{AddTorrent, TorrentSource, TransmissionClient};
|
use arr_dl::{AddTorrent, TorrentSource, TransmissionClient};
|
||||||
use arr_indexer::{ProwlarrClient, SearchRelease, SearchRequest};
|
use arr_indexer::{ProwlarrClient, SearchRelease, SearchRequest};
|
||||||
|
use arr_meta::TmdbClient;
|
||||||
|
|
||||||
use crate::reconcile::{Action, ActionFuture, Outcome};
|
use crate::reconcile::{Action, ActionFuture, Outcome};
|
||||||
|
|
||||||
@@ -79,6 +81,10 @@ pub enum GrabError {
|
|||||||
Policy(#[from] arr_db::PolicyError),
|
Policy(#[from] arr_db::PolicyError),
|
||||||
#[error("prowlarr: {0}")]
|
#[error("prowlarr: {0}")]
|
||||||
Prowlarr(#[from] arr_indexer::Error),
|
Prowlarr(#[from] arr_indexer::Error),
|
||||||
|
#[error("TMDB: {0}")]
|
||||||
|
Metadata(#[from] arr_meta::Error),
|
||||||
|
#[error("movie {0} has an invalid TMDB id")]
|
||||||
|
InvalidTmdbId(i64),
|
||||||
#[error("transmission: {0}")]
|
#[error("transmission: {0}")]
|
||||||
Transmission(#[from] arr_dl::Error),
|
Transmission(#[from] arr_dl::Error),
|
||||||
#[error("release {name}: {source}")]
|
#[error("release {name}: {source}")]
|
||||||
@@ -105,6 +111,7 @@ pub struct GrabAction {
|
|||||||
transmission: TransmissionClient,
|
transmission: TransmissionClient,
|
||||||
download_dir: PathBuf,
|
download_dir: PathBuf,
|
||||||
seeding: SeedingRules,
|
seeding: SeedingRules,
|
||||||
|
tmdb: Option<Arc<TmdbClient>>,
|
||||||
indexers: tokio::sync::RwLock<IndexerCache>,
|
indexers: tokio::sync::RwLock<IndexerCache>,
|
||||||
/// [`INDEXER_DISCOVERY_TIMEOUT`], overridden by tests that cannot wait
|
/// [`INDEXER_DISCOVERY_TIMEOUT`], overridden by tests that cannot wait
|
||||||
/// out the real one. Mirrors `ReconcileLoop`'s action timeout override.
|
/// out the real one. Mirrors `ReconcileLoop`'s action timeout override.
|
||||||
@@ -124,11 +131,18 @@ impl GrabAction {
|
|||||||
transmission,
|
transmission,
|
||||||
download_dir,
|
download_dir,
|
||||||
seeding,
|
seeding,
|
||||||
|
tmdb: None,
|
||||||
indexers: tokio::sync::RwLock::new(IndexerCache::default()),
|
indexers: tokio::sync::RwLock::new(IndexerCache::default()),
|
||||||
discovery_timeout: INDEXER_DISCOVERY_TIMEOUT,
|
discovery_timeout: INDEXER_DISCOVERY_TIMEOUT,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[must_use]
|
||||||
|
pub fn with_tmdb(mut self, tmdb: Arc<TmdbClient>) -> Self {
|
||||||
|
self.tmdb = Some(tmdb);
|
||||||
|
self
|
||||||
|
}
|
||||||
|
|
||||||
async fn tick(&self, database: &Db) -> Result<Vec<Outcome>, GrabError> {
|
async fn tick(&self, database: &Db) -> Result<Vec<Outcome>, GrabError> {
|
||||||
let mut outcomes = self.track_sent_grabs(database).await?;
|
let mut outcomes = self.track_sent_grabs(database).await?;
|
||||||
let gaps = pending_movies(database).await?;
|
let gaps = pending_movies(database).await?;
|
||||||
@@ -136,13 +150,26 @@ impl GrabAction {
|
|||||||
return Ok(outcomes);
|
return Ok(outcomes);
|
||||||
}
|
}
|
||||||
|
|
||||||
let searchable = self.searchable_indexers().await?;
|
|
||||||
if searchable.is_empty() {
|
|
||||||
tracing::warn!("no indexer advertises a text search; nothing can be grabbed");
|
|
||||||
return Ok(outcomes);
|
|
||||||
}
|
|
||||||
|
|
||||||
for movie in gaps {
|
for movie in gaps {
|
||||||
|
let (movie_id, title) = (movie.id, movie.title.clone());
|
||||||
|
let movie = match self.refresh_metadata(database, movie).await {
|
||||||
|
Ok(Some(movie)) => movie,
|
||||||
|
Ok(None) => continue,
|
||||||
|
// A title's metadata is not the rest of the tick's problem,
|
||||||
|
// same as a grab failure below.
|
||||||
|
Err(error) => {
|
||||||
|
tracing::error!(movie_id, title, %error, "metadata refresh failed");
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
if !search_due(&movie) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
let searchable = self.searchable_indexers().await?;
|
||||||
|
if searchable.is_empty() {
|
||||||
|
tracing::warn!("no indexer advertises a text search; nothing can be grabbed");
|
||||||
|
return Ok(outcomes);
|
||||||
|
}
|
||||||
match self.grab_one(database, &movie, &searchable).await {
|
match self.grab_one(database, &movie, &searchable).await {
|
||||||
Ok(Some(outcome)) => outcomes.push(outcome),
|
Ok(Some(outcome)) => outcomes.push(outcome),
|
||||||
Ok(None) => {}
|
Ok(None) => {}
|
||||||
@@ -159,6 +186,93 @@ impl GrabAction {
|
|||||||
Ok(outcomes)
|
Ok(outcomes)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn refresh_metadata(
|
||||||
|
&self,
|
||||||
|
database: &Db,
|
||||||
|
movie: PendingMovie,
|
||||||
|
) -> Result<Option<PendingMovie>, GrabError> {
|
||||||
|
let Some(tmdb) = &self.tmdb else {
|
||||||
|
return Ok(Some(movie));
|
||||||
|
};
|
||||||
|
if !metadata_refresh_due(movie.metadata_refreshed_at.as_deref()) {
|
||||||
|
let released = is_digitally_released(movie.digital_release.as_deref());
|
||||||
|
return Ok(released.then_some(movie));
|
||||||
|
}
|
||||||
|
let tmdb_id =
|
||||||
|
u32::try_from(movie.tmdb_id).map_err(|_| GrabError::InvalidTmdbId(movie.id))?;
|
||||||
|
let metadata = tmdb.movie(tmdb_id).await?;
|
||||||
|
let title = metadata.title.clone();
|
||||||
|
let year = metadata.year().map(i64::from);
|
||||||
|
let original_language =
|
||||||
|
(!metadata.original_language.is_empty()).then_some(metadata.original_language.clone());
|
||||||
|
let digital_release = metadata.digital_release.map(|date| date.to_string());
|
||||||
|
let title_ref = title.as_str();
|
||||||
|
let original_language_ref = original_language.as_deref();
|
||||||
|
let digital_release_ref = digital_release.as_deref();
|
||||||
|
let changed = sqlx::query!(
|
||||||
|
r#"UPDATE movies
|
||||||
|
SET title = ?, year = ?, original_language = ?, digital_release = ?,
|
||||||
|
metadata_refreshed_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now'),
|
||||||
|
search_attempts = 0, last_searched_at = NULL,
|
||||||
|
updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
|
||||||
|
WHERE id = ? AND (
|
||||||
|
title IS NOT ? OR year IS NOT ? OR original_language IS NOT ?
|
||||||
|
OR digital_release IS NOT ?
|
||||||
|
)"#,
|
||||||
|
title_ref,
|
||||||
|
year,
|
||||||
|
original_language_ref,
|
||||||
|
digital_release_ref,
|
||||||
|
movie.id,
|
||||||
|
title_ref,
|
||||||
|
year,
|
||||||
|
original_language_ref,
|
||||||
|
digital_release_ref,
|
||||||
|
)
|
||||||
|
.execute(database.pool())
|
||||||
|
.await?
|
||||||
|
.rows_affected()
|
||||||
|
!= 0;
|
||||||
|
if changed {
|
||||||
|
tracing::info!(
|
||||||
|
movie_id = movie.id,
|
||||||
|
"metadata changed; reset targeted search backoff"
|
||||||
|
);
|
||||||
|
} else {
|
||||||
|
// Still stamp the refresh even when nothing changed, or the TTL
|
||||||
|
// gate above never engages and every tick pays for TMDB again.
|
||||||
|
sqlx::query!(
|
||||||
|
"UPDATE movies SET metadata_refreshed_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
|
||||||
|
WHERE id = ?",
|
||||||
|
movie.id
|
||||||
|
)
|
||||||
|
.execute(database.pool())
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
|
if !metadata.is_digitally_released(chrono::Utc::now().date_naive()) {
|
||||||
|
tracing::debug!(
|
||||||
|
movie_id = movie.id,
|
||||||
|
"digital release has not happened; skipping targeted search"
|
||||||
|
);
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
|
Ok(Some(PendingMovie {
|
||||||
|
id: movie.id,
|
||||||
|
tmdb_id: movie.tmdb_id,
|
||||||
|
title,
|
||||||
|
year,
|
||||||
|
original_language,
|
||||||
|
search_attempts: if changed { 0 } else { movie.search_attempts },
|
||||||
|
last_searched_at: if changed {
|
||||||
|
None
|
||||||
|
} else {
|
||||||
|
movie.last_searched_at
|
||||||
|
},
|
||||||
|
digital_release,
|
||||||
|
metadata_refreshed_at: None,
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
/// Move grabs Transmission reports as complete out of `sent`.
|
/// Move grabs Transmission reports as complete out of `sent`.
|
||||||
///
|
///
|
||||||
/// Transmission is authoritative and its view is rebuilt on every tick
|
/// Transmission is authoritative and its view is rebuilt on every tick
|
||||||
@@ -284,8 +398,6 @@ impl GrabAction {
|
|||||||
}
|
}
|
||||||
releases.extend(search.releases);
|
releases.extend(search.releases);
|
||||||
}
|
}
|
||||||
record_search(database, movie.id).await?;
|
|
||||||
|
|
||||||
let mut candidates = Vec::new();
|
let mut candidates = Vec::new();
|
||||||
for release in releases {
|
for release in releases {
|
||||||
let stored = store_release(
|
let stored = store_release(
|
||||||
@@ -354,6 +466,7 @@ impl GrabAction {
|
|||||||
let Some(winner) = candidates.into_iter().find(|candidate| {
|
let Some(winner) = candidates.into_iter().find(|candidate| {
|
||||||
!blacklist.blocks_candidate(&candidate.name, &candidate.download_url)
|
!blacklist.blocks_candidate(&candidate.name, &candidate.download_url)
|
||||||
}) else {
|
}) else {
|
||||||
|
record_search(database, movie.id).await?;
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
movie_id = movie.id,
|
movie_id = movie.id,
|
||||||
title = movie.title,
|
title = movie.title,
|
||||||
@@ -362,23 +475,49 @@ impl GrabAction {
|
|||||||
return Ok(None);
|
return Ok(None);
|
||||||
};
|
};
|
||||||
|
|
||||||
|
self.send_winner(database, movie, &loaded, &blacklist, winner)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Add the winning release to Transmission and record the grab.
|
||||||
|
///
|
||||||
|
/// The search still counts as an attempt (§6.2) on every exit that is
|
||||||
|
/// not a completed grab — a Transmission error or a blacklisted-infohash
|
||||||
|
/// drop must not leave the same release to repeat next tick with no
|
||||||
|
/// backoff.
|
||||||
|
async fn send_winner(
|
||||||
|
&self,
|
||||||
|
database: &Db,
|
||||||
|
movie: &PendingMovie,
|
||||||
|
loaded: &MoviePolicy,
|
||||||
|
blacklist: &Blacklist,
|
||||||
|
winner: Eligible,
|
||||||
|
) -> Result<Option<Outcome>, GrabError> {
|
||||||
let seeding = self.seeding.for_indexer(winner.indexer_id);
|
let seeding = self.seeding.for_indexer(winner.indexer_id);
|
||||||
let added = self
|
let added = match self
|
||||||
.transmission
|
.transmission
|
||||||
.add_torrent(AddTorrent {
|
.add_torrent(AddTorrent {
|
||||||
source: torrent_source(&winner.download_url),
|
source: torrent_source(&winner.download_url),
|
||||||
label: label(&loaded),
|
label: label(loaded),
|
||||||
download_dir: self.download_dir.clone(),
|
download_dir: self.download_dir.clone(),
|
||||||
seed_ratio_limit: seeding.ratio,
|
seed_ratio_limit: seeding.ratio,
|
||||||
seed_idle_limit_minutes: seeding.idle_minutes,
|
seed_idle_limit_minutes: seeding.idle_minutes,
|
||||||
})
|
})
|
||||||
.await?;
|
.await
|
||||||
|
{
|
||||||
|
Ok(added) => added,
|
||||||
|
Err(error) => {
|
||||||
|
record_search(database, movie.id).await?;
|
||||||
|
return Err(error.into());
|
||||||
|
}
|
||||||
|
};
|
||||||
let infohash = added.hash.to_ascii_lowercase();
|
let infohash = added.hash.to_ascii_lowercase();
|
||||||
|
|
||||||
// §6.3's second key. A `.torrent` link hides its infohash until
|
// §6.3's second key. A `.torrent` link hides its infohash until
|
||||||
// Transmission has fetched it, so the same blacklisted torrent can
|
// Transmission has fetched it, so the same blacklisted torrent can
|
||||||
// reach here under a new name.
|
// reach here under a new name.
|
||||||
if blacklist.blocks_infohash(&infohash) {
|
if blacklist.blocks_infohash(&infohash) {
|
||||||
|
record_search(database, movie.id).await?;
|
||||||
self.drop_blacklisted_torrent(database, movie, &winner, &added)
|
self.drop_blacklisted_torrent(database, movie, &winner, &added)
|
||||||
.await?;
|
.await?;
|
||||||
return Ok(None);
|
return Ok(None);
|
||||||
@@ -499,11 +638,16 @@ impl Action for GrabAction {
|
|||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
struct PendingMovie {
|
struct PendingMovie {
|
||||||
id: i64,
|
id: i64,
|
||||||
|
tmdb_id: i64,
|
||||||
title: String,
|
title: String,
|
||||||
year: Option<i64>,
|
year: Option<i64>,
|
||||||
/// §5.2's language rules are expressed against this, and guessing it is
|
/// §5.2's language rules are expressed against this, and guessing it is
|
||||||
/// worse than not grabbing.
|
/// worse than not grabbing.
|
||||||
original_language: Option<String>,
|
original_language: Option<String>,
|
||||||
|
search_attempts: i64,
|
||||||
|
last_searched_at: Option<String>,
|
||||||
|
digital_release: Option<String>,
|
||||||
|
metadata_refreshed_at: Option<String>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The eligible view of a stored release, ranked for selection.
|
/// The eligible view of a stored release, ranked for selection.
|
||||||
@@ -526,9 +670,14 @@ async fn pending_movies(database: &Db) -> Result<Vec<PendingMovie>, GrabError> {
|
|||||||
let rows = sqlx::query!(
|
let rows = sqlx::query!(
|
||||||
r#"
|
r#"
|
||||||
SELECT m.id AS "id!: i64",
|
SELECT m.id AS "id!: i64",
|
||||||
|
m.tmdb_id AS "tmdb_id!: i64",
|
||||||
m.title AS "title!: String",
|
m.title AS "title!: String",
|
||||||
m.year,
|
m.year,
|
||||||
m.original_language
|
m.original_language,
|
||||||
|
m.search_attempts AS "search_attempts!: i64",
|
||||||
|
m.last_searched_at,
|
||||||
|
m.digital_release,
|
||||||
|
m.metadata_refreshed_at
|
||||||
FROM movies m
|
FROM movies m
|
||||||
WHERE m.wanted = 1
|
WHERE m.wanted = 1
|
||||||
AND m.blocked = 0
|
AND m.blocked = 0
|
||||||
@@ -553,13 +702,54 @@ async fn pending_movies(database: &Db) -> Result<Vec<PendingMovie>, GrabError> {
|
|||||||
.into_iter()
|
.into_iter()
|
||||||
.map(|row| PendingMovie {
|
.map(|row| PendingMovie {
|
||||||
id: row.id,
|
id: row.id,
|
||||||
|
tmdb_id: row.tmdb_id,
|
||||||
title: row.title,
|
title: row.title,
|
||||||
year: row.year,
|
year: row.year,
|
||||||
original_language: row.original_language,
|
original_language: row.original_language,
|
||||||
|
search_attempts: row.search_attempts,
|
||||||
|
last_searched_at: row.last_searched_at,
|
||||||
|
digital_release: row.digital_release,
|
||||||
|
metadata_refreshed_at: row.metadata_refreshed_at,
|
||||||
})
|
})
|
||||||
.collect())
|
.collect())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn search_due(movie: &PendingMovie) -> bool {
|
||||||
|
let Some(last_searched_at) = &movie.last_searched_at else {
|
||||||
|
return true;
|
||||||
|
};
|
||||||
|
let Ok(last_searched_at) = chrono::DateTime::parse_from_rfc3339(last_searched_at) else {
|
||||||
|
return true;
|
||||||
|
};
|
||||||
|
let backoff = match movie.search_attempts {
|
||||||
|
1 => chrono::TimeDelta::hours(1),
|
||||||
|
2 => chrono::TimeDelta::hours(6),
|
||||||
|
3 => chrono::TimeDelta::days(1),
|
||||||
|
4 => chrono::TimeDelta::days(3),
|
||||||
|
_ => chrono::TimeDelta::days(7),
|
||||||
|
};
|
||||||
|
last_searched_at.with_timezone(&chrono::Utc) + backoff <= chrono::Utc::now()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn metadata_refresh_due(metadata_refreshed_at: Option<&str>) -> bool {
|
||||||
|
let Some(refreshed_at) = metadata_refreshed_at else {
|
||||||
|
return true;
|
||||||
|
};
|
||||||
|
let Ok(refreshed_at) = chrono::DateTime::parse_from_rfc3339(refreshed_at) else {
|
||||||
|
return true;
|
||||||
|
};
|
||||||
|
// Independent of the search backoff (§6.2): a title stuck on a day-long
|
||||||
|
// backoff, or one with no digital release date yet, must not cost a
|
||||||
|
// TMDB call every tick.
|
||||||
|
refreshed_at.with_timezone(&chrono::Utc) + chrono::TimeDelta::hours(6) <= chrono::Utc::now()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn is_digitally_released(digital_release: Option<&str>) -> bool {
|
||||||
|
digital_release
|
||||||
|
.and_then(|date| date.parse::<chrono::NaiveDate>().ok())
|
||||||
|
.is_some_and(|date| date <= chrono::Utc::now().date_naive())
|
||||||
|
}
|
||||||
|
|
||||||
/// Cache the classified release and associate it with the title.
|
/// Cache the classified release and associate it with the title.
|
||||||
///
|
///
|
||||||
/// Returns the candidate only when the release is eligible: automatic
|
/// Returns the candidate only when the release is eligible: automatic
|
||||||
@@ -865,6 +1055,29 @@ mod tests {
|
|||||||
</item>
|
</item>
|
||||||
</channel></rss>"#;
|
</channel></rss>"#;
|
||||||
|
|
||||||
|
const EMPTY_RSS: &str = "<rss><channel></channel></rss>";
|
||||||
|
const RELEASED_METADATA: &str = r#"{
|
||||||
|
"id": 693134,
|
||||||
|
"title": "Dune Part Two",
|
||||||
|
"original_title": "Dune: Part Two",
|
||||||
|
"original_language": "en",
|
||||||
|
"origin_country": ["US"],
|
||||||
|
"release_date": "2024-02-27",
|
||||||
|
"release_dates": {"results": [{"release_dates": [{
|
||||||
|
"type": 4, "release_date": "2024-04-16T00:00:00.000Z"
|
||||||
|
}]}]}
|
||||||
|
}"#;
|
||||||
|
|
||||||
|
const UNRELEASED_METADATA: &str = r#"{
|
||||||
|
"id": 693134,
|
||||||
|
"title": "Dune Part Two",
|
||||||
|
"original_title": "Dune: Part Two",
|
||||||
|
"original_language": "en",
|
||||||
|
"origin_country": ["US"],
|
||||||
|
"release_date": "2024-02-27",
|
||||||
|
"release_dates": {"results": []}
|
||||||
|
}"#;
|
||||||
|
|
||||||
async fn prowlarr() -> MockServer {
|
async fn prowlarr() -> MockServer {
|
||||||
let server = MockServer::start().await;
|
let server = MockServer::start().await;
|
||||||
Mock::given(method("GET"))
|
Mock::given(method("GET"))
|
||||||
@@ -892,6 +1105,44 @@ mod tests {
|
|||||||
server
|
server
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn empty_prowlarr() -> MockServer {
|
||||||
|
let server = MockServer::start().await;
|
||||||
|
Mock::given(method("GET"))
|
||||||
|
.and(path("/api/v1/indexer"))
|
||||||
|
.respond_with(
|
||||||
|
ResponseTemplate::new(200)
|
||||||
|
.set_body_json(json!([{"id": 7, "name": "tracker", "enable": true}])),
|
||||||
|
)
|
||||||
|
.mount(&server)
|
||||||
|
.await;
|
||||||
|
Mock::given(method("GET"))
|
||||||
|
.and(path("/7/api"))
|
||||||
|
.and(query_param("t", "caps"))
|
||||||
|
.respond_with(ResponseTemplate::new(200).set_body_string(
|
||||||
|
r#"<caps><searching><search available="yes" supportedParams="q"/></searching></caps>"#,
|
||||||
|
))
|
||||||
|
.mount(&server)
|
||||||
|
.await;
|
||||||
|
Mock::given(method("GET"))
|
||||||
|
.and(path("/7/api"))
|
||||||
|
.and(query_param("t", "search"))
|
||||||
|
.respond_with(ResponseTemplate::new(200).set_body_string(EMPTY_RSS))
|
||||||
|
.mount(&server)
|
||||||
|
.await;
|
||||||
|
server
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn tmdb(metadata: &str) -> MockServer {
|
||||||
|
let server = MockServer::start().await;
|
||||||
|
Mock::given(method("GET"))
|
||||||
|
.and(path("/3/movie/693134"))
|
||||||
|
.and(query_param("append_to_response", "release_dates"))
|
||||||
|
.respond_with(ResponseTemplate::new(200).set_body_string(metadata))
|
||||||
|
.mount(&server)
|
||||||
|
.await;
|
||||||
|
server
|
||||||
|
}
|
||||||
|
|
||||||
async fn transmission() -> (MockServer, FakeTransmission) {
|
async fn transmission() -> (MockServer, FakeTransmission) {
|
||||||
let server = MockServer::start().await;
|
let server = MockServer::start().await;
|
||||||
let fake = FakeTransmission::default();
|
let fake = FakeTransmission::default();
|
||||||
@@ -932,6 +1183,35 @@ mod tests {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn action_with_tmdb(
|
||||||
|
prowlarr: &MockServer,
|
||||||
|
transmission: &MockServer,
|
||||||
|
metadata: &MockServer,
|
||||||
|
) -> GrabAction {
|
||||||
|
action(prowlarr, transmission).with_tmdb(Arc::new(
|
||||||
|
TmdbClient::builder("key")
|
||||||
|
.base_url(format!("{}/3/", metadata.uri()))
|
||||||
|
.build()
|
||||||
|
.unwrap(),
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn targeted_searches(indexer: &MockServer) -> usize {
|
||||||
|
indexer
|
||||||
|
.received_requests()
|
||||||
|
.await
|
||||||
|
.unwrap()
|
||||||
|
.into_iter()
|
||||||
|
.filter(|request| {
|
||||||
|
request.url.path() == "/7/api"
|
||||||
|
&& request
|
||||||
|
.url
|
||||||
|
.query_pairs()
|
||||||
|
.any(|(name, value)| name == "t" && value == "search")
|
||||||
|
})
|
||||||
|
.count()
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn seeding_rules_select_by_prowlarr_indexer_id() {
|
fn seeding_rules_select_by_prowlarr_indexer_id() {
|
||||||
let default = SeedingLimits {
|
let default = SeedingLimits {
|
||||||
@@ -1079,7 +1359,7 @@ mod tests {
|
|||||||
.fetch_one(database.pool())
|
.fetch_one(database.pool())
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
assert_eq!(attempts, 1);
|
assert_eq!(attempts, 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Transmission is authoritative (§8): a completed torrent moves its grab
|
/// Transmission is authoritative (§8): a completed torrent moves its grab
|
||||||
@@ -1231,6 +1511,152 @@ mod tests {
|
|||||||
assert!(grabs(&database).await.is_empty());
|
assert!(grabs(&database).await.is_empty());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn unreleased_movies_never_call_an_indexer() {
|
||||||
|
let (_dir, database) = wanted_movie().await;
|
||||||
|
let indexer = prowlarr().await;
|
||||||
|
let metadata = tmdb(UNRELEASED_METADATA).await;
|
||||||
|
let (downloader, _fake) = transmission().await;
|
||||||
|
let action = action_with_tmdb(&indexer, &downloader, &metadata);
|
||||||
|
|
||||||
|
for _ in 0..3 {
|
||||||
|
action.tick(&database).await.unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
|
assert!(indexer.received_requests().await.unwrap().is_empty());
|
||||||
|
let attempts: i64 = sqlx::query_scalar("SELECT search_attempts FROM movies")
|
||||||
|
.fetch_one(database.pool())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(attempts, 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn unsuccessful_searches_follow_the_backoff_schedule() {
|
||||||
|
let (_dir, database) = wanted_movie().await;
|
||||||
|
let indexer = empty_prowlarr().await;
|
||||||
|
let metadata = tmdb(RELEASED_METADATA).await;
|
||||||
|
let (downloader, _fake) = transmission().await;
|
||||||
|
let action = action_with_tmdb(&indexer, &downloader, &metadata);
|
||||||
|
|
||||||
|
action.tick(&database).await.unwrap();
|
||||||
|
assert_eq!(targeted_searches(&indexer).await, 1);
|
||||||
|
action.tick(&database).await.unwrap();
|
||||||
|
assert_eq!(targeted_searches(&indexer).await, 1);
|
||||||
|
|
||||||
|
for (delay, expected_searches) in [("-1 hour", 2), ("-6 hours", 3), ("-1 day", 4)] {
|
||||||
|
sqlx::query(
|
||||||
|
"UPDATE movies SET last_searched_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now', ?)",
|
||||||
|
)
|
||||||
|
.bind(delay)
|
||||||
|
.execute(database.pool())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
action.tick(&database).await.unwrap();
|
||||||
|
assert_eq!(targeted_searches(&indexer).await, expected_searches);
|
||||||
|
}
|
||||||
|
|
||||||
|
sqlx::query(
|
||||||
|
"UPDATE movies SET search_attempts = 5, last_searched_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now', '-6 days')",
|
||||||
|
)
|
||||||
|
.execute(database.pool())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
action.tick(&database).await.unwrap();
|
||||||
|
assert_eq!(targeted_searches(&indexer).await, 4);
|
||||||
|
|
||||||
|
sqlx::query(
|
||||||
|
"UPDATE movies SET last_searched_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now', '-7 days')",
|
||||||
|
)
|
||||||
|
.execute(database.pool())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
action.tick(&database).await.unwrap();
|
||||||
|
assert_eq!(targeted_searches(&indexer).await, 5);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A title stuck on backoff must not pay for TMDB on every tick: the
|
||||||
|
/// metadata refresh has its own TTL, independent of the search backoff.
|
||||||
|
///
|
||||||
|
/// A fresh `GrabAction` (and so a fresh `TmdbClient`) is built for every
|
||||||
|
/// tick, as a restarted process would, so the only thing that can be
|
||||||
|
/// suppressing a real TMDB request is the persisted
|
||||||
|
/// `metadata_refreshed_at` gate rather than the client's own in-process
|
||||||
|
/// response cache.
|
||||||
|
#[tokio::test]
|
||||||
|
async fn metadata_refresh_is_throttled_by_its_own_ttl() {
|
||||||
|
let (_dir, database) = wanted_movie().await;
|
||||||
|
let indexer = empty_prowlarr().await;
|
||||||
|
let metadata = tmdb(RELEASED_METADATA).await;
|
||||||
|
let (downloader, _fake) = transmission().await;
|
||||||
|
|
||||||
|
action_with_tmdb(&indexer, &downloader, &metadata)
|
||||||
|
.tick(&database)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(metadata.received_requests().await.unwrap().len(), 1);
|
||||||
|
assert_eq!(targeted_searches(&indexer).await, 1);
|
||||||
|
|
||||||
|
// Due for another search attempt, but the metadata refresh is not
|
||||||
|
// due yet: TMDB is not called again, and the stored digital release
|
||||||
|
// still gates the search correctly.
|
||||||
|
sqlx::query(
|
||||||
|
"UPDATE movies SET last_searched_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now', '-2 hours')",
|
||||||
|
)
|
||||||
|
.execute(database.pool())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
action_with_tmdb(&indexer, &downloader, &metadata)
|
||||||
|
.tick(&database)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(metadata.received_requests().await.unwrap().len(), 1);
|
||||||
|
assert_eq!(targeted_searches(&indexer).await, 2);
|
||||||
|
|
||||||
|
// Past the refresh TTL: the next due attempt pays for TMDB again.
|
||||||
|
sqlx::query(
|
||||||
|
"UPDATE movies SET last_searched_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now', '-7 hours'),
|
||||||
|
metadata_refreshed_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now', '-7 hours')",
|
||||||
|
)
|
||||||
|
.execute(database.pool())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
action_with_tmdb(&indexer, &downloader, &metadata)
|
||||||
|
.tick(&database)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(metadata.received_requests().await.unwrap().len(), 2);
|
||||||
|
assert_eq!(targeted_searches(&indexer).await, 3);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn metadata_changes_reset_a_title_backoff() {
|
||||||
|
let (_dir, database) = wanted_movie().await;
|
||||||
|
sqlx::query(
|
||||||
|
"UPDATE movies SET search_attempts = 4, last_searched_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')",
|
||||||
|
)
|
||||||
|
.execute(database.pool())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let indexer = empty_prowlarr().await;
|
||||||
|
let metadata = tmdb(&RELEASED_METADATA.replace("Dune Part Two", "Dune: Part Two")).await;
|
||||||
|
let (downloader, _fake) = transmission().await;
|
||||||
|
|
||||||
|
action_with_tmdb(&indexer, &downloader, &metadata)
|
||||||
|
.tick(&database)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
assert_eq!(targeted_searches(&indexer).await, 1);
|
||||||
|
let (title, attempts): (String, i64) =
|
||||||
|
sqlx::query_as("SELECT title, search_attempts FROM movies")
|
||||||
|
.fetch_one(database.pool())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(title, "Dune: Part Two");
|
||||||
|
assert_eq!(attempts, 1);
|
||||||
|
}
|
||||||
|
|
||||||
/// Prowlarr probes capabilities one indexer at a time, so paying for
|
/// Prowlarr probes capabilities one indexer at a time, so paying for
|
||||||
/// discovery every 30 s would leave the tick no room to search.
|
/// discovery every 30 s would leave the tick no room to search.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
|
|||||||
@@ -94,10 +94,19 @@ async fn run() -> Result<(), Error> {
|
|||||||
database.migrate().await?;
|
database.migrate().await?;
|
||||||
|
|
||||||
let transmission = arr_dl::TransmissionClient::new(&config.transmission_url)?;
|
let transmission = arr_dl::TransmissionClient::new(&config.transmission_url)?;
|
||||||
|
let tmdb = if let Some(key) = &config.tmdb_api_key {
|
||||||
|
let mut client = TmdbClient::builder(key.clone());
|
||||||
|
if let Some(url) = &config.tmdb_url {
|
||||||
|
client = client.base_url(url.clone());
|
||||||
|
}
|
||||||
|
Some(Arc::new(client.build()?))
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
};
|
||||||
let mut reconcile = ReconcileLoop::new(database.clone());
|
let mut reconcile = ReconcileLoop::new(database.clone());
|
||||||
// Without a Prowlarr key nothing can be searched, so the grab lane stays
|
// Without a Prowlarr key nothing can be searched, so the grab lane stays
|
||||||
// unregistered rather than failing a tick every 30 seconds.
|
// unregistered rather than failing a tick every 30 seconds.
|
||||||
if let Some(key) = config.prowlarr_api_key.clone() {
|
if let (Some(key), Some(tmdb)) = (config.prowlarr_api_key.clone(), tmdb.clone()) {
|
||||||
let prowlarr = arr_indexer::ProwlarrClient::new(config.prowlarr_url.clone(), key)?;
|
let prowlarr = arr_indexer::ProwlarrClient::new(config.prowlarr_url.clone(), key)?;
|
||||||
reconcile = reconcile.register(
|
reconcile = reconcile.register(
|
||||||
Tick::Reconcile,
|
Tick::Reconcile,
|
||||||
@@ -124,10 +133,11 @@ async fn run() -> Result<(), Error> {
|
|||||||
})
|
})
|
||||||
.collect(),
|
.collect(),
|
||||||
),
|
),
|
||||||
),
|
)
|
||||||
|
.with_tmdb(tmdb),
|
||||||
);
|
);
|
||||||
} else {
|
} else {
|
||||||
tracing::warn!("no Prowlarr API key: nothing will be grabbed");
|
tracing::warn!("Prowlarr or TMDB is not configured: nothing will be grabbed");
|
||||||
}
|
}
|
||||||
// Grab before import, so a download that completes on this tick is
|
// Grab before import, so a download that completes on this tick is
|
||||||
// imported on this tick.
|
// imported on this tick.
|
||||||
@@ -140,12 +150,8 @@ async fn run() -> Result<(), Error> {
|
|||||||
// Jellyseerr's Radarr shim (DESIGN.md §9.4) reads the same database and
|
// Jellyseerr's Radarr shim (DESIGN.md §9.4) reads the same database and
|
||||||
// needs its own TMDB client for `movie/lookup`.
|
// needs its own TMDB client for `movie/lookup`.
|
||||||
let mut compat = CompatState::new(database.clone());
|
let mut compat = CompatState::new(database.clone());
|
||||||
if let Some(key) = &config.tmdb_api_key {
|
if let Some(tmdb) = tmdb {
|
||||||
let mut tmdb = TmdbClient::builder(key.clone());
|
compat = compat.with_tmdb(tmdb);
|
||||||
if let Some(tmdb_url) = &config.tmdb_url {
|
|
||||||
tmdb = tmdb.base_url(tmdb_url.clone());
|
|
||||||
}
|
|
||||||
compat = compat.with_tmdb(Arc::new(tmdb.build()?));
|
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut upstreams = Upstreams::new(config.prowlarr_url, config.transmission_url)
|
let mut upstreams = Upstreams::new(config.prowlarr_url, config.transmission_url)
|
||||||
|
|||||||
@@ -42,10 +42,10 @@ impl Action for ReaperAction {
|
|||||||
"torrent-reaper"
|
"torrent-reaper"
|
||||||
}
|
}
|
||||||
|
|
||||||
fn run<'a>(&'a self, _database: &'a Db) -> ActionFuture<'a> {
|
fn run<'a>(&'a self, database: &'a Db) -> ActionFuture<'a> {
|
||||||
Box::pin(async move {
|
Box::pin(async move {
|
||||||
let labels = sqlx::query_as::<_, (String, String)>("SELECT kind, audience FROM roots")
|
let labels = sqlx::query_as::<_, (String, String)>("SELECT kind, audience FROM roots")
|
||||||
.fetch_all(_database.pool())
|
.fetch_all(database.pool())
|
||||||
.await
|
.await
|
||||||
.map_err(|error| Box::new(error) as crate::reconcile::ActionError)?
|
.map_err(|error| Box::new(error) as crate::reconcile::ActionError)?
|
||||||
.into_iter()
|
.into_iter()
|
||||||
|
|||||||
@@ -0,0 +1,4 @@
|
|||||||
|
-- TMDB's earliest digital release date is the gate for targeted searches
|
||||||
|
-- (§6.2). Keep it with title metadata so a newly announced date resets the
|
||||||
|
-- search backoff.
|
||||||
|
ALTER TABLE movies ADD COLUMN digital_release TEXT;
|
||||||
@@ -0,0 +1,5 @@
|
|||||||
|
-- Bounds how often targeted search pays for a TMDB call per movie,
|
||||||
|
-- independent of the search backoff schedule (§6.2): a title on a
|
||||||
|
-- day-long backoff, or one with no digital release date yet, must not
|
||||||
|
-- cost a TMDB call every 30 s tick.
|
||||||
|
ALTER TABLE movies ADD COLUMN metadata_refreshed_at TEXT;
|
||||||
Reference in New Issue
Block a user