1use chrono::{DateTime, Utc};
7use ragfs_core::VectorStore;
8use serde::{Deserialize, Serialize};
9use std::fs;
10use std::path::{Path, PathBuf};
11use std::sync::Arc;
12use tokio::sync::{RwLock, mpsc};
13use tracing::{debug, info, warn};
14use uuid::Uuid;
15
16use crate::safety::{HistoryOperation, SafetyManager, UndoData};
17
18#[derive(Debug, Clone, Serialize, Deserialize)]
20#[serde(tag = "op", rename_all = "snake_case")]
21pub enum Operation {
22 Create { path: PathBuf, content: String },
24 Delete { path: PathBuf },
26 Move { src: PathBuf, dst: PathBuf },
28 Copy { src: PathBuf, dst: PathBuf },
30 Write {
32 path: PathBuf,
33 content: String,
34 #[serde(default)]
35 append: bool,
36 },
37 Mkdir { path: PathBuf },
39 Symlink { target: PathBuf, link: PathBuf },
41}
42
43#[derive(Debug, Clone, Serialize, Deserialize)]
45pub struct BatchRequest {
46 pub operations: Vec<Operation>,
48 #[serde(default)]
50 pub atomic: bool,
51 #[serde(default)]
53 pub dry_run: bool,
54}
55
56#[derive(Debug, Clone, Serialize, Deserialize)]
58pub struct OperationResult {
59 pub id: Uuid,
61 pub success: bool,
63 pub operation: String,
65 pub path: PathBuf,
67 #[serde(skip_serializing_if = "Option::is_none")]
69 pub error: Option<String>,
70 pub timestamp: DateTime<Utc>,
72 pub indexed: bool,
74 #[serde(skip_serializing_if = "Option::is_none")]
76 pub undo_id: Option<Uuid>,
77 #[serde(skip)]
79 trash_id: Option<Uuid>,
80}
81
82impl OperationResult {
83 pub fn success(operation: &str, path: PathBuf, indexed: bool) -> Self {
85 Self {
86 id: Uuid::new_v4(),
87 success: true,
88 operation: operation.to_string(),
89 path,
90 error: None,
91 timestamp: Utc::now(),
92 indexed,
93 undo_id: None,
94 trash_id: None,
95 }
96 }
97
98 pub fn failure(operation: &str, path: PathBuf, error: String) -> Self {
100 Self {
101 id: Uuid::new_v4(),
102 success: false,
103 operation: operation.to_string(),
104 path,
105 error: Some(error),
106 timestamp: Utc::now(),
107 indexed: false,
108 undo_id: None,
109 trash_id: None,
110 }
111 }
112}
113
114#[derive(Debug, Clone, Serialize, Deserialize)]
116pub struct BatchResult {
117 pub id: Uuid,
119 pub success: bool,
121 pub total: usize,
123 pub succeeded: usize,
125 pub failed: usize,
127 pub results: Vec<OperationResult>,
129 pub timestamp: DateTime<Utc>,
131 #[serde(skip_serializing_if = "Option::is_none")]
133 pub rollback_performed: Option<bool>,
134 #[serde(skip_serializing_if = "Option::is_none")]
136 pub rollback_details: Option<RollbackDetails>,
137}
138
139#[derive(Debug, Clone, Serialize, Deserialize)]
141pub struct RollbackDetails {
142 pub rolled_back: usize,
144 pub rollback_failures: usize,
146 pub errors: Vec<RollbackError>,
148}
149
150#[derive(Debug, Clone, Serialize, Deserialize)]
152pub struct RollbackError {
153 pub operation_index: usize,
155 pub error: String,
157}
158
159#[derive(Debug, Clone)]
162enum RollbackData {
163 Create { created_path: PathBuf },
165 Delete {
167 original_path: PathBuf,
168 trash_id: Option<Uuid>,
169 content_backup: Option<Vec<u8>>,
170 },
171 Move { src: PathBuf, dst: PathBuf },
173 Copy { copied_path: PathBuf },
175 Write {
177 path: PathBuf,
178 previous_content: Option<Vec<u8>>,
179 file_existed: bool,
180 },
181 Mkdir { created_path: PathBuf },
183 Symlink { link_path: PathBuf },
185}
186
187#[derive(Debug, Clone)]
189struct JournalEntry {
190 operation_index: usize,
192 rollback_data: RollbackData,
194}
195
196pub struct OpsManager {
198 source: PathBuf,
200 store: Option<Arc<dyn VectorStore>>,
202 reindex_sender: Option<mpsc::Sender<PathBuf>>,
204 safety_manager: Option<Arc<SafetyManager>>,
206 last_result: Arc<RwLock<Option<OperationResult>>>,
208 last_batch_result: Arc<RwLock<Option<BatchResult>>>,
210}
211
212impl OpsManager {
213 pub fn new(
215 source: PathBuf,
216 store: Option<Arc<dyn VectorStore>>,
217 reindex_sender: Option<mpsc::Sender<PathBuf>>,
218 ) -> Self {
219 Self {
220 source,
221 store,
222 reindex_sender,
223 safety_manager: None,
224 last_result: Arc::new(RwLock::new(None)),
225 last_batch_result: Arc::new(RwLock::new(None)),
226 }
227 }
228
229 pub fn with_safety(
231 source: PathBuf,
232 store: Option<Arc<dyn VectorStore>>,
233 reindex_sender: Option<mpsc::Sender<PathBuf>>,
234 safety_manager: Arc<SafetyManager>,
235 ) -> Self {
236 Self {
237 source,
238 store,
239 reindex_sender,
240 safety_manager: Some(safety_manager),
241 last_result: Arc::new(RwLock::new(None)),
242 last_batch_result: Arc::new(RwLock::new(None)),
243 }
244 }
245
246 pub fn set_safety_manager(&mut self, safety_manager: Arc<SafetyManager>) {
248 self.safety_manager = Some(safety_manager);
249 }
250
251 fn log_to_history(
254 &self,
255 operation: HistoryOperation,
256 undo_data: Option<UndoData>,
257 ) -> Option<Uuid> {
258 self.safety_manager
259 .as_ref()
260 .map(|safety| safety.log_success(operation, undo_data))
261 }
262
263 fn log_failure_to_history(&self, operation: HistoryOperation, error: &str) {
265 if let Some(ref safety) = self.safety_manager {
266 safety.log_failure(operation, error.to_string());
267 }
268 }
269
270 fn resolve_path(&self, path: &PathBuf) -> Result<PathBuf, String> {
273 ragfs_core::resolve_under_root(&self.source, path)
274 }
275
276 fn jail_symlink_target(
279 &self,
280 target: &PathBuf,
281 resolved_link: &Path,
282 ) -> Result<PathBuf, String> {
283 let validation_target = if target.is_absolute() {
284 target.clone()
285 } else if let Some(parent) = resolved_link.parent() {
286 parent.join(target)
287 } else {
288 target.clone()
289 };
290 self.resolve_path(&validation_target)
291 }
292
293 async fn fail_and_store(
294 &self,
295 operation: &str,
296 path: PathBuf,
297 error: String,
298 ) -> OperationResult {
299 let result = OperationResult::failure(operation, path, error);
300 *self.last_result.write().await = Some(result.clone());
301 result
302 }
303
304 async fn trigger_reindex(&self, path: &PathBuf) -> bool {
306 if let Some(ref sender) = self.reindex_sender {
307 match sender.send(path.clone()).await {
308 Ok(()) => {
309 debug!("Reindex triggered for {:?}", path);
310 true
311 }
312 Err(e) => {
313 warn!("Failed to trigger reindex for {:?}: {}", path, e);
314 false
315 }
316 }
317 } else {
318 false
319 }
320 }
321
322 async fn delete_from_store(&self, path: &PathBuf) {
324 if let Some(ref store) = self.store
325 && let Err(e) = store.delete_by_file_path(path).await
326 {
327 warn!("Failed to delete {:?} from store: {e}", path);
328 }
329 }
330
331 async fn update_store_path(&self, from: &PathBuf, to: &PathBuf) {
333 if let Some(ref store) = self.store
334 && let Err(e) = store.update_file_path(from, to).await
335 {
336 warn!("Failed to update path in store {:?} -> {:?}: {e}", from, to);
337 }
338 }
339
340 async fn execute_with_rollback(
343 &self,
344 op: &Operation,
345 ) -> (OperationResult, Option<RollbackData>) {
346 match op {
347 Operation::Create { path, content } => {
348 let result = self.create(path, content).await;
349 let rollback = if result.success {
350 self.resolve_path(path)
351 .ok()
352 .map(|created_path| RollbackData::Create { created_path })
353 } else {
354 None
355 };
356 (result, rollback)
357 }
358 Operation::Delete { path } => {
359 let resolved = match self.resolve_path(path) {
360 Ok(p) => p,
361 Err(_) => return (self.delete(path).await, None),
362 };
363 let content_backup = if resolved.exists() && resolved.is_file() {
365 fs::read(&resolved).ok()
366 } else {
367 None
368 };
369
370 let result = self.delete(path).await;
371 let rollback = if result.success {
372 Some(RollbackData::Delete {
373 original_path: resolved,
374 trash_id: result.trash_id,
375 content_backup: if result.trash_id.is_none() {
376 content_backup
377 } else {
378 None
379 },
380 })
381 } else {
382 None
383 };
384 (result, rollback)
385 }
386 Operation::Move { src, dst } => {
387 let result = self.move_file(src, dst).await;
388 let rollback = if result.success {
389 match (self.resolve_path(src), self.resolve_path(dst)) {
390 (Ok(resolved_src), Ok(resolved_dst)) => Some(RollbackData::Move {
391 src: resolved_dst,
392 dst: resolved_src,
393 }),
394 _ => None,
395 }
396 } else {
397 None
398 };
399 (result, rollback)
400 }
401 Operation::Copy { src, dst } => {
402 let result = self.copy(src, dst).await;
403 let rollback = if result.success {
404 self.resolve_path(dst)
405 .ok()
406 .map(|copied_path| RollbackData::Copy { copied_path })
407 } else {
408 None
409 };
410 (result, rollback)
411 }
412 Operation::Write {
413 path,
414 content,
415 append,
416 } => {
417 let (file_existed, previous_content, resolved) = match self.resolve_path(path) {
418 Ok(resolved) => {
419 let file_existed = resolved.exists();
420 let previous_content = if file_existed {
421 fs::read(&resolved).ok()
422 } else {
423 None
424 };
425 (file_existed, previous_content, Some(resolved))
426 }
427 Err(_) => (false, None, None),
428 };
429
430 let result = self.write(path, content, *append).await;
431 let rollback = if result.success {
432 resolved.map(|path| RollbackData::Write {
433 path,
434 previous_content,
435 file_existed,
436 })
437 } else {
438 None
439 };
440 (result, rollback)
441 }
442 Operation::Mkdir { path } => {
443 let result = self.mkdir(path).await;
444 let rollback = if result.success {
445 self.resolve_path(path)
446 .ok()
447 .map(|created_path| RollbackData::Mkdir { created_path })
448 } else {
449 None
450 };
451 (result, rollback)
452 }
453 Operation::Symlink { target, link } => {
454 let result = self.symlink(target, link).await;
455 let rollback = if result.success {
456 self.resolve_path(link)
457 .ok()
458 .map(|link_path| RollbackData::Symlink { link_path })
459 } else {
460 None
461 };
462 (result, rollback)
463 }
464 }
465 }
466
467 async fn rollback_operation(&self, rollback_data: &RollbackData) -> Result<(), String> {
469 match rollback_data {
470 RollbackData::Create { created_path } => {
471 if created_path.exists() {
473 fs::remove_file(created_path)
474 .map_err(|e| format!("Failed to rollback create: {e}"))?;
475 self.delete_from_store(created_path).await;
477 }
478 Ok(())
479 }
480 RollbackData::Delete {
481 original_path,
482 trash_id,
483 content_backup,
484 } => {
485 if let Some(id) = trash_id
487 && let Some(ref safety) = self.safety_manager
488 {
489 return safety
490 .restore(*id)
491 .await
492 .map(|_| ())
493 .map_err(|e| format!("Failed to restore from trash: {e}"));
494 }
495 if let Some(content) = content_backup {
497 if let Some(parent) = original_path.parent() {
498 fs::create_dir_all(parent)
499 .map_err(|e| format!("Failed to create parent dir: {e}"))?;
500 }
501 fs::write(original_path, content)
502 .map_err(|e| format!("Failed to restore content: {e}"))?;
503 self.trigger_reindex(original_path).await;
504 Ok(())
505 } else {
506 Err("Cannot rollback delete: no trash entry or content backup".into())
507 }
508 }
509 RollbackData::Move { src, dst } => {
510 if src.exists() {
512 if let Some(parent) = dst.parent() {
513 fs::create_dir_all(parent)
514 .map_err(|e| format!("Failed to create parent dir: {e}"))?;
515 }
516 fs::rename(src, dst).map_err(|e| format!("Failed to rollback move: {e}"))?;
517 self.update_store_path(src, dst).await;
518 }
519 Ok(())
520 }
521 RollbackData::Copy { copied_path } => {
522 if copied_path.exists() {
524 fs::remove_file(copied_path)
525 .map_err(|e| format!("Failed to rollback copy: {e}"))?;
526 self.delete_from_store(copied_path).await;
527 }
528 Ok(())
529 }
530 RollbackData::Write {
531 path,
532 previous_content,
533 file_existed,
534 } => {
535 if *file_existed {
536 if let Some(content) = previous_content {
537 fs::write(path, content)
538 .map_err(|e| format!("Failed to rollback write: {e}"))?;
539 } else {
540 return Err("Cannot rollback write: no previous content saved".into());
541 }
542 } else {
543 if path.exists() {
545 fs::remove_file(path)
546 .map_err(|e| format!("Failed to rollback write (delete): {e}"))?;
547 self.delete_from_store(path).await;
548 }
549 }
550 self.trigger_reindex(path).await;
551 Ok(())
552 }
553 RollbackData::Mkdir { created_path } => {
554 if created_path.exists() && created_path.is_dir() {
556 fs::remove_dir(created_path)
557 .map_err(|e| format!("Failed to rollback mkdir: {e}"))?;
558 }
559 Ok(())
560 }
561 RollbackData::Symlink { link_path } => {
562 if link_path.exists() || link_path.symlink_metadata().is_ok() {
564 fs::remove_file(link_path)
565 .map_err(|e| format!("Failed to rollback symlink: {e}"))?;
566 }
567 Ok(())
568 }
569 }
570 }
571
572 async fn perform_rollback(&self, journal: &[JournalEntry]) -> RollbackDetails {
574 let mut rolled_back = 0;
575 let mut rollback_failures = 0;
576 let mut errors = Vec::new();
577
578 for entry in journal.iter().rev() {
580 match self.rollback_operation(&entry.rollback_data).await {
581 Ok(()) => {
582 rolled_back += 1;
583 info!("Rolled back operation index {}", entry.operation_index);
584 }
585 Err(e) => {
586 rollback_failures += 1;
587 warn!(
588 "Failed to rollback operation index {}: {}",
589 entry.operation_index, e
590 );
591 errors.push(RollbackError {
592 operation_index: entry.operation_index,
593 error: e,
594 });
595 }
596 }
597 }
598
599 RollbackDetails {
600 rolled_back,
601 rollback_failures,
602 errors,
603 }
604 }
605
606 pub async fn create(&self, path: &PathBuf, content: &str) -> OperationResult {
608 let resolved = match self.resolve_path(path) {
609 Ok(p) => p,
610 Err(e) => return self.fail_and_store("create", path.clone(), e).await,
611 };
612 debug!("ops::create {:?}", resolved);
613
614 if resolved.exists() {
616 return OperationResult::failure("create", path.clone(), "File already exists".into());
617 }
618
619 if let Some(parent) = resolved.parent()
621 && !parent.exists()
622 && let Err(e) = fs::create_dir_all(parent)
623 {
624 return OperationResult::failure(
625 "create",
626 path.clone(),
627 format!("Failed to create parent directory: {e}"),
628 );
629 }
630
631 match fs::write(&resolved, content) {
633 Ok(()) => {
634 info!("Created file: {:?}", resolved);
635 let indexed = self.trigger_reindex(&resolved).await;
636 let undo_id = self.log_to_history(
637 HistoryOperation::Create {
638 path: resolved.clone(),
639 },
640 Some(UndoData::Create { path: resolved }),
641 );
642 let mut result = OperationResult::success("create", path.clone(), indexed);
643 result.undo_id = undo_id;
644
645 *self.last_result.write().await = Some(result.clone());
646 result
647 }
648 Err(e) => {
649 let error_msg = format!("Failed to create file: {e}");
650 self.log_failure_to_history(
651 HistoryOperation::Create { path: resolved },
652 &error_msg,
653 );
654
655 let result = OperationResult::failure("create", path.clone(), error_msg);
656 *self.last_result.write().await = Some(result.clone());
657 result
658 }
659 }
660 }
661
662 pub async fn delete(&self, path: &PathBuf) -> OperationResult {
665 let resolved = match self.resolve_path(path) {
666 Ok(p) => p,
667 Err(e) => return self.fail_and_store("delete", path.clone(), e).await,
668 };
669 debug!("ops::delete {:?}", resolved);
670
671 if !resolved.exists() {
672 return OperationResult::failure("delete", path.clone(), "File not found".into());
673 }
674
675 if resolved.is_dir() {
676 return OperationResult::failure(
677 "delete",
678 path.clone(),
679 "Cannot delete directory with .delete, use rmdir".into(),
680 );
681 }
682
683 self.delete_from_store(&resolved).await;
685
686 if let Some(ref safety) = self.safety_manager {
688 match safety.soft_delete(&resolved).await {
689 Ok(entry) => {
690 info!("Soft deleted file: {:?} -> trash/{}", resolved, entry.id);
691
692 let undo_id = self.log_to_history(
693 HistoryOperation::Delete {
694 path: resolved,
695 trash_id: Some(entry.id),
696 },
697 Some(UndoData::Delete { trash_id: entry.id }),
698 );
699
700 let mut result = OperationResult::success("delete", path.clone(), false);
701 result.undo_id = undo_id;
702 result.trash_id = Some(entry.id);
703 *self.last_result.write().await = Some(result.clone());
704 return result;
705 }
706 Err(e) => {
707 warn!("Soft delete failed, falling back to hard delete: {e}");
708 }
709 }
710 }
711
712 match fs::remove_file(&resolved) {
714 Ok(()) => {
715 info!("Hard deleted file: {:?}", resolved);
716
717 let undo_id = self.log_to_history(
718 HistoryOperation::Delete {
719 path: resolved,
720 trash_id: None,
721 },
722 None, );
724
725 let mut result = OperationResult::success("delete", path.clone(), false);
726 result.undo_id = undo_id;
727 *self.last_result.write().await = Some(result.clone());
728 result
729 }
730 Err(e) => {
731 let error_msg = format!("Failed to delete file: {e}");
732 self.log_failure_to_history(
733 HistoryOperation::Delete {
734 path: resolved,
735 trash_id: None,
736 },
737 &error_msg,
738 );
739
740 let result = OperationResult::failure("delete", path.clone(), error_msg);
741 *self.last_result.write().await = Some(result.clone());
742 result
743 }
744 }
745 }
746
747 pub async fn move_file(&self, src: &PathBuf, dst: &PathBuf) -> OperationResult {
749 let resolved_src = match self.resolve_path(src) {
750 Ok(p) => p,
751 Err(e) => return self.fail_and_store("move", src.clone(), e).await,
752 };
753 let resolved_dst = match self.resolve_path(dst) {
754 Ok(p) => p,
755 Err(e) => return self.fail_and_store("move", src.clone(), e).await,
756 };
757 debug!("ops::move {:?} -> {:?}", resolved_src, resolved_dst);
758
759 if !resolved_src.exists() {
760 return OperationResult::failure("move", src.clone(), "Source file not found".into());
761 }
762
763 if resolved_dst.exists() {
764 return OperationResult::failure(
765 "move",
766 src.clone(),
767 "Destination already exists".into(),
768 );
769 }
770
771 if let Some(parent) = resolved_dst.parent()
773 && !parent.exists()
774 && let Err(e) = fs::create_dir_all(parent)
775 {
776 return OperationResult::failure(
777 "move",
778 src.clone(),
779 format!("Failed to create parent directory: {e}"),
780 );
781 }
782
783 match fs::rename(&resolved_src, &resolved_dst) {
785 Ok(()) => {
786 info!("Moved: {:?} -> {:?}", resolved_src, resolved_dst);
787 self.update_store_path(&resolved_src, &resolved_dst).await;
788
789 let undo_id = self.log_to_history(
790 HistoryOperation::Move {
791 src: resolved_src.clone(),
792 dst: resolved_dst.clone(),
793 },
794 Some(UndoData::Move {
795 src: resolved_src,
796 dst: resolved_dst,
797 }),
798 );
799
800 let mut result = OperationResult::success("move", dst.clone(), true);
801 result.undo_id = undo_id;
802 *self.last_result.write().await = Some(result.clone());
803 result
804 }
805 Err(e) => {
806 let error_msg = format!("Failed to move file: {e}");
807 self.log_failure_to_history(
808 HistoryOperation::Move {
809 src: resolved_src,
810 dst: resolved_dst,
811 },
812 &error_msg,
813 );
814
815 let result = OperationResult::failure("move", src.clone(), error_msg);
816 *self.last_result.write().await = Some(result.clone());
817 result
818 }
819 }
820 }
821
822 pub async fn copy(&self, src: &PathBuf, dst: &PathBuf) -> OperationResult {
824 let resolved_src = match self.resolve_path(src) {
825 Ok(p) => p,
826 Err(e) => return self.fail_and_store("copy", src.clone(), e).await,
827 };
828 let resolved_dst = match self.resolve_path(dst) {
829 Ok(p) => p,
830 Err(e) => return self.fail_and_store("copy", src.clone(), e).await,
831 };
832 debug!("ops::copy {:?} -> {:?}", resolved_src, resolved_dst);
833
834 if !resolved_src.exists() {
835 return OperationResult::failure("copy", src.clone(), "Source file not found".into());
836 }
837
838 if resolved_dst.exists() {
839 return OperationResult::failure(
840 "copy",
841 src.clone(),
842 "Destination already exists".into(),
843 );
844 }
845
846 if let Some(parent) = resolved_dst.parent()
848 && !parent.exists()
849 && let Err(e) = fs::create_dir_all(parent)
850 {
851 return OperationResult::failure(
852 "copy",
853 src.clone(),
854 format!("Failed to create parent directory: {e}"),
855 );
856 }
857
858 match fs::copy(&resolved_src, &resolved_dst) {
860 Ok(_) => {
861 info!("Copied: {:?} -> {:?}", resolved_src, resolved_dst);
862 let indexed = self.trigger_reindex(&resolved_dst).await;
863
864 let undo_id = self.log_to_history(
865 HistoryOperation::Copy {
866 src: resolved_src,
867 dst: resolved_dst.clone(),
868 },
869 Some(UndoData::Copy { path: resolved_dst }),
870 );
871
872 let mut result = OperationResult::success("copy", dst.clone(), indexed);
873 result.undo_id = undo_id;
874 *self.last_result.write().await = Some(result.clone());
875 result
876 }
877 Err(e) => {
878 let error_msg = format!("Failed to copy file: {e}");
879 self.log_failure_to_history(
880 HistoryOperation::Copy {
881 src: resolved_src,
882 dst: resolved_dst,
883 },
884 &error_msg,
885 );
886
887 let result = OperationResult::failure("copy", src.clone(), error_msg);
888 *self.last_result.write().await = Some(result.clone());
889 result
890 }
891 }
892 }
893
894 pub async fn write(&self, path: &PathBuf, content: &str, append: bool) -> OperationResult {
896 let resolved = match self.resolve_path(path) {
897 Ok(p) => p,
898 Err(e) => return self.fail_and_store("write", path.clone(), e).await,
899 };
900 debug!("ops::write {:?} (append={})", resolved, append);
901
902 let write_result = if append {
903 use std::io::Write;
904 fs::OpenOptions::new()
905 .create(true)
906 .append(true)
907 .open(&resolved)
908 .and_then(|mut f| f.write_all(content.as_bytes()))
909 } else {
910 fs::write(&resolved, content)
911 };
912
913 match write_result {
914 Ok(()) => {
915 info!("Wrote to file: {:?}", resolved);
916 let indexed = self.trigger_reindex(&resolved).await;
917
918 let undo_id = self.log_to_history(
919 HistoryOperation::Write {
920 path: resolved,
921 append,
922 },
923 None, );
925
926 let mut result = OperationResult::success("write", path.clone(), indexed);
927 result.undo_id = undo_id;
928 *self.last_result.write().await = Some(result.clone());
929 result
930 }
931 Err(e) => {
932 let error_msg = format!("Failed to write file: {e}");
933 self.log_failure_to_history(
934 HistoryOperation::Write {
935 path: resolved,
936 append,
937 },
938 &error_msg,
939 );
940
941 let result = OperationResult::failure("write", path.clone(), error_msg);
942 *self.last_result.write().await = Some(result.clone());
943 result
944 }
945 }
946 }
947
948 pub async fn mkdir(&self, path: &PathBuf) -> OperationResult {
950 let resolved = match self.resolve_path(path) {
951 Ok(p) => p,
952 Err(e) => return self.fail_and_store("mkdir", path.clone(), e).await,
953 };
954 debug!("ops::mkdir {:?}", resolved);
955
956 if resolved.exists() {
957 return OperationResult::failure("mkdir", path.clone(), "Path already exists".into());
958 }
959
960 match fs::create_dir_all(&resolved) {
961 Ok(()) => {
962 info!("Created directory: {:?}", resolved);
963
964 let result = OperationResult::success("mkdir", path.clone(), false);
966
967 *self.last_result.write().await = Some(result.clone());
968 result
969 }
970 Err(e) => {
971 let error_msg = format!("Failed to create directory: {e}");
972 let result = OperationResult::failure("mkdir", path.clone(), error_msg);
973 *self.last_result.write().await = Some(result.clone());
974 result
975 }
976 }
977 }
978
979 #[cfg(unix)]
981 pub async fn symlink(&self, target: &PathBuf, link: &PathBuf) -> OperationResult {
982 let resolved_link = match self.resolve_path(link) {
983 Ok(p) => p,
984 Err(e) => return self.fail_and_store("symlink", link.clone(), e).await,
985 };
986 if let Err(e) = self.jail_symlink_target(target, &resolved_link) {
987 return self.fail_and_store("symlink", link.clone(), e).await;
988 }
989 debug!("ops::symlink {:?} -> {:?}", resolved_link, target);
990
991 if resolved_link.exists() {
992 return OperationResult::failure(
993 "symlink",
994 link.clone(),
995 "Link path already exists".into(),
996 );
997 }
998
999 if let Some(parent) = resolved_link.parent()
1001 && !parent.exists()
1002 && let Err(e) = fs::create_dir_all(parent)
1003 {
1004 return OperationResult::failure(
1005 "symlink",
1006 link.clone(),
1007 format!("Failed to create parent directory: {e}"),
1008 );
1009 }
1010
1011 match std::os::unix::fs::symlink(target, &resolved_link) {
1013 Ok(()) => {
1014 info!("Created symlink: {:?} -> {:?}", resolved_link, target);
1015
1016 let result = OperationResult::success("symlink", link.clone(), false);
1017
1018 *self.last_result.write().await = Some(result.clone());
1019 result
1020 }
1021 Err(e) => {
1022 let error_msg = format!("Failed to create symlink: {e}");
1023 let result = OperationResult::failure("symlink", link.clone(), error_msg);
1024 *self.last_result.write().await = Some(result.clone());
1025 result
1026 }
1027 }
1028 }
1029
1030 #[cfg(not(unix))]
1032 pub async fn symlink(&self, _target: &PathBuf, link: &PathBuf) -> OperationResult {
1033 OperationResult::failure(
1034 "symlink",
1035 link.clone(),
1036 "Symlinks are not supported on this platform".into(),
1037 )
1038 }
1039
1040 pub async fn execute_operation(&self, op: &Operation) -> OperationResult {
1042 match op {
1043 Operation::Create { path, content } => self.create(path, content).await,
1044 Operation::Delete { path } => self.delete(path).await,
1045 Operation::Move { src, dst } => self.move_file(src, dst).await,
1046 Operation::Copy { src, dst } => self.copy(src, dst).await,
1047 Operation::Write {
1048 path,
1049 content,
1050 append,
1051 } => self.write(path, content, *append).await,
1052 Operation::Mkdir { path } => self.mkdir(path).await,
1053 Operation::Symlink { target, link } => self.symlink(target, link).await,
1054 }
1055 }
1056
1057 pub async fn batch(&self, request: BatchRequest) -> BatchResult {
1059 let batch_id = Uuid::new_v4();
1060 let total = request.operations.len();
1061 let mut results = Vec::with_capacity(total);
1062 let mut succeeded = 0;
1063 let mut failed = 0;
1064 let mut journal: Vec<JournalEntry> = Vec::new();
1065
1066 debug!(
1067 "ops::batch {} operations (atomic={}, dry_run={})",
1068 total, request.atomic, request.dry_run
1069 );
1070
1071 if request.dry_run {
1072 for op in &request.operations {
1074 let result = self.validate_operation(op);
1075 if result.success {
1076 succeeded += 1;
1077 } else {
1078 failed += 1;
1079 }
1080 results.push(result);
1081 }
1082
1083 let batch_result = BatchResult {
1084 id: batch_id,
1085 success: failed == 0,
1086 total,
1087 succeeded,
1088 failed,
1089 results,
1090 timestamp: Utc::now(),
1091 rollback_performed: None,
1092 rollback_details: None,
1093 };
1094
1095 *self.last_batch_result.write().await = Some(batch_result.clone());
1096 return batch_result;
1097 }
1098
1099 for (index, op) in request.operations.iter().enumerate() {
1101 let (result, rollback_data) = self.execute_with_rollback(op).await;
1102
1103 if result.success {
1104 succeeded += 1;
1105 if request.atomic
1107 && let Some(rd) = rollback_data
1108 {
1109 journal.push(JournalEntry {
1110 operation_index: index,
1111 rollback_data: rd,
1112 });
1113 }
1114 results.push(result);
1115 } else {
1116 failed += 1;
1117 results.push(result);
1118
1119 if request.atomic && !journal.is_empty() {
1120 info!(
1122 "Batch failed at operation {}, rolling back {} operations",
1123 index,
1124 journal.len()
1125 );
1126 let rollback_details = self.perform_rollback(&journal).await;
1127 let rollback_success = rollback_details.rollback_failures == 0;
1128
1129 let batch_result = BatchResult {
1130 id: batch_id,
1131 success: false,
1132 total,
1133 succeeded,
1134 failed,
1135 results,
1136 timestamp: Utc::now(),
1137 rollback_performed: Some(rollback_success),
1138 rollback_details: Some(rollback_details),
1139 };
1140
1141 *self.last_batch_result.write().await = Some(batch_result.clone());
1142 return batch_result;
1143 } else if request.atomic {
1144 let batch_result = BatchResult {
1146 id: batch_id,
1147 success: false,
1148 total,
1149 succeeded,
1150 failed,
1151 results,
1152 timestamp: Utc::now(),
1153 rollback_performed: None,
1154 rollback_details: None,
1155 };
1156
1157 *self.last_batch_result.write().await = Some(batch_result.clone());
1158 return batch_result;
1159 }
1160 }
1162 }
1163
1164 let batch_result = BatchResult {
1165 id: batch_id,
1166 success: failed == 0,
1167 total,
1168 succeeded,
1169 failed,
1170 results,
1171 timestamp: Utc::now(),
1172 rollback_performed: None,
1173 rollback_details: None,
1174 };
1175
1176 *self.last_batch_result.write().await = Some(batch_result.clone());
1177 batch_result
1178 }
1179
1180 fn validate_operation(&self, op: &Operation) -> OperationResult {
1182 match op {
1183 Operation::Create { path, .. } => {
1184 let resolved = match self.resolve_path(path) {
1185 Ok(p) => p,
1186 Err(e) => return OperationResult::failure("create", path.clone(), e),
1187 };
1188 if resolved.exists() {
1189 OperationResult::failure("create", path.clone(), "File already exists".into())
1190 } else {
1191 OperationResult::success("create", path.clone(), false)
1192 }
1193 }
1194 Operation::Delete { path } => {
1195 let resolved = match self.resolve_path(path) {
1196 Ok(p) => p,
1197 Err(e) => return OperationResult::failure("delete", path.clone(), e),
1198 };
1199 if !resolved.exists() {
1200 OperationResult::failure("delete", path.clone(), "File not found".into())
1201 } else if resolved.is_dir() {
1202 OperationResult::failure(
1203 "delete",
1204 path.clone(),
1205 "Cannot delete directory".into(),
1206 )
1207 } else {
1208 OperationResult::success("delete", path.clone(), false)
1209 }
1210 }
1211 Operation::Move { src, dst } => {
1212 let resolved_src = match self.resolve_path(src) {
1213 Ok(p) => p,
1214 Err(e) => return OperationResult::failure("move", src.clone(), e),
1215 };
1216 let resolved_dst = match self.resolve_path(dst) {
1217 Ok(p) => p,
1218 Err(e) => return OperationResult::failure("move", src.clone(), e),
1219 };
1220 if !resolved_src.exists() {
1221 OperationResult::failure("move", src.clone(), "Source not found".into())
1222 } else if resolved_dst.exists() {
1223 OperationResult::failure(
1224 "move",
1225 src.clone(),
1226 "Destination already exists".into(),
1227 )
1228 } else {
1229 OperationResult::success("move", dst.clone(), false)
1230 }
1231 }
1232 Operation::Copy { src, dst } => {
1233 let resolved_src = match self.resolve_path(src) {
1234 Ok(p) => p,
1235 Err(e) => return OperationResult::failure("copy", src.clone(), e),
1236 };
1237 let resolved_dst = match self.resolve_path(dst) {
1238 Ok(p) => p,
1239 Err(e) => return OperationResult::failure("copy", src.clone(), e),
1240 };
1241 if !resolved_src.exists() {
1242 OperationResult::failure("copy", src.clone(), "Source not found".into())
1243 } else if resolved_dst.exists() {
1244 OperationResult::failure(
1245 "copy",
1246 src.clone(),
1247 "Destination already exists".into(),
1248 )
1249 } else {
1250 OperationResult::success("copy", dst.clone(), false)
1251 }
1252 }
1253 Operation::Write { path, .. } => match self.resolve_path(path) {
1254 Ok(_) => OperationResult::success("write", path.clone(), false),
1255 Err(e) => OperationResult::failure("write", path.clone(), e),
1256 },
1257 Operation::Mkdir { path } => {
1258 let resolved = match self.resolve_path(path) {
1259 Ok(p) => p,
1260 Err(e) => return OperationResult::failure("mkdir", path.clone(), e),
1261 };
1262 if resolved.exists() {
1263 OperationResult::failure("mkdir", path.clone(), "Path already exists".into())
1264 } else {
1265 OperationResult::success("mkdir", path.clone(), false)
1266 }
1267 }
1268 Operation::Symlink { target, link } => {
1269 let resolved_link = match self.resolve_path(link) {
1270 Ok(p) => p,
1271 Err(e) => return OperationResult::failure("symlink", link.clone(), e),
1272 };
1273 if let Err(e) = self.jail_symlink_target(target, &resolved_link) {
1274 return OperationResult::failure("symlink", link.clone(), e);
1275 }
1276 if resolved_link.exists() {
1277 OperationResult::failure(
1278 "symlink",
1279 link.clone(),
1280 "Link path already exists".into(),
1281 )
1282 } else {
1283 OperationResult::success("symlink", link.clone(), false)
1284 }
1285 }
1286 }
1287 }
1288
1289 pub async fn get_last_result(&self) -> Vec<u8> {
1291 let result = self.last_result.read().await;
1292 if let Some(r) = &*result {
1293 serde_json::to_string_pretty(r)
1294 .unwrap_or_else(|_| "{}".to_string())
1295 .into_bytes()
1296 } else {
1297 let empty = serde_json::json!({
1298 "message": "No operations performed yet"
1299 });
1300 serde_json::to_string_pretty(&empty)
1301 .unwrap_or_default()
1302 .into_bytes()
1303 }
1304 }
1305
1306 #[allow(dead_code)]
1308 pub async fn get_last_batch_result(&self) -> Vec<u8> {
1309 let result = self.last_batch_result.read().await;
1310 if let Some(r) = &*result {
1311 serde_json::to_string_pretty(r)
1312 .unwrap_or_else(|_| "{}".to_string())
1313 .into_bytes()
1314 } else {
1315 let empty = serde_json::json!({
1316 "message": "No batch operations performed yet"
1317 });
1318 serde_json::to_string_pretty(&empty)
1319 .unwrap_or_default()
1320 .into_bytes()
1321 }
1322 }
1323
1324 pub async fn parse_and_create(&self, input: &str) -> OperationResult {
1327 let parts: Vec<&str> = input.splitn(2, '\n').collect();
1328 if parts.len() < 2 {
1329 return OperationResult::failure(
1330 "create",
1331 PathBuf::new(),
1332 "Invalid format. Expected: path\\ncontent".into(),
1333 );
1334 }
1335 let path = PathBuf::from(parts[0].trim());
1336 let content = parts[1];
1337 self.create(&path, content).await
1338 }
1339
1340 pub async fn parse_and_delete(&self, input: &str) -> OperationResult {
1343 let path = PathBuf::from(input.trim());
1344 if path.as_os_str().is_empty() {
1345 return OperationResult::failure(
1346 "delete",
1347 PathBuf::new(),
1348 "Invalid format. Expected: path".into(),
1349 );
1350 }
1351 self.delete(&path).await
1352 }
1353
1354 pub async fn parse_and_move(&self, input: &str) -> OperationResult {
1357 let parts: Vec<&str> = input.splitn(2, '\n').collect();
1358 if parts.len() < 2 {
1359 return OperationResult::failure(
1360 "move",
1361 PathBuf::new(),
1362 "Invalid format. Expected: src\\ndst".into(),
1363 );
1364 }
1365 let src = PathBuf::from(parts[0].trim());
1366 let dst = PathBuf::from(parts[1].trim());
1367 self.move_file(&src, &dst).await
1368 }
1369
1370 pub async fn parse_and_batch(&self, input: &str) -> BatchResult {
1372 match serde_json::from_str::<BatchRequest>(input) {
1373 Ok(request) => self.batch(request).await,
1374 Err(e) => BatchResult {
1375 id: Uuid::new_v4(),
1376 success: false,
1377 total: 0,
1378 succeeded: 0,
1379 failed: 1,
1380 results: vec![OperationResult::failure(
1381 "batch",
1382 PathBuf::new(),
1383 format!("Invalid JSON: {e}"),
1384 )],
1385 timestamp: Utc::now(),
1386 rollback_performed: None,
1387 rollback_details: None,
1388 },
1389 }
1390 }
1391}
1392
1393#[cfg(test)]
1394mod tests {
1395 use super::*;
1396 use tempfile::TempDir;
1397
1398 fn create_test_manager() -> (OpsManager, TempDir) {
1399 let temp_dir = TempDir::new().unwrap();
1400 let manager = OpsManager::new(temp_dir.path().to_path_buf(), None, None);
1401 (manager, temp_dir)
1402 }
1403
1404 #[tokio::test]
1405 async fn test_create_file() {
1406 let (manager, _temp) = create_test_manager();
1407 let path = PathBuf::from("test.txt");
1408
1409 let result = manager.create(&path, "Hello, World!").await;
1410
1411 assert!(result.success);
1412 assert_eq!(result.operation, "create");
1413 assert!(manager.resolve_path(&path).unwrap().exists());
1414 }
1415
1416 #[tokio::test]
1417 async fn test_create_file_already_exists() {
1418 let (manager, temp) = create_test_manager();
1419 let path = PathBuf::from("existing.txt");
1420 fs::write(temp.path().join("existing.txt"), "content").unwrap();
1421
1422 let result = manager.create(&path, "new content").await;
1423
1424 assert!(!result.success);
1425 assert!(result.error.unwrap().contains("already exists"));
1426 }
1427
1428 #[tokio::test]
1429 async fn test_delete_file() {
1430 let (manager, temp) = create_test_manager();
1431 let path = PathBuf::from("to_delete.txt");
1432 fs::write(temp.path().join("to_delete.txt"), "content").unwrap();
1433
1434 let result = manager.delete(&path).await;
1435
1436 assert!(result.success);
1437 assert!(!manager.resolve_path(&path).unwrap().exists());
1438 }
1439
1440 #[tokio::test]
1441 async fn test_delete_nonexistent() {
1442 let (manager, _temp) = create_test_manager();
1443 let path = PathBuf::from("nonexistent.txt");
1444
1445 let result = manager.delete(&path).await;
1446
1447 assert!(!result.success);
1448 assert!(result.error.unwrap().contains("not found"));
1449 }
1450
1451 #[tokio::test]
1452 async fn test_move_file() {
1453 let (manager, temp) = create_test_manager();
1454 let src = PathBuf::from("source.txt");
1455 let dst = PathBuf::from("dest.txt");
1456 fs::write(temp.path().join("source.txt"), "content").unwrap();
1457
1458 let result = manager.move_file(&src, &dst).await;
1459
1460 assert!(result.success);
1461 assert!(!manager.resolve_path(&src).unwrap().exists());
1462 assert!(manager.resolve_path(&dst).unwrap().exists());
1463 }
1464
1465 #[tokio::test]
1466 async fn test_copy_file() {
1467 let (manager, temp) = create_test_manager();
1468 let src = PathBuf::from("original.txt");
1469 let dst = PathBuf::from("copy.txt");
1470 fs::write(temp.path().join("original.txt"), "content").unwrap();
1471
1472 let result = manager.copy(&src, &dst).await;
1473
1474 assert!(result.success);
1475 assert!(manager.resolve_path(&src).unwrap().exists());
1476 assert!(manager.resolve_path(&dst).unwrap().exists());
1477 }
1478
1479 #[tokio::test]
1480 async fn test_write_file() {
1481 let (manager, _temp) = create_test_manager();
1482 let path = PathBuf::from("write_test.txt");
1483
1484 let result = manager.write(&path, "content", false).await;
1485
1486 assert!(result.success);
1487 let content = fs::read_to_string(manager.resolve_path(&path).unwrap()).unwrap();
1488 assert_eq!(content, "content");
1489 }
1490
1491 #[tokio::test]
1492 async fn test_write_append() {
1493 let (manager, temp) = create_test_manager();
1494 let path = PathBuf::from("append_test.txt");
1495 fs::write(temp.path().join("append_test.txt"), "first").unwrap();
1496
1497 let result = manager.write(&path, " second", true).await;
1498
1499 assert!(result.success);
1500 let content = fs::read_to_string(manager.resolve_path(&path).unwrap()).unwrap();
1501 assert_eq!(content, "first second");
1502 }
1503
1504 #[tokio::test]
1505 async fn test_batch_operations() {
1506 let (manager, _temp) = create_test_manager();
1507
1508 let request = BatchRequest {
1509 operations: vec![
1510 Operation::Create {
1511 path: PathBuf::from("file1.txt"),
1512 content: "content1".to_string(),
1513 },
1514 Operation::Create {
1515 path: PathBuf::from("file2.txt"),
1516 content: "content2".to_string(),
1517 },
1518 ],
1519 atomic: false,
1520 dry_run: false,
1521 };
1522
1523 let result = manager.batch(request).await;
1524
1525 assert!(result.success);
1526 assert_eq!(result.total, 2);
1527 assert_eq!(result.succeeded, 2);
1528 assert_eq!(result.failed, 0);
1529 }
1530
1531 #[tokio::test]
1532 async fn test_batch_dry_run() {
1533 let (manager, _temp) = create_test_manager();
1534
1535 let request = BatchRequest {
1536 operations: vec![Operation::Create {
1537 path: PathBuf::from("dry_run.txt"),
1538 content: "content".to_string(),
1539 }],
1540 atomic: false,
1541 dry_run: true,
1542 };
1543
1544 let result = manager.batch(request).await;
1545
1546 assert!(result.success);
1547 assert!(
1549 !manager
1550 .resolve_path(&PathBuf::from("dry_run.txt"))
1551 .unwrap()
1552 .exists()
1553 );
1554 }
1555
1556 #[tokio::test]
1557 async fn test_parse_and_create() {
1558 let (manager, _temp) = create_test_manager();
1559
1560 let result = manager.parse_and_create("test.txt\nHello!").await;
1561
1562 assert!(result.success);
1563 let content =
1564 fs::read_to_string(manager.resolve_path(&PathBuf::from("test.txt")).unwrap()).unwrap();
1565 assert_eq!(content, "Hello!");
1566 }
1567
1568 #[tokio::test]
1569 async fn test_parse_and_move() {
1570 let (manager, temp) = create_test_manager();
1571 fs::write(temp.path().join("src.txt"), "content").unwrap();
1572
1573 let result = manager.parse_and_move("src.txt\ndst.txt").await;
1574
1575 assert!(result.success);
1576 }
1577
1578 #[tokio::test]
1579 async fn test_parse_and_batch() {
1580 let (manager, _temp) = create_test_manager();
1581
1582 let json = r#"{"operations":[{"op":"create","path":"batch.txt","content":"test"}]}"#;
1583 let result = manager.parse_and_batch(json).await;
1584
1585 assert!(result.success);
1586 }
1587
1588 #[tokio::test]
1589 async fn test_last_result() {
1590 let (manager, _temp) = create_test_manager();
1591
1592 let initial = manager.get_last_result().await;
1594 let initial_str = String::from_utf8(initial).unwrap();
1595 assert!(initial_str.contains("No operations"));
1596
1597 manager.create(&PathBuf::from("test.txt"), "content").await;
1599 let after = manager.get_last_result().await;
1600 let after_str = String::from_utf8(after).unwrap();
1601 assert!(after_str.contains("create"));
1602 assert!(after_str.contains("success"));
1603 }
1604
1605 #[test]
1606 fn test_operation_result_success() {
1607 let result = OperationResult::success("test", PathBuf::from("/test"), true);
1608 assert!(result.success);
1609 assert!(result.error.is_none());
1610 assert!(
1611 result.undo_id.is_none(),
1612 "undo_id is assigned from SafetyManager history, not generated here"
1613 );
1614 }
1615
1616 #[test]
1617 fn test_operation_result_failure() {
1618 let result = OperationResult::failure("test", PathBuf::from("/test"), "error".into());
1619 assert!(!result.success);
1620 assert!(result.error.is_some());
1621 assert!(result.undo_id.is_none());
1622 }
1623
1624 #[tokio::test]
1625 async fn test_atomic_batch_rollback_on_failure() {
1626 let (manager, temp) = create_test_manager();
1627
1628 fs::write(temp.path().join("existing.txt"), "exists").unwrap();
1630
1631 let request = BatchRequest {
1632 operations: vec![
1633 Operation::Create {
1634 path: PathBuf::from("new.txt"),
1635 content: "content".to_string(),
1636 },
1637 Operation::Create {
1638 path: PathBuf::from("existing.txt"), content: "content".to_string(),
1640 },
1641 ],
1642 atomic: true,
1643 dry_run: false,
1644 };
1645
1646 let result = manager.batch(request).await;
1647
1648 assert!(!result.success);
1649 assert_eq!(result.succeeded, 1);
1650 assert_eq!(result.failed, 1);
1651 assert_eq!(result.rollback_performed, Some(true));
1652 assert!(result.rollback_details.is_some());
1653 let details = result.rollback_details.unwrap();
1654 assert_eq!(details.rolled_back, 1);
1655 assert_eq!(details.rollback_failures, 0);
1656 assert!(!temp.path().join("new.txt").exists());
1658 }
1659
1660 #[tokio::test]
1661 async fn test_atomic_batch_rollback_move_operations() {
1662 let (manager, temp) = create_test_manager();
1663
1664 fs::write(temp.path().join("file1.txt"), "content1").unwrap();
1665
1666 let request = BatchRequest {
1667 operations: vec![
1668 Operation::Move {
1669 src: PathBuf::from("file1.txt"),
1670 dst: PathBuf::from("moved1.txt"),
1671 },
1672 Operation::Delete {
1673 path: PathBuf::from("nonexistent.txt"), },
1675 ],
1676 atomic: true,
1677 dry_run: false,
1678 };
1679
1680 let result = manager.batch(request).await;
1681
1682 assert!(!result.success);
1683 assert_eq!(result.rollback_performed, Some(true));
1684 assert!(temp.path().join("file1.txt").exists());
1686 assert!(!temp.path().join("moved1.txt").exists());
1687 }
1688
1689 #[tokio::test]
1690 async fn test_atomic_batch_rollback_write_restores_content() {
1691 let (manager, temp) = create_test_manager();
1692
1693 fs::write(temp.path().join("existing.txt"), "original content").unwrap();
1694
1695 let request = BatchRequest {
1696 operations: vec![
1697 Operation::Write {
1698 path: PathBuf::from("existing.txt"),
1699 content: "modified content".to_string(),
1700 append: false,
1701 },
1702 Operation::Delete {
1703 path: PathBuf::from("nonexistent.txt"), },
1705 ],
1706 atomic: true,
1707 dry_run: false,
1708 };
1709
1710 let result = manager.batch(request).await;
1711
1712 assert!(!result.success);
1713 assert_eq!(result.rollback_performed, Some(true));
1714 let content = fs::read_to_string(temp.path().join("existing.txt")).unwrap();
1716 assert_eq!(content, "original content");
1717 }
1718
1719 #[tokio::test]
1720 async fn test_non_atomic_batch_no_rollback() {
1721 let (manager, temp) = create_test_manager();
1722
1723 let request = BatchRequest {
1724 operations: vec![
1725 Operation::Create {
1726 path: PathBuf::from("file1.txt"),
1727 content: "content".to_string(),
1728 },
1729 Operation::Delete {
1730 path: PathBuf::from("nonexistent.txt"), },
1732 Operation::Create {
1733 path: PathBuf::from("file2.txt"),
1734 content: "content".to_string(),
1735 },
1736 ],
1737 atomic: false, dry_run: false,
1739 };
1740
1741 let result = manager.batch(request).await;
1742
1743 assert!(!result.success);
1744 assert!(result.rollback_performed.is_none());
1745 assert!(temp.path().join("file1.txt").exists());
1747 assert!(temp.path().join("file2.txt").exists());
1748 }
1749
1750 #[tokio::test]
1751 async fn test_atomic_batch_first_op_fails_no_rollback_needed() {
1752 let (manager, temp) = create_test_manager();
1753
1754 fs::write(temp.path().join("existing.txt"), "exists").unwrap();
1756
1757 let request = BatchRequest {
1758 operations: vec![
1759 Operation::Create {
1760 path: PathBuf::from("existing.txt"), content: "content".to_string(),
1762 },
1763 Operation::Create {
1764 path: PathBuf::from("new.txt"),
1765 content: "content".to_string(),
1766 },
1767 ],
1768 atomic: true,
1769 dry_run: false,
1770 };
1771
1772 let result = manager.batch(request).await;
1773
1774 assert!(!result.success);
1775 assert_eq!(result.succeeded, 0);
1776 assert_eq!(result.failed, 1);
1777 assert!(result.rollback_performed.is_none());
1779 assert!(result.rollback_details.is_none());
1780 assert!(!temp.path().join("new.txt").exists());
1782 }
1783
1784 #[tokio::test]
1785 async fn test_atomic_batch_copy_rollback() {
1786 let (manager, temp) = create_test_manager();
1787
1788 fs::write(temp.path().join("source.txt"), "source content").unwrap();
1789
1790 let request = BatchRequest {
1791 operations: vec![
1792 Operation::Copy {
1793 src: PathBuf::from("source.txt"),
1794 dst: PathBuf::from("copied.txt"),
1795 },
1796 Operation::Delete {
1797 path: PathBuf::from("nonexistent.txt"), },
1799 ],
1800 atomic: true,
1801 dry_run: false,
1802 };
1803
1804 let result = manager.batch(request).await;
1805
1806 assert!(!result.success);
1807 assert_eq!(result.rollback_performed, Some(true));
1808 assert!(temp.path().join("source.txt").exists());
1810 assert!(!temp.path().join("copied.txt").exists());
1811 }
1812
1813 #[tokio::test]
1814 async fn test_mkdir() {
1815 let (manager, temp) = create_test_manager();
1816 let path = PathBuf::from("new_directory");
1817
1818 let result = manager.mkdir(&path).await;
1819
1820 assert!(result.success);
1821 assert_eq!(result.operation, "mkdir");
1822 assert!(temp.path().join("new_directory").is_dir());
1823 }
1824
1825 #[tokio::test]
1826 async fn test_mkdir_nested() {
1827 let (manager, temp) = create_test_manager();
1828 let path = PathBuf::from("parent/child/grandchild");
1829
1830 let result = manager.mkdir(&path).await;
1831
1832 assert!(result.success);
1833 assert!(temp.path().join("parent/child/grandchild").is_dir());
1834 }
1835
1836 #[tokio::test]
1837 async fn test_mkdir_already_exists() {
1838 let (manager, temp) = create_test_manager();
1839 let path = PathBuf::from("existing_dir");
1840 fs::create_dir(temp.path().join("existing_dir")).unwrap();
1841
1842 let result = manager.mkdir(&path).await;
1843
1844 assert!(!result.success);
1845 assert!(result.error.unwrap().contains("already exists"));
1846 }
1847
1848 #[tokio::test]
1849 #[cfg(unix)]
1850 async fn test_symlink() {
1851 let (manager, temp) = create_test_manager();
1852 let target = PathBuf::from("target_file.txt");
1853 let link = PathBuf::from("link_to_target");
1854
1855 fs::write(temp.path().join("target_file.txt"), "content").unwrap();
1857
1858 let result = manager.symlink(&target, &link).await;
1859
1860 assert!(result.success);
1861 assert_eq!(result.operation, "symlink");
1862 let link_path = temp.path().join("link_to_target");
1863 assert!(link_path.symlink_metadata().is_ok());
1864 assert_eq!(
1865 std::fs::read_link(&link_path).unwrap(),
1866 PathBuf::from("target_file.txt")
1867 );
1868 }
1869
1870 #[tokio::test]
1871 #[cfg(unix)]
1872 async fn test_symlink_relative_target_from_link_parent() {
1873 let (manager, temp) = create_test_manager();
1874 fs::create_dir(temp.path().join("sub")).unwrap();
1875 fs::write(temp.path().join("ok.txt"), "ok").unwrap();
1876
1877 let target = PathBuf::from("../ok.txt");
1878 let link = PathBuf::from("sub/link");
1879 let result = manager.symlink(&target, &link).await;
1880
1881 assert!(result.success, "{:?}", result.error);
1882 assert_eq!(
1883 std::fs::read_link(temp.path().join("sub/link")).unwrap(),
1884 PathBuf::from("../ok.txt")
1885 );
1886 }
1887
1888 #[tokio::test]
1889 #[cfg(unix)]
1890 async fn test_symlink_rejects_parent_relative_escape_via_existing_link() {
1891 let (manager, temp) = create_test_manager();
1892 let outside = TempDir::new().unwrap();
1893 fs::create_dir(temp.path().join("sub")).unwrap();
1894 std::os::unix::fs::symlink(outside.path(), temp.path().join("sub/out")).unwrap();
1895
1896 let target = PathBuf::from("out/secret.txt");
1897 let link = PathBuf::from("sub/link");
1898 let result = manager.symlink(&target, &link).await;
1899
1900 assert!(!result.success);
1901 assert!(
1902 result.error.as_ref().unwrap().contains("escapes"),
1903 "{:?}",
1904 result.error
1905 );
1906 assert!(!temp.path().join("sub/link").exists());
1907 }
1908
1909 #[tokio::test]
1910 #[cfg(unix)]
1911 async fn test_symlink_already_exists() {
1912 let (manager, temp) = create_test_manager();
1913 let target = PathBuf::from("target.txt");
1914 let link = PathBuf::from("existing_link");
1915
1916 fs::write(temp.path().join("existing_link"), "content").unwrap();
1918
1919 let result = manager.symlink(&target, &link).await;
1920
1921 assert!(!result.success);
1922 assert!(result.error.unwrap().contains("already exists"));
1923 }
1924
1925 #[tokio::test]
1926 async fn test_atomic_batch_mkdir_rollback() {
1927 let (manager, temp) = create_test_manager();
1928
1929 let request = BatchRequest {
1930 operations: vec![
1931 Operation::Mkdir {
1932 path: PathBuf::from("new_dir"),
1933 },
1934 Operation::Delete {
1935 path: PathBuf::from("nonexistent.txt"), },
1937 ],
1938 atomic: true,
1939 dry_run: false,
1940 };
1941
1942 let result = manager.batch(request).await;
1943
1944 assert!(!result.success);
1945 assert_eq!(result.rollback_performed, Some(true));
1946 assert!(!temp.path().join("new_dir").exists());
1948 }
1949
1950 fn create_test_manager_with_safety() -> (OpsManager, Arc<SafetyManager>, TempDir, TempDir) {
1951 let source_dir = TempDir::new().unwrap();
1952 let data_dir = TempDir::new().unwrap();
1953 let safety = Arc::new(SafetyManager::new(
1954 &source_dir.path().to_path_buf(),
1955 Some(crate::safety::SafetyConfig {
1956 data_dir: data_dir.path().to_path_buf(),
1957 trash_retention_days: 7,
1958 soft_delete: true,
1959 }),
1960 ));
1961 let manager = OpsManager::with_safety(
1962 source_dir.path().to_path_buf(),
1963 None,
1964 None,
1965 Arc::clone(&safety),
1966 );
1967 (manager, safety, source_dir, data_dir)
1968 }
1969
1970 #[tokio::test]
1971 async fn test_resolve_path_rejects_relative_escape() {
1972 let (manager, _temp) = create_test_manager();
1973 let err = manager
1974 .resolve_path(&PathBuf::from("../escape.txt"))
1975 .unwrap_err();
1976 assert!(err.contains("escapes"), "{err}");
1977 }
1978
1979 #[tokio::test]
1980 async fn test_resolve_path_rejects_absolute_outside_root() {
1981 let (manager, _temp) = create_test_manager();
1982 let outside = TempDir::new().unwrap();
1983 let err = manager
1984 .resolve_path(&outside.path().join("secret.txt"))
1985 .unwrap_err();
1986 assert!(err.contains("escapes"), "{err}");
1987 }
1988
1989 #[tokio::test]
1990 async fn test_create_rejects_path_jail_escape() {
1991 let (manager, temp) = create_test_manager();
1992 let result = manager
1993 .create(&PathBuf::from("../jailbreak.txt"), "nope")
1994 .await;
1995 assert!(!result.success);
1996 assert!(result.error.unwrap().contains("escapes"));
1997 assert!(!temp.path().parent().unwrap().join("jailbreak.txt").exists());
1998 }
1999
2000 #[tokio::test]
2001 async fn test_create_allows_absolute_path_inside_root() {
2002 let (manager, temp) = create_test_manager();
2003 let inside = temp.path().join("inside.txt");
2004 let result = manager.create(&inside, "ok").await;
2005 assert!(result.success);
2006 assert!(inside.exists());
2007 }
2008
2009 #[tokio::test]
2010 #[cfg(unix)]
2011 async fn test_create_rejects_symlink_escape() {
2012 let (manager, temp) = create_test_manager();
2013 let outside = TempDir::new().unwrap();
2014 std::os::unix::fs::symlink(outside.path(), temp.path().join("out")).unwrap();
2015
2016 let result = manager
2017 .create(&PathBuf::from("out/evil.txt"), "pwned")
2018 .await;
2019 assert!(!result.success);
2020 assert!(result.error.unwrap().contains("escapes"));
2021 assert!(!outside.path().join("evil.txt").exists());
2022 }
2023
2024 #[tokio::test]
2025 async fn test_create_undo_id_is_history_id_and_undo_uses_trash() {
2026 let (manager, safety, source_dir, _data) = create_test_manager_with_safety();
2027 let path = PathBuf::from("created.txt");
2028
2029 let result = manager.create(&path, "hello undo").await;
2030 assert!(result.success);
2031 let undo_id = result.undo_id.expect("create must expose history undo_id");
2032
2033 let history = safety.find_operation(undo_id).expect("history entry");
2034 assert_eq!(history.id, undo_id);
2035 assert!(history.reversible);
2036
2037 let file = source_dir.path().join("created.txt");
2038 assert!(file.exists());
2039
2040 let msg = safety.undo(undo_id).await.unwrap();
2041 assert!(msg.contains("trash"), "{msg}");
2042 assert!(!file.exists(), "undo create must not permanently delete");
2043
2044 let trash = safety.list_trash().await;
2045 assert_eq!(trash.len(), 1);
2046 let content = safety.get_trash_content(trash[0].id).unwrap();
2047 assert_eq!(content, b"hello undo");
2048 }
2049
2050 #[tokio::test]
2051 async fn test_delete_undo_id_is_history_id_not_trash_id() {
2052 let (manager, safety, source_dir, _data) = create_test_manager_with_safety();
2053 let path = PathBuf::from("to_soft_delete.txt");
2054 fs::write(source_dir.path().join("to_soft_delete.txt"), "keep me").unwrap();
2055
2056 let result = manager.delete(&path).await;
2057 assert!(result.success);
2058 let undo_id = result.undo_id.expect("delete must expose history undo_id");
2059
2060 let history = safety.find_operation(undo_id).expect("history entry");
2061 assert_eq!(history.id, undo_id);
2062 let trash_id = match history.undo_data {
2063 Some(UndoData::Delete { trash_id }) => trash_id,
2064 other => panic!("expected delete undo data, got {other:?}"),
2065 };
2066 assert_ne!(
2067 undo_id, trash_id,
2068 "undo_id must be the history id, not the trash id"
2069 );
2070 assert!(!source_dir.path().join("to_soft_delete.txt").exists());
2071
2072 safety.undo(undo_id).await.unwrap();
2073 assert!(source_dir.path().join("to_soft_delete.txt").exists());
2074 let content = fs::read_to_string(source_dir.path().join("to_soft_delete.txt")).unwrap();
2075 assert_eq!(content, "keep me");
2076 }
2077
2078 #[tokio::test]
2079 async fn test_move_undo_id_is_history_id() {
2080 let (manager, safety, source_dir, _data) = create_test_manager_with_safety();
2081 fs::write(source_dir.path().join("src.txt"), "moved").unwrap();
2082
2083 let result = manager
2084 .move_file(&PathBuf::from("src.txt"), &PathBuf::from("dst.txt"))
2085 .await;
2086 assert!(result.success);
2087 let undo_id = result.undo_id.expect("move must expose history undo_id");
2088 assert_eq!(safety.find_operation(undo_id).unwrap().id, undo_id);
2089
2090 safety.undo(undo_id).await.unwrap();
2091 assert!(source_dir.path().join("src.txt").exists());
2092 assert!(!source_dir.path().join("dst.txt").exists());
2093 }
2094
2095 #[tokio::test]
2096 async fn test_result_json_omits_internal_trash_id() {
2097 let (manager, _safety, source_dir, _data) = create_test_manager_with_safety();
2098 fs::write(source_dir.path().join("gone.txt"), "x").unwrap();
2099 manager.delete(&PathBuf::from("gone.txt")).await;
2100 let json = String::from_utf8(manager.get_last_result().await).unwrap();
2101 assert!(json.contains("undo_id"));
2102 assert!(!json.contains("trash_id"));
2103 }
2104}