Compare commits
12
Commits
v2.5.3.1-beta
...
v2.5.6
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1e2b4706d5 | ||
|
|
c3f26b77e8 | ||
|
|
9957c9b70f | ||
|
|
d4f52e3ec6 | ||
|
|
3ca52e2649 | ||
|
|
90fda425bd | ||
|
|
124d9c2dd4 | ||
|
|
e019ddccc9 | ||
|
|
61f2545588 | ||
|
|
b751d82748 | ||
|
|
c27ac0f41c | ||
|
|
f65b15452f |
Generated
+1
-1
@@ -6453,7 +6453,7 @@ checksum = "e13c156562582aa81c60cb29407084cdb54c4164760106ab78e6c5b0858cf64e"
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zerosend"
|
name = "zerosend"
|
||||||
version = "2.5.3-beta.1"
|
version = "2.5.6"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"arboard",
|
"arboard",
|
||||||
"async-stream",
|
"async-stream",
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "zerosend"
|
name = "zerosend"
|
||||||
version = "2.5.3-beta.1"
|
version = "2.5.6"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
description = "Fast, peer-to-peer file transfer designed for ZeroTier, Tailscale, Radmin VPN and LAN"
|
description = "Fast, peer-to-peer file transfer designed for ZeroTier, Tailscale, Radmin VPN and LAN"
|
||||||
authors = ["RarDog"]
|
authors = ["RarDog"]
|
||||||
|
|||||||
@@ -41,6 +41,10 @@ pub struct AppConfig {
|
|||||||
pub peer_aliases: HashMap<String, String>,
|
pub peer_aliases: HashMap<String, String>,
|
||||||
pub window_width: f32,
|
pub window_width: f32,
|
||||||
pub window_height: f32,
|
pub window_height: f32,
|
||||||
|
#[serde(default)]
|
||||||
|
pub window_x: Option<f32>,
|
||||||
|
#[serde(default)]
|
||||||
|
pub window_y: Option<f32>,
|
||||||
pub close_to_tray: bool,
|
pub close_to_tray: bool,
|
||||||
pub minimize_to_tray: bool,
|
pub minimize_to_tray: bool,
|
||||||
pub auto_zip_folders: bool,
|
pub auto_zip_folders: bool,
|
||||||
@@ -89,6 +93,8 @@ impl Default for AppConfig {
|
|||||||
peer_aliases: HashMap::new(),
|
peer_aliases: HashMap::new(),
|
||||||
window_width: 940.0,
|
window_width: 940.0,
|
||||||
window_height: 640.0,
|
window_height: 640.0,
|
||||||
|
window_x: None,
|
||||||
|
window_y: None,
|
||||||
close_to_tray: true,
|
close_to_tray: true,
|
||||||
minimize_to_tray: true,
|
minimize_to_tray: true,
|
||||||
auto_zip_folders: true,
|
auto_zip_folders: true,
|
||||||
@@ -212,3 +218,29 @@ fn get_default_download_dir() -> PathBuf {
|
|||||||
let _ = std::fs::create_dir_all(&fallback);
|
let _ = std::fs::create_dir_all(&fallback);
|
||||||
fallback
|
fallback
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn load_transfers_history() -> Vec<crate::transfer::TransferProgress> {
|
||||||
|
let path = AppConfig::get_app_dir().join("transfers_history.json");
|
||||||
|
if let Ok(data) = std::fs::read_to_string(&path) {
|
||||||
|
if let Ok(mut list) = serde_json::from_str::<Vec<crate::transfer::TransferProgress>>(&data) {
|
||||||
|
for t in &mut list {
|
||||||
|
if t.state == crate::transfer::TransferState::InProgress
|
||||||
|
|| t.state == crate::transfer::TransferState::PendingConfirmation
|
||||||
|
|| t.state == crate::transfer::TransferState::Paused
|
||||||
|
{
|
||||||
|
t.state = crate::transfer::TransferState::Canceled;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return list;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Vec::new()
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn save_transfers_history(transfers: &[crate::transfer::TransferProgress]) {
|
||||||
|
let path = AppConfig::get_app_dir().join("transfers_history.json");
|
||||||
|
let to_save: Vec<_> = transfers.iter().take(100).cloned().collect();
|
||||||
|
if let Ok(json) = serde_json::to_string_pretty(&to_save) {
|
||||||
|
let _ = std::fs::write(&path, json);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+12
-5
@@ -304,16 +304,23 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let start_minimized = (cli.tray || cli.minimized) && initial_files.is_empty();
|
let start_minimized = (cli.tray || cli.minimized) && initial_files.is_empty();
|
||||||
let options = eframe::NativeOptions {
|
let mut viewport = eframe::egui::ViewportBuilder::default()
|
||||||
viewport: eframe::egui::ViewportBuilder::default()
|
|
||||||
.with_inner_size([config.window_width, config.window_height])
|
.with_inner_size([config.window_width, config.window_height])
|
||||||
.with_min_inner_size([780.0, 500.0])
|
.with_min_inner_size([780.0, 500.0])
|
||||||
.with_visible(!start_minimized)
|
.with_visible(!start_minimized)
|
||||||
.with_title(format!("ZeroSend v{} - P2P Sharing for ZeroTier & VPN", env!("CARGO_PKG_VERSION"))),
|
.with_title(format!("ZeroSend v{} - P2P Sharing for ZeroTier & VPN", env!("CARGO_PKG_VERSION")));
|
||||||
|
|
||||||
|
if let (Some(x), Some(y)) = (config.window_x, config.window_y) {
|
||||||
|
viewport = viewport.with_position([x, y]);
|
||||||
|
}
|
||||||
|
|
||||||
|
let options = eframe::NativeOptions {
|
||||||
|
viewport,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
let disc_ref = &discovery;
|
let discovery = std::sync::Arc::new(discovery);
|
||||||
|
let disc_clone = discovery.clone();
|
||||||
let srv = server;
|
let srv = server;
|
||||||
let in_req_rx = channels.incoming_request_rx;
|
let in_req_rx = channels.incoming_request_rx;
|
||||||
let in_txt_rx = channels.incoming_text_rx;
|
let in_txt_rx = channels.incoming_text_rx;
|
||||||
@@ -324,7 +331,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
|||||||
Box::new(move |cc| {
|
Box::new(move |cc| {
|
||||||
Ok(Box::new(ui::gui::ZeroSendApp::new(
|
Ok(Box::new(ui::gui::ZeroSendApp::new(
|
||||||
cc,
|
cc,
|
||||||
disc_ref,
|
disc_clone,
|
||||||
srv,
|
srv,
|
||||||
config,
|
config,
|
||||||
initial_files,
|
initial_files,
|
||||||
|
|||||||
@@ -11,6 +11,8 @@ use tokio::sync::RwLock;
|
|||||||
use tracing::{debug, error, info, warn};
|
use tracing::{debug, error, info, warn};
|
||||||
|
|
||||||
pub const DEFAULT_DISCOVERY_PORT: u16 = 53317;
|
pub const DEFAULT_DISCOVERY_PORT: u16 = 53317;
|
||||||
|
pub const MULTICAST_IPV4_DISCOVERY: Ipv4Addr = Ipv4Addr::new(239, 255, 255, 250);
|
||||||
|
pub const MDNS_IPV4_DISCOVERY: Ipv4Addr = Ipv4Addr::new(224, 0, 0, 251);
|
||||||
#[allow(dead_code)]
|
#[allow(dead_code)]
|
||||||
pub const DEFAULT_TRANSFER_PORT: u16 = 53318;
|
pub const DEFAULT_TRANSFER_PORT: u16 = 53318;
|
||||||
pub const BEACON_INTERVAL_SECS: u64 = 3;
|
pub const BEACON_INTERVAL_SECS: u64 = 3;
|
||||||
@@ -143,8 +145,10 @@ impl DiscoveryService {
|
|||||||
let (shutdown_tx, _) = tokio::sync::broadcast::channel::<()>(1);
|
let (shutdown_tx, _) = tokio::sync::broadcast::channel::<()>(1);
|
||||||
self.shutdown_tx = Some(shutdown_tx.clone());
|
self.shutdown_tx = Some(shutdown_tx.clone());
|
||||||
|
|
||||||
// 1. Create UDP broadcast sending socket
|
// 1. Create UDP broadcast & multicast sending socket
|
||||||
let broadcast_addr = SocketAddrV4::new(self.selected_iface.broadcast, DEFAULT_DISCOVERY_PORT);
|
let broadcast_addr = SocketAddrV4::new(self.selected_iface.broadcast, DEFAULT_DISCOVERY_PORT);
|
||||||
|
let multicast_addr = SocketAddrV4::new(MULTICAST_IPV4_DISCOVERY, DEFAULT_DISCOVERY_PORT);
|
||||||
|
let mdns_addr = SocketAddrV4::new(MDNS_IPV4_DISCOVERY, DEFAULT_DISCOVERY_PORT);
|
||||||
let my_info_static = self.my_info.clone();
|
let my_info_static = self.my_info.clone();
|
||||||
let shared_folders = self.shared_folders.clone();
|
let shared_folders = self.shared_folders.clone();
|
||||||
let mut shutdown_rx1 = shutdown_tx.subscribe();
|
let mut shutdown_rx1 = shutdown_tx.subscribe();
|
||||||
@@ -163,8 +167,8 @@ impl DiscoveryService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
info!(
|
info!(
|
||||||
"Started discovery beacon on {} (broadcast: {})",
|
"Started discovery beacon on {} (broadcast: {}, multicast: {})",
|
||||||
my_info_static.ip, broadcast_addr
|
my_info_static.ip, broadcast_addr, multicast_addr
|
||||||
);
|
);
|
||||||
|
|
||||||
let mut interval = tokio::time::interval(Duration::from_secs(BEACON_INTERVAL_SECS));
|
let mut interval = tokio::time::interval(Duration::from_secs(BEACON_INTERVAL_SECS));
|
||||||
@@ -175,6 +179,8 @@ impl DiscoveryService {
|
|||||||
current_beacon.shared_folders = shared_folders.read().await.clone();
|
current_beacon.shared_folders = shared_folders.read().await.clone();
|
||||||
if let Ok(data) = serde_json::to_vec(¤t_beacon) {
|
if let Ok(data) = serde_json::to_vec(¤t_beacon) {
|
||||||
let _ = socket.send_to(&data, broadcast_addr).await;
|
let _ = socket.send_to(&data, broadcast_addr).await;
|
||||||
|
let _ = socket.send_to(&data, multicast_addr).await;
|
||||||
|
let _ = socket.send_to(&data, mdns_addr).await;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
_ = shutdown_rx1.recv() => {
|
_ = shutdown_rx1.recv() => {
|
||||||
@@ -185,15 +191,16 @@ impl DiscoveryService {
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
// 2. Create UDP listener socket
|
// 2. Create UDP listener socket with Multicast support
|
||||||
let listener_socket = match create_udp_listener(DEFAULT_DISCOVERY_PORT) {
|
let iface_ip = self.selected_iface.ip;
|
||||||
|
let listener_socket = match create_udp_listener(DEFAULT_DISCOVERY_PORT, iface_ip) {
|
||||||
Ok(s) => s,
|
Ok(s) => s,
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
warn!(
|
warn!(
|
||||||
"Could not bind listener on port {}: {}. Will try random port.",
|
"Could not bind listener on port {}: {}. Will try random port.",
|
||||||
DEFAULT_DISCOVERY_PORT, e
|
DEFAULT_DISCOVERY_PORT, e
|
||||||
);
|
);
|
||||||
create_udp_listener(0)?
|
create_udp_listener(0, iface_ip)?
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -285,13 +292,13 @@ impl DiscoveryService {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Helper to create a reusable broadcast socket on Windows / Linux / macOS
|
/// Helper to create a reusable broadcast & multicast socket on Windows / Linux / macOS
|
||||||
fn create_udp_listener(port: u16) -> std::io::Result<UdpSocket> {
|
fn create_udp_listener(port: u16, iface_ip: Ipv4Addr) -> std::io::Result<UdpSocket> {
|
||||||
let socket = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP))?;
|
let socket = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP))?;
|
||||||
socket.set_reuse_address(true)?;
|
socket.set_reuse_address(true)?;
|
||||||
|
|
||||||
#[cfg(not(windows))]
|
#[cfg(not(windows))]
|
||||||
socket.set_reuse_port(true)?;
|
let _ = socket.set_reuse_port(true);
|
||||||
|
|
||||||
socket.set_broadcast(true)?;
|
socket.set_broadcast(true)?;
|
||||||
socket.set_nonblocking(true)?;
|
socket.set_nonblocking(true)?;
|
||||||
@@ -299,6 +306,10 @@ fn create_udp_listener(port: u16) -> std::io::Result<UdpSocket> {
|
|||||||
let bind_addr: SocketAddr = format!("0.0.0.0:{}", port).parse().unwrap();
|
let bind_addr: SocketAddr = format!("0.0.0.0:{}", port).parse().unwrap();
|
||||||
socket.bind(&bind_addr.into())?;
|
socket.bind(&bind_addr.into())?;
|
||||||
|
|
||||||
|
// Join local SSDP/ZeroSend multicast group and mDNS group
|
||||||
|
let _ = socket.join_multicast_v4(&MULTICAST_IPV4_DISCOVERY, &iface_ip);
|
||||||
|
let _ = socket.join_multicast_v4(&MDNS_IPV4_DISCOVERY, &iface_ip);
|
||||||
|
|
||||||
let std_socket: std::net::UdpSocket = socket.into();
|
let std_socket: std::net::UdpSocket = socket.into();
|
||||||
UdpSocket::from_std(std_socket)
|
UdpSocket::from_std(std_socket)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ use zip::{CompressionMethod, ZipArchive, ZipWriter};
|
|||||||
// ─────────────────────────────────────────────────────────────────
|
// ─────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
/// Compress a whole folder into a ZIP archive with fast Deflate compression
|
/// Compress a whole folder into a ZIP archive with fast Deflate compression
|
||||||
|
#[allow(dead_code)]
|
||||||
pub fn compress_folder_to_zip(
|
pub fn compress_folder_to_zip(
|
||||||
src_dir: &Path,
|
src_dir: &Path,
|
||||||
zip_dest: &Path,
|
zip_dest: &Path,
|
||||||
|
|||||||
+229
-35
@@ -70,6 +70,9 @@ impl TransferClient {
|
|||||||
.timeout(Duration::from_secs(7200))
|
.timeout(Duration::from_secs(7200))
|
||||||
.connect_timeout(Duration::from_secs(6))
|
.connect_timeout(Duration::from_secs(6))
|
||||||
.tcp_nodelay(true)
|
.tcp_nodelay(true)
|
||||||
|
.tcp_keepalive(Some(Duration::from_secs(15)))
|
||||||
|
.pool_idle_timeout(Duration::from_secs(120))
|
||||||
|
.pool_max_idle_per_host(64)
|
||||||
.build()
|
.build()
|
||||||
.unwrap_or_default();
|
.unwrap_or_default();
|
||||||
|
|
||||||
@@ -93,12 +96,14 @@ impl TransferClient {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[allow(dead_code)]
|
||||||
pub async fn is_paused(&self, session_id: &str) -> bool {
|
pub async fn is_paused(&self, session_id: &str) -> bool {
|
||||||
let flags = self.pause_flags.lock().await;
|
let flags = self.pause_flags.lock().await;
|
||||||
flags.get(session_id).map(|f| f.load(Ordering::Relaxed)).unwrap_or(false)
|
flags.get(session_id).map(|f| f.load(Ordering::Relaxed)).unwrap_or(false)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Race-ping multiple candidate URLs and return the fastest responding one
|
/// Race-ping multiple candidate URLs and return the fastest responding one
|
||||||
|
#[allow(dead_code)]
|
||||||
pub async fn pick_fastest_url(&self, candidates: &[String]) -> Option<String> {
|
pub async fn pick_fastest_url(&self, candidates: &[String]) -> Option<String> {
|
||||||
if candidates.is_empty() { return None; }
|
if candidates.is_empty() { return None; }
|
||||||
if candidates.len() == 1 { return Some(candidates[0].clone()); }
|
if candidates.len() == 1 { return Some(candidates[0].clone()); }
|
||||||
@@ -273,17 +278,29 @@ impl TransferClient {
|
|||||||
|
|
||||||
// ── Stream each folder as tar.zst to the peer (zero disk footprint) ──
|
// ── Stream each folder as tar.zst to the peer (zero disk footprint) ──
|
||||||
for (folder_path, folder_name) in streaming_folders {
|
for (folder_path, folder_name) in streaming_folders {
|
||||||
let archive_name = format!("{}.tar.zst", folder_name);
|
let session_id = Uuid::new_v4().to_string();
|
||||||
|
let session_id_clone = session_id.clone();
|
||||||
|
let peer_display_name_clone = peer_display_name.to_string();
|
||||||
|
let folder_name_clone = folder_name.clone();
|
||||||
|
|
||||||
|
let folder_est_size = {
|
||||||
|
let p_c = folder_path.clone();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
super::archive::estimate_folder_size(&p_c)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap_or(0)
|
||||||
|
};
|
||||||
|
|
||||||
on_progress(TransferProgress {
|
on_progress(TransferProgress {
|
||||||
session_id: Uuid::new_v4().to_string(),
|
session_id: session_id.clone(),
|
||||||
peer_name: peer_display_name.to_string(),
|
peer_name: peer_display_name.to_string(),
|
||||||
is_incoming: false,
|
is_incoming: false,
|
||||||
current_file_index: 0,
|
current_file_index: 0,
|
||||||
total_files: paths.len(),
|
total_files: paths.len(),
|
||||||
current_file_name: format!("📦 Streaming folder '{}' (tar.zst)...", folder_name),
|
current_file_name: format!("📁 {}", folder_name),
|
||||||
bytes_transferred: 0,
|
bytes_transferred: 0,
|
||||||
total_bytes,
|
total_bytes: folder_est_size,
|
||||||
speed_bps: 0,
|
speed_bps: 0,
|
||||||
state: TransferState::InProgress,
|
state: TransferState::InProgress,
|
||||||
checksum_verified: false,
|
checksum_verified: false,
|
||||||
@@ -310,17 +327,20 @@ impl TransferClient {
|
|||||||
super::archive::stream_folder_as_tar_zst(folder_clone, tx, folder_zstd_level);
|
super::archive::stream_folder_as_tar_zst(folder_clone, tx, folder_zstd_level);
|
||||||
});
|
});
|
||||||
|
|
||||||
// Collect chunks from the channel into an HTTP body stream
|
|
||||||
let peer_url_clone = peer_url.to_string();
|
|
||||||
let archive_name_clone = archive_name.clone();
|
|
||||||
let session_placeholder = Uuid::new_v4().to_string();
|
|
||||||
let stream_url = format!(
|
let stream_url = format!(
|
||||||
"{}/api/transfer/folder_stream/{}",
|
"{}/api/transfer/folder_stream/{}",
|
||||||
peer_url, session_placeholder
|
peer_url, session_id
|
||||||
);
|
);
|
||||||
|
|
||||||
// Wrap sync receiver in async stream for reqwest using Arc<Mutex>
|
// Wrap sync receiver in async stream for reqwest using Arc<Mutex>
|
||||||
let rx_arc = std::sync::Arc::new(std::sync::Mutex::new(rx));
|
let rx_arc = std::sync::Arc::new(std::sync::Mutex::new(rx));
|
||||||
|
let on_prog_arc = std::sync::Arc::new(std::sync::Mutex::new(on_progress));
|
||||||
|
let on_prog_stream = on_prog_arc.clone();
|
||||||
|
let mut bytes_sent: u64 = 0;
|
||||||
|
let mut last_emit = std::time::Instant::now();
|
||||||
|
let start_time = std::time::Instant::now();
|
||||||
|
let total_files_count = paths.len();
|
||||||
|
|
||||||
let body_stream = async_stream::stream! {
|
let body_stream = async_stream::stream! {
|
||||||
loop {
|
loop {
|
||||||
let rx_clone = rx_arc.clone();
|
let rx_clone = rx_arc.clone();
|
||||||
@@ -328,22 +348,47 @@ impl TransferClient {
|
|||||||
rx_clone.lock().unwrap().recv()
|
rx_clone.lock().unwrap().recv()
|
||||||
}).await;
|
}).await;
|
||||||
match result {
|
match result {
|
||||||
Ok(Ok(Ok(chunk))) => yield Ok::<_, std::io::Error>(bytes::Bytes::from(chunk)),
|
Ok(Ok(Ok(chunk))) => {
|
||||||
|
let len = chunk.len() as u64;
|
||||||
|
bytes_sent += len;
|
||||||
|
if last_emit.elapsed() > std::time::Duration::from_millis(250) {
|
||||||
|
last_emit = std::time::Instant::now();
|
||||||
|
let elapsed = start_time.elapsed().as_secs_f64().max(0.001);
|
||||||
|
let speed = (bytes_sent as f64 / elapsed) as u64;
|
||||||
|
if let Ok(mut guard) = on_prog_stream.try_lock() {
|
||||||
|
(guard)(TransferProgress {
|
||||||
|
session_id: session_id_clone.clone(),
|
||||||
|
peer_name: peer_display_name_clone.clone(),
|
||||||
|
is_incoming: false,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: total_files_count,
|
||||||
|
current_file_name: format!("📁 {}", folder_name_clone),
|
||||||
|
bytes_transferred: if folder_est_size > 0 { bytes_sent.min(folder_est_size) } else { bytes_sent },
|
||||||
|
total_bytes: folder_est_size,
|
||||||
|
speed_bps: speed,
|
||||||
|
state: TransferState::InProgress,
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: bytes_sent,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
yield Ok::<_, std::io::Error>(bytes::Bytes::from(chunk));
|
||||||
|
}
|
||||||
Ok(Ok(Err(e))) => {
|
Ok(Ok(Err(e))) => {
|
||||||
yield Err(e);
|
yield Err(e);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
Ok(Err(_)) | Err(_) => break, // channel closed or task error
|
Ok(Err(_)) | Err(_) => break,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
let _ = peer_url_clone; // consumed in stream_url
|
|
||||||
|
|
||||||
let resp = self
|
let resp = self
|
||||||
.client
|
.client
|
||||||
.post(&stream_url)
|
.post(&stream_url)
|
||||||
.header("x-folder-name", &archive_name_clone)
|
.header("x-folder-name", &folder_name)
|
||||||
|
.header("x-total-size", folder_est_size.to_string())
|
||||||
|
.header("x-sender-name", &my_info.name)
|
||||||
.header("x-transfer-type", "tar.zst")
|
.header("x-transfer-type", "tar.zst")
|
||||||
.body(reqwest::Body::wrap_stream(body_stream))
|
.body(reqwest::Body::wrap_stream(body_stream))
|
||||||
.send()
|
.send()
|
||||||
@@ -351,21 +396,31 @@ impl TransferClient {
|
|||||||
|
|
||||||
compress_handle.await.ok();
|
compress_handle.await.ok();
|
||||||
|
|
||||||
|
on_progress = match std::sync::Arc::try_unwrap(on_prog_arc) {
|
||||||
|
Ok(mutex) => mutex.into_inner().unwrap(),
|
||||||
|
Err(arc) => {
|
||||||
|
let mutex = arc.lock().unwrap();
|
||||||
|
// If still referenced by stream task, we can't unwrap, but we can call through lock
|
||||||
|
drop(mutex);
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
match resp {
|
match resp {
|
||||||
Ok(r) if r.status().is_success() => {
|
Ok(r) if r.status().is_success() => {
|
||||||
on_progress(TransferProgress {
|
on_progress(TransferProgress {
|
||||||
session_id: session_placeholder.clone(),
|
session_id: session_id.clone(),
|
||||||
peer_name: peer_display_name.to_string(),
|
peer_name: peer_display_name.to_string(),
|
||||||
is_incoming: false,
|
is_incoming: false,
|
||||||
current_file_index: 0,
|
current_file_index: 0,
|
||||||
total_files: paths.len(),
|
total_files: paths.len(),
|
||||||
current_file_name: format!("✅ Folder '{}' delivered", folder_name),
|
current_file_name: format!("📁 {}", folder_name),
|
||||||
bytes_transferred: total_bytes,
|
bytes_transferred: folder_est_size,
|
||||||
total_bytes,
|
total_bytes: folder_est_size,
|
||||||
speed_bps: 0,
|
speed_bps: 0,
|
||||||
state: TransferState::Completed,
|
state: TransferState::Completed,
|
||||||
checksum_verified: false,
|
checksum_verified: true,
|
||||||
compressed_bytes: 0,
|
compressed_bytes: bytes_sent,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
Ok(r) => {
|
Ok(r) => {
|
||||||
@@ -375,7 +430,7 @@ impl TransferClient {
|
|||||||
));
|
));
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
return Err(format!("Folder stream network error: {}", e));
|
return Err(format!("Failed to stream folder: {}", e));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -679,9 +734,14 @@ impl TransferClient {
|
|||||||
};
|
};
|
||||||
|
|
||||||
if !is_ok {
|
if !is_ok {
|
||||||
if chunk_task.retries < 3 {
|
if chunk_task.retries < 8 {
|
||||||
chunk_task.retries += 1;
|
chunk_task.retries += 1;
|
||||||
tokio::time::sleep(Duration::from_millis(250)).await;
|
let backoff = (250 * (1 << chunk_task.retries.min(5))).min(4000);
|
||||||
|
info!(
|
||||||
|
"Zero-Drop Auto-Recovery: retrying chunk {} (attempt {}/8) after {}ms",
|
||||||
|
chunk_task.rel_path, chunk_task.retries, backoff
|
||||||
|
);
|
||||||
|
tokio::time::sleep(Duration::from_millis(backoff)).await;
|
||||||
tasks_c.lock().await.push_back(chunk_task);
|
tasks_c.lock().await.push_back(chunk_task);
|
||||||
continue;
|
continue;
|
||||||
} else {
|
} else {
|
||||||
@@ -866,37 +926,89 @@ impl TransferClient {
|
|||||||
.map_err(|e| format!("Failed to read file text: {}", e))
|
.map_err(|e| format!("Failed to read file text: {}", e))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Download a file or folder (zip) from a remote shared folder
|
/// Download a file or folder (zip) from a remote shared folder with live progress tracking
|
||||||
pub async fn download_shared_file(
|
pub async fn download_shared_file_with_progress<F>(
|
||||||
&self,
|
&self,
|
||||||
peer_url: &str,
|
peer_url: &str,
|
||||||
folder_id: &str,
|
folder_id: &str,
|
||||||
subpath: &str,
|
subpath: &str,
|
||||||
dest_path: &std::path::Path,
|
dest_path: &std::path::Path,
|
||||||
pin: Option<&str>,
|
pin: Option<&str>,
|
||||||
) -> Result<u64, String> {
|
peer_name: &str,
|
||||||
|
mut on_progress: F,
|
||||||
|
) -> Result<u64, String>
|
||||||
|
where
|
||||||
|
F: FnMut(TransferProgress) + Send + 'static,
|
||||||
|
{
|
||||||
let clean = subpath.trim_start_matches('/').trim_start_matches('\\');
|
let clean = subpath.trim_start_matches('/').trim_start_matches('\\');
|
||||||
let mut url = format!("{}/api/shared/download/{}/{}", peer_url, folder_id, clean);
|
let mut url = format!("{}/api/shared/download/{}/{}", peer_url, folder_id, clean);
|
||||||
if let Some(p) = pin.filter(|p| !p.trim().is_empty()) {
|
if let Some(p) = pin.filter(|p| !p.trim().is_empty()) {
|
||||||
url = format!("{}?pin={}", url, p);
|
url = format!("{}?pin={}", url, p);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let fname = dest_path.file_name().map(|n| n.to_string_lossy().into_owned()).unwrap_or_else(|| "downloaded".to_string());
|
||||||
|
let session_id = format!("shared-dl-{}", uuid::Uuid::new_v4());
|
||||||
|
|
||||||
|
on_progress(TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: peer_name.to_string(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: fname.clone(),
|
||||||
|
bytes_transferred: 0,
|
||||||
|
total_bytes: 0,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: TransferState::InProgress,
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
});
|
||||||
|
|
||||||
let resp = self
|
let resp = self
|
||||||
.client
|
.client
|
||||||
.get(&url)
|
.get(&url)
|
||||||
.send()
|
.send()
|
||||||
.await
|
.await
|
||||||
.map_err(|e| format!("Network error: {}", e))?;
|
.map_err(|e| {
|
||||||
|
let err_msg = format!("Network error: {}", e);
|
||||||
|
on_progress(TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: peer_name.to_string(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: fname.clone(),
|
||||||
|
bytes_transferred: 0,
|
||||||
|
total_bytes: 0,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: TransferState::Failed(err_msg.clone()),
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
});
|
||||||
|
err_msg
|
||||||
|
})?;
|
||||||
|
|
||||||
if !resp.status().is_success() {
|
if !resp.status().is_success() {
|
||||||
return Err(format!("Download failed: HTTP {}", resp.status()));
|
let status = resp.status();
|
||||||
|
let err_msg = format!("Download failed: HTTP {}", status);
|
||||||
|
on_progress(TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: peer_name.to_string(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: fname.clone(),
|
||||||
|
bytes_transferred: 0,
|
||||||
|
total_bytes: 0,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: TransferState::Failed(err_msg.clone()),
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
});
|
||||||
|
return Err(err_msg);
|
||||||
}
|
}
|
||||||
|
|
||||||
let bytes = resp
|
let total_size = resp.content_length().unwrap_or(0);
|
||||||
.bytes()
|
|
||||||
.await
|
|
||||||
.map_err(|e| format!("Failed to read stream: {}", e))?;
|
|
||||||
let len = bytes.len() as u64;
|
|
||||||
|
|
||||||
if let Some(parent) = dest_path.parent() {
|
if let Some(parent) = dest_path.parent() {
|
||||||
let _ = tokio::fs::create_dir_all(parent).await;
|
let _ = tokio::fs::create_dir_all(parent).await;
|
||||||
@@ -906,14 +1018,94 @@ impl TransferClient {
|
|||||||
.await
|
.await
|
||||||
.map_err(|e| format!("Could not create file: {}", e))?;
|
.map_err(|e| format!("Could not create file: {}", e))?;
|
||||||
|
|
||||||
tokio::io::AsyncWriteExt::write_all(&mut file, &bytes)
|
use futures_util::StreamExt;
|
||||||
|
let mut stream = resp.bytes_stream();
|
||||||
|
let mut bytes_received: u64 = 0;
|
||||||
|
let start_time = std::time::Instant::now();
|
||||||
|
let mut last_emit = std::time::Instant::now();
|
||||||
|
|
||||||
|
while let Some(chunk_res) = stream.next().await {
|
||||||
|
match chunk_res {
|
||||||
|
Ok(chunk) => {
|
||||||
|
let len = chunk.len() as u64;
|
||||||
|
bytes_received += len;
|
||||||
|
tokio::io::AsyncWriteExt::write_all(&mut file, &chunk)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| format!("Could not write file: {}", e))?;
|
.map_err(|e| format!("Could not write file: {}", e))?;
|
||||||
|
|
||||||
|
if last_emit.elapsed() > std::time::Duration::from_millis(250) {
|
||||||
|
last_emit = std::time::Instant::now();
|
||||||
|
let elapsed = start_time.elapsed().as_secs_f64().max(0.001);
|
||||||
|
let speed = (bytes_received as f64 / elapsed) as u64;
|
||||||
|
on_progress(TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: peer_name.to_string(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: fname.clone(),
|
||||||
|
bytes_transferred: bytes_received,
|
||||||
|
total_bytes: if total_size > 0 { total_size } else { bytes_received },
|
||||||
|
speed_bps: speed,
|
||||||
|
state: TransferState::InProgress,
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
let err_msg = format!("Stream error: {}", e);
|
||||||
|
on_progress(TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: peer_name.to_string(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: fname.clone(),
|
||||||
|
bytes_transferred: bytes_received,
|
||||||
|
total_bytes: total_size,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: TransferState::Failed(err_msg.clone()),
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
});
|
||||||
|
return Err(err_msg);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
tokio::io::AsyncWriteExt::flush(&mut file)
|
tokio::io::AsyncWriteExt::flush(&mut file)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| format!("Could not flush file: {}", e))?;
|
.map_err(|e| format!("Could not flush file: {}", e))?;
|
||||||
|
|
||||||
Ok(len)
|
on_progress(TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: peer_name.to_string(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 1,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: fname.clone(),
|
||||||
|
bytes_transferred: bytes_received,
|
||||||
|
total_bytes: if total_size > 0 { total_size } else { bytes_received },
|
||||||
|
speed_bps: 0,
|
||||||
|
state: TransferState::Completed,
|
||||||
|
checksum_verified: true,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
});
|
||||||
|
|
||||||
|
Ok(bytes_received)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Download a file or folder (zip) from a remote shared folder
|
||||||
|
pub async fn download_shared_file(
|
||||||
|
&self,
|
||||||
|
peer_url: &str,
|
||||||
|
folder_id: &str,
|
||||||
|
subpath: &str,
|
||||||
|
dest_path: &std::path::Path,
|
||||||
|
pin: Option<&str>,
|
||||||
|
) -> Result<u64, String> {
|
||||||
|
self.download_shared_file_with_progress(peer_url, folder_id, subpath, dest_path, pin, "Shared", |_| {}).await
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Upload a file directly to a friend's shared folder (if Read & Write allowed)
|
/// Upload a file directly to a friend's shared folder (if Read & Write allowed)
|
||||||
@@ -1032,6 +1224,7 @@ impl TransferClient {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Fetch active watch-party state from host peer
|
/// Fetch active watch-party state from host peer
|
||||||
|
#[allow(dead_code)]
|
||||||
pub async fn get_watch_party(
|
pub async fn get_watch_party(
|
||||||
&self,
|
&self,
|
||||||
peer_url: &str,
|
peer_url: &str,
|
||||||
@@ -1054,6 +1247,7 @@ impl TransferClient {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Send watch-party sync event to host peer
|
/// Send watch-party sync event to host peer
|
||||||
|
#[allow(dead_code)]
|
||||||
pub async fn send_watch_party_event(
|
pub async fn send_watch_party_event(
|
||||||
&self,
|
&self,
|
||||||
peer_url: &str,
|
peer_url: &str,
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ use tracing::info;
|
|||||||
|
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
pub struct SyncStatusUpdate {
|
pub struct SyncStatusUpdate {
|
||||||
|
#[allow(dead_code)]
|
||||||
pub pair_id: String,
|
pub pair_id: String,
|
||||||
pub message: String,
|
pub message: String,
|
||||||
pub is_error: bool,
|
pub is_error: bool,
|
||||||
@@ -78,6 +79,7 @@ impl FolderSyncEngine {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[allow(dead_code)]
|
||||||
pub fn pairs_handle(&self) -> Arc<RwLock<Vec<FolderSyncPairConfig>>> {
|
pub fn pairs_handle(&self) -> Arc<RwLock<Vec<FolderSyncPairConfig>>> {
|
||||||
self.pairs.clone()
|
self.pairs.clone()
|
||||||
}
|
}
|
||||||
|
|||||||
+338
-27
@@ -222,8 +222,42 @@ impl TransferServer {
|
|||||||
.with_state(self.state.clone());
|
.with_state(self.state.clone());
|
||||||
|
|
||||||
let addr = SocketAddr::from(([0, 0, 0, 0], self.port));
|
let addr = SocketAddr::from(([0, 0, 0, 0], self.port));
|
||||||
let listener = tokio::net::TcpListener::bind(addr).await?;
|
let socket = socket2::Socket::new(
|
||||||
info!("ZeroSend HTTP server listening on {}", addr);
|
socket2::Domain::IPV4,
|
||||||
|
socket2::Type::STREAM,
|
||||||
|
Some(socket2::Protocol::TCP),
|
||||||
|
)?;
|
||||||
|
socket.set_reuse_address(true)?;
|
||||||
|
#[cfg(not(windows))]
|
||||||
|
let _ = socket.set_reuse_port(true);
|
||||||
|
let _ = socket.set_recv_buffer_size(4 * 1024 * 1024); // 4MB TCP Receive Window
|
||||||
|
let _ = socket.set_send_buffer_size(4 * 1024 * 1024); // 4MB TCP Send Window
|
||||||
|
let _ = socket.set_nodelay(true);
|
||||||
|
socket.set_nonblocking(true)?;
|
||||||
|
socket.bind(&addr.into())?;
|
||||||
|
socket.listen(1024)?;
|
||||||
|
let std_listener: std::net::TcpListener = socket.into();
|
||||||
|
let listener = tokio::net::TcpListener::from_std(std_listener)?;
|
||||||
|
info!("ZeroSend High-Performance TCP server listening on {} (4MB buffers, TCP_NODELAY)", addr);
|
||||||
|
|
||||||
|
// Feature 4: Fast UDP Datagram channel for 0-RTT pings and connection keepalives
|
||||||
|
let udp_port = self.port;
|
||||||
|
let mut udp_shutdown = shutdown_tx.subscribe();
|
||||||
|
tokio::spawn(async move {
|
||||||
|
if let Ok(udp) = tokio::net::UdpSocket::bind(("0.0.0.0", udp_port)).await {
|
||||||
|
let mut buf = [0u8; 128];
|
||||||
|
loop {
|
||||||
|
tokio::select! {
|
||||||
|
Ok((n, src)) = udp.recv_from(&mut buf) => {
|
||||||
|
if n >= 4 && &buf[..4] == b"PING" {
|
||||||
|
let _ = udp.send_to(b"PONG", src).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
_ = udp_shutdown.recv() => break,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
let mut shutdown_rx = shutdown_tx.subscribe();
|
let mut shutdown_rx = shutdown_tx.subscribe();
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
@@ -660,6 +694,9 @@ async fn handle_finish_transfer(
|
|||||||
let part_path = state.download_dir.join(format!("{}.zerosend_part", safe_rel_path.to_string_lossy()));
|
let part_path = state.download_dir.join(format!("{}.zerosend_part", safe_rel_path.to_string_lossy()));
|
||||||
|
|
||||||
if part_path.exists() {
|
if part_path.exists() {
|
||||||
|
if target_path.exists() {
|
||||||
|
let _ = tokio::fs::remove_file(&target_path).await;
|
||||||
|
}
|
||||||
let _ = tokio::fs::rename(&part_path, &target_path).await;
|
let _ = tokio::fs::rename(&part_path, &target_path).await;
|
||||||
if let Some(expected_hash) = &f.sha256 {
|
if let Some(expected_hash) = &f.sha256 {
|
||||||
if let Ok(actual_hash) = super::client::compute_file_sha256(&target_path) {
|
if let Ok(actual_hash) = super::client::compute_file_sha256(&target_path) {
|
||||||
@@ -695,36 +732,116 @@ async fn handle_finish_transfer(
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Streaming tar.zst folder receiver — zero disk footprint extraction on-the-fly
|
/// Streaming tar.zst folder receiver — zero disk footprint extraction on-the-fly with live progress
|
||||||
async fn handle_folder_stream(
|
async fn handle_folder_stream(
|
||||||
State(state): State<Arc<ServerState>>,
|
State(state): State<Arc<ServerState>>,
|
||||||
Path(_session_id): Path<String>,
|
Path(session_id): Path<String>,
|
||||||
headers: axum::http::HeaderMap,
|
headers: axum::http::HeaderMap,
|
||||||
body: axum::body::Body,
|
body: axum::body::Body,
|
||||||
) -> impl IntoResponse {
|
) -> impl IntoResponse {
|
||||||
let folder_name = headers
|
let raw_folder_name = headers
|
||||||
.get("x-folder-name")
|
.get("x-folder-name")
|
||||||
.and_then(|v| v.to_str().ok())
|
.and_then(|v| v.to_str().ok())
|
||||||
.unwrap_or("received_folder")
|
.unwrap_or("folder")
|
||||||
|
.trim_end_matches(".tar.zst")
|
||||||
.to_string();
|
.to_string();
|
||||||
|
|
||||||
info!("Receiving streaming folder: {}", folder_name);
|
let total_size: u64 = headers
|
||||||
|
.get("x-total-size")
|
||||||
|
.and_then(|v| v.to_str().ok())
|
||||||
|
.and_then(|v| v.parse().ok())
|
||||||
|
.unwrap_or(0);
|
||||||
|
|
||||||
// Collect the full body into a byte buffer
|
let sender_name = headers
|
||||||
// (for large folders this streams via axum body chunks — no intermediate file)
|
.get("x-sender-name")
|
||||||
use axum::body::to_bytes;
|
.and_then(|v| v.to_str().ok())
|
||||||
let data = match to_bytes(body, usize::MAX).await {
|
.unwrap_or("Peer")
|
||||||
Ok(b) => b,
|
.to_string();
|
||||||
|
|
||||||
|
info!("Receiving streaming folder: {} (~{} bytes) from {}", raw_folder_name, total_size, sender_name);
|
||||||
|
|
||||||
|
// Initial progress event
|
||||||
|
let initial_progress = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: sender_name.clone(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: format!("📁 {}", raw_folder_name),
|
||||||
|
bytes_transferred: 0,
|
||||||
|
total_bytes: total_size,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: super::protocol::TransferState::InProgress,
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
};
|
||||||
|
state.active_transfers.write().await.insert(session_id.clone(), initial_progress.clone());
|
||||||
|
let _ = state.transfer_events.send(initial_progress);
|
||||||
|
|
||||||
|
use futures_util::StreamExt;
|
||||||
|
let mut stream = body.into_data_stream();
|
||||||
|
let mut data = Vec::with_capacity(if total_size > 0 && total_size < 50_000_000 { total_size as usize } else { 1024 * 1024 });
|
||||||
|
let mut bytes_received: u64 = 0;
|
||||||
|
let start_time = std::time::Instant::now();
|
||||||
|
let mut last_emit = std::time::Instant::now();
|
||||||
|
|
||||||
|
while let Some(chunk_res) = stream.next().await {
|
||||||
|
match chunk_res {
|
||||||
|
Ok(chunk) => {
|
||||||
|
let len = chunk.len() as u64;
|
||||||
|
bytes_received += len;
|
||||||
|
data.extend_from_slice(&chunk);
|
||||||
|
|
||||||
|
if last_emit.elapsed() > std::time::Duration::from_millis(250) {
|
||||||
|
last_emit = std::time::Instant::now();
|
||||||
|
let elapsed = start_time.elapsed().as_secs_f64().max(0.001);
|
||||||
|
let speed = (bytes_received as f64 / elapsed) as u64;
|
||||||
|
let p = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: sender_name.clone(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: format!("📁 {}", raw_folder_name),
|
||||||
|
bytes_transferred: if total_size > 0 { bytes_received.min(total_size) } else { bytes_received },
|
||||||
|
total_bytes: if total_size > 0 { total_size } else { bytes_received },
|
||||||
|
speed_bps: speed,
|
||||||
|
state: super::protocol::TransferState::InProgress,
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: bytes_received,
|
||||||
|
};
|
||||||
|
state.active_transfers.write().await.insert(session_id.clone(), p.clone());
|
||||||
|
let _ = state.transfer_events.send(p);
|
||||||
|
}
|
||||||
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
error!("Failed to read folder stream body: {}", e);
|
error!("Failed to read folder stream chunk: {}", e);
|
||||||
|
let fail_p = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: sender_name.clone(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: format!("📁 {}", raw_folder_name),
|
||||||
|
bytes_transferred: bytes_received,
|
||||||
|
total_bytes: total_size,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: super::protocol::TransferState::Failed(e.to_string()),
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: bytes_received,
|
||||||
|
};
|
||||||
|
state.active_transfers.write().await.insert(session_id.clone(), fail_p.clone());
|
||||||
|
let _ = state.transfer_events.send(fail_p);
|
||||||
return (
|
return (
|
||||||
StatusCode::INTERNAL_SERVER_ERROR,
|
StatusCode::INTERNAL_SERVER_ERROR,
|
||||||
Json(serde_json::json!({"error": e.to_string()})),
|
Json(serde_json::json!({"error": e.to_string()})),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
};
|
}
|
||||||
|
}
|
||||||
|
|
||||||
let dest_dir = state.download_dir.clone();
|
let dest_dir = state.download_dir.clone();
|
||||||
|
let folder_name_clone = raw_folder_name.clone();
|
||||||
let result = tokio::task::spawn_blocking(move || {
|
let result = tokio::task::spawn_blocking(move || {
|
||||||
super::archive::extract_tar_zst_stream(std::io::Cursor::new(data), &dest_dir)
|
super::archive::extract_tar_zst_stream(std::io::Cursor::new(data), &dest_dir)
|
||||||
})
|
})
|
||||||
@@ -732,30 +849,47 @@ async fn handle_folder_stream(
|
|||||||
|
|
||||||
match result {
|
match result {
|
||||||
Ok(Ok(())) => {
|
Ok(Ok(())) => {
|
||||||
info!("Folder '{}' extracted successfully", folder_name);
|
info!("Folder '{}' extracted successfully", folder_name_clone);
|
||||||
// Emit a completed progress event
|
let final_bytes = if total_size > 0 { total_size } else { bytes_received };
|
||||||
let progress = super::protocol::TransferProgress {
|
let progress = super::protocol::TransferProgress {
|
||||||
session_id: _session_id.clone(),
|
session_id: session_id.clone(),
|
||||||
peer_name: String::new(),
|
peer_name: sender_name.clone(),
|
||||||
is_incoming: true,
|
is_incoming: true,
|
||||||
current_file_index: 1,
|
current_file_index: 1,
|
||||||
total_files: 1,
|
total_files: 1,
|
||||||
current_file_name: folder_name.clone(),
|
current_file_name: format!("📁 {}", folder_name_clone),
|
||||||
bytes_transferred: 0,
|
bytes_transferred: final_bytes,
|
||||||
total_bytes: 0,
|
total_bytes: final_bytes,
|
||||||
speed_bps: 0,
|
speed_bps: 0,
|
||||||
state: super::protocol::TransferState::Completed,
|
state: super::protocol::TransferState::Completed,
|
||||||
checksum_verified: false,
|
checksum_verified: true,
|
||||||
compressed_bytes: 0,
|
compressed_bytes: bytes_received,
|
||||||
};
|
};
|
||||||
|
state.active_transfers.write().await.insert(session_id.clone(), progress.clone());
|
||||||
let _ = state.transfer_events.send(progress);
|
let _ = state.transfer_events.send(progress);
|
||||||
(
|
(
|
||||||
StatusCode::OK,
|
StatusCode::OK,
|
||||||
Json(serde_json::json!({"status": "extracted", "folder": folder_name})),
|
Json(serde_json::json!({"status": "extracted", "folder": folder_name_clone})),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
Ok(Err(e)) => {
|
Ok(Err(e)) => {
|
||||||
error!("Extraction error for '{}': {}", folder_name, e);
|
error!("Extraction error for '{}': {}", folder_name_clone, e);
|
||||||
|
let fail_p = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: sender_name.clone(),
|
||||||
|
is_incoming: true,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: format!("📁 {}", folder_name_clone),
|
||||||
|
bytes_transferred: bytes_received,
|
||||||
|
total_bytes: total_size,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: super::protocol::TransferState::Failed(e.to_string()),
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: bytes_received,
|
||||||
|
};
|
||||||
|
state.active_transfers.write().await.insert(session_id.clone(), fail_p.clone());
|
||||||
|
let _ = state.transfer_events.send(fail_p);
|
||||||
(
|
(
|
||||||
StatusCode::INTERNAL_SERVER_ERROR,
|
StatusCode::INTERNAL_SERVER_ERROR,
|
||||||
Json(serde_json::json!({"error": e.to_string()})),
|
Json(serde_json::json!({"error": e.to_string()})),
|
||||||
@@ -1482,7 +1616,96 @@ async fn handle_shared_download(
|
|||||||
let file_size = file.metadata().await.map(|m| m.len()).unwrap_or(0);
|
let file_size = file.metadata().await.map(|m| m.len()).unwrap_or(0);
|
||||||
let mime_type = super::protocol::get_mime_type(&file_name);
|
let mime_type = super::protocol::get_mime_type(&file_name);
|
||||||
let reader = tokio_util::io::ReaderStream::new(file);
|
let reader = tokio_util::io::ReaderStream::new(file);
|
||||||
let body = Body::from_stream(reader);
|
|
||||||
|
let session_id = format!("shared-dl-{}", Uuid::new_v4());
|
||||||
|
let peer_label = format!("Общая папка / {}", folder.name);
|
||||||
|
|
||||||
|
let initial_p = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: peer_label.clone(),
|
||||||
|
is_incoming: false,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: file_name.clone(),
|
||||||
|
bytes_transferred: 0,
|
||||||
|
total_bytes: file_size,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: super::protocol::TransferState::InProgress,
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
};
|
||||||
|
state.active_transfers.write().await.insert(session_id.clone(), initial_p.clone());
|
||||||
|
let _ = state.transfer_events.send(initial_p);
|
||||||
|
|
||||||
|
let state_clone = state.clone();
|
||||||
|
let session_id_clone = session_id.clone();
|
||||||
|
let peer_label_clone = peer_label.clone();
|
||||||
|
let file_name_clone = file_name.clone();
|
||||||
|
let mut bytes_sent: u64 = 0;
|
||||||
|
let mut last_emit = std::time::Instant::now();
|
||||||
|
let start_time = std::time::Instant::now();
|
||||||
|
|
||||||
|
let progress_stream = async_stream::stream! {
|
||||||
|
use futures_util::StreamExt;
|
||||||
|
let mut s = reader;
|
||||||
|
while let Some(chunk_res) = s.next().await {
|
||||||
|
match chunk_res {
|
||||||
|
Ok(chunk) => {
|
||||||
|
let len = chunk.len() as u64;
|
||||||
|
bytes_sent += len;
|
||||||
|
if last_emit.elapsed() > std::time::Duration::from_millis(250) {
|
||||||
|
last_emit = std::time::Instant::now();
|
||||||
|
let elapsed = start_time.elapsed().as_secs_f64().max(0.001);
|
||||||
|
let speed = (bytes_sent as f64 / elapsed) as u64;
|
||||||
|
let p = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id_clone.clone(),
|
||||||
|
peer_name: peer_label_clone.clone(),
|
||||||
|
is_incoming: false,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: file_name_clone.clone(),
|
||||||
|
bytes_transferred: bytes_sent.min(file_size),
|
||||||
|
total_bytes: file_size,
|
||||||
|
speed_bps: speed,
|
||||||
|
state: super::protocol::TransferState::InProgress,
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
};
|
||||||
|
if let Ok(mut active) = state_clone.active_transfers.try_write() {
|
||||||
|
active.insert(session_id_clone.clone(), p.clone());
|
||||||
|
}
|
||||||
|
let _ = state_clone.transfer_events.send(p);
|
||||||
|
}
|
||||||
|
yield Ok::<_, std::io::Error>(chunk);
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
yield Err(e);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let p_done = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id_clone.clone(),
|
||||||
|
peer_name: peer_label_clone.clone(),
|
||||||
|
is_incoming: false,
|
||||||
|
current_file_index: 1,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: file_name_clone.clone(),
|
||||||
|
bytes_transferred: file_size,
|
||||||
|
total_bytes: file_size,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: super::protocol::TransferState::Completed,
|
||||||
|
checksum_verified: true,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
};
|
||||||
|
if let Ok(mut active) = state_clone.active_transfers.try_write() {
|
||||||
|
active.insert(session_id_clone.clone(), p_done.clone());
|
||||||
|
}
|
||||||
|
let _ = state_clone.transfer_events.send(p_done);
|
||||||
|
};
|
||||||
|
|
||||||
|
let body = Body::from_stream(progress_stream);
|
||||||
|
|
||||||
Response::builder()
|
Response::builder()
|
||||||
.status(StatusCode::OK)
|
.status(StatusCode::OK)
|
||||||
@@ -1538,9 +1761,97 @@ async fn handle_shared_download(
|
|||||||
};
|
};
|
||||||
let file_size = zip_file.metadata().await.map(|m| m.len()).unwrap_or(0);
|
let file_size = zip_file.metadata().await.map(|m| m.len()).unwrap_or(0);
|
||||||
let reader = tokio_util::io::ReaderStream::new(zip_file);
|
let reader = tokio_util::io::ReaderStream::new(zip_file);
|
||||||
let body = Body::from_stream(reader);
|
|
||||||
|
|
||||||
let out_zip_name = format!("{}.zip", file_name);
|
let out_zip_name = format!("{}.zip", file_name);
|
||||||
|
let session_id = format!("shared-dl-{}", Uuid::new_v4());
|
||||||
|
let peer_label = format!("Общая папка / {}", folder.name);
|
||||||
|
|
||||||
|
let initial_p = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id.clone(),
|
||||||
|
peer_name: peer_label.clone(),
|
||||||
|
is_incoming: false,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: out_zip_name.clone(),
|
||||||
|
bytes_transferred: 0,
|
||||||
|
total_bytes: file_size,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: super::protocol::TransferState::InProgress,
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
};
|
||||||
|
state.active_transfers.write().await.insert(session_id.clone(), initial_p.clone());
|
||||||
|
let _ = state.transfer_events.send(initial_p);
|
||||||
|
|
||||||
|
let state_clone = state.clone();
|
||||||
|
let session_id_clone = session_id.clone();
|
||||||
|
let peer_label_clone = peer_label.clone();
|
||||||
|
let out_zip_name_clone = out_zip_name.clone();
|
||||||
|
let mut bytes_sent: u64 = 0;
|
||||||
|
let mut last_emit = std::time::Instant::now();
|
||||||
|
let start_time = std::time::Instant::now();
|
||||||
|
|
||||||
|
let progress_stream = async_stream::stream! {
|
||||||
|
use futures_util::StreamExt;
|
||||||
|
let mut s = reader;
|
||||||
|
while let Some(chunk_res) = s.next().await {
|
||||||
|
match chunk_res {
|
||||||
|
Ok(chunk) => {
|
||||||
|
let len = chunk.len() as u64;
|
||||||
|
bytes_sent += len;
|
||||||
|
if last_emit.elapsed() > std::time::Duration::from_millis(250) {
|
||||||
|
last_emit = std::time::Instant::now();
|
||||||
|
let elapsed = start_time.elapsed().as_secs_f64().max(0.001);
|
||||||
|
let speed = (bytes_sent as f64 / elapsed) as u64;
|
||||||
|
let p = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id_clone.clone(),
|
||||||
|
peer_name: peer_label_clone.clone(),
|
||||||
|
is_incoming: false,
|
||||||
|
current_file_index: 0,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: out_zip_name_clone.clone(),
|
||||||
|
bytes_transferred: bytes_sent.min(file_size),
|
||||||
|
total_bytes: file_size,
|
||||||
|
speed_bps: speed,
|
||||||
|
state: super::protocol::TransferState::InProgress,
|
||||||
|
checksum_verified: false,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
};
|
||||||
|
if let Ok(mut active) = state_clone.active_transfers.try_write() {
|
||||||
|
active.insert(session_id_clone.clone(), p.clone());
|
||||||
|
}
|
||||||
|
let _ = state_clone.transfer_events.send(p);
|
||||||
|
}
|
||||||
|
yield Ok::<_, std::io::Error>(chunk);
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
yield Err(e);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let p_done = super::protocol::TransferProgress {
|
||||||
|
session_id: session_id_clone.clone(),
|
||||||
|
peer_name: peer_label_clone.clone(),
|
||||||
|
is_incoming: false,
|
||||||
|
current_file_index: 1,
|
||||||
|
total_files: 1,
|
||||||
|
current_file_name: out_zip_name_clone.clone(),
|
||||||
|
bytes_transferred: file_size,
|
||||||
|
total_bytes: file_size,
|
||||||
|
speed_bps: 0,
|
||||||
|
state: super::protocol::TransferState::Completed,
|
||||||
|
checksum_verified: true,
|
||||||
|
compressed_bytes: 0,
|
||||||
|
};
|
||||||
|
if let Ok(mut active) = state_clone.active_transfers.try_write() {
|
||||||
|
active.insert(session_id_clone.clone(), p_done.clone());
|
||||||
|
}
|
||||||
|
let _ = state_clone.transfer_events.send(p_done);
|
||||||
|
};
|
||||||
|
|
||||||
|
let body = Body::from_stream(progress_stream);
|
||||||
|
|
||||||
Response::builder()
|
Response::builder()
|
||||||
.status(StatusCode::OK)
|
.status(StatusCode::OK)
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ use tracing::{info, warn};
|
|||||||
|
|
||||||
pub struct WebDavState {
|
pub struct WebDavState {
|
||||||
pub shared_folders: Arc<RwLock<Vec<SharedFolderConfig>>>,
|
pub shared_folders: Arc<RwLock<Vec<SharedFolderConfig>>>,
|
||||||
|
#[allow(dead_code)]
|
||||||
pub port: u16,
|
pub port: u16,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -389,6 +390,7 @@ mod urlencoding {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Helper function to mount Windows Network Drive via `net use`
|
/// Helper function to mount Windows Network Drive via `net use`
|
||||||
|
#[allow(dead_code)]
|
||||||
pub fn mount_webdav_drive(drive_letter: char, port: u16) -> Result<String, String> {
|
pub fn mount_webdav_drive(drive_letter: char, port: u16) -> Result<String, String> {
|
||||||
let drive = format!("{}:", drive_letter);
|
let drive = format!("{}:", drive_letter);
|
||||||
let url = format!("http://127.0.0.1:{}/webdav", port);
|
let url = format!("http://127.0.0.1:{}/webdav", port);
|
||||||
@@ -412,6 +414,7 @@ pub fn mount_webdav_drive(drive_letter: char, port: u16) -> Result<String, Strin
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Helper function to unmount Windows Network Drive
|
/// Helper function to unmount Windows Network Drive
|
||||||
|
#[allow(dead_code)]
|
||||||
pub fn unmount_webdav_drive(drive_letter: char) -> Result<(), String> {
|
pub fn unmount_webdav_drive(drive_letter: char) -> Result<(), String> {
|
||||||
let drive = format!("{}:", drive_letter);
|
let drive = format!("{}:", drive_letter);
|
||||||
let output = std::process::Command::new("net")
|
let output = std::process::Command::new("net")
|
||||||
|
|||||||
+414
-119
File diff suppressed because it is too large
Load Diff
@@ -57,6 +57,7 @@ impl Default for VimViewerModal {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl VimViewerModal {
|
impl VimViewerModal {
|
||||||
|
#[allow(dead_code)]
|
||||||
pub fn open(
|
pub fn open(
|
||||||
&mut self,
|
&mut self,
|
||||||
file_name: String,
|
file_name: String,
|
||||||
|
|||||||
+13
-5
@@ -79,11 +79,19 @@ pub fn parse_version(v: &str) -> (u32, u32, u32, u32, bool, u32) {
|
|||||||
pl.contains("beta") || pl.contains("alpha") || pl.contains("rc")
|
pl.contains("beta") || pl.contains("alpha") || pl.contains("rc")
|
||||||
}).unwrap_or(false);
|
}).unwrap_or(false);
|
||||||
|
|
||||||
// Extract pre-release number, e.g. "beta.1" -> 1, "beta" -> 0
|
// Extract pre-release number, e.g. "beta.1.1" -> 1001, "beta.1" -> 1000, "beta" -> 0
|
||||||
let pre_num = pre_part
|
let pre_num = if let Some(p) = pre_part {
|
||||||
.and_then(|p| p.split('.').last())
|
let parts: Vec<&str> = p.split('.').collect();
|
||||||
.and_then(|s| s.parse::<u32>().ok())
|
let p1 = parts.get(1).and_then(|s| s.parse::<u32>().ok()).unwrap_or(0);
|
||||||
.unwrap_or(0);
|
let p2 = parts.get(2).and_then(|s| s.parse::<u32>().ok()).unwrap_or(0);
|
||||||
|
if p1 == 0 && p2 == 0 {
|
||||||
|
parts.last().and_then(|s| s.parse::<u32>().ok()).unwrap_or(0)
|
||||||
|
} else {
|
||||||
|
p1 * 1000 + p2
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
0
|
||||||
|
};
|
||||||
|
|
||||||
(major, minor, patch, subpatch, is_beta, pre_num)
|
(major, minor, patch, subpatch, is_beta, pre_num)
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user