From e1d35f81bc9f79dc2da4cbce9abe79f3a390ca30 Mon Sep 17 00:00:00 2001 From: Alexander Date: Tue, 21 Jul 2026 17:47:54 +0200 Subject: [PATCH] Fix torad hanging on add of torrent --- crates/torad/src/torrents.rs | 122 ++++++++++++++++++++++++----------- 1 file changed, 83 insertions(+), 39 deletions(-) diff --git a/crates/torad/src/torrents.rs b/crates/torad/src/torrents.rs index a89e6cd..800f82f 100644 --- a/crates/torad/src/torrents.rs +++ b/crates/torad/src/torrents.rs @@ -1,4 +1,4 @@ -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::path::PathBuf; use std::sync::Arc; use std::time::Duration; @@ -22,6 +22,7 @@ use crate::db::{self, TorrentRow, TorrentState}; use crate::source::SourceResolver; const POLL_INTERVAL: Duration = Duration::from_secs(2); +const ADD_TIMEOUT: Duration = Duration::from_secs(60); const BYTES_PER_MIB: f64 = 1_048_576.0; pub enum ResolveError { @@ -36,6 +37,7 @@ pub struct TorrentManager { default_output_dir: PathBuf, source: SourceResolver, tracked: Mutex>>, + adding: Mutex>, } impl TorrentManager { @@ -56,6 +58,7 @@ impl TorrentManager { default_output_dir: download_dir, source, tracked: Mutex::new(HashMap::new()), + adding: Mutex::new(HashSet::new()), })) } @@ -170,13 +173,13 @@ impl TorrentManager { let mut interval = tokio::time::interval(POLL_INTERVAL); loop { interval.tick().await; - self.pick_up_pending().await; + Arc::clone(&self).pick_up_pending().await; self.report_progress().await; } }); } - async fn pick_up_pending(&self) { + async fn pick_up_pending(self: Arc) { let pending = match db::list_pending(&self.pool).await { Ok(rows) => rows, Err(err) => { @@ -188,46 +191,74 @@ impl TorrentManager { return; } - let mut tracked = self.tracked.lock().await; - for row in pending { - if tracked.contains_key(&row.id) { - continue; + // Mark in-flight adds under brief locks so the poller doesn't re-spawn + // the same torrent on the next tick. The actual `add_torrent` call runs + // in a detached task WITHOUT holding either lock; otherwise a slow add + // (e.g. a magnet with no peers) would block every RPC that touches + // `tracked` indefinitely. + let to_start: Vec = { + let tracked = self.tracked.lock().await; + let mut adding = self.adding.lock().await; + let selected = pending + .into_iter() + .filter(|row| !tracked.contains_key(&row.id) && !adding.contains(&row.id)) + .collect::>(); + for row in &selected { + adding.insert(row.id); } - let options = AddTorrentOptions { - output_folder: Some(row.output_path.clone()), - overwrite: true, - ..Default::default() - }; - let add = match AddTorrent::from_cli_argument(&row.source) { - Ok(add) => add, - Err(err) => { - error!(id = %row.id, source = %row.source, error = %err, "failed to parse torrent source"); - continue; - } - }; - match self.session.add_torrent(add, Some(options)).await { - Ok(response) => match response.into_handle() { - Some(handle) => { - info!(id = %row.id, info_hash = %row.info_hash, "torrent added to session"); - tracked.insert(row.id, handle); + selected + }; + + for row in to_start { + let manager = Arc::clone(&self); + tokio::spawn(async move { + let options = AddTorrentOptions { + output_folder: Some(row.output_path.clone()), + overwrite: true, + ..Default::default() + }; + let add = match AddTorrent::from_cli_argument(&row.source) { + Ok(add) => add, + Err(err) => { + error!(id = %row.id, source = %row.source, error = %err, "failed to parse torrent source"); + let message = err.to_string(); + persist_error(&manager.pool, row.id, &message).await; + manager.adding.lock().await.remove(&row.id); + return; } - None => warn!(id = %row.id, "add_torrent returned no handle"), - }, - Err(err) => { - error!(id = %row.id, error = %err, "failed to add torrent to session"); - let message = err.to_string(); - let progress = db::Progress { - name: None, - total_bytes: 0, - downloaded_bytes: 0, - state: TorrentState::Error, - error_message: Some(message.as_str()), - }; - if let Err(err) = db::update_progress(&self.pool, row.id, progress).await { - error!(error = %err, "failed to persist torrent add error"); + }; + + let add_result = tokio::time::timeout( + ADD_TIMEOUT, + manager.session.add_torrent(add, Some(options)), + ) + .await; + match add_result { + Ok(Ok(response)) => match response.into_handle() { + Some(handle) => { + info!(id = %row.id, info_hash = %row.info_hash, "torrent added to session"); + let mut tracked = manager.tracked.lock().await; + tracked.insert(row.id, handle); + } + None => warn!(id = %row.id, "add_torrent returned no handle"), + }, + Ok(Err(err)) => { + error!(id = %row.id, error = %err, "failed to add torrent to session"); + let message = err.to_string(); + persist_error(&manager.pool, row.id, &message).await; + } + Err(_) => { + warn!(id = %row.id, timeout = ?ADD_TIMEOUT, "timed out adding torrent to session"); + persist_error( + &manager.pool, + row.id, + "metadata fetch timed out (no peers or unreachable trackers)", + ) + .await; } } - } + manager.adding.lock().await.remove(&row.id); + }); } } @@ -270,6 +301,19 @@ struct LiveTorrentStats { seeds: u32, } +async fn persist_error(pool: &PgPool, id: Uuid, message: &str) { + let progress = db::Progress { + name: None, + total_bytes: 0, + downloaded_bytes: 0, + state: TorrentState::Error, + error_message: Some(message), + }; + if let Err(err) = db::update_progress(pool, id, progress).await { + error!(id = %id, error = %err, "failed to persist torrent add error"); + } +} + fn extract_live_stats(tracked: &HashMap>, id: &Uuid) -> LiveTorrentStats { let Some(handle) = tracked.get(id) else { return LiveTorrentStats::default();