diff --git a/Cargo.lock b/Cargo.lock index 4591012..bca7f7d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -111,6 +111,7 @@ dependencies = [ name = "arr-indexer" version = "0.1.0" dependencies = [ + "chrono", "quick-xml", "reqwest", "serde", diff --git a/crates/arr-indexer/Cargo.toml b/crates/arr-indexer/Cargo.toml index db37235..7503622 100644 --- a/crates/arr-indexer/Cargo.toml +++ b/crates/arr-indexer/Cargo.toml @@ -7,13 +7,14 @@ repository.workspace = true publish = false [dependencies] +chrono.workspace = true quick-xml.workspace = true reqwest.workspace = true serde.workspace = true thiserror.workspace = true +tokio.workspace = true [dev-dependencies] -tokio.workspace = true wiremock.workspace = true [lints] diff --git a/crates/arr-indexer/src/lib.rs b/crates/arr-indexer/src/lib.rs index a09133e..3c6614d 100644 --- a/crates/arr-indexer/src/lib.rs +++ b/crates/arr-indexer/src/lib.rs @@ -1,5 +1,9 @@ //! arr-indexer — see DESIGN.md. +mod search; + +pub use search::{IndexerSearch, SearchError, SearchRelease, SearchRequest}; + use std::{collections::BTreeSet, time::Duration}; use quick_xml::events::Event; diff --git a/crates/arr-indexer/src/search.rs b/crates/arr-indexer/src/search.rs new file mode 100644 index 0000000..e8e6d6a --- /dev/null +++ b/crates/arr-indexer/src/search.rs @@ -0,0 +1,623 @@ +use std::{borrow::Cow, time::SystemTime}; + +use chrono::DateTime; +use quick_xml::{ + events::{BytesStart, Event}, + name::LocalName, + Reader, +}; +use reqwest::StatusCode; +use thiserror::Error; +use tokio::task::JoinSet; + +use crate::ProwlarrClient; + +/// One operation against an indexer's Torznab endpoint. +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum SearchRequest { + /// The indexer's RSS feed: a Torznab text search without a query. + Rss, + Text { + query: String, + }, + Movie { + imdb_id: String, + }, + Tv { + tvdb_id: u64, + season: Option, + episode: Option, + }, +} + +/// Release metadata returned directly by one Torznab indexer. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct SearchRelease { + pub indexer_id: i64, + pub guid: String, + pub name: String, + pub size: Option, + pub seeders: Option, + pub publish_date: Option, + pub download_url: String, +} + +/// Results from one indexer in a multi-indexer search. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct IndexerSearch { + pub indexer_id: i64, + pub releases: Vec, + pub error: Option, +} + +/// A failure isolated to one indexer's search response. +#[derive(Clone, Debug, Eq, Error, PartialEq)] +pub enum SearchError { + #[error("search request failed (status: {status:?})")] + Request { status: Option }, + #[error("search response was not valid Torznab XML")] + InvalidResponse, + #[error("tracker returned Torznab error {code:?}: {description}")] + Torznab { + code: Option, + description: String, + }, + #[error("search task did not complete")] + Task, +} + +impl ProwlarrClient { + /// Searches each indexer independently, preserving failures beside successes. + pub async fn search_indexers( + &self, + indexer_ids: &[i64], + request: &SearchRequest, + ) -> Vec { + let mut tasks = JoinSet::new(); + for (position, &indexer_id) in indexer_ids.iter().enumerate() { + let client = self.clone(); + let request = request.clone(); + tasks.spawn(async move { + ( + position, + indexer_id, + client.search_indexer(indexer_id, &request).await, + ) + }); + } + + let mut searches = vec![None; indexer_ids.len()]; + while let Some(result) = tasks.join_next().await { + if let Ok((position, indexer_id, result)) = result { + searches[position] = Some(search_result(indexer_id, result)); + } + } + + searches + .into_iter() + .enumerate() + .map(|(position, search)| { + search.unwrap_or_else(|| IndexerSearch { + indexer_id: indexer_ids[position], + releases: Vec::new(), + error: Some(SearchError::Task), + }) + }) + .collect() + } + + /// Runs one search against one indexer's Torznab endpoint. + /// + /// # Errors + /// + /// Returns an error for an HTTP failure, invalid XML, or Torznab error response. + pub async fn search_indexer( + &self, + indexer_id: i64, + request: &SearchRequest, + ) -> Result, SearchError> { + let mut parameters = vec![ + ("apikey".to_owned(), self.api_key.clone()), + ("t".to_owned(), request.operation().to_owned()), + ]; + request.add_parameters(&mut parameters); + + let response = self + .client + .get(format!("{}/{indexer_id}/api", self.base_url)) + .query(¶meters) + .send() + .await + .map_err(|error| SearchError::Request { + status: error.status(), + })? + .error_for_status() + .map_err(|error| SearchError::Request { + status: error.status(), + })?; + let body = response + .bytes() + .await + .map_err(|error| SearchError::Request { + status: error.status(), + })?; + + parse_releases(indexer_id, &body) + } +} + +impl SearchRequest { + fn operation(&self) -> &'static str { + match self { + Self::Rss | Self::Text { .. } => "search", + Self::Movie { .. } => "movie", + Self::Tv { .. } => "tvsearch", + } + } + + fn add_parameters(&self, parameters: &mut Vec<(String, String)>) { + match self { + Self::Rss => {} + Self::Text { query } => parameters.push(("q".to_owned(), query.clone())), + Self::Movie { imdb_id } => { + parameters.push(("imdbid".to_owned(), imdb_id.clone())); + } + Self::Tv { + tvdb_id, + season, + episode, + } => { + parameters.push(("tvdbid".to_owned(), tvdb_id.to_string())); + if let Some(season) = season { + parameters.push(("season".to_owned(), season.to_string())); + } + if let Some(episode) = episode { + parameters.push(("ep".to_owned(), episode.to_string())); + } + } + } + } +} + +#[derive(Default)] +struct ReleaseBuilder { + title: Option, + guid: Option, + size: Option, + seeders: Option, + publish_date: Option, + link: Option, + enclosure_url: Option, +} + +impl ReleaseBuilder { + fn finish(self, indexer_id: i64) -> Option { + let name = cleaned(self.title)?; + let download_url = cleaned(self.enclosure_url.or(self.link))?; + let guid = cleaned(self.guid).unwrap_or_else(|| download_url.clone()); + + Some(SearchRelease { + indexer_id, + guid, + name, + size: self.size, + seeders: self.seeders, + publish_date: self.publish_date, + download_url, + }) + } + + fn set_text(&mut self, field: Field, value: &str) { + match field { + Field::Title => append_text(&mut self.title, value), + Field::Guid => append_text(&mut self.guid, value), + Field::Link => append_text(&mut self.link, value), + Field::PublishDate => { + self.publish_date = parse_date(value.trim()).or(self.publish_date); + } + Field::Size => self.size = value.trim().parse().ok().or(self.size), + Field::Seeders => self.seeders = value.trim().parse().ok().or(self.seeders), + } + } +} + +#[derive(Clone, Copy)] +enum Field { + Title, + Guid, + Link, + PublishDate, + Size, + Seeders, +} + +fn parse_releases(indexer_id: i64, body: &[u8]) -> Result, SearchError> { + let mut reader = Reader::from_reader(body); + let mut root = None; + let mut item = None; + let mut field = None; + let mut releases = Vec::new(); + let mut depth = 0_u32; + + loop { + match reader + .read_event() + .map_err(|_| SearchError::InvalidResponse)? + { + Event::Start(element) => { + depth = depth.checked_add(1).ok_or(SearchError::InvalidResponse)?; + let name = element.local_name(); + if root.is_none() { + root = Some(name.as_ref().to_vec()); + if name.as_ref() == b"error" { + return Err(parse_torznab_error(&reader, &element)); + } + } + if name.as_ref() == b"item" { + item = Some(ReleaseBuilder::default()); + field = None; + } else if let Some(builder) = item.as_mut() { + read_element_attributes(&reader, builder, &element); + field = field_for(name); + } + } + Event::Empty(element) => { + let name = element.local_name(); + if root.is_none() { + root = Some(name.as_ref().to_vec()); + if name.as_ref() == b"error" { + return Err(parse_torznab_error(&reader, &element)); + } + } + if let Some(builder) = item.as_mut() { + read_element_attributes(&reader, builder, &element); + } + } + Event::Text(text) => { + if let (Some(builder), Some(field)) = (item.as_mut(), field) { + let value = text.unescape().map_err(|_| SearchError::InvalidResponse)?; + builder.set_text(field, &value); + } + } + Event::CData(text) => { + if let (Some(builder), Some(field)) = (item.as_mut(), field) { + let value = text.decode().map_err(|_| SearchError::InvalidResponse)?; + builder.set_text(field, &value); + } + } + Event::End(element) => { + depth = depth.checked_sub(1).ok_or(SearchError::InvalidResponse)?; + if element.local_name().as_ref() == b"item" { + if let Some(release) = + item.take().and_then(|builder| builder.finish(indexer_id)) + { + releases.push(release); + } + } + field = None; + } + Event::Eof if depth == 0 => break, + Event::Eof => return Err(SearchError::InvalidResponse), + _ => {} + } + } + + match root.as_deref() { + Some(b"rss" | b"feed") => Ok(releases), + _ => Err(SearchError::InvalidResponse), + } +} + +fn read_element_attributes( + reader: &Reader<&[u8]>, + builder: &mut ReleaseBuilder, + element: &BytesStart<'_>, +) { + match element.local_name().as_ref() { + b"attr" => { + let name = attribute(reader, element, b"name"); + let value = attribute(reader, element, b"value"); + match name.as_deref() { + Some("size") => { + builder.size = value.and_then(|value| value.parse().ok()).or(builder.size); + } + Some("seeders" | "seed") => { + builder.seeders = value + .and_then(|value| value.parse().ok()) + .or(builder.seeders); + } + _ => {} + } + } + b"enclosure" => { + if let Some(url) = attribute(reader, element, b"url") { + builder.enclosure_url = Some(url); + } + if builder.size.is_none() { + builder.size = + attribute(reader, element, b"length").and_then(|value| value.parse().ok()); + } + } + _ => {} + } +} + +fn attribute(reader: &Reader<&[u8]>, element: &BytesStart<'_>, name: &[u8]) -> Option { + for attribute in element.attributes() { + let Ok(attribute) = attribute else { + continue; + }; + if attribute.key.local_name().as_ref() == name { + return attribute + .decode_and_unescape_value(reader.decoder()) + .ok() + .map(Cow::into_owned); + } + } + None +} + +fn field_for(name: LocalName<'_>) -> Option { + match name.as_ref() { + b"title" => Some(Field::Title), + b"guid" => Some(Field::Guid), + b"link" => Some(Field::Link), + b"pubDate" | b"published" => Some(Field::PublishDate), + b"size" => Some(Field::Size), + b"seeders" | b"seed" => Some(Field::Seeders), + _ => None, + } +} + +fn parse_torznab_error(reader: &Reader<&[u8]>, element: &BytesStart<'_>) -> SearchError { + let code = attribute(reader, element, b"code").and_then(|code| code.parse().ok()); + let description = attribute(reader, element, b"description").unwrap_or_default(); + SearchError::Torznab { code, description } +} + +fn append_text(target: &mut Option, value: &str) { + if value.is_empty() { + return; + } + target.get_or_insert_with(String::new).push_str(value); +} + +fn cleaned(value: Option) -> Option { + value + .map(|value| value.trim().to_owned()) + .filter(|value| !value.is_empty()) +} + +fn search_result( + indexer_id: i64, + result: Result, SearchError>, +) -> IndexerSearch { + let (releases, error) = match result { + Ok(releases) => (releases, None), + Err(error) => (Vec::new(), Some(error)), + }; + IndexerSearch { + indexer_id, + releases, + error, + } +} + +fn parse_date(value: &str) -> Option { + DateTime::parse_from_rfc2822(value) + .or_else(|_| DateTime::parse_from_rfc3339(value)) + .ok() + .map(Into::into) +} + +#[cfg(test)] +mod tests { + use super::{parse_releases, SearchError, SearchRequest}; + use crate::ProwlarrClient; + use wiremock::{ + matchers::{method, path, query_param, query_param_is_missing}, + Mock, MockServer, ResponseTemplate, + }; + + const API_KEY: &str = "test-api-key"; + const RESPONSE: &str = include_str!("../tests/fixtures/beyond-hd-search.xml"); + + #[tokio::test] + async fn sends_each_torznab_request_shape() { + let server = MockServer::start().await; + mount_request(&server, 1, &[('t', "search")], Some("q")).await; + mount_request(&server, 2, &[('t', "search"), ('q', "dune part two")], None).await; + mount_request(&server, 3, &[('t', "movie"), ('i', "tt15239678")], None).await; + mount_request( + &server, + 4, + &[('t', "tvsearch"), ('v', "371980"), ('s', "2"), ('e', "3")], + None, + ) + .await; + + let client = ProwlarrClient::new(server.uri(), API_KEY).expect("client is valid"); + client + .search_indexer(1, &SearchRequest::Rss) + .await + .expect("RSS response is valid"); + client + .search_indexer( + 2, + &SearchRequest::Text { + query: "dune part two".to_owned(), + }, + ) + .await + .expect("text response is valid"); + client + .search_indexer( + 3, + &SearchRequest::Movie { + imdb_id: "tt15239678".to_owned(), + }, + ) + .await + .expect("movie response is valid"); + client + .search_indexer( + 4, + &SearchRequest::Tv { + tvdb_id: 371_980, + season: Some(2), + episode: Some(3), + }, + ) + .await + .expect("TV response is valid"); + } + + #[tokio::test] + async fn parses_recorded_variants_and_isolates_broken_xml() { + let server = MockServer::start().await; + mount_response( + &server, + 3, + include_str!("../tests/fixtures/beyond-hd-search.xml"), + ) + .await; + mount_response( + &server, + 17, + include_str!("../tests/fixtures/torrentleech-search.xml"), + ) + .await; + mount_response( + &server, + 23, + include_str!("../tests/fixtures/iptorrents-search.xml"), + ) + .await; + mount_response( + &server, + 29, + include_str!("../tests/fixtures/broken-search.xml"), + ) + .await; + + let searches = ProwlarrClient::new(server.uri(), API_KEY) + .expect("client is valid") + .search_indexers( + &[3, 17, 23, 29], + &SearchRequest::Text { + query: "example".to_owned(), + }, + ) + .await; + + assert_eq!(searches.len(), 4); + assert_eq!(searches[0].releases[0].indexer_id, 3); + assert_eq!( + searches[0].releases[0].name, + "Example.Movie.2026.2160p.WEB-DL.DDP5.1.H.265-GROUP" + ); + assert_eq!(searches[0].releases[0].guid, "bhd-redacted-1"); + assert_eq!( + searches[0].releases[0].download_url, + "https://indexer.invalid/download/bhd-redacted-1" + ); + assert_eq!(searches[0].releases[0].size, Some(23_622_320_128)); + assert_eq!(searches[0].releases[0].seeders, Some(41)); + assert!(searches[0].releases[0].publish_date.is_some()); + + assert_eq!(searches[1].releases.len(), 2); + assert_eq!(searches[1].releases[0].size, Some(4_294_967_296)); + assert_eq!(searches[1].releases[0].seeders, Some(7)); + assert_eq!( + searches[1].releases[0].download_url, + "https://indexer.invalid/download/tl-redacted-1?token=redacted" + ); + assert_eq!(searches[1].releases[1].size, None); + assert_eq!(searches[1].releases[1].publish_date, None); + + assert_eq!(searches[2].releases[0].size, Some(9_126_805_504)); + assert_eq!(searches[2].releases[0].seeders, Some(18)); + assert_eq!( + searches[2].releases[0].guid, + "https://indexer.invalid/details/ipt-redacted-1" + ); + assert_eq!( + searches[2].releases[0].download_url, + "https://indexer.invalid/download/ipt-redacted-1" + ); + assert_eq!(searches[2].error, None); + + assert!(searches[3].releases.is_empty()); + assert_eq!(searches[3].error, Some(SearchError::InvalidResponse)); + } + + #[test] + fn preserves_split_text_and_skips_undownloadable_items() { + let body = br#" + + + + Movie <![CDATA[& More]]> 2026 + split-title + https://indexer.invalid/details/split-title + + + + No download + no-download + + + + "#; + + let releases = parse_releases(3, body).expect("feed structure is valid"); + + assert_eq!(releases.len(), 1); + assert_eq!(releases[0].name, "Movie & More 2026"); + assert_eq!( + releases[0].download_url, + "https://indexer.invalid/download/split-title" + ); + } + + async fn mount_request( + server: &MockServer, + id: i64, + parameters: &[(char, &str)], + missing: Option<&str>, + ) { + let mut mock = Mock::given(method("GET")) + .and(path(format!("/{id}/api"))) + .and(query_param("apikey", API_KEY)); + for &(parameter, value) in parameters { + let name = match parameter { + 't' => "t", + 'q' => "q", + 'i' => "imdbid", + 'v' => "tvdbid", + 's' => "season", + 'e' => "ep", + _ => unreachable!("test parameter is known"), + }; + mock = mock.and(query_param(name, value)); + } + if let Some(name) = missing { + mock = mock.and(query_param_is_missing(name)); + } + mock.respond_with(ResponseTemplate::new(200).set_body_raw(RESPONSE, "application/xml")) + .expect(1) + .mount(server) + .await; + } + + async fn mount_response(server: &MockServer, id: i64, body: &str) { + Mock::given(method("GET")) + .and(path(format!("/{id}/api"))) + .and(query_param("apikey", API_KEY)) + .and(query_param("t", "search")) + .and(query_param("q", "example")) + .respond_with(ResponseTemplate::new(200).set_body_raw(body, "application/xml")) + .mount(server) + .await; + } +} diff --git a/crates/arr-indexer/tests/fixtures/beyond-hd-search.xml b/crates/arr-indexer/tests/fixtures/beyond-hd-search.xml new file mode 100644 index 0000000..5ea5050 --- /dev/null +++ b/crates/arr-indexer/tests/fixtures/beyond-hd-search.xml @@ -0,0 +1,15 @@ + + + + Beyond-HD + + <![CDATA[Example.Movie.2026.2160p.WEB-DL.DDP5.1.H.265-GROUP]]> + bhd-redacted-1 + https://indexer.invalid/download/bhd-redacted-1 + Sat, 22 Aug 2026 12:30:00 +0000 + + + + + + diff --git a/crates/arr-indexer/tests/fixtures/broken-search.xml b/crates/arr-indexer/tests/fixtures/broken-search.xml new file mode 100644 index 0000000..566a5d8 --- /dev/null +++ b/crates/arr-indexer/tests/fixtures/broken-search.xml @@ -0,0 +1,6 @@ + + + + + Broken tracker response + broken-1 diff --git a/crates/arr-indexer/tests/fixtures/iptorrents-search.xml b/crates/arr-indexer/tests/fixtures/iptorrents-search.xml new file mode 100644 index 0000000..72dded7 --- /dev/null +++ b/crates/arr-indexer/tests/fixtures/iptorrents-search.xml @@ -0,0 +1,14 @@ + + + + IPTorrents + + Example.Movie.2024.1080p.BluRay.x264-GROUP + https://indexer.invalid/details/ipt-redacted-1 + https://indexer.invalid/download/ipt-redacted-1 + Fri, 21 Aug 2026 20:05:10 GMT + 9126805504 + 18 + + + diff --git a/crates/arr-indexer/tests/fixtures/torrentleech-search.xml b/crates/arr-indexer/tests/fixtures/torrentleech-search.xml new file mode 100644 index 0000000..6e6c4b1 --- /dev/null +++ b/crates/arr-indexer/tests/fixtures/torrentleech-search.xml @@ -0,0 +1,21 @@ + + + + TorrentLeech + + Example Show S02E03 1080p WEB H264-GROUP + tl-redacted-1 + + 2026-08-22T11:15:00Z + + + + + Example Show S02E04 1080p WEB H264-GROUP + tl-redacted-2 + https://indexer.invalid/download/tl-redacted-2 + not supplied by tracker + + + +