@@ -167,33 +167,38 @@ impl FileWatcher {
167167 let config = self . config . clone ( ) ;
168168 let watched_paths = self . watched_paths . clone ( ) ;
169169
170- tokio:: spawn ( async move {
171- let mut last_flush = std:: time:: Instant :: now ( ) ;
172-
173- while let Ok ( event) = rx. recv ( ) {
174- match event {
175- Ok ( event) => {
176- if Self :: should_ignore_event ( & event, & watched_paths) . await {
177- continue ;
178- }
179-
180- if let Some ( file_event) = Self :: convert_event ( & event) {
181- {
182- let mut buffer = lock_event_buffer ( & event_buffer) ;
183- buffer. push ( file_event) ;
184- }
185-
186- let now = std:: time:: Instant :: now ( ) ;
187- if now. duration_since ( last_flush) . as_millis ( ) as u64
188- >= config. debounce_interval_ms
189- {
190- Self :: flush_events_static ( & event_buffer, & emitter_arc) . await ;
191- last_flush = now;
170+ // Run on a dedicated blocking thread to avoid starving the async runtime.
171+ // True debounce: accumulate events, then flush once the stream goes quiet for
172+ // `debounce_interval_ms`. A 50 ms poll interval keeps latency low even for
173+ // single-event bursts (e.g. one `fs::write` from an agentic tool).
174+ tokio:: task:: spawn_blocking ( move || {
175+ let rt = tokio:: runtime:: Handle :: current ( ) ;
176+ let debounce = std:: time:: Duration :: from_millis ( config. debounce_interval_ms ) ;
177+ let poll = std:: time:: Duration :: from_millis ( 50 ) ;
178+ let mut last_event_time: Option < std:: time:: Instant > = None ;
179+
180+ loop {
181+ match rx. recv_timeout ( poll) {
182+ Ok ( Ok ( event) ) => {
183+ let ignore =
184+ rt. block_on ( Self :: should_ignore_event ( & event, & watched_paths) ) ;
185+ if !ignore {
186+ if let Some ( file_event) = Self :: convert_event ( & event) {
187+ lock_event_buffer ( & event_buffer) . push ( file_event) ;
188+ last_event_time = Some ( std:: time:: Instant :: now ( ) ) ;
192189 }
193190 }
194191 }
195- Err ( e) => {
196- eprintln ! ( "Watch error: {:?}" , e) ;
192+ Ok ( Err ( e) ) => eprintln ! ( "Watch error: {:?}" , e) ,
193+ Err ( std:: sync:: mpsc:: RecvTimeoutError :: Timeout ) => { }
194+ Err ( std:: sync:: mpsc:: RecvTimeoutError :: Disconnected ) => break ,
195+ }
196+
197+ // Flush only after events have been quiet for the debounce window.
198+ if let Some ( t) = last_event_time {
199+ if t. elapsed ( ) >= debounce {
200+ rt. block_on ( Self :: flush_events_static ( & event_buffer, & emitter_arc) ) ;
201+ last_event_time = None ;
197202 }
198203 }
199204 }
0 commit comments