From f6faf17a98d7c559ea29a14c040706d67d8fab15 Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Sat, 19 Sep 2026 11:27:00 +0200 Subject: [PATCH] fix(desktop): folder-move descendant events + chunk_concurrency i64 (#151, #169) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit #151: when a directory is moved/renamed, the OS emits Remove+Create events for every descendant in addition to the directory's own events. The event blocker only covered the directory's exact paths, so child removes were propagated as remote deletes — children landed in the recycle bin while the folder moved without them. - EventBlocker: TTL-scoped path-prefix suppression (register_prefix); rename() suppresses Remove(source/*) + Create(target/*) for 15s around the remote call, covering the in-flight window - process_fs_events: deterministic group order (name->create->modify-> remove) instead of HashMap order so same-batch removes observe the post-move world - tests: prefix descendant/self/sibling matching, expiry, exact-consume #169: chunk_concurrency parsed as i32 overflowed on values above i32::MAX, failing serde for the whole listing response and rendering the sync root unusable. Widened to i64; consumer already clamps max(1). Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../cloudreve-api/src/models/explorer.rs | 2 +- .../cloudreve-sync/src/drive/commands.rs | 32 ++++- .../cloudreve-sync/src/drive/event_blocker.rs | 120 +++++++++++++++++- 3 files changed, 151 insertions(+), 3 deletions(-) diff --git a/desktop/crates/cloudreve-api/src/models/explorer.rs b/desktop/crates/cloudreve-api/src/models/explorer.rs index 12594d54..3b806261 100644 --- a/desktop/crates/cloudreve-api/src/models/explorer.rs +++ b/desktop/crates/cloudreve-api/src/models/explorer.rs @@ -193,7 +193,7 @@ pub struct StoragePolicy { #[serde(skip_serializing_if = "Option::is_none")] pub children: Option>, #[serde(skip_serializing_if = "Option::is_none")] - pub chunk_concurrency: Option, + pub chunk_concurrency: Option, #[serde(skip_serializing_if = "Option::is_none")] pub encryption: Option, #[serde(skip_serializing_if = "Option::is_none")] diff --git a/desktop/crates/cloudreve-sync/src/drive/commands.rs b/desktop/crates/cloudreve-sync/src/drive/commands.rs index 93fb5600..7951bb4f 100644 --- a/desktop/crates/cloudreve-sync/src/drive/commands.rs +++ b/desktop/crates/cloudreve-sync/src/drive/commands.rs @@ -37,6 +37,7 @@ use std::{ collections::HashMap, ops::Range, path::{Path, PathBuf}, + time::Duration, }; use tokio::sync::oneshot::Sender; use uuid::Uuid; @@ -522,6 +523,19 @@ impl Mount { return Ok(()); } + // If the source is a directory, the OS emits Remove/Create events for + // every descendant in addition to the directory's own events. Suppress + // the subtree for a short window — registered before the remote call so + // events racing the in-flight move are covered as well — so child + // events are not propagated as remote deletes/creates. + if source.is_dir() { + let ttl = Duration::from_secs(15); + self.event_blocker + .register_prefix(&EventKind::Remove(RemoveKind::Any), source.clone(), ttl); + self.event_blocker + .register_prefix(&EventKind::Create(CreateKind::Any), target.clone(), ttl); + } + // if target and src under the same dir, trigger rename call let target_parent = target.parent().context("root cannot be moved")?; let source_parent = source.parent().context("root cannot be moved")?; @@ -539,6 +553,7 @@ impl Mount { .await { Ok(_) => { + // Block the modify name events for rename (From for source, To for target) self.event_blocker.register_once( &EventKind::Modify(ModifyKind::Name(RenameMode::From)), @@ -574,6 +589,7 @@ impl Mount { .await { Ok(_) => { + // Block remove event for source and create event for target self.event_blocker .register_once(&EventKind::Remove(RemoveKind::Any), source.clone()); @@ -589,7 +605,21 @@ impl Mount { } pub async fn process_fs_events(&self, events: GroupedFsEvents) -> Result<()> { - for (event_kind, events) in events { + // Process groups in a deterministic order: renames/moves first, then + // creates and modifications, removes last. A directory move can emit + // per-descendant remove events in the same batch; committing the move + // remotely before handling removes prevents child URIs from being + // deleted on the server. + let mut groups: Vec<(EventKind, Vec)> = events.into_iter().collect(); + groups.sort_by_key(|(kind, _)| match kind { + EventKind::Modify(ModifyKind::Name(_)) => 0, + EventKind::Create(_) => 1, + EventKind::Modify(_) => 2, + EventKind::Remove(_) => 3, + _ => 4, + }); + + for (event_kind, events) in groups { // Filter out events that were pre-registered by rename operations let filtered_events = self.event_blocker.filter_events(events, &event_kind); diff --git a/desktop/crates/cloudreve-sync/src/drive/event_blocker.rs b/desktop/crates/cloudreve-sync/src/drive/event_blocker.rs index 4ed0d7ae..6d41125b 100644 --- a/desktop/crates/cloudreve-sync/src/drive/event_blocker.rs +++ b/desktop/crates/cloudreve-sync/src/drive/event_blocker.rs @@ -4,6 +4,7 @@ use std::{ collections::HashMap, path::PathBuf, sync::{Arc, Mutex}, + time::{Duration, Instant}, }; /// A key for identifying blocked events, consisting of a normalized EventKind and a path. @@ -52,11 +53,24 @@ impl From<&EventKind> for NormalizedEventKind { /// When a rename operation is processed, it may trigger additional filesystem events /// (like Remove for the source and Create for the target). These events should be /// blocked to avoid duplicate processing. +/// +/// A path-prefix suppression entry: blocks every event of `kind` whose path +/// lies under `prefix` until `expires_at`. Used to swallow the per-descendant +/// Remove/Create events the OS emits when a directory is moved or renamed. +#[derive(Debug)] +struct PrefixBlock { + kind: NormalizedEventKind, + prefix: PathBuf, + expires_at: Instant, +} + #[derive(Debug, Clone, Default)] pub struct EventBlocker { /// Map of blocked event keys to their remaining block count. /// When count reaches 0, the entry is removed. blocked: Arc>>, + /// Path-prefix suppressions with a fixed expiry. + prefix_blocked: Arc>>, } impl EventBlocker { @@ -64,6 +78,7 @@ impl EventBlocker { pub fn new() -> Self { Self { blocked: Arc::new(Mutex::new(HashMap::new())), + prefix_blocked: Arc::new(Mutex::new(Vec::new())), } } @@ -91,6 +106,35 @@ impl EventBlocker { self.register(kind, path, 1); } + /// Registers a path-prefix suppression for `ttl`. + /// + /// While the entry is alive, every event of `kind` whose path is `prefix` + /// itself or a descendant of it is blocked. Unlike [register], entries are + /// not consumed per event — they live until expiry, since the number of + /// descendant events the OS emits for a moved directory is not knowable. + /// + /// Keep the TTL tight: a genuine operation inside the tree during the + /// window is swallowed as well. + pub fn register_prefix(&self, kind: &EventKind, prefix: PathBuf, ttl: Duration) { + let mut blocked = self.prefix_blocked.lock().unwrap(); + blocked.retain(|b| b.expires_at > Instant::now()); + blocked.push(PrefixBlock { + kind: NormalizedEventKind::from(kind), + prefix, + expires_at: Instant::now() + ttl, + }); + } + + /// Returns true if `path` is covered by a live prefix suppression of `kind`. + fn is_prefix_blocked(&self, kind: &EventKind, path: &PathBuf) -> bool { + let normalized = NormalizedEventKind::from(kind); + let mut blocked = self.prefix_blocked.lock().unwrap(); + blocked.retain(|b| b.expires_at > Instant::now()); + blocked + .iter() + .any(|b| b.kind == normalized && path.starts_with(&b.prefix)) + } + /// Checks if an event should be blocked and decrements the counter if so. /// /// # Arguments @@ -100,6 +144,16 @@ impl EventBlocker { /// # Returns /// `true` if the event should be blocked (was pre-registered), `false` otherwise pub fn should_block(&self, kind: &EventKind, path: &PathBuf) -> bool { + if self.is_prefix_blocked(kind, path) { + tracing::debug!( + target: "drive::event_blocker", + kind = ?kind, + path = %path.display(), + "Blocked event under suppressed prefix" + ); + return true; + } + let key = BlockKey { kind: NormalizedEventKind::from(kind), path: path.clone(), @@ -156,12 +210,15 @@ impl EventBlocker { pub fn clear(&self) { let mut blocked = self.blocked.lock().unwrap(); blocked.clear(); + let mut prefix_blocked = self.prefix_blocked.lock().unwrap(); + prefix_blocked.clear(); } /// Returns the number of currently registered event blocks. pub fn len(&self) -> usize { let blocked = self.blocked.lock().unwrap(); - blocked.len() + let prefix_blocked = self.prefix_blocked.lock().unwrap(); + blocked.len() + prefix_blocked.len() } /// Returns true if there are no registered event blocks. @@ -169,3 +226,64 @@ impl EventBlocker { self.len() == 0 } } + +#[cfg(test)] +mod tests { + use super::*; + use notify_debouncer_full::notify::event::RemoveKind; + + #[test] + fn prefix_blocks_descendant_and_self() { + let blocker = EventBlocker::new(); + let dir = PathBuf::from("/sync/dir"); + blocker.register_prefix( + &EventKind::Remove(RemoveKind::Any), + dir.clone(), + Duration::from_secs(60), + ); + + assert!(blocker.should_block( + &EventKind::Remove(RemoveKind::Any), + &dir.join("child/file.txt") + )); + assert!(blocker.should_block(&EventKind::Remove(RemoveKind::Any), &dir)); + // Different kind and different subtree are not blocked + assert!(!blocker.should_block( + &EventKind::Create(notify_debouncer_full::notify::event::CreateKind::Any), + &dir.join("child/file.txt") + )); + assert!(!blocker.should_block( + &EventKind::Remove(RemoveKind::Any), + &PathBuf::from("/sync/other/file.txt") + )); + // Sibling sharing a name prefix must not match ("dir2" vs "dir") + assert!(!blocker.should_block( + &EventKind::Remove(RemoveKind::Any), + &PathBuf::from("/sync/dir2/file.txt") + )); + } + + #[test] + fn prefix_expires() { + let blocker = EventBlocker::new(); + let dir = PathBuf::from("/sync/dir"); + blocker.register_prefix( + &EventKind::Remove(RemoveKind::Any), + dir.clone(), + Duration::from_millis(0), + ); + assert!(!blocker.should_block( + &EventKind::Remove(RemoveKind::Any), + &dir.join("child") + )); + } + + #[test] + fn exact_block_consumed_once() { + let blocker = EventBlocker::new(); + let path = PathBuf::from("/sync/file.txt"); + blocker.register_once(&EventKind::Remove(RemoveKind::Any), path.clone()); + assert!(blocker.should_block(&EventKind::Remove(RemoveKind::Any), &path)); + assert!(!blocker.should_block(&EventKind::Remove(RemoveKind::Any), &path)); + } +}