144 lines
8.0 KiB
Rust
Executable File
144 lines
8.0 KiB
Rust
Executable File
//! File transfer infrastructure — Phase 14.
|
|
//! TransferManager and FileTransfer record types.
|
|
|
|
use crate::core::message::ChatMessage;
|
|
use crate::core::protocol::ProtocolType;
|
|
use crate::engine::mux::StreamId;
|
|
use chrono::{DateTime, Utc};
|
|
use dashmap::DashMap;
|
|
use sha2::{Digest, Sha256};
|
|
use std::path::Path;
|
|
use tokio::sync::mpsc;
|
|
|
|
pub mod engine;
|
|
#[allow(unused_imports)]
|
|
pub use engine::{TransferHeader, TransferProgress, TransferResult, format_progress_bar, receive_file, send_file};
|
|
|
|
/// Generate a new unique transfer ID.
|
|
pub fn new_transfer_id() -> String {
|
|
use std::time::{SystemTime, UNIX_EPOCH};
|
|
format!("xfer-{:x}", SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_nanos())
|
|
}
|
|
|
|
pub type TransferId = String;
|
|
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub enum TransferDirection { Send, Receive }
|
|
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub enum TransferState { Pending, Active, Complete, Failed, Cancelled }
|
|
|
|
#[derive(Debug, Clone)]
|
|
pub struct FileTransfer {
|
|
pub id: TransferId, pub direction: TransferDirection, pub state: TransferState,
|
|
pub protocol: ProtocolType, pub peer: String, pub filename: String,
|
|
pub local_path: std::path::PathBuf, pub file_size: u64, pub bytes_transferred: u64,
|
|
pub sha256: Option<String>, pub stream_id: Option<StreamId>,
|
|
pub started_at: Option<DateTime<Utc>>, pub finished_at: Option<DateTime<Utc>>, pub error: Option<String>,
|
|
}
|
|
|
|
impl FileTransfer {
|
|
pub fn new_send(protocol: ProtocolType, peer: &str, filepath: &Path, tx: &mpsc::Sender<ChatMessage>) -> anyhow::Result<(Self, TransferId)> {
|
|
let filename = filepath.file_name().and_then(|n| n.to_str()).unwrap_or("unknown").to_owned();
|
|
let local_path = std::fs::canonicalize(filepath)?;
|
|
let file_size = std::fs::metadata(&local_path)?.len();
|
|
let id = new_transfer_id();
|
|
let transfer = Self { id: id.clone(), direction: TransferDirection::Send, state: TransferState::Pending, protocol, peer: peer.to_owned(), filename, local_path, file_size, bytes_transferred: 0, sha256: None, stream_id: None, started_at: None, finished_at: None, error: None };
|
|
let _ = tx.try_send(ChatMessage::notice(protocol, peer, &format!("Transfer queued: {} ({}B)", transfer.filename, file_size)));
|
|
Ok((transfer, id))
|
|
}
|
|
pub fn compute_hash(path: &Path) -> anyhow::Result<String> {
|
|
let mut f = std::fs::File::open(path)?; let mut h = Sha256::new(); std::io::copy(&mut f, &mut h)?;
|
|
Ok(format!("{:x}", h.finalize()))
|
|
}
|
|
pub fn progress_percent(&self) -> f64 { if self.file_size == 0 { 0.0 } else { (self.bytes_transferred as f64 / self.file_size as f64) * 100.0 } }
|
|
pub fn progress_str(&self) -> String {
|
|
let p = self.progress_percent(); let icon = match self.state { TransferState::Pending=>"⏳",TransferState::Active=>"▶",TransferState::Complete=>"✓",TransferState::Failed=>"✗",TransferState::Cancelled=>"⊘" };
|
|
format!("{icon} {p:.1}% ({}/{}) {}", human_bytes(self.bytes_transferred), human_bytes(self.file_size), self.filename)
|
|
}
|
|
pub fn eta_secs(&self) -> Option<f64> {
|
|
if self.state != TransferState::Active || self.bytes_transferred == 0 { return None; }
|
|
let started = self.started_at?; let elapsed = (Utc::now() - started).num_seconds() as f64; if elapsed <= 0.0 { return None; }
|
|
Some((self.file_size - self.bytes_transferred) as f64 / (self.bytes_transferred as f64 / elapsed))
|
|
}
|
|
}
|
|
|
|
fn human_bytes(b: u64) -> String {
|
|
if b >= 1_073_741_824 { format!("{:.2} GiB", b as f64 / 1_073_741_824.0) }
|
|
else if b >= 1_048_576 { format!("{:.2} MiB", b as f64 / 1_048_576.0) }
|
|
else if b >= 1024 { format!("{:.2} KiB", b as f64 / 1024.0) }
|
|
else { format!("{b} B") }
|
|
}
|
|
|
|
pub struct TransferManager { transfers: DashMap<TransferId, FileTransfer>, tx: mpsc::Sender<ChatMessage> }
|
|
|
|
// DashMap doesn't implement Clone directly; we use Arc internally in practice.
|
|
// For testing convenience, provide a method to get a handle sharing the same map.
|
|
impl TransferManager {
|
|
pub fn new(tx: mpsc::Sender<ChatMessage>) -> Self { Self { transfers: DashMap::new(), tx } }
|
|
/// Get a clone-like handle for sharing across tasks (in real usage, wrap in Arc).
|
|
pub fn clone_ref(&self) -> Self {
|
|
Self { transfers: self.transfers.clone(), tx: self.tx.clone() }
|
|
}
|
|
pub fn queue_send(&self, protocol: ProtocolType, peer: &str, filepath: &Path) -> anyhow::Result<TransferId> {
|
|
let (transfer, id) = FileTransfer::new_send(protocol, peer, filepath, &self.tx)?;
|
|
self.transfers.insert(id.clone(), transfer); Ok(id)
|
|
}
|
|
pub fn queue_receive(&self, id: TransferId, protocol: ProtocolType, peer: &str, filename: &str, size: u64, save_path: &Path) {
|
|
let id_for_struct = id.clone();
|
|
self.transfers.insert(id, FileTransfer { id: id_for_struct, direction: TransferDirection::Receive, state: TransferState::Pending, protocol, peer: peer.to_owned(), filename: filename.to_owned(), local_path: save_path.to_path_buf(), file_size: size, bytes_transferred: 0, sha256: None, stream_id: None, started_at: None, finished_at: None, error: None });
|
|
}
|
|
pub fn get(&self, id: &TransferId) -> Option<FileTransfer> { self.transfers.get(id).map(|r| r.clone()) }
|
|
pub fn update_state(&self, id: &TransferId, state: TransferState) {
|
|
if let Some(mut t) = self.transfers.get_mut(id) {
|
|
t.state = state;
|
|
if matches!(state, TransferState::Active) && t.started_at.is_none() { t.started_at = Some(Utc::now()); }
|
|
if matches!(state, TransferState::Complete | TransferState::Failed | TransferState::Cancelled) { t.finished_at = Some(Utc::now()); }
|
|
}
|
|
}
|
|
pub fn update_progress(&self, id: &TransferId, bytes: u64) { if let Some(mut t) = self.transfers.get_mut(id) { t.bytes_transferred = bytes; } }
|
|
pub fn cancel(&self, id: &TransferId) -> bool { self.update_state(id, TransferState::Cancelled); true }
|
|
pub fn list_all(&self) -> Vec<FileTransfer> { self.transfers.iter().map(|r| r.clone()).collect() }
|
|
pub fn list_active(&self) -> Vec<FileTransfer> { self.transfers.iter().filter(|r| matches!(r.value().state, TransferState::Active | TransferState::Pending)).map(|r| r.clone()).collect() }
|
|
/// Get the top N downloads (Receive direction) sorted by progress descending.
|
|
pub fn top_downloads(&self, n: usize) -> Vec<FileTransfer> {
|
|
let mut dl: Vec<FileTransfer> = self.transfers.iter()
|
|
.filter(|r| r.value().direction == TransferDirection::Receive
|
|
&& matches!(r.value().state, TransferState::Active | TransferState::Pending))
|
|
.map(|r| r.clone())
|
|
.collect();
|
|
dl.sort_by(|a, b| b.bytes_transferred.cmp(&a.bytes_transferred));
|
|
dl.truncate(n);
|
|
dl
|
|
}
|
|
/// Get the top N uploads (Send direction) sorted by progress descending.
|
|
pub fn top_uploads(&self, n: usize) -> Vec<FileTransfer> {
|
|
let mut ul: Vec<FileTransfer> = self.transfers.iter()
|
|
.filter(|r| r.value().direction == TransferDirection::Send
|
|
&& matches!(r.value().state, TransferState::Active | TransferState::Pending))
|
|
.map(|r| r.clone())
|
|
.collect();
|
|
ul.sort_by(|a, b| b.bytes_transferred.cmp(&a.bytes_transferred));
|
|
ul.truncate(n);
|
|
ul
|
|
}
|
|
/// Get total active transfer counts (downloads, uploads).
|
|
pub fn transfer_counts(&self) -> (usize, usize) {
|
|
let (mut dl, mut ul) = (0usize, 0usize);
|
|
for r in self.transfers.iter() {
|
|
if matches!(r.value().state, TransferState::Active | TransferState::Pending) {
|
|
match r.value().direction {
|
|
TransferDirection::Receive => dl += 1,
|
|
TransferDirection::Send => ul += 1,
|
|
}
|
|
}
|
|
}
|
|
(dl, ul)
|
|
}
|
|
pub fn remove(&self, id: &TransferId) -> bool {
|
|
if let Some(t) = self.transfers.get(id) { if matches!(t.state, TransferState::Complete | TransferState::Failed | TransferState::Cancelled) { drop(t); self.transfers.remove(id); return true; } }
|
|
false
|
|
}
|
|
}
|
|
|