Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c3f26b77e8 | ||
|
|
9957c9b70f | ||
|
|
d4f52e3ec6 | ||
|
|
3ca52e2649 | ||
|
|
90fda425bd | ||
|
|
124d9c2dd4 | ||
|
|
e019ddccc9 | ||
|
|
61f2545588 | ||
|
|
b751d82748 | ||
|
|
c27ac0f41c | ||
|
|
f65b15452f | ||
|
|
686eb2fb6f | ||
|
|
0867b85018 | ||
|
|
a6a32bb23f |
Generated
+118
-3
@@ -440,6 +440,28 @@ dependencies = [
|
||||
"windows-sys 0.61.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "async-stream"
|
||||
version = "0.3.6"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "0b5a71a6f37880a80d1d7f19efd781e4b5de42c88f0722cc13bcb6cc2cfe8476"
|
||||
dependencies = [
|
||||
"async-stream-impl",
|
||||
"futures-core",
|
||||
"pin-project-lite",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "async-stream-impl"
|
||||
version = "0.3.6"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "c7c24de15d275a1ecfd47a380fb4d5ec9bfe0933f309ed5e705b775596a3574d"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
"syn 2.0.119",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "async-task"
|
||||
version = "4.7.1"
|
||||
@@ -632,6 +654,19 @@ dependencies = [
|
||||
"serde_core",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "blake3"
|
||||
version = "1.8.7"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "6d9e454fc11f76977dc803893aff6304ed33d6a26efae8696573bea74baa27ae"
|
||||
dependencies = [
|
||||
"arrayvec",
|
||||
"cc",
|
||||
"cfg-if",
|
||||
"constant_time_eq",
|
||||
"cpufeatures 0.3.0",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "block"
|
||||
version = "0.1.6"
|
||||
@@ -1012,6 +1047,12 @@ dependencies = [
|
||||
"windows-sys 0.59.0",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "constant_time_eq"
|
||||
version = "0.4.2"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3d52eff69cd5e647efe296129160853a42795992097e8af39800e1060caeea9b"
|
||||
|
||||
[[package]]
|
||||
name = "core-foundation"
|
||||
version = "0.9.4"
|
||||
@@ -1071,6 +1112,15 @@ dependencies = [
|
||||
"libc",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "cpufeatures"
|
||||
version = "0.3.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "8b2a41393f66f16b0823bb79094d54ac5fbd34ab292ddafb9a0456ac9f87d201"
|
||||
dependencies = [
|
||||
"libc",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "crc32fast"
|
||||
version = "1.5.1"
|
||||
@@ -1544,6 +1594,16 @@ dependencies = [
|
||||
"rustc_version",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "filetime"
|
||||
version = "0.2.29"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "5c287a33c7f0a620c38e641e7f60827713987b3c0f26e8ddc9462cc69cf75759"
|
||||
dependencies = [
|
||||
"cfg-if",
|
||||
"libc",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "find-msvc-tools"
|
||||
version = "0.1.11"
|
||||
@@ -4301,7 +4361,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a978451301f4db1d02937a4ab3ccce137717b81826e79b7d49ffe3244a13c3b8"
|
||||
dependencies = [
|
||||
"cfg-if",
|
||||
"cpufeatures",
|
||||
"cpufeatures 0.2.17",
|
||||
"digest",
|
||||
]
|
||||
|
||||
@@ -4312,7 +4372,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283"
|
||||
dependencies = [
|
||||
"cfg-if",
|
||||
"cpufeatures",
|
||||
"cpufeatures 0.2.17",
|
||||
"digest",
|
||||
]
|
||||
|
||||
@@ -4651,6 +4711,17 @@ dependencies = [
|
||||
"version-compare",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "tar"
|
||||
version = "0.4.46"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3f6221d9a6003c78398e3b239969f352578258df48c8eb051caadae0015bc840"
|
||||
dependencies = [
|
||||
"filetime",
|
||||
"libc",
|
||||
"xattr",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "target-lexicon"
|
||||
version = "0.12.16"
|
||||
@@ -6090,6 +6161,16 @@ version = "0.13.2"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "ea6fc2961e4ef194dcbfe56bb845534d0dc8098940c7e5c012a258bfec6701bd"
|
||||
|
||||
[[package]]
|
||||
name = "xattr"
|
||||
version = "1.6.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "32e45ad4206f6d2479085147f02bc2ef834ac85886624a23575ae137c8aa8156"
|
||||
dependencies = [
|
||||
"libc",
|
||||
"rustix 1.1.4",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "xcursor"
|
||||
version = "0.3.11"
|
||||
@@ -6372,10 +6453,13 @@ checksum = "e13c156562582aa81c60cb29407084cdb54c4164760106ab78e6c5b0858cf64e"
|
||||
|
||||
[[package]]
|
||||
name = "zerosend"
|
||||
version = "2.5.2"
|
||||
version = "2.5.5"
|
||||
dependencies = [
|
||||
"arboard",
|
||||
"async-stream",
|
||||
"axum",
|
||||
"blake3",
|
||||
"bytes",
|
||||
"chrono",
|
||||
"clap",
|
||||
"crossterm",
|
||||
@@ -6395,6 +6479,7 @@ dependencies = [
|
||||
"serde_json",
|
||||
"sha2",
|
||||
"socket2 0.5.10",
|
||||
"tar",
|
||||
"tokio",
|
||||
"tokio-util",
|
||||
"tower",
|
||||
@@ -6409,6 +6494,7 @@ dependencies = [
|
||||
"winreg",
|
||||
"winres",
|
||||
"zip",
|
||||
"zstd",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -6459,6 +6545,7 @@ dependencies = [
|
||||
"memchr",
|
||||
"thiserror 2.0.20",
|
||||
"zopfli",
|
||||
"zstd",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -6479,6 +6566,34 @@ dependencies = [
|
||||
"simd-adler32",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "zstd"
|
||||
version = "0.13.3"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "e91ee311a569c327171651566e07972200e76fcfe2242a4fa446149a3881c08a"
|
||||
dependencies = [
|
||||
"zstd-safe",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "zstd-safe"
|
||||
version = "7.2.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "8f49c4d5f0abb602a93fb8736af2a4f4dd9512e36f7f570d66e65ff867ed3b9d"
|
||||
dependencies = [
|
||||
"zstd-sys",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "zstd-sys"
|
||||
version = "2.0.16+zstd.1.5.7"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "91e19ebc2adc8f83e43039e79776e3fda8ca919132d68a1fed6a5faca2683748"
|
||||
dependencies = [
|
||||
"cc",
|
||||
"pkg-config",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "zune-core"
|
||||
version = "0.5.3"
|
||||
|
||||
+9
-2
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "zerosend"
|
||||
version = "2.5.2"
|
||||
version = "2.5.5"
|
||||
edition = "2021"
|
||||
description = "Fast, peer-to-peer file transfer designed for ZeroTier, Tailscale, Radmin VPN and LAN"
|
||||
authors = ["RarDog"]
|
||||
@@ -10,6 +10,8 @@ authors = ["RarDog"]
|
||||
tokio = { version = "1.43", features = ["full"] }
|
||||
tokio-util = { version = "0.7", features = ["io", "codec"] }
|
||||
futures-util = "0.3"
|
||||
async-stream = "0.3"
|
||||
bytes = "1.9"
|
||||
axum = { version = "0.8", features = ["multipart"] }
|
||||
tower = "0.5"
|
||||
tower-http = { version = "0.6", features = ["cors", "trace", "fs"] }
|
||||
@@ -19,6 +21,11 @@ socket2 = { version = "0.5", features = ["all"] }
|
||||
network-interface = "2.0"
|
||||
ipnetwork = "0.20"
|
||||
|
||||
# Compression & Archiving
|
||||
zstd = "0.13"
|
||||
tar = "0.4"
|
||||
zip = { version = "2.2", default-features = false, features = ["deflate", "zstd"] }
|
||||
|
||||
# GUI & Dialogs
|
||||
eframe = "0.29"
|
||||
rfd = "0.15"
|
||||
@@ -26,7 +33,6 @@ qrcode = "0.14"
|
||||
arboard = "3.4"
|
||||
image = { version = "0.25", default-features = false, features = ["png"] }
|
||||
tray-icon = "0.19"
|
||||
zip = { version = "2.2", default-features = false, features = ["deflate"] }
|
||||
|
||||
# CLI & TUI
|
||||
clap = { version = "4.5", features = ["derive"] }
|
||||
@@ -43,6 +49,7 @@ sha2 = "0.10"
|
||||
hex = "0.4"
|
||||
walkdir = "2.5"
|
||||
human_bytes = "0.4"
|
||||
blake3 = "1.5"
|
||||
|
||||
# Windows Native Registry & Sound
|
||||
[target.'cfg(windows)'.dependencies]
|
||||
|
||||
@@ -212,3 +212,29 @@ fn get_default_download_dir() -> PathBuf {
|
||||
let _ = std::fs::create_dir_all(&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);
|
||||
}
|
||||
}
|
||||
|
||||
+6
-2
@@ -192,6 +192,9 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
TransferState::Rejected => {
|
||||
pb.abandon_with_message("⛔ Transfer rejected by recipient.");
|
||||
}
|
||||
TransferState::Paused => {
|
||||
pb.set_message("⏸ Paused");
|
||||
}
|
||||
TransferState::Canceled => {
|
||||
pb.abandon_with_message("🛑 Canceled.");
|
||||
}
|
||||
@@ -310,7 +313,8 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let disc_ref = &discovery;
|
||||
let discovery = std::sync::Arc::new(discovery);
|
||||
let disc_clone = discovery.clone();
|
||||
let srv = server;
|
||||
let in_req_rx = channels.incoming_request_rx;
|
||||
let in_txt_rx = channels.incoming_text_rx;
|
||||
@@ -321,7 +325,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
Box::new(move |cc| {
|
||||
Ok(Box::new(ui::gui::ZeroSendApp::new(
|
||||
cc,
|
||||
disc_ref,
|
||||
disc_clone,
|
||||
srv,
|
||||
config,
|
||||
initial_files,
|
||||
|
||||
@@ -11,6 +11,8 @@ use tokio::sync::RwLock;
|
||||
use tracing::{debug, error, info, warn};
|
||||
|
||||
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)]
|
||||
pub const DEFAULT_TRANSFER_PORT: u16 = 53318;
|
||||
pub const BEACON_INTERVAL_SECS: u64 = 3;
|
||||
@@ -143,8 +145,10 @@ impl DiscoveryService {
|
||||
let (shutdown_tx, _) = tokio::sync::broadcast::channel::<()>(1);
|
||||
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 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 shared_folders = self.shared_folders.clone();
|
||||
let mut shutdown_rx1 = shutdown_tx.subscribe();
|
||||
@@ -163,8 +167,8 @@ impl DiscoveryService {
|
||||
}
|
||||
|
||||
info!(
|
||||
"Started discovery beacon on {} (broadcast: {})",
|
||||
my_info_static.ip, broadcast_addr
|
||||
"Started discovery beacon on {} (broadcast: {}, multicast: {})",
|
||||
my_info_static.ip, broadcast_addr, multicast_addr
|
||||
);
|
||||
|
||||
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();
|
||||
if let Ok(data) = serde_json::to_vec(¤t_beacon) {
|
||||
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() => {
|
||||
@@ -185,15 +191,16 @@ impl DiscoveryService {
|
||||
}
|
||||
});
|
||||
|
||||
// 2. Create UDP listener socket
|
||||
let listener_socket = match create_udp_listener(DEFAULT_DISCOVERY_PORT) {
|
||||
// 2. Create UDP listener socket with Multicast support
|
||||
let iface_ip = self.selected_iface.ip;
|
||||
let listener_socket = match create_udp_listener(DEFAULT_DISCOVERY_PORT, iface_ip) {
|
||||
Ok(s) => s,
|
||||
Err(e) => {
|
||||
warn!(
|
||||
"Could not bind listener on port {}: {}. Will try random port.",
|
||||
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
|
||||
fn create_udp_listener(port: u16) -> std::io::Result<UdpSocket> {
|
||||
/// Helper to create a reusable broadcast & multicast socket on Windows / Linux / macOS
|
||||
fn create_udp_listener(port: u16, iface_ip: Ipv4Addr) -> std::io::Result<UdpSocket> {
|
||||
let socket = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP))?;
|
||||
socket.set_reuse_address(true)?;
|
||||
|
||||
#[cfg(not(windows))]
|
||||
socket.set_reuse_port(true)?;
|
||||
let _ = socket.set_reuse_port(true);
|
||||
|
||||
socket.set_broadcast(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();
|
||||
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();
|
||||
UdpSocket::from_std(std_socket)
|
||||
}
|
||||
|
||||
@@ -6,7 +6,12 @@ use walkdir::WalkDir;
|
||||
use zip::write::SimpleFileOptions;
|
||||
use zip::{CompressionMethod, ZipArchive, ZipWriter};
|
||||
|
||||
// ─────────────────────────────────────────────────────────────────
|
||||
// Legacy Zip helpers (kept for backward-compat & single-file zips)
|
||||
// ─────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Compress a whole folder into a ZIP archive with fast Deflate compression
|
||||
#[allow(dead_code)]
|
||||
pub fn compress_folder_to_zip(
|
||||
src_dir: &Path,
|
||||
zip_dest: &Path,
|
||||
@@ -97,3 +102,234 @@ pub fn extract_zip_archive(
|
||||
info!("Extraction completed successfully for {:?}", zip_path);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ─────────────────────────────────────────────────────────────────
|
||||
// Streaming Tar + Zstd (zero disk footprint for folder transfers)
|
||||
// ─────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Returns an estimated uncompressed byte count for the folder
|
||||
/// (used only to display progress; actual compressed size differs).
|
||||
pub fn estimate_folder_size(dir: &Path) -> u64 {
|
||||
WalkDir::new(dir)
|
||||
.into_iter()
|
||||
.filter_map(|e| e.ok())
|
||||
.filter(|e| e.path().is_file())
|
||||
.map(|e| e.metadata().map(|m| m.len()).unwrap_or(0))
|
||||
.sum()
|
||||
}
|
||||
|
||||
/// Returns the optimal zstd compression level for a given filename.
|
||||
/// Returns 0 for already-compressed formats (store only), 1 otherwise (ultra-fast).
|
||||
pub fn optimal_zstd_level_for_file(name: &str) -> i32 {
|
||||
let ext = name.rsplit('.').next().unwrap_or("").to_lowercase();
|
||||
match ext.as_str() {
|
||||
// Already compressed — don't waste CPU
|
||||
"jpg" | "jpeg" | "png" | "gif" | "webp" | "avif" | "heic" |
|
||||
"mp4" | "mkv" | "avi" | "mov" | "wmv" | "flv" | "webm" |
|
||||
"mp3" | "aac" | "ogg" | "flac" | "m4a" | "opus" |
|
||||
"zip" | "rar" | "7z" | "gz" | "bz2" | "xz" | "zst" | "lz4" |
|
||||
"pdf" | "docx" | "xlsx" | "pptx" => 0,
|
||||
// Compressible text/code/data
|
||||
_ => 1,
|
||||
}
|
||||
}
|
||||
|
||||
/// Estimates the compressibility of a folder: returns true if most content is compressible.
|
||||
pub fn folder_is_compressible(dir: &Path) -> bool {
|
||||
let mut compressible = 0u64;
|
||||
let mut total = 0u64;
|
||||
for entry in WalkDir::new(dir).into_iter().filter_map(|e| e.ok()) {
|
||||
if entry.path().is_file() {
|
||||
let size = entry.metadata().map(|m| m.len()).unwrap_or(0);
|
||||
total += size;
|
||||
let name = entry.file_name().to_string_lossy();
|
||||
if optimal_zstd_level_for_file(&name) > 0 {
|
||||
compressible += size;
|
||||
}
|
||||
}
|
||||
}
|
||||
if total == 0 { return true; }
|
||||
(compressible as f64 / total as f64) > 0.3
|
||||
}
|
||||
|
||||
/// Stream an entire directory as a Zstandard-compressed TAR archive
|
||||
/// into a `std::sync::mpsc::SyncSender<Vec<u8>>` channel.
|
||||
///
|
||||
/// Run this inside `tokio::task::spawn_blocking` so it does not stall
|
||||
/// the async executor. The caller receives raw compressed bytes and
|
||||
/// can forward them directly over the network without any temp file.
|
||||
pub fn stream_folder_as_tar_zst(
|
||||
src_dir: PathBuf,
|
||||
tx: std::sync::mpsc::SyncSender<std::io::Result<Vec<u8>>>,
|
||||
zstd_level: i32, // 1 = ultra-fast, 3 = balanced
|
||||
) {
|
||||
info!("Streaming folder {:?} as tar.zst (level={})", src_dir, zstd_level);
|
||||
|
||||
let writer = ChannelWriter::new(tx.clone(), 256 * 1024);
|
||||
|
||||
let enc = match zstd::stream::write::Encoder::new(writer, zstd_level) {
|
||||
Ok(e) => e,
|
||||
Err(e) => {
|
||||
let _ = tx.send(Err(std::io::Error::other(e.to_string())));
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let mut tar = tar::Builder::new(enc);
|
||||
let parent = src_dir.parent().unwrap_or(&src_dir);
|
||||
|
||||
for entry in WalkDir::new(&src_dir).into_iter() {
|
||||
let entry = match entry {
|
||||
Ok(e) => e,
|
||||
Err(e) => {
|
||||
error!("WalkDir error: {}", e);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let path = entry.path();
|
||||
let rel_path = match path.strip_prefix(parent) {
|
||||
Ok(r) => r,
|
||||
Err(_) => path,
|
||||
};
|
||||
|
||||
if rel_path.as_os_str().is_empty() || rel_path == Path::new(".") {
|
||||
continue;
|
||||
}
|
||||
|
||||
if path.is_dir() {
|
||||
let _ = tar.append_dir(rel_path, path);
|
||||
} else if path.is_file() {
|
||||
match File::open(path) {
|
||||
Ok(mut f) => {
|
||||
let _ = tar.append_file(rel_path, &mut f);
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Cannot open {:?}: {}", path, e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let enc = match tar.into_inner() {
|
||||
Ok(e) => e,
|
||||
Err(e) => {
|
||||
let _ = tx.send(Err(std::io::Error::other(e.to_string())));
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
if let Err(e) = enc.finish() {
|
||||
let _ = tx.send(Err(std::io::Error::other(e.to_string())));
|
||||
return;
|
||||
}
|
||||
|
||||
info!("Folder streaming completed for {:?}", src_dir);
|
||||
}
|
||||
|
||||
/// Extract a `tar.zst` byte stream directly into `dest_dir`,
|
||||
/// with path-traversal protection.
|
||||
///
|
||||
/// Call this inside `spawn_blocking`.
|
||||
pub fn extract_tar_zst_stream<R: Read>(
|
||||
reader: R,
|
||||
dest_dir: &Path,
|
||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
info!("Extracting tar.zst stream into {:?}", dest_dir);
|
||||
|
||||
let decoder = zstd::stream::read::Decoder::new(reader)?;
|
||||
let mut archive = tar::Archive::new(decoder);
|
||||
|
||||
for entry in archive.entries()? {
|
||||
let mut entry = entry?;
|
||||
let entry_path = entry.path()?.into_owned();
|
||||
|
||||
let safe_rel: PathBuf = entry_path
|
||||
.components()
|
||||
.filter(|c| matches!(c, std::path::Component::Normal(_)))
|
||||
.collect();
|
||||
|
||||
if safe_rel.as_os_str().is_empty() {
|
||||
continue;
|
||||
}
|
||||
|
||||
let outpath = dest_dir.join(&safe_rel);
|
||||
|
||||
match entry.header().entry_type() {
|
||||
tar::EntryType::Directory => {
|
||||
let _ = std::fs::create_dir_all(&outpath);
|
||||
}
|
||||
tar::EntryType::Regular => {
|
||||
if let Some(p) = outpath.parent() {
|
||||
let _ = std::fs::create_dir_all(p);
|
||||
}
|
||||
if let Ok(mut out) = File::create(&outpath) {
|
||||
if let Err(e) = std::io::copy(&mut entry, &mut out) {
|
||||
error!("Write error for {:?}: {}", outpath, e);
|
||||
}
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
info!("Extraction completed into {:?}", dest_dir);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ─────────────────────────────────────────────────────────────────
|
||||
// Internal helpers
|
||||
// ─────────────────────────────────────────────────────────────────
|
||||
|
||||
/// A `Write` adapter that accumulates bytes and flushes chunks through
|
||||
/// a `SyncSender<Result<Vec<u8>>>` when the buffer is full.
|
||||
struct ChannelWriter {
|
||||
tx: std::sync::mpsc::SyncSender<std::io::Result<Vec<u8>>>,
|
||||
buf: Vec<u8>,
|
||||
chunk_size: usize,
|
||||
}
|
||||
|
||||
impl ChannelWriter {
|
||||
fn new(
|
||||
tx: std::sync::mpsc::SyncSender<std::io::Result<Vec<u8>>>,
|
||||
chunk_size: usize,
|
||||
) -> Self {
|
||||
Self {
|
||||
tx,
|
||||
buf: Vec::with_capacity(chunk_size),
|
||||
chunk_size,
|
||||
}
|
||||
}
|
||||
|
||||
fn flush_buf(&mut self) -> std::io::Result<()> {
|
||||
if self.buf.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
let chunk = std::mem::replace(&mut self.buf, Vec::with_capacity(self.chunk_size));
|
||||
self.tx
|
||||
.send(Ok(chunk))
|
||||
.map_err(|_| std::io::Error::new(std::io::ErrorKind::BrokenPipe, "receiver dropped"))
|
||||
}
|
||||
}
|
||||
|
||||
impl Write for ChannelWriter {
|
||||
fn write(&mut self, data: &[u8]) -> std::io::Result<usize> {
|
||||
self.buf.extend_from_slice(data);
|
||||
if self.buf.len() >= self.chunk_size {
|
||||
self.flush_buf()?;
|
||||
}
|
||||
Ok(data.len())
|
||||
}
|
||||
|
||||
fn flush(&mut self) -> std::io::Result<()> {
|
||||
self.flush_buf()
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for ChannelWriter {
|
||||
fn drop(&mut self) {
|
||||
let _ = self.flush_buf();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
+449
-60
@@ -1,4 +1,5 @@
|
||||
use super::archive::compress_folder_to_zip;
|
||||
#[allow(unused_imports)]
|
||||
use super::archive;
|
||||
use super::protocol::{
|
||||
FileMetadata, GrepResponse, PrepareTransferRequest, PrepareTransferResponse,
|
||||
SharedFileSaveRequest, SharedFileSaveResponse, SharedFolderInfo, SharedFolderTreeResponse,
|
||||
@@ -9,8 +10,9 @@ use super::text_analyzer::{analyze_text, AnalyzedText};
|
||||
use crate::network::DiscoveryBeacon;
|
||||
use reqwest::Client;
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::collections::VecDeque;
|
||||
use std::collections::{HashMap, VecDeque};
|
||||
use std::path::PathBuf;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
use tokio::fs::File;
|
||||
@@ -59,6 +61,7 @@ struct ChunkTask {
|
||||
|
||||
pub struct TransferClient {
|
||||
client: Client,
|
||||
pub pause_flags: Arc<Mutex<HashMap<String, Arc<AtomicBool>>>>,
|
||||
}
|
||||
|
||||
impl TransferClient {
|
||||
@@ -67,10 +70,63 @@ impl TransferClient {
|
||||
.timeout(Duration::from_secs(7200))
|
||||
.connect_timeout(Duration::from_secs(6))
|
||||
.tcp_nodelay(true)
|
||||
.tcp_keepalive(Some(Duration::from_secs(15)))
|
||||
.pool_idle_timeout(Duration::from_secs(120))
|
||||
.pool_max_idle_per_host(64)
|
||||
.build()
|
||||
.unwrap_or_default();
|
||||
|
||||
Self { client }
|
||||
Self {
|
||||
client,
|
||||
pause_flags: Arc::new(Mutex::new(HashMap::new())),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn pause_transfer(&self, session_id: &str) {
|
||||
let flags = self.pause_flags.lock().await;
|
||||
if let Some(flag) = flags.get(session_id) {
|
||||
flag.store(true, Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn resume_transfer(&self, session_id: &str) {
|
||||
let flags = self.pause_flags.lock().await;
|
||||
if let Some(flag) = flags.get(session_id) {
|
||||
flag.store(false, Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub async fn is_paused(&self, session_id: &str) -> bool {
|
||||
let flags = self.pause_flags.lock().await;
|
||||
flags.get(session_id).map(|f| f.load(Ordering::Relaxed)).unwrap_or(false)
|
||||
}
|
||||
|
||||
/// 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> {
|
||||
if candidates.is_empty() { return None; }
|
||||
if candidates.len() == 1 { return Some(candidates[0].clone()); }
|
||||
|
||||
let (tx, mut rx) = tokio::sync::mpsc::channel(candidates.len());
|
||||
|
||||
for url in candidates {
|
||||
let client = self.client.clone();
|
||||
let tx = tx.clone();
|
||||
let url = url.clone();
|
||||
tokio::spawn(async move {
|
||||
let start = Instant::now();
|
||||
if let Ok(resp) = client.get(&format!("{}/api/info", url)).send().await {
|
||||
if resp.status().is_success() {
|
||||
let latency = start.elapsed().as_millis() as u64;
|
||||
let _ = tx.send((url, latency)).await;
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
drop(tx);
|
||||
|
||||
rx.recv().await.map(|(url, _)| url)
|
||||
}
|
||||
|
||||
/// Ping a remote peer to check health and get beacon info
|
||||
@@ -160,8 +216,8 @@ impl TransferClient {
|
||||
{
|
||||
let mut file_entries: Vec<(PathBuf, String, u64)> = Vec::new();
|
||||
let mut total_bytes = 0u64;
|
||||
let temp_dir = std::env::temp_dir().join("ZeroSend_Zip");
|
||||
let _ = std::fs::create_dir_all(&temp_dir);
|
||||
// Folders that need streaming tar.zst transfer (no disk zip)
|
||||
let mut streaming_folders: Vec<(PathBuf, String)> = Vec::new();
|
||||
|
||||
for p in paths {
|
||||
if !p.exists() {
|
||||
@@ -175,35 +231,18 @@ impl TransferClient {
|
||||
.map(|n| n.to_string_lossy().into_owned())
|
||||
.unwrap_or_else(|| "folder".to_string());
|
||||
|
||||
let zip_filename = format!("{}.zip", folder_name);
|
||||
let temp_zip_path = temp_dir.join(&zip_filename);
|
||||
// Estimate uncompressed size for progress display
|
||||
let estimated_size = {
|
||||
let p_clone = p.clone();
|
||||
tokio::task::spawn_blocking(move || {
|
||||
super::archive::estimate_folder_size(&p_clone)
|
||||
})
|
||||
.await
|
||||
.unwrap_or(0)
|
||||
};
|
||||
|
||||
// Notify UI that folder is being compressed into ZIP archive
|
||||
on_progress(TransferProgress {
|
||||
session_id: Uuid::new_v4().to_string(),
|
||||
peer_name: peer_display_name.to_string(),
|
||||
is_incoming: false,
|
||||
current_file_index: 0,
|
||||
total_files: paths.len(),
|
||||
current_file_name: format!("📦 Compressing folder '{}' into ZIP...", folder_name),
|
||||
bytes_transferred: 0,
|
||||
total_bytes: 0,
|
||||
speed_bps: 0,
|
||||
state: TransferState::InProgress,
|
||||
checksum_verified: false,
|
||||
});
|
||||
|
||||
let p_clone = p.clone();
|
||||
let temp_zip_clone = temp_zip_path.clone();
|
||||
let zip_size = tokio::task::spawn_blocking(move || {
|
||||
compress_folder_to_zip(&p_clone, &temp_zip_clone)
|
||||
})
|
||||
.await
|
||||
.map_err(|e| format!("Compression task error: {}", e))?
|
||||
.map_err(|e| format!("Failed to compress folder {:?}: {}", p, e))?;
|
||||
|
||||
total_bytes += zip_size;
|
||||
file_entries.push((temp_zip_path, zip_filename, zip_size));
|
||||
total_bytes += estimated_size;
|
||||
streaming_folders.push((p.clone(), folder_name));
|
||||
} else {
|
||||
let base_parent = p.parent().unwrap_or(p);
|
||||
for entry in WalkDir::new(p).into_iter().filter_map(|e| e.ok()) {
|
||||
@@ -237,10 +276,177 @@ impl TransferClient {
|
||||
}
|
||||
}
|
||||
|
||||
// ── Stream each folder as tar.zst to the peer (zero disk footprint) ──
|
||||
for (folder_path, folder_name) in streaming_folders {
|
||||
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 {
|
||||
session_id: session_id.clone(),
|
||||
peer_name: peer_display_name.to_string(),
|
||||
is_incoming: false,
|
||||
current_file_index: 0,
|
||||
total_files: paths.len(),
|
||||
current_file_name: format!("📁 {}", folder_name),
|
||||
bytes_transferred: 0,
|
||||
total_bytes: folder_est_size,
|
||||
speed_bps: 0,
|
||||
state: TransferState::InProgress,
|
||||
checksum_verified: false,
|
||||
compressed_bytes: 0,
|
||||
});
|
||||
|
||||
// Feature 1: Determine optimal compression level for entire folder
|
||||
let is_compressible = {
|
||||
let folder_clone2 = folder_path.clone();
|
||||
tokio::task::spawn_blocking(move || {
|
||||
super::archive::folder_is_compressible(&folder_clone2)
|
||||
})
|
||||
.await
|
||||
.unwrap_or(true)
|
||||
};
|
||||
let folder_zstd_level = if is_compressible { 1 } else { 0 };
|
||||
|
||||
// Build a bounded channel: blocking producer → async consumer → HTTP body
|
||||
let (tx, rx) = std::sync::mpsc::sync_channel::<std::io::Result<Vec<u8>>>(32);
|
||||
let folder_clone = folder_path.clone();
|
||||
|
||||
// Spawn blocking tar+zstd compressor thread
|
||||
let compress_handle = tokio::task::spawn_blocking(move || {
|
||||
super::archive::stream_folder_as_tar_zst(folder_clone, tx, folder_zstd_level);
|
||||
});
|
||||
|
||||
let stream_url = format!(
|
||||
"{}/api/transfer/folder_stream/{}",
|
||||
peer_url, session_id
|
||||
);
|
||||
|
||||
// Wrap sync receiver in async stream for reqwest using Arc<Mutex>
|
||||
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! {
|
||||
loop {
|
||||
let rx_clone = rx_arc.clone();
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
rx_clone.lock().unwrap().recv()
|
||||
}).await;
|
||||
match result {
|
||||
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))) => {
|
||||
yield Err(e);
|
||||
break;
|
||||
}
|
||||
Ok(Err(_)) | Err(_) => break,
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let resp = self
|
||||
.client
|
||||
.post(&stream_url)
|
||||
.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")
|
||||
.body(reqwest::Body::wrap_stream(body_stream))
|
||||
.send()
|
||||
.await;
|
||||
|
||||
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 {
|
||||
Ok(r) if r.status().is_success() => {
|
||||
on_progress(TransferProgress {
|
||||
session_id: session_id.clone(),
|
||||
peer_name: peer_display_name.to_string(),
|
||||
is_incoming: false,
|
||||
current_file_index: 0,
|
||||
total_files: paths.len(),
|
||||
current_file_name: format!("📁 {}", folder_name),
|
||||
bytes_transferred: folder_est_size,
|
||||
total_bytes: folder_est_size,
|
||||
speed_bps: 0,
|
||||
state: TransferState::Completed,
|
||||
checksum_verified: true,
|
||||
compressed_bytes: bytes_sent,
|
||||
});
|
||||
}
|
||||
Ok(r) => {
|
||||
return Err(format!(
|
||||
"Folder stream rejected by peer: HTTP {}",
|
||||
r.status()
|
||||
));
|
||||
}
|
||||
Err(e) => {
|
||||
return Err(format!("Failed to stream folder: {}", e));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if file_entries.is_empty() && paths.iter().all(|p| p.is_dir() && auto_zip_folders) {
|
||||
// All items were streaming folders — done
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if file_entries.is_empty() {
|
||||
return Err("No files selected for transfer".to_string());
|
||||
}
|
||||
|
||||
// Feature 7: Priority queue - sort files by size ascending so small files finish first
|
||||
file_entries.sort_by_key(|(_, _, size)| *size);
|
||||
|
||||
let file_metas: Vec<FileMetadata> = file_entries
|
||||
.iter()
|
||||
.map(|(path, rel_path, size)| {
|
||||
@@ -281,6 +487,7 @@ impl TransferClient {
|
||||
speed_bps: 0,
|
||||
state: TransferState::PendingConfirmation,
|
||||
checksum_verified: false,
|
||||
compressed_bytes: 0,
|
||||
});
|
||||
|
||||
let prepare_url = format!("{}/api/transfer/prepare", peer_url);
|
||||
@@ -311,6 +518,7 @@ impl TransferClient {
|
||||
speed_bps: 0,
|
||||
state: TransferState::Rejected,
|
||||
checksum_verified: false,
|
||||
compressed_bytes: 0,
|
||||
});
|
||||
|
||||
return Err(reason);
|
||||
@@ -365,9 +573,16 @@ impl TransferClient {
|
||||
}
|
||||
}
|
||||
|
||||
let pause_flag = {
|
||||
let mut flags = self.pause_flags.lock().await;
|
||||
let flag = Arc::new(AtomicBool::new(false));
|
||||
flags.insert(session_id.clone(), flag.clone());
|
||||
flag
|
||||
};
|
||||
|
||||
let rate_limiter = RateLimiter::new((max_upload_speed_mbps as u64) * 1024 * 1024);
|
||||
let tasks_mutex = Arc::new(Mutex::new(chunk_tasks));
|
||||
let (chunk_tx, mut chunk_rx) = tokio::sync::mpsc::unbounded_channel::<(usize, String)>();
|
||||
let (chunk_tx, mut chunk_rx) = tokio::sync::mpsc::unbounded_channel::<(usize, usize, String)>();
|
||||
let (progress_update_tx, mut progress_update_rx) = tokio::sync::mpsc::unbounded_channel::<TransferProgress>();
|
||||
|
||||
let s_id = session_id.clone();
|
||||
@@ -375,15 +590,17 @@ impl TransferClient {
|
||||
let t_files = file_metas.len();
|
||||
let p_tx = progress_update_tx.clone();
|
||||
|
||||
// Background progress tracker with smoothed EMA speed calculation
|
||||
// Background progress tracker with smoothed EMA speed calculation & compressed_bytes
|
||||
tokio::spawn(async move {
|
||||
let mut total_sent = already_transferred;
|
||||
let mut total_wire = already_transferred;
|
||||
let mut last_bytes = total_sent;
|
||||
let mut last_time = Instant::now();
|
||||
let mut smoothed_speed = 0.0f64;
|
||||
|
||||
while let Some((len, fname)) = chunk_rx.recv().await {
|
||||
total_sent += len as u64;
|
||||
while let Some((raw_len, wire_len, fname)) = chunk_rx.recv().await {
|
||||
total_sent += raw_len as u64;
|
||||
total_wire += wire_len as u64;
|
||||
|
||||
if last_time.elapsed().as_millis() >= 350 {
|
||||
let elapsed_secs = last_time.elapsed().as_secs_f64();
|
||||
@@ -414,6 +631,7 @@ impl TransferClient {
|
||||
speed_bps,
|
||||
state: TransferState::InProgress,
|
||||
checksum_verified: false,
|
||||
compressed_bytes: total_wire,
|
||||
});
|
||||
|
||||
last_time = Instant::now();
|
||||
@@ -434,6 +652,7 @@ impl TransferClient {
|
||||
speed_bps: 0,
|
||||
state: TransferState::InProgress,
|
||||
checksum_verified: false,
|
||||
compressed_bytes: already_transferred,
|
||||
});
|
||||
|
||||
let start_time = Instant::now();
|
||||
@@ -446,6 +665,7 @@ impl TransferClient {
|
||||
let chunk_url = format!("{}/api/transfer/chunk/{}", peer_url, session_id);
|
||||
let limiter_c = rate_limiter.clone();
|
||||
let tx_c = chunk_tx.clone();
|
||||
let pause_c = pause_flag.clone();
|
||||
|
||||
worker_set.spawn(async move {
|
||||
loop {
|
||||
@@ -472,16 +692,41 @@ impl TransferClient {
|
||||
return Err(format!("Failed to read chunk from {:?}: {}", chunk_task.file_path, e));
|
||||
}
|
||||
|
||||
limiter_c.throttle(buf.len()).await;
|
||||
// Feature 1: Smart compression check per file type
|
||||
let is_compressible = super::archive::optimal_zstd_level_for_file(&chunk_task.rel_path) > 0;
|
||||
let (payload, is_zstd) = if is_compressible && buf.len() >= 1024 {
|
||||
if let Ok(compressed) = zstd::bulk::compress(&buf, 1) {
|
||||
if compressed.len() + 64 < buf.len() {
|
||||
(compressed, true)
|
||||
} else {
|
||||
(buf, false)
|
||||
}
|
||||
} else {
|
||||
(buf, false)
|
||||
}
|
||||
} else {
|
||||
(buf, false)
|
||||
};
|
||||
|
||||
let resp = client_c
|
||||
let wire_len = payload.len();
|
||||
limiter_c.throttle(wire_len).await;
|
||||
|
||||
// Feature 5: Pause support in worker loop
|
||||
while pause_c.load(Ordering::Relaxed) {
|
||||
tokio::time::sleep(Duration::from_millis(200)).await;
|
||||
}
|
||||
|
||||
let mut req = client_c
|
||||
.post(&chunk_url)
|
||||
.header("x-file-name", &chunk_task.rel_path)
|
||||
.header("x-chunk-offset", chunk_task.offset.to_string())
|
||||
.header("x-chunk-length", chunk_task.length.to_string())
|
||||
.body(buf)
|
||||
.send()
|
||||
.await;
|
||||
.header("x-chunk-length", chunk_task.length.to_string());
|
||||
|
||||
if is_zstd {
|
||||
req = req.header("x-compression", "zstd");
|
||||
}
|
||||
|
||||
let resp = req.body(payload).send().await;
|
||||
|
||||
let is_ok = match &resp {
|
||||
Ok(r) => r.status().is_success(),
|
||||
@@ -489,9 +734,14 @@ impl TransferClient {
|
||||
};
|
||||
|
||||
if !is_ok {
|
||||
if chunk_task.retries < 3 {
|
||||
if chunk_task.retries < 8 {
|
||||
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);
|
||||
continue;
|
||||
} else {
|
||||
@@ -503,7 +753,7 @@ impl TransferClient {
|
||||
}
|
||||
}
|
||||
|
||||
let _ = tx_c.send((chunk_task.length, chunk_task.rel_path));
|
||||
let _ = tx_c.send((chunk_task.length, wire_len, chunk_task.rel_path));
|
||||
}
|
||||
Ok::<(), String>(())
|
||||
});
|
||||
@@ -530,6 +780,7 @@ impl TransferClient {
|
||||
speed_bps: 0,
|
||||
state: TransferState::Failed(e.clone()),
|
||||
checksum_verified: false,
|
||||
compressed_bytes: 0,
|
||||
});
|
||||
return Err(e);
|
||||
}
|
||||
@@ -550,6 +801,9 @@ impl TransferClient {
|
||||
on_progress(progress);
|
||||
}
|
||||
|
||||
// Cleanup pause flags
|
||||
self.pause_flags.lock().await.remove(&session_id);
|
||||
|
||||
// Finalize transfer on receiver (atomic rename + zip extraction)
|
||||
let finish_url = format!("{}/api/transfer/finish/{}", peer_url, session_id);
|
||||
let finish_resp = self.client.post(&finish_url).send().await;
|
||||
@@ -577,6 +831,7 @@ impl TransferClient {
|
||||
speed_bps: avg_speed,
|
||||
state: TransferState::Completed,
|
||||
checksum_verified: true,
|
||||
compressed_bytes: total_bytes,
|
||||
});
|
||||
|
||||
Ok(())
|
||||
@@ -671,37 +926,89 @@ impl TransferClient {
|
||||
.map_err(|e| format!("Failed to read file text: {}", e))
|
||||
}
|
||||
|
||||
/// Download a file or folder (zip) from a remote shared folder
|
||||
pub async fn download_shared_file(
|
||||
/// Download a file or folder (zip) from a remote shared folder with live progress tracking
|
||||
pub async fn download_shared_file_with_progress<F>(
|
||||
&self,
|
||||
peer_url: &str,
|
||||
folder_id: &str,
|
||||
subpath: &str,
|
||||
dest_path: &std::path::Path,
|
||||
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 mut url = format!("{}/api/shared/download/{}/{}", peer_url, folder_id, clean);
|
||||
if let Some(p) = pin.filter(|p| !p.trim().is_empty()) {
|
||||
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
|
||||
.client
|
||||
.get(&url)
|
||||
.send()
|
||||
.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() {
|
||||
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
|
||||
.bytes()
|
||||
.await
|
||||
.map_err(|e| format!("Failed to read stream: {}", e))?;
|
||||
let len = bytes.len() as u64;
|
||||
let total_size = resp.content_length().unwrap_or(0);
|
||||
|
||||
if let Some(parent) = dest_path.parent() {
|
||||
let _ = tokio::fs::create_dir_all(parent).await;
|
||||
@@ -711,14 +1018,94 @@ impl TransferClient {
|
||||
.await
|
||||
.map_err(|e| format!("Could not create file: {}", e))?;
|
||||
|
||||
tokio::io::AsyncWriteExt::write_all(&mut file, &bytes)
|
||||
.await
|
||||
.map_err(|e| format!("Could not write file: {}", e))?;
|
||||
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
|
||||
.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)
|
||||
.await
|
||||
.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)
|
||||
@@ -837,6 +1224,7 @@ impl TransferClient {
|
||||
}
|
||||
|
||||
/// Fetch active watch-party state from host peer
|
||||
#[allow(dead_code)]
|
||||
pub async fn get_watch_party(
|
||||
&self,
|
||||
peer_url: &str,
|
||||
@@ -859,6 +1247,7 @@ impl TransferClient {
|
||||
}
|
||||
|
||||
/// Send watch-party sync event to host peer
|
||||
#[allow(dead_code)]
|
||||
pub async fn send_watch_party_event(
|
||||
&self,
|
||||
peer_url: &str,
|
||||
|
||||
+121
-10
@@ -10,6 +10,7 @@ use tracing::info;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct SyncStatusUpdate {
|
||||
#[allow(dead_code)]
|
||||
pub pair_id: String,
|
||||
pub message: String,
|
||||
pub is_error: bool,
|
||||
@@ -22,6 +23,49 @@ pub struct FolderSyncEngine {
|
||||
status_tx: mpsc::UnboundedSender<SyncStatusUpdate>,
|
||||
}
|
||||
|
||||
// ── Sync cache structures ──────────────────────────────────────────────────
|
||||
|
||||
#[derive(serde::Serialize, serde::Deserialize, Default)]
|
||||
struct SyncCache {
|
||||
files: HashMap<String, CacheEntry>,
|
||||
}
|
||||
|
||||
#[derive(serde::Serialize, serde::Deserialize, Clone)]
|
||||
struct CacheEntry {
|
||||
size: u64,
|
||||
mtime: u64,
|
||||
blake3: Option<String>,
|
||||
}
|
||||
|
||||
fn load_sync_cache(local_path: &Path) -> SyncCache {
|
||||
let cache_path = local_path.join(".zerosend_sync_cache.json");
|
||||
std::fs::read_to_string(&cache_path)
|
||||
.ok()
|
||||
.and_then(|s| serde_json::from_str(&s).ok())
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
fn save_sync_cache(local_path: &Path, cache: &SyncCache) {
|
||||
let cache_path = local_path.join(".zerosend_sync_cache.json");
|
||||
if let Ok(s) = serde_json::to_string(cache) {
|
||||
let _ = std::fs::write(&cache_path, s);
|
||||
}
|
||||
}
|
||||
|
||||
fn file_blake3(path: &std::path::Path) -> Option<String> {
|
||||
let mut file = std::fs::File::open(path).ok()?;
|
||||
let mut hasher = blake3::Hasher::new();
|
||||
let mut buf = [0u8; 65536];
|
||||
loop {
|
||||
let n = std::io::Read::read(&mut file, &mut buf).ok()?;
|
||||
if n == 0 { break; }
|
||||
hasher.update(&buf[..n]);
|
||||
}
|
||||
Some(hasher.finalize().to_hex().to_string())
|
||||
}
|
||||
|
||||
// ── Engine ─────────────────────────────────────────────────────────────────
|
||||
|
||||
impl FolderSyncEngine {
|
||||
pub fn new(
|
||||
client: Arc<TransferClient>,
|
||||
@@ -35,6 +79,7 @@ impl FolderSyncEngine {
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub fn pairs_handle(&self) -> Arc<RwLock<Vec<FolderSyncPairConfig>>> {
|
||||
self.pairs.clone()
|
||||
}
|
||||
@@ -109,17 +154,31 @@ impl FolderSyncEngine {
|
||||
let _ = tokio::fs::create_dir_all(local_path).await;
|
||||
}
|
||||
|
||||
// Load sync cache (size/mtime/hash from last sync)
|
||||
let cache = load_sync_cache(local_path);
|
||||
|
||||
// Fetch remote directory tree
|
||||
let remote_tree = client
|
||||
.fetch_shared_tree(peer_url, remote_folder_id, None, None)
|
||||
.await?;
|
||||
|
||||
// 1. Scan local files
|
||||
let mut local_files: HashMap<String, (u64, u64)> = HashMap::new(); // rel_path -> (size, mtime)
|
||||
if let Ok(entries) = walkdir_local(local_path) {
|
||||
for (rel, size, mtime) in entries {
|
||||
local_files.insert(rel, (size, mtime));
|
||||
// 1. Scan local files with blake3 hash support
|
||||
// Returns Vec<(rel_path, size, mtime, Option<hash>)>
|
||||
let local_files_vec = {
|
||||
let local_path_buf = local_path.to_path_buf();
|
||||
tokio::task::spawn_blocking(move || walkdir_local(&local_path_buf))
|
||||
.await
|
||||
.map_err(|e| format!("walkdir task error: {}", e))?
|
||||
.map_err(|e| format!("walkdir error: {}", e))?
|
||||
};
|
||||
|
||||
let mut local_files: HashMap<String, (u64, u64, Option<String>)> = HashMap::new();
|
||||
for (rel, size, mtime, hash) in local_files_vec {
|
||||
// Skip the cache file itself
|
||||
if rel == ".zerosend_sync_cache.json" {
|
||||
continue;
|
||||
}
|
||||
local_files.insert(rel, (size, mtime, hash));
|
||||
}
|
||||
|
||||
// 2. Map remote files
|
||||
@@ -131,12 +190,36 @@ impl FolderSyncEngine {
|
||||
}
|
||||
|
||||
let mut synced_count = 0;
|
||||
let mut updated_cache = SyncCache { files: cache.files.clone() };
|
||||
|
||||
// A. Upload newly created or changed local files to remote
|
||||
if !remote_tree.read_only {
|
||||
for (rel, (size, _)) in &local_files {
|
||||
for (rel, (size, mtime, hash)) in &local_files {
|
||||
let need_upload = match remote_files.get(rel) {
|
||||
Some(remote_size) => *remote_size != *size,
|
||||
Some(remote_size) => {
|
||||
if *remote_size != *size {
|
||||
true // size differs — definitely upload
|
||||
} else {
|
||||
// Same size: check cache to detect content-only changes
|
||||
if let Some(cached) = cache.files.get(rel) {
|
||||
if cached.size == *size && cached.mtime == *mtime {
|
||||
// size AND mtime unchanged since last sync — skip
|
||||
false
|
||||
} else if cached.size == *size {
|
||||
// mtime changed but size same — compare blake3
|
||||
match (hash, &cached.blake3) {
|
||||
(Some(h), Some(c)) => h != c, // upload if hash changed
|
||||
_ => true, // no hash available, upload to be safe
|
||||
}
|
||||
} else {
|
||||
true // cached size differs — upload
|
||||
}
|
||||
} else {
|
||||
// No cache entry — upload to be safe
|
||||
true
|
||||
}
|
||||
}
|
||||
},
|
||||
None => true, // New local file
|
||||
};
|
||||
|
||||
@@ -153,8 +236,21 @@ impl FolderSyncEngine {
|
||||
)
|
||||
.await
|
||||
{
|
||||
// Update cache entry after successful upload
|
||||
updated_cache.files.insert(rel.clone(), CacheEntry {
|
||||
size: *size,
|
||||
mtime: *mtime,
|
||||
blake3: hash.clone(),
|
||||
});
|
||||
synced_count += 1;
|
||||
}
|
||||
} else {
|
||||
// File unchanged — update cache entry to keep mtime fresh
|
||||
updated_cache.files.insert(rel.clone(), CacheEntry {
|
||||
size: *size,
|
||||
mtime: *mtime,
|
||||
blake3: hash.clone(),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -162,7 +258,7 @@ impl FolderSyncEngine {
|
||||
// B. Download newly created remote files to local
|
||||
for (rel, remote_size) in &remote_files {
|
||||
let need_download = match local_files.get(rel) {
|
||||
Some((local_size, _)) => *local_size != *remote_size,
|
||||
Some((local_size, _, _)) => *local_size != *remote_size,
|
||||
None => true, // New remote file
|
||||
};
|
||||
|
||||
@@ -177,11 +273,16 @@ impl FolderSyncEngine {
|
||||
.download_shared_file(peer_url, remote_folder_id, rel, &target_file_path, None)
|
||||
.await
|
||||
{
|
||||
// Invalidate cache entry since we downloaded new content
|
||||
updated_cache.files.remove(rel);
|
||||
synced_count += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Save updated cache
|
||||
save_sync_cache(local_path, &updated_cache);
|
||||
|
||||
if synced_count > 0 {
|
||||
let _ = status_tx.send(SyncStatusUpdate {
|
||||
pair_id: pair_id.to_string(),
|
||||
@@ -195,7 +296,10 @@ impl FolderSyncEngine {
|
||||
}
|
||||
}
|
||||
|
||||
fn walkdir_local(base: &Path) -> Result<Vec<(String, u64, u64)>, std::io::Error> {
|
||||
/// Walk local directory and return (rel_path, size, mtime, Option<blake3_hash>).
|
||||
/// Hash is only computed for files < 500MB to avoid stalling.
|
||||
fn walkdir_local(base: &Path) -> Result<Vec<(String, u64, u64, Option<String>)>, std::io::Error> {
|
||||
const MAX_HASH_SIZE: u64 = 500 * 1024 * 1024; // 500 MB
|
||||
let mut results = Vec::new();
|
||||
for entry in walkdir::WalkDir::new(base).into_iter().filter_map(|e| e.ok()) {
|
||||
if entry.file_type().is_file() {
|
||||
@@ -209,7 +313,14 @@ fn walkdir_local(base: &Path) -> Result<Vec<(String, u64, u64)>, std::io::Error>
|
||||
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
|
||||
.map(|d| d.as_secs())
|
||||
.unwrap_or(0);
|
||||
results.push((rel_str, size, mtime));
|
||||
|
||||
let hash = if size < MAX_HASH_SIZE {
|
||||
file_blake3(entry.path())
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
results.push((rel_str, size, mtime, hash));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -51,6 +51,7 @@ pub enum PrepareTransferResponse {
|
||||
pub enum TransferState {
|
||||
PendingConfirmation,
|
||||
InProgress,
|
||||
Paused,
|
||||
Completed,
|
||||
Failed(String),
|
||||
Rejected,
|
||||
@@ -71,6 +72,8 @@ pub struct TransferProgress {
|
||||
pub state: TransferState,
|
||||
#[serde(default)]
|
||||
pub checksum_verified: bool,
|
||||
#[serde(default)]
|
||||
pub compressed_bytes: u64, // actual bytes sent over network (after compression)
|
||||
}
|
||||
|
||||
impl TransferProgress {
|
||||
|
||||
+431
-6
@@ -1,3 +1,5 @@
|
||||
#[allow(unused_imports)]
|
||||
use super::archive;
|
||||
use super::archive::extract_zip_archive;
|
||||
use super::protocol::{
|
||||
FileMetadata, GrepMatch, GrepResponse, PrepareTransferRequest, PrepareTransferResponse,
|
||||
@@ -186,6 +188,10 @@ impl TransferServer {
|
||||
"/api/transfer/finish/{session_id}",
|
||||
post(handle_finish_transfer),
|
||||
)
|
||||
.route(
|
||||
"/api/transfer/folder_stream/{session_id}",
|
||||
post(handle_folder_stream).layer(DefaultBodyLimit::disable()),
|
||||
)
|
||||
.route(
|
||||
"/api/transfer/upload/{session_id}",
|
||||
post(handle_upload).layer(DefaultBodyLimit::disable()),
|
||||
@@ -216,8 +222,42 @@ impl TransferServer {
|
||||
.with_state(self.state.clone());
|
||||
|
||||
let addr = SocketAddr::from(([0, 0, 0, 0], self.port));
|
||||
let listener = tokio::net::TcpListener::bind(addr).await?;
|
||||
info!("ZeroSend HTTP server listening on {}", addr);
|
||||
let socket = socket2::Socket::new(
|
||||
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();
|
||||
tokio::spawn(async move {
|
||||
@@ -416,6 +456,7 @@ async fn handle_prepare(
|
||||
speed_bps: 0,
|
||||
state: TransferState::InProgress,
|
||||
checksum_verified: false,
|
||||
compressed_bytes: 0,
|
||||
};
|
||||
|
||||
state.active_transfers.write().await.insert(session_id.clone(), progress.clone());
|
||||
@@ -462,6 +503,7 @@ async fn handle_prepare(
|
||||
speed_bps: 0,
|
||||
state: TransferState::InProgress,
|
||||
checksum_verified: false,
|
||||
compressed_bytes: 0,
|
||||
};
|
||||
|
||||
state.active_transfers.write().await.insert(session_id.clone(), progress.clone());
|
||||
@@ -545,8 +587,36 @@ async fn handle_upload_chunk(
|
||||
);
|
||||
}
|
||||
|
||||
let chunk_len = body.len();
|
||||
if let Err(e) = file.write_all(&body).await {
|
||||
let raw_length: usize = headers
|
||||
.get("x-chunk-length")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.and_then(|v| v.parse().ok())
|
||||
.unwrap_or(body.len());
|
||||
|
||||
let is_zstd = headers
|
||||
.get("x-compression")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.map(|v| v.eq_ignore_ascii_case("zstd"))
|
||||
.unwrap_or(false);
|
||||
|
||||
let decompressed_data = if is_zstd {
|
||||
match zstd::bulk::decompress(&body, raw_length.max(16 * 1024 * 1024)) {
|
||||
Ok(decomp) => decomp,
|
||||
Err(e) => {
|
||||
error!("Zstd decompression failed for {}: {}", file_name, e);
|
||||
return (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(serde_json::json!({"error": format!("Zstd decompress error: {}", e)})),
|
||||
);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
body.to_vec()
|
||||
};
|
||||
|
||||
let chunk_len = decompressed_data.len();
|
||||
let wire_bytes = body.len() as u64;
|
||||
if let Err(e) = file.write_all(&decompressed_data).await {
|
||||
error!("Write error on {:?} at {}: {}", part_path, offset, e);
|
||||
return (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
@@ -561,6 +631,7 @@ async fn handle_upload_chunk(
|
||||
let mut transfers = state.active_transfers.write().await;
|
||||
if let Some(t) = transfers.get_mut(&session_id) {
|
||||
t.bytes_transferred = (t.bytes_transferred + chunk_len as u64).min(t.total_bytes);
|
||||
t.compressed_bytes += wire_bytes;
|
||||
t.current_file_name = file_name;
|
||||
t.state = TransferState::InProgress;
|
||||
|
||||
@@ -623,6 +694,9 @@ async fn handle_finish_transfer(
|
||||
let part_path = state.download_dir.join(format!("{}.zerosend_part", safe_rel_path.to_string_lossy()));
|
||||
|
||||
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;
|
||||
if let Some(expected_hash) = &f.sha256 {
|
||||
if let Ok(actual_hash) = super::client::compute_file_sha256(&target_path) {
|
||||
@@ -658,6 +732,179 @@ async fn handle_finish_transfer(
|
||||
)
|
||||
}
|
||||
|
||||
/// Streaming tar.zst folder receiver — zero disk footprint extraction on-the-fly with live progress
|
||||
async fn handle_folder_stream(
|
||||
State(state): State<Arc<ServerState>>,
|
||||
Path(session_id): Path<String>,
|
||||
headers: axum::http::HeaderMap,
|
||||
body: axum::body::Body,
|
||||
) -> impl IntoResponse {
|
||||
let raw_folder_name = headers
|
||||
.get("x-folder-name")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.unwrap_or("folder")
|
||||
.trim_end_matches(".tar.zst")
|
||||
.to_string();
|
||||
|
||||
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);
|
||||
|
||||
let sender_name = headers
|
||||
.get("x-sender-name")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.unwrap_or("Peer")
|
||||
.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) => {
|
||||
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 (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(serde_json::json!({"error": e.to_string()})),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let dest_dir = state.download_dir.clone();
|
||||
let folder_name_clone = raw_folder_name.clone();
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
super::archive::extract_tar_zst_stream(std::io::Cursor::new(data), &dest_dir)
|
||||
})
|
||||
.await;
|
||||
|
||||
match result {
|
||||
Ok(Ok(())) => {
|
||||
info!("Folder '{}' extracted successfully", folder_name_clone);
|
||||
let final_bytes = if total_size > 0 { total_size } else { bytes_received };
|
||||
let progress = super::protocol::TransferProgress {
|
||||
session_id: session_id.clone(),
|
||||
peer_name: sender_name.clone(),
|
||||
is_incoming: true,
|
||||
current_file_index: 1,
|
||||
total_files: 1,
|
||||
current_file_name: format!("📁 {}", folder_name_clone),
|
||||
bytes_transferred: final_bytes,
|
||||
total_bytes: final_bytes,
|
||||
speed_bps: 0,
|
||||
state: super::protocol::TransferState::Completed,
|
||||
checksum_verified: true,
|
||||
compressed_bytes: bytes_received,
|
||||
};
|
||||
state.active_transfers.write().await.insert(session_id.clone(), progress.clone());
|
||||
let _ = state.transfer_events.send(progress);
|
||||
(
|
||||
StatusCode::OK,
|
||||
Json(serde_json::json!({"status": "extracted", "folder": folder_name_clone})),
|
||||
)
|
||||
}
|
||||
Ok(Err(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,
|
||||
Json(serde_json::json!({"error": e.to_string()})),
|
||||
)
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Spawn blocking error: {}", e);
|
||||
(
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(serde_json::json!({"error": e.to_string()})),
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_upload(
|
||||
State(state): State<Arc<ServerState>>,
|
||||
Path(session_id): Path<String>,
|
||||
@@ -896,6 +1143,7 @@ async fn handle_web_upload(
|
||||
speed_bps: 0,
|
||||
state: 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);
|
||||
@@ -1368,7 +1616,96 @@ async fn handle_shared_download(
|
||||
let file_size = file.metadata().await.map(|m| m.len()).unwrap_or(0);
|
||||
let mime_type = super::protocol::get_mime_type(&file_name);
|
||||
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()
|
||||
.status(StatusCode::OK)
|
||||
@@ -1424,9 +1761,97 @@ async fn handle_shared_download(
|
||||
};
|
||||
let file_size = zip_file.metadata().await.map(|m| m.len()).unwrap_or(0);
|
||||
let reader = tokio_util::io::ReaderStream::new(zip_file);
|
||||
let body = Body::from_stream(reader);
|
||||
|
||||
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()
|
||||
.status(StatusCode::OK)
|
||||
|
||||
@@ -14,6 +14,7 @@ use tracing::{info, warn};
|
||||
|
||||
pub struct WebDavState {
|
||||
pub shared_folders: Arc<RwLock<Vec<SharedFolderConfig>>>,
|
||||
#[allow(dead_code)]
|
||||
pub port: u16,
|
||||
}
|
||||
|
||||
@@ -389,6 +390,7 @@ mod urlencoding {
|
||||
}
|
||||
|
||||
/// 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> {
|
||||
let drive = format!("{}:", drive_letter);
|
||||
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
|
||||
#[allow(dead_code)]
|
||||
pub fn unmount_webdav_drive(drive_letter: char) -> Result<(), String> {
|
||||
let drive = format!("{}:", drive_letter);
|
||||
let output = std::process::Command::new("net")
|
||||
|
||||
+455
-137
File diff suppressed because it is too large
Load Diff
@@ -479,6 +479,7 @@ fn draw_ui(f: &mut Frame, app: &mut AppState) {
|
||||
TransferState::Completed => "✅ Completed".to_string(),
|
||||
TransferState::Failed(e) => format!("❌ Failed: {}", e),
|
||||
TransferState::Rejected => "⛔ Rejected".to_string(),
|
||||
TransferState::Paused => "⏸ Paused".to_string(),
|
||||
TransferState::Canceled => "🛑 Canceled".to_string(),
|
||||
};
|
||||
|
||||
|
||||
@@ -57,6 +57,7 @@ impl Default for VimViewerModal {
|
||||
}
|
||||
|
||||
impl VimViewerModal {
|
||||
#[allow(dead_code)]
|
||||
pub fn open(
|
||||
&mut self,
|
||||
file_name: String,
|
||||
|
||||
+72
-13
@@ -43,6 +43,8 @@ pub struct UpdateInfo {
|
||||
pub download_url: Option<String>,
|
||||
pub gitea_url: String,
|
||||
pub asset_size: u64,
|
||||
pub is_beta: bool,
|
||||
pub pre_number: u32,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
@@ -54,21 +56,72 @@ pub enum UpdateProgress {
|
||||
Error(String),
|
||||
}
|
||||
|
||||
/// Parse semantic version string (e.g. "v2.4.1", "2.4.2-beta") into (u32, u32, u32)
|
||||
pub fn parse_version(v: &str) -> (u32, u32, u32) {
|
||||
/// Parse semantic version string into (major, minor, patch, subpatch, is_beta, pre_num).
|
||||
/// Handles formats like "v2.5.3", "2.5.3-beta", "2.5.3-beta.1", "2.5.3.1"
|
||||
pub fn parse_version(v: &str) -> (u32, u32, u32, u32, bool, u32) {
|
||||
let clean = v.trim().trim_start_matches('v').trim_start_matches('V');
|
||||
let parts: Vec<&str> = clean.split('.').collect();
|
||||
let major = parts.first().and_then(|s| s.split('-').next()).and_then(|s| s.parse::<u32>().ok()).unwrap_or(0);
|
||||
let minor = parts.get(1).and_then(|s| s.split('-').next()).and_then(|s| s.parse::<u32>().ok()).unwrap_or(0);
|
||||
let patch = parts.get(2).and_then(|s| s.split('-').next()).and_then(|s| s.parse::<u32>().ok()).unwrap_or(0);
|
||||
(major, minor, patch)
|
||||
|
||||
// Split off pre-release suffix (everything after the first '-')
|
||||
let (numeric_part, pre_part) = if let Some(idx) = clean.find('-') {
|
||||
(&clean[..idx], Some(&clean[idx + 1..]))
|
||||
} else {
|
||||
(clean, None)
|
||||
};
|
||||
|
||||
let num_parts: Vec<&str> = numeric_part.split('.').collect();
|
||||
let major = num_parts.first().and_then(|s| s.parse::<u32>().ok()).unwrap_or(0);
|
||||
let minor = num_parts.get(1).and_then(|s| s.parse::<u32>().ok()).unwrap_or(0);
|
||||
let patch = num_parts.get(2).and_then(|s| s.parse::<u32>().ok()).unwrap_or(0);
|
||||
let subpatch = num_parts.get(3).and_then(|s| s.parse::<u32>().ok()).unwrap_or(0);
|
||||
|
||||
let is_beta = pre_part.map(|p| {
|
||||
let pl = p.to_lowercase();
|
||||
pl.contains("beta") || pl.contains("alpha") || pl.contains("rc")
|
||||
}).unwrap_or(false);
|
||||
|
||||
// Extract pre-release number, e.g. "beta.1.1" -> 1001, "beta.1" -> 1000, "beta" -> 0
|
||||
let pre_num = if let Some(p) = pre_part {
|
||||
let parts: Vec<&str> = p.split('.').collect();
|
||||
let p1 = parts.get(1).and_then(|s| s.parse::<u32>().ok()).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)
|
||||
}
|
||||
|
||||
/// Returns true if latest version is strictly newer than current version
|
||||
/// Returns true if a version string contains "beta", "alpha", or "rc"
|
||||
pub fn is_beta_version(v: &str) -> bool {
|
||||
let lower = v.to_lowercase();
|
||||
lower.contains("beta") || lower.contains("alpha") || lower.contains("rc")
|
||||
}
|
||||
|
||||
/// Returns true if latest version is strictly newer than current version.
|
||||
/// Rules: stable > beta of same numeric version; higher pre_num beats lower.
|
||||
pub fn is_newer_version(latest: &str, current: &str) -> bool {
|
||||
let (l_maj, l_min, l_pat) = parse_version(latest);
|
||||
let (c_maj, c_min, c_pat) = parse_version(current);
|
||||
(l_maj, l_min, l_pat) > (c_maj, c_min, c_pat)
|
||||
let (l_maj, l_min, l_pat, l_sub, l_beta, l_pre) = parse_version(latest);
|
||||
let (c_maj, c_min, c_pat, c_sub, c_beta, c_pre) = parse_version(current);
|
||||
|
||||
// Compare numeric components first
|
||||
let l_num = (l_maj, l_min, l_pat, l_sub);
|
||||
let c_num = (c_maj, c_min, c_pat, c_sub);
|
||||
|
||||
if l_num != c_num {
|
||||
return l_num > c_num;
|
||||
}
|
||||
|
||||
// Same numeric version: stable > beta
|
||||
match (l_beta, c_beta) {
|
||||
(false, true) => true, // latest is stable, current is beta → update
|
||||
(true, false) => false, // latest is beta, current is stable → no update
|
||||
_ => l_pre > c_pre, // both same prerelease type: compare pre number
|
||||
}
|
||||
}
|
||||
|
||||
/// Check Gitea repository for available releases or newer version tags
|
||||
@@ -110,13 +163,15 @@ pub async fn check_gitea_update(api_base: Option<&str>) -> Result<UpdateInfo, St
|
||||
|
||||
return Ok(UpdateInfo {
|
||||
current_version,
|
||||
latest_version: latest_ver,
|
||||
latest_version: latest_ver.clone(),
|
||||
has_update,
|
||||
release_name: release.name,
|
||||
release_notes: release.body,
|
||||
download_url,
|
||||
gitea_url,
|
||||
asset_size,
|
||||
is_beta: is_beta_version(&latest_ver),
|
||||
pre_number: parse_version(&latest_ver).5,
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -149,13 +204,15 @@ pub async fn check_gitea_update(api_base: Option<&str>) -> Result<UpdateInfo, St
|
||||
|
||||
return Ok(UpdateInfo {
|
||||
current_version,
|
||||
latest_version: latest_ver,
|
||||
latest_version: latest_ver.clone(),
|
||||
has_update,
|
||||
release_name: release.name,
|
||||
release_notes: release.body,
|
||||
download_url,
|
||||
gitea_url,
|
||||
asset_size,
|
||||
is_beta: is_beta_version(&latest_ver),
|
||||
pre_number: parse_version(&latest_ver).5,
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -182,6 +239,8 @@ pub async fn check_gitea_update(api_base: Option<&str>) -> Result<UpdateInfo, St
|
||||
download_url: None,
|
||||
gitea_url: GITEA_DEFAULT_REPO_URL.to_string(),
|
||||
asset_size: 0,
|
||||
is_beta: is_beta_version(parsed),
|
||||
pre_number: parse_version(parsed).5,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user