Skip to main content

vtcode_core/tools/
edited_file_monitor.rs

1use crate::audit::{FileConflictAuditEvent, FileConflictAuditLog};
2use crate::config::PermissionsConfig;
3use crate::tools::file_ops::{build_diff_preview, diff_preview_error_skip};
4use anyhow::{Context, Result};
5use notify::{EventKind, RecommendedWatcher, RecursiveMode, Watcher};
6use parking_lot::Mutex;
7use serde_json::{Value, json};
8use std::collections::{HashMap, HashSet, VecDeque};
9use std::path::{Path, PathBuf};
10use std::sync::Arc;
11use std::sync::Mutex as StdMutex;
12use std::sync::atomic::{AtomicU64, Ordering};
13use std::sync::mpsc::{self, Receiver, RecvTimeoutError};
14use std::thread;
15use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
16use tokio::sync::Notify;
17use vtcode_commons::utils::calculate_sha256;
18use vtcode_commons::{canonicalize, workspace_relative_display};
19
20pub const FILE_CONFLICT_OVERRIDE_ARG: &str = "__vtcode_conflict_override";
21pub const FILE_CONFLICT_DETECTED_FIELD: &str = "conflict_detected";
22pub const FILE_CONFLICT_PATH_FIELD: &str = "conflict_path";
23
24#[derive(Clone, Debug, PartialEq, Eq)]
25pub struct FileSnapshot {
26    pub exists: bool,
27    pub size_bytes: u64,
28    pub modified_millis: Option<u128>,
29    pub sha256: String,
30    pub text_content: Option<String>,
31}
32
33impl FileSnapshot {
34    fn fingerprint(&self) -> String {
35        let modified = self
36            .modified_millis
37            .map(|value| value.to_string())
38            .unwrap_or_else(|| "none".to_string());
39        format!("exists={};size={};modified={};sha256={}", self.exists, self.size_bytes, modified, self.sha256)
40    }
41
42    fn to_json(&self) -> Value {
43        json!({
44            "exists": self.exists,
45            "size_bytes": self.size_bytes,
46            "modified_millis": self.modified_millis,
47            "sha256": self.sha256,
48        })
49    }
50
51    fn same_contents(&self, other: &Self) -> bool {
52        self.exists == other.exists && self.size_bytes == other.size_bytes && self.sha256 == other.sha256
53    }
54
55    fn same_identity(&self, other: &Self) -> bool {
56        self.same_contents(other) && self.modified_millis == other.modified_millis
57    }
58
59    fn from_identity_value(value: &Value) -> Option<Self> {
60        let object = value.as_object()?;
61        let modified_millis = match object.get("modified_millis") {
62            Some(Value::Null) | None => None,
63            Some(value) => value.as_u64().map(u128::from).or_else(|| {
64                value.as_i64().filter(|millis| *millis >= 0).map(|millis| {
65                    #[allow(
66                        clippy::cast_sign_loss,
67                        reason = "Intentional compatibility, platform, or test-only suppression."
68                    )]
69                    let val: u64 = millis as u64;
70                    val as u128
71                }) // safe: filtered non-negative
72            }),
73        };
74
75        Some(Self {
76            exists: object.get("exists")?.as_bool()?,
77            size_bytes: object.get("size_bytes")?.as_u64()?,
78            modified_millis,
79            sha256: object.get("sha256")?.as_str()?.to_string(),
80            text_content: None,
81        })
82    }
83
84    fn from_text_content(content: &str) -> Self {
85        let bytes = content.as_bytes();
86        let sha256 = calculate_sha256(bytes);
87
88        Self {
89            exists: true,
90            size_bytes: bytes.len() as u64,
91            modified_millis: None,
92            sha256,
93            text_content: Some(content.to_string()),
94        }
95    }
96
97    fn missing() -> Self {
98        Self {
99            exists: false,
100            size_bytes: 0,
101            modified_millis: None,
102            sha256: String::new(),
103            text_content: None,
104        }
105    }
106}
107
108#[derive(Clone, Debug)]
109pub struct FileConflict {
110    pub path: PathBuf,
111    pub read_snapshot: Option<FileSnapshot>,
112    pub disk_snapshot: Option<FileSnapshot>,
113    pub intended_content: Option<String>,
114    pub emit_hitl_notification: bool,
115}
116
117#[derive(Clone, Debug, PartialEq, Eq)]
118pub struct TrackedPathFreshness {
119    pub path: PathBuf,
120    pub is_stale: bool,
121    pub fingerprint: Option<String>,
122}
123
124impl FileConflict {
125    pub fn to_tool_output(&self, workspace_root: &Path) -> Value {
126        let display_path = workspace_relative_display(workspace_root, &self.path);
127        let disk_content = self.disk_snapshot.as_ref().and_then(|snapshot| snapshot.text_content.clone());
128        let diff_preview = match (&disk_content, &self.intended_content) {
129            (Some(before), Some(after)) => build_diff_preview(&display_path, Some(before), after),
130            (None, Some(_)) => diff_preview_error_skip("binary_or_non_utf8_disk_content", None),
131            _ => diff_preview_error_skip("missing_intended_content", None),
132        };
133
134        json!({
135            "success": true,
136            FILE_CONFLICT_DETECTED_FIELD: true,
137            FILE_CONFLICT_PATH_FIELD: display_path,
138            "message": "File changed on disk since the agent last read it.",
139            "resolution": "pending",
140            "emit_hitl_notification": self.emit_hitl_notification,
141            "disk_content": disk_content,
142            "intended_content": self.intended_content,
143            "read_snapshot": self.read_snapshot.as_ref().map(FileSnapshot::to_json),
144            "disk_snapshot": self.disk_snapshot.as_ref().map(FileSnapshot::to_json),
145            "diff_preview": diff_preview,
146        })
147    }
148}
149
150#[derive(Clone, Debug)]
151struct StaleConflictState {
152    notification_emitted: bool,
153}
154
155#[derive(Clone, Debug)]
156struct TrackedFileState {
157    last_read_snapshot: Option<FileSnapshot>,
158    last_known_disk_snapshot: Option<FileSnapshot>,
159    last_agent_write_snapshot: Option<FileSnapshot>,
160    active_mutation: Option<u64>,
161    pending_mutations: VecDeque<u64>,
162    stale_conflict: Option<StaleConflictState>,
163    notify: Arc<Notify>,
164}
165
166impl TrackedFileState {
167    fn new() -> Self {
168        Self {
169            last_read_snapshot: None,
170            last_known_disk_snapshot: None,
171            last_agent_write_snapshot: None,
172            active_mutation: None,
173            pending_mutations: VecDeque::new(),
174            stale_conflict: None,
175            notify: Arc::new(Notify::new()),
176        }
177    }
178}
179
180#[derive(Default)]
181struct MonitorState {
182    tracked_files: HashMap<PathBuf, TrackedFileState>,
183    watched_parents: HashSet<PathBuf>,
184}
185
186impl MonitorState {
187    fn tracked_entry(&mut self, path: &Path) -> &mut TrackedFileState {
188        self.tracked_files
189            .entry(path.to_path_buf())
190            .or_insert_with(TrackedFileState::new)
191    }
192}
193
194struct EditedFileMonitorInner {
195    state: Mutex<MonitorState>,
196    watcher: StdMutex<Option<RecommendedWatcher>>,
197    audit_log: Mutex<Option<FileConflictAuditLog>>,
198    next_mutation_id: AtomicU64,
199    debounce_duration: Duration,
200}
201
202#[derive(Clone)]
203pub struct EditedFileMonitor {
204    inner: Arc<EditedFileMonitorInner>,
205}
206
207pub struct MutationLease {
208    inner: Arc<EditedFileMonitorInner>,
209    path: PathBuf,
210    mutation_id: u64,
211    released: bool,
212}
213
214struct PendingMutationGuard {
215    inner: Arc<EditedFileMonitorInner>,
216    path: PathBuf,
217    mutation_id: u64,
218    armed: bool,
219}
220
221impl PendingMutationGuard {
222    fn new(inner: Arc<EditedFileMonitorInner>, path: PathBuf, mutation_id: u64) -> Self {
223        Self { inner, path, mutation_id, armed: true }
224    }
225
226    fn disarm(&mut self) {
227        self.armed = false;
228    }
229}
230
231impl Drop for PendingMutationGuard {
232    fn drop(&mut self) {
233        if !self.armed {
234            return;
235        }
236
237        let mut state = self.inner.state.lock();
238        let Some(entry) = state.tracked_files.get_mut(&self.path) else {
239            return;
240        };
241
242        if remove_pending_mutation(entry, self.mutation_id) {
243            entry.notify.notify_waiters();
244        }
245    }
246}
247
248impl MutationLease {
249    pub fn path(&self) -> &Path {
250        &self.path
251    }
252}
253
254impl Drop for MutationLease {
255    fn drop(&mut self) {
256        if self.released {
257            return;
258        }
259        let mut state = self.inner.state.lock();
260        if let Some(entry) = state.tracked_files.get_mut(&self.path) {
261            let mut released_active = false;
262            if entry.active_mutation == Some(self.mutation_id) {
263                entry.active_mutation = None;
264                released_active = true;
265            }
266            let removed_pending = remove_pending_mutation(entry, self.mutation_id);
267            if released_active || removed_pending {
268                entry.notify.notify_waiters();
269            }
270        }
271        self.released = true;
272    }
273}
274
275impl EditedFileMonitor {
276    pub fn new() -> Self {
277        let (event_tx, event_rx) = mpsc::channel::<PathBuf>();
278        let watcher = RecommendedWatcher::new(
279            move |event: notify::Result<notify::Event>| {
280                let Ok(event) = event else {
281                    return;
282                };
283
284                if !matches!(event.kind, EventKind::Modify(_) | EventKind::Create(_) | EventKind::Remove(_)) {
285                    return;
286                }
287
288                for path in event.paths {
289                    let _ = event_tx.send(path);
290                }
291            },
292            notify::Config::default(),
293        )
294        .ok();
295
296        let inner = Arc::new(EditedFileMonitorInner {
297            state: Mutex::new(MonitorState::default()),
298            watcher: StdMutex::new(watcher),
299            audit_log: Mutex::new(None),
300            next_mutation_id: AtomicU64::new(1),
301            debounce_duration: Duration::from_millis(250),
302        });
303
304        let monitor = Self { inner: Arc::clone(&inner) };
305        spawn_event_loop(inner, event_rx);
306        monitor
307    }
308
309    pub fn apply_permissions_config(&self, permissions: &PermissionsConfig) {
310        let mut audit_log = self.inner.audit_log.lock();
311        if !permissions.audit_enabled {
312            *audit_log = None;
313            return;
314        }
315
316        let audit_dir = vtcode_commons::paths::expand_tilde(&permissions.audit_directory);
317        match FileConflictAuditLog::new(audit_dir) {
318            Ok(log) => *audit_log = Some(log),
319            Err(err) => {
320                tracing::warn!(error = %err, "Failed to initialize file conflict audit log");
321                *audit_log = None;
322            }
323        }
324    }
325
326    pub async fn track_read(&self, path: &Path) -> Result<()> {
327        let path = normalize_event_path(path);
328        let snapshot = snapshot_path_async(path.clone()).await?;
329        self.record_read_snapshot(&path, snapshot)
330    }
331
332    pub async fn accept_disk_version(&self, path: &Path) -> Result<()> {
333        let path = normalize_event_path(path);
334        let snapshot = snapshot_path_async(path.clone()).await?;
335        self.record_read_snapshot(&path, snapshot)
336    }
337
338    pub fn record_read_snapshot(&self, path: &Path, snapshot: FileSnapshot) -> Result<()> {
339        let path = normalize_event_path(path);
340        {
341            let mut state = self.inner.state.lock();
342            let entry = state.tracked_entry(&path);
343            entry.last_read_snapshot = Some(snapshot.clone());
344            entry.last_known_disk_snapshot = Some(snapshot);
345            entry.stale_conflict = None;
346        }
347        self.watch_parent(&path);
348        Ok(())
349    }
350
351    pub fn record_read_text(&self, path: &Path, content: &str) -> Result<()> {
352        self.record_read_snapshot(path, FileSnapshot::from_text_content(content))
353    }
354
355    pub fn record_agent_write_snapshot(&self, path: &Path, snapshot: FileSnapshot) -> Result<()> {
356        let path = normalize_event_path(path);
357        {
358            let mut state = self.inner.state.lock();
359            let entry = state.tracked_entry(&path);
360            let snap_for_disk = snapshot.clone();
361            entry.last_read_snapshot = Some(snap_for_disk);
362            entry.last_known_disk_snapshot = Some(snapshot.clone());
363            entry.last_agent_write_snapshot = Some(snapshot);
364            entry.stale_conflict = None;
365        }
366        self.watch_parent(&path);
367        Ok(())
368    }
369
370    pub fn record_agent_write_text(&self, path: &Path, content: &str) -> Result<()> {
371        self.record_agent_write_snapshot(path, FileSnapshot::from_text_content(content))
372    }
373
374    pub fn record_agent_removal(&self, path: &Path) -> Result<()> {
375        self.record_agent_write_snapshot(path, FileSnapshot::missing())
376    }
377
378    pub async fn tracked_read_text(&self, path: &Path) -> Option<String> {
379        let path = normalize_event_path(path);
380        let state = self.inner.state.lock();
381        state
382            .tracked_files
383            .get(&path)
384            .and_then(|entry| entry.last_read_snapshot.as_ref())
385            .and_then(|snapshot| snapshot.text_content.clone())
386    }
387
388    pub fn tracked_paths(&self) -> Vec<PathBuf> {
389        let state = self.inner.state.lock();
390        let mut paths = state.tracked_files.keys().cloned().collect::<Vec<_>>();
391        paths.sort();
392        paths
393    }
394
395    pub fn stale_tracked_paths(&self) -> Vec<PathBuf> {
396        self.tracked_path_freshness()
397            .into_iter()
398            .filter_map(|entry| entry.is_stale.then_some(entry.path))
399            .collect()
400    }
401
402    pub fn tracked_path_freshness(&self) -> Vec<TrackedPathFreshness> {
403        let state = self.inner.state.lock();
404        let mut entries = state
405            .tracked_files
406            .iter()
407            .map(|(path, entry)| {
408                let is_stale = entry
409                    .last_read_snapshot
410                    .as_ref()
411                    .zip(entry.last_known_disk_snapshot.as_ref())
412                    .is_some_and(|(read_snapshot, disk_snapshot)| !read_snapshot.same_contents(disk_snapshot));
413                TrackedPathFreshness {
414                    path: path.clone(),
415                    is_stale,
416                    fingerprint: entry.last_known_disk_snapshot.as_ref().map(FileSnapshot::fingerprint),
417                }
418            })
419            .collect::<Vec<_>>();
420        entries.sort_by(|left, right| left.path.cmp(&right.path));
421        entries
422    }
423
424    pub fn tracked_path_fingerprint(&self, path: &Path) -> Option<String> {
425        let path = normalize_event_path(path);
426        let state = self.inner.state.lock();
427        state
428            .tracked_files
429            .get(&path)
430            .and_then(|entry| entry.last_known_disk_snapshot.as_ref())
431            .map(FileSnapshot::fingerprint)
432    }
433
434    pub async fn acquire_mutation(&self, path: &Path) -> MutationLease {
435        let path = normalize_event_path(path);
436        let mutation_id = self.inner.next_mutation_id.fetch_add(1, Ordering::SeqCst);
437        let mut pending_guard = PendingMutationGuard::new(Arc::clone(&self.inner), path.clone(), mutation_id);
438
439        loop {
440            let notify = {
441                let mut state = self.inner.state.lock();
442                let entry = state.tracked_entry(&path);
443
444                if entry.active_mutation.is_none() {
445                    if let Some(front) = entry.pending_mutations.front() {
446                        if *front == mutation_id {
447                            let _ = entry.pending_mutations.pop_front();
448                            entry.active_mutation = Some(mutation_id);
449                            pending_guard.disarm();
450                            return MutationLease {
451                                inner: Arc::clone(&self.inner),
452                                path,
453                                mutation_id,
454                                released: false,
455                            };
456                        }
457                    } else {
458                        entry.active_mutation = Some(mutation_id);
459                        pending_guard.disarm();
460                        return MutationLease {
461                            inner: Arc::clone(&self.inner),
462                            path,
463                            mutation_id,
464                            released: false,
465                        };
466                    }
467                }
468
469                if !entry.pending_mutations.iter().any(|pending_id| *pending_id == mutation_id) {
470                    entry.pending_mutations.push_back(mutation_id);
471                }
472
473                entry.notify.clone()
474            };
475
476            notify.notified().await;
477        }
478    }
479
480    pub async fn detect_conflict(
481        &self,
482        path: &Path,
483        intended_content: Option<String>,
484        approved_snapshot: Option<FileSnapshot>,
485    ) -> Result<Option<FileConflict>> {
486        let path = normalize_event_path(path);
487        let current_snapshot = snapshot_path_async(path.clone()).await?;
488        let mut should_audit = false;
489
490        let maybe_conflict = {
491            let mut state = self.inner.state.lock();
492            let entry = state.tracked_entry(&path);
493            entry.last_known_disk_snapshot = Some(current_snapshot.clone());
494
495            if entry
496                .last_agent_write_snapshot
497                .as_ref()
498                .is_some_and(|snapshot| snapshot.same_contents(&current_snapshot))
499            {
500                entry.last_agent_write_snapshot = None;
501                entry.last_read_snapshot = Some(current_snapshot);
502                entry.stale_conflict = None;
503                return Ok(None);
504            }
505
506            let Some(read_snapshot) = entry.last_read_snapshot.clone() else {
507                return Ok(None);
508            };
509
510            if read_snapshot.same_contents(&current_snapshot) {
511                entry.stale_conflict = None;
512                return Ok(None);
513            }
514
515            if approved_snapshot
516                .as_ref()
517                .is_some_and(|snapshot| snapshot.same_identity(&current_snapshot))
518            {
519                entry.stale_conflict = None;
520                return Ok(None);
521            }
522
523            let emit_hitl_notification = match entry.stale_conflict.as_mut() {
524                Some(existing) => {
525                    let emit = !existing.notification_emitted;
526                    existing.notification_emitted = true;
527                    emit
528                }
529                None => {
530                    entry.stale_conflict = Some(StaleConflictState { notification_emitted: true });
531                    should_audit = true;
532                    true
533                }
534            };
535
536            Some(FileConflict {
537                path: path.clone(),
538                read_snapshot: Some(read_snapshot),
539                disk_snapshot: Some(current_snapshot.clone()),
540                intended_content,
541                emit_hitl_notification,
542            })
543        };
544
545        if should_audit {
546            self.record_conflict_audit(&path, &current_snapshot, "pre_write_conflict");
547        }
548
549        Ok(maybe_conflict)
550    }
551
552    #[cfg(test)]
553    pub async fn debug_process_path_change(&self, path: &Path) -> Result<()> {
554        self.process_path_change(path.to_path_buf())
555    }
556
557    fn watch_parent(&self, path: &Path) {
558        let Some(parent) = path.parent().map(Path::to_path_buf) else {
559            return;
560        };
561
562        let should_watch = {
563            let mut state = self.inner.state.lock();
564            state.watched_parents.insert(parent.clone())
565        };
566
567        if !should_watch {
568            return;
569        }
570
571        let Ok(mut watcher) = self.inner.watcher.lock() else {
572            return;
573        };
574        let Some(watcher) = watcher.as_mut() else {
575            return;
576        };
577
578        if let Err(err) = watcher.watch(&parent, RecursiveMode::NonRecursive) {
579            tracing::warn!(path = %parent.display(), error = %err, "Failed to watch edited-file parent directory");
580        }
581    }
582
583    fn record_conflict_audit(&self, path: &Path, snapshot: &FileSnapshot, reason: &str) {
584        let event = FileConflictAuditEvent {
585            timestamp: chrono::Local::now(),
586            path: path.to_path_buf(),
587            reason: reason.to_string(),
588            file_exists: snapshot.exists,
589            size_bytes: snapshot.exists.then_some(snapshot.size_bytes),
590            sha256: snapshot.exists.then_some(snapshot.sha256.clone()),
591        };
592
593        if let Some(log) = self.inner.audit_log.lock().as_mut()
594            && let Err(err) = log.record(&event)
595        {
596            tracing::warn!(error = %err, "Failed to record file conflict audit event");
597        }
598    }
599
600    fn process_path_change(&self, path: PathBuf) -> Result<()> {
601        let tracked_path = normalize_event_path(&path);
602        let snapshot = snapshot_path_sync(&tracked_path)
603            .with_context(|| format!("Failed to snapshot externally modified file {}", tracked_path.display()))?;
604
605        let mut should_audit = false;
606
607        {
608            let mut state = self.inner.state.lock();
609            let Some(entry) = state.tracked_files.get_mut(&tracked_path) else {
610                return Ok(());
611            };
612
613            if entry
614                .last_known_disk_snapshot
615                .as_ref()
616                .is_some_and(|known| known.same_contents(&snapshot))
617            {
618                return Ok(());
619            }
620
621            entry.last_known_disk_snapshot = Some(snapshot.clone());
622
623            if entry
624                .last_agent_write_snapshot
625                .as_ref()
626                .is_some_and(|known| known.same_contents(&snapshot))
627            {
628                entry.last_read_snapshot = Some(snapshot);
629                entry.last_agent_write_snapshot = None;
630                entry.stale_conflict = None;
631                return Ok(());
632            }
633
634            if entry
635                .last_read_snapshot
636                .as_ref()
637                .is_some_and(|known| known.same_contents(&snapshot))
638            {
639                entry.stale_conflict = None;
640                return Ok(());
641            }
642
643            if entry
644                .last_read_snapshot
645                .as_ref()
646                .is_some_and(|read_snapshot| !read_snapshot.same_contents(&snapshot))
647            {
648                should_audit = true;
649                match entry.stale_conflict.as_mut() {
650                    Some(_) => {}
651                    None => {
652                        entry.stale_conflict = Some(StaleConflictState { notification_emitted: false });
653                    }
654                }
655            }
656        }
657
658        if should_audit {
659            self.record_conflict_audit(&tracked_path, &snapshot, "watcher_detected_external_change");
660        }
661
662        Ok(())
663    }
664}
665
666impl Default for EditedFileMonitor {
667    fn default() -> Self {
668        Self::new()
669    }
670}
671
672fn spawn_event_loop(inner: Arc<EditedFileMonitorInner>, event_rx: Receiver<PathBuf>) {
673    thread::spawn(move || {
674        let mut pending = HashMap::<PathBuf, Instant>::new();
675
676        loop {
677            match event_rx.recv_timeout(Duration::from_millis(50)) {
678                Ok(path) => {
679                    pending.insert(normalize_event_path(&path), Instant::now());
680                }
681                Err(RecvTimeoutError::Timeout) => {}
682                Err(RecvTimeoutError::Disconnected) => break,
683            }
684
685            let now = Instant::now();
686            let ready = pending
687                .iter()
688                .filter_map(|(path, observed_at)| {
689                    (now.duration_since(*observed_at) >= inner.debounce_duration).then_some(path.clone())
690                })
691                .collect::<Vec<_>>();
692
693            for path in ready {
694                pending.remove(&path);
695                let monitor = EditedFileMonitor { inner: Arc::clone(&inner) };
696                if let Err(err) = monitor.process_path_change(path.clone()) {
697                    tracing::debug!(path = %path.display(), error = %err, "Edited-file watcher refresh failed");
698                }
699            }
700        }
701    });
702}
703
704fn normalize_event_path(path: &Path) -> PathBuf {
705    canonicalize(path).unwrap_or_else(|_| {
706        path.parent()
707            .and_then(|parent| canonicalize(parent).ok())
708            .and_then(|parent| path.file_name().map(|name| parent.join(name)))
709            .unwrap_or_else(|| path.to_path_buf())
710    })
711}
712
713async fn snapshot_path_async(path: PathBuf) -> Result<FileSnapshot> {
714    tokio::task::spawn_blocking(move || snapshot_path_sync(&path))
715        .await
716        .context("Failed to join file snapshot task")?
717}
718
719fn snapshot_path_sync(path: &Path) -> Result<FileSnapshot> {
720    let metadata = match std::fs::metadata(path) {
721        Ok(metadata) => metadata,
722        Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
723            return Ok(FileSnapshot {
724                exists: false,
725                size_bytes: 0,
726                modified_millis: None,
727                sha256: String::new(),
728                text_content: None,
729            });
730        }
731        Err(err) => {
732            return Err(err).with_context(|| format!("Failed to read metadata for {}", path.display()));
733        }
734    };
735
736    if !metadata.is_file() {
737        return Ok(FileSnapshot {
738            exists: false,
739            size_bytes: 0,
740            modified_millis: metadata.modified().ok().and_then(system_time_to_millis),
741            sha256: String::new(),
742            text_content: None,
743        });
744    }
745
746    let bytes = std::fs::read(path).with_context(|| format!("Failed to read file bytes for {}", path.display()))?;
747    let sha256 = calculate_sha256(&bytes);
748
749    Ok(FileSnapshot {
750        exists: true,
751        size_bytes: metadata.len(),
752        modified_millis: metadata.modified().ok().and_then(system_time_to_millis),
753        sha256,
754        text_content: String::from_utf8(bytes).ok(),
755    })
756}
757
758fn system_time_to_millis(time: SystemTime) -> Option<u128> {
759    time.duration_since(UNIX_EPOCH).ok().map(|duration| duration.as_millis())
760}
761
762pub fn conflict_override_snapshot(args: &Value) -> Option<FileSnapshot> {
763    args.get(FILE_CONFLICT_OVERRIDE_ARG).and_then(FileSnapshot::from_identity_value)
764}
765
766fn remove_pending_mutation(entry: &mut TrackedFileState, mutation_id: u64) -> bool {
767    let Some(index) = entry.pending_mutations.iter().position(|pending_id| *pending_id == mutation_id) else {
768        return false;
769    };
770
771    let _ = entry.pending_mutations.remove(index);
772    true
773}
774
775#[cfg(test)]
776mod tests {
777    use super::*;
778    use tempfile::TempDir;
779
780    #[tokio::test]
781    async fn detects_same_mtime_same_size_hash_mismatch() -> Result<()> {
782        let temp = TempDir::new()?;
783        let file = temp.path().join("sample.txt");
784        std::fs::write(&file, "before\n")?;
785
786        let first = snapshot_path_async(file.clone()).await?;
787        std::fs::write(&file, "after!\n")?;
788        let second = snapshot_path_async(file.clone()).await?;
789
790        assert_eq!(first.size_bytes, second.size_bytes);
791        assert_ne!(first.sha256, second.sha256);
792        Ok(())
793    }
794
795    #[tokio::test]
796    async fn suppresses_self_write_event() -> Result<()> {
797        let temp = TempDir::new()?;
798        let file = temp.path().join("sample.txt");
799        std::fs::write(&file, "before\n")?;
800        let monitor = EditedFileMonitor::new();
801
802        monitor.track_read(&file).await?;
803        std::fs::write(&file, "after\n")?;
804        monitor.record_agent_write_text(&file, "after\n")?;
805        monitor.debug_process_path_change(&file).await?;
806
807        assert!(
808            monitor
809                .detect_conflict(&file, Some("after\n".to_string()), None)
810                .await?
811                .is_none()
812        );
813        Ok(())
814    }
815
816    #[tokio::test]
817    async fn expected_agent_snapshot_does_not_hide_external_follow_up_change() -> Result<()> {
818        let temp = TempDir::new()?;
819        let file = temp.path().join("sample.txt");
820        std::fs::write(&file, "before\n")?;
821        let monitor = EditedFileMonitor::new();
822
823        monitor.track_read(&file).await?;
824        std::fs::write(&file, "after\n")?;
825        monitor.record_agent_write_text(&file, "after\n")?;
826
827        std::fs::write(&file, "after formatted\n")?;
828        monitor.debug_process_path_change(&file).await?;
829
830        let conflict = monitor.detect_conflict(&file, Some("agent\n".to_string()), None).await?;
831        assert!(conflict.is_some());
832        Ok(())
833    }
834
835    #[tokio::test]
836    async fn queues_mutations_in_order() -> Result<()> {
837        let temp = TempDir::new()?;
838        let file = temp.path().join("sample.txt");
839        std::fs::write(&file, "hello\n")?;
840        let monitor = Arc::new(EditedFileMonitor::new());
841        let normalized_file = normalize_event_path(&file);
842
843        let first = monitor.acquire_mutation(&file).await;
844        let monitor_clone = Arc::clone(&monitor);
845        let file_clone = file.clone();
846        let waiter = tokio::spawn(async move {
847            let lease = monitor_clone.acquire_mutation(&file_clone).await;
848            lease.path().to_path_buf()
849        });
850
851        tokio::time::sleep(Duration::from_millis(50)).await;
852        drop(first);
853        let acquired = waiter.await?;
854        assert_eq!(acquired, normalized_file);
855        Ok(())
856    }
857
858    #[tokio::test]
859    async fn detects_external_change_after_read() -> Result<()> {
860        let temp = TempDir::new()?;
861        let file = temp.path().join("sample.txt");
862        std::fs::write(&file, "before\n")?;
863        let monitor = EditedFileMonitor::new();
864
865        monitor.track_read(&file).await?;
866        std::fs::write(&file, "external\n")?;
867        monitor.debug_process_path_change(&file).await?;
868
869        let conflict = monitor.detect_conflict(&file, Some("agent\n".to_string()), None).await?;
870        assert!(conflict.is_some());
871        Ok(())
872    }
873
874    #[tokio::test]
875    async fn cancelled_waiter_does_not_block_following_mutation() -> Result<()> {
876        let temp = TempDir::new()?;
877        let file = temp.path().join("sample.txt");
878        std::fs::write(&file, "hello\n")?;
879        let monitor = Arc::new(EditedFileMonitor::new());
880        let normalized_file = normalize_event_path(&file);
881
882        let first = monitor.acquire_mutation(&file).await;
883
884        let pending_monitor = Arc::clone(&monitor);
885        let pending_file = file.clone();
886        let pending = tokio::spawn(async move {
887            let _lease = pending_monitor.acquire_mutation(&pending_file).await;
888        });
889        tokio::time::sleep(Duration::from_millis(50)).await;
890        pending.abort();
891
892        let next_monitor = Arc::clone(&monitor);
893        let next_file = file.clone();
894        let next = tokio::spawn(async move {
895            let lease = next_monitor.acquire_mutation(&next_file).await;
896            lease.path().to_path_buf()
897        });
898
899        drop(first);
900
901        let acquired = tokio::time::timeout(Duration::from_secs(1), next).await??;
902        assert_eq!(acquired, normalized_file);
903        Ok(())
904    }
905
906    #[tokio::test]
907    async fn clears_conflict_state_when_disk_returns_to_read_snapshot() -> Result<()> {
908        let temp = TempDir::new()?;
909        let file = temp.path().join("sample.txt");
910        std::fs::write(&file, "before\n")?;
911        let monitor = EditedFileMonitor::new();
912
913        monitor.track_read(&file).await?;
914        std::fs::write(&file, "external one\n")?;
915        monitor.debug_process_path_change(&file).await?;
916
917        let first_conflict = monitor
918            .detect_conflict(&file, Some("agent\n".to_string()), None)
919            .await?
920            .expect("expected initial conflict");
921        assert!(first_conflict.emit_hitl_notification);
922
923        std::fs::write(&file, "before\n")?;
924        monitor.debug_process_path_change(&file).await?;
925        assert!(
926            monitor
927                .detect_conflict(&file, Some("agent\n".to_string()), None)
928                .await?
929                .is_none()
930        );
931
932        std::fs::write(&file, "external two\n")?;
933        monitor.debug_process_path_change(&file).await?;
934        let second_conflict = monitor
935            .detect_conflict(&file, Some("agent\n".to_string()), None)
936            .await?
937            .expect("expected renewed conflict");
938        assert!(second_conflict.emit_hitl_notification);
939        Ok(())
940    }
941
942    #[tokio::test]
943    async fn stale_path_queries_expose_external_change_without_clearing_state() -> Result<()> {
944        let temp = TempDir::new()?;
945        let file = temp.path().join("sample.txt");
946        std::fs::write(&file, "before\n")?;
947        let monitor = EditedFileMonitor::new();
948
949        monitor.track_read(&file).await?;
950        let tracked_paths = monitor.tracked_paths();
951        assert_eq!(tracked_paths.len(), 1);
952        let tracked_path = tracked_paths[0].clone();
953        assert!(monitor.stale_tracked_paths().is_empty());
954
955        std::fs::write(&file, "external\n")?;
956        monitor.debug_process_path_change(&file).await?;
957
958        let freshness = monitor.tracked_path_freshness();
959        assert_eq!(freshness.len(), 1);
960        assert_eq!(freshness[0].path, tracked_path);
961        assert!(freshness[0].is_stale);
962        assert!(freshness[0].fingerprint.is_some());
963        assert_eq!(monitor.stale_tracked_paths(), vec![freshness[0].path.clone()]);
964        assert_eq!(monitor.tracked_path_fingerprint(&freshness[0].path), freshness[0].fingerprint);
965        Ok(())
966    }
967
968    #[tokio::test]
969    async fn override_requires_matching_approved_snapshot() -> Result<()> {
970        let temp = TempDir::new()?;
971        let file = temp.path().join("sample.txt");
972        std::fs::write(&file, "before\n")?;
973        let monitor = EditedFileMonitor::new();
974
975        monitor.track_read(&file).await?;
976        std::fs::write(&file, "external one\n")?;
977        let approved_snapshot = snapshot_path_async(file.clone()).await?;
978
979        std::fs::write(&file, "external two\n")?;
980
981        let conflict = monitor
982            .detect_conflict(&file, Some("agent\n".to_string()), Some(approved_snapshot))
983            .await?;
984        assert!(conflict.is_some());
985        Ok(())
986    }
987
988    #[test]
989    fn normalizes_missing_event_paths_via_canonical_parent() -> Result<()> {
990        let temp = TempDir::new()?;
991        let missing = temp.path().join("missing.txt");
992        let canonical_parent = canonicalize(temp.path())?;
993
994        assert_eq!(normalize_event_path(&missing), canonical_parent.join("missing.txt"));
995        Ok(())
996    }
997}