diff --git a/.gitignore b/.gitignore index 3040f88..b18b79a 100644 --- a/.gitignore +++ b/.gitignore @@ -1,5 +1,9 @@ node_modules/ -dist/ +dist/*.tar.gz +dist/*.zip +dist/*.deb +dist/*.rpm +dist/pkg/ target/ .tauri/ src/target/ diff --git a/README.md b/README.md index bf7eb70..796d62e 100644 --- a/README.md +++ b/README.md @@ -10,7 +10,7 @@ FTP/SFTP клиент на Tauri 2: бэкенд на Rust, фронтенд —

Wherry main screen -
Welcome-экран с сохранёнными сайтами и историей подключений +
Welcome-экран с сохранёнными сайтами и историей подключений (тема Midnight, нестандартная — в настройках доступно 6 тем)

## Стек diff --git a/dist/.SRCINFO b/dist/.SRCINFO new file mode 100644 index 0000000..8b0fe0c --- /dev/null +++ b/dist/.SRCINFO @@ -0,0 +1,22 @@ +pkgbase = wherry + pkgdesc = A modern dual-pane file manager with SFTP/FTP/FTPS support + pkgver = 0.1.0 + pkgrel = 1 + url = https://github.com/loki5512344/Wherry + install = + arch = x86_64 + license = GPL3 + makedepends = cargo + makedepends = rust + makedepends = patchelf + depends = webkit2gtk-4.1 + depends = gtk3 + depends = libsoup-3.0 + depends = librsvg + depends = libappindicator-gtk3 + depends = libdbus-1 + depends = openssl + source = wherry-0.1.0.tar.gz::https://github.com/loki5512344/Wherry/archive/refs/tags/v0.1.0.tar.gz + sha256sums = SKIP + +pkgname = wherry diff --git a/dist/PKGBUILD b/dist/PKGBUILD new file mode 100644 index 0000000..a5afdc0 --- /dev/null +++ b/dist/PKGBUILD @@ -0,0 +1,43 @@ +# Maintainer: loki5512344 +# Contributor: Your name + +pkgname=wherry +pkgver=0.1.0 +pkgrel=1 +pkgdesc="A modern dual-pane file manager with SFTP/FTP/FTPS support" +arch=('x86_64') +url='https://github.com/loki5512344/Wherry' +license=('GPL3') +depends=( + 'webkit2gtk-4.1' + 'gtk3' + 'libsoup-3.0' + 'librsvg' + 'libappindicator-gtk3' + 'libdbus-1' + 'openssl' +) +makedepends=( + 'cargo' + 'rust' + 'patchelf' +) +source=("$pkgname-$pkgver.tar.gz::https://github.com/loki5512344/Wherry/archive/refs/tags/v$pkgver.tar.gz") +sha256sums=('SKIP') + +prepare() { + cd "$srcdir/$pkgname-$pkgver/src" +} + +build() { + cd "$srcdir/$pkgname-$pkgver/src" + export CARGO_TARGET_DIR="$srcdir/target" + cargo build --release --frozen +} + +package() { + cd "$srcdir/$pkgname-$pkgver/src" + install -Dm755 target/release/wherry "$pkgdir/usr/bin/wherry" + install -Dm644 icons/128x128.png "$pkgdir/usr/share/pixmaps/wherry.png" + install -Dm644 dist/wherry.desktop "$pkgdir/usr/share/applications/wherry.desktop" +} diff --git a/dist/wherry.desktop b/dist/wherry.desktop new file mode 100644 index 0000000..952dc02 --- /dev/null +++ b/dist/wherry.desktop @@ -0,0 +1,9 @@ +[Desktop Entry] +Type=Application +Name=Wherry +Comment=FTP/SFTP client with dual-pane file manager +Exec=wherry +Icon=wherry +Categories=Network;FileTransfer;Utility; +Terminal=false +StartupNotify=true diff --git a/src/commands.rs b/src/commands.rs index f47fd5b..7ca16a3 100644 --- a/src/commands.rs +++ b/src/commands.rs @@ -3,23 +3,23 @@ use std::sync::Arc; use std::sync::atomic::{AtomicU32, Ordering}; -use tauri::State; use tauri::Manager; +use tauri::State; -use crate::domain::connection::{ConnectionParams, Protocol}; -use crate::domain::error::AppError; -use crate::domain::file_entry::FileEntry; -use crate::domain::site::Site; -use crate::domain::transfer::{TaskState, TransferKind, TransferTask}; -use crate::domain::window_state::WindowState; +use crate::domain::{ + ConnectionParams, FileEntry, Protocol, Site, TaskState, TransferKind, TransferTask, +}; +use crate::error::AppError; use crate::fs::remote::RemoteRegistry; use crate::protocols::{ RemoteFs, ftp::{FtpClient, FtpsClient}, sftp::SftpClient, }; -use crate::storage::db::{self, HistoryRow}; -use crate::transfer::queue::TransferQueue; +use crate::settings; +use crate::storage::{self, HistoryRow}; +use crate::transfers::queue::TransferQueue; +use crate::window::WindowState; pub struct AppState { pub db: Arc>, @@ -46,7 +46,7 @@ pub async fn connect( .lock() .ok() .and_then(|conn| { - db::find_history_conn_id(&conn, ¶ms.host, params.port, ¶ms.username) + storage::find_history_conn_id(&conn, ¶ms.host, params.port, ¶ms.username) .ok() .flatten() }) @@ -98,7 +98,7 @@ pub async fn connect( registry.insert(params.id.clone(), fs); if let Ok(conn) = db.lock() { - let _ = db::add_history_entry( + let _ = storage::add_history_entry( &conn, ¶ms.host, params.port, @@ -266,7 +266,7 @@ pub fn remove_task(state: State<'_, AppState>, id: String) { pub fn set_max_concurrent(state: State<'_, AppState>, n: u32) { state.max_concurrent.store(n.max(1), Ordering::Relaxed); if let Ok(conn) = state.db.lock() { - db::set_u32(&conn, "max_concurrent_transfers", n.max(1)); + settings::set_u32(&conn, "max_concurrent_transfers", n.max(1)); } } @@ -278,7 +278,7 @@ pub fn list_sites(state: State<'_, AppState>) -> Result, AppError> { .db .lock() .map_err(|e| AppError::Internal(e.to_string()))?; - Ok(db::get_sites(&conn)?) + Ok(storage::get_sites(&conn)?) } #[tauri::command] @@ -287,7 +287,7 @@ pub fn save_site(state: State<'_, AppState>, site: Site) -> Result<(), AppError> .db .lock() .map_err(|e| AppError::Internal(e.to_string()))?; - Ok(db::save_site(&conn, &site)?) + Ok(storage::save_site(&conn, &site)?) } #[tauri::command] @@ -296,7 +296,7 @@ pub fn delete_site(state: State<'_, AppState>, id: String) -> Result<(), AppErro .db .lock() .map_err(|e| AppError::Internal(e.to_string()))?; - Ok(db::delete_site(&conn, &id)?) + Ok(storage::delete_site(&conn, &id)?) } #[tauri::command] @@ -305,7 +305,7 @@ pub fn list_bookmarks(state: State<'_, AppState>) -> Result, id: i64) -> Result<(), AppErr .db .lock() .map_err(|e| AppError::Internal(e.to_string()))?; - Ok(db::remove_bookmark(&conn, id)?) + Ok(settings::remove_bookmark(&conn, id)?) } #[tauri::command] @@ -336,7 +336,7 @@ pub fn list_history(state: State<'_, AppState>) -> Result, AppEr .db .lock() .map_err(|e| AppError::Internal(e.to_string()))?; - Ok(db::get_history(&conn)?) + Ok(storage::get_history(&conn)?) } #[tauri::command] @@ -345,7 +345,7 @@ pub fn clear_history(state: State<'_, AppState>) -> Result<(), AppError> { .db .lock() .map_err(|e| AppError::Internal(e.to_string()))?; - Ok(db::clear_history(&conn)?) + Ok(storage::clear_history(&conn)?) } #[tauri::command] @@ -359,7 +359,9 @@ pub fn find_history_conn_id( .db .lock() .map_err(|e| AppError::Internal(e.to_string()))?; - Ok(db::find_history_conn_id(&conn, &host, port, &username)?) + Ok(storage::find_history_conn_id( + &conn, &host, port, &username, + )?) } #[tauri::command] @@ -368,19 +370,21 @@ pub fn get_pref(state: State<'_, AppState>, key: String) -> Result, key: String, value: String) -> Result<(), AppError> { - if key == "auto_clear_completed_secs" && let Ok(secs) = value.parse::() { + if key == "auto_clear_completed_secs" + && let Ok(secs) = value.parse::() + { state.auto_clear_secs.store(secs, Ordering::Relaxed); } let conn = state .db .lock() .map_err(|e| AppError::Internal(e.to_string()))?; - Ok(db::set_setting(&conn, &key, &value)?) + Ok(settings::set_setting(&conn, &key, &value)?) } /// Удалить пароль из сохранённого сайта (очистить поле password в БД). @@ -454,7 +458,7 @@ pub fn save_window_state( let key = format!("window_state_{}", window_state.label); let json = serde_json::to_string(&window_state).map_err(|e| AppError::Internal(e.to_string()))?; - db::set_setting(&conn, &key, &json)?; + settings::set_setting(&conn, &key, &json)?; Ok(()) } @@ -468,7 +472,7 @@ pub fn load_window_state( .lock() .map_err(|e| AppError::Internal(e.to_string()))?; let key = format!("window_state_{}", label); - match db::get_setting(&conn, &key) { + match settings::get_setting(&conn, &key) { Some(json) => { let ws: WindowState = serde_json::from_str(&json).map_err(|e| AppError::Internal(e.to_string()))?; @@ -484,17 +488,14 @@ pub async fn new_window( state: State<'_, AppState>, ) -> Result<(), String> { let label = format!("browser-{}", uuid::Uuid::new_v4()); - let window = tauri::WebviewWindowBuilder::new( - &app_handle, - &label, - tauri::WebviewUrl::App("/".into()), - ) - .title("Wherry") - .inner_size(1280.0, 800.0) - .min_inner_size(900.0, 600.0) - .center() - .build() - .map_err(|e| e.to_string())?; + let window = + tauri::WebviewWindowBuilder::new(&app_handle, &label, tauri::WebviewUrl::App("/".into())) + .title("Wherry") + .inner_size(1280.0, 800.0) + .min_inner_size(900.0, 600.0) + .center() + .build() + .map_err(|e| e.to_string())?; let db = state.db.clone(); let app_handle = app_handle.clone(); @@ -516,12 +517,10 @@ pub(crate) fn save_window_state_internal( ) { let label = window.label().to_string(); let pos = window.outer_position().ok(); - let size = window - .outer_size() - .unwrap_or(tauri::PhysicalSize { - width: 1280, - height: 800, - }); + let size = window.outer_size().unwrap_or(tauri::PhysicalSize { + width: 1280, + height: 800, + }); let maximized = window.is_maximized().unwrap_or(false); let ws = WindowState { @@ -536,7 +535,7 @@ pub(crate) fn save_window_state_internal( if let Ok(conn) = db.lock() { let key = format!("window_state_{}", label); if let Ok(json) = serde_json::to_string(&ws) { - let _ = db::set_setting(&conn, &key, &json); + let _ = settings::set_setting(&conn, &key, &json); } } } diff --git a/src/domain.rs b/src/domain.rs new file mode 100644 index 0000000..c61ae92 --- /dev/null +++ b/src/domain.rs @@ -0,0 +1,376 @@ +use serde::{Deserialize, Serialize}; +use std::fmt; +use uuid::Uuid; + +// ── Protocol & ConnectionParams ── +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(rename_all = "lowercase")] +pub enum Protocol { + Sftp, + Ftp, + Ftps, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct ConnectionParams { + pub id: String, + pub label: String, + pub protocol: Protocol, + pub host: String, + pub port: u16, + pub username: String, + /// None = use keychain + pub password: Option, + pub key_path: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(rename_all = "lowercase")] +pub enum ConnectionStatus { + Connected, + Disconnected, + Connecting, + Error(String), +} + +// ── EntryKind & FileEntry ── +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(rename_all = "lowercase")] +pub enum EntryKind { + File, + Dir, + Symlink, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct FileEntry { + pub name: String, + pub path: String, + pub kind: EntryKind, + pub size: Option, + pub modified: Option, // unix timestamp + pub permissions: Option, +} + +// ── Site ── +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct Site { + pub id: String, + pub name: String, + pub protocol: Protocol, + pub host: String, + pub port: u16, + pub username: String, + pub password: Option, + pub key_path: Option, + pub folder: Option, + pub note: Option, +} + +// ── TransferKind, TaskState, TransferTask ── +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(rename_all = "lowercase")] +pub enum TransferKind { + Upload, + Download, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(rename_all = "lowercase")] +pub enum TaskState { + Queued, + Running, + Paused, + Cancelled, + Completed, + Failed(String), + Retrying(u32), +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct TransferTask { + pub id: String, + pub kind: TransferKind, + pub connection_id: String, + pub local_path: String, + pub remote_path: String, + pub file_name: String, + pub total_bytes: u64, + pub transferred_bytes: u64, + pub state: TaskState, + /// bytes/sec, updated with throttle + pub speed: Option, + pub eta_secs: Option, +} + +impl fmt::Display for TaskState { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + TaskState::Queued => write!(f, "queued"), + TaskState::Running => write!(f, "running"), + TaskState::Paused => write!(f, "paused"), + TaskState::Cancelled => write!(f, "cancelled"), + TaskState::Completed => write!(f, "completed"), + TaskState::Failed(e) => write!(f, "failed: {}", e), + TaskState::Retrying(n) => write!(f, "retrying({})", n), + } + } +} + +impl TransferTask { + pub fn new( + kind: TransferKind, + connection_id: String, + local_path: String, + remote_path: String, + file_name: String, + total_bytes: u64, + ) -> Self { + Self { + id: Uuid::new_v4().to_string(), + kind, + connection_id, + local_path, + remote_path, + file_name, + total_bytes, + transferred_bytes: 0, + state: TaskState::Queued, + speed: None, + eta_secs: None, + } + } + + pub fn progress_pct(&self) -> f64 { + if self.total_bytes == 0 { + return 0.0; + } + (self.transferred_bytes as f64 / self.total_bytes as f64) * 100.0 + } +} + +#[cfg(test)] +mod tests { + use super::*; + + // from connection.rs + #[test] + fn test_protocol_serde() { + let sftp = Protocol::Sftp; + let json = serde_json::to_string(&sftp).unwrap(); + assert_eq!(json, "\"sftp\""); + let deserialized: Protocol = serde_json::from_str(&json).unwrap(); + assert_eq!(deserialized, Protocol::Sftp); + + let ftp = Protocol::Ftp; + let json = serde_json::to_string(&ftp).unwrap(); + assert_eq!(json, "\"ftp\""); + + let ftps = Protocol::Ftps; + let json = serde_json::to_string(&ftps).unwrap(); + assert_eq!(json, "\"ftps\""); + } + + #[test] + fn test_connection_params() { + let params = ConnectionParams { + id: "test-id".into(), + label: "Test".into(), + protocol: Protocol::Sftp, + host: "example.com".into(), + port: 22, + username: "user".into(), + password: Some("pass".into()), + key_path: None, + }; + assert_eq!(params.id, "test-id"); + assert_eq!(params.protocol, Protocol::Sftp); + assert_eq!(params.port, 22); + assert!(params.password.is_some()); + assert!(params.key_path.is_none()); + } + + #[test] + fn test_connection_status_serde() { + let json = serde_json::to_string(&ConnectionStatus::Connected).unwrap(); + assert_eq!(json, "\"connected\""); + + let json = serde_json::to_string(&ConnectionStatus::Error("timeout".into())).unwrap(); + assert_eq!(json, "{\"error\":\"timeout\"}"); + } + + // from file_entry.rs + #[test] + fn test_entry_kind_serde() { + assert_eq!(serde_json::to_string(&EntryKind::File).unwrap(), "\"file\""); + assert_eq!(serde_json::to_string(&EntryKind::Dir).unwrap(), "\"dir\""); + assert_eq!( + serde_json::to_string(&EntryKind::Symlink).unwrap(), + "\"symlink\"" + ); + } + + #[test] + fn test_file_entry() { + let entry = FileEntry { + name: "test.txt".into(), + path: "/tmp/test.txt".into(), + kind: EntryKind::File, + size: Some(1024), + modified: Some(1234567890), + permissions: Some("rw-r--r--".into()), + }; + assert_eq!(entry.name, "test.txt"); + assert_eq!(entry.kind, EntryKind::File); + assert_eq!(entry.size, Some(1024)); + } + + #[test] + fn test_file_entry_serde() { + let entry = FileEntry { + name: "f".into(), + path: "/f".into(), + kind: EntryKind::Dir, + size: None, + modified: None, + permissions: None, + }; + let json = serde_json::to_string(&entry).unwrap(); + let deserialized: FileEntry = serde_json::from_str(&json).unwrap(); + assert_eq!(deserialized.name, "f"); + assert_eq!(deserialized.kind, EntryKind::Dir); + } + + // from site.rs + #[test] + fn test_site_serde() { + let site = Site { + id: "site-1".into(), + name: "My Server".into(), + protocol: Protocol::Sftp, + host: "example.com".into(), + port: 22, + username: "admin".into(), + password: Some("secret".into()), + key_path: None, + folder: Some("/remote".into()), + note: Some("my note".into()), + }; + let json = serde_json::to_string(&site).unwrap(); + let deserialized: Site = serde_json::from_str(&json).unwrap(); + assert_eq!(deserialized.id, site.id); + assert_eq!(deserialized.name, site.name); + assert_eq!(deserialized.protocol, site.protocol); + assert_eq!(deserialized.host, site.host); + assert_eq!(deserialized.port, site.port); + assert_eq!(deserialized.username, site.username); + assert_eq!(deserialized.password, site.password); + assert_eq!(deserialized.folder, site.folder); + assert_eq!(deserialized.note, site.note); + } + + #[test] + fn test_site_minimal() { + let site = Site { + id: "site-2".into(), + name: "Minimal".into(), + protocol: Protocol::Ftp, + host: "ftp.example.com".into(), + port: 21, + username: "user".into(), + password: None, + key_path: None, + folder: None, + note: None, + }; + let json = serde_json::to_string(&site).unwrap(); + let deserialized: Site = serde_json::from_str(&json).unwrap(); + assert_eq!(deserialized.protocol, Protocol::Ftp); + assert!(deserialized.password.is_none()); + } + + // from transfer.rs + #[test] + fn test_transfer_task_new() { + let task = TransferTask::new( + TransferKind::Download, + "conn-1".into(), + "/local".into(), + "/remote".into(), + "file.txt".into(), + 1000, + ); + assert!(!task.id.is_empty()); + assert_eq!(task.kind, TransferKind::Download); + assert_eq!(task.total_bytes, 1000); + assert_eq!(task.transferred_bytes, 0); + assert_eq!(task.state, TaskState::Queued); + assert!(task.speed.is_none()); + assert!(task.eta_secs.is_none()); + } + + #[test] + fn test_progress_pct() { + let mut task = TransferTask::new( + TransferKind::Upload, + "conn-1".into(), + "/local".into(), + "/remote".into(), + "file.txt".into(), + 200, + ); + assert_eq!(task.progress_pct(), 0.0); + task.transferred_bytes = 50; + assert!((task.progress_pct() - 25.0).abs() < f64::EPSILON); + task.transferred_bytes = 200; + assert!((task.progress_pct() - 100.0).abs() < f64::EPSILON); + } + + #[test] + fn test_progress_pct_zero_total() { + let task = TransferTask::new( + TransferKind::Upload, + "conn-1".into(), + "/local".into(), + "/remote".into(), + "file.txt".into(), + 0, + ); + assert_eq!(task.progress_pct(), 0.0); + } + + #[test] + fn test_task_state_display() { + assert_eq!(TaskState::Queued.to_string(), "queued"); + assert_eq!(TaskState::Running.to_string(), "running"); + assert_eq!(TaskState::Paused.to_string(), "paused"); + assert_eq!(TaskState::Cancelled.to_string(), "cancelled"); + assert_eq!(TaskState::Completed.to_string(), "completed"); + assert_eq!(TaskState::Failed("err".into()).to_string(), "failed: err"); + assert_eq!(TaskState::Retrying(3).to_string(), "retrying(3)"); + } + + #[test] + fn test_transfer_kind_serde() { + assert_eq!( + serde_json::to_string(&TransferKind::Upload).unwrap(), + "\"upload\"" + ); + assert_eq!( + serde_json::to_string(&TransferKind::Download).unwrap(), + "\"download\"" + ); + } + + #[test] + fn test_task_state_serde() { + let json = serde_json::to_string(&TaskState::Retrying(2)).unwrap(); + assert_eq!(json, "{\"retrying\":2}"); + let deserialized: TaskState = serde_json::from_str(&json).unwrap(); + assert_eq!(deserialized, TaskState::Retrying(2)); + } +} diff --git a/src/domain/connection.rs b/src/domain/connection.rs deleted file mode 100644 index 1aa5eaf..0000000 --- a/src/domain/connection.rs +++ /dev/null @@ -1,82 +0,0 @@ -use serde::{Deserialize, Serialize}; - -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(rename_all = "lowercase")] -pub enum Protocol { - Sftp, - Ftp, - Ftps, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct ConnectionParams { - pub id: String, - pub label: String, - pub protocol: Protocol, - pub host: String, - pub port: u16, - pub username: String, - /// None = use keychain - pub password: Option, - pub key_path: Option, -} - -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(rename_all = "lowercase")] -pub enum ConnectionStatus { - Connected, - Disconnected, - Connecting, - Error(String), -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_protocol_serde() { - let sftp = Protocol::Sftp; - let json = serde_json::to_string(&sftp).unwrap(); - assert_eq!(json, "\"sftp\""); - let deserialized: Protocol = serde_json::from_str(&json).unwrap(); - assert_eq!(deserialized, Protocol::Sftp); - - let ftp = Protocol::Ftp; - let json = serde_json::to_string(&ftp).unwrap(); - assert_eq!(json, "\"ftp\""); - - let ftps = Protocol::Ftps; - let json = serde_json::to_string(&ftps).unwrap(); - assert_eq!(json, "\"ftps\""); - } - - #[test] - fn test_connection_params() { - let params = ConnectionParams { - id: "test-id".into(), - label: "Test".into(), - protocol: Protocol::Sftp, - host: "example.com".into(), - port: 22, - username: "user".into(), - password: Some("pass".into()), - key_path: None, - }; - assert_eq!(params.id, "test-id"); - assert_eq!(params.protocol, Protocol::Sftp); - assert_eq!(params.port, 22); - assert!(params.password.is_some()); - assert!(params.key_path.is_none()); - } - - #[test] - fn test_connection_status_serde() { - let json = serde_json::to_string(&ConnectionStatus::Connected).unwrap(); - assert_eq!(json, "\"connected\""); - - let json = serde_json::to_string(&ConnectionStatus::Error("timeout".into())).unwrap(); - assert_eq!(json, "{\"error\":\"timeout\"}"); - } -} diff --git a/src/domain/file_entry.rs b/src/domain/file_entry.rs deleted file mode 100644 index 407415a..0000000 --- a/src/domain/file_entry.rs +++ /dev/null @@ -1,66 +0,0 @@ -use serde::{Deserialize, Serialize}; - -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(rename_all = "lowercase")] -pub enum EntryKind { - File, - Dir, - Symlink, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct FileEntry { - pub name: String, - pub path: String, - pub kind: EntryKind, - pub size: Option, - pub modified: Option, // unix timestamp - pub permissions: Option, -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_entry_kind_serde() { - assert_eq!(serde_json::to_string(&EntryKind::File).unwrap(), "\"file\""); - assert_eq!(serde_json::to_string(&EntryKind::Dir).unwrap(), "\"dir\""); - assert_eq!( - serde_json::to_string(&EntryKind::Symlink).unwrap(), - "\"symlink\"" - ); - } - - #[test] - fn test_file_entry() { - let entry = FileEntry { - name: "test.txt".into(), - path: "/tmp/test.txt".into(), - kind: EntryKind::File, - size: Some(1024), - modified: Some(1234567890), - permissions: Some("rw-r--r--".into()), - }; - assert_eq!(entry.name, "test.txt"); - assert_eq!(entry.kind, EntryKind::File); - assert_eq!(entry.size, Some(1024)); - } - - #[test] - fn test_file_entry_serde() { - let entry = FileEntry { - name: "f".into(), - path: "/f".into(), - kind: EntryKind::Dir, - size: None, - modified: None, - permissions: None, - }; - let json = serde_json::to_string(&entry).unwrap(); - let deserialized: FileEntry = serde_json::from_str(&json).unwrap(); - assert_eq!(deserialized.name, "f"); - assert_eq!(deserialized.kind, EntryKind::Dir); - } -} diff --git a/src/domain/mod.rs b/src/domain/mod.rs deleted file mode 100644 index c27971b..0000000 --- a/src/domain/mod.rs +++ /dev/null @@ -1,6 +0,0 @@ -pub mod connection; -pub mod error; -pub mod file_entry; -pub mod site; -pub mod transfer; -pub mod window_state; diff --git a/src/domain/site.rs b/src/domain/site.rs deleted file mode 100644 index 30812dd..0000000 --- a/src/domain/site.rs +++ /dev/null @@ -1,70 +0,0 @@ -use crate::domain::connection::Protocol; -use serde::{Deserialize, Serialize}; - -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct Site { - pub id: String, - pub name: String, - pub protocol: Protocol, - pub host: String, - pub port: u16, - pub username: String, - pub password: Option, - pub key_path: Option, - pub folder: Option, - pub note: Option, -} - -#[cfg(test)] -mod tests { - use super::*; - use crate::domain::connection::Protocol; - - #[test] - fn test_site_serde() { - let site = Site { - id: "site-1".into(), - name: "My Server".into(), - protocol: Protocol::Sftp, - host: "example.com".into(), - port: 22, - username: "admin".into(), - password: Some("secret".into()), - key_path: None, - folder: Some("/remote".into()), - note: Some("my note".into()), - }; - let json = serde_json::to_string(&site).unwrap(); - let deserialized: Site = serde_json::from_str(&json).unwrap(); - assert_eq!(deserialized.id, site.id); - assert_eq!(deserialized.name, site.name); - assert_eq!(deserialized.protocol, site.protocol); - assert_eq!(deserialized.host, site.host); - assert_eq!(deserialized.port, site.port); - assert_eq!(deserialized.username, site.username); - assert_eq!(deserialized.password, site.password); - assert_eq!(deserialized.folder, site.folder); - assert_eq!(deserialized.note, site.note); - } - - #[test] - fn test_site_minimal() { - let site = Site { - id: "site-2".into(), - name: "Minimal".into(), - protocol: Protocol::Ftp, - host: "ftp.example.com".into(), - port: 21, - username: "user".into(), - password: None, - key_path: None, - folder: None, - note: None, - }; - let json = serde_json::to_string(&site).unwrap(); - let deserialized: Site = serde_json::from_str(&json).unwrap(); - assert_eq!(deserialized.protocol, Protocol::Ftp); - assert!(deserialized.password.is_none()); - } -} diff --git a/src/domain/transfer.rs b/src/domain/transfer.rs deleted file mode 100644 index 3a9f363..0000000 --- a/src/domain/transfer.rs +++ /dev/null @@ -1,170 +0,0 @@ -use serde::{Deserialize, Serialize}; -use std::fmt; -use uuid::Uuid; - -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(rename_all = "lowercase")] -pub enum TransferKind { - Upload, - Download, -} - -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(rename_all = "lowercase")] -pub enum TaskState { - Queued, - Running, - Paused, - Cancelled, - Completed, - Failed(String), - Retrying(u32), -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct TransferTask { - pub id: String, - pub kind: TransferKind, - pub connection_id: String, - pub local_path: String, - pub remote_path: String, - pub file_name: String, - pub total_bytes: u64, - pub transferred_bytes: u64, - pub state: TaskState, - /// bytes/sec, updated with throttle - pub speed: Option, - pub eta_secs: Option, -} - -impl fmt::Display for TaskState { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - match self { - TaskState::Queued => write!(f, "queued"), - TaskState::Running => write!(f, "running"), - TaskState::Paused => write!(f, "paused"), - TaskState::Cancelled => write!(f, "cancelled"), - TaskState::Completed => write!(f, "completed"), - TaskState::Failed(e) => write!(f, "failed: {}", e), - TaskState::Retrying(n) => write!(f, "retrying({})", n), - } - } -} - -impl TransferTask { - pub fn new( - kind: TransferKind, - connection_id: String, - local_path: String, - remote_path: String, - file_name: String, - total_bytes: u64, - ) -> Self { - Self { - id: Uuid::new_v4().to_string(), - kind, - connection_id, - local_path, - remote_path, - file_name, - total_bytes, - transferred_bytes: 0, - state: TaskState::Queued, - speed: None, - eta_secs: None, - } - } - - pub fn progress_pct(&self) -> f64 { - if self.total_bytes == 0 { - return 0.0; - } - (self.transferred_bytes as f64 / self.total_bytes as f64) * 100.0 - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_transfer_task_new() { - let task = TransferTask::new( - TransferKind::Download, - "conn-1".into(), - "/local".into(), - "/remote".into(), - "file.txt".into(), - 1000, - ); - assert!(!task.id.is_empty()); - assert_eq!(task.kind, TransferKind::Download); - assert_eq!(task.total_bytes, 1000); - assert_eq!(task.transferred_bytes, 0); - assert_eq!(task.state, TaskState::Queued); - assert!(task.speed.is_none()); - assert!(task.eta_secs.is_none()); - } - - #[test] - fn test_progress_pct() { - let mut task = TransferTask::new( - TransferKind::Upload, - "conn-1".into(), - "/local".into(), - "/remote".into(), - "file.txt".into(), - 200, - ); - assert_eq!(task.progress_pct(), 0.0); - task.transferred_bytes = 50; - assert!((task.progress_pct() - 25.0).abs() < f64::EPSILON); - task.transferred_bytes = 200; - assert!((task.progress_pct() - 100.0).abs() < f64::EPSILON); - } - - #[test] - fn test_progress_pct_zero_total() { - let task = TransferTask::new( - TransferKind::Upload, - "conn-1".into(), - "/local".into(), - "/remote".into(), - "file.txt".into(), - 0, - ); - assert_eq!(task.progress_pct(), 0.0); - } - - #[test] - fn test_task_state_display() { - assert_eq!(TaskState::Queued.to_string(), "queued"); - assert_eq!(TaskState::Running.to_string(), "running"); - assert_eq!(TaskState::Paused.to_string(), "paused"); - assert_eq!(TaskState::Cancelled.to_string(), "cancelled"); - assert_eq!(TaskState::Completed.to_string(), "completed"); - assert_eq!(TaskState::Failed("err".into()).to_string(), "failed: err"); - assert_eq!(TaskState::Retrying(3).to_string(), "retrying(3)"); - } - - #[test] - fn test_transfer_kind_serde() { - assert_eq!( - serde_json::to_string(&TransferKind::Upload).unwrap(), - "\"upload\"" - ); - assert_eq!( - serde_json::to_string(&TransferKind::Download).unwrap(), - "\"download\"" - ); - } - - #[test] - fn test_task_state_serde() { - let json = serde_json::to_string(&TaskState::Retrying(2)).unwrap(); - assert_eq!(json, "{\"retrying\":2}"); - let deserialized: TaskState = serde_json::from_str(&json).unwrap(); - assert_eq!(deserialized, TaskState::Retrying(2)); - } -} diff --git a/src/domain/error.rs b/src/error.rs similarity index 100% rename from src/domain/error.rs rename to src/error.rs diff --git a/src/fs/local.rs b/src/fs/local.rs index 73e89e9..6d24869 100644 --- a/src/fs/local.rs +++ b/src/fs/local.rs @@ -1,4 +1,4 @@ -use crate::domain::file_entry::{EntryKind, FileEntry}; +use crate::domain::{EntryKind, FileEntry}; use anyhow::Result; use std::fs; diff --git a/src/fs/remote.rs b/src/fs/remote.rs index c83fbf8..7e06c74 100644 --- a/src/fs/remote.rs +++ b/src/fs/remote.rs @@ -27,7 +27,7 @@ impl RemoteRegistry { #[cfg(test)] mod tests { use super::*; - use crate::domain::file_entry::{EntryKind, FileEntry}; + use crate::domain::{EntryKind, FileEntry}; use crate::protocols::ProgressAction; struct MockFs; diff --git a/src/lib.rs b/src/lib.rs index aefc1af..b070e1e 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,15 +1,18 @@ pub mod commands; pub mod domain; +pub mod error; pub mod fs; pub mod i18n; pub mod protocols; +pub mod settings; pub mod storage; -pub mod transfer; +pub mod transfers; +pub mod window; use crate::commands::AppState; -use crate::domain::window_state::WindowState; use crate::fs::remote::RemoteRegistry; -use crate::transfer::queue::TransferQueue; +use crate::transfers::queue::TransferQueue; +use crate::window::WindowState; use std::sync::Arc; use std::sync::atomic::AtomicU32; use tauri::Manager; @@ -27,12 +30,12 @@ pub fn run() { std::fs::create_dir_all(parent).ok(); } let conn = rusqlite::Connection::open(&db_path).expect("failed to open wherry database"); - storage::db::init_tables(&conn).expect("failed to init db tables"); + storage::init_tables(&conn).expect("failed to init db tables"); let db = Arc::new(std::sync::Mutex::new(conn)); let max_concurrent = { let c = db.lock().unwrap(); - Arc::new(AtomicU32::new(storage::db::get_u32( + Arc::new(AtomicU32::new(settings::get_u32( &c, "max_concurrent_transfers", DEFAULT_MAX_CONCURRENT, @@ -41,7 +44,7 @@ pub fn run() { let auto_clear_secs = { let c = db.lock().unwrap(); - Arc::new(AtomicU32::new(storage::db::get_u32( + Arc::new(AtomicU32::new(settings::get_u32( &c, "auto_clear_completed_secs", 0, @@ -62,7 +65,7 @@ pub fn run() { }) .setup(move |app| { let handle = app.handle().clone(); - transfer::worker::spawn_worker( + transfers::worker::spawn_worker( queue.clone(), registry.clone(), tauri::async_runtime::handle().inner().clone(), @@ -78,16 +81,14 @@ pub fn run() { let label = window.label().to_string(); let key = format!("window_state_{}", label); if let Ok(conn) = db.lock() - && let Some(json) = storage::db::get_setting(&conn, &key) + && let Some(json) = settings::get_setting(&conn, &key) && let Ok(ws) = serde_json::from_str::(&json) { drop(conn); if let (Some(x), Some(y)) = (ws.x, ws.y) { - let _ = window - .set_position(tauri::PhysicalPosition::new(x, y)); + let _ = window.set_position(tauri::PhysicalPosition::new(x, y)); } - let _ = window - .set_size(tauri::PhysicalSize::new(ws.width, ws.height)); + let _ = window.set_size(tauri::PhysicalSize::new(ws.width, ws.height)); if ws.maximized { let _ = window.maximize(); } diff --git a/src/protocols/ftp.rs b/src/protocols/ftp.rs new file mode 100644 index 0000000..b508eb9 --- /dev/null +++ b/src/protocols/ftp.rs @@ -0,0 +1,311 @@ +use anyhow::Result; +use async_trait::async_trait; +use futures_lite::io::{AsyncReadExt, AsyncWriteExt}; +use std::path::Path; +use std::str::FromStr; +use std::time::SystemTime; +use suppaftp::async_native_tls::TlsConnector; +use suppaftp::{ + AsyncFtpStream, AsyncNativeTlsConnector, AsyncNativeTlsFtpStream, list::File as FtpListFile, +}; +use tokio::fs::File; +use tokio::io::{AsyncReadExt as TokioReadExt, AsyncWriteExt as TokioWriteExt}; +use tokio::sync::Mutex; + +use crate::domain::{EntryKind, FileEntry}; +use crate::protocols::{ProgressAction, RemoteFs}; + +fn parse_list_entry(line: &str) -> Option { + let file = FtpListFile::from_str(line).ok()?; + let kind = if file.is_directory() { + EntryKind::Dir + } else if file.is_symlink() { + EntryKind::Symlink + } else { + EntryKind::File + }; + + let modified = file + .modified() + .duration_since(SystemTime::UNIX_EPOCH) + .ok() + .map(|d| d.as_secs() as i64); + + Some(FileEntry { + name: file.name().to_string(), + path: file.name().to_string(), + kind, + size: Some(file.size() as u64), + modified, + permissions: None, + }) +} + +pub struct FtpClient { + pub conn: Mutex>, +} + +pub struct FtpsClient { + pub conn: Mutex>, +} + +impl FtpClient { + pub async fn connect(host: &str, port: u16, user: &str, pass: &str) -> Result { + let addr = format!("{}:{}", host, port); + let mut stream = AsyncFtpStream::connect(&addr) + .await + .map_err(|e| anyhow::anyhow!("FTP connect failed: {}", e))?; + stream + .login(user, pass) + .await + .map_err(|e| anyhow::anyhow!("FTP login failed: {}", e))?; + Ok(Self { + conn: Mutex::new(Some(stream)), + }) + } +} + +impl FtpsClient { + pub async fn connect(host: &str, port: u16, user: &str, pass: &str) -> Result { + let addr = format!("{}:{}", host, port); + let stream = AsyncNativeTlsFtpStream::connect(&addr) + .await + .map_err(|e| anyhow::anyhow!("FTPS connect failed: {}", e))?; + let mut stream = stream + .into_secure(AsyncNativeTlsConnector::from(TlsConnector::new()), host) + .await + .map_err(|e| anyhow::anyhow!("FTPS TLS upgrade failed: {}", e))?; + stream + .login(user, pass) + .await + .map_err(|e| anyhow::anyhow!("FTPS login failed: {}", e))?; + Ok(Self { + conn: Mutex::new(Some(stream)), + }) + } +} + +const CHUNK_SIZE: usize = 64 * 1024; + +type ProgressCb = Option ProgressAction + Send>>; + +fn check(cb: &ProgressCb, total: u64) -> Result<()> { + if let Some(cb) = cb { + match cb(total) { + ProgressAction::Continue => {} + ProgressAction::Cancel => return Err(anyhow::anyhow!("cancelled")), + ProgressAction::Pause => return Err(anyhow::anyhow!("paused")), + } + } + Ok(()) +} + +macro_rules! gen_transfer { + ($upload:ident, $download:ident, $stream:ty) => { + async fn $upload( + conn: &mut $stream, + local: &str, + remote: &str, + on_progress: ProgressCb, + ) -> Result<()> { + let mut file = File::open(local) + .await + .map_err(|e| anyhow::anyhow!("cannot open local file: {}", e))?; + + let mut dstream = conn + .put_with_stream(remote) + .await + .map_err(|e| anyhow::anyhow!("FTP upload stream 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; + } + dstream + .write_all(&buf[..n]) + .await + .map_err(|e| anyhow::anyhow!("FTP upload write failed: {}", e))?; + total += n as u64; + check(&on_progress, total)?; + } + + conn.finalize_put_stream(dstream) + .await + .map_err(|e| anyhow::anyhow!("FTP upload finalize failed: {}", e))?; + Ok(()) + } + + async fn $download( + conn: &mut $stream, + remote: &str, + local: &str, + on_progress: ProgressCb, + ) -> Result<()> { + let mut dstream = conn + .retr_as_stream(remote) + .await + .map_err(|e| anyhow::anyhow!("FTP download stream failed: {}", e))?; + + let mut file = File::create(local) + .await + .map_err(|e| anyhow::anyhow!("cannot create local file: {}", e))?; + + let mut buf = vec![0u8; CHUNK_SIZE]; + let mut total = 0u64; + loop { + let n = dstream + .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; + check(&on_progress, total)?; + } + + conn.finalize_retr_stream(dstream) + .await + .map_err(|e| anyhow::anyhow!("FTP download finalize failed: {}", e))?; + Ok(()) + } + }; +} + +gen_transfer!(upload_ftp, download_ftp, AsyncFtpStream); +gen_transfer!(upload_ftps, download_ftps, AsyncNativeTlsFtpStream); + +macro_rules! impl_remote_fs { + ($client:ty, $label:literal, $upload:ident, $download:ident) => { + #[async_trait] + impl RemoteFs for $client { + async fn list(&self, path: &str) -> Result> { + let mut guard = self.conn.lock().await; + let conn = guard + .as_mut() + .ok_or_else(|| anyhow::anyhow!("not connected"))?; + let lines = conn + .list(Some(path)) + .await + .map_err(|e| anyhow::anyhow!(concat!($label, " list failed: {}"), e))?; + + let mut entries: Vec = lines + .iter() + .filter_map(|line| parse_list_entry(line)) + .collect(); + + entries.sort_by(|a, b| match (&a.kind, &b.kind) { + (EntryKind::Dir, EntryKind::Dir) => a.name.cmp(&b.name), + (EntryKind::Dir, _) => std::cmp::Ordering::Less, + (_, EntryKind::Dir) => std::cmp::Ordering::Greater, + _ => a.name.cmp(&b.name), + }); + + Ok(entries) + } + + async fn upload_with_progress( + &self, + local: &str, + remote: &str, + on_progress: Option ProgressAction + Send>>, + ) -> Result<()> { + let mut guard = self.conn.lock().await; + let conn = guard + .as_mut() + .ok_or_else(|| anyhow::anyhow!("not connected"))?; + $upload(conn, local, remote, on_progress).await + } + + 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() + .ok_or_else(|| anyhow::anyhow!("not connected"))?; + $download(conn, remote, local, on_progress).await + } + + async fn mkdir(&self, path: &str) -> Result<()> { + let mut guard = self.conn.lock().await; + let conn = guard + .as_mut() + .ok_or_else(|| anyhow::anyhow!("not connected"))?; + conn.mkdir(path) + .await + .map_err(|e| anyhow::anyhow!(concat!($label, " mkdir failed: {}"), e))?; + Ok(()) + } + + async fn rename(&self, from: &str, to: &str) -> Result<()> { + let mut guard = self.conn.lock().await; + let conn = guard + .as_mut() + .ok_or_else(|| anyhow::anyhow!("not connected"))?; + conn.rename(from, to) + .await + .map_err(|e| anyhow::anyhow!(concat!($label, " rename failed: {}"), e))?; + Ok(()) + } + + async fn delete(&self, path: &str) -> Result<()> { + let mut guard = self.conn.lock().await; + let conn = guard + .as_mut() + .ok_or_else(|| anyhow::anyhow!("not connected"))?; + if conn.rm(path).await.is_err() { + conn.rmdir(path) + .await + .map_err(|e| anyhow::anyhow!(concat!($label, " delete failed: {}"), e))?; + } + Ok(()) + } + + async fn stat(&self, path: &str) -> Result { + let mut guard = self.conn.lock().await; + let conn = guard + .as_mut() + .ok_or_else(|| anyhow::anyhow!("not connected"))?; + let size = conn + .size(path) + .await + .map_err(|e| anyhow::anyhow!(concat!($label, " stat failed: {}"), e))?; + let modified = conn + .mdtm(path) + .await + .map(|dt| dt.and_utc().timestamp()) + .ok(); + + let name = Path::new(path) + .file_name() + .map(|s| s.to_string_lossy().to_string()) + .unwrap_or_else(|| path.to_string()); + + Ok(FileEntry { + name, + path: path.to_string(), + kind: EntryKind::File, + size: Some(size as u64), + modified, + permissions: None, + }) + } + } + }; +} + +impl_remote_fs!(FtpClient, "FTP", upload_ftp, download_ftp); +impl_remote_fs!(FtpsClient, "FTPS", upload_ftps, download_ftps); diff --git a/src/protocols/ftp/mod.rs b/src/protocols/ftp/mod.rs deleted file mode 100644 index 785bd68..0000000 --- a/src/protocols/ftp/mod.rs +++ /dev/null @@ -1,84 +0,0 @@ -use anyhow::Result; -use std::str::FromStr; -use std::time::SystemTime; -use tokio::sync::Mutex; - -use suppaftp::async_native_tls::TlsConnector; -use suppaftp::{ - AsyncFtpStream, AsyncNativeTlsConnector, AsyncNativeTlsFtpStream, list::File as FtpListFile, -}; - -use crate::domain::file_entry::{EntryKind, FileEntry}; - -mod ops; -mod transfer; - -pub struct FtpClient { - pub conn: Mutex>, -} - -pub struct FtpsClient { - pub conn: Mutex>, -} - -pub(super) fn parse_list_entry(line: &str) -> Option { - let file = FtpListFile::from_str(line).ok()?; - let kind = if file.is_directory() { - EntryKind::Dir - } else if file.is_symlink() { - EntryKind::Symlink - } else { - EntryKind::File - }; - - let modified = file - .modified() - .duration_since(SystemTime::UNIX_EPOCH) - .ok() - .map(|d| d.as_secs() as i64); - - Some(FileEntry { - name: file.name().to_string(), - path: file.name().to_string(), - kind, - size: Some(file.size() as u64), - modified, - permissions: None, - }) -} - -impl FtpClient { - pub async fn connect(host: &str, port: u16, user: &str, pass: &str) -> Result { - let addr = format!("{}:{}", host, port); - let mut stream = AsyncFtpStream::connect(&addr) - .await - .map_err(|e| anyhow::anyhow!("FTP connect failed: {}", e))?; - stream - .login(user, pass) - .await - .map_err(|e| anyhow::anyhow!("FTP login failed: {}", e))?; - Ok(Self { - conn: Mutex::new(Some(stream)), - }) - } -} - -impl FtpsClient { - pub async fn connect(host: &str, port: u16, user: &str, pass: &str) -> Result { - let addr = format!("{}:{}", host, port); - let stream = AsyncNativeTlsFtpStream::connect(&addr) - .await - .map_err(|e| anyhow::anyhow!("FTPS connect failed: {}", e))?; - let mut stream = stream - .into_secure(AsyncNativeTlsConnector::from(TlsConnector::new()), host) - .await - .map_err(|e| anyhow::anyhow!("FTPS TLS upgrade failed: {}", e))?; - stream - .login(user, pass) - .await - .map_err(|e| anyhow::anyhow!("FTPS login failed: {}", e))?; - Ok(Self { - conn: Mutex::new(Some(stream)), - }) - } -} diff --git a/src/protocols/ftp/ops.rs b/src/protocols/ftp/ops.rs deleted file mode 100644 index 0bb33be..0000000 --- a/src/protocols/ftp/ops.rs +++ /dev/null @@ -1,133 +0,0 @@ -//! Реализация `RemoteFs` для FTP. Чанкованные upload/download — в [`super::transfer`]. -use anyhow::Result; -use async_trait::async_trait; - -use super::{FtpClient, FtpsClient, parse_list_entry, transfer}; -use crate::domain::file_entry::{EntryKind, FileEntry}; -use crate::protocols::{ProgressAction, RemoteFs}; - -macro_rules! impl_remote_fs { - ($client:ty, $label:literal, $upload:ident, $download:ident) => { - #[async_trait] - impl RemoteFs for $client { - async fn list(&self, path: &str) -> Result> { - let mut guard = self.conn.lock().await; - let conn = guard - .as_mut() - .ok_or_else(|| anyhow::anyhow!("not connected"))?; - let lines = conn - .list(Some(path)) - .await - .map_err(|e| anyhow::anyhow!(concat!($label, " list failed: {}"), e))?; - - let mut entries: Vec = lines - .iter() - .filter_map(|line| parse_list_entry(line)) - .collect(); - - entries.sort_by(|a, b| match (&a.kind, &b.kind) { - (EntryKind::Dir, EntryKind::Dir) => a.name.cmp(&b.name), - (EntryKind::Dir, _) => std::cmp::Ordering::Less, - (_, EntryKind::Dir) => std::cmp::Ordering::Greater, - _ => a.name.cmp(&b.name), - }); - - Ok(entries) - } - - async fn upload_with_progress( - &self, - local: &str, - remote: &str, - on_progress: Option ProgressAction + Send>>, - ) -> Result<()> { - let mut guard = self.conn.lock().await; - let conn = guard - .as_mut() - .ok_or_else(|| anyhow::anyhow!("not connected"))?; - transfer::$upload(conn, local, remote, on_progress).await - } - - 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() - .ok_or_else(|| anyhow::anyhow!("not connected"))?; - transfer::$download(conn, remote, local, on_progress).await - } - - async fn mkdir(&self, path: &str) -> Result<()> { - let mut guard = self.conn.lock().await; - let conn = guard - .as_mut() - .ok_or_else(|| anyhow::anyhow!("not connected"))?; - conn.mkdir(path) - .await - .map_err(|e| anyhow::anyhow!(concat!($label, " mkdir failed: {}"), e))?; - Ok(()) - } - - async fn rename(&self, from: &str, to: &str) -> Result<()> { - let mut guard = self.conn.lock().await; - let conn = guard - .as_mut() - .ok_or_else(|| anyhow::anyhow!("not connected"))?; - conn.rename(from, to) - .await - .map_err(|e| anyhow::anyhow!(concat!($label, " rename failed: {}"), e))?; - Ok(()) - } - - async fn delete(&self, path: &str) -> Result<()> { - let mut guard = self.conn.lock().await; - let conn = guard - .as_mut() - .ok_or_else(|| anyhow::anyhow!("not connected"))?; - if conn.rm(path).await.is_err() { - conn.rmdir(path) - .await - .map_err(|e| anyhow::anyhow!(concat!($label, " delete failed: {}"), e))?; - } - Ok(()) - } - - async fn stat(&self, path: &str) -> Result { - let mut guard = self.conn.lock().await; - let conn = guard - .as_mut() - .ok_or_else(|| anyhow::anyhow!("not connected"))?; - let size = conn - .size(path) - .await - .map_err(|e| anyhow::anyhow!(concat!($label, " stat failed: {}"), e))?; - let modified = conn - .mdtm(path) - .await - .map(|dt| dt.and_utc().timestamp()) - .ok(); - - let name = std::path::Path::new(path) - .file_name() - .map(|s| s.to_string_lossy().to_string()) - .unwrap_or_else(|| path.to_string()); - - Ok(FileEntry { - name, - path: path.to_string(), - kind: EntryKind::File, - size: Some(size as u64), - modified, - permissions: None, - }) - } - } - }; -} - -impl_remote_fs!(FtpClient, "FTP", upload_ftp, download_ftp); -impl_remote_fs!(FtpsClient, "FTPS", upload_ftps, download_ftps); diff --git a/src/protocols/ftp/transfer.rs b/src/protocols/ftp/transfer.rs deleted file mode 100644 index e8c3e17..0000000 --- a/src/protocols/ftp/transfer.rs +++ /dev/null @@ -1,107 +0,0 @@ -//! Чанкованная загрузка/выгрузка по FTP с колбэком прогресса (pause/cancel). -use anyhow::Result; -use futures_lite::io::{AsyncReadExt, AsyncWriteExt}; -use suppaftp::{AsyncFtpStream, AsyncNativeTlsFtpStream}; -use tokio::fs::File; -use tokio::io::{AsyncReadExt as TokioReadExt, AsyncWriteExt as TokioWriteExt}; - -use crate::protocols::ProgressAction; - -const CHUNK_SIZE: usize = 64 * 1024; - -type ProgressCb = Option ProgressAction + Send>>; - -fn check(cb: &ProgressCb, total: u64) -> Result<()> { - if let Some(cb) = cb { - match cb(total) { - ProgressAction::Continue => {} - ProgressAction::Cancel => return Err(anyhow::anyhow!("cancelled")), - ProgressAction::Pause => return Err(anyhow::anyhow!("paused")), - } - } - Ok(()) -} - -macro_rules! gen_transfer { - ($upload:ident, $download:ident, $stream:ty) => { - pub(super) async fn $upload( - conn: &mut $stream, - local: &str, - remote: &str, - on_progress: ProgressCb, - ) -> Result<()> { - let mut file = File::open(local) - .await - .map_err(|e| anyhow::anyhow!("cannot open local file: {}", e))?; - - let mut dstream = conn - .put_with_stream(remote) - .await - .map_err(|e| anyhow::anyhow!("FTP upload stream 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; - } - dstream - .write_all(&buf[..n]) - .await - .map_err(|e| anyhow::anyhow!("FTP upload write failed: {}", e))?; - total += n as u64; - check(&on_progress, total)?; - } - - conn.finalize_put_stream(dstream) - .await - .map_err(|e| anyhow::anyhow!("FTP upload finalize failed: {}", e))?; - Ok(()) - } - - pub(super) async fn $download( - conn: &mut $stream, - remote: &str, - local: &str, - on_progress: ProgressCb, - ) -> Result<()> { - let mut dstream = conn - .retr_as_stream(remote) - .await - .map_err(|e| anyhow::anyhow!("FTP download stream failed: {}", e))?; - - let mut file = File::create(local) - .await - .map_err(|e| anyhow::anyhow!("cannot create local file: {}", e))?; - - let mut buf = vec![0u8; CHUNK_SIZE]; - let mut total = 0u64; - loop { - let n = dstream - .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; - check(&on_progress, total)?; - } - - conn.finalize_retr_stream(dstream) - .await - .map_err(|e| anyhow::anyhow!("FTP download finalize failed: {}", e))?; - Ok(()) - } - }; -} - -gen_transfer!(upload_ftp, download_ftp, AsyncFtpStream); -gen_transfer!(upload_ftps, download_ftps, AsyncNativeTlsFtpStream); diff --git a/src/protocols/mod.rs b/src/protocols/mod.rs index 0f2a5e6..582a682 100644 --- a/src/protocols/mod.rs +++ b/src/protocols/mod.rs @@ -1,8 +1,10 @@ pub mod ftp; -pub use ftp::FtpsClient; pub mod sftp; -use crate::domain::file_entry::FileEntry; +pub use ftp::{FtpClient, FtpsClient}; +pub use sftp::SftpClient; + +use crate::domain::FileEntry; use anyhow::Result; use async_trait::async_trait; diff --git a/src/protocols/sftp.rs b/src/protocols/sftp.rs new file mode 100644 index 0000000..77a68df --- /dev/null +++ b/src/protocols/sftp.rs @@ -0,0 +1,398 @@ +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 crate::domain::{EntryKind, FileEntry}; +use crate::protocols::{ProgressAction, RemoteFs}; + +fn known_hosts_path() -> Option { + dirs::home_dir().map(|p| p.join(".ssh/known_hosts")) +} + +fn fingerprint_hex(session: &Session) -> String { + session + .host_key_hash(HashType::Sha256) + .map(|h| { + h.iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(":") + }) + .unwrap_or_else(|| "unknown".into()) +} + +fn verify_host_key(session: &Session, host: &str, port: u16) -> Result<()> { + let mut known = session.known_hosts().context("known_hosts init failed")?; + + if let Some(ref path) = known_hosts_path() { + let _ = known.read_file(path, KnownHostFileKind::OpenSSH); + } + + let (key, key_type) = session + .host_key() + .context("no host key received from server")?; + + match known.check_port(host, port, key) { + CheckResult::Match => {} + CheckResult::Mismatch => { + bail!( + "SSH host key mismatch for {}!\n\ + The server's host key has changed since the last connection.\n\ + This could mean someone is intercepting the connection (MITM attack).\n\ + Fingerprint (SHA256): {}", + host, + fingerprint_hex(session) + ); + } + CheckResult::NotFound => { + tracing::info!( + "Unknown host key for {}, adding to known_hosts (SHA256: {})", + host, + fingerprint_hex(session) + ); + known + .add(host, key, "wherry", key_type.into()) + .context("failed to add host key to known_hosts")?; + if let Some(ref path) = known_hosts_path() + && let Err(e) = known.write_file(path, KnownHostFileKind::OpenSSH) + { + tracing::warn!("failed to write known_hosts: {}", e); + } + } + CheckResult::Failure => { + bail!("known_hosts check failed for {}", host); + } + } + + Ok(()) +} + +pub struct SftpClient { + session: Arc>, +} + +impl SftpClient { + pub fn connect_password(host: &str, port: u16, user: &str, password: &str) -> Result { + let session = handshake(host, port)?; + session + .userauth_password(user, password) + .context("SSH password auth failed")?; + Ok(Self { + session: Arc::new(Mutex::new(session)), + }) + } + + pub fn connect_key( + host: &str, + port: u16, + user: &str, + key_path: &str, + passphrase: Option<&str>, + ) -> Result { + let session = handshake(host, port)?; + session + .userauth_pubkey_file(user, None, Path::new(key_path), passphrase) + .with_context(|| format!("SSH key auth failed ({})", key_path))?; + Ok(Self { + session: Arc::new(Mutex::new(session)), + }) + } + + pub fn connect_auto(host: &str, port: u16, user: &str) -> Result { + let session = handshake(host, port)?; + + if session.userauth_agent(user).is_ok() && session.authenticated() { + return Ok(Self { + session: Arc::new(Mutex::new(session)), + }); + } + + let ssh_dir = dirs::home_dir() + .context("Cannot resolve home directory")? + .join(".ssh"); + let mut tried = vec!["ssh-agent".to_string()]; + for name in ["id_ed25519", "id_ecdsa", "id_rsa"] { + let key = ssh_dir.join(name); + if !key.is_file() { + continue; + } + if session.userauth_pubkey_file(user, None, &key, None).is_ok() + && session.authenticated() + { + return Ok(Self { + session: Arc::new(Mutex::new(session)), + }); + } + tried.push(format!("~/.ssh/{}", name)); + } + + anyhow::bail!( + "SSH auth failed: no password given, and none of [{}] were accepted. \ + Pick a key file manually or enter a password.", + tried.join(", ") + ) + } +} + +fn handshake(host: &str, port: u16) -> Result { + let tcp = TcpStream::connect(format!("{}:{}", host, port)).context("TCP connect failed")?; + let mut session = Session::new().context("SSH session init failed")?; + session.set_tcp_stream(tcp); + session.handshake().context("SSH handshake failed")?; + verify_host_key(&session, host, port)?; + Ok(session) +} + +const CHUNK_SIZE: usize = 64 * 1024; + +type ProgressCb = Option ProgressAction + Send>>; + +fn check(cb: &ProgressCb, total: u64) -> Result<()> { + if let Some(cb) = cb { + match cb(total) { + ProgressAction::Continue => {} + ProgressAction::Cancel => return Err(anyhow::anyhow!("cancelled")), + ProgressAction::Pause => return Err(anyhow::anyhow!("paused")), + } + } + Ok(()) +} + +async fn upload( + session: Arc>, + local: String, + remote: String, + on_progress: ProgressCb, +) -> Result<()> { + 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; + } + remote_file + .write_all(&buf[..n]) + .context("write remote file failed")?; + total += n as u64; + check(&on_progress, total)?; + } + Ok(()) + }) + .await + .map_err(|e| anyhow::anyhow!("SFTP upload task failed: {}", e))? +} + +async fn download( + session: Arc>, + remote: String, + local: String, + on_progress: ProgressCb, +) -> Result<()> { + 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")?; + + 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; + check(&on_progress, total)?; + } + Ok(()) + }) + .await + .map_err(|e| anyhow::anyhow!("SFTP download task failed: {}", e))? +} + +#[async_trait] +impl RemoteFs for SftpClient { + async fn list(&self, path: &str) -> Result> { + 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, + } + }) + .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<()> { + upload( + self.session.clone(), + local.to_string(), + remote.to_string(), + on_progress, + ) + .await + } + + async fn download_with_progress( + &self, + remote: &str, + local: &str, + on_progress: Option ProgressAction + Send>>, + ) -> Result<()> { + download( + self.session.clone(), + remote.to_string(), + local.to_string(), + on_progress, + ) + .await + } + + async fn mkdir(&self, path: &str) -> Result<()> { + 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 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 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 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, + }) + }) + .await + .map_err(|e| anyhow::anyhow!("SFTP stat task failed: {}", e))? + } +} diff --git a/src/protocols/sftp/hostkey.rs b/src/protocols/sftp/hostkey.rs deleted file mode 100644 index 26eb050..0000000 --- a/src/protocols/sftp/hostkey.rs +++ /dev/null @@ -1,66 +0,0 @@ -//! Проверка SSH host key через ~/.ssh/known_hosts (TOFU: неизвестный ключ добавляется). -use anyhow::{Context, Result, bail}; -use ssh2::{CheckResult, HashType, KnownHostFileKind, Session}; -use std::path::PathBuf; - -fn known_hosts_path() -> Option { - dirs::home_dir().map(|p| p.join(".ssh/known_hosts")) -} - -fn fingerprint_hex(session: &Session) -> String { - session - .host_key_hash(HashType::Sha256) - .map(|h| { - h.iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(":") - }) - .unwrap_or_else(|| "unknown".into()) -} - -pub(super) fn verify_host_key(session: &Session, host: &str, port: u16) -> Result<()> { - let mut known = session.known_hosts().context("known_hosts init failed")?; - - if let Some(ref path) = known_hosts_path() { - let _ = known.read_file(path, KnownHostFileKind::OpenSSH); - } - - let (key, key_type) = session - .host_key() - .context("no host key received from server")?; - - match known.check_port(host, port, key) { - CheckResult::Match => {} // known and verified - CheckResult::Mismatch => { - bail!( - "SSH host key mismatch for {}!\n\ - The server's host key has changed since the last connection.\n\ - This could mean someone is intercepting the connection (MITM attack).\n\ - Fingerprint (SHA256): {}", - host, - fingerprint_hex(session) - ); - } - CheckResult::NotFound => { - tracing::info!( - "Unknown host key for {}, adding to known_hosts (SHA256: {})", - host, - fingerprint_hex(session) - ); - known - .add(host, key, "wherry", key_type.into()) - .context("failed to add host key to known_hosts")?; - if let Some(ref path) = known_hosts_path() - && let Err(e) = known.write_file(path, KnownHostFileKind::OpenSSH) - { - tracing::warn!("failed to write known_hosts: {}", e); - } - } - CheckResult::Failure => { - bail!("known_hosts check failed for {}", host); - } - } - - Ok(()) -} diff --git a/src/protocols/sftp/mod.rs b/src/protocols/sftp/mod.rs deleted file mode 100644 index 711311e..0000000 --- a/src/protocols/sftp/mod.rs +++ /dev/null @@ -1,93 +0,0 @@ -use anyhow::{Context, Result}; -use ssh2::Session; -use std::net::TcpStream; -use std::path::Path; -use std::sync::{Arc, Mutex}; - -mod hostkey; -mod ops; -mod transfer; - -pub struct SftpClient { - session: Arc>, -} - -impl SftpClient { - /// Подключение по паролю - pub fn connect_password(host: &str, port: u16, user: &str, password: &str) -> Result { - let session = handshake(host, port)?; - session - .userauth_password(user, password) - .context("SSH password auth failed")?; - Ok(Self { - session: Arc::new(Mutex::new(session)), - }) - } - - /// Подключение по явно выбранному ключу. `passphrase` — пароль из формы - /// (если задан): им расшифровывается защищённый паролем ключ. - pub fn connect_key( - host: &str, - port: u16, - user: &str, - key_path: &str, - passphrase: Option<&str>, - ) -> Result { - let session = handshake(host, port)?; - session - .userauth_pubkey_file(user, None, Path::new(key_path), passphrase) - .with_context(|| format!("SSH key auth failed ({})", key_path))?; - Ok(Self { - session: Arc::new(Mutex::new(session)), - }) - } - - /// Автоматический подбор SSH-аутентификации, когда пароль не задан и ключ - /// не выбран: сперва ssh-agent, затем стандартные ключи из ~/.ssh. - pub fn connect_auto(host: &str, port: u16, user: &str) -> Result { - let session = handshake(host, port)?; - - // 1) ssh-agent — покрывает ключи с passphrase, добавленные в агент. - if session.userauth_agent(user).is_ok() && session.authenticated() { - return Ok(Self { - session: Arc::new(Mutex::new(session)), - }); - } - - // 2) стандартные ключи в ~/.ssh (без passphrase). - let ssh_dir = dirs::home_dir() - .context("Cannot resolve home directory")? - .join(".ssh"); - let mut tried = vec!["ssh-agent".to_string()]; - for name in ["id_ed25519", "id_ecdsa", "id_rsa"] { - let key = ssh_dir.join(name); - if !key.is_file() { - continue; - } - if session.userauth_pubkey_file(user, None, &key, None).is_ok() - && session.authenticated() - { - return Ok(Self { - session: Arc::new(Mutex::new(session)), - }); - } - tried.push(format!("~/.ssh/{}", name)); - } - - anyhow::bail!( - "SSH auth failed: no password given, and none of [{}] were accepted. \ - Pick a key file manually or enter a password.", - tried.join(", ") - ) - } -} - -/// TCP + SSH handshake + проверка host key (общий пролог обоих способов аутентификации). -fn handshake(host: &str, port: u16) -> Result { - let tcp = TcpStream::connect(format!("{}:{}", host, port)).context("TCP connect failed")?; - let mut session = Session::new().context("SSH session init failed")?; - session.set_tcp_stream(tcp); - session.handshake().context("SSH handshake failed")?; - hostkey::verify_host_key(&session, host, port)?; - Ok(session) -} diff --git a/src/protocols/sftp/ops.rs b/src/protocols/sftp/ops.rs deleted file mode 100644 index 13054e3..0000000 --- a/src/protocols/sftp/ops.rs +++ /dev/null @@ -1,169 +0,0 @@ -//! Реализация `RemoteFs` для SFTP: всё через `spawn_blocking` (ssh2 синхронный). -//! Чанкованные upload/download вынесены в [`super::transfer`]. -use anyhow::{Context, Result}; -use async_trait::async_trait; -use std::path::Path; - -use super::{SftpClient, transfer}; -use crate::domain::file_entry::{EntryKind, FileEntry}; -use crate::protocols::{ProgressAction, RemoteFs}; - -#[async_trait] -impl RemoteFs for SftpClient { - async fn list(&self, path: &str) -> Result> { - 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, - } - }) - .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<()> { - transfer::upload( - self.session.clone(), - local.to_string(), - remote.to_string(), - on_progress, - ) - .await - } - - async fn download_with_progress( - &self, - remote: &str, - local: &str, - on_progress: Option ProgressAction + Send>>, - ) -> Result<()> { - transfer::download( - self.session.clone(), - remote.to_string(), - local.to_string(), - on_progress, - ) - .await - } - - async fn mkdir(&self, path: &str) -> Result<()> { - 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 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 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 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, - }) - }) - .await - .map_err(|e| anyhow::anyhow!("SFTP stat task failed: {}", e))? - } -} diff --git a/src/protocols/sftp/transfer.rs b/src/protocols/sftp/transfer.rs deleted file mode 100644 index 148353c..0000000 --- a/src/protocols/sftp/transfer.rs +++ /dev/null @@ -1,100 +0,0 @@ -//! Чанкованная загрузка/выгрузка по SFTP с колбэком прогресса (pause/cancel). -use anyhow::{Context, Result}; -use ssh2::Session; -use std::io::{Read, Write}; -use std::path::Path; -use std::sync::{Arc, Mutex}; - -use crate::protocols::ProgressAction; - -const CHUNK_SIZE: usize = 64 * 1024; - -type ProgressCb = Option ProgressAction + Send>>; - -/// Проверяет действие прогресса; возвращает Err для pause/cancel (прерывает цикл). -fn check(cb: &ProgressCb, total: u64) -> Result<()> { - if let Some(cb) = cb { - match cb(total) { - ProgressAction::Continue => {} - ProgressAction::Cancel => return Err(anyhow::anyhow!("cancelled")), - ProgressAction::Pause => return Err(anyhow::anyhow!("paused")), - } - } - Ok(()) -} - -pub(super) async fn upload( - session: Arc>, - local: String, - remote: String, - on_progress: ProgressCb, -) -> Result<()> { - 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; - } - remote_file - .write_all(&buf[..n]) - .context("write remote file failed")?; - total += n as u64; - check(&on_progress, total)?; - } - Ok(()) - }) - .await - .map_err(|e| anyhow::anyhow!("SFTP upload task failed: {}", e))? -} - -pub(super) async fn download( - session: Arc>, - remote: String, - local: String, - on_progress: ProgressCb, -) -> Result<()> { - 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")?; - - 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; - check(&on_progress, total)?; - } - Ok(()) - }) - .await - .map_err(|e| anyhow::anyhow!("SFTP download task failed: {}", e))? -} diff --git a/src/settings.rs b/src/settings.rs new file mode 100644 index 0000000..e54b568 --- /dev/null +++ b/src/settings.rs @@ -0,0 +1,119 @@ +use anyhow::Result; +use rusqlite::{Connection, params}; + +// ── App settings ── + +pub fn get_setting(conn: &Connection, key: &str) -> Option { + conn.query_row( + "SELECT value FROM app_settings WHERE key = ?1", + params![key], + |row| row.get(0), + ) + .ok() +} + +pub fn set_setting(conn: &Connection, key: &str, value: &str) -> Result<()> { + conn.execute( + "INSERT INTO app_settings (key, value) VALUES (?1, ?2) + ON CONFLICT(key) DO UPDATE SET value = excluded.value", + params![key, value], + )?; + Ok(()) +} + +pub fn get_bool(conn: &Connection, key: &str, default: bool) -> bool { + get_setting(conn, key) + .and_then(|v| v.parse().ok()) + .unwrap_or(default) +} + +pub fn get_u32(conn: &Connection, key: &str, default: u32) -> u32 { + get_setting(conn, key) + .and_then(|v| v.parse().ok()) + .unwrap_or(default) +} + +pub fn set_bool(conn: &Connection, key: &str, value: bool) { + let _ = set_setting(conn, key, &value.to_string()); +} + +pub fn set_u32(conn: &Connection, key: &str, value: u32) { + let _ = set_setting(conn, key, &value.to_string()); +} + +// ── Bookmarks ── + +pub fn get_bookmarks(conn: &Connection) -> Result> { + let mut stmt = conn.prepare("SELECT id, name, path FROM bookmarks ORDER BY id")?; + let rows = stmt + .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))? + .filter_map(|r| r.ok()) + .collect(); + Ok(rows) +} + +pub fn add_bookmark(conn: &Connection, name: &str, path: &str) -> Result { + conn.execute( + "INSERT INTO bookmarks (name, path) VALUES (?1, ?2)", + params![name, path], + )?; + Ok(conn.last_insert_rowid()) +} + +pub fn remove_bookmark(conn: &Connection, id: i64) -> Result<()> { + conn.execute("DELETE FROM bookmarks WHERE id = ?1", params![id])?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn test_db() -> rusqlite::Connection { + let conn = rusqlite::Connection::open_in_memory().unwrap(); + crate::storage::init_tables(&conn).unwrap(); + conn + } + + #[test] + fn test_add_get_remove_bookmark() { + let conn = test_db(); + let id = add_bookmark(&conn, "Projects", "/home/user/projects").unwrap(); + assert!(id > 0); + + let bookmarks = get_bookmarks(&conn).unwrap(); + assert_eq!(bookmarks.len(), 1); + assert_eq!(bookmarks[0].1, "Projects"); + assert_eq!(bookmarks[0].2, "/home/user/projects"); + + remove_bookmark(&conn, id).unwrap(); + assert!(get_bookmarks(&conn).unwrap().is_empty()); + } + + #[test] + fn test_settings_string() { + let conn = test_db(); + assert!(get_setting(&conn, "theme").is_none()); + set_setting(&conn, "theme", "dark").unwrap(); + assert_eq!(get_setting(&conn, "theme").unwrap(), "dark"); + set_setting(&conn, "theme", "light").unwrap(); + assert_eq!(get_setting(&conn, "theme").unwrap(), "light"); + } + + #[test] + fn test_settings_bool() { + let conn = test_db(); + assert!(!get_bool(&conn, "show_hidden", false)); + assert!(get_bool(&conn, "show_hidden", true)); + set_bool(&conn, "show_hidden", true); + assert!(get_bool(&conn, "show_hidden", false)); + } + + #[test] + fn test_settings_u32() { + let conn = test_db(); + assert_eq!(get_u32(&conn, "timeout", 30), 30); + set_u32(&conn, "timeout", 60); + assert_eq!(get_u32(&conn, "timeout", 30), 60); + } +} diff --git a/src/storage/db/mod.rs b/src/storage.rs similarity index 57% rename from src/storage/db/mod.rs rename to src/storage.rs index e74f389..13daaca 100644 --- a/src/storage/db/mod.rs +++ b/src/storage.rs @@ -1,21 +1,8 @@ -//! SQLite-хранилище: схема + операции по таблицам (bookmarks/history/sites). -use crate::domain::connection::Protocol; +use crate::domain::{Protocol, Site}; use anyhow::Result; -use rusqlite::Connection; +use rusqlite::{Connection, params}; -mod bookmarks; -mod history; -mod settings; -mod sites; - -pub use bookmarks::{add_bookmark, get_bookmarks, remove_bookmark}; -pub use history::{ - add_history_entry, clear_history, find_history_conn_id, get_history, HistoryRow, -}; -pub use settings::{get_bool, get_setting, get_u32, set_bool, set_setting, set_u32}; -pub use sites::{delete_site, get_sites, save_site}; - -pub(super) fn protocol_to_str(protocol: &Protocol) -> &'static str { +fn protocol_to_str(protocol: &Protocol) -> &'static str { match protocol { Protocol::Sftp => "sftp", Protocol::Ftp => "ftp", @@ -23,7 +10,7 @@ pub(super) fn protocol_to_str(protocol: &Protocol) -> &'static str { } } -pub(super) fn protocol_from_str(s: &str) -> Protocol { +fn protocol_from_str(s: &str) -> Protocol { match s { "sftp" => Protocol::Sftp, "ftps" => Protocol::Ftps, @@ -70,7 +57,6 @@ pub fn init_tables(conn: &Connection) -> Result<()> { ", )?; - // Миграции let _ = conn.execute("ALTER TABLE connection_history ADD COLUMN conn_id TEXT", []); let _ = conn.execute( "ALTER TABLE connection_history ADD COLUMN protocol TEXT NOT NULL DEFAULT 'sftp'", @@ -81,8 +67,6 @@ pub fn init_tables(conn: &Connection) -> Result<()> { [], ); let _ = conn.execute("ALTER TABLE sites ADD COLUMN password TEXT", []); - // На случай старых записей без уникальности (host,port,username) — оставляем - // только самую свежую строку на каждую цель перед созданием индекса. conn.execute( "DELETE FROM connection_history WHERE id NOT IN ( SELECT MAX(id) FROM connection_history GROUP BY host, port, username @@ -98,10 +82,152 @@ pub fn init_tables(conn: &Connection) -> Result<()> { Ok(()) } +// ── Sites ── + +pub fn get_sites(conn: &Connection) -> Result> { + let mut stmt = conn.prepare( + "SELECT id, name, protocol, host, port, username, password, key_path, folder, note FROM sites ORDER BY name", + )?; + let sites = stmt + .query_map([], |row| { + let proto_str: String = row.get(2)?; + let protocol = protocol_from_str(&proto_str); + Ok(Site { + id: row.get(0)?, + name: row.get(1)?, + protocol, + host: row.get(3)?, + port: row.get(4)?, + username: row.get(5)?, + password: row.get(6)?, + key_path: row.get(7)?, + folder: row.get(8)?, + note: row.get(9)?, + }) + })? + .filter_map(|r| r.ok()) + .collect(); + Ok(sites) +} + +pub fn save_site(conn: &Connection, site: &Site) -> Result<()> { + let proto = protocol_to_str(&site.protocol); + conn.execute( + "INSERT OR REPLACE INTO sites + (id, name, protocol, host, port, username, password, key_path, folder, note) + VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10)", + params![ + site.id, + site.name, + proto, + site.host, + site.port, + site.username, + site.password, + site.key_path, + site.folder, + site.note + ], + )?; + Ok(()) +} + +pub fn delete_site(conn: &Connection, id: &str) -> Result<()> { + conn.execute("DELETE FROM sites WHERE id = ?1", params![id])?; + Ok(()) +} + +// ── History ── + +pub type HistoryRow = ( + String, + u16, + String, + String, + String, + Protocol, + Option, +); + +pub fn find_history_conn_id( + conn: &Connection, + host: &str, + port: u16, + username: &str, +) -> Result> { + let id: Option = conn + .query_row( + "SELECT conn_id FROM connection_history WHERE host = ?1 AND port = ?2 AND username = ?3", + params![host, port, username], + |row| row.get(0), + ) + .ok(); + Ok(id) +} + +#[allow(clippy::too_many_arguments)] +pub fn add_history_entry( + conn: &Connection, + host: &str, + port: u16, + username: &str, + conn_id: &str, + protocol: &Protocol, + key_path: Option<&str>, +) -> Result<()> { + use chrono::Local; + let connected_at = Local::now().format("%H:%M %d.%m").to_string(); + let proto = protocol_to_str(protocol); + conn.execute( + "INSERT INTO connection_history (host, port, username, connected_at, conn_id, protocol, key_path) + VALUES (?1,?2,?3,?4,?5,?6,?7) + ON CONFLICT(host, port, username) DO UPDATE SET + connected_at = excluded.connected_at, + protocol = excluded.protocol, + key_path = excluded.key_path, + conn_id = COALESCE(connection_history.conn_id, excluded.conn_id)", + params![host, port, username, connected_at, conn_id, proto, key_path], + )?; + conn.execute( + "DELETE FROM connection_history WHERE id NOT IN ( + SELECT id FROM connection_history ORDER BY id DESC LIMIT 20 + )", + [], + )?; + Ok(()) +} + +pub fn clear_history(conn: &Connection) -> Result<()> { + conn.execute("DELETE FROM connection_history", [])?; + Ok(()) +} + +pub fn get_history(conn: &Connection) -> Result> { + let mut stmt = conn.prepare( + "SELECT host, port, username, connected_at, conn_id, protocol, key_path + FROM connection_history ORDER BY id DESC LIMIT 20", + )?; + let rows = stmt + .query_map([], |row| { + let proto_str: String = row.get(5)?; + Ok(( + row.get(0)?, + row.get(1)?, + row.get(2)?, + row.get(3)?, + row.get(4)?, + protocol_from_str(&proto_str), + row.get(6)?, + )) + })? + .filter_map(|r| r.ok()) + .collect(); + Ok(rows) +} + #[cfg(test)] mod tests { use super::*; - use crate::domain::site::Site; fn test_db() -> rusqlite::Connection { let conn = rusqlite::Connection::open_in_memory().unwrap(); @@ -204,51 +330,8 @@ mod tests { let protocol = Protocol::Sftp; add_history_entry(&conn, "example.com", 22, "admin", "conn-1", &protocol, None).unwrap(); add_history_entry(&conn, "example.com", 22, "admin", "conn-2", &protocol, None).unwrap(); - // same target keeps only the latest; conn_id preserved from first insert let history = get_history(&conn).unwrap(); assert_eq!(history.len(), 1); assert_eq!(history[0].4, "conn-1"); } - - #[test] - fn test_add_get_remove_bookmark() { - let conn = test_db(); - let id = add_bookmark(&conn, "Projects", "/home/user/projects").unwrap(); - assert!(id > 0); - - let bookmarks = get_bookmarks(&conn).unwrap(); - assert_eq!(bookmarks.len(), 1); - assert_eq!(bookmarks[0].1, "Projects"); - assert_eq!(bookmarks[0].2, "/home/user/projects"); - - remove_bookmark(&conn, id).unwrap(); - assert!(get_bookmarks(&conn).unwrap().is_empty()); - } - - #[test] - fn test_settings_string() { - let conn = test_db(); - assert!(get_setting(&conn, "theme").is_none()); - set_setting(&conn, "theme", "dark").unwrap(); - assert_eq!(get_setting(&conn, "theme").unwrap(), "dark"); - set_setting(&conn, "theme", "light").unwrap(); - assert_eq!(get_setting(&conn, "theme").unwrap(), "light"); - } - - #[test] - fn test_settings_bool() { - let conn = test_db(); - assert!(!get_bool(&conn, "show_hidden", false)); - assert!(get_bool(&conn, "show_hidden", true)); - set_bool(&conn, "show_hidden", true); - assert!(get_bool(&conn, "show_hidden", false)); - } - - #[test] - fn test_settings_u32() { - let conn = test_db(); - assert_eq!(get_u32(&conn, "timeout", 30), 30); - set_u32(&conn, "timeout", 60); - assert_eq!(get_u32(&conn, "timeout", 30), 60); - } } diff --git a/src/storage/db/bookmarks.rs b/src/storage/db/bookmarks.rs deleted file mode 100644 index 3d92dea..0000000 --- a/src/storage/db/bookmarks.rs +++ /dev/null @@ -1,25 +0,0 @@ -//! Таблица закладок локальных папок. -use anyhow::Result; -use rusqlite::{params, Connection}; - -pub fn get_bookmarks(conn: &Connection) -> Result> { - let mut stmt = conn.prepare("SELECT id, name, path FROM bookmarks ORDER BY id")?; - let rows = stmt - .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))? - .filter_map(|r| r.ok()) - .collect(); - Ok(rows) -} - -pub fn add_bookmark(conn: &Connection, name: &str, path: &str) -> Result { - conn.execute( - "INSERT INTO bookmarks (name, path) VALUES (?1, ?2)", - params![name, path], - )?; - Ok(conn.last_insert_rowid()) -} - -pub fn remove_bookmark(conn: &Connection, id: i64) -> Result<()> { - conn.execute("DELETE FROM bookmarks WHERE id = ?1", params![id])?; - Ok(()) -} diff --git a/src/storage/db/history.rs b/src/storage/db/history.rs deleted file mode 100644 index f3b136c..0000000 --- a/src/storage/db/history.rs +++ /dev/null @@ -1,98 +0,0 @@ -//! Таблица истории подключений (последние 20 целей, уникальных по host/port/user). -use anyhow::Result; -use rusqlite::{params, Connection}; - -use super::{protocol_from_str, protocol_to_str}; -use crate::domain::connection::Protocol; - -/// Одна запись истории: host, port, username, время, conn_id (стабильный id — -/// под ним же лежит пароль в keychain), протокол, key_path. -pub type HistoryRow = ( - String, - u16, - String, - String, - String, - Protocol, - Option, -); - -/// Ищет уже существующий conn_id для этой цели (host,port,username), чтобы -/// повторные подключения писали пароль в keychain под тем же id, а не заводили новый. -pub fn find_history_conn_id( - conn: &Connection, - host: &str, - port: u16, - username: &str, -) -> Result> { - let id: Option = conn - .query_row( - "SELECT conn_id FROM connection_history WHERE host = ?1 AND port = ?2 AND username = ?3", - params![host, port, username], - |row| row.get(0), - ) - .ok(); - Ok(id) -} - -#[allow(clippy::too_many_arguments)] -pub fn add_history_entry( - conn: &Connection, - host: &str, - port: u16, - username: &str, - conn_id: &str, - protocol: &Protocol, - key_path: Option<&str>, -) -> Result<()> { - use chrono::Local; - let connected_at = Local::now().format("%H:%M %d.%m").to_string(); - let proto = protocol_to_str(protocol); - conn.execute( - "INSERT INTO connection_history (host, port, username, connected_at, conn_id, protocol, key_path) - VALUES (?1,?2,?3,?4,?5,?6,?7) - ON CONFLICT(host, port, username) DO UPDATE SET - connected_at = excluded.connected_at, - protocol = excluded.protocol, - key_path = excluded.key_path, - conn_id = COALESCE(connection_history.conn_id, excluded.conn_id)", - params![host, port, username, connected_at, conn_id, proto, key_path], - )?; - // Храним только последние 20 записей - conn.execute( - "DELETE FROM connection_history WHERE id NOT IN ( - SELECT id FROM connection_history ORDER BY id DESC LIMIT 20 - )", - [], - )?; - Ok(()) -} - -/// Полностью очищает историю подключений (Settings → History → Clear All). -pub fn clear_history(conn: &Connection) -> Result<()> { - conn.execute("DELETE FROM connection_history", [])?; - Ok(()) -} - -pub fn get_history(conn: &Connection) -> Result> { - let mut stmt = conn.prepare( - "SELECT host, port, username, connected_at, conn_id, protocol, key_path - FROM connection_history ORDER BY id DESC LIMIT 20", - )?; - let rows = stmt - .query_map([], |row| { - let proto_str: String = row.get(5)?; - Ok(( - row.get(0)?, - row.get(1)?, - row.get(2)?, - row.get(3)?, - row.get(4)?, - protocol_from_str(&proto_str), - row.get(6)?, - )) - })? - .filter_map(|r| r.ok()) - .collect(); - Ok(rows) -} diff --git a/src/storage/db/settings.rs b/src/storage/db/settings.rs deleted file mode 100644 index 0dee418..0000000 --- a/src/storage/db/settings.rs +++ /dev/null @@ -1,43 +0,0 @@ -//! Таблица произвольных настроек приложения (ключ-значение) — используется -//! диалогом Settings для того, что должно переживать перезапуск (в отличие -//! от видимости панелей, которая остаётся только на время сессии). -use anyhow::Result; -use rusqlite::{params, Connection}; - -pub fn get_setting(conn: &Connection, key: &str) -> Option { - conn.query_row( - "SELECT value FROM app_settings WHERE key = ?1", - params![key], - |row| row.get(0), - ) - .ok() -} - -pub fn set_setting(conn: &Connection, key: &str, value: &str) -> Result<()> { - conn.execute( - "INSERT INTO app_settings (key, value) VALUES (?1, ?2) - ON CONFLICT(key) DO UPDATE SET value = excluded.value", - params![key, value], - )?; - Ok(()) -} - -pub fn get_bool(conn: &Connection, key: &str, default: bool) -> bool { - get_setting(conn, key) - .and_then(|v| v.parse().ok()) - .unwrap_or(default) -} - -pub fn get_u32(conn: &Connection, key: &str, default: u32) -> u32 { - get_setting(conn, key) - .and_then(|v| v.parse().ok()) - .unwrap_or(default) -} - -pub fn set_bool(conn: &Connection, key: &str, value: bool) { - let _ = set_setting(conn, key, &value.to_string()); -} - -pub fn set_u32(conn: &Connection, key: &str, value: u32) { - let _ = set_setting(conn, key, &value.to_string()); -} diff --git a/src/storage/db/sites.rs b/src/storage/db/sites.rs deleted file mode 100644 index 4afa140..0000000 --- a/src/storage/db/sites.rs +++ /dev/null @@ -1,59 +0,0 @@ -//! Таблица сохранённых серверов (Sites). -use anyhow::Result; -use rusqlite::{params, Connection}; - -use super::{protocol_from_str, protocol_to_str}; -use crate::domain::site::Site; - -pub fn get_sites(conn: &Connection) -> Result> { - let mut stmt = conn.prepare( - "SELECT id, name, protocol, host, port, username, password, key_path, folder, note FROM sites ORDER BY name", - )?; - let sites = stmt - .query_map([], |row| { - let proto_str: String = row.get(2)?; - let protocol = protocol_from_str(&proto_str); - Ok(Site { - id: row.get(0)?, - name: row.get(1)?, - protocol, - host: row.get(3)?, - port: row.get(4)?, - username: row.get(5)?, - password: row.get(6)?, - key_path: row.get(7)?, - folder: row.get(8)?, - note: row.get(9)?, - }) - })? - .filter_map(|r| r.ok()) - .collect(); - Ok(sites) -} - -pub fn save_site(conn: &Connection, site: &Site) -> Result<()> { - let proto = protocol_to_str(&site.protocol); - conn.execute( - "INSERT OR REPLACE INTO sites - (id, name, protocol, host, port, username, password, key_path, folder, note) - VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10)", - params![ - site.id, - site.name, - proto, - site.host, - site.port, - site.username, - site.password, - site.key_path, - site.folder, - site.note - ], - )?; - Ok(()) -} - -pub fn delete_site(conn: &Connection, id: &str) -> Result<()> { - conn.execute("DELETE FROM sites WHERE id = ?1", params![id])?; - Ok(()) -} diff --git a/src/storage/mod.rs b/src/storage/mod.rs deleted file mode 100644 index dec1023..0000000 --- a/src/storage/mod.rs +++ /dev/null @@ -1 +0,0 @@ -pub mod db; diff --git a/src/tests/commands_test.rs b/src/tests/commands_test.rs index 6ce666b..d5cd51d 100644 --- a/src/tests/commands_test.rs +++ b/src/tests/commands_test.rs @@ -1,14 +1,16 @@ //! Integration-level smoke tests: verify module imports and command function //! signatures compile correctly. Does not exercise Tauri runtime. -use wherry_lib::domain::connection::{ConnectionParams, Protocol}; -use wherry_lib::domain::error::AppError; -use wherry_lib::domain::file_entry::{EntryKind, FileEntry}; -use wherry_lib::domain::site::Site; -use wherry_lib::domain::transfer::{TaskState, TransferKind, TransferTask}; -use wherry_lib::storage::db::{ - add_bookmark, add_history_entry, clear_history, delete_site, find_history_conn_id, - get_bookmarks, get_bool, get_history, get_setting, get_sites, get_u32, init_tables, - remove_bookmark, save_site, set_bool, set_setting, set_u32, HistoryRow, +use wherry_lib::domain::{ + ConnectionParams, EntryKind, FileEntry, Protocol, Site, TaskState, TransferKind, TransferTask, +}; +use wherry_lib::error::AppError; +use wherry_lib::settings::{ + add_bookmark, get_bookmarks, get_bool, get_setting, get_u32, remove_bookmark, set_bool, + set_setting, set_u32, +}; +use wherry_lib::storage::{ + HistoryRow, add_history_entry, clear_history, delete_site, find_history_conn_id, get_history, + get_sites, init_tables, save_site, }; #[test] diff --git a/src/transfer/manager.rs b/src/transfer/manager.rs deleted file mode 100644 index e495ae0..0000000 --- a/src/transfer/manager.rs +++ /dev/null @@ -1,21 +0,0 @@ -use crate::fs::remote::RemoteRegistry; -use crate::transfer::queue::TransferQueue; -use std::sync::Arc; - -pub struct TransferManager { - pub queue: TransferQueue, - pub registry: Arc, -} - -impl TransferManager { - pub fn new(registry: Arc) -> Arc { - Arc::new(Self { - queue: TransferQueue::default(), - registry, - }) - } - - pub fn registry(&self) -> &RemoteRegistry { - &self.registry - } -} diff --git a/src/transfer/mod.rs b/src/transfer/mod.rs deleted file mode 100644 index 4b2be47..0000000 --- a/src/transfer/mod.rs +++ /dev/null @@ -1,5 +0,0 @@ -pub mod manager; -pub mod progress; -pub mod queue; -pub mod task; -pub mod worker; diff --git a/src/transfer/progress.rs b/src/transfer/progress.rs deleted file mode 100644 index 7b64f79..0000000 --- a/src/transfer/progress.rs +++ /dev/null @@ -1,70 +0,0 @@ -use std::time::{Duration, Instant}; - -/// Throttle прогресс-событий — максимум 10 в секунду. -/// Пользователь разницы не увидит, а тысяч emit'ов не будет. -pub struct ProgressThrottle { - last_emit: Instant, - interval: Duration, -} - -impl Default for ProgressThrottle { - fn default() -> Self { - Self { - last_emit: Instant::now() - Duration::from_secs(1), - interval: Duration::from_millis(100), // 10 fps - } - } -} - -impl ProgressThrottle { - pub fn should_emit(&mut self) -> bool { - let now = Instant::now(); - if now.duration_since(self.last_emit) >= self.interval { - self.last_emit = now; - true - } else { - false - } - } - - /// Всегда emit при завершении/ошибке - pub fn force(&mut self) -> bool { - self.last_emit = Instant::now() - self.interval; - true - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_first_call_emits() { - let mut throttle = ProgressThrottle::default(); - assert!(throttle.should_emit()); - } - - #[test] - fn test_too_soon_does_not_emit() { - let mut throttle = ProgressThrottle::default(); - throttle.should_emit(); - assert!(!throttle.should_emit()); - } - - #[test] - fn test_force_resets() { - let mut throttle = ProgressThrottle::default(); - throttle.should_emit(); - assert!(throttle.force()); - assert!(throttle.should_emit()); - } - - #[test] - fn test_interval_elapsed_emits() { - let mut throttle = ProgressThrottle { - last_emit: Instant::now() - Duration::from_millis(200), - interval: Duration::from_millis(100), - }; - assert!(throttle.should_emit()); - } -} diff --git a/src/transfer/task.rs b/src/transfer/task.rs deleted file mode 100644 index e5e9dfe..0000000 --- a/src/transfer/task.rs +++ /dev/null @@ -1,19 +0,0 @@ -use crate::domain::transfer::{TransferKind, TransferTask}; - -pub fn new_task( - kind: TransferKind, - connection_id: impl Into, - local_path: impl Into, - remote_path: impl Into, - file_name: impl Into, - total_bytes: u64, -) -> TransferTask { - TransferTask::new( - kind, - connection_id.into(), - local_path.into(), - remote_path.into(), - file_name.into(), - total_bytes, - ) -} diff --git a/src/transfers/mod.rs b/src/transfers/mod.rs new file mode 100644 index 0000000..6568a6a --- /dev/null +++ b/src/transfers/mod.rs @@ -0,0 +1,5 @@ +pub mod queue; +pub mod worker; + +pub use queue::TransferQueue; +pub use worker::spawn_worker; diff --git a/src/transfer/queue.rs b/src/transfers/queue.rs similarity index 89% rename from src/transfer/queue.rs rename to src/transfers/queue.rs index 4d139a2..7b86680 100644 --- a/src/transfer/queue.rs +++ b/src/transfers/queue.rs @@ -1,7 +1,25 @@ -use crate::domain::transfer::{TaskState, TransferTask}; +use crate::domain::{TaskState, TransferKind, TransferTask}; use std::collections::VecDeque; use std::sync::{Arc, Mutex}; +pub fn new_task( + kind: TransferKind, + connection_id: impl Into, + local_path: impl Into, + remote_path: impl Into, + file_name: impl Into, + total_bytes: u64, +) -> TransferTask { + TransferTask::new( + kind, + connection_id.into(), + local_path.into(), + remote_path.into(), + file_name.into(), + total_bytes, + ) +} + #[derive(Clone, Default)] pub struct TransferQueue { inner: Arc>>, @@ -55,7 +73,7 @@ impl TransferQueue { #[cfg(test)] mod tests { use super::*; - use crate::domain::transfer::{TaskState, TransferKind, TransferTask}; + use crate::domain::{TaskState, TransferKind, TransferTask}; fn make_task(id: &str, total: u64) -> TransferTask { TransferTask { diff --git a/src/transfer/worker.rs b/src/transfers/worker.rs similarity index 76% rename from src/transfer/worker.rs rename to src/transfers/worker.rs index 90a654d..99d6f7a 100644 --- a/src/transfer/worker.rs +++ b/src/transfers/worker.rs @@ -1,19 +1,67 @@ use std::collections::HashMap; use std::sync::Arc; -use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering}; use std::sync::Mutex; +use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering}; use std::time::{Duration, Instant}; use tauri::{AppHandle, Emitter}; use tokio::time::sleep; -use crate::domain::transfer::{TaskState, TransferKind, TransferTask}; +use crate::domain::{TaskState, TransferKind, TransferTask}; use crate::fs::remote::RemoteRegistry; use crate::protocols::ProgressAction; -use crate::transfer::progress::ProgressThrottle; -use crate::transfer::queue::TransferQueue; +use crate::transfers::queue::TransferQueue; const POLL_INTERVAL_MS: u64 = 200; +pub struct ProgressThrottle { + last_emit: Instant, + interval: Duration, +} + +impl Default for ProgressThrottle { + fn default() -> Self { + Self { + last_emit: Instant::now() - Duration::from_secs(1), + interval: Duration::from_millis(100), + } + } +} + +impl ProgressThrottle { + pub fn should_emit(&mut self) -> bool { + let now = Instant::now(); + if now.duration_since(self.last_emit) >= self.interval { + self.last_emit = now; + true + } else { + false + } + } + + pub fn force(&mut self) -> bool { + self.last_emit = Instant::now() - self.interval; + true + } +} + +pub struct TransferManager { + pub queue: TransferQueue, + pub registry: Arc, +} + +impl TransferManager { + pub fn new(registry: Arc) -> Arc { + Arc::new(Self { + queue: TransferQueue::default(), + registry, + }) + } + + pub fn registry(&self) -> &RemoteRegistry { + &self.registry + } +} + #[derive(Clone, serde::Serialize)] struct ProgressPayload { id: String, @@ -52,9 +100,6 @@ fn emit_state(app: &AppHandle, id: &str, state: TaskState) { ); } -/// Уменьшает счётчик активных передач при выходе из скоупа — в том числе при -/// панике внутри задачи (Drop отрабатывает и во время размотки стека), иначе -/// упавшая задача навсегда съедала бы один слот параллелизма. struct InFlightGuard(Arc); impl Drop for InFlightGuard { fn drop(&mut self) { @@ -62,10 +107,6 @@ impl Drop for InFlightGuard { } } -/// Запускает диспетчер передач: раз в тик подбирает из очереди задачи со -/// статусом `Queued` и запускает их параллельно, пока число активных передач -/// меньше `max_concurrent` — значение читается заново на каждой итерации, так -/// что Settings → Transfers меняет параллелизм на лету, без перезапуска. pub fn spawn_worker( queue: TransferQueue, registry: Arc, @@ -122,7 +163,6 @@ pub fn spawn_worker( }); } -/// Выполняет одну задачу передачи целиком и записывает итог в очередь. async fn run_transfer( task: TransferTask, queue: TransferQueue, @@ -143,8 +183,6 @@ async fn run_transfer( } }; - // Progress callback: троттлит обновления очереди до ~10 FPS, считает - // текущую скорость и реагирует на отмену/паузу, выставленные пользователем. let queue_for_progress = queue.clone(); let task_id_for_progress = task.id.clone(); let app_for_progress = app.clone(); @@ -199,7 +237,10 @@ async fn run_transfer( Ok(()) => { queue.update_state(&task.id, TaskState::Completed); queue.update_progress(&task.id, task.total_bytes, 0); - completed_at.lock().unwrap().insert(task.id.clone(), Instant::now()); + completed_at + .lock() + .unwrap() + .insert(task.id.clone(), Instant::now()); emit_progress(&app, &queue, &task.id); emit_state(&app, &task.id, TaskState::Completed); } @@ -224,3 +265,38 @@ async fn run_transfer( } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_first_call_emits() { + let mut throttle = ProgressThrottle::default(); + assert!(throttle.should_emit()); + } + + #[test] + fn test_too_soon_does_not_emit() { + let mut throttle = ProgressThrottle::default(); + throttle.should_emit(); + assert!(!throttle.should_emit()); + } + + #[test] + fn test_force_resets() { + let mut throttle = ProgressThrottle::default(); + throttle.should_emit(); + assert!(throttle.force()); + assert!(throttle.should_emit()); + } + + #[test] + fn test_interval_elapsed_emits() { + let mut throttle = ProgressThrottle { + last_emit: Instant::now() - Duration::from_millis(200), + interval: Duration::from_millis(100), + }; + assert!(throttle.should_emit()); + } +} diff --git a/src/domain/window_state.rs b/src/window.rs similarity index 100% rename from src/domain/window_state.rs rename to src/window.rs