Add torrents-csv and YTS as TPB search fallbacks
apibay has been timing out and the configured 1337x mirrors are all Cloudflare 521, so general-content search was failing closed. TPB stays primary; on failure or empty results, movies try YTS then torrents-csv then 1337x, and TV tries torrents-csv then 1337x. Both new sources are JSON hash-to-magnet, same grab shape as TPB. Prepend a working 1337x mirror (1337xx.to) to the default ring.
This commit is contained in:
parent
7ab28d30a7
commit
6059d77065
9 changed files with 600 additions and 171 deletions
|
|
@ -355,6 +355,26 @@ async fn background_loop(
|
|||
) {
|
||||
error!(error = %e, "failed to register tpb source row");
|
||||
}
|
||||
if let Err(e) = conn.execute(
|
||||
"INSERT OR IGNORE INTO source (id, name, kind, base_url, poll_interval_secs, enabled)
|
||||
VALUES (4, 'torrents-csv', 'scrape', ?1, ?2, 1)",
|
||||
rusqlite::params![
|
||||
config.sources.torrents_csv_url,
|
||||
config.sources.search_poll_interval_secs
|
||||
],
|
||||
) {
|
||||
error!(error = %e, "failed to register torrents-csv source row");
|
||||
}
|
||||
if let Err(e) = conn.execute(
|
||||
"INSERT OR IGNORE INTO source (id, name, kind, base_url, poll_interval_secs, enabled)
|
||||
VALUES (5, 'yts', 'scrape', ?1, ?2, 1)",
|
||||
rusqlite::params![
|
||||
config.sources.yts_api_url,
|
||||
config.sources.search_poll_interval_secs
|
||||
],
|
||||
) {
|
||||
error!(error = %e, "failed to register yts source row");
|
||||
}
|
||||
}
|
||||
|
||||
// A transient failure here (network blip during the one-time model
|
||||
|
|
@ -377,6 +397,21 @@ async fn background_loop(
|
|||
let scrape_source =
|
||||
sources::scrape::ScrapeSource::new(config.sources.torrent_1337x_mirrors.clone());
|
||||
let tpb_source = sources::tpb::TpbSource::new(config.sources.tpb_api_url.clone());
|
||||
let torrents_csv_source =
|
||||
sources::torrents_csv::TorrentsCsvSource::new(config.sources.torrents_csv_url.clone());
|
||||
let yts_source = sources::yts::YtsSource::new(config.sources.yts_api_url.clone());
|
||||
let search_sources = scheduler::SearchSources {
|
||||
tpb: &tpb_source,
|
||||
tpb_id: 3,
|
||||
torrents_csv: &torrents_csv_source,
|
||||
torrents_csv_id: 4,
|
||||
yts: &yts_source,
|
||||
yts_id: 5,
|
||||
scrape: &scrape_source,
|
||||
scrape_id: 2,
|
||||
nyaa_search: &nyaa_source,
|
||||
nyaa_id: 1,
|
||||
};
|
||||
|
||||
let mut grab_ticker = tokio::time::interval(std::time::Duration::from_secs(
|
||||
config.sources.grab_poll_interval_secs,
|
||||
|
|
@ -501,9 +536,7 @@ async fn background_loop(
|
|||
let conn = conn.lock().await;
|
||||
scheduler::run_search_cycle(
|
||||
&conn,
|
||||
&tpb_source, 3,
|
||||
&scrape_source, 2,
|
||||
&nyaa_source, 1,
|
||||
&search_sources,
|
||||
&mut title_matcher,
|
||||
&qbit, &config.qbit.category,
|
||||
config.sources.search_budget_per_cycle,
|
||||
|
|
@ -562,9 +595,7 @@ async fn background_loop(
|
|||
let conn = conn.lock().await;
|
||||
scheduler::run_upgrade_cycle(
|
||||
&conn,
|
||||
&tpb_source, 3,
|
||||
&scrape_source, 2,
|
||||
&nyaa_source, 1,
|
||||
&search_sources,
|
||||
&mut title_matcher,
|
||||
&qbit, &config.qbit.category,
|
||||
config.sources.upgrade_budget_per_cycle,
|
||||
|
|
@ -679,9 +710,7 @@ async fn background_loop(
|
|||
Ok(targets) => scheduler::execute_search_targets(
|
||||
&conn,
|
||||
&targets,
|
||||
&tpb_source, 3,
|
||||
&scrape_source, 2,
|
||||
&nyaa_source, 1,
|
||||
&search_sources,
|
||||
&mut title_matcher,
|
||||
&qbit, &config.qbit.category,
|
||||
).await,
|
||||
|
|
@ -705,9 +734,7 @@ async fn background_loop(
|
|||
&conn,
|
||||
media_item_id,
|
||||
episode_id,
|
||||
&tpb_source, 3,
|
||||
&scrape_source, 2,
|
||||
&nyaa_source, 1,
|
||||
&search_sources,
|
||||
).await
|
||||
};
|
||||
let _ = reply.send(result);
|
||||
|
|
@ -1133,6 +1160,22 @@ async fn debug_search_show(config: &Config, title: &str) -> Result<()> {
|
|||
config.sources.search_poll_interval_secs
|
||||
],
|
||||
)?;
|
||||
conn.execute(
|
||||
"INSERT OR IGNORE INTO source (id, name, kind, base_url, poll_interval_secs, enabled)
|
||||
VALUES (4, 'torrents-csv', 'scrape', ?1, ?2, 1)",
|
||||
rusqlite::params![
|
||||
config.sources.torrents_csv_url,
|
||||
config.sources.search_poll_interval_secs
|
||||
],
|
||||
)?;
|
||||
conn.execute(
|
||||
"INSERT OR IGNORE INTO source (id, name, kind, base_url, poll_interval_secs, enabled)
|
||||
VALUES (5, 'yts', 'scrape', ?1, ?2, 1)",
|
||||
rusqlite::params![
|
||||
config.sources.yts_api_url,
|
||||
config.sources.search_poll_interval_secs
|
||||
],
|
||||
)?;
|
||||
|
||||
let qbit = QbitClient::new(config.qbit.base_url.clone())?;
|
||||
if !config.qbit.username.is_empty() {
|
||||
|
|
@ -1144,6 +1187,21 @@ async fn debug_search_show(config: &Config, title: &str) -> Result<()> {
|
|||
let scrape_source =
|
||||
sources::scrape::ScrapeSource::new(config.sources.torrent_1337x_mirrors.clone());
|
||||
let nyaa_source = sources::rss::RssSource::new(config.sources.nyaa_rss_url.clone());
|
||||
let torrents_csv_source =
|
||||
sources::torrents_csv::TorrentsCsvSource::new(config.sources.torrents_csv_url.clone());
|
||||
let yts_source = sources::yts::YtsSource::new(config.sources.yts_api_url.clone());
|
||||
let search_sources = scheduler::SearchSources {
|
||||
tpb: &tpb_source,
|
||||
tpb_id: 3,
|
||||
torrents_csv: &torrents_csv_source,
|
||||
torrents_csv_id: 4,
|
||||
yts: &yts_source,
|
||||
yts_id: 5,
|
||||
scrape: &scrape_source,
|
||||
scrape_id: 2,
|
||||
nyaa_search: &nyaa_source,
|
||||
nyaa_id: 1,
|
||||
};
|
||||
|
||||
let targets = scheduler::enumerate_search_targets_for_media_item(&conn, media_item_id)?;
|
||||
println!("{} missing episode(s)/movie for {title:?}", targets.len());
|
||||
|
|
@ -1151,12 +1209,7 @@ async fn debug_search_show(config: &Config, title: &str) -> Result<()> {
|
|||
let stats = scheduler::execute_search_targets(
|
||||
&conn,
|
||||
&targets,
|
||||
&tpb_source,
|
||||
3,
|
||||
&scrape_source,
|
||||
2,
|
||||
&nyaa_source,
|
||||
1,
|
||||
&search_sources,
|
||||
&mut title_matcher,
|
||||
&qbit,
|
||||
&config.qbit.category,
|
||||
|
|
|
|||
|
|
@ -652,7 +652,8 @@ async fn process_item(
|
|||
return Ok(ProcessOutcome::QueuedForReview);
|
||||
}
|
||||
|
||||
let release_score = scoring::score(&parsed, item.seeders.unwrap_or(0), parsed.has_hdr, &profile);
|
||||
let release_score =
|
||||
scoring::score(&parsed, item.seeders.unwrap_or(0), parsed.has_hdr, &profile);
|
||||
let existing_best = match (episode_id, season_pack_number) {
|
||||
(Some(eid), _) => best_existing_score(conn, eid)?,
|
||||
(None, Some(season)) => best_existing_season_pack_score(conn, media_item.id, season)?,
|
||||
|
|
@ -911,8 +912,8 @@ pub async fn run_grab_cycle(
|
|||
Ok(stats)
|
||||
}
|
||||
|
||||
// --- Search-driven acquisition (1337x for general TV/movies, nyaa search
|
||||
// for anime movies) ---
|
||||
// --- Search-driven acquisition (TPB + YTS/csv/1337x fallbacks for
|
||||
// general TV/movies, nyaa search for anime movies) ---
|
||||
//
|
||||
// Unlike the feed-based path above, there's no natural stream of "new"
|
||||
// items to dedup against — the recurring cost here is the *search itself*,
|
||||
|
|
@ -921,18 +922,12 @@ pub async fn run_grab_cycle(
|
|||
// single cycle; cadence backs off exponentially (6h, 12h, 24h, 48h, 96h,
|
||||
// capped at a week) the more times it's been searched without success.
|
||||
|
||||
/// Other general-content sources considered and rejected (live-tested
|
||||
/// 2026-07-12, not just assumed) before landing on TPB as primary:
|
||||
/// - **YTS** (`yts.mx`): DNS doesn't resolve at all. Every known mirror
|
||||
/// (`yts.am`, `yts.ag`, `yts.lt`, `yts.pe`) either 301s in a loop or drops
|
||||
/// the query and lands on a bare homepage. The whole mirror network looks
|
||||
/// dead, not just one domain — re-check before assuming a fix is quick.
|
||||
/// - **EZTV** (`eztv.re`): redirects to `eztvx.to`, which fails to connect
|
||||
/// outright (TLS/connection error, not a slow response). Also
|
||||
/// Cloudflare-fronted, so even if connectivity is restored it carries the
|
||||
/// same risk profile 1337x does.
|
||||
/// If revisiting either, re-verify connectivity first — this isn't a
|
||||
/// permanent architectural decision, just what was true when checked.
|
||||
/// General-content fallbacks live-tested 2026-08-16 (TPB/apibay was
|
||||
/// timing out; every configured 1337x mirror returned Cloudflare 521):
|
||||
/// - **torrents-csv** and **YTS** (`yts.lt` API — `yts.mx` still does not
|
||||
/// resolve) are JSON hash-to-magnet sources, same grab shape as TPB.
|
||||
/// - **EZTV**'s JSON API is up but IMDb-id only; name search is a
|
||||
/// Cloudflare challenge. Not wired — breadarr has TMDB/TVDB, not IMDb.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
|
||||
enum SearchRoute {
|
||||
/// Primary general-content (movies + non-anime TV) route — a JSON API,
|
||||
|
|
@ -947,6 +942,22 @@ enum SearchRoute {
|
|||
NyaaSearch,
|
||||
}
|
||||
|
||||
/// The search-driven sources plus their `source` table ids — one bundle
|
||||
/// so `execute_search_targets` / `fetch_candidates` / the cycle runners
|
||||
/// don't each grow another four arguments every time a fallback is added.
|
||||
pub struct SearchSources<'a> {
|
||||
pub tpb: &'a sources::tpb::TpbSource,
|
||||
pub tpb_id: i64,
|
||||
pub torrents_csv: &'a sources::torrents_csv::TorrentsCsvSource,
|
||||
pub torrents_csv_id: i64,
|
||||
pub yts: &'a sources::yts::YtsSource,
|
||||
pub yts_id: i64,
|
||||
pub scrape: &'a sources::scrape::ScrapeSource,
|
||||
pub scrape_id: i64,
|
||||
pub nyaa_search: &'a sources::rss::RssSource,
|
||||
pub nyaa_id: i64,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct SearchTarget {
|
||||
media_item_id: i64,
|
||||
|
|
@ -1565,35 +1576,16 @@ pub struct SearchCycleStats {
|
|||
/// single broad query can dump into the review queue.
|
||||
const MAX_RESULTS_PER_SEARCH: usize = 15;
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn run_search_cycle(
|
||||
conn: &Connection,
|
||||
tpb: &sources::tpb::TpbSource,
|
||||
tpb_source_id: i64,
|
||||
scrape: &sources::scrape::ScrapeSource,
|
||||
scrape_source_id: i64,
|
||||
nyaa_search: &sources::rss::RssSource,
|
||||
nyaa_source_id: i64,
|
||||
sources: &SearchSources<'_>,
|
||||
matcher: &mut TitleMatcher,
|
||||
qbit: &QbitClient,
|
||||
qbit_category: &str,
|
||||
budget: usize,
|
||||
) -> Result<SearchCycleStats> {
|
||||
let targets = enumerate_search_targets(conn, budget)?;
|
||||
execute_search_targets(
|
||||
conn,
|
||||
&targets,
|
||||
tpb,
|
||||
tpb_source_id,
|
||||
scrape,
|
||||
scrape_source_id,
|
||||
nyaa_search,
|
||||
nyaa_source_id,
|
||||
matcher,
|
||||
qbit,
|
||||
qbit_category,
|
||||
)
|
||||
.await
|
||||
execute_search_targets(conn, &targets, sources, matcher, qbit, qbit_category).await
|
||||
}
|
||||
|
||||
/// Upgrade-search counterpart to `run_search_cycle`: same execution engine
|
||||
|
|
@ -1601,15 +1593,9 @@ pub async fn run_search_cycle(
|
|||
/// ones. `min_gain` is threaded onto every target via
|
||||
/// `enumerate_upgrade_targets`, which is what routes `process_item` into
|
||||
/// its upgrade-eligibility path instead of the normal missing-content one.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn run_upgrade_cycle(
|
||||
conn: &Connection,
|
||||
tpb: &sources::tpb::TpbSource,
|
||||
tpb_source_id: i64,
|
||||
scrape: &sources::scrape::ScrapeSource,
|
||||
scrape_source_id: i64,
|
||||
nyaa_search: &sources::rss::RssSource,
|
||||
nyaa_source_id: i64,
|
||||
sources: &SearchSources<'_>,
|
||||
matcher: &mut TitleMatcher,
|
||||
qbit: &QbitClient,
|
||||
qbit_category: &str,
|
||||
|
|
@ -1617,20 +1603,7 @@ pub async fn run_upgrade_cycle(
|
|||
min_gain: f32,
|
||||
) -> Result<SearchCycleStats> {
|
||||
let targets = enumerate_upgrade_targets(conn, budget, min_gain)?;
|
||||
execute_search_targets(
|
||||
conn,
|
||||
&targets,
|
||||
tpb,
|
||||
tpb_source_id,
|
||||
scrape,
|
||||
scrape_source_id,
|
||||
nyaa_search,
|
||||
nyaa_source_id,
|
||||
matcher,
|
||||
qbit,
|
||||
qbit_category,
|
||||
)
|
||||
.await
|
||||
execute_search_targets(conn, &targets, sources, matcher, qbit, qbit_category).await
|
||||
}
|
||||
|
||||
/// Every currently-missing episode/movie for one specific `media_item`,
|
||||
|
|
@ -1728,16 +1701,10 @@ pub fn find_media_item_id_by_title(conn: &Connection, title: &str) -> Result<Opt
|
|||
.map_err(Into::into)
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn execute_search_targets(
|
||||
conn: &Connection,
|
||||
targets: &[SearchTarget],
|
||||
tpb: &sources::tpb::TpbSource,
|
||||
tpb_source_id: i64,
|
||||
scrape: &sources::scrape::ScrapeSource,
|
||||
scrape_source_id: i64,
|
||||
nyaa_search: &sources::rss::RssSource,
|
||||
nyaa_source_id: i64,
|
||||
sources: &SearchSources<'_>,
|
||||
matcher: &mut TitleMatcher,
|
||||
qbit: &QbitClient,
|
||||
qbit_category: &str,
|
||||
|
|
@ -1753,8 +1720,8 @@ pub async fn execute_search_targets(
|
|||
// without this, a 10-episode backlog fires 10 indistinguishable
|
||||
// requests at a single-domain API with no mirror fallback, which is
|
||||
// exactly the kind of pattern that gets a source rate-limited. Caches
|
||||
// which source actually answered too (the TPB→1337x fallback can mean
|
||||
// two different targets with the same query string were served by two
|
||||
// which source actually answered too (the TPB → YTS/csv/1337x chain
|
||||
// can mean two targets with the same query were served by two
|
||||
// different sources), so a cache hit still attributes dedup/grab
|
||||
// records to the right source id.
|
||||
let mut query_cache: std::collections::HashMap<
|
||||
|
|
@ -1763,12 +1730,6 @@ pub async fn execute_search_targets(
|
|||
> = std::collections::HashMap::new();
|
||||
|
||||
for (i, target) in targets.iter().enumerate() {
|
||||
let (primary, primary_id): (&dyn ReleaseSource, i64) = match target.route {
|
||||
SearchRoute::Tpb => (tpb, tpb_source_id),
|
||||
SearchRoute::X1337 => (scrape, scrape_source_id),
|
||||
SearchRoute::NyaaSearch => (nyaa_search, nyaa_source_id),
|
||||
};
|
||||
|
||||
let cache_key = (target.route, target.query.clone());
|
||||
let (items, source_id) =
|
||||
if let Some((cached_items, cached_source_id)) = query_cache.get(&cache_key) {
|
||||
|
|
@ -1779,34 +1740,11 @@ pub async fn execute_search_targets(
|
|||
tokio::time::sleep(std::time::Duration::from_secs(jitter_secs)).await;
|
||||
}
|
||||
|
||||
let primary_result = primary.fetch(Some(&target.query)).await;
|
||||
// TPB is the primary route for general content, but a fetch
|
||||
// *failure* there (not just "no relevant results") falls back
|
||||
// to 1337x for the same query before giving up — keeps the
|
||||
// mirror-rotation/cooldown machinery built for 1337x as real
|
||||
// resilience rather than dead code, just no longer the first
|
||||
// choice given TPB's better precision. `result_source_id`
|
||||
// tracks which source actually produced whatever we end up
|
||||
// with, since dedup (`is_seen`/`mark_seen`) and grab records
|
||||
// are keyed by source id — attributing a 1337x-sourced guid to
|
||||
// TPB's source id would silently break dedup between the two.
|
||||
let (fetch_result, result_source_id) = match primary_result {
|
||||
Err(e) if matches!(target.route, SearchRoute::Tpb) => {
|
||||
tracing::warn!(
|
||||
query = %target.query,
|
||||
error = %e,
|
||||
"TPB search failed, falling back to 1337x"
|
||||
);
|
||||
(scrape.fetch(Some(&target.query)).await, scrape_source_id)
|
||||
}
|
||||
other => (other, primary_id),
|
||||
};
|
||||
|
||||
match fetch_result {
|
||||
Ok(items) => {
|
||||
match fetch_for_target(target, sources).await {
|
||||
Ok(pair) => {
|
||||
consecutive_fetch_errors = 0;
|
||||
query_cache.insert(cache_key, (items.clone(), result_source_id));
|
||||
(items, result_source_id)
|
||||
query_cache.insert(cache_key, pair.clone());
|
||||
pair
|
||||
}
|
||||
Err(e) => {
|
||||
consecutive_fetch_errors += 1;
|
||||
|
|
@ -1937,14 +1875,96 @@ pub async fn execute_search_targets(
|
|||
Ok(stats)
|
||||
}
|
||||
|
||||
fn source_route_name(route: SearchRoute) -> &'static str {
|
||||
match route {
|
||||
SearchRoute::Tpb => "tpb",
|
||||
SearchRoute::X1337 => "1337x",
|
||||
SearchRoute::NyaaSearch => "nyaa",
|
||||
fn source_name_for_id(id: i64) -> &'static str {
|
||||
match id {
|
||||
1 => "nyaa",
|
||||
2 => "1337x",
|
||||
3 => "tpb",
|
||||
4 => "torrents-csv",
|
||||
5 => "yts",
|
||||
_ => "unknown",
|
||||
}
|
||||
}
|
||||
|
||||
/// TPB first; on failure *or* empty results, walk the supplement chain
|
||||
/// (torrents-csv for everything, YTS for movies, 1337x last). Empty is
|
||||
/// treated as "try the next one" so a live-but-empty apibay doesn't hide
|
||||
/// a title that YTS/csv actually has. Dedup/grabs use whichever source
|
||||
/// actually answered.
|
||||
async fn fetch_for_target(
|
||||
target: &SearchTarget,
|
||||
sources: &SearchSources<'_>,
|
||||
) -> Result<(Vec<RawReleaseItem>, i64)> {
|
||||
match target.route {
|
||||
SearchRoute::NyaaSearch => Ok((
|
||||
sources.nyaa_search.fetch(Some(&target.query)).await?,
|
||||
sources.nyaa_id,
|
||||
)),
|
||||
SearchRoute::X1337 => Ok((
|
||||
sources.scrape.fetch(Some(&target.query)).await?,
|
||||
sources.scrape_id,
|
||||
)),
|
||||
SearchRoute::Tpb => fetch_general_content(target, sources).await,
|
||||
}
|
||||
}
|
||||
|
||||
async fn fetch_general_content(
|
||||
target: &SearchTarget,
|
||||
sources: &SearchSources<'_>,
|
||||
) -> Result<(Vec<RawReleaseItem>, i64)> {
|
||||
let is_movie = target.episode_id.is_none();
|
||||
let mut attempts: Vec<(&dyn ReleaseSource, i64, &'static str)> =
|
||||
vec![(sources.tpb, sources.tpb_id, "tpb")];
|
||||
if is_movie {
|
||||
attempts.push((sources.yts, sources.yts_id, "yts"));
|
||||
}
|
||||
attempts.push((
|
||||
sources.torrents_csv,
|
||||
sources.torrents_csv_id,
|
||||
"torrents-csv",
|
||||
));
|
||||
attempts.push((sources.scrape, sources.scrape_id, "1337x"));
|
||||
|
||||
let mut last_err: Option<anyhow::Error> = None;
|
||||
let mut any_ok = false;
|
||||
for (src, id, name) in attempts {
|
||||
match src.fetch(Some(&target.query)).await {
|
||||
Ok(items) if !items.is_empty() => {
|
||||
if name != "tpb" {
|
||||
tracing::info!(
|
||||
query = %target.query,
|
||||
source = name,
|
||||
n = items.len(),
|
||||
"search fallback produced results"
|
||||
);
|
||||
}
|
||||
return Ok((items, id));
|
||||
}
|
||||
Ok(_) => {
|
||||
any_ok = true;
|
||||
tracing::debug!(
|
||||
query = %target.query,
|
||||
source = name,
|
||||
"search source returned no results"
|
||||
);
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
query = %target.query,
|
||||
source = name,
|
||||
error = %e,
|
||||
"search source failed"
|
||||
);
|
||||
last_err = Some(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
if any_ok {
|
||||
return Ok((Vec::new(), sources.tpb_id));
|
||||
}
|
||||
Err(last_err.unwrap_or_else(|| anyhow::anyhow!("all general-content sources failed")))
|
||||
}
|
||||
|
||||
/// Fetches and scores (or gate-rejects) candidates for one search target —
|
||||
/// the same evaluation `execute_search_targets` does automatically, minus
|
||||
/// the grab decision, surfaced instead for a human to choose from. Used by
|
||||
|
|
@ -1954,29 +1974,18 @@ fn source_route_name(route: SearchRoute) -> &'static str {
|
|||
/// episode/movie right now (already owned, unmonitored, or mid-grab) —
|
||||
/// same "nothing to do" cases `enumerate_search_targets_for_media_item`
|
||||
/// already excludes.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn fetch_candidates(
|
||||
conn: &Connection,
|
||||
media_item_id: i64,
|
||||
episode_id: Option<i64>,
|
||||
tpb: &sources::tpb::TpbSource,
|
||||
tpb_source_id: i64,
|
||||
scrape: &sources::scrape::ScrapeSource,
|
||||
scrape_source_id: i64,
|
||||
nyaa_search: &sources::rss::RssSource,
|
||||
nyaa_source_id: i64,
|
||||
sources: &SearchSources<'_>,
|
||||
) -> Result<Vec<breadarr_shared::dto::ReleaseCandidate>> {
|
||||
let targets = enumerate_search_targets_for_media_item(conn, media_item_id)?;
|
||||
let Some(target) = targets.into_iter().find(|t| t.episode_id == episode_id) else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
|
||||
let (source, source_id): (&dyn ReleaseSource, i64) = match target.route {
|
||||
SearchRoute::Tpb => (tpb, tpb_source_id),
|
||||
SearchRoute::X1337 => (scrape, scrape_source_id),
|
||||
SearchRoute::NyaaSearch => (nyaa_search, nyaa_source_id),
|
||||
};
|
||||
let items = source.fetch(Some(&target.query)).await?;
|
||||
let (items, source_id) = fetch_for_target(&target, sources).await?;
|
||||
|
||||
let media_item = get_media_item(conn, media_item_id)?;
|
||||
let anime = match media_item.tvdb_id {
|
||||
|
|
@ -2029,7 +2038,7 @@ pub async fn fetch_candidates(
|
|||
link: item.link.clone(),
|
||||
guid: item.guid.clone(),
|
||||
source_id,
|
||||
source_name: source_route_name(target.route).to_string(),
|
||||
source_name: source_name_for_id(source_id).to_string(),
|
||||
seeders: item.seeders,
|
||||
leechers: item.leechers,
|
||||
size_bytes: item.size_bytes,
|
||||
|
|
|
|||
|
|
@ -1,6 +1,8 @@
|
|||
pub mod rss;
|
||||
pub mod scrape;
|
||||
pub mod torrents_csv;
|
||||
pub mod tpb;
|
||||
pub mod yts;
|
||||
|
||||
use anyhow::Result;
|
||||
use async_trait::async_trait;
|
||||
|
|
@ -25,6 +27,39 @@ pub trait ReleaseSource {
|
|||
async fn fetch(&self, query: Option<&str>) -> Result<Vec<RawReleaseItem>>;
|
||||
}
|
||||
|
||||
/// Trackers attached to every magnet we synthesize from an info-hash
|
||||
/// (TPB, torrents-csv, YTS). Same set the TPB client has used since it
|
||||
/// landed — qBittorrent needs *some* announce list or the torrent sits
|
||||
/// hash-only until DHT finds peers.
|
||||
const MAGNET_TRACKERS: &[&str] = &[
|
||||
"udp://tracker.opentrackr.org:1337/announce",
|
||||
"udp://open.stealth.si:80/announce",
|
||||
"udp://tracker.torrent.eu.org:451/announce",
|
||||
"udp://tracker.openbittorrent.com:6969/announce",
|
||||
"udp://exodus.desync.com:6969/announce",
|
||||
];
|
||||
|
||||
/// A valid BitTorrent v1 info_hash: 40 hex chars or 32 base32 chars — same
|
||||
/// shape `qbit::extract_btih` accepts out of a magnet URI. Shared by every
|
||||
/// hash-to-magnet source so a malformed value can't silently produce a
|
||||
/// magnet the grab path then fails to parse back.
|
||||
pub(crate) fn is_valid_info_hash(hash: &str) -> bool {
|
||||
(hash.len() == 40 && hash.bytes().all(|b| b.is_ascii_hexdigit()))
|
||||
|| (hash.len() == 32
|
||||
&& hash
|
||||
.bytes()
|
||||
.all(|b| matches!(b, b'2'..=b'7' | b'a'..=b'z' | b'A'..=b'Z')))
|
||||
}
|
||||
|
||||
pub(crate) fn build_magnet(info_hash: &str, name: &str) -> String {
|
||||
let mut magnet = format!("magnet:?xt=urn:btih:{info_hash}&dn={}", urlencode(name));
|
||||
for t in MAGNET_TRACKERS {
|
||||
magnet.push_str("&tr=");
|
||||
magnet.push_str(&urlencode(t));
|
||||
}
|
||||
magnet
|
||||
}
|
||||
|
||||
pub(crate) fn urlencode(s: &str) -> String {
|
||||
s.chars()
|
||||
.map(|c| {
|
||||
|
|
|
|||
130
breadarrd/src/sources/torrents_csv.rs
Normal file
130
breadarrd/src/sources/torrents_csv.rs
Normal file
|
|
@ -0,0 +1,130 @@
|
|||
use anyhow::{Context, Result};
|
||||
use async_trait::async_trait;
|
||||
use serde::Deserialize;
|
||||
|
||||
use super::{build_magnet, is_valid_info_hash, urlencode, RawReleaseItem, ReleaseSource};
|
||||
|
||||
/// Public JSON search over the torrents.csv DHT dump. Same grab shape as
|
||||
/// TPB (info-hash → magnet, no HTML): used as the first general-content
|
||||
/// fallback when apibay is down or returns nothing. Seeders are scrape
|
||||
/// snapshots, not live tracker data, so a high number can still stall —
|
||||
/// the existing seeder gate still applies.
|
||||
pub struct TorrentsCsvSource {
|
||||
api_url: String,
|
||||
client: reqwest::Client,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct CsvResponse {
|
||||
#[serde(default)]
|
||||
torrents: Vec<CsvTorrent>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct CsvTorrent {
|
||||
infohash: String,
|
||||
name: String,
|
||||
size_bytes: Option<u64>,
|
||||
seeders: Option<i64>,
|
||||
leechers: Option<i64>,
|
||||
}
|
||||
|
||||
impl TorrentsCsvSource {
|
||||
pub fn new(api_url: impl Into<String>) -> Self {
|
||||
Self {
|
||||
api_url: api_url.into(),
|
||||
client: reqwest::Client::builder()
|
||||
.timeout(std::time::Duration::from_secs(30))
|
||||
.build()
|
||||
.expect("reqwest client build"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn as_u32(n: Option<i64>) -> Option<u32> {
|
||||
n.and_then(|v| u32::try_from(v.max(0)).ok())
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl ReleaseSource for TorrentsCsvSource {
|
||||
async fn fetch(&self, query: Option<&str>) -> Result<Vec<RawReleaseItem>> {
|
||||
let Some(query) = query else {
|
||||
anyhow::bail!(
|
||||
"TorrentsCsvSource requires a search query (this is a search-driven source, not a feed)"
|
||||
);
|
||||
};
|
||||
let sep = if self.api_url.contains('?') { '&' } else { '?' };
|
||||
let url = format!(
|
||||
"{}{sep}q={}&size=25",
|
||||
self.api_url.trim_end_matches('/'),
|
||||
urlencode(query)
|
||||
);
|
||||
let parsed: CsvResponse = self
|
||||
.client
|
||||
.get(&url)
|
||||
.send()
|
||||
.await
|
||||
.with_context(|| format!("request to {url} failed"))?
|
||||
.error_for_status()
|
||||
.with_context(|| format!("{url} returned an error status"))?
|
||||
.json()
|
||||
.await
|
||||
.context("failed to parse torrents-csv response as JSON")?;
|
||||
|
||||
Ok(parsed
|
||||
.torrents
|
||||
.into_iter()
|
||||
.filter(|t| is_valid_info_hash(&t.infohash))
|
||||
.map(|t| RawReleaseItem {
|
||||
title: t.name.clone(),
|
||||
link: build_magnet(&t.infohash, &t.name),
|
||||
guid: t.infohash,
|
||||
size_bytes: t.size_bytes,
|
||||
seeders: as_u32(t.seeders),
|
||||
leechers: as_u32(t.leechers),
|
||||
})
|
||||
.collect())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn parses_a_real_captured_response() {
|
||||
let body = r#"{
|
||||
"torrents": [
|
||||
{
|
||||
"infohash": "ed0da850c273e3e15a819bdcbbf418bc85107ec8",
|
||||
"name": "Dune (2021) [1080p] [WEBRip]",
|
||||
"size_bytes": 2947989023,
|
||||
"seeders": 746,
|
||||
"leechers": 38
|
||||
},
|
||||
{
|
||||
"infohash": "not-a-hash",
|
||||
"name": "garbage",
|
||||
"size_bytes": 1,
|
||||
"seeders": 0,
|
||||
"leechers": 0
|
||||
}
|
||||
]
|
||||
}"#;
|
||||
let parsed: CsvResponse = serde_json::from_str(body).unwrap();
|
||||
let items: Vec<_> = parsed
|
||||
.torrents
|
||||
.into_iter()
|
||||
.filter(|t| is_valid_info_hash(&t.infohash))
|
||||
.collect();
|
||||
assert_eq!(items.len(), 1);
|
||||
assert_eq!(items[0].name, "Dune (2021) [1080p] [WEBRip]");
|
||||
assert_eq!(items[0].seeders, Some(746));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn empty_payload_is_no_results_not_an_error() {
|
||||
let parsed: CsvResponse = serde_json::from_str(r#"{"torrents":[]}"#).unwrap();
|
||||
assert!(parsed.torrents.is_empty());
|
||||
}
|
||||
}
|
||||
|
|
@ -2,18 +2,7 @@ use anyhow::{Context, Result};
|
|||
use async_trait::async_trait;
|
||||
use serde::Deserialize;
|
||||
|
||||
use super::{urlencode, RawReleaseItem, ReleaseSource};
|
||||
|
||||
/// A valid BitTorrent v1 info_hash: 40 hex chars or 32 base32 chars — same
|
||||
/// shape `qbit::extract_btih` accepts out of a magnet URI. Checked before
|
||||
/// building a magnet from `info_hash` at all: apibay is normally reliable,
|
||||
/// but a malformed value would otherwise silently produce a magnet
|
||||
/// `extract_btih` can't parse back out, downgrading that grab to the slow
|
||||
/// ~30s `torrents/info` polling path with no visible error anywhere.
|
||||
fn is_valid_info_hash(hash: &str) -> bool {
|
||||
(hash.len() == 40 && hash.bytes().all(|b| b.is_ascii_hexdigit()))
|
||||
|| (hash.len() == 32 && hash.bytes().all(|b| matches!(b, b'2'..=b'7' | b'a'..=b'z' | b'A'..=b'Z')))
|
||||
}
|
||||
use super::{build_magnet, is_valid_info_hash, urlencode, RawReleaseItem, ReleaseSource};
|
||||
|
||||
/// A community-run JSON API mirror of The Pirate Bay's search — unlike
|
||||
/// 1337x, this is a genuine machine-readable API (not HTML scraping), and
|
||||
|
|
@ -38,23 +27,6 @@ struct TpbResult {
|
|||
size: String,
|
||||
}
|
||||
|
||||
const TRACKERS: &[&str] = &[
|
||||
"udp://tracker.opentrackr.org:1337/announce",
|
||||
"udp://open.stealth.si:80/announce",
|
||||
"udp://tracker.torrent.eu.org:451/announce",
|
||||
"udp://tracker.openbittorrent.com:6969/announce",
|
||||
"udp://exodus.desync.com:6969/announce",
|
||||
];
|
||||
|
||||
fn build_magnet(info_hash: &str, name: &str) -> String {
|
||||
let mut magnet = format!("magnet:?xt=urn:btih:{info_hash}&dn={}", urlencode(name));
|
||||
for t in TRACKERS {
|
||||
magnet.push_str("&tr=");
|
||||
magnet.push_str(&urlencode(t));
|
||||
}
|
||||
magnet
|
||||
}
|
||||
|
||||
impl TpbSource {
|
||||
pub fn new(api_url: impl Into<String>) -> Self {
|
||||
Self {
|
||||
|
|
@ -154,7 +126,9 @@ mod tests {
|
|||
|
||||
#[test]
|
||||
fn is_valid_info_hash_accepts_both_real_shapes() {
|
||||
assert!(is_valid_info_hash("8F87C7C186172F17E35F4512BB1A3E93B614ADED")); // 40 hex
|
||||
assert!(is_valid_info_hash(
|
||||
"8F87C7C186172F17E35F4512BB1A3E93B614ADED"
|
||||
)); // 40 hex
|
||||
assert!(is_valid_info_hash("abcdefghijklmnopqrstuvwxyz234567")); // 32 base32
|
||||
}
|
||||
|
||||
|
|
@ -167,7 +141,11 @@ mod tests {
|
|||
fn is_valid_info_hash_rejects_malformed_values() {
|
||||
assert!(!is_valid_info_hash(""));
|
||||
assert!(!is_valid_info_hash("too-short"));
|
||||
assert!(!is_valid_info_hash("not-a-hex-string-at-all-nope!!!!!!!!!!!!")); // 40 chars, non-hex
|
||||
assert!(!is_valid_info_hash("8F87C7C186172F17E35F4512BB1A3E93B614ADE")); // 39 hex chars
|
||||
assert!(!is_valid_info_hash(
|
||||
"not-a-hex-string-at-all-nope!!!!!!!!!!!!"
|
||||
)); // 40 chars, non-hex
|
||||
assert!(!is_valid_info_hash(
|
||||
"8F87C7C186172F17E35F4512BB1A3E93B614ADE"
|
||||
)); // 39 hex chars
|
||||
}
|
||||
}
|
||||
|
|
|
|||
194
breadarrd/src/sources/yts.rs
Normal file
194
breadarrd/src/sources/yts.rs
Normal file
|
|
@ -0,0 +1,194 @@
|
|||
use anyhow::{Context, Result};
|
||||
use async_trait::async_trait;
|
||||
use serde::Deserialize;
|
||||
|
||||
use super::{build_magnet, is_valid_info_hash, urlencode, RawReleaseItem, ReleaseSource};
|
||||
|
||||
/// YTS movie API. `yts.mx` itself no longer resolves (checked 2026-07-12
|
||||
/// and again 2026-08-16); the `yts.lt` / `yts.am` hosts still serve the
|
||||
/// v2 JSON API, which is why the default URL is a working mirror rather
|
||||
/// than the brand domain. Movies only — each hit expands into one
|
||||
/// `RawReleaseItem` per quality so the scorer sees 720p/1080p/2160p as
|
||||
/// distinct candidates, same as if they were separate TPB rows.
|
||||
pub struct YtsSource {
|
||||
api_url: String,
|
||||
client: reqwest::Client,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct YtsResponse {
|
||||
status: String,
|
||||
data: Option<YtsData>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Default)]
|
||||
struct YtsData {
|
||||
#[serde(default)]
|
||||
movies: Vec<YtsMovie>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct YtsMovie {
|
||||
title: String,
|
||||
year: Option<i64>,
|
||||
#[serde(default)]
|
||||
torrents: Vec<YtsTorrent>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct YtsTorrent {
|
||||
hash: String,
|
||||
quality: Option<String>,
|
||||
#[serde(rename = "type")]
|
||||
source_type: Option<String>,
|
||||
video_codec: Option<String>,
|
||||
seeds: Option<i64>,
|
||||
peers: Option<i64>,
|
||||
size_bytes: Option<u64>,
|
||||
}
|
||||
|
||||
impl YtsSource {
|
||||
pub fn new(api_url: impl Into<String>) -> Self {
|
||||
Self {
|
||||
api_url: api_url.into(),
|
||||
client: reqwest::Client::builder()
|
||||
.timeout(std::time::Duration::from_secs(30))
|
||||
.build()
|
||||
.expect("reqwest client build"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn as_u32(n: Option<i64>) -> Option<u32> {
|
||||
n.and_then(|v| u32::try_from(v.max(0)).ok())
|
||||
}
|
||||
|
||||
/// Builds a release title the existing parser can read quality/source/codec
|
||||
/// out of — YTS stores those as structured fields, not in `title`.
|
||||
fn release_title(movie: &YtsMovie, torrent: &YtsTorrent) -> String {
|
||||
let mut title = movie.title.clone();
|
||||
if let Some(year) = movie.year {
|
||||
title.push_str(&format!(" ({year})"));
|
||||
}
|
||||
for part in [
|
||||
torrent.quality.as_deref(),
|
||||
torrent.source_type.as_deref(),
|
||||
torrent.video_codec.as_deref(),
|
||||
]
|
||||
.into_iter()
|
||||
.flatten()
|
||||
{
|
||||
if !part.is_empty() {
|
||||
title.push_str(&format!(" [{part}]"));
|
||||
}
|
||||
}
|
||||
title
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl ReleaseSource for YtsSource {
|
||||
async fn fetch(&self, query: Option<&str>) -> Result<Vec<RawReleaseItem>> {
|
||||
let Some(query) = query else {
|
||||
anyhow::bail!(
|
||||
"YtsSource requires a search query (this is a search-driven source, not a feed)"
|
||||
);
|
||||
};
|
||||
let sep = if self.api_url.contains('?') { '&' } else { '?' };
|
||||
let url = format!(
|
||||
"{}{sep}query_term={}&limit=20&sort_by=seeds",
|
||||
self.api_url.trim_end_matches('/'),
|
||||
urlencode(query)
|
||||
);
|
||||
let parsed: YtsResponse = self
|
||||
.client
|
||||
.get(&url)
|
||||
.send()
|
||||
.await
|
||||
.with_context(|| format!("request to {url} failed"))?
|
||||
.error_for_status()
|
||||
.with_context(|| format!("{url} returned an error status"))?
|
||||
.json()
|
||||
.await
|
||||
.context("failed to parse YTS response as JSON")?;
|
||||
anyhow::ensure!(
|
||||
parsed.status == "ok",
|
||||
"YTS returned status {:?}",
|
||||
parsed.status
|
||||
);
|
||||
|
||||
let movies = parsed.data.unwrap_or_default().movies;
|
||||
let mut items = Vec::new();
|
||||
for movie in movies {
|
||||
for torrent in &movie.torrents {
|
||||
if !is_valid_info_hash(&torrent.hash) {
|
||||
continue;
|
||||
}
|
||||
let title = release_title(&movie, torrent);
|
||||
items.push(RawReleaseItem {
|
||||
title: title.clone(),
|
||||
link: build_magnet(&torrent.hash, &title),
|
||||
guid: torrent.hash.clone(),
|
||||
size_bytes: torrent.size_bytes,
|
||||
seeders: as_u32(torrent.seeds),
|
||||
leechers: as_u32(torrent.peers),
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok(items)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
const SAMPLE: &str = r#"{
|
||||
"status": "ok",
|
||||
"data": {
|
||||
"movie_count": 1,
|
||||
"movies": [{
|
||||
"title": "Dune: Part One",
|
||||
"year": 2021,
|
||||
"torrents": [
|
||||
{
|
||||
"hash": "DEB6929BEEB09ADCBD14DC4D6081F7E6B297B88C",
|
||||
"quality": "1080p",
|
||||
"type": "web",
|
||||
"video_codec": "x264",
|
||||
"seeds": 12,
|
||||
"peers": 3,
|
||||
"size_bytes": 2147483648
|
||||
},
|
||||
{
|
||||
"hash": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
|
||||
"quality": "720p",
|
||||
"type": "bluray",
|
||||
"video_codec": "x265",
|
||||
"seeds": 4,
|
||||
"peers": 1,
|
||||
"size_bytes": 1073741824
|
||||
}
|
||||
]
|
||||
}]
|
||||
}
|
||||
}"#;
|
||||
|
||||
#[test]
|
||||
fn expands_one_movie_into_per_quality_rows() {
|
||||
let parsed: YtsResponse = serde_json::from_str(SAMPLE).unwrap();
|
||||
assert_eq!(parsed.status, "ok");
|
||||
let movie = &parsed.data.unwrap().movies[0];
|
||||
assert_eq!(movie.torrents.len(), 2);
|
||||
let t0 = release_title(movie, &movie.torrents[0]);
|
||||
assert_eq!(t0, "Dune: Part One (2021) [1080p] [web] [x264]");
|
||||
let t1 = release_title(movie, &movie.torrents[1]);
|
||||
assert_eq!(t1, "Dune: Part One (2021) [720p] [bluray] [x265]");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_movies_array_is_empty_not_an_error() {
|
||||
let parsed: YtsResponse =
|
||||
serde_json::from_str(r#"{"status":"ok","data":{"movie_count":0}}"#).unwrap();
|
||||
assert!(parsed.data.unwrap_or_default().movies.is_empty());
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue