Fix torad hanging on add of torrent
This commit is contained in:
@@ -1,4 +1,4 @@
|
|||||||
use std::collections::HashMap;
|
use std::collections::{HashMap, HashSet};
|
||||||
use std::path::PathBuf;
|
use std::path::PathBuf;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
@@ -22,6 +22,7 @@ use crate::db::{self, TorrentRow, TorrentState};
|
|||||||
use crate::source::SourceResolver;
|
use crate::source::SourceResolver;
|
||||||
|
|
||||||
const POLL_INTERVAL: Duration = Duration::from_secs(2);
|
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;
|
const BYTES_PER_MIB: f64 = 1_048_576.0;
|
||||||
|
|
||||||
pub enum ResolveError {
|
pub enum ResolveError {
|
||||||
@@ -36,6 +37,7 @@ pub struct TorrentManager {
|
|||||||
default_output_dir: PathBuf,
|
default_output_dir: PathBuf,
|
||||||
source: SourceResolver,
|
source: SourceResolver,
|
||||||
tracked: Mutex<HashMap<Uuid, Arc<ManagedTorrent>>>,
|
tracked: Mutex<HashMap<Uuid, Arc<ManagedTorrent>>>,
|
||||||
|
adding: Mutex<HashSet<Uuid>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TorrentManager {
|
impl TorrentManager {
|
||||||
@@ -56,6 +58,7 @@ impl TorrentManager {
|
|||||||
default_output_dir: download_dir,
|
default_output_dir: download_dir,
|
||||||
source,
|
source,
|
||||||
tracked: Mutex::new(HashMap::new()),
|
tracked: Mutex::new(HashMap::new()),
|
||||||
|
adding: Mutex::new(HashSet::new()),
|
||||||
}))
|
}))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -170,13 +173,13 @@ impl TorrentManager {
|
|||||||
let mut interval = tokio::time::interval(POLL_INTERVAL);
|
let mut interval = tokio::time::interval(POLL_INTERVAL);
|
||||||
loop {
|
loop {
|
||||||
interval.tick().await;
|
interval.tick().await;
|
||||||
self.pick_up_pending().await;
|
Arc::clone(&self).pick_up_pending().await;
|
||||||
self.report_progress().await;
|
self.report_progress().await;
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn pick_up_pending(&self) {
|
async fn pick_up_pending(self: Arc<Self>) {
|
||||||
let pending = match db::list_pending(&self.pool).await {
|
let pending = match db::list_pending(&self.pool).await {
|
||||||
Ok(rows) => rows,
|
Ok(rows) => rows,
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
@@ -188,46 +191,74 @@ impl TorrentManager {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut tracked = self.tracked.lock().await;
|
// Mark in-flight adds under brief locks so the poller doesn't re-spawn
|
||||||
for row in pending {
|
// the same torrent on the next tick. The actual `add_torrent` call runs
|
||||||
if tracked.contains_key(&row.id) {
|
// in a detached task WITHOUT holding either lock; otherwise a slow add
|
||||||
continue;
|
// (e.g. a magnet with no peers) would block every RPC that touches
|
||||||
|
// `tracked` indefinitely.
|
||||||
|
let to_start: Vec<TorrentRow> = {
|
||||||
|
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::<Vec<_>>();
|
||||||
|
for row in &selected {
|
||||||
|
adding.insert(row.id);
|
||||||
}
|
}
|
||||||
let options = AddTorrentOptions {
|
selected
|
||||||
output_folder: Some(row.output_path.clone()),
|
};
|
||||||
overwrite: true,
|
|
||||||
..Default::default()
|
for row in to_start {
|
||||||
};
|
let manager = Arc::clone(&self);
|
||||||
let add = match AddTorrent::from_cli_argument(&row.source) {
|
tokio::spawn(async move {
|
||||||
Ok(add) => add,
|
let options = AddTorrentOptions {
|
||||||
Err(err) => {
|
output_folder: Some(row.output_path.clone()),
|
||||||
error!(id = %row.id, source = %row.source, error = %err, "failed to parse torrent source");
|
overwrite: true,
|
||||||
continue;
|
..Default::default()
|
||||||
}
|
};
|
||||||
};
|
let add = match AddTorrent::from_cli_argument(&row.source) {
|
||||||
match self.session.add_torrent(add, Some(options)).await {
|
Ok(add) => add,
|
||||||
Ok(response) => match response.into_handle() {
|
Err(err) => {
|
||||||
Some(handle) => {
|
error!(id = %row.id, source = %row.source, error = %err, "failed to parse torrent source");
|
||||||
info!(id = %row.id, info_hash = %row.info_hash, "torrent added to session");
|
let message = err.to_string();
|
||||||
tracked.insert(row.id, handle);
|
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) => {
|
let add_result = tokio::time::timeout(
|
||||||
error!(id = %row.id, error = %err, "failed to add torrent to session");
|
ADD_TIMEOUT,
|
||||||
let message = err.to_string();
|
manager.session.add_torrent(add, Some(options)),
|
||||||
let progress = db::Progress {
|
)
|
||||||
name: None,
|
.await;
|
||||||
total_bytes: 0,
|
match add_result {
|
||||||
downloaded_bytes: 0,
|
Ok(Ok(response)) => match response.into_handle() {
|
||||||
state: TorrentState::Error,
|
Some(handle) => {
|
||||||
error_message: Some(message.as_str()),
|
info!(id = %row.id, info_hash = %row.info_hash, "torrent added to session");
|
||||||
};
|
let mut tracked = manager.tracked.lock().await;
|
||||||
if let Err(err) = db::update_progress(&self.pool, row.id, progress).await {
|
tracked.insert(row.id, handle);
|
||||||
error!(error = %err, "failed to persist torrent add error");
|
}
|
||||||
|
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,
|
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<Uuid, Arc<ManagedTorrent>>, id: &Uuid) -> LiveTorrentStats {
|
fn extract_live_stats(tracked: &HashMap<Uuid, Arc<ManagedTorrent>>, id: &Uuid) -> LiveTorrentStats {
|
||||||
let Some(handle) = tracked.get(id) else {
|
let Some(handle) = tracked.get(id) else {
|
||||||
return LiveTorrentStats::default();
|
return LiveTorrentStats::default();
|
||||||
|
|||||||
Reference in New Issue
Block a user