diff --git a/crates/models/external/src/lib.rs b/crates/models/external/src/lib.rs index 24bc5dcc4c0..cba83177d01 100644 --- a/crates/models/external/src/lib.rs +++ b/crates/models/external/src/lib.rs @@ -76,6 +76,20 @@ pub mod audiobookshelf { pub struct ListResponse { pub results: Vec, } + + #[derive(Debug, Serialize, Deserialize)] + #[serde(rename_all = "camelCase")] + pub struct MediaProgress { + pub is_finished: bool, + pub episode_id: Option, + pub library_item_id: Option, + } + + #[derive(Debug, Serialize, Deserialize)] + #[serde(rename_all = "camelCase")] + pub struct MeResponse { + pub media_progress: Vec, + } } pub mod plex { diff --git a/crates/services/importer/audiobookshelf/src/lib.rs b/crates/services/importer/audiobookshelf/src/lib.rs index 60e7450918a..30cab6691e0 100644 --- a/crates/services/importer/audiobookshelf/src/lib.rs +++ b/crates/services/importer/audiobookshelf/src/lib.rs @@ -1,8 +1,8 @@ -use std::{result::Result as StdResult, sync::Arc}; +use std::{collections::HashSet, result::Result as StdResult, sync::Arc, time::Duration}; -use anyhow::Result; +use anyhow::{Context, Result}; use application_utils::get_podcast_episode_number_by_name; -use common_utils::{get_http_client_with_tls_config, ryot_log}; +use common_utils::{USER_AGENT_STR, ryot_log, sleep_for_n_seconds}; use data_encoding::BASE64; use dependent_entity_utils::commit_metadata; use dependent_models::{ImportCompletedItem, ImportOrExportMetadataItem, ImportResult}; @@ -18,8 +18,8 @@ use media_models::{ }; use openlibrary_provider::OpenlibraryService; use reqwest::{ - Client, - header::{AUTHORIZATION, HeaderValue}, + Client, ClientBuilder, + header::{AUTHORIZATION, HeaderMap, HeaderValue, USER_AGENT}, }; use supporting_service::SupportingService; @@ -28,6 +28,7 @@ struct ImportServices<'a> { hardcover_service: &'a HardcoverService, google_books_service: &'a GoogleBooksService, open_library_service: &'a OpenlibraryService, + finished_podcast_episodes: &'a HashSet<(String, String)>, } pub async fn import( @@ -40,19 +41,31 @@ pub async fn import( let mut completed = vec![]; let mut failed = vec![]; let url = format!("{}/api", input.api_url); - let client = get_http_client_with_tls_config( - Some(vec![( - AUTHORIZATION, - HeaderValue::from_str(&format!("Bearer {}", input.api_key)).unwrap(), - )]), + let client = build_client( + &input.api_key, input.allow_insecure_connections.unwrap_or(false), ); + let me = client + .get(format!("{url}/me")) + .send() + .await? + .json::() + .await + .unwrap(); + let finished_podcast_episodes = me + .media_progress + .into_iter() + .filter(|progress| progress.is_finished) + .filter_map(|progress| Some((progress.library_item_id?, progress.episode_id?))) + .collect::>(); + let services = ImportServices { ss, hardcover_service, google_books_service, open_library_service, + finished_podcast_episodes: &finished_podcast_episodes, }; let libraries_resp = client @@ -80,7 +93,7 @@ pub async fn import( let results: Vec<_> = stream::iter(finished_items.results.into_iter().enumerate()) .map(|(idx, item)| process_item(idx, item, len, &client, &url, &services)) - .buffer_unordered(5) + .buffer_unordered(1) .collect() .await; @@ -120,73 +133,71 @@ async fn process_item( _ => { return Err(ImportFailedItem { identifier: title, + lot: Some(MediaLot::Book), step: ImportFailStep::InputTransformation, error: Some("No Google Books ID found".to_string()), - ..Default::default() }); } }, _ => { return Err(ImportFailedItem { identifier: title, + lot: Some(MediaLot::Book), error: Some("No ISBN found".to_string()), step: ImportFailStep::InputTransformation, - ..Default::default() }); } } } else if let Some(asin) = metadata.asin.clone() { (asin, MediaLot::AudioBook, MediaSource::Audible, None) } else if let Some(itunes_id) = metadata.itunes_id.clone() { - let item_details = get_item_details(client, url, &item.id, None) - .await - .map_err(|e| ImportFailedItem { - identifier: title.clone(), - error: Some(e.to_string()), - step: ImportFailStep::ItemDetailsFromSource, - ..Default::default() - })?; + let item_details = + get_item_details(client, url, &item.id) + .await + .map_err(|e| ImportFailedItem { + identifier: title.clone(), + error: Some(format!("item_id={}: {e:#}", item.id)), + lot: Some(MediaLot::Podcast), + step: ImportFailStep::ItemDetailsFromSource, + })?; match item_details.media.and_then(|m| m.episodes) { Some(episodes) => { let lot = MediaLot::Podcast; let source = MediaSource::Itunes; let mut to_return = vec![]; for episode in episodes { - ryot_log!(debug, "Importing episode {:?}", episode.title); - let episode_details = - get_item_details(client, url, &item.id, Some(episode.id.unwrap())) - .await - .map_err(|e| ImportFailedItem { - identifier: title.clone(), - error: Some(e.to_string()), - step: ImportFailStep::ItemDetailsFromSource, - ..Default::default() - })?; - if let Some(true) = - episode_details.user_media_progress.map(|u| u.is_finished) + let Some(episode_id) = episode.id.clone() else { + continue; + }; + if !services + .finished_podcast_episodes + .contains(&(item.id.clone(), episode_id)) { - let (podcast, _) = commit_metadata( - PartialMetadataWithoutId { - lot, - source, - identifier: itunes_id.clone(), - ..Default::default() - }, - services.ss, - Some(true), - ) - .await - .map_err(|e| ImportFailedItem { - identifier: title.clone(), - error: Some(e.to_string()), - step: ImportFailStep::ItemDetailsFromSource, + continue; + } + ryot_log!(debug, "Importing episode {:?}", episode.title); + let (podcast, _) = commit_metadata( + PartialMetadataWithoutId { + lot, + source, + identifier: itunes_id.clone(), ..Default::default() - })?; - if let Some(pe) = podcast.podcast_specifics.and_then(|p| { - get_podcast_episode_number_by_name(&p, &episode.title) - }) { - to_return.push(pe); - } + }, + services.ss, + Some(true), + ) + .await + .map_err(|e| ImportFailedItem { + identifier: title.clone(), + error: Some(format!("item_id={}: {e:#}", item.id)), + lot: Some(MediaLot::Podcast), + step: ImportFailStep::DatabaseCommit, + })?; + if let Some(pe) = podcast + .podcast_specifics + .and_then(|p| get_podcast_episode_number_by_name(&p, &episode.title)) + { + to_return.push(pe); } } (itunes_id, lot, source, Some(to_return)) @@ -196,7 +207,10 @@ async fn process_item( identifier: title, lot: Some(MediaLot::Podcast), step: ImportFailStep::ItemDetailsFromSource, - error: Some("No episodes found for podcast".to_string()), + error: Some(format!( + "item_id={}: No episodes found for podcast", + item.id + )), }); } } @@ -205,8 +219,11 @@ async fn process_item( return Err(ImportFailedItem { identifier: title, step: ImportFailStep::InputTransformation, - error: Some("No ASIN, ISBN or iTunes ID found".to_string()), - ..Default::default() + lot: None, + error: Some(format!( + "item_id={}: No ASIN, ISBN or iTunes ID found", + item.id + )), }); }; let mut seen_history = vec![]; @@ -234,22 +251,72 @@ async fn process_item( })) } +const MAX_ITEM_DETAILS_ATTEMPTS: u32 = 3; + async fn get_item_details( client: &Client, url: &str, id: &str, - episode: Option, ) -> Result { - let mut query = serde_json::json!({ "expanded": "1", "include": "progress" }); - if let Some(episode) = episode { - query["episode"] = serde_json::json!(episode); + let query = serde_json::json!({ "expanded": "1" }); + let request_url = format!("{url}/items/{id}"); + let mut last_error = None; + for attempt in 0..MAX_ITEM_DETAILS_ATTEMPTS { + match fetch_item_details(client, &request_url, &query).await { + Ok(item) => return Ok(item), + Err(error) => { + last_error = Some(error); + if attempt + 1 < MAX_ITEM_DETAILS_ATTEMPTS { + let sleep_time = u64::pow(2, attempt + 1); + ryot_log!( + debug, + "Attempt {}/{} failed for {}, retrying in {}s", + attempt + 1, + MAX_ITEM_DETAILS_ATTEMPTS, + request_url, + sleep_time + ); + sleep_for_n_seconds(sleep_time).await; + } + } + } } - let item = client - .get(format!("{url}/items/{id}")) - .query(&query) + Err(last_error.unwrap()) +} + +async fn fetch_item_details( + client: &Client, + request_url: &str, + query: &serde_json::Value, +) -> Result { + let response = client + .get(request_url) + .query(query) .send() - .await? - .json::() - .await?; - Ok(item) + .await + .with_context(|| format!("Network error requesting audiobookshelf API at {request_url}"))?; + let status = response.status(); + let body = response.text().await.unwrap_or_default(); + if !status.is_success() { + anyhow::bail!("Audiobookshelf returned HTTP {status} (url: {request_url}): {body}"); + } + serde_json::from_str::(&body).with_context(|| { + format!("Failed to parse JSON from Audiobookshelf response for {request_url}") + }) +} + +// Longer timeout than the shared client since podcast item lookups can be slow. +fn build_client(bearer_token: &str, danger_accept_invalid_certs: bool) -> Client { + let mut headers = HeaderMap::new(); + headers.insert(USER_AGENT, HeaderValue::from_static(USER_AGENT_STR)); + headers.insert( + AUTHORIZATION, + HeaderValue::from_str(&format!("Bearer {bearer_token}")).unwrap(), + ); + ClientBuilder::new() + .default_headers(headers) + .timeout(Duration::from_secs(60)) + .danger_accept_invalid_certs(danger_accept_invalid_certs) + .build() + .unwrap() }