Files
ZeroSend/src/transfer/folder_sync.rs
T

218 lines
7.2 KiB
Rust

use crate::transfer::client::TransferClient;
use crate::transfer::protocol::FolderSyncPairConfig;
use chrono::Local;
use std::collections::HashMap;
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::{mpsc, RwLock};
use tracing::info;
#[derive(Debug, Clone)]
pub struct SyncStatusUpdate {
pub pair_id: String,
pub message: String,
pub is_error: bool,
pub timestamp: String,
}
pub struct FolderSyncEngine {
client: Arc<TransferClient>,
pairs: Arc<RwLock<Vec<FolderSyncPairConfig>>>,
status_tx: mpsc::UnboundedSender<SyncStatusUpdate>,
}
impl FolderSyncEngine {
pub fn new(
client: Arc<TransferClient>,
pairs: Vec<FolderSyncPairConfig>,
status_tx: mpsc::UnboundedSender<SyncStatusUpdate>,
) -> Self {
Self {
client,
pairs: Arc::new(RwLock::new(pairs)),
status_tx,
}
}
pub fn pairs_handle(&self) -> Arc<RwLock<Vec<FolderSyncPairConfig>>> {
self.pairs.clone()
}
pub async fn start_background_sync(
self: Arc<Self>,
peers_map: Arc<RwLock<HashMap<String, crate::network::PeerInfo>>>,
) {
info!("Folder auto-sync engine started.");
loop {
tokio::time::sleep(Duration::from_secs(8)).await;
let pairs = self.pairs.read().await.clone();
for pair in pairs {
if !pair.is_active {
continue;
}
// Check if peer is currently online
let peer_opt = {
let peers = peers_map.read().await;
peers.get(&pair.remote_peer_id).cloned()
};
let peer = match peer_opt {
Some(p) if p.is_online => p,
_ => continue,
};
let peer_url = peer.transfer_url();
let client = self.client.clone();
let status_tx = self.status_tx.clone();
let pair_id = pair.id.clone();
let pair_name = pair.name.clone();
let local_path = pair.local_path.clone();
let remote_folder_id = pair.remote_folder_id.clone();
tokio::spawn(async move {
if let Err(e) = Self::sync_pair(
&client,
&peer_url,
&local_path,
&remote_folder_id,
&pair_id,
&pair_name,
&status_tx,
)
.await
{
let _ = status_tx.send(SyncStatusUpdate {
pair_id,
message: format!("Sync error: {}", e),
is_error: true,
timestamp: Local::now().format("%H:%M:%S").to_string(),
});
}
});
}
}
}
async fn sync_pair(
client: &TransferClient,
peer_url: &str,
local_path: &Path,
remote_folder_id: &str,
pair_id: &str,
pair_name: &str,
status_tx: &mpsc::UnboundedSender<SyncStatusUpdate>,
) -> Result<(), String> {
if !local_path.exists() {
let _ = tokio::fs::create_dir_all(local_path).await;
}
// 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));
}
}
// 2. Map remote files
let mut remote_files: HashMap<String, u64> = HashMap::new();
for entry in &remote_tree.entries {
if !entry.is_dir {
remote_files.insert(entry.relative_path.clone(), entry.size);
}
}
let mut synced_count = 0;
// A. Upload newly created or changed local files to remote
if !remote_tree.read_only {
for (rel, (size, _)) in &local_files {
let need_upload = match remote_files.get(rel) {
Some(remote_size) => *remote_size != *size,
None => true, // New local file
};
if need_upload {
let file_full_path = local_path.join(rel);
let subpath = Path::new(rel).parent().and_then(|p| p.to_str()).unwrap_or("");
if let Ok(_) = client
.upload_to_shared_folder(
peer_url,
remote_folder_id,
subpath,
&file_full_path,
None,
)
.await
{
synced_count += 1;
}
}
}
}
// 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,
None => true, // New remote file
};
if need_download {
let target_file_path = local_path.join(rel);
if let Some(parent) = target_file_path.parent() {
let _ = tokio::fs::create_dir_all(parent).await;
}
// Download file
if let Ok(_) = client
.download_shared_file(peer_url, remote_folder_id, rel, &target_file_path, None)
.await
{
synced_count += 1;
}
}
}
if synced_count > 0 {
let _ = status_tx.send(SyncStatusUpdate {
pair_id: pair_id.to_string(),
message: format!("Synced {} item(s) in '{}'", synced_count, pair_name),
is_error: false,
timestamp: Local::now().format("%H:%M:%S").to_string(),
});
}
Ok(())
}
}
fn walkdir_local(base: &Path) -> Result<Vec<(String, u64, u64)>, std::io::Error> {
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() {
if let Ok(rel) = entry.path().strip_prefix(base) {
let rel_str = rel.to_string_lossy().replace('\\', "/");
let meta = entry.metadata().ok();
let size = meta.as_ref().map(|m| m.len()).unwrap_or(0);
let mtime = meta
.as_ref()
.and_then(|m| m.modified().ok())
.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));
}
}
}
Ok(results)
}