@ -4,6 +4,7 @@ use std::{
collections ::HashMap ,
collections ::HashMap ,
path ::PathBuf ,
path ::PathBuf ,
sync ::{ Arc , Mutex } ,
sync ::{ Arc , Mutex } ,
time ::{ Duration , Instant } ,
} ;
} ;
/// A key for identifying blocked events, consisting of a normalized EventKind and a path.
/// 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
/// 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
/// (like Remove for the source and Create for the target). These events should be
/// blocked to avoid duplicate processing.
/// 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) ]
#[ derive(Debug, Clone, Default) ]
pub struct EventBlocker {
pub struct EventBlocker {
/// Map of blocked event keys to their remaining block count.
/// Map of blocked event keys to their remaining block count.
/// When count reaches 0, the entry is removed.
/// When count reaches 0, the entry is removed.
blocked : Arc < Mutex < HashMap < BlockKey , usize > > > ,
blocked : Arc < Mutex < HashMap < BlockKey , usize > > > ,
/// Path-prefix suppressions with a fixed expiry.
prefix_blocked : Arc < Mutex < Vec < PrefixBlock > > > ,
}
}
impl EventBlocker {
impl EventBlocker {
@ -64,6 +78,7 @@ impl EventBlocker {
pub fn new ( ) -> Self {
pub fn new ( ) -> Self {
Self {
Self {
blocked : Arc ::new ( Mutex ::new ( HashMap ::new ( ) ) ) ,
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 ) ;
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.
/// Checks if an event should be blocked and decrements the counter if so.
///
///
/// # Arguments
/// # Arguments
@ -100,6 +144,16 @@ impl EventBlocker {
/// # Returns
/// # Returns
/// `true` if the event should be blocked (was pre-registered), `false` otherwise
/// `true` if the event should be blocked (was pre-registered), `false` otherwise
pub fn should_block ( & self , kind : & EventKind , path : & PathBuf ) -> bool {
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 {
let key = BlockKey {
kind : NormalizedEventKind ::from ( kind ) ,
kind : NormalizedEventKind ::from ( kind ) ,
path : path . clone ( ) ,
path : path . clone ( ) ,
@ -156,12 +210,15 @@ impl EventBlocker {
pub fn clear ( & self ) {
pub fn clear ( & self ) {
let mut blocked = self . blocked . lock ( ) . unwrap ( ) ;
let mut blocked = self . blocked . lock ( ) . unwrap ( ) ;
blocked . clear ( ) ;
blocked . clear ( ) ;
let mut prefix_blocked = self . prefix_blocked . lock ( ) . unwrap ( ) ;
prefix_blocked . clear ( ) ;
}
}
/// Returns the number of currently registered event blocks.
/// Returns the number of currently registered event blocks.
pub fn len ( & self ) -> usize {
pub fn len ( & self ) -> usize {
let blocked = self . blocked . lock ( ) . unwrap ( ) ;
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.
/// Returns true if there are no registered event blocks.
@ -169,3 +226,64 @@ impl EventBlocker {
self . len ( ) = = 0
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 ) ) ;
}
}