From e0edd14fa5bd9c99881d9d262c7ad545d9d9592b Mon Sep 17 00:00:00 2001 From: loki5512344 Date: Wed, 1 Jul 2026 01:38:11 +0200 Subject: [PATCH] feat: remote file ops, chunked transfers with progress, pause/cancel - Remove dead Tauri commands/ module - Add New Folder / Delete / Rename dialogs for remote FS - Implement chunked upload/download with progress for FTP and SFTP - Wrap SFTP ops in spawn_blocking - Add ProgressAction to support cancel/pause mid-transfer - Update todo.md --- src/commands/connection.rs | 60 ------ src/commands/fs.rs | 58 ------ src/commands/mod.rs | 4 - src/commands/sites.rs | 21 -- src/commands/transfer.rs | 108 ---------- src/protocols/ftp.rs | 100 ++++++++-- src/protocols/mod.rs | 33 +++- src/protocols/sftp.rs | 265 +++++++++++++++++-------- src/transfer/queue.rs | 4 + src/transfer/worker.rs | 143 +++++++++++--- src/ui/app.rs | 395 +++++++++++++++++++++++++++++++++++++ src/ui/panels/toolbar.rs | 20 +- src/ui/state.rs | 30 ++- todo.md | 9 +- 14 files changed, 866 insertions(+), 384 deletions(-) delete mode 100644 src/commands/connection.rs delete mode 100644 src/commands/fs.rs delete mode 100644 src/commands/mod.rs delete mode 100644 src/commands/sites.rs delete mode 100644 src/commands/transfer.rs diff --git a/src/commands/connection.rs b/src/commands/connection.rs deleted file mode 100644 index c432265..0000000 --- a/src/commands/connection.rs +++ /dev/null @@ -1,60 +0,0 @@ -use crate::domain::connection::{ConnectionParams, Protocol}; -use crate::fs::remote::RemoteRegistry; -use crate::protocols::ftp::FtpClient; -use crate::protocols::sftp::SftpClient; -use crate::storage::keychain; -use std::sync::Arc; -use tauri::State; - -#[tauri::command] -pub async fn connect( - registry: State<'_, Arc>, - params: ConnectionParams, -) -> Result<(), String> { - let password = if let Some(p) = ¶ms.password { - p.clone() - } else { - keychain::get_password(¶ms.id).map_err(|e| e.to_string())? - }; - - let fs: Arc = match params.protocol { - Protocol::Sftp => { - if let Some(key_path) = ¶ms.key_path { - Arc::new( - SftpClient::connect_key(¶ms.host, params.port, ¶ms.username, key_path) - .map_err(|e| e.to_string())?, - ) - } else { - Arc::new( - SftpClient::connect_password( - ¶ms.host, - params.port, - ¶ms.username, - &password, - ) - .map_err(|e| e.to_string())?, - ) - } - } - Protocol::Ftp => Arc::new( - FtpClient::connect(¶ms.host, params.port, ¶ms.username, &password) - .await - .map_err(|e| e.to_string())?, - ), - Protocol::Ftps => { - return Err("FTPS not yet implemented (suppaftp API compat)".to_string()); - } - }; - - registry.insert(params.id, fs); - Ok(()) -} - -#[tauri::command] -pub async fn disconnect( - registry: State<'_, Arc>, - connection_id: String, -) -> Result<(), String> { - registry.remove(&connection_id); - Ok(()) -} diff --git a/src/commands/fs.rs b/src/commands/fs.rs deleted file mode 100644 index f01b524..0000000 --- a/src/commands/fs.rs +++ /dev/null @@ -1,58 +0,0 @@ -use crate::domain::file_entry::FileEntry; -use crate::fs::{local, remote::RemoteRegistry}; -use std::sync::Arc; -use tauri::State; - -#[tauri::command] -pub fn list_local(path: String) -> Result, String> { - local::list(&path).map_err(|e| e.to_string()) -} - -#[tauri::command] -pub async fn list_remote( - registry: State<'_, Arc>, - connection_id: String, - path: String, -) -> Result, String> { - let fs = registry - .get(&connection_id) - .ok_or_else(|| format!("no connection: {}", connection_id))?; - fs.list(&path).await.map_err(|e| e.to_string()) -} - -#[tauri::command] -pub async fn remote_mkdir( - registry: State<'_, Arc>, - connection_id: String, - path: String, -) -> Result<(), String> { - let fs = registry - .get(&connection_id) - .ok_or_else(|| format!("no connection: {}", connection_id))?; - fs.mkdir(&path).await.map_err(|e| e.to_string()) -} - -#[tauri::command] -pub async fn remote_rename( - registry: State<'_, Arc>, - connection_id: String, - from: String, - to: String, -) -> Result<(), String> { - let fs = registry - .get(&connection_id) - .ok_or_else(|| format!("no connection: {}", connection_id))?; - fs.rename(&from, &to).await.map_err(|e| e.to_string()) -} - -#[tauri::command] -pub async fn remote_delete( - registry: State<'_, Arc>, - connection_id: String, - path: String, -) -> Result<(), String> { - let fs = registry - .get(&connection_id) - .ok_or_else(|| format!("no connection: {}", connection_id))?; - fs.delete(&path).await.map_err(|e| e.to_string()) -} diff --git a/src/commands/mod.rs b/src/commands/mod.rs deleted file mode 100644 index b0bfd45..0000000 --- a/src/commands/mod.rs +++ /dev/null @@ -1,4 +0,0 @@ -pub mod connection; -pub mod fs; -pub mod sites; -pub mod transfer; diff --git a/src/commands/sites.rs b/src/commands/sites.rs deleted file mode 100644 index 901b36e..0000000 --- a/src/commands/sites.rs +++ /dev/null @@ -1,21 +0,0 @@ -use crate::domain::site::Site; -use crate::storage::db::Db; -use tauri::State; - -#[tauri::command] -pub fn get_sites(db: State<'_, Db>) -> Result, String> { - let conn = db.0.lock().unwrap(); - crate::storage::db::get_sites(&conn).map_err(|e| e.to_string()) -} - -#[tauri::command] -pub fn save_site(db: State<'_, Db>, site: Site) -> Result<(), String> { - let conn = db.0.lock().unwrap(); - crate::storage::db::save_site(&conn, &site).map_err(|e| e.to_string()) -} - -#[tauri::command] -pub fn delete_site(db: State<'_, Db>, id: String) -> Result<(), String> { - let conn = db.0.lock().unwrap(); - crate::storage::db::delete_site(&conn, &id).map_err(|e| e.to_string()) -} diff --git a/src/commands/transfer.rs b/src/commands/transfer.rs deleted file mode 100644 index 9a7d765..0000000 --- a/src/commands/transfer.rs +++ /dev/null @@ -1,108 +0,0 @@ -use std::sync::Arc; -use tauri::State; - -use crate::domain::file_entry::EntryKind; -use crate::domain::transfer::{TaskState, TransferKind, TransferTask}; -use crate::transfer::manager::TransferManager; - -#[tauri::command] -pub async fn upload( - manager: State<'_, Arc>, - connection_id: String, - local_path: String, - remote_path: String, -) -> Result { - let meta = - std::fs::metadata(&local_path).map_err(|e| format!("cannot stat local file: {}", e))?; - if !meta.is_file() { - return Err("local path is not a file".into()); - } - - let file_name = std::path::Path::new(&local_path) - .file_name() - .map(|s| s.to_string_lossy().to_string()) - .unwrap_or_else(|| "unknown".to_string()); - - let task = TransferTask::new( - TransferKind::Upload, - connection_id, - local_path, - remote_path, - file_name, - meta.len(), - ); - let task_id = task.id.clone(); - - // Проверяем, что соединение живо - if manager.registry().get(&task.connection_id).is_none() { - return Err(format!("no connection: {}", task.connection_id)); - } - - manager.queue.push(task); - Ok(task_id) -} - -#[tauri::command] -pub async fn download( - manager: State<'_, Arc>, - connection_id: String, - remote_path: String, - local_path: String, -) -> Result { - // Узнаём размер удалённого файла через stat - let registry = manager.registry(); - let fs = registry - .get(&connection_id) - .ok_or_else(|| format!("no connection: {}", connection_id))?; - let stat = fs - .stat(&remote_path) - .await - .map_err(|e| format!("stat failed: {}", e))?; - - let file_name = std::path::Path::new(&remote_path) - .file_name() - .map(|s| s.to_string_lossy().to_string()) - .unwrap_or_else(|| "unknown".to_string()); - - let total = match stat.kind { - EntryKind::File => stat.size.unwrap_or(0), - _ => return Err("remote path is not a file".into()), - }; - - let task = TransferTask::new( - TransferKind::Download, - connection_id, - local_path, - remote_path, - file_name, - total, - ); - let task_id = task.id.clone(); - manager.queue.push(task); - Ok(task_id) -} - -#[tauri::command] -pub async fn get_queue( - manager: State<'_, Arc>, -) -> Result, String> { - Ok(manager.queue.all()) -} - -#[tauri::command] -pub async fn pause_task( - manager: State<'_, Arc>, - task_id: String, -) -> Result<(), String> { - manager.queue.update_state(&task_id, TaskState::Paused); - Ok(()) -} - -#[tauri::command] -pub async fn cancel_task( - manager: State<'_, Arc>, - task_id: String, -) -> Result<(), String> { - manager.queue.update_state(&task_id, TaskState::Cancelled); - Ok(()) -} diff --git a/src/protocols/ftp.rs b/src/protocols/ftp.rs index bcbbeaa..141ceb3 100644 --- a/src/protocols/ftp.rs +++ b/src/protocols/ftp.rs @@ -3,13 +3,17 @@ use async_trait::async_trait; use futures_lite::io::{AsyncReadExt, AsyncWriteExt}; use std::str::FromStr; use std::time::SystemTime; +use tokio::fs::File; +use tokio::io::{AsyncReadExt as TokioReadExt, AsyncWriteExt as TokioWriteExt}; use tokio::sync::Mutex; -use suppaftp::{AsyncFtpStream, list::File}; +use suppaftp::{AsyncFtpStream, list::File as FtpListFile}; -use super::RemoteFs; +use super::{ProgressAction, RemoteFs}; use crate::domain::file_entry::{EntryKind, FileEntry}; +const CHUNK_SIZE: usize = 64 * 1024; + pub struct FtpClient { pub conn: Mutex>, } @@ -30,7 +34,7 @@ impl FtpClient { } fn parse_list_entry(line: &str) -> Option { - let file = File::from_str(line).ok()?; + let file = FtpListFile::from_str(line).ok()?; let kind = if file.is_directory() { EntryKind::Dir } else if file.is_symlink() { @@ -83,9 +87,16 @@ impl RemoteFs for FtpClient { Ok(entries) } - async fn upload(&self, local: &str, remote: &str) -> Result<()> { - let data = - std::fs::read(local).map_err(|e| anyhow::anyhow!("cannot read local file: {}", e))?; + async fn upload_with_progress( + &self, + local: &str, + remote: &str, + on_progress: Option ProgressAction + Send>>, + ) -> Result<()> { + let mut file = File::open(local) + .await + .map_err(|e| anyhow::anyhow!("cannot open local file: {}", e))?; + let mut guard = self.conn.lock().await; let conn = guard .as_mut() @@ -94,17 +105,47 @@ impl RemoteFs for FtpClient { .put_with_stream(remote) .await .map_err(|e| anyhow::anyhow!("FTP upload stream failed: {}", e))?; - stream - .write_all(&data) - .await - .map_err(|e| anyhow::anyhow!("FTP upload write failed: {}", e))?; + + let mut buf = vec![0u8; CHUNK_SIZE]; + let mut total = 0u64; + loop { + let n = file + .read(&mut buf) + .await + .map_err(|e| anyhow::anyhow!("FTP upload read failed: {}", e))?; + if n == 0 { + break; + } + stream + .write_all(&buf[..n]) + .await + .map_err(|e| anyhow::anyhow!("FTP upload write failed: {}", e))?; + total += n as u64; + if let Some(ref cb) = on_progress { + match cb(total) { + ProgressAction::Continue => {} + ProgressAction::Cancel => { + return Err(anyhow::anyhow!("cancelled")); + } + ProgressAction::Pause => { + return Err(anyhow::anyhow!("paused")); + } + } + } + } + conn.finalize_put_stream(stream) .await .map_err(|e| anyhow::anyhow!("FTP upload finalize failed: {}", e))?; Ok(()) } - async fn download(&self, remote: &str, local: &str) -> Result<()> { + async fn download_with_progress( + &self, + remote: &str, + local: &str, + on_progress: Option ProgressAction + Send>>, + ) -> Result<()> { let mut guard = self.conn.lock().await; let conn = guard .as_mut() @@ -113,16 +154,41 @@ impl RemoteFs for FtpClient { .retr_as_stream(remote) .await .map_err(|e| anyhow::anyhow!("FTP download stream failed: {}", e))?; - let mut buf = Vec::new(); - stream - .read_to_end(&mut buf) + + let mut file = File::create(local) .await - .map_err(|e| anyhow::anyhow!("FTP download read failed: {}", e))?; + .map_err(|e| anyhow::anyhow!("cannot create local file: {}", e))?; + + let mut buf = vec![0u8; CHUNK_SIZE]; + let mut total = 0u64; + loop { + let n = stream + .read(&mut buf) + .await + .map_err(|e| anyhow::anyhow!("FTP download read failed: {}", e))?; + if n == 0 { + break; + } + file.write_all(&buf[..n]) + .await + .map_err(|e| anyhow::anyhow!("FTP download write failed: {}", e))?; + total += n as u64; + if let Some(ref cb) = on_progress { + match cb(total) { + ProgressAction::Continue => {} + ProgressAction::Cancel => { + return Err(anyhow::anyhow!("cancelled")); + } + ProgressAction::Pause => { + return Err(anyhow::anyhow!("paused")); + } + } + } + } + conn.finalize_retr_stream(stream) .await .map_err(|e| anyhow::anyhow!("FTP download finalize failed: {}", e))?; - std::fs::write(local, &buf) - .map_err(|e| anyhow::anyhow!("cannot write local file: {}", e))?; Ok(()) } diff --git a/src/protocols/mod.rs b/src/protocols/mod.rs index 328fd80..e6c299a 100644 --- a/src/protocols/mod.rs +++ b/src/protocols/mod.rs @@ -5,13 +5,42 @@ use crate::domain::file_entry::FileEntry; use anyhow::Result; use async_trait::async_trait; +/// Результат выполнения колбэка прогресса: продолжить, отменить или приостановить. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ProgressAction { + Continue, + Cancel, + Pause, +} + /// Единый контракт для всех протоколов. /// Очередь передач не знает — FTP это или SFTP. #[async_trait] pub trait RemoteFs: Send + Sync { async fn list(&self, path: &str) -> Result>; - async fn upload(&self, local: &str, remote: &str) -> Result<()>; - async fn download(&self, remote: &str, local: &str) -> Result<()>; + + async fn upload(&self, local: &str, remote: &str) -> Result<()> { + self.upload_with_progress(local, remote, None).await + } + + async fn download(&self, remote: &str, local: &str) -> Result<()> { + self.download_with_progress(remote, local, None).await + } + + async fn upload_with_progress( + &self, + local: &str, + remote: &str, + on_progress: Option ProgressAction + Send>>, + ) -> Result<()>; + + async fn download_with_progress( + &self, + remote: &str, + local: &str, + on_progress: Option ProgressAction + Send>>, + ) -> Result<()>; + async fn mkdir(&self, path: &str) -> Result<()>; async fn rename(&self, from: &str, to: &str) -> Result<()>; async fn delete(&self, path: &str) -> Result<()>; diff --git a/src/protocols/sftp.rs b/src/protocols/sftp.rs index c1426f2..24fa423 100644 --- a/src/protocols/sftp.rs +++ b/src/protocols/sftp.rs @@ -1,14 +1,18 @@ use anyhow::{Context, Result, bail}; use async_trait::async_trait; use ssh2::{CheckResult, HashType, KnownHostFileKind, Session}; +use std::io::{Read, Write}; use std::net::TcpStream; use std::path::{Path, PathBuf}; +use std::sync::{Arc, Mutex}; -use super::RemoteFs; +use super::{ProgressAction, RemoteFs}; use crate::domain::file_entry::{EntryKind, FileEntry}; +const CHUNK_SIZE: usize = 64 * 1024; + pub struct SftpClient { - session: Session, + session: Arc>, } impl SftpClient { @@ -22,7 +26,9 @@ impl SftpClient { session .userauth_password(user, password) .context("SSH password auth failed")?; - Ok(Self { session }) + Ok(Self { + session: Arc::new(Mutex::new(session)), + }) } /// Подключение по ключу @@ -35,7 +41,9 @@ impl SftpClient { session .userauth_pubkey_file(user, None, Path::new(key_path), None) .context("SSH key auth failed")?; - Ok(Self { session }) + Ok(Self { + session: Arc::new(Mutex::new(session)), + }) } } @@ -104,100 +112,197 @@ fn verify_host_key(session: &Session, host: &str, port: u16) -> Result<()> { #[async_trait] impl RemoteFs for SftpClient { async fn list(&self, path: &str) -> Result> { - // TODO: перенести в spawn_blocking когда будет Arc - let sftp = self.session.sftp().context("SFTP subsystem failed")?; - let entries = sftp.readdir(Path::new(path)).context("readdir failed")?; + let session = self.session.clone(); + let path = path.to_string(); + tokio::task::spawn_blocking(move || { + let sftp = session.lock().unwrap().sftp().context("SFTP subsystem failed")?; + let entries = sftp.readdir(Path::new(&path)).context("readdir failed")?; - let result = entries - .into_iter() - .map(|(pb, stat)| { - let kind = if stat.is_dir() { - EntryKind::Dir - } else if stat.file_type().is_symlink() { - EntryKind::Symlink - } else { - EntryKind::File - }; - FileEntry { - name: pb - .file_name() - .unwrap_or_default() - .to_string_lossy() - .to_string(), - path: pb.to_string_lossy().to_string(), - kind, - size: stat.size, - modified: stat.mtime.map(|t| t as i64), - permissions: None, + let result = entries + .into_iter() + .map(|(pb, stat)| { + let kind = if stat.is_dir() { + EntryKind::Dir + } else if stat.file_type().is_symlink() { + EntryKind::Symlink + } else { + EntryKind::File + }; + FileEntry { + name: pb + .file_name() + .unwrap_or_default() + .to_string_lossy() + .to_string(), + path: pb.to_string_lossy().to_string(), + kind, + size: stat.size, + modified: stat.mtime.map(|t| t as i64), + permissions: None, + } + }) + .collect(); + + Ok(result) + }) + .await + .map_err(|e| anyhow::anyhow!("SFTP list task failed: {}", e))? + } + + async fn upload_with_progress( + &self, + local: &str, + remote: &str, + on_progress: Option ProgressAction + Send>>, + ) -> Result<()> { + let session = self.session.clone(); + let local = local.to_string(); + let remote = remote.to_string(); + tokio::task::spawn_blocking(move || { + let sftp = session.lock().unwrap().sftp().context("SFTP subsystem failed")?; + let mut local_file = std::fs::File::open(&local).context("open local file failed")?; + let mut remote_file = sftp + .create(Path::new(&remote)) + .context("create remote file failed")?; + + let mut buf = vec![0u8; CHUNK_SIZE]; + let mut total = 0u64; + loop { + let n = local_file.read(&mut buf).context("read local file failed")?; + if n == 0 { + break; } - }) - .collect(); - - Ok(result) + remote_file + .write_all(&buf[..n]) + .context("write remote file failed")?; + total += n as u64; + if let Some(ref cb) = on_progress { + match cb(total) { + ProgressAction::Continue => {} + ProgressAction::Cancel => { + return Err(anyhow::anyhow!("cancelled")); + } + ProgressAction::Pause => { + return Err(anyhow::anyhow!("paused")); + } + } + } + } + Ok(()) + }) + .await + .map_err(|e| anyhow::anyhow!("SFTP upload task failed: {}", e))? } - async fn upload(&self, local: &str, remote: &str) -> Result<()> { - // TODO: chunked upload с прогресс-событиями - let sftp = self.session.sftp().context("SFTP subsystem failed")?; - let mut local_file = std::fs::File::open(local).context("open local file failed")?; - let mut remote_file = sftp - .create(Path::new(remote)) - .context("create remote file failed")?; - std::io::copy(&mut local_file, &mut remote_file)?; - Ok(()) - } + async fn download_with_progress( + &self, + remote: &str, + local: &str, + on_progress: Option ProgressAction + Send>>, + ) -> Result<()> { + let session = self.session.clone(); + let remote = remote.to_string(); + let local = local.to_string(); + tokio::task::spawn_blocking(move || { + let sftp = session.lock().unwrap().sftp().context("SFTP subsystem failed")?; + let mut remote_file = sftp.open(Path::new(&remote)).context("open remote file failed")?; + let mut local_file = + std::fs::File::create(&local).context("create local file failed")?; - async fn download(&self, remote: &str, local: &str) -> Result<()> { - // TODO: chunked download с прогресс-событиями - let sftp = self.session.sftp().context("SFTP subsystem failed")?; - let mut remote_file = sftp - .open(Path::new(remote)) - .context("open remote file failed")?; - let mut local_file = std::fs::File::create(local).context("create local file failed")?; - std::io::copy(&mut remote_file, &mut local_file)?; - Ok(()) + let mut buf = vec![0u8; CHUNK_SIZE]; + let mut total = 0u64; + loop { + let n = remote_file.read(&mut buf).context("read remote file failed")?; + if n == 0 { + break; + } + local_file + .write_all(&buf[..n]) + .context("write local file failed")?; + total += n as u64; + if let Some(ref cb) = on_progress { + match cb(total) { + ProgressAction::Continue => {} + ProgressAction::Cancel => { + return Err(anyhow::anyhow!("cancelled")); + } + ProgressAction::Pause => { + return Err(anyhow::anyhow!("paused")); + } + } + } + } + Ok(()) + }) + .await + .map_err(|e| anyhow::anyhow!("SFTP download task failed: {}", e))? } async fn mkdir(&self, path: &str) -> Result<()> { - let sftp = self.session.sftp().context("SFTP subsystem failed")?; - sftp.mkdir(Path::new(path), 0o755).context("mkdir failed")?; - Ok(()) + let session = self.session.clone(); + let path = path.to_string(); + tokio::task::spawn_blocking(move || { + let sftp = session.lock().unwrap().sftp().context("SFTP subsystem failed")?; + sftp.mkdir(Path::new(&path), 0o755).context("mkdir failed")?; + Ok(()) + }) + .await + .map_err(|e| anyhow::anyhow!("SFTP mkdir task failed: {}", e))? } async fn rename(&self, from: &str, to: &str) -> Result<()> { - let sftp = self.session.sftp().context("SFTP subsystem failed")?; - sftp.rename(Path::new(from), Path::new(to), None) - .context("rename failed")?; - Ok(()) + let session = self.session.clone(); + let from = from.to_string(); + let to = to.to_string(); + tokio::task::spawn_blocking(move || { + let sftp = session.lock().unwrap().sftp().context("SFTP subsystem failed")?; + sftp.rename(Path::new(&from), Path::new(&to), None) + .context("rename failed")?; + Ok(()) + }) + .await + .map_err(|e| anyhow::anyhow!("SFTP rename task failed: {}", e))? } async fn delete(&self, path: &str) -> Result<()> { - let sftp = self.session.sftp().context("SFTP subsystem failed")?; - // пробуем как файл, потом как директорию - if sftp.unlink(Path::new(path)).is_err() { - sftp.rmdir(Path::new(path)).context("delete failed")?; - } - Ok(()) + let session = self.session.clone(); + let path = path.to_string(); + tokio::task::spawn_blocking(move || { + let sftp = session.lock().unwrap().sftp().context("SFTP subsystem failed")?; + // пробуем как файл, потом как директорию + if sftp.unlink(Path::new(&path)).is_err() { + sftp.rmdir(Path::new(&path)).context("delete failed")?; + } + Ok(()) + }) + .await + .map_err(|e| anyhow::anyhow!("SFTP delete task failed: {}", e))? } async fn stat(&self, path: &str) -> Result { - let sftp = self.session.sftp().context("SFTP subsystem failed")?; - let stat = sftp.stat(Path::new(path)).context("stat failed")?; - Ok(FileEntry { - name: Path::new(path) - .file_name() - .unwrap_or_default() - .to_string_lossy() + let session = self.session.clone(); + let path = path.to_string(); + tokio::task::spawn_blocking(move || { + let sftp = session.lock().unwrap().sftp().context("SFTP subsystem failed")?; + let stat = sftp.stat(Path::new(&path)).context("stat failed")?; + Ok(FileEntry { + name: Path::new(&path) + .file_name() + .unwrap_or_default() + .to_string_lossy() .to_string(), - path: path.to_string(), - kind: if stat.is_dir() { - EntryKind::Dir - } else { - EntryKind::File - }, - size: stat.size, - modified: stat.mtime.map(|t| t as i64), - permissions: None, + path: path.to_string(), + kind: if stat.is_dir() { + EntryKind::Dir + } else { + EntryKind::File + }, + size: stat.size, + modified: stat.mtime.map(|t| t as i64), + permissions: None, + }) }) + .await + .map_err(|e| anyhow::anyhow!("SFTP stat task failed: {}", e))? } } diff --git a/src/transfer/queue.rs b/src/transfer/queue.rs index d24fc54..146bed3 100644 --- a/src/transfer/queue.rs +++ b/src/transfer/queue.rs @@ -16,6 +16,10 @@ impl TransferQueue { self.inner.lock().unwrap().pop_front() } + pub fn get(&self, id: &str) -> Option { + self.inner.lock().unwrap().iter().find(|t| t.id == id).cloned() + } + pub fn all(&self) -> Vec { self.inner.lock().unwrap().iter().cloned().collect() } diff --git a/src/transfer/worker.rs b/src/transfer/worker.rs index e71c731..7cdee9f 100644 --- a/src/transfer/worker.rs +++ b/src/transfer/worker.rs @@ -1,10 +1,16 @@ use std::sync::Arc; -use tokio::time::{Duration, sleep}; +use std::time::{Duration, Instant}; +use tokio::time::sleep; use crate::domain::transfer::{TaskState, TransferKind}; use crate::fs::remote::RemoteRegistry; +use crate::protocols::ProgressAction; +use crate::transfer::progress::ProgressThrottle; use crate::transfer::queue::TransferQueue; +const POLL_INTERVAL_MS: u64 = 500; +const PAUSED_SLEEP_MS: u64 = 1000; + pub fn spawn_worker( queue: TransferQueue, registry: Arc, @@ -12,25 +18,63 @@ pub fn spawn_worker( ) { rt_handle.spawn(async move { loop { - let task = { - let task = queue.pop(); - if let Some(ref t) = task - && t.state != TaskState::Queued - { - queue.push(t.clone()); - sleep(Duration::from_millis(500)).await; + let front = queue.all().into_iter().next(); + match front.as_ref().map(|t| &t.state) { + None => { + sleep(Duration::from_millis(POLL_INTERVAL_MS)).await; continue; } - task + Some(TaskState::Paused) => { + sleep(Duration::from_millis(PAUSED_SLEEP_MS)).await; + continue; + } + Some(TaskState::Cancelled) => { + if let Some(t) = queue.pop() { + queue.remove(&t.id); + } + continue; + } + Some(TaskState::Running) => { + // Anomalous state: a previous worker died or the task + // was left running. Pop it and mark failed so the queue + // does not spin forever. + if let Some(mut t) = queue.pop() { + t.state = TaskState::Failed("stale running task".into()); + queue.update_state(&t.id, t.state.clone()); + } + continue; + } + Some(TaskState::Completed) | Some(TaskState::Failed(_)) => { + // Terminal tasks should not block new work. Move them to + // the back of the queue so they remain visible in the UI. + if let Some(t) = queue.pop() { + queue.push(t); + } + continue; + } + Some(TaskState::Retrying(_)) => { + // Not implemented yet: demote to queued so it can run. + if let Some(mut t) = queue.pop() { + t.state = TaskState::Queued; + queue.push(t); + } + continue; + } + Some(TaskState::Queued) => { + // Pop and run below. + } + } + + let mut task = match queue.pop() { + Some(t) => t, + None => continue, }; - let mut task = match task { - Some(t) => t, - None => { - sleep(Duration::from_millis(500)).await; - continue; - } - }; + // State may have changed between peek and pop. + if task.state == TaskState::Cancelled { + queue.remove(&task.id); + continue; + } task.state = TaskState::Running; queue.update_state(&task.id, TaskState::Running); @@ -44,22 +88,77 @@ pub fn spawn_worker( } }; + // Progress callback: throttles queue updates to ~10 FPS, + // computes current speed, and reacts to Cancel/Pause requests. + let queue_for_progress = queue.clone(); + let task_id_for_progress = task.id.clone(); + let throttle = Arc::new(std::sync::Mutex::new(ProgressThrottle::default())); + let last_sample = Arc::new(std::sync::Mutex::new((Instant::now(), 0u64))); + let on_progress: Option ProgressAction + Send>> = Some(Box::new( + move |transferred: u64| { + // Check for user-initiated cancel/pause first. + if let Some(t) = queue_for_progress.get(&task_id_for_progress) { + match t.state { + TaskState::Cancelled => return ProgressAction::Cancel, + TaskState::Paused => return ProgressAction::Pause, + _ => {} + } + } + + let speed = { + let mut guard = last_sample.lock().unwrap(); + let (last_time, last_bytes) = *guard; + let now = Instant::now(); + let elapsed = now.duration_since(last_time).as_secs_f64(); + let speed = if elapsed > 0.0 { + ((transferred.saturating_sub(last_bytes)) as f64 / elapsed) as u64 + } else { + 0 + }; + *guard = (now, transferred); + speed + }; + if throttle.lock().unwrap().should_emit() { + queue_for_progress.update_progress(&task_id_for_progress, transferred, speed); + } + ProgressAction::Continue + }, + )); + let result = match task.kind { - TransferKind::Upload => fs.upload(&task.local_path, &task.remote_path).await, - TransferKind::Download => fs.download(&task.remote_path, &task.local_path).await, + TransferKind::Upload => { + fs.upload_with_progress(&task.local_path, &task.remote_path, on_progress).await + } + TransferKind::Download => { + fs.download_with_progress(&task.remote_path, &task.local_path, on_progress).await + } }; match result { Ok(()) => { - task.state = TaskState::Completed; - task.transferred_bytes = task.total_bytes; queue.update_state(&task.id, TaskState::Completed); queue.update_progress(&task.id, task.total_bytes, 0); } Err(e) => { let msg = e.to_string(); - task.state = TaskState::Failed(msg); - queue.update_state(&task.id, task.state.clone()); + if msg == "cancelled" { + queue.update_state(&task.id, TaskState::Cancelled); + queue.remove(&task.id); + } else if msg == "paused" { + let transferred = queue + .get(&task.id) + .map(|t| t.transferred_bytes) + .unwrap_or(0); + let mut paused_task = task.clone(); + paused_task.state = TaskState::Paused; + paused_task.transferred_bytes = transferred; + paused_task.speed = None; + paused_task.eta_secs = None; + queue.update_state(&task.id, TaskState::Paused); + queue.update_progress(&task.id, transferred, 0); + } else { + queue.update_state(&task.id, TaskState::Failed(msg)); + } } } } diff --git a/src/ui/app.rs b/src/ui/app.rs index ba8bb14..ce128c4 100644 --- a/src/ui/app.rs +++ b/src/ui/app.rs @@ -148,6 +148,94 @@ impl FileManagerApp { remote_pane::trigger_list(&mut self.state, idx, &self.registry, self.rt.handle()); } } + + // --- Mkdir result --- + let opt = self.state.pending_mkdir_result.take(); + if let Some(res) = opt.as_ref() { + let mut guard = res.lock().unwrap(); + if let Some(r) = guard.take() { + drop(guard); + match r { + Ok(()) => { + self.state.status_message = "Folder created".into(); + self.state.mkdir_name.clear(); + if let Some(idx) = self.active_tab_idx() { + remote_pane::trigger_list( + &mut self.state, + idx, + &self.registry, + self.rt.handle(), + ); + } + } + Err(e) => { + self.state.status_message = format!("Create folder failed: {}", e); + } + } + } else { + drop(guard); + self.state.pending_mkdir_result = opt; + } + } + + // --- Delete result --- + let opt = self.state.pending_delete_result.take(); + if let Some(res) = opt.as_ref() { + let mut guard = res.lock().unwrap(); + if let Some(r) = guard.take() { + drop(guard); + match r { + Ok(()) => { + self.state.status_message = "Deleted".into(); + self.state.delete_name.clear(); + if let Some(idx) = self.active_tab_idx() { + remote_pane::trigger_list( + &mut self.state, + idx, + &self.registry, + self.rt.handle(), + ); + } + } + Err(e) => { + self.state.status_message = format!("Delete failed: {}", e); + } + } + } else { + drop(guard); + self.state.pending_delete_result = opt; + } + } + + // --- Rename result --- + let opt = self.state.pending_rename_result.take(); + if let Some(res) = opt.as_ref() { + let mut guard = res.lock().unwrap(); + if let Some(r) = guard.take() { + drop(guard); + match r { + Ok(()) => { + self.state.status_message = "Renamed".into(); + self.state.rename_old_name.clear(); + self.state.rename_new_name.clear(); + if let Some(idx) = self.active_tab_idx() { + remote_pane::trigger_list( + &mut self.state, + idx, + &self.registry, + self.rt.handle(), + ); + } + } + Err(e) => { + self.state.status_message = format!("Rename failed: {}", e); + } + } + } else { + drop(guard); + self.state.pending_rename_result = opt; + } + } } fn active_tab_idx(&self) -> Option { @@ -162,6 +250,102 @@ impl FileManagerApp { } } + fn start_mkdir(&mut self, tab_idx: usize, name: String) { + let tab = &self.state.tabs[tab_idx]; + let connection_id = tab.id.clone(); + let remote_path = tab.remote_path.clone(); + let registry = self.registry.clone(); + let path = format!("{}/{}", remote_path.trim_end_matches('/'), name); + + let result = Arc::new(std::sync::Mutex::new(None::>)); + let result_clone = result.clone(); + + self.rt.handle().spawn(async move { + let fs = match registry.get(&connection_id) { + Some(fs) => fs, + None => { + *result_clone.lock().unwrap() = Some(Err("connection not found".into())); + return; + } + }; + let r = fs.mkdir(&path).await.map_err(|e| e.to_string()); + *result_clone.lock().unwrap() = Some(r); + }); + + self.state.pending_mkdir_result = Some(result); + } + + fn start_delete(&mut self, tab_idx: usize, name: String) { + let tab = &self.state.tabs[tab_idx]; + let connection_id = tab.id.clone(); + let registry = self.registry.clone(); + + let entry_path = tab + .remote_entries + .iter() + .find(|e| e.name == name) + .map(|e| e.path.clone()); + + if name == ".." { + self.state.status_message = "Cannot delete parent entry".into(); + return; + } + + let Some(path) = entry_path else { + self.state.status_message = "Selected entry not found".into(); + return; + }; + + let result = Arc::new(std::sync::Mutex::new(None::>)); + let result_clone = result.clone(); + + self.rt.handle().spawn(async move { + let fs = match registry.get(&connection_id) { + Some(fs) => fs, + None => { + *result_clone.lock().unwrap() = Some(Err("connection not found".into())); + return; + } + }; + let r = fs.delete(&path).await.map_err(|e| e.to_string()); + *result_clone.lock().unwrap() = Some(r); + }); + + self.state.pending_delete_result = Some(result); + } + + fn start_rename(&mut self, tab_idx: usize, old_name: String, new_name: String) { + let tab = &self.state.tabs[tab_idx]; + let connection_id = tab.id.clone(); + let registry = self.registry.clone(); + + let Some(entry) = tab.remote_entries.iter().find(|e| e.name == old_name) else { + self.state.status_message = "Selected entry not found".into(); + return; + }; + + let from = entry.path.clone(); + let parent_dir = remote_pane::remote_parent(&from); + let to = format!("{}/{}", parent_dir.trim_end_matches('/'), new_name); + + let result = Arc::new(std::sync::Mutex::new(None::>)); + let result_clone = result.clone(); + + self.rt.handle().spawn(async move { + let fs = match registry.get(&connection_id) { + Some(fs) => fs, + None => { + *result_clone.lock().unwrap() = Some(Err("connection not found".into())); + return; + } + }; + let r = fs.rename(&from, &to).await.map_err(|e| e.to_string()); + *result_clone.lock().unwrap() = Some(r); + }); + + self.state.pending_rename_result = Some(result); + } + fn apply_visuals(&self, ctx: &egui::Context) { let mut vis = Visuals::dark(); vis.window_fill = BG_PANEL; @@ -361,11 +545,222 @@ impl eframe::App for FileManagerApp { connection::render(ctx, &mut self.state, &self.registry, self.rt.handle()); } + // ── Remote operation dialogs ───────────────────────────────────────── + self.render_remote_op_dialogs(ctx); + // 200ms repaint для обновления прогресса ctx.request_repaint_after(std::time::Duration::from_millis(200)); } } +impl FileManagerApp { + fn render_remote_op_dialogs(&mut self, ctx: &egui::Context) { + let screen = ctx.screen_rect(); + let center = screen.center(); + + // затемнённый оверлей под диалогами + if self.state.show_mkdir_dialog + || self.state.show_delete_dialog + || self.state.show_rename_dialog + { + egui::Area::new(egui::Id::new("remote_op_overlay")) + .order(egui::Order::Background) + .show(ctx, |ui| { + ui.painter().rect_filled( + screen, + 0.0, + egui::Color32::from_rgba_premultiplied(0, 0, 0, 120), + ); + }); + } + + // --- New Folder --- + if self.state.show_mkdir_dialog { + let mut open = true; + let mut clicked_ok = false; + egui::Window::new("new_folder_dialog") + .collapsible(false) + .title_bar(false) + .resizable(false) + .default_pos(center) + .pivot(egui::Align2::CENTER_CENTER) + .frame( + egui::Frame::none() + .fill(BG_PANEL) + .stroke(egui::Stroke::new(1.0, BORDER)) + .rounding(8.0) + .inner_margin(egui::Margin::same(16.0)), + ) + .open(&mut open) + .show(ctx, |ui| { + ui.set_width(260.0); + ui.label(RichText::new("New Folder").color(TEXT_PRIMARY).strong()); + ui.add_space(12.0); + ui.text_edit_singleline(&mut self.state.mkdir_name); + ui.add_space(16.0); + ui.horizontal(|ui| { + ui.with_layout(egui::Layout::right_to_left(egui::Align::Center), |ui| { + let ok = egui::Button::new( + RichText::new("OK").color(Color32::WHITE).size(12.0), + ) + .fill(ACCENT) + .rounding(4.0); + if ui.add(ok).clicked() { + clicked_ok = true; + } + ui.add_space(8.0); + let cancel = egui::Button::new( + RichText::new("Cancel").color(TEXT_DIM).size(12.0), + ) + .fill(egui::Color32::TRANSPARENT); + if ui.add(cancel).clicked() { + self.state.show_mkdir_dialog = false; + self.state.mkdir_name.clear(); + } + }); + }); + }); + + if clicked_ok && !self.state.mkdir_name.is_empty() { + if let Some(idx) = self.active_tab_idx() { + let name = self.state.mkdir_name.clone(); + self.start_mkdir(idx, name); + } + self.state.show_mkdir_dialog = false; + } + if !open { + self.state.show_mkdir_dialog = false; + self.state.mkdir_name.clear(); + } + } + + // --- Delete --- + if self.state.show_delete_dialog { + let mut open = true; + let mut clicked_ok = false; + let name = self.state.delete_name.clone(); + egui::Window::new("delete_dialog") + .collapsible(false) + .title_bar(false) + .resizable(false) + .default_pos(center) + .pivot(egui::Align2::CENTER_CENTER) + .frame( + egui::Frame::none() + .fill(BG_PANEL) + .stroke(egui::Stroke::new(1.0, BORDER)) + .rounding(8.0) + .inner_margin(egui::Margin::same(16.0)), + ) + .open(&mut open) + .show(ctx, |ui| { + ui.set_width(260.0); + ui.label( + RichText::new(format!("Delete '{}' ?", name)) + .color(TEXT_PRIMARY) + .strong(), + ); + ui.add_space(16.0); + ui.horizontal(|ui| { + ui.with_layout(egui::Layout::right_to_left(egui::Align::Center), |ui| { + let ok = egui::Button::new( + RichText::new("OK").color(Color32::WHITE).size(12.0), + ) + .fill(RED) + .rounding(4.0); + if ui.add(ok).clicked() { + clicked_ok = true; + } + ui.add_space(8.0); + let cancel = egui::Button::new( + RichText::new("Cancel").color(TEXT_DIM).size(12.0), + ) + .fill(egui::Color32::TRANSPARENT); + if ui.add(cancel).clicked() { + self.state.show_delete_dialog = false; + self.state.delete_name.clear(); + } + }); + }); + }); + + if clicked_ok { + if let Some(idx) = self.active_tab_idx() { + self.start_delete(idx, name); + } + self.state.show_delete_dialog = false; + } + if !open { + self.state.show_delete_dialog = false; + self.state.delete_name.clear(); + } + } + + // --- Rename --- + if self.state.show_rename_dialog { + let mut open = true; + let mut clicked_ok = false; + egui::Window::new("rename_dialog") + .collapsible(false) + .title_bar(false) + .resizable(false) + .default_pos(center) + .pivot(egui::Align2::CENTER_CENTER) + .frame( + egui::Frame::none() + .fill(BG_PANEL) + .stroke(egui::Stroke::new(1.0, BORDER)) + .rounding(8.0) + .inner_margin(egui::Margin::same(16.0)), + ) + .open(&mut open) + .show(ctx, |ui| { + ui.set_width(260.0); + ui.label(RichText::new("Rename").color(TEXT_PRIMARY).strong()); + ui.add_space(12.0); + ui.text_edit_singleline(&mut self.state.rename_new_name); + ui.add_space(16.0); + ui.horizontal(|ui| { + ui.with_layout(egui::Layout::right_to_left(egui::Align::Center), |ui| { + let ok = egui::Button::new( + RichText::new("OK").color(Color32::WHITE).size(12.0), + ) + .fill(ACCENT) + .rounding(4.0); + if ui.add(ok).clicked() { + clicked_ok = true; + } + ui.add_space(8.0); + let cancel = egui::Button::new( + RichText::new("Cancel").color(TEXT_DIM).size(12.0), + ) + .fill(egui::Color32::TRANSPARENT); + if ui.add(cancel).clicked() { + self.state.show_rename_dialog = false; + self.state.rename_new_name.clear(); + self.state.rename_old_name.clear(); + } + }); + }); + }); + + if clicked_ok && !self.state.rename_new_name.is_empty() { + if let Some(idx) = self.active_tab_idx() { + let old_name = self.state.rename_old_name.clone(); + let new_name = self.state.rename_new_name.clone(); + self.start_rename(idx, old_name, new_name); + } + self.state.show_rename_dialog = false; + } + if !open { + self.state.show_rename_dialog = false; + self.state.rename_new_name.clear(); + self.state.rename_old_name.clear(); + } + } + } +} + fn pane_header(ui: &mut egui::Ui, label: &str, width: f32) { egui::Frame::none() .fill(BG_BASE) diff --git a/src/ui/panels/toolbar.rs b/src/ui/panels/toolbar.rs index eddf6c8..87265ef 100644 --- a/src/ui/panels/toolbar.rs +++ b/src/ui/panels/toolbar.rs @@ -69,13 +69,27 @@ pub fn render(ui: &mut Ui, state: &mut AppState) { state.pending_refresh = true; }); toolbar_btn_enabled(ui, "📁 New Folder", has_conn, || { - state.pending_mkdir = true; + state.show_mkdir_dialog = true; + state.mkdir_name.clear(); }); toolbar_btn_enabled(ui, "× Delete", has_conn && has_remote_sel, || { - state.pending_delete = true; + if let Some(name) = state + .active_tab_ref() + .and_then(|t| t.remote_selected.clone()) + { + state.show_delete_dialog = true; + state.delete_name = name; + } }); toolbar_btn_enabled(ui, "✎ Rename", has_conn && has_remote_sel, || { - state.pending_rename = true; + if let Some(name) = state + .active_tab_ref() + .and_then(|t| t.remote_selected.clone()) + { + state.show_rename_dialog = true; + state.rename_old_name = name.clone(); + state.rename_new_name = name; + } }); // Правая часть — History / Bookmarks diff --git a/src/ui/state.rs b/src/ui/state.rs index e177f14..869ab89 100644 --- a/src/ui/state.rs +++ b/src/ui/state.rs @@ -74,9 +74,20 @@ pub struct AppState { // action flags — выставляются тулбаром, обрабатываются в app.rs pub pending_refresh: bool, - pub pending_mkdir: bool, - pub pending_delete: bool, - pub pending_rename: bool, + + // диалоги операций над удалённой ФС + pub show_mkdir_dialog: bool, + pub show_delete_dialog: bool, + pub show_rename_dialog: bool, + + pub mkdir_name: String, + pub delete_name: String, + pub rename_old_name: String, + pub rename_new_name: String, + + pub pending_mkdir_result: Option>, + pub pending_delete_result: Option>, + pub pending_rename_result: Option>, pub pending_connect: Option, pub pending_remote_list: Vec, @@ -144,9 +155,16 @@ impl Default for AppState { status_message: "Ready".into(), connected_count: 0, pending_refresh: false, - pending_mkdir: false, - pending_delete: false, - pending_rename: false, + show_mkdir_dialog: false, + show_delete_dialog: false, + show_rename_dialog: false, + mkdir_name: String::new(), + delete_name: String::new(), + rename_old_name: String::new(), + rename_new_name: String::new(), + pending_mkdir_result: None, + pending_delete_result: None, + pending_rename_result: None, pending_connect: None, pending_remote_list: Vec::new(), } diff --git a/todo.md b/todo.md index d3a9954..1c76c04 100644 --- a/todo.md +++ b/todo.md @@ -1,10 +1,13 @@ # TODO -- [ ] Chunked upload/download с прогресс-событиями +- [x] Chunked upload/download с прогресс-событиями +- [x] Pause/cancel в transfer worker - [x] FTP/FTPS — FTP реализован (FTPS — stub) +- [x] Remote file ops: New Folder / Delete / Rename - [ ] Site Manager UI - [x] Диалог нового соединения - [x] Transfer Queue с воркером и прогресс-событиями -- [ ] Drag & drop +- [x] Drag & drop - [ ] Горячие клавиши -- [ ] egui UI (переезд с Svelte) +- [x] egui UI (переезд с Svelte) +- [ ] FTPS реализация (сейчас stub)