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 }) }),
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(¤t_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(¤t_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(¤t_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, ¤t_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}