From 1ffc0417ae656ea3c82a03b0a540072fcf999356 Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Sat, 19 Sep 2026 11:16:48 +0200 Subject: [PATCH] feat(desktop): port delayed sync for growing files from upstream PR#39 Upload tasks optionally wait until a file's size/mtime signature is stable for a configured delay (off/10s/30s/60s/120s/300s), so large or slowly-written files aren't uploaded mid-write. Experimental setting, default off. Ported from https://github.com/cloudreve/desktop/pull/39 with hardening on top: - bound total stability wait at max(10min, 4x delay) so a perpetually growing file cannot starve the 2-worker upload queue - added the missing experimental/delaySync locale keys for all 9 languages the upstream PR left out - added unit tests for file-signature capture (file/dir/missing) Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- desktop/crates/cloudreve-sync/src/config.rs | 18 ++ .../crates/cloudreve-sync/src/tasks/queue.rs | 221 +++++++++++++++++- desktop/src-tauri/src/commands.rs | 10 + desktop/src-tauri/src/lib.rs | 1 + desktop/ui/public/locales/de/common.json | 3 + desktop/ui/public/locales/en-US/common.json | 3 + desktop/ui/public/locales/es/common.json | 3 + desktop/ui/public/locales/fr/common.json | 3 + desktop/ui/public/locales/it/common.json | 3 + desktop/ui/public/locales/ja/common.json | 3 + desktop/ui/public/locales/ko/common.json | 3 + desktop/ui/public/locales/pl/common.json | 3 + desktop/ui/public/locales/ru/common.json | 3 + desktop/ui/public/locales/zh-CN/common.json | 3 + desktop/ui/public/locales/zh-TW/common.json | 3 + .../ui/src/pages/settings/GeneralSection.tsx | 44 +++- 16 files changed, 324 insertions(+), 3 deletions(-) diff --git a/desktop/crates/cloudreve-sync/src/config.rs b/desktop/crates/cloudreve-sync/src/config.rs index bdb426ba..a9813804 100644 --- a/desktop/crates/cloudreve-sync/src/config.rs +++ b/desktop/crates/cloudreve-sync/src/config.rs @@ -60,6 +60,8 @@ pub struct AppConfig { pub log_level: LogLevel, /// Maximum number of log files to keep pub log_max_files: usize, + /// Delay before starting upload after file stops changing (seconds). 0 disables the delay. + pub sync_delay_seconds: u64, /// Language/locale setting (e.g., "en-US", "zh-CN"). None means use system default. pub language: Option, } @@ -74,6 +76,7 @@ impl Default for AppConfig { log_to_file: true, log_level: LogLevel::Debug, log_max_files: 5, + sync_delay_seconds: 0, language: None, } } @@ -281,6 +284,21 @@ impl ConfigManager { }) } + /// Get the sync delay in seconds + pub fn sync_delay_seconds(&self) -> u64 { + self.config + .read() + .map(|c| c.sync_delay_seconds) + .unwrap_or(0) + } + + /// Set the sync delay in seconds + pub fn set_sync_delay_seconds(&self, seconds: u64) -> Result<()> { + self.update(|config| { + config.sync_delay_seconds = seconds; + }) + } + /// Get the language setting pub fn language(&self) -> Option { self.config.read().ok().and_then(|c| c.language.clone()) diff --git a/desktop/crates/cloudreve-sync/src/tasks/queue.rs b/desktop/crates/cloudreve-sync/src/tasks/queue.rs index 05f53df5..745416a4 100644 --- a/desktop/crates/cloudreve-sync/src/tasks/queue.rs +++ b/desktop/crates/cloudreve-sync/src/tasks/queue.rs @@ -1,3 +1,4 @@ +use crate::config::ConfigManager; use crate::inventory::{InventoryDb, NewTaskRecord, TaskRecord, TaskStatus, TaskUpdate}; use crate::tasks::download::DownloadTask; use crate::tasks::types::{TaskKind, TaskPayload, TaskProgress}; @@ -6,14 +7,17 @@ use anyhow::{Context, Result, anyhow}; use cloudreve_api::Client; use dashmap::DashMap; use serde_json::Value; -use std::path::PathBuf; +use std::io::ErrorKind; +use std::path::{Path, PathBuf}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::time::{Duration, SystemTime}; use tokio::sync::{ Mutex, Notify, Semaphore, mpsc::{self, UnboundedReceiver, UnboundedSender}, }; use tokio::task::JoinHandle; +use tokio::time::{Instant, sleep}; use tracing::{debug, error, info, warn}; use uuid::Uuid; @@ -48,6 +52,12 @@ pub struct TaskQueue { task_paths: DashMap, } +#[derive(Debug, Clone, PartialEq, Eq)] +struct FileChangeSignature { + size: u64, + modified: SystemTime, +} + impl TaskQueue { pub async fn new( drive_id: impl Into, @@ -501,6 +511,11 @@ impl TaskQueue { match &task.payload.kind { TaskKind::Upload => { + match self.wait_for_upload_stability(task).await? { + TaskRunState::Cancelled => return Ok(TaskRunState::Cancelled), + TaskRunState::Completed => {} + } + let mut task_executor = UploadTask::new( self.inventory.clone(), self.cr_client.clone(), @@ -567,6 +582,174 @@ impl TaskQueue { Ok(TaskRunState::Completed) } + fn is_task_marked_inactive(&self, task_id: &str) -> bool { + match self.inventory.get_task_status(task_id) { + Ok(Some(TaskStatus::Cancelled)) => true, + Ok(Some(status)) => !status.is_active(), + Err(err) => { + warn!( + target: "tasks::queue", + drive = %self.drive_id, + task_id = %task_id, + error = %err, + "Failed to check task status while waiting for delayed upload" + ); + false + } + _ => false, + } + } + + fn capture_file_signature(path: &Path) -> Result> { + let metadata = match std::fs::metadata(path) { + Ok(metadata) => metadata, + Err(err) if err.kind() == ErrorKind::NotFound => return Ok(None), + Err(err) => { + return Err(err) + .with_context(|| format!("failed to read metadata for {}", path.display())); + } + }; + if !metadata.is_file() { + return Ok(None); + } + + let modified = metadata + .modified() + .with_context(|| format!("failed to read modified time for {}", path.display()))?; + + Ok(Some(FileChangeSignature { + size: metadata.len(), + modified, + })) + } + + async fn wait_for_upload_stability(&self, task: &QueuedTask) -> Result { + let delay_seconds = ConfigManager::try_get() + .map(|config| config.sync_delay_seconds()) + .unwrap_or(0); + + if delay_seconds == 0 { + return Ok(TaskRunState::Completed); + } + + let local_path = &task.payload.local_path; + let mut signature = match Self::capture_file_signature(local_path.as_path()) { + Ok(Some(signature)) => signature, + Ok(None) => return Ok(TaskRunState::Completed), + Err(err) => { + warn!( + target: "tasks::queue", + drive = %self.drive_id, + task_id = %task.task_id, + path = %local_path.display(), + error = %err, + "Failed to capture initial file signature, skipping delayed upload wait" + ); + return Ok(TaskRunState::Completed); + } + }; + + info!( + target: "tasks::queue", + drive = %self.drive_id, + task_id = %task.task_id, + path = %local_path.display(), + delay_seconds = delay_seconds, + "Waiting for file to stop changing before upload" + ); + + let required_stable = Duration::from_secs(delay_seconds); + // A continuously-growing file must not block a worker slot forever: + // give up waiting after 4x the configured delay (min 10 minutes) and + // upload whatever is on disk, rather than starving the queue. + let max_total_wait = required_stable + .saturating_mul(4) + .max(Duration::from_secs(600)); + let wait_started = Instant::now(); + let poll_interval = if delay_seconds <= 30 { + Duration::from_secs(1) + } else if delay_seconds <= 120 { + Duration::from_secs(2) + } else { + Duration::from_secs(3) + }; + let status_check_interval = Duration::from_secs(5); + let mut unchanged_since = Instant::now(); + let mut next_status_check = Instant::now(); + + while unchanged_since.elapsed() < required_stable { + if wait_started.elapsed() >= max_total_wait { + warn!( + target: "tasks::queue", + drive = %self.drive_id, + task_id = %task.task_id, + path = %local_path.display(), + "File still changing after maximum wait, proceeding with upload" + ); + return Ok(TaskRunState::Completed); + } + + if self.cancel_requested.load(Ordering::SeqCst) { + debug!( + target: "tasks::queue", + drive = %self.drive_id, + task_id = %task.task_id, + path = %local_path.display(), + "Task cancelled while waiting for upload delay" + ); + return Ok(TaskRunState::Cancelled); + } + + if Instant::now() >= next_status_check { + if self.is_task_marked_inactive(&task.task_id) { + debug!( + target: "tasks::queue", + drive = %self.drive_id, + task_id = %task.task_id, + path = %local_path.display(), + "Task became inactive while waiting for upload delay" + ); + return Ok(TaskRunState::Cancelled); + } + next_status_check = Instant::now() + status_check_interval; + } + + sleep(poll_interval).await; + + let current_signature = match Self::capture_file_signature(local_path.as_path()) { + Ok(Some(signature)) => signature, + Ok(None) => return Ok(TaskRunState::Completed), + Err(err) => { + warn!( + target: "tasks::queue", + drive = %self.drive_id, + task_id = %task.task_id, + path = %local_path.display(), + error = %err, + "Failed to capture file signature while waiting, continuing upload" + ); + return Ok(TaskRunState::Completed); + } + }; + + if current_signature != signature { + signature = current_signature; + unchanged_since = Instant::now(); + } + } + + debug!( + target: "tasks::queue", + drive = %self.drive_id, + task_id = %task.task_id, + path = %local_path.display(), + delay_seconds = delay_seconds, + "File remained unchanged for configured delay, proceeding with upload" + ); + + Ok(TaskRunState::Completed) + } + #[allow(dead_code)] async fn wait_for_idle(&self) { while self.inflight.load(Ordering::SeqCst) > 0 { @@ -716,3 +899,39 @@ pub struct QueuedTask { pub task_id: String, pub payload: TaskPayload, } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn file_signature_captures_size_and_mtime() { + let dir = tempfile::tempdir().unwrap(); + let file = dir.path().join("data.bin"); + std::fs::write(&file, b"hello").unwrap(); + + let signature = TaskQueue::capture_file_signature(&file) + .unwrap() + .expect("file should produce a signature"); + assert_eq!(signature.size, 5); + + std::fs::write(&file, b"hello world").unwrap(); + let updated = TaskQueue::capture_file_signature(&file) + .unwrap() + .unwrap(); + assert_ne!(signature, updated); + } + + #[test] + fn file_signature_none_for_dir_and_missing() { + let dir = tempfile::tempdir().unwrap(); + assert!(TaskQueue::capture_file_signature(dir.path()) + .unwrap() + .is_none()); + assert!(TaskQueue::capture_file_signature( + &dir.path().join("missing.bin") + ) + .unwrap() + .is_none()); + } +} diff --git a/desktop/src-tauri/src/commands.rs b/desktop/src-tauri/src/commands.rs index c03abbe6..ecfd3d5a 100644 --- a/desktop/src-tauri/src/commands.rs +++ b/desktop/src-tauri/src/commands.rs @@ -1187,6 +1187,7 @@ pub async fn get_general_settings() -> CommandResult { log_to_file: config.log_to_file, log_level: config.log_level.as_str().to_string(), log_max_files: config.log_max_files, + sync_delay_seconds: config.sync_delay_seconds, log_dir: ConfigManager::get_log_dir().display().to_string(), language: config.language, }) @@ -1200,6 +1201,7 @@ pub struct GeneralSettings { pub log_to_file: bool, pub log_level: String, pub log_max_files: usize, + pub sync_delay_seconds: u64, pub log_dir: String, pub language: Option, } @@ -1231,6 +1233,14 @@ pub async fn set_log_max_files(max_files: usize) -> CommandResult<()> { .map_err(|e| e.to_string()) } +/// Set sync delay in seconds +#[tauri::command] +pub async fn set_sync_delay_seconds(seconds: u64) -> CommandResult<()> { + ConfigManager::get() + .set_sync_delay_seconds(seconds) + .map_err(|e| e.to_string()) +} + /// Set language setting and update rust_i18n locale #[tauri::command] pub async fn set_language(app: AppHandle, language: Option) -> CommandResult<()> { diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index e86d77fb..2020205f 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -474,6 +474,7 @@ pub fn run() { commands::set_log_to_file, commands::set_log_level, commands::set_log_max_files, + commands::set_sync_delay_seconds, commands::set_language, commands::open_log_folder, ]) diff --git a/desktop/ui/public/locales/de/common.json b/desktop/ui/public/locales/de/common.json index b01f81c1..6ec830f2 100644 --- a/desktop/ui/public/locales/de/common.json +++ b/desktop/ui/public/locales/de/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "Detailgrad der Protokollausgabe festlegen (Neustart erforderlich)", "logMaxFiles": "Maximale Protokolldateien", "logMaxFilesDescription": "Anzahl der aufzubewahrenden Protokolldateien (Neustart erforderlich)", + "experimentalFeatures": "Experimentelle Funktionen", + "delaySync": "Synchronisierung verzögern", + "delaySyncDescription": "Die Synchronisierungsaufgabe wartet, bis die Datei nicht mehr geändert wird", "storage": "Speicher", "openFolder": "Ordner öffnen", "openSite": "Website öffnen", diff --git a/desktop/ui/public/locales/en-US/common.json b/desktop/ui/public/locales/en-US/common.json index 9c6a7b65..296a2ea9 100644 --- a/desktop/ui/public/locales/en-US/common.json +++ b/desktop/ui/public/locales/en-US/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "Set the verbosity of log output (restart required)", "logMaxFiles": "Max log files", "logMaxFilesDescription": "Number of log files to keep (restart required)", + "experimentalFeatures": "Experimental features", + "delaySync": "Delay sync", + "delaySyncDescription": "Sync task wait until file stop changing", "storage": "Storage", "openFolder": "Open folder", "openSite": "Open site", diff --git a/desktop/ui/public/locales/es/common.json b/desktop/ui/public/locales/es/common.json index 459baf94..0081f32e 100644 --- a/desktop/ui/public/locales/es/common.json +++ b/desktop/ui/public/locales/es/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "Establecer el nivel de detalle de la salida del registro (requiere reinicio)", "logMaxFiles": "Máximo de archivos de registro", "logMaxFilesDescription": "Número de archivos de registro a conservar (requiere reinicio)", + "experimentalFeatures": "Funciones experimentales", + "delaySync": "Retrasar sincronización", + "delaySyncDescription": "La tarea de sincronización espera hasta que el archivo deje de cambiar", "storage": "Almacenamiento", "openFolder": "Abrir carpeta", "openSite": "Abrir sitio", diff --git a/desktop/ui/public/locales/fr/common.json b/desktop/ui/public/locales/fr/common.json index e778c7b2..16b8310c 100644 --- a/desktop/ui/public/locales/fr/common.json +++ b/desktop/ui/public/locales/fr/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "Définir le niveau de détail de la sortie des journaux (redémarrage requis)", "logMaxFiles": "Nombre maximum de fichiers journaux", "logMaxFilesDescription": "Nombre de fichiers journaux à conserver (redémarrage requis)", + "experimentalFeatures": "Fonctionnalités expérimentales", + "delaySync": "Retarder la synchronisation", + "delaySyncDescription": "La tâche de synchronisation attend que le fichier cesse d'être modifié", "storage": "Stockage", "openFolder": "Ouvrir le dossier", "openSite": "Ouvrir le site", diff --git a/desktop/ui/public/locales/it/common.json b/desktop/ui/public/locales/it/common.json index 0b3ab779..3bb814a5 100644 --- a/desktop/ui/public/locales/it/common.json +++ b/desktop/ui/public/locales/it/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "Imposta il livello di dettaglio dell'output dei log (riavvio necessario)", "logMaxFiles": "Numero massimo file di log", "logMaxFilesDescription": "Numero di file di log da conservare (riavvio necessario)", + "experimentalFeatures": "Funzionalità sperimentali", + "delaySync": "Ritarda sincronizzazione", + "delaySyncDescription": "L'attività di sincronizzazione attende che il file smetta di cambiare", "storage": "Archiviazione", "openFolder": "Apri cartella", "openSite": "Apri sito", diff --git a/desktop/ui/public/locales/ja/common.json b/desktop/ui/public/locales/ja/common.json index 43a1603b..25bcd436 100644 --- a/desktop/ui/public/locales/ja/common.json +++ b/desktop/ui/public/locales/ja/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "ログ出力の詳細度を設定(再起動が必要)", "logMaxFiles": "最大ログファイル数", "logMaxFilesDescription": "保持するログファイルの数(再起動が必要)", + "experimentalFeatures": "実験的な機能", + "delaySync": "同期を遅延", + "delaySyncDescription": "同期タスクはファイルの変更が止まるまで待機します", "storage": "ストレージ", "openFolder": "フォルダを開く", "openSite": "サイトを開く", diff --git a/desktop/ui/public/locales/ko/common.json b/desktop/ui/public/locales/ko/common.json index e91d818c..0a1ef64c 100644 --- a/desktop/ui/public/locales/ko/common.json +++ b/desktop/ui/public/locales/ko/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "로그 출력의 세부 수준 설정 (재시작 필요)", "logMaxFiles": "최대 로그 파일 수", "logMaxFilesDescription": "보관할 로그 파일 수 (재시작 필요)", + "experimentalFeatures": "실험적 기능", + "delaySync": "동기화 지연", + "delaySyncDescription": "동기화 작업은 파일 변경이 멈출 때까지 기다립니다", "storage": "저장소", "openFolder": "폴더 열기", "openSite": "사이트 열기", diff --git a/desktop/ui/public/locales/pl/common.json b/desktop/ui/public/locales/pl/common.json index 6ded0c5c..df2bb978 100644 --- a/desktop/ui/public/locales/pl/common.json +++ b/desktop/ui/public/locales/pl/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "Ustaw szczegółowość dziennika (wymaga restartu)", "logMaxFiles": "Maksymalna liczba plików dziennika", "logMaxFilesDescription": "Liczba plików dziennika do przechowywania (wymaga restartu)", + "experimentalFeatures": "Funkcje eksperymentalne", + "delaySync": "Opóźnij synchronizację", + "delaySyncDescription": "Zadanie synchronizacji czeka, aż plik przestanie się zmieniać", "storage": "Pamięć", "openFolder": "Otwórz folder", "openSite": "Otwórz witrynę", diff --git a/desktop/ui/public/locales/ru/common.json b/desktop/ui/public/locales/ru/common.json index 4e7bee22..c965ea1e 100644 --- a/desktop/ui/public/locales/ru/common.json +++ b/desktop/ui/public/locales/ru/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "Установить уровень детализации журнала (требуется перезапуск)", "logMaxFiles": "Максимум файлов журнала", "logMaxFilesDescription": "Количество сохраняемых файлов журнала (требуется перезапуск)", + "experimentalFeatures": "Экспериментальные функции", + "delaySync": "Задержка синхронизации", + "delaySyncDescription": "Задача синхронизации ждёт, пока файл перестанет изменяться", "storage": "Хранилище", "openFolder": "Открыть папку", "openSite": "Открыть сайт", diff --git a/desktop/ui/public/locales/zh-CN/common.json b/desktop/ui/public/locales/zh-CN/common.json index 9f06eaca..92e43e8d 100644 --- a/desktop/ui/public/locales/zh-CN/common.json +++ b/desktop/ui/public/locales/zh-CN/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "设置日志输出的详细程度(需要重启)", "logMaxFiles": "最大日志文件数", "logMaxFilesDescription": "保留的日志文件数量(需要重启)", + "experimentalFeatures": "实验性功能", + "delaySync": "延迟同步", + "delaySyncDescription": "同步任务等待文件停止更改后再上传", "storage": "存储空间", "openFolder": "打开文件夹", "openSite": "打开站点", diff --git a/desktop/ui/public/locales/zh-TW/common.json b/desktop/ui/public/locales/zh-TW/common.json index 9c2bd52b..331e95f3 100644 --- a/desktop/ui/public/locales/zh-TW/common.json +++ b/desktop/ui/public/locales/zh-TW/common.json @@ -88,6 +88,9 @@ "logLevelDescription": "設定日誌輸出的詳細程度(需要重新啟動)", "logMaxFiles": "最大日誌檔案數", "logMaxFilesDescription": "保留的日誌檔案數量(需要重新啟動)", + "experimentalFeatures": "實驗性功能", + "delaySync": "延遲同步", + "delaySyncDescription": "同步任務會等待檔案停止變更", "storage": "儲存空間", "openFolder": "開啟資料夾", "openSite": "開啟網站", diff --git a/desktop/ui/src/pages/settings/GeneralSection.tsx b/desktop/ui/src/pages/settings/GeneralSection.tsx index 51fdfbfa..d2d76671 100644 --- a/desktop/ui/src/pages/settings/GeneralSection.tsx +++ b/desktop/ui/src/pages/settings/GeneralSection.tsx @@ -171,15 +171,16 @@ function SettingActionItem({ interface SettingsGroupProps { title: string; + titleColor?: string; children: React.ReactNode; } -function SettingsGroup({ title, children }: SettingsGroupProps) { +function SettingsGroup({ title, titleColor = "text.secondary", children }: SettingsGroupProps) { return ( {title} @@ -209,6 +210,7 @@ interface GeneralSettings { log_to_file: boolean; log_level: string; log_max_files: number; + sync_delay_seconds: number; log_dir: string; language: string | null; } @@ -228,6 +230,15 @@ const MAX_FILES_OPTIONS = [ { value: "10", label: "10" }, ]; +const DELAY_SYNC_OPTIONS = [ + { value: "0", label: "off" }, + { value: "10", label: "10s" }, + { value: "30", label: "30s" }, + { value: "60", label: "1min" }, + { value: "120", label: "2min" }, + { value: "300", label: "5min" }, +]; + export default function GeneralSection() { const { t, i18n } = useTranslation(); const [autoStart, setAutoStart] = useState(true); @@ -237,6 +248,7 @@ export default function GeneralSection() { const [logToFile, setLogToFile] = useState(true); const [logLevel, setLogLevel] = useState("info"); const [logMaxFiles, setLogMaxFiles] = useState(5); + const [syncDelaySeconds, setSyncDelaySeconds] = useState(0); const [logDir, setLogDir] = useState(""); const [language, setLanguage] = useState(null); const [loading, setLoading] = useState(true); @@ -255,6 +267,7 @@ export default function GeneralSection() { setLogToFile(settings.log_to_file); setLogLevel(settings.log_level); setLogMaxFiles(settings.log_max_files); + setSyncDelaySeconds(settings.sync_delay_seconds); setLogDir(settings.log_dir); setLanguage(settings.language); } catch (error) { @@ -347,6 +360,18 @@ export default function GeneralSection() { } }; + const handleSyncDelayChange = async (value: string) => { + const seconds = parseInt(value, 10); + const previousValue = syncDelaySeconds; + setSyncDelaySeconds(seconds); + try { + await invoke("set_sync_delay_seconds", { seconds }); + } catch (error) { + console.error("Failed to change sync delay:", error); + setSyncDelaySeconds(previousValue); + } + }; + const handleOpenLogFolder = async () => { try { await invoke("open_log_folder"); @@ -469,6 +494,21 @@ export default function GeneralSection() { isLast={true} /> + + + + ); }