refactor(backend): clean up sse mapping and handler logic

This commit is contained in:
spinline
2026-02-03 21:45:24 +03:00
parent c3431db35f
commit 251da58a82
2 changed files with 213 additions and 232 deletions

View File

@@ -1,35 +1,98 @@
use crate::xmlrpc::{parse_multicall_response, RtorrentClient, XmlRpcError};
use crate::AppState;
use axum::extract::State;
use axum::response::sse::{Event, Sse};
use futures::stream::{self, Stream};
use shared::{AppEvent, Torrent, TorrentStatus};
use std::convert::Infallible;
use tokio_stream::StreamExt;
// Helper (should be moved to utils)
fn parse_size(s: &str) -> i64 {
s.parse().unwrap_or(0)
// Constants for rTorrent fields to ensure query and parser stay in sync
const RTORRENT_FIELDS: &[&str] = &[
"", // 0: default (ignored)
"main", // 1: view
"d.hash=", // 0 -> row index starts after view
"d.name=", // 1
"d.size_bytes=", // 2
"d.bytes_done=", // 3
"d.down.rate=", // 4
"d.up.rate=", // 5
"d.state=", // 6
"d.complete=", // 7
"d.message=", // 8
"d.left_bytes=", // 9
"d.creation_date=", // 10
"d.hashing=", // 11
];
fn parse_long(s: Option<&String>) -> i64 {
s.map(|v| v.parse().unwrap_or(0)).unwrap_or(0)
}
fn parse_string(s: Option<&String>) -> String {
s.cloned().unwrap_or_default()
}
/// Converts a raw row of strings from rTorrent XML-RPC into a generic Torrent struct
fn from_rtorrent_row(row: Vec<String>) -> Torrent {
// Indexes correspond to the params list below (excluding the first two view/target args)
let hash = parse_string(row.get(0));
let name = parse_string(row.get(1));
let size = parse_long(row.get(2));
let completed = parse_long(row.get(3));
let down_rate = parse_long(row.get(4));
let up_rate = parse_long(row.get(5));
let state = parse_long(row.get(6));
let is_complete = parse_long(row.get(7));
let message = parse_string(row.get(8));
let left_bytes = parse_long(row.get(9));
let added_date = parse_long(row.get(10));
let is_hashing = parse_long(row.get(11));
let percent_complete = if size > 0 {
(completed as f64 / size as f64) * 100.0
} else {
0.0
};
// Status Logic
let status = if !message.is_empty() {
TorrentStatus::Error
} else if is_hashing != 0 {
TorrentStatus::Checking
} else if state == 0 {
TorrentStatus::Paused
} else if is_complete != 0 {
TorrentStatus::Seeding
} else {
TorrentStatus::Downloading
};
// ETA Logic (seconds)
let eta = if down_rate > 0 && left_bytes > 0 {
left_bytes / down_rate
} else {
0
};
Torrent {
hash,
name,
size,
completed,
down_rate,
up_rate,
eta,
percent_complete,
status,
error_message: message,
added_date,
}
}
pub async fn fetch_torrents(client: &RtorrentClient) -> Result<Vec<Torrent>, XmlRpcError> {
// d.multicall2("", "main", ...)
let params = vec![
"",
"main",
"d.hash=",
"d.name=",
"d.size_bytes=",
"d.bytes_done=",
"d.down.rate=",
"d.up.rate=",
"d.state=", // 6
"d.complete=", // 7
"d.message=", // 8
"d.left_bytes=", // 9
"d.creation_date=", // 10
"d.hashing=", // 11
];
let xml = client.call("d.multicall2", &params).await?;
let xml = client.call("d.multicall2", RTORRENT_FIELDS).await?;
if xml.trim().is_empty() {
return Err(XmlRpcError::Parse("Empty response from SCGI".to_string()));
@@ -37,75 +100,11 @@ pub async fn fetch_torrents(client: &RtorrentClient) -> Result<Vec<Torrent>, Xml
let rows = parse_multicall_response(&xml)?;
let torrents = rows
.into_iter()
.map(|row| {
// row map indexes:
// 0: hash, 1: name, 2: size, 3: completed, 4: down_rate, 5: up_rate
// 6: state, 7: complete, 8: message, 9: left_bytes, 10: added, 11: hashing
let hash = row.get(0).cloned().unwrap_or_default();
let name = row.get(1).cloned().unwrap_or_default();
let size = parse_size(row.get(2).unwrap_or(&"0".to_string()));
let completed = parse_size(row.get(3).unwrap_or(&"0".to_string()));
let down_rate = parse_size(row.get(4).unwrap_or(&"0".to_string()));
let up_rate = parse_size(row.get(5).unwrap_or(&"0".to_string()));
let state = parse_size(row.get(6).unwrap_or(&"0".to_string()));
let is_complete = parse_size(row.get(7).unwrap_or(&"0".to_string()));
let message = row.get(8).cloned().unwrap_or_default();
let left_bytes = parse_size(row.get(9).unwrap_or(&"0".to_string()));
let added_date = parse_size(row.get(10).unwrap_or(&"0".to_string()));
let is_hashing = parse_size(row.get(11).unwrap_or(&"0".to_string()));
let percent_complete = if size > 0 {
(completed as f64 / size as f64) * 100.0
} else {
0.0
};
// Status Logic
let status = if !message.is_empty() {
TorrentStatus::Error
} else if is_hashing != 0 {
TorrentStatus::Checking
} else if state == 0 {
TorrentStatus::Paused
} else if is_complete != 0 {
TorrentStatus::Seeding
} else {
TorrentStatus::Downloading
};
// ETA Logic (seconds)
let eta = if down_rate > 0 && left_bytes > 0 {
left_bytes / down_rate
} else {
0
};
Torrent {
hash,
name,
size,
completed,
down_rate,
up_rate,
eta,
percent_complete,
status,
error_message: message,
added_date,
}
})
.collect();
let torrents = rows.into_iter().map(from_rtorrent_row).collect();
Ok(torrents)
}
use crate::AppState;
use axum::extract::State; // Import from crate root
pub async fn sse_handler(
State(state): State<AppState>,
) -> Sse<impl Stream<Item = Result<Event, Infallible>>> {