use super::archive::compress_folder_to_zip; use super::protocol::{ FileMetadata, GrepResponse, PrepareTransferRequest, PrepareTransferResponse, SharedFileSaveRequest, SharedFileSaveResponse, SharedFolderInfo, SharedFolderTreeResponse, SharedUploadResponse, TextMessage, TransferProgress, TransferState, WatchPartySyncEvent, }; use super::rate_limiter::RateLimiter; use super::text_analyzer::{analyze_text, AnalyzedText}; use crate::network::DiscoveryBeacon; use reqwest::Client; use sha2::{Digest, Sha256}; use std::collections::VecDeque; use std::path::PathBuf; use std::sync::Arc; use std::time::{Duration, Instant}; use tokio::fs::File; use tokio::io::{AsyncReadExt, AsyncSeekExt}; use tokio::sync::Mutex; use tracing::info; use uuid::Uuid; use walkdir::WalkDir; pub fn compute_file_sha256(path: &std::path::Path) -> Result { let mut file = std::fs::File::open(path)?; let mut hasher = Sha256::new(); let mut buffer = [0u8; 64 * 1024]; loop { let n = std::io::Read::read(&mut file, &mut buffer)?; if n == 0 { break; } hasher.update(&buffer[..n]); } Ok(hex::encode(hasher.finalize())) } const NUM_PARALLEL_STREAMS: usize = 4; // 4 concurrent HTTP/TCP streams fn calculate_adaptive_chunk_size(file_size: u64) -> usize { if file_size < 5 * 1024 * 1024 { 512 * 1024 // 512 KB for small files < 5 MB } else if file_size < 50 * 1024 * 1024 { 2 * 1024 * 1024 // 2 MB for medium files 5-50 MB } else if file_size < 500 * 1024 * 1024 { 4 * 1024 * 1024 // 4 MB for large files 50-500 MB } else { 8 * 1024 * 1024 // 8 MB for huge files >= 500 MB } } #[derive(Clone)] struct ChunkTask { file_path: PathBuf, rel_path: String, offset: u64, length: usize, retries: usize, } pub struct TransferClient { client: Client, } impl TransferClient { pub fn new() -> Self { let client = Client::builder() .timeout(Duration::from_secs(7200)) .connect_timeout(Duration::from_secs(6)) .tcp_nodelay(true) .build() .unwrap_or_default(); Self { client } } /// Ping a remote peer to check health and get beacon info pub async fn ping_peer(&self, ip: &str, port: u16) -> Result { let url = format!("http://{}:{}/api/info", ip, port); let resp = self .client .get(&url) .send() .await .map_err(|e| format!("Failed to connect to {}: {}", url, e))?; if !resp.status().is_success() { return Err(format!("Peer returned status: {}", resp.status())); } let beacon = resp .json::() .await .map_err(|e| format!("Failed to parse response: {}", e))?; Ok(beacon) } /// Measure real-time round-trip latency in milliseconds pub async fn measure_ping(&self, ip: &str, port: u16) -> Option { let start = Instant::now(); if self.ping_peer(ip, port).await.is_ok() { Some(start.elapsed().as_millis() as u64) } else { None } } /// Send quick text / clipboard content to a peer pub async fn send_text( &self, peer_url: &str, my_info: &DiscoveryBeacon, text: &str, pin: Option, ) -> Result { let analyzed = analyze_text(text); let msg = TextMessage { message_id: Uuid::new_v4().to_string(), sender_id: my_info.peer_id.clone(), sender_name: my_info.name.clone(), sender_os: my_info.os.clone(), content: analyzed.content.clone(), content_type: analyzed.content_type.clone(), pin, }; let url = format!("{}/api/transfer/text", peer_url); let resp = self .client .post(&url) .json(&msg) .send() .await .map_err(|e| format!("Failed to contact {}: {}", url, e))?; if !resp.status().is_success() { if resp.status() == reqwest::StatusCode::FORBIDDEN { return Err("Recipient rejected text (Incorrect or missing PIN code)".to_string()); } return Err(format!("Server returned status: {}", resp.status())); } Ok(analyzed) } /// Send files/folders with 4-Stream Parallel Chunking, Adaptive Chunk Sizing, and Auto-Retry Fault Tolerance pub async fn send_files( &self, peer_url: &str, peer_display_name: &str, my_info: &DiscoveryBeacon, paths: &[PathBuf], auto_zip_folders: bool, max_upload_speed_mbps: u32, pin: Option, mut on_progress: F, ) -> Result<(), String> where F: FnMut(TransferProgress) + Send + 'static, { 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); for p in paths { if !p.exists() { return Err(format!("Path does not exist: {:?}", p)); } if p.is_dir() { if auto_zip_folders { let folder_name = p .file_name() .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); // 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)); } else { let base_parent = p.parent().unwrap_or(p); for entry in WalkDir::new(p).into_iter().filter_map(|e| e.ok()) { let entry_path = entry.path(); if entry_path.is_file() { let rel_path = entry_path .strip_prefix(base_parent) .unwrap_or(entry_path) .to_string_lossy() .replace('\\', "/"); let size = entry.metadata().map(|m| m.len()).unwrap_or(0); total_bytes += size; file_entries.push((entry_path.to_path_buf(), rel_path, size)); } } } } else { let file_name = p .file_name() .map(|f| f.to_string_lossy().into_owned()) .unwrap_or_else(|| "unnamed".to_string()); let size = tokio::fs::metadata(p) .await .map(|m| m.len()) .unwrap_or(0); total_bytes += size; file_entries.push((p.clone(), file_name, size)); } } if file_entries.is_empty() { return Err("No files selected for transfer".to_string()); } let file_metas: Vec = file_entries .iter() .map(|(path, rel_path, size)| { let sha256 = compute_file_sha256(path).ok(); FileMetadata { id: Uuid::new_v4().to_string(), name: rel_path.clone(), relative_path: rel_path.clone(), size: *size, is_dir: false, sha256, } }) .collect(); let session_id = Uuid::new_v4().to_string(); let prepare_req = PrepareTransferRequest { session_id: session_id.clone(), sender_id: my_info.peer_id.clone(), sender_name: my_info.name.clone(), sender_os: my_info.os.clone(), files: file_metas.clone(), total_size: total_bytes, pin, }; // Notify UI: Pending confirmation on_progress(TransferProgress { session_id: session_id.clone(), peer_name: peer_display_name.to_string(), is_incoming: false, current_file_index: 0, total_files: file_metas.len(), current_file_name: file_metas.first().map(|f| f.name.clone()).unwrap_or_default(), bytes_transferred: 0, total_bytes, speed_bps: 0, state: TransferState::PendingConfirmation, checksum_verified: false, }); let prepare_url = format!("{}/api/transfer/prepare", peer_url); let resp = self .client .post(&prepare_url) .json(&prepare_req) .send() .await .map_err(|e| format!("Failed to contact peer at {}: {}", prepare_url, e))?; if !resp.status().is_success() { let reason = if resp.status() == reqwest::StatusCode::FORBIDDEN { format!("Transfer rejected by {} (Check PIN code or permissions)", peer_display_name) } else { format!("Peer responded with status {}", resp.status()) }; on_progress(TransferProgress { session_id: session_id.clone(), peer_name: peer_display_name.to_string(), is_incoming: false, current_file_index: 0, total_files: file_metas.len(), current_file_name: String::new(), bytes_transferred: 0, total_bytes, speed_bps: 0, state: TransferState::Rejected, checksum_verified: false, }); return Err(reason); } let prep_resp: PrepareTransferResponse = resp .json() .await .map_err(|e| format!("Invalid prepare response: {}", e))?; let (session_id, received_offsets) = match prep_resp { PrepareTransferResponse::Accepted { session_id: accepted_id, received_offsets } => { info!( "Transfer session {} accepted by {} with {} resumable offset(s)", accepted_id, peer_display_name, received_offsets.len() ); (accepted_id, received_offsets) } PrepareTransferResponse::Rejected { reason } => { return Err(format!("Rejected: {}", reason)); } }; // Build queue of chunk tasks with dynamic adaptive chunk sizing let mut chunk_tasks = VecDeque::new(); let mut already_transferred = 0u64; for (file_path, rel_path, size) in &file_entries { let start_offset = received_offsets.get(rel_path).copied().unwrap_or(0); if start_offset > 0 && start_offset < *size { already_transferred += start_offset; info!("Resuming {} from byte {}", rel_path, start_offset); } let mut curr_offset = start_offset; if curr_offset >= *size { continue; } let chunk_size = calculate_adaptive_chunk_size(*size); while curr_offset < *size { let chunk_len = ((*size - curr_offset) as usize).min(chunk_size); chunk_tasks.push_back(ChunkTask { file_path: file_path.clone(), rel_path: rel_path.clone(), offset: curr_offset, length: chunk_len, retries: 0, }); curr_offset += chunk_len as u64; } } 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 (progress_update_tx, mut progress_update_rx) = tokio::sync::mpsc::unbounded_channel::(); let s_id = session_id.clone(); let p_name = peer_display_name.to_string(); let t_files = file_metas.len(); let p_tx = progress_update_tx.clone(); // Background progress tracker with smoothed EMA speed calculation tokio::spawn(async move { let mut total_sent = 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; if last_time.elapsed().as_millis() >= 350 { let elapsed_secs = last_time.elapsed().as_secs_f64(); let delta = total_sent.saturating_sub(last_bytes); let instant_speed = if elapsed_secs > 0.0 { (delta as f64) / elapsed_secs } else { 0.0 }; smoothed_speed = if smoothed_speed <= 0.0 { instant_speed } else { (smoothed_speed * 0.60) + (instant_speed * 0.40) }; let speed_bps = smoothed_speed as u64; let _ = p_tx.send(TransferProgress { session_id: s_id.clone(), peer_name: p_name.clone(), is_incoming: false, current_file_index: 0, total_files: t_files, current_file_name: fname, bytes_transferred: total_sent, total_bytes, speed_bps, state: TransferState::InProgress, checksum_verified: false, }); last_time = Instant::now(); last_bytes = total_sent; } } }); on_progress(TransferProgress { session_id: session_id.clone(), peer_name: peer_display_name.to_string(), is_incoming: false, current_file_index: 0, total_files: file_metas.len(), current_file_name: file_metas.first().map(|f| f.name.clone()).unwrap_or_default(), bytes_transferred: already_transferred, total_bytes, speed_bps: 0, state: TransferState::InProgress, checksum_verified: false, }); let start_time = Instant::now(); let mut worker_set = tokio::task::JoinSet::new(); // Spawn parallel HTTP chunk streaming worker pool with auto-retry resilience for _ in 0..NUM_PARALLEL_STREAMS { let tasks_c = tasks_mutex.clone(); let client_c = self.client.clone(); let chunk_url = format!("{}/api/transfer/chunk/{}", peer_url, session_id); let limiter_c = rate_limiter.clone(); let tx_c = chunk_tx.clone(); worker_set.spawn(async move { loop { let mut task = { let mut guard = tasks_c.lock().await; guard.pop_front() }; let Some(mut chunk_task) = task.take() else { break; }; let mut file = match File::open(&chunk_task.file_path).await { Ok(f) => f, Err(e) => return Err(format!("Failed to open file {:?}: {}", chunk_task.file_path, e)), }; if let Err(e) = file.seek(std::io::SeekFrom::Start(chunk_task.offset)).await { return Err(format!("Failed to seek in {:?}: {}", chunk_task.file_path, e)); } let mut buf = vec![0u8; chunk_task.length]; if let Err(e) = file.read_exact(&mut buf).await { return Err(format!("Failed to read chunk from {:?}: {}", chunk_task.file_path, e)); } limiter_c.throttle(buf.len()).await; let resp = 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; let is_ok = match &resp { Ok(r) => r.status().is_success(), Err(_) => false, }; if !is_ok { if chunk_task.retries < 3 { chunk_task.retries += 1; tokio::time::sleep(Duration::from_millis(250)).await; tasks_c.lock().await.push_back(chunk_task); continue; } else { let err_msg = match resp { Ok(r) => format!("Server returned error status {}", r.status()), Err(e) => format!("Network transfer error: {}", e), }; return Err(err_msg); } } let _ = tx_c.send((chunk_task.length, chunk_task.rel_path)); } Ok::<(), String>(()) }); } drop(chunk_tx); // Drop main thread copy // Wait for all workers while continuously forwarding progress updates to UI while !worker_set.is_empty() { tokio::select! { res = worker_set.join_next() => { if let Some(join_result) = res { match join_result { Ok(Ok(())) => {} Ok(Err(e)) => { on_progress(TransferProgress { session_id: session_id.clone(), peer_name: peer_display_name.to_string(), is_incoming: false, current_file_index: 0, total_files: file_metas.len(), current_file_name: String::new(), bytes_transferred: 0, total_bytes, speed_bps: 0, state: TransferState::Failed(e.clone()), checksum_verified: false, }); return Err(e); } Err(join_err) => { return Err(format!("Worker thread panicked: {}", join_err)); } } } } Some(progress) = progress_update_rx.recv() => { on_progress(progress); } } } // Drain remaining progress updates while let Ok(progress) = progress_update_rx.try_recv() { on_progress(progress); } // 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; if let Err(e) = finish_resp { info!("Warning: Finalize call response: {}", e); } let elapsed = start_time.elapsed().as_secs_f64(); let actual_uploaded = total_bytes.saturating_sub(already_transferred); let avg_speed = if elapsed > 0.0 { (actual_uploaded as f64 / elapsed) as u64 } else { 0 }; on_progress(TransferProgress { session_id, peer_name: peer_display_name.to_string(), is_incoming: false, current_file_index: file_metas.len(), total_files: file_metas.len(), current_file_name: format!("{} file(s) delivered", file_metas.len()), bytes_transferred: total_bytes, total_bytes, speed_bps: avg_speed, state: TransferState::Completed, checksum_verified: true, }); Ok(()) } /// Fetch list of shared folders from a remote peer #[allow(dead_code)] pub async fn fetch_shared_folders( &self, peer_url: &str, pin: Option<&str>, ) -> Result, String> { let mut url = format!("{}/api/shared/list", peer_url); if let Some(p) = pin.filter(|p| !p.trim().is_empty()) { url = format!("{}?pin={}", url, p); } let resp = self .client .get(&url) .send() .await .map_err(|e| format!("Network error: {}", e))?; if !resp.status().is_success() { return Err(format!("Server returned HTTP {}", resp.status())); } resp.json::>() .await .map_err(|e| format!("Failed to parse shared folders: {}", e)) } /// Fetch file/folder tree of a shared folder pub async fn fetch_shared_tree( &self, peer_url: &str, folder_id: &str, subpath: Option<&str>, pin: Option<&str>, ) -> Result { let url = format!("{}/api/shared/tree/{}", peer_url, folder_id); let mut req = self.client.get(&url); if let Some(p) = subpath.filter(|s| !s.trim().is_empty()) { req = req.query(&[("path", p)]); } if let Some(p) = pin.filter(|p| !p.trim().is_empty()) { req = req.query(&[("pin", p)]); } let resp = req .send() .await .map_err(|e| format!("Network error: {}", e))?; if !resp.status().is_success() { return Err(format!("Server returned HTTP {}", resp.status())); } resp.json::() .await .map_err(|e| format!("Failed to parse folder tree: {}", e)) } /// Fetch raw text of a file (used for Vim Code Viewer) pub async fn fetch_shared_file_text( &self, peer_url: &str, folder_id: &str, subpath: &str, pin: Option<&str>, ) -> Result { let clean = subpath.trim_start_matches('/').trim_start_matches('\\'); let mut url = format!("{}/api/shared/stream/{}/{}", peer_url, folder_id, clean); if let Some(p) = pin.filter(|p| !p.trim().is_empty()) { url = format!("{}?pin={}", url, p); } let resp = self .client .get(&url) .send() .await .map_err(|e| format!("Network error: {}", e))?; if !resp.status().is_success() { return Err(format!("Server returned HTTP {}", resp.status())); } resp.text() .await .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( &self, peer_url: &str, folder_id: &str, subpath: &str, dest_path: &std::path::Path, pin: Option<&str>, ) -> Result { 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 resp = self .client .get(&url) .send() .await .map_err(|e| format!("Network error: {}", e))?; if !resp.status().is_success() { return Err(format!("Download failed: HTTP {}", resp.status())); } let bytes = resp .bytes() .await .map_err(|e| format!("Failed to read stream: {}", e))?; let len = bytes.len() as u64; if let Some(parent) = dest_path.parent() { let _ = tokio::fs::create_dir_all(parent).await; } let mut file = tokio::fs::File::create(dest_path) .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))?; tokio::io::AsyncWriteExt::flush(&mut file) .await .map_err(|e| format!("Could not flush file: {}", e))?; Ok(len) } /// Upload a file directly to a friend's shared folder (if Read & Write allowed) pub async fn upload_to_shared_folder( &self, peer_url: &str, folder_id: &str, subpath: &str, file_path: &std::path::Path, pin: Option<&str>, ) -> Result { let clean = subpath.trim_start_matches('/').trim_start_matches('\\'); let mut url = format!("{}/api/shared/upload/{}/{}", peer_url, folder_id, clean); if let Some(p) = pin.filter(|p| !p.trim().is_empty()) { url = format!("{}?pin={}", url, p); } let file_name = file_path .file_name() .map(|f| f.to_string_lossy().into_owned()) .unwrap_or_else(|| "file".to_string()); let file_bytes = tokio::fs::read(file_path) .await .map_err(|e| format!("Could not read local file: {}", e))?; let part = reqwest::multipart::Part::bytes(file_bytes).file_name(file_name); let form = reqwest::multipart::Form::new().part("file", part); let resp = self .client .post(&url) .multipart(form) .send() .await .map_err(|e| format!("Upload error: {}", e))?; if !resp.status().is_success() { return Err(format!("Upload failed: HTTP {}", resp.status())); } resp.json::() .await .map_err(|e| format!("Failed to parse upload response: {}", e)) } /// Save updated text/code back to a peer's writable shared folder pub async fn save_shared_file( &self, peer_url: &str, folder_id: &str, relative_path: &str, content: &str, pin: Option<&str>, ) -> Result { let clean = relative_path.trim_start_matches('/').trim_start_matches('\\'); let url = format!("{}/api/shared/save/{}/{}", peer_url, folder_id, clean); let req_body = SharedFileSaveRequest { content: content.to_string(), pin: pin.map(|p| p.to_string()), }; let resp = self .client .post(&url) .json(&req_body) .send() .await .map_err(|e| format!("Save error: {}", e))?; if !resp.status().is_success() { return Err(format!("Save failed: HTTP {}", resp.status())); } resp.json::() .await .map_err(|e| format!("Failed to parse save response: {}", e)) } /// Full-text search across remote shared folder pub async fn grep_shared_folder( &self, peer_url: &str, folder_id: &str, query: &str, pin: Option<&str>, subpath: Option<&str>, case_sensitive: bool, ) -> Result { let mut url = format!("{}/api/shared/grep/{}?query={}", peer_url, folder_id, urlencoding::encode(query)); if let Some(p) = pin.filter(|p| !p.trim().is_empty()) { url = format!("{}&pin={}", url, urlencoding::encode(p)); } if let Some(sub) = subpath.filter(|s| !s.trim().is_empty()) { url = format!("{}&subpath={}", url, urlencoding::encode(sub)); } if case_sensitive { url = format!("{}&case_sensitive=true", url); } let resp = self .client .get(&url) .send() .await .map_err(|e| format!("Search error: {}", e))?; if !resp.status().is_success() { return Err(format!("Search failed: HTTP {}", resp.status())); } resp.json::() .await .map_err(|e| format!("Failed to parse search response: {}", e)) } /// Fetch active watch-party state from host peer pub async fn get_watch_party( &self, peer_url: &str, ) -> Result, String> { let url = format!("{}/api/shared/watch-party", peer_url); let resp = self .client .get(&url) .send() .await .map_err(|e| format!("Watch party sync error: {}", e))?; if !resp.status().is_success() { return Err(format!("HTTP {}", resp.status())); } resp.json::>() .await .map_err(|e| format!("Failed to parse watch party event: {}", e)) } /// Send watch-party sync event to host peer pub async fn send_watch_party_event( &self, peer_url: &str, event: &WatchPartySyncEvent, ) -> Result<(), String> { let url = format!("{}/api/shared/watch-party", peer_url); let resp = self .client .post(&url) .json(event) .send() .await .map_err(|e| format!("Failed to broadcast watch party event: {}", e))?; if !resp.status().is_success() { return Err(format!("Broadcast failed: HTTP {}", resp.status())); } Ok(()) } } mod urlencoding { pub fn encode(s: &str) -> String { url::form_urlencoded::byte_serialize(s.as_bytes()).collect() } }