feat(desktop): Yoda integration patches on v8.1.1
- macOS-style window controls via gtk-decoration-layout - torrent file preview and context menu improvements - single-instance relay for new torrents/magnet links
This commit is contained in:
parent
00b9748516
commit
99219f90bc
8 changed files with 405 additions and 13 deletions
|
|
@ -4,10 +4,14 @@
|
|||
mod config;
|
||||
|
||||
use std::{
|
||||
env,
|
||||
fs::{File, OpenOptions},
|
||||
io::{BufReader, BufWriter},
|
||||
path::Path,
|
||||
io::{BufReader, BufWriter, Read, Write},
|
||||
os::unix::net::{UnixListener, UnixStream},
|
||||
path::{Path, PathBuf},
|
||||
process::Command,
|
||||
sync::Arc,
|
||||
thread,
|
||||
};
|
||||
|
||||
use anyhow::Context;
|
||||
|
|
@ -26,10 +30,17 @@ use librqbit::{
|
|||
};
|
||||
use parking_lot::RwLock;
|
||||
use serde::Serialize;
|
||||
use tauri::{AppHandle, Emitter};
|
||||
use tracing::{error, error_span, info, warn};
|
||||
|
||||
const ERR_NOT_CONFIGURED: ApiError =
|
||||
ApiError::new_from_text(StatusCode::FAILED_DEPENDENCY, "not configured");
|
||||
const TORRENTS_CHANGED_EVENT: &str = "rqbit-desktop-torrents-changed";
|
||||
|
||||
#[derive(Clone, Serialize)]
|
||||
struct TorrentsChangedPayload {
|
||||
id: usize,
|
||||
}
|
||||
|
||||
struct StateShared {
|
||||
config: config::RqbitDesktopConfig,
|
||||
|
|
@ -39,6 +50,7 @@ struct StateShared {
|
|||
struct State {
|
||||
config_filename: String,
|
||||
shared: Arc<RwLock<Option<StateShared>>>,
|
||||
pending_launch_inputs: Arc<RwLock<Vec<String>>>,
|
||||
init_logging: InitLoggingResult,
|
||||
}
|
||||
|
||||
|
|
@ -194,6 +206,7 @@ impl State {
|
|||
return Self {
|
||||
config_filename,
|
||||
shared,
|
||||
pending_launch_inputs: Arc::new(RwLock::new(Vec::new())),
|
||||
init_logging,
|
||||
};
|
||||
}
|
||||
|
|
@ -202,6 +215,7 @@ impl State {
|
|||
config_filename,
|
||||
init_logging,
|
||||
shared: Arc::new(RwLock::new(None)),
|
||||
pending_launch_inputs: Arc::new(RwLock::new(Vec::new())),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -244,6 +258,211 @@ impl State {
|
|||
}
|
||||
}
|
||||
|
||||
fn is_torrent_launch_input(input: &str) -> bool {
|
||||
input.starts_with("magnet:")
|
||||
|| input.starts_with("http://")
|
||||
|| input.starts_with("https://")
|
||||
|| input.starts_with("file://")
|
||||
|| Path::new(input)
|
||||
.extension()
|
||||
.and_then(|extension| extension.to_str())
|
||||
.is_some_and(|extension| extension.eq_ignore_ascii_case("torrent"))
|
||||
}
|
||||
|
||||
fn collect_launch_inputs() -> Vec<String> {
|
||||
env::args()
|
||||
.skip(1)
|
||||
.filter(|arg| is_torrent_launch_input(arg))
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn desktop_ipc_socket_path() -> PathBuf {
|
||||
let runtime_dir = env::var_os("XDG_RUNTIME_DIR")
|
||||
.map(PathBuf::from)
|
||||
.unwrap_or_else(env::temp_dir);
|
||||
let user = env::var("USER").unwrap_or_else(|_| "user".to_string());
|
||||
runtime_dir.join(format!("rqbit-desktop-{user}.sock"))
|
||||
}
|
||||
|
||||
fn forward_inputs_to_existing(socket_path: &Path, inputs: &[String]) -> anyhow::Result<()> {
|
||||
let mut stream = UnixStream::connect(socket_path)?;
|
||||
let payload = serde_json::to_vec(inputs)?;
|
||||
stream.write_all(&payload)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn bind_ipc_listener(socket_path: &Path) -> std::io::Result<UnixListener> {
|
||||
if UnixStream::connect(socket_path).is_err() {
|
||||
let _ = std::fs::remove_file(socket_path);
|
||||
}
|
||||
|
||||
if let Some(parent) = socket_path.parent() {
|
||||
std::fs::create_dir_all(parent)?;
|
||||
}
|
||||
|
||||
UnixListener::bind(socket_path)
|
||||
}
|
||||
|
||||
fn spawn_ipc_listener(
|
||||
socket_path: PathBuf,
|
||||
shared: Arc<RwLock<Option<StateShared>>>,
|
||||
pending: Arc<RwLock<Vec<String>>>,
|
||||
app_handle: AppHandle,
|
||||
) {
|
||||
let listener = match bind_ipc_listener(&socket_path) {
|
||||
Ok(listener) => listener,
|
||||
Err(err) => {
|
||||
warn!(path=%socket_path.display(), error=%err, "couldn't bind desktop IPC socket");
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let handle = tokio::runtime::Handle::current();
|
||||
|
||||
thread::spawn(move || {
|
||||
for stream in listener.incoming() {
|
||||
let shared = shared.clone();
|
||||
let pending = pending.clone();
|
||||
let app_handle = app_handle.clone();
|
||||
match stream {
|
||||
Ok(mut stream) => {
|
||||
let mut payload = String::new();
|
||||
if let Err(err) = stream.read_to_string(&mut payload) {
|
||||
warn!(error=%err, "couldn't read desktop IPC payload");
|
||||
continue;
|
||||
}
|
||||
|
||||
let inputs = match serde_json::from_str::<Vec<String>>(&payload) {
|
||||
Ok(inputs) => inputs,
|
||||
Err(err) => {
|
||||
warn!(error=%err, "couldn't parse desktop IPC payload");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
handle.spawn(async move {
|
||||
add_launch_inputs(shared, pending, inputs, Some(app_handle)).await;
|
||||
});
|
||||
}
|
||||
Err(err) => warn!(error=%err, "desktop IPC accept failed"),
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
async fn add_launch_inputs(
|
||||
shared: Arc<RwLock<Option<StateShared>>>,
|
||||
pending: Arc<RwLock<Vec<String>>>,
|
||||
inputs: Vec<String>,
|
||||
app_handle: Option<AppHandle>,
|
||||
) {
|
||||
for input in inputs {
|
||||
match add_launch_input(shared.clone(), pending.clone(), &input).await {
|
||||
Ok(Some(id)) => {
|
||||
if let Some(app_handle) = &app_handle {
|
||||
if let Err(err) =
|
||||
app_handle.emit(TORRENTS_CHANGED_EVENT, TorrentsChangedPayload { id })
|
||||
{
|
||||
warn!(input=%input, error=%err, "couldn't emit desktop torrent change event");
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(None) => {}
|
||||
Err(err) => {
|
||||
warn!(input=%input, error=%err, "couldn't add torrent from desktop launch input");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn add_launch_input(
|
||||
shared: Arc<RwLock<Option<StateShared>>>,
|
||||
pending: Arc<RwLock<Vec<String>>>,
|
||||
input: &str,
|
||||
) -> anyhow::Result<Option<usize>> {
|
||||
let api = {
|
||||
let g = shared.read();
|
||||
g.as_ref().and_then(|state| state.api.as_ref()).cloned()
|
||||
};
|
||||
let Some(api) = api else {
|
||||
pending.write().push(input.to_owned());
|
||||
info!(input=%input, "queued desktop launch input until rqbit is configured");
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
let opts = AddTorrentOptions {
|
||||
overwrite: true,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let torrent = if input.starts_with("magnet:")
|
||||
|| input.starts_with("http://")
|
||||
|| input.starts_with("https://")
|
||||
{
|
||||
AddTorrent::Url(input.to_string().into())
|
||||
} else {
|
||||
let path = path_from_launch_input(input)?;
|
||||
let bytes = std::fs::read(&path)
|
||||
.with_context(|| format!("couldn't read torrent file {}", path.display()))?;
|
||||
AddTorrent::TorrentFileBytes(bytes.into())
|
||||
};
|
||||
|
||||
api.api_add_torrent(torrent, Some(opts))
|
||||
.await
|
||||
.map(|response| response.id)
|
||||
.map_err(|err| anyhow::anyhow!("{err:?}"))
|
||||
}
|
||||
|
||||
async fn drain_pending_launch_inputs(
|
||||
shared: Arc<RwLock<Option<StateShared>>>,
|
||||
pending: Arc<RwLock<Vec<String>>>,
|
||||
app_handle: Option<AppHandle>,
|
||||
) {
|
||||
let inputs = {
|
||||
let mut pending = pending.write();
|
||||
if pending.is_empty() {
|
||||
return;
|
||||
}
|
||||
pending.drain(..).collect::<Vec<_>>()
|
||||
};
|
||||
|
||||
add_launch_inputs(shared, pending, inputs, app_handle).await;
|
||||
}
|
||||
|
||||
fn path_from_launch_input(input: &str) -> anyhow::Result<PathBuf> {
|
||||
if let Some(uri_path) = input.strip_prefix("file://localhost/") {
|
||||
return Ok(PathBuf::from(format!("/{}", percent_decode(uri_path)?)));
|
||||
}
|
||||
|
||||
if let Some(uri_path) = input.strip_prefix("file://") {
|
||||
return Ok(PathBuf::from(percent_decode(uri_path)?));
|
||||
}
|
||||
|
||||
Ok(PathBuf::from(input))
|
||||
}
|
||||
|
||||
fn percent_decode(value: &str) -> anyhow::Result<String> {
|
||||
let bytes = value.as_bytes();
|
||||
let mut decoded = Vec::with_capacity(bytes.len());
|
||||
let mut index = 0;
|
||||
|
||||
while index < bytes.len() {
|
||||
if bytes[index] == b'%' {
|
||||
let Some(hex) = bytes.get(index + 1..index + 3) else {
|
||||
anyhow::bail!("invalid percent-encoded path");
|
||||
};
|
||||
let hex = std::str::from_utf8(hex).context("invalid percent-encoded path")?;
|
||||
decoded.push(u8::from_str_radix(hex, 16).context("invalid percent-encoded path")?);
|
||||
index += 3;
|
||||
} else {
|
||||
decoded.push(bytes[index]);
|
||||
index += 1;
|
||||
}
|
||||
}
|
||||
|
||||
String::from_utf8(decoded).context("invalid UTF-8 path")
|
||||
}
|
||||
|
||||
#[derive(Default, Serialize)]
|
||||
struct CurrentState {
|
||||
config: Option<RqbitDesktopConfig>,
|
||||
|
|
@ -269,10 +488,18 @@ fn config_current(state: tauri::State<'_, State>) -> CurrentState {
|
|||
|
||||
#[tauri::command]
|
||||
async fn config_change(
|
||||
app_handle: AppHandle,
|
||||
state: tauri::State<'_, State>,
|
||||
config: RqbitDesktopConfig,
|
||||
) -> Result<EmptyJsonResponse, ApiError> {
|
||||
state.configure(config).await.map(|_| EmptyJsonResponse {})
|
||||
state.configure(config).await?;
|
||||
drain_pending_launch_inputs(
|
||||
state.shared.clone(),
|
||||
state.pending_launch_inputs.clone(),
|
||||
Some(app_handle),
|
||||
)
|
||||
.await;
|
||||
Ok(EmptyJsonResponse {})
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
|
|
@ -374,6 +601,37 @@ async fn stats(state: tauri::State<'_, State>) -> Result<SessionStatsSnapshot, A
|
|||
Ok(state.api()?.api_session_stats())
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
async fn open_output_folder(path: String) -> Result<EmptyJsonResponse, ApiError> {
|
||||
let path = PathBuf::from(path);
|
||||
let metadata = std::fs::metadata(&path)
|
||||
.with_context(|| format!("couldn't access output folder {}", path.display()))
|
||||
.map_err(|e| ApiError::new_from_anyhow(StatusCode::BAD_REQUEST, e))?;
|
||||
|
||||
if !metadata.is_dir() {
|
||||
return Err(ApiError::new_from_anyhow(
|
||||
StatusCode::BAD_REQUEST,
|
||||
anyhow::anyhow!("output path is not a folder: {}", path.display()),
|
||||
));
|
||||
}
|
||||
|
||||
let status = Command::new("xdg-open")
|
||||
.arg(&path)
|
||||
.status()
|
||||
.or_else(|_| Command::new("gio").arg("open").arg(&path).status())
|
||||
.with_context(|| format!("couldn't open output folder {}", path.display()))
|
||||
.map_err(|e| ApiError::new_from_anyhow(StatusCode::INTERNAL_SERVER_ERROR, e))?;
|
||||
|
||||
if !status.success() {
|
||||
return Err(ApiError::new_from_anyhow(
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
anyhow::anyhow!("file manager exited with status {status}"),
|
||||
));
|
||||
}
|
||||
|
||||
Ok(EmptyJsonResponse {})
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
fn get_version() -> &'static str {
|
||||
env!("CARGO_PKG_VERSION")
|
||||
|
|
@ -393,11 +651,41 @@ async fn start() {
|
|||
Err(e) => warn!("failed increasing open file limit: {:#}", e),
|
||||
};
|
||||
|
||||
let launch_inputs = collect_launch_inputs();
|
||||
let socket_path = desktop_ipc_socket_path();
|
||||
|
||||
if !launch_inputs.is_empty() && forward_inputs_to_existing(&socket_path, &launch_inputs).is_ok()
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
let state = State::new(init_logging_result).await;
|
||||
let shared = state.shared.clone();
|
||||
let pending = state.pending_launch_inputs.clone();
|
||||
|
||||
tauri::Builder::default()
|
||||
.plugin(tauri_plugin_shell::init())
|
||||
.manage(state)
|
||||
.setup(move |app| {
|
||||
let app_handle = app.handle().clone();
|
||||
spawn_ipc_listener(
|
||||
socket_path.clone(),
|
||||
shared.clone(),
|
||||
pending.clone(),
|
||||
app_handle.clone(),
|
||||
);
|
||||
|
||||
if !launch_inputs.is_empty() {
|
||||
tauri::async_runtime::spawn(add_launch_inputs(
|
||||
shared.clone(),
|
||||
pending.clone(),
|
||||
launch_inputs.clone(),
|
||||
Some(app_handle),
|
||||
));
|
||||
}
|
||||
|
||||
Ok(())
|
||||
})
|
||||
.invoke_handler(tauri::generate_handler![
|
||||
torrents_list,
|
||||
torrent_details,
|
||||
|
|
@ -409,6 +697,7 @@ async fn start() {
|
|||
torrent_action_start,
|
||||
torrent_action_configure,
|
||||
torrent_create_from_base64_file,
|
||||
open_output_folder,
|
||||
stats,
|
||||
get_version,
|
||||
config_default,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue