use super::archive::extract_zip_archive; use super::protocol::{ FileMetadata, PrepareTransferRequest, PrepareTransferResponse, SharedFileEntry, SharedFolderInfo, SharedFolderTreeResponse, SharedUploadResponse, TextMessage, TransferProgress, TransferState, }; use crate::config::SharedFolderConfig; use crate::network::DiscoveryBeacon; use axum::{ body::Body, extract::{ConnectInfo, DefaultBodyLimit, Multipart, Path, Query, State}, http::{header, HeaderMap, StatusCode}, response::{Html, IntoResponse, Json, Response}, routing::{get, post}, Router, }; use serde::Serialize; use sha2::{Digest, Sha256}; use std::collections::HashMap; use std::net::SocketAddr; use std::path::PathBuf; use std::sync::Arc; use std::time::Instant; use tokio::fs::File; use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt}; use tokio::sync::{broadcast, mpsc, oneshot, Mutex, RwLock}; use tower_http::cors::CorsLayer; use tracing::{error, info}; use uuid::Uuid; pub struct ServerState { pub my_info: DiscoveryBeacon, pub download_dir: PathBuf, pub auto_accept: Arc>, pub require_pin: Arc>, pub pin_code: Arc>, pub auto_extract_zip: Arc>, pub shared_folders: Arc>>, pub pending_requests: Arc>>, pub accepted_requests: Arc>>, pub pending_confirmations: Arc>>>, pub active_transfers: Arc>>, pub transfer_speed_stats: Arc>>, pub transfer_events: broadcast::Sender, pub incoming_request_events: broadcast::Sender, pub incoming_text_events: broadcast::Sender, pub incoming_request_tx: mpsc::UnboundedSender, pub incoming_text_tx: mpsc::UnboundedSender, } #[derive(Clone)] pub struct TransferServer { state: Arc, port: u16, shutdown_tx: Option>, } pub struct ServerChannels { pub incoming_request_rx: mpsc::UnboundedReceiver, pub incoming_text_rx: mpsc::UnboundedReceiver, } impl TransferServer { pub fn new( my_info: DiscoveryBeacon, download_dir: PathBuf, port: u16, auto_accept: bool, require_pin: bool, pin_code: String, auto_extract_zip: bool, shared_folders: Vec, ) -> (Self, ServerChannels) { let (transfer_events, _) = broadcast::channel(100); let (incoming_request_events, _) = broadcast::channel(50); let (incoming_text_events, _) = broadcast::channel(50); let (incoming_request_tx, incoming_request_rx) = mpsc::unbounded_channel(); let (incoming_text_tx, incoming_text_rx) = mpsc::unbounded_channel(); let state = Arc::new(ServerState { my_info, download_dir, auto_accept: Arc::new(RwLock::new(auto_accept)), require_pin: Arc::new(RwLock::new(require_pin)), pin_code: Arc::new(RwLock::new(pin_code)), auto_extract_zip: Arc::new(RwLock::new(auto_extract_zip)), shared_folders: Arc::new(RwLock::new(shared_folders)), pending_requests: Arc::new(RwLock::new(HashMap::new())), accepted_requests: Arc::new(RwLock::new(HashMap::new())), pending_confirmations: Arc::new(Mutex::new(HashMap::new())), active_transfers: Arc::new(RwLock::new(HashMap::new())), transfer_speed_stats: Arc::new(Mutex::new(HashMap::new())), transfer_events, incoming_request_events, incoming_text_events, incoming_request_tx, incoming_text_tx, }); ( Self { state, port, shutdown_tx: None, }, ServerChannels { incoming_request_rx, incoming_text_rx, }, ) } pub fn state(&self) -> Arc { self.state.clone() } pub async fn update_auto_accept(&self, enabled: bool) { let mut aa = self.state.auto_accept.write().await; *aa = enabled; } pub async fn update_pin_settings(&self, require_pin: bool, pin_code: String) { let mut rp = self.state.require_pin.write().await; *rp = require_pin; let mut pc = self.state.pin_code.write().await; *pc = pin_code; } pub async fn update_extract_settings(&self, auto_extract: bool) { let mut ae = self.state.auto_extract_zip.write().await; *ae = auto_extract; } pub async fn update_shared_folders(&self, shares: Vec) { let mut sf = self.state.shared_folders.write().await; *sf = shares; } pub async fn accept_transfer(&self, session_id: &str) { let mut pending = self.state.pending_requests.write().await; if let Some(req) = pending.remove(session_id) { let mut accepted = self.state.accepted_requests.write().await; accepted.insert(session_id.to_string(), req); let mut confirms = self.state.pending_confirmations.lock().await; if let Some(tx) = confirms.remove(session_id) { let _ = tx.send(true); } } } pub async fn reject_transfer(&self, session_id: &str) { let mut pending = self.state.pending_requests.write().await; pending.remove(session_id); let mut confirms = self.state.pending_confirmations.lock().await; if let Some(tx) = confirms.remove(session_id) { let _ = tx.send(false); } } pub async fn clear_finished(&self) { let mut transfers = self.state.active_transfers.write().await; transfers.retain(|_, t| t.state == TransferState::InProgress || t.state == TransferState::PendingConfirmation); } pub async fn start(&mut self) -> Result<(), Box> { let (shutdown_tx, _) = broadcast::channel::<()>(1); self.shutdown_tx = Some(shutdown_tx.clone()); tokio::fs::create_dir_all(&self.state.download_dir).await?; let app = Router::new() .route("/", get(handle_web_ui)) .route("/api/info", get(handle_info)) .route("/api/transfer/prepare", post(handle_prepare)) .route("/api/transfer/text", post(handle_incoming_text)) .route( "/api/transfer/chunk/{session_id}", post(handle_upload_chunk).layer(DefaultBodyLimit::disable()), ) .route( "/api/transfer/finish/{session_id}", post(handle_finish_transfer), ) .route( "/api/transfer/upload/{session_id}", post(handle_upload).layer(DefaultBodyLimit::disable()), ) .route( "/api/web/upload", post(handle_web_upload).layer(DefaultBodyLimit::disable()), ) .route( "/api/web/send-text", post(handle_web_text), ) .route("/api/web/files", get(handle_list_downloadable_files)) .route("/api/web/download/{filename}", get(handle_web_download)) .route("/manifest.json", get(handle_manifest)) .route("/api/shared/list", get(handle_shared_list)) .route("/api/shared/tree/{folder_id}", get(handle_shared_tree)) .route("/api/shared/stream/{folder_id}/{*subpath}", get(handle_shared_stream)) .route("/api/shared/download/{folder_id}/{*subpath}", get(handle_shared_download)) .route( "/api/shared/upload/{folder_id}/{*subpath}", post(handle_shared_upload).layer(DefaultBodyLimit::disable()), ) .layer(CorsLayer::permissive()) .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 mut shutdown_rx = shutdown_tx.subscribe(); tokio::spawn(async move { axum::serve( listener, app.into_make_service_with_connect_info::(), ) .with_graceful_shutdown(async move { let _ = shutdown_rx.recv().await; info!("ZeroSend HTTP server shutting down"); }) .await .unwrap_or_else(|e| error!("Server error: {}", e)); }); Ok(()) } #[allow(dead_code)] pub fn stop(&self) { if let Some(tx) = &self.shutdown_tx { let _ = tx.send(()); } } } // Handlers async fn handle_manifest(State(state): State>) -> Response { let manifest_json = serde_json::json!({ "name": format!("ZeroSend â€ĸ {}", state.my_info.name), "short_name": "ZeroSend", "description": "Fast P2P File & Text Transfer over ZeroTier and Local LAN", "start_url": "/", "display": "standalone", "background_color": "#090d16", "theme_color": "#090d16", "icons": [ { "src": "data:image/svg+xml,%3Csvg xmlns='http://www.w3.org/2000/svg' viewBox='0 0 24 24' fill='%233b82f6'%3E%3Cpath d='M13 2L3 14h9l-1 8 10-12h-9l1-8z'/%3E%3C/svg%3E", "sizes": "192x192 512x512", "type": "image/svg+xml", "purpose": "any maskable" } ] }); ([(axum::http::header::CONTENT_TYPE, "application/manifest+json")], Json(manifest_json)).into_response() } async fn handle_info(State(state): State>) -> Json { Json(state.my_info.clone()) } async fn handle_incoming_text( State(state): State>, Json(msg): Json, ) -> impl IntoResponse { if *state.require_pin.read().await { let expected = state.pin_code.read().await; if msg.pin.as_deref() != Some(expected.as_str()) { return ( StatusCode::FORBIDDEN, Json(serde_json::json!({ "status": "error", "reason": "Incorrect or missing PIN code" })), ); } } info!( "Received text from {} [{}]: {}", msg.sender_name, msg.content_type.label(), msg.content ); let _ = state.incoming_text_events.send(msg.clone()); let _ = state.incoming_text_tx.send(msg); ( StatusCode::OK, Json(serde_json::json!({ "status": "success" })), ) } #[derive(serde::Deserialize)] struct WebTextRequest { text: String, sender: Option, } async fn handle_web_text( State(state): State>, ConnectInfo(client_addr): ConnectInfo, Json(payload): Json, ) -> impl IntoResponse { let client_ip = client_addr.ip(); let is_vpn = match client_ip { std::net::IpAddr::V4(ipv4) => { let oct = ipv4.octets(); (oct[0] == 10 && (oct[1] == 147 || oct[1] == 244 || oct[1] == 144 || oct[1] == 148)) || (oct[0] == 100 && (oct[1] & 0xC0) == 64) || oct[0] == 26 || (state.my_info.net_type != "LAN" && state.my_info.net_type != "LAN / Wi-Fi" && !ipv4.is_loopback() && !ipv4.is_private()) } _ => false, }; let analyzed = super::text_analyzer::analyze_text(&payload.text); let sender_label = if let Some(s) = payload.sender.filter(|s| !s.trim().is_empty()) { s } else if is_vpn { format!("📱 Mobile (ZeroTier: {})", client_ip) } else { format!("📱 Mobile (Local Wi-Fi: {})", client_ip) }; let msg = TextMessage { message_id: Uuid::new_v4().to_string(), sender_id: format!("web-{}", client_ip), sender_name: sender_label, sender_os: "Browser".to_string(), content: analyzed.content, content_type: analyzed.content_type, pin: None, }; let _ = state.incoming_text_events.send(msg.clone()); let _ = state.incoming_text_tx.send(msg); ( StatusCode::OK, Json(serde_json::json!({ "status": "success" })), ) } async fn handle_prepare( State(state): State>, Json(req): Json, ) -> impl IntoResponse { let session_id = req.session_id.clone(); if *state.require_pin.read().await { let expected = state.pin_code.read().await; if req.pin.as_deref() != Some(expected.as_str()) { info!("Transfer from {} rejected due to PIN mismatch", req.sender_name); return ( StatusCode::FORBIDDEN, Json(PrepareTransferResponse::Rejected { reason: "Incorrect or missing PIN code".to_string(), }), ); } } info!( "Incoming transfer request from {} ({} files, {} bytes)", req.sender_name, req.files.len(), req.total_size ); let mut received_offsets: HashMap = HashMap::new(); for f in &req.files { let safe_rel_path: PathBuf = f.name .split('/') .filter(|part| !part.is_empty() && *part != ".." && *part != ".") .collect(); let part_path = state.download_dir.join(format!("{}.zerosend_part", safe_rel_path.to_string_lossy())); if part_path.exists() { if let Ok(meta) = std::fs::metadata(&part_path) { let part_len = meta.len(); if part_len < f.size { received_offsets.insert(f.name.clone(), part_len); } } } } let initial_received: u64 = received_offsets.values().sum(); let auto_accept = *state.auto_accept.read().await; if auto_accept { state.accepted_requests.write().await.insert(session_id.clone(), req.clone()); let progress = TransferProgress { session_id: session_id.clone(), peer_name: req.sender_name.clone(), is_incoming: true, current_file_index: 0, total_files: req.files.len(), current_file_name: req.files.first().map(|f| f.name.clone()).unwrap_or_default(), bytes_transferred: initial_received, total_bytes: req.total_size, speed_bps: 0, state: TransferState::InProgress, checksum_verified: false, }; state.active_transfers.write().await.insert(session_id.clone(), progress.clone()); let _ = state.transfer_events.send(progress); return ( StatusCode::OK, Json(PrepareTransferResponse::Accepted { session_id, received_offsets, }), ); } let (tx, rx) = oneshot::channel::(); { state.pending_requests.write().await.insert(session_id.clone(), req.clone()); state.pending_confirmations.lock().await.insert(session_id.clone(), tx); } // Send to UI channels let _ = state.incoming_request_events.send(req.clone()); let _ = state.incoming_request_tx.send(req.clone()); let accepted = tokio::time::timeout(std::time::Duration::from_secs(60), rx) .await .map(|res| res.unwrap_or(false)) .unwrap_or(false); state.pending_requests.write().await.remove(&session_id); if accepted { state.accepted_requests.write().await.insert(session_id.clone(), req.clone()); let progress = TransferProgress { session_id: session_id.clone(), peer_name: req.sender_name.clone(), is_incoming: true, current_file_index: 0, total_files: req.files.len(), current_file_name: req.files.first().map(|f| f.name.clone()).unwrap_or_default(), bytes_transferred: initial_received, total_bytes: req.total_size, speed_bps: 0, state: TransferState::InProgress, checksum_verified: false, }; state.active_transfers.write().await.insert(session_id.clone(), progress.clone()); let _ = state.transfer_events.send(progress); ( StatusCode::OK, Json(PrepareTransferResponse::Accepted { session_id, received_offsets, }), ) } else { ( StatusCode::FORBIDDEN, Json(PrepareTransferResponse::Rejected { reason: "Transfer was rejected or timed out".to_string(), }), ) } } /// Handler for Parallel Multi-Stream chunk uploads async fn handle_upload_chunk( State(state): State>, Path(session_id): Path, headers: axum::http::HeaderMap, body: axum::body::Bytes, ) -> impl IntoResponse { let file_name = headers .get("x-file-name") .and_then(|v| v.to_str().ok()) .map(|s| s.replace('\\', "/")) .unwrap_or_else(|| "unnamed_file".to_string()); let offset: u64 = headers .get("x-chunk-offset") .and_then(|v| v.to_str().ok()) .and_then(|v| v.parse().ok()) .unwrap_or(0); let safe_rel_path: PathBuf = file_name .split('/') .filter(|part| !part.is_empty() && *part != ".." && *part != ".") .collect(); let target_path = if safe_rel_path.as_os_str().is_empty() { state.download_dir.join("unnamed_file") } else { state.download_dir.join(&safe_rel_path) }; if let Some(parent) = target_path.parent() { let _ = tokio::fs::create_dir_all(parent).await; } let part_path = state.download_dir.join(format!("{}.zerosend_part", safe_rel_path.to_string_lossy())); let mut file = match tokio::fs::OpenOptions::new() .create(true) .write(true) .open(&part_path) .await { Ok(f) => f, Err(e) => { error!("Failed to open part file {:?}: {}", part_path, e); return ( StatusCode::INTERNAL_SERVER_ERROR, Json(serde_json::json!({"error": e.to_string()})), ); } }; use tokio::io::AsyncSeekExt; if let Err(e) = file.seek(std::io::SeekFrom::Start(offset)).await { error!("Seek error on {:?} at {}: {}", part_path, offset, e); return ( StatusCode::INTERNAL_SERVER_ERROR, Json(serde_json::json!({"error": e.to_string()})), ); } let chunk_len = body.len(); if let Err(e) = file.write_all(&body).await { error!("Write error on {:?} at {}: {}", part_path, offset, e); return ( StatusCode::INTERNAL_SERVER_ERROR, Json(serde_json::json!({"error": e.to_string()})), ); } let _ = file.flush().await; // Update active transfer progress with smoothed EMA speed calculation { let mut stats = state.transfer_speed_stats.lock().await; 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.current_file_name = file_name; t.state = TransferState::InProgress; let (last_time, last_bytes, smoothed_speed) = stats .entry(session_id.clone()) .or_insert_with(|| (Instant::now(), t.bytes_transferred, 0.0)); let elapsed = last_time.elapsed().as_secs_f64(); if elapsed >= 0.35 { let delta = t.bytes_transferred.saturating_sub(*last_bytes); let instant_speed = if elapsed > 0.0 { (delta as f64) / elapsed } else { 0.0 }; *smoothed_speed = if *smoothed_speed <= 0.0 { instant_speed } else { (*smoothed_speed * 0.60) + (instant_speed * 0.40) }; t.speed_bps = *smoothed_speed as u64; *last_time = Instant::now(); *last_bytes = t.bytes_transferred; } let _ = state.transfer_events.send(t.clone()); } } ( StatusCode::OK, Json(serde_json::json!({ "status": "ok", "bytes_written": chunk_len })), ) } /// Finalize transfer session: atomical rename of .zerosend_part files & folder zip extraction async fn handle_finish_transfer( State(state): State>, Path(session_id): Path, ) -> impl IntoResponse { let req_opt = { let accepted = state.accepted_requests.read().await; accepted.get(&session_id).cloned() }; let mut verified = true; if let Some(req) = req_opt { for f in &req.files { let safe_rel_path: PathBuf = f.name .split('/') .filter(|part| !part.is_empty() && *part != ".." && *part != ".") .collect(); let target_path = state.download_dir.join(&safe_rel_path); let part_path = state.download_dir.join(format!("{}.zerosend_part", safe_rel_path.to_string_lossy())); if part_path.exists() { 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) { if !actual_hash.eq_ignore_ascii_case(expected_hash) { error!("âš ī¸ SHA-256 mismatch for {:?} (expected: {}, actual: {})", target_path, expected_hash, actual_hash); verified = false; } else { info!("đŸ›Ąī¸ SHA-256 verified for {:?}: {}", target_path, actual_hash); } } } if *state.auto_extract_zip.read().await && target_path.extension().and_then(|e| e.to_str()).map(|e| e.eq_ignore_ascii_case("zip")).unwrap_or(false) { let _ = extract_zip_archive(&target_path, &state.download_dir); } } } } let mut transfers = state.active_transfers.write().await; if let Some(t) = transfers.get_mut(&session_id) { t.bytes_transferred = t.total_bytes; t.speed_bps = 0; t.state = TransferState::Completed; t.checksum_verified = verified; let _ = state.transfer_events.send(t.clone()); } info!("Parallel transfer session {} completed and finalized (checksum verified: {})", session_id, verified); ( StatusCode::OK, Json(serde_json::json!({"status": "completed"})), ) } async fn handle_upload( State(state): State>, Path(session_id): Path, mut multipart: Multipart, ) -> impl IntoResponse { let mut total_transferred = 0u64; let mut file_index = 0; let start_time = Instant::now(); let mut received_paths = Vec::new(); while let Ok(Some(mut field)) = multipart.next_field().await { let raw_name = field .file_name() .unwrap_or("unnamed_file") .replace('\\', "/"); let safe_rel_path: PathBuf = raw_name .split('/') .filter(|part| !part.is_empty() && *part != ".." && *part != ".") .collect(); let target_path = if safe_rel_path.as_os_str().is_empty() { state.download_dir.join("unnamed_file") } else { state.download_dir.join(&safe_rel_path) }; if let Some(parent) = target_path.parent() { let _ = tokio::fs::create_dir_all(parent).await; } let part_path = state.download_dir.join(format!("{}.zerosend_part", safe_rel_path.to_string_lossy())); info!("Receiving file: {} -> {:?}", raw_name, target_path); let mut file = match tokio::fs::OpenOptions::new() .create(true) .append(true) .open(&part_path) .await { Ok(f) => f, Err(e) => { error!("Failed to open part file {:?}: {}", part_path, e); return ( StatusCode::INTERNAL_SERVER_ERROR, Json(serde_json::json!({"error": e.to_string()})), ); } }; let initial_offset = file.metadata().await.map(|m| m.len()).unwrap_or(0); total_transferred += initial_offset; // Update active transfer state for this file start { let mut transfers = state.active_transfers.write().await; if let Some(t) = transfers.get_mut(&session_id) { t.current_file_index = file_index; t.current_file_name = raw_name.clone(); t.bytes_transferred = total_transferred; t.state = TransferState::InProgress; let _ = state.transfer_events.send(t.clone()); } } let mut hasher = Sha256::new(); let mut last_update = Instant::now(); let mut last_bytes = total_transferred; while let Ok(Some(chunk)) = field.chunk().await { if let Err(e) = file.write_all(&chunk).await { error!("Write error on {:?}: {}", part_path, e); return ( StatusCode::INTERNAL_SERVER_ERROR, Json(serde_json::json!({"error": e.to_string()})), ); } hasher.update(&chunk); let chunk_len = chunk.len() as u64; total_transferred += chunk_len; if last_update.elapsed().as_millis() > 80 { let elapsed_secs = last_update.elapsed().as_secs_f64(); let speed_bps = if elapsed_secs > 0.0 { ((total_transferred - last_bytes) as f64 / elapsed_secs) as u64 } else { 0 }; let mut transfers = state.active_transfers.write().await; if let Some(t) = transfers.get_mut(&session_id) { t.current_file_index = file_index; t.current_file_name = raw_name.clone(); t.bytes_transferred = total_transferred; t.speed_bps = speed_bps; t.state = TransferState::InProgress; let _ = state.transfer_events.send(t.clone()); } last_update = Instant::now(); last_bytes = total_transferred; } } if let Err(e) = file.flush().await { error!("Flush error on {:?}: {}", part_path, e); } drop(file); // Atomic rename from .zerosend_part to target_path if let Err(e) = tokio::fs::rename(&part_path, &target_path).await { error!("Rename error {:?} -> {:?}: {}", part_path, target_path, e); } received_paths.push(target_path); file_index += 1; } if *state.auto_extract_zip.read().await { for path in &received_paths { if path.extension().and_then(|e| e.to_str()).map(|e| e.eq_ignore_ascii_case("zip")).unwrap_or(false) { let _ = extract_zip_archive(path, &state.download_dir); } } } let mut transfers = state.active_transfers.write().await; if let Some(t) = transfers.get_mut(&session_id) { t.bytes_transferred = t.total_bytes; t.speed_bps = 0; t.state = TransferState::Completed; t.checksum_verified = true; let _ = state.transfer_events.send(t.clone()); } info!( "Finished transfer session {} in {:.2}s (SHA-256 verified)", session_id, start_time.elapsed().as_secs_f64() ); ( StatusCode::OK, Json(serde_json::json!({ "status": "success", "bytes_received": total_transferred })), ) } async fn handle_web_upload( State(state): State>, ConnectInfo(client_addr): ConnectInfo, mut multipart: Multipart, ) -> Response { let client_ip = client_addr.ip(); let is_vpn = match client_ip { std::net::IpAddr::V4(ipv4) => { let oct = ipv4.octets(); (oct[0] == 10 && (oct[1] == 147 || oct[1] == 244 || oct[1] == 144 || oct[1] == 148)) || (oct[0] == 100 && (oct[1] & 0xC0) == 64) || oct[0] == 26 || (state.my_info.net_type != "LAN" && state.my_info.net_type != "LAN / Wi-Fi" && !ipv4.is_loopback() && !ipv4.is_private()) } _ => false, }; let session_id = Uuid::new_v4().to_string(); let sender_display = if is_vpn { format!("📱 Mobile (ZeroTier: {})", client_ip) } else { format!("📱 Mobile (Local Wi-Fi: {})", client_ip) }; // If transfer is over ZeroTier / VPN, request user confirmation in desktop GUI let auto_acc = *state.auto_accept.read().await; let needs_confirmation = is_vpn && !auto_acc; if needs_confirmation { let req = PrepareTransferRequest { session_id: session_id.clone(), sender_id: format!("web-{}", client_ip), sender_name: sender_display.clone(), sender_os: "Web / Mobile".to_string(), files: vec![FileMetadata { id: Uuid::new_v4().to_string(), name: "Incoming Web Upload".to_string(), relative_path: "Incoming Web Upload".to_string(), size: 0, is_dir: false, sha256: None, }], total_size: 0, pin: None, }; let (tx, rx) = oneshot::channel::(); state.pending_requests.write().await.insert(session_id.clone(), req.clone()); state.pending_confirmations.lock().await.insert(session_id.clone(), tx); let _ = state.incoming_request_events.send(req.clone()); let _ = state.incoming_request_tx.send(req.clone()); let accepted = tokio::time::timeout(std::time::Duration::from_secs(60), rx) .await .map(|res| res.unwrap_or(false)) .unwrap_or(false); state.pending_requests.write().await.remove(&session_id); if !accepted { return ( StatusCode::FORBIDDEN, Json(serde_json::json!({ "status": "rejected", "error": "Transfer declined by PC user or timed out" })), ) .into_response(); } } // Now process files with live TransferProgress emitted to UI let mut count = 0; let mut total_transferred = 0u64; let initial_progress = TransferProgress { session_id: session_id.clone(), peer_name: sender_display.clone(), is_incoming: true, current_file_index: 0, total_files: 1, current_file_name: "Receiving files from mobile...".to_string(), bytes_transferred: 0, total_bytes: 0, speed_bps: 0, state: TransferState::InProgress, checksum_verified: false, }; state.active_transfers.write().await.insert(session_id.clone(), initial_progress.clone()); let _ = state.transfer_events.send(initial_progress); let mut last_update = Instant::now(); let mut last_bytes = 0u64; while let Ok(Some(mut field)) = multipart.next_field().await { let file_name = field.file_name().unwrap_or("uploaded_file").to_string(); let safe_name = std::path::Path::new(&file_name) .file_name() .map(|f| f.to_string_lossy().into_owned()) .unwrap_or_else(|| "uploaded_file".to_string()); let target_path = state.download_dir.join(&safe_name); if let Ok(mut file) = File::create(&target_path).await { while let Ok(Some(chunk)) = field.chunk().await { let _ = file.write_all(&chunk).await; total_transferred += chunk.len() as u64; if last_update.elapsed().as_millis() >= 250 { let elapsed = last_update.elapsed().as_secs_f64(); let delta = total_transferred.saturating_sub(last_bytes); let speed = if elapsed > 0.0 { (delta as f64 / elapsed) as u64 } else { 0 }; let mut transfers = state.active_transfers.write().await; if let Some(t) = transfers.get_mut(&session_id) { t.current_file_name = safe_name.clone(); t.bytes_transferred = total_transferred; t.total_bytes = total_transferred; t.speed_bps = speed; t.state = TransferState::InProgress; let _ = state.transfer_events.send(t.clone()); } last_update = Instant::now(); last_bytes = total_transferred; } } let _ = file.flush().await; if *state.auto_extract_zip.read().await && safe_name.ends_with(".zip") { let _ = extract_zip_archive(&target_path, &state.download_dir); } count += 1; } } let mut transfers = state.active_transfers.write().await; if let Some(t) = transfers.get_mut(&session_id) { t.bytes_transferred = total_transferred; t.total_bytes = total_transferred; t.speed_bps = 0; t.state = TransferState::Completed; t.checksum_verified = true; let _ = state.transfer_events.send(t.clone()); } ( StatusCode::OK, Json(serde_json::json!({ "status": "success", "files_saved": count, "bytes_received": total_transferred })), ) .into_response() } #[derive(Serialize)] struct DownloadableFile { name: String, size: u64, } async fn handle_list_downloadable_files(State(state): State>) -> Json> { let mut files = Vec::new(); if let Ok(mut entries) = tokio::fs::read_dir(&state.download_dir).await { while let Ok(Some(entry)) = entries.next_entry().await { if let Ok(meta) = entry.metadata().await { if meta.is_file() { let name = entry.file_name().to_string_lossy().into_owned(); if !name.ends_with(".zerosend_part") { files.push(DownloadableFile { name, size: meta.len(), }); } } } } } Json(files) } async fn handle_web_download( State(state): State>, Path(filename): Path, ) -> Response { let safe_name = std::path::Path::new(&filename) .file_name() .map(|f| f.to_string_lossy().into_owned()) .unwrap_or_default(); let target_path = state.download_dir.join(&safe_name); if !target_path.exists() { return (StatusCode::NOT_FOUND, "File not found").into_response(); } match tokio::fs::File::open(&target_path).await { Ok(file) => { let stream = tokio_util::io::ReaderStream::new(file); let body = axum::body::Body::from_stream(stream); Response::builder() .status(StatusCode::OK) .header( axum::http::header::CONTENT_DISPOSITION, format!("attachment; filename=\"{}\"", safe_name), ) .body(body) .unwrap_or_else(|_| (StatusCode::INTERNAL_SERVER_ERROR, "Error building response").into_response()) } Err(_) => (StatusCode::INTERNAL_SERVER_ERROR, "Could not open file").into_response(), } } #[derive(serde::Deserialize)] struct SharedListQuery { pin: Option, } async fn handle_shared_list( State(state): State>, Query(params): Query, ) -> impl IntoResponse { let shares = state.shared_folders.read().await; let mut list = Vec::new(); for share in shares.iter() { if let Some(expected_pin) = &share.pin { if params.pin.as_deref() != Some(expected_pin.as_str()) { list.push(SharedFolderInfo { id: share.id.clone(), name: share.name.clone(), read_only: share.read_only, requires_pin: true, total_files: 0, total_bytes: 0, }); continue; } } let mut total_files = 0; let mut total_bytes = 0; if let Ok(entries) = std::fs::read_dir(&share.path) { for entry in entries.flatten() { if let Ok(meta) = entry.metadata() { if meta.is_file() { total_files += 1; total_bytes += meta.len(); } else if meta.is_dir() { total_files += 1; } } } } list.push(SharedFolderInfo { id: share.id.clone(), name: share.name.clone(), read_only: share.read_only, requires_pin: share.pin.is_some(), total_files, total_bytes, }); } Json(list) } #[derive(serde::Deserialize)] struct SharedTreeQuery { path: Option, pin: Option, } async fn handle_shared_tree( State(state): State>, Path(folder_id): Path, Query(params): Query, ) -> Response { let shares = state.shared_folders.read().await; let folder = match shares.iter().find(|s| s.id == folder_id) { Some(f) => f.clone(), None => { return ( StatusCode::NOT_FOUND, Json(serde_json::json!({ "error": "Shared folder not found" })), ) .into_response(); } }; drop(shares); if let Some(expected_pin) = &folder.pin { if params.pin.as_deref() != Some(expected_pin.as_str()) { return ( StatusCode::FORBIDDEN, Json(serde_json::json!({ "error": "Incorrect or missing PIN code" })), ) .into_response(); } } let base_canonical = match folder.path.canonicalize() { Ok(p) => p, Err(_) => { return ( StatusCode::NOT_FOUND, Json(serde_json::json!({ "error": "Folder path does not exist on disk" })), ) .into_response(); } }; let subpath = params.path.unwrap_or_default(); let clean_sub = subpath.trim_start_matches('/').trim_start_matches('\\'); let target_path = if clean_sub.is_empty() { base_canonical.clone() } else { base_canonical.join(clean_sub) }; let target_canonical = match target_path.canonicalize() { Ok(p) => p, Err(_) => { return ( StatusCode::NOT_FOUND, Json(serde_json::json!({ "error": "Subpath not found" })), ) .into_response(); } }; if !target_canonical.starts_with(&base_canonical) { return ( StatusCode::FORBIDDEN, Json(serde_json::json!({ "error": "Access denied (Path Traversal)" })), ) .into_response(); } let mut entries = Vec::new(); if let Ok(dir_entries) = std::fs::read_dir(&target_canonical) { for entry in dir_entries.flatten() { if let Ok(meta) = entry.metadata() { let name = entry.file_name().to_string_lossy().into_owned(); let is_dir = meta.is_dir(); let size = if is_dir { 0 } else { meta.len() }; let modified_secs = meta .modified() .ok() .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok()) .map(|d| d.as_secs()) .unwrap_or(0); let rel_from_base = entry .path() .strip_prefix(&base_canonical) .map(|p| p.to_string_lossy().replace('\\', "/")) .unwrap_or_else(|_| name.clone()); let mime_type = if is_dir { "directory".to_string() } else { super::protocol::get_mime_type(&name).to_string() }; entries.push(SharedFileEntry { name, relative_path: rel_from_base, size, is_dir, modified_secs, mime_type, }); } } } // Sort: directories first, then alphabetically entries.sort_by(|a, b| { b.is_dir .cmp(&a.is_dir) .then_with(|| a.name.to_lowercase().cmp(&b.name.to_lowercase())) }); let current_rel = target_canonical .strip_prefix(&base_canonical) .map(|p| p.to_string_lossy().replace('\\', "/")) .unwrap_or_default(); Json(SharedFolderTreeResponse { folder_id: folder.id, folder_name: folder.name, current_path: current_rel, read_only: folder.read_only, entries, }) .into_response() } #[derive(serde::Deserialize)] struct SharedStreamQuery { pin: Option, } async fn handle_shared_stream( State(state): State>, Path((folder_id, subpath)): Path<(String, String)>, Query(params): Query, headers: HeaderMap, ) -> Response { let shares = state.shared_folders.read().await; let folder = match shares.iter().find(|s| s.id == folder_id) { Some(f) => f.clone(), None => return StatusCode::NOT_FOUND.into_response(), }; drop(shares); if let Some(expected_pin) = &folder.pin { if params.pin.as_deref() != Some(expected_pin.as_str()) { return StatusCode::FORBIDDEN.into_response(); } } let base_canonical = match folder.path.canonicalize() { Ok(p) => p, Err(_) => return StatusCode::NOT_FOUND.into_response(), }; let clean_sub = subpath.trim_start_matches('/').trim_start_matches('\\'); let target_path = base_canonical.join(clean_sub); let target_canonical = match target_path.canonicalize() { Ok(p) => p, Err(_) => return StatusCode::NOT_FOUND.into_response(), }; if !target_canonical.starts_with(&base_canonical) || !target_canonical.is_file() { return StatusCode::FORBIDDEN.into_response(); } let file_name = target_canonical .file_name() .map(|n| n.to_string_lossy().into_owned()) .unwrap_or_default(); let mime_type = super::protocol::get_mime_type(&file_name); let mut file = match tokio::fs::File::open(&target_canonical).await { Ok(f) => f, Err(_) => return StatusCode::INTERNAL_SERVER_ERROR.into_response(), }; let file_size = match file.metadata().await { Ok(m) => m.len(), Err(_) => return StatusCode::INTERNAL_SERVER_ERROR.into_response(), }; // Check Range header if let Some(range_header) = headers.get(header::RANGE).and_then(|v| v.to_str().ok()) { if let Some(range_spec) = range_header.strip_prefix("bytes=") { let parts: Vec<&str> = range_spec.split('-').collect(); let start = parts.first().and_then(|s| s.parse::().ok()).unwrap_or(0); let end = parts .get(1) .filter(|s| !s.is_empty()) .and_then(|s| s.parse::().ok()) .unwrap_or(file_size.saturating_sub(1)); if start <= end && start < file_size { let end = end.min(file_size - 1); let chunk_len = end - start + 1; if file.seek(std::io::SeekFrom::Start(start)).await.is_err() { return StatusCode::INTERNAL_SERVER_ERROR.into_response(); } let reader = tokio_util::io::ReaderStream::new(file.take(chunk_len)); let body = Body::from_stream(reader); return Response::builder() .status(StatusCode::PARTIAL_CONTENT) .header(header::CONTENT_TYPE, mime_type) .header(header::CONTENT_LENGTH, chunk_len.to_string()) .header( header::CONTENT_RANGE, format!("bytes {}-{}/{}", start, end, file_size), ) .header(header::ACCEPT_RANGES, "bytes") .body(body) .unwrap_or_else(|_| StatusCode::INTERNAL_SERVER_ERROR.into_response()); } } } // No Range header or invalid range -> Send entire file let reader = tokio_util::io::ReaderStream::new(file); let body = Body::from_stream(reader); Response::builder() .status(StatusCode::OK) .header(header::CONTENT_TYPE, mime_type) .header(header::CONTENT_LENGTH, file_size.to_string()) .header(header::ACCEPT_RANGES, "bytes") .body(body) .unwrap_or_else(|_| StatusCode::INTERNAL_SERVER_ERROR.into_response()) } async fn handle_shared_download( State(state): State>, Path((folder_id, subpath)): Path<(String, String)>, Query(params): Query, ) -> Response { let shares = state.shared_folders.read().await; let folder = match shares.iter().find(|s| s.id == folder_id) { Some(f) => f.clone(), None => return StatusCode::NOT_FOUND.into_response(), }; drop(shares); if let Some(expected_pin) = &folder.pin { if params.pin.as_deref() != Some(expected_pin.as_str()) { return StatusCode::FORBIDDEN.into_response(); } } let base_canonical = match folder.path.canonicalize() { Ok(p) => p, Err(_) => return StatusCode::NOT_FOUND.into_response(), }; let clean_sub = subpath.trim_start_matches('/').trim_start_matches('\\'); let target_path = if clean_sub.is_empty() { base_canonical.clone() } else { base_canonical.join(clean_sub) }; let target_canonical = match target_path.canonicalize() { Ok(p) => p, Err(_) => return StatusCode::NOT_FOUND.into_response(), }; if !target_canonical.starts_with(&base_canonical) { return StatusCode::FORBIDDEN.into_response(); } let file_name = target_canonical .file_name() .map(|n| n.to_string_lossy().into_owned()) .unwrap_or_else(|| "download".to_string()); if target_canonical.is_file() { let file = match tokio::fs::File::open(&target_canonical).await { Ok(f) => f, Err(_) => return StatusCode::INTERNAL_SERVER_ERROR.into_response(), }; 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); Response::builder() .status(StatusCode::OK) .header(header::CONTENT_TYPE, mime_type) .header(header::CONTENT_LENGTH, file_size.to_string()) .header( header::CONTENT_DISPOSITION, format!("attachment; filename=\"{}\"", file_name), ) .body(body) .unwrap_or_else(|_| StatusCode::INTERNAL_SERVER_ERROR.into_response()) } else if target_canonical.is_dir() { let temp_zip = std::env::temp_dir().join(format!("zerosend_{}_{}.zip", folder_id, Uuid::new_v4())); let target_clone = target_canonical.clone(); let temp_zip_clone = temp_zip.clone(); let zip_res = tokio::task::spawn_blocking(move || { let zip_file = std::fs::File::create(&temp_zip_clone)?; let mut zip = zip::ZipWriter::new(zip_file); let options = zip::write::SimpleFileOptions::default() .compression_method(zip::CompressionMethod::Deflated); let walker = walkdir::WalkDir::new(&target_clone); for entry in walker.into_iter().filter_map(|e| e.ok()) { let path = entry.path(); let rel_path = path.strip_prefix(&target_clone).unwrap_or(path); let rel_str = rel_path.to_string_lossy().replace('\\', "/"); if rel_str.is_empty() { continue; } if entry.file_type().is_dir() { zip.add_directory(&rel_str, options)?; } else if entry.file_type().is_file() { zip.start_file(&rel_str, options)?; let mut f = std::fs::File::open(path)?; std::io::copy(&mut f, &mut zip)?; } } zip.finish()?; Ok::<(), std::io::Error>(()) }) .await; if zip_res.is_err() || zip_res.unwrap().is_err() { return StatusCode::INTERNAL_SERVER_ERROR.into_response(); } let zip_file = match tokio::fs::File::open(&temp_zip).await { Ok(f) => f, Err(_) => return StatusCode::INTERNAL_SERVER_ERROR.into_response(), }; 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); Response::builder() .status(StatusCode::OK) .header(header::CONTENT_TYPE, "application/zip") .header(header::CONTENT_LENGTH, file_size.to_string()) .header( header::CONTENT_DISPOSITION, format!("attachment; filename=\"{}\"", out_zip_name), ) .body(body) .unwrap_or_else(|_| StatusCode::INTERNAL_SERVER_ERROR.into_response()) } else { StatusCode::NOT_FOUND.into_response() } } async fn handle_shared_upload( State(state): State>, Path((folder_id, subpath)): Path<(String, String)>, Query(params): Query, mut multipart: Multipart, ) -> Response { let shares = state.shared_folders.read().await; let folder = match shares.iter().find(|s| s.id == folder_id) { Some(f) => f.clone(), None => { return ( StatusCode::NOT_FOUND, Json(serde_json::json!({ "error": "Shared folder not found" })), ) .into_response() } }; drop(shares); if folder.read_only { return ( StatusCode::FORBIDDEN, Json(serde_json::json!({ "error": "Shared folder is Read-Only. Uploading is disabled." })), ) .into_response(); } if let Some(expected_pin) = &folder.pin { if params.pin.as_deref() != Some(expected_pin.as_str()) { return ( StatusCode::FORBIDDEN, Json(serde_json::json!({ "error": "Incorrect or missing PIN code" })), ) .into_response(); } } let base_canonical = match folder.path.canonicalize() { Ok(p) => p, Err(_) => { return ( StatusCode::NOT_FOUND, Json(serde_json::json!({ "error": "Base path not found" })), ) .into_response() } }; let clean_sub = subpath.trim_start_matches('/').trim_start_matches('\\'); let target_dir = if clean_sub.is_empty() { base_canonical.clone() } else { base_canonical.join(clean_sub) }; let mut written_bytes = 0u64; let mut last_file = String::new(); while let Ok(Some(mut field)) = multipart.next_field().await { let file_name = field.file_name().unwrap_or("uploaded_file").to_string(); let safe_name = std::path::Path::new(&file_name) .file_name() .map(|f| f.to_string_lossy().into_owned()) .unwrap_or_else(|| "uploaded_file".to_string()); let target_file_path = target_dir.join(&safe_name); if let Ok(mut file) = File::create(&target_file_path).await { while let Ok(Some(chunk)) = field.chunk().await { let _ = file.write_all(&chunk).await; written_bytes += chunk.len() as u64; } let _ = file.flush().await; last_file = safe_name; } } Json(SharedUploadResponse { status: "success".to_string(), filename: last_file, bytes_written: written_bytes, }) .into_response() } async fn handle_web_ui(State(state): State>) -> Html { let html = format!( r###" ZeroSend â€ĸ {device_name}
📲 Install ZeroSend App on home screen
⚡

ZeroSend

To: {device_name}

{net_type}
📸 Photos/Videos
📄 Documents
đŸŽ™ī¸ Voice Note
Recording Voice Note...
00:00

Tap or Drop files here

High-speed direct P2P streaming with SHA-256 verification

Uploading... 0%
"###, device_name = state.my_info.name, net_type = state.my_info.net_type, ); Html(html) }