Skip to main content

ragfs_fuse/
semantic.rs

1//! Semantic operations for intelligent file management.
2//!
3//! This module provides AI-powered file operations based on vector embeddings:
4//! - File organization by topic/similarity
5//! - Duplicate detection
6//! - Cleanup analysis
7//! - Similar file discovery
8//!
9//! All operations follow a Propose-Review-Apply pattern for safety.
10
11use chrono::{DateTime, Utc};
12use ragfs_core::{
13    Chunk, DistanceMetric, Embedder, EmbeddingConfig, FileRecord, SearchQuery, VectorStore,
14};
15use serde::{Deserialize, Serialize};
16use std::collections::{HashMap, HashSet};
17use std::fs;
18use std::path::{Path, PathBuf};
19use std::sync::Arc;
20use tokio::sync::RwLock;
21use tracing::{debug, info, warn};
22use uuid::Uuid;
23
24/// Request to organize files in a directory.
25#[derive(Debug, Clone, Serialize, Deserialize)]
26pub struct OrganizeRequest {
27    /// Directory scope (relative to source root)
28    pub scope: PathBuf,
29    /// Organization strategy
30    pub strategy: OrganizeStrategy,
31    /// Maximum number of groups to create
32    #[serde(default = "default_max_groups")]
33    pub max_groups: usize,
34    /// Minimum similarity threshold for grouping (0.0-1.0)
35    #[serde(default = "default_similarity_threshold")]
36    pub similarity_threshold: f32,
37}
38
39fn default_max_groups() -> usize {
40    10
41}
42
43fn default_similarity_threshold() -> f32 {
44    0.7
45}
46
47/// Strategy for organizing files.
48#[derive(Debug, Clone, Serialize, Deserialize)]
49#[serde(rename_all = "snake_case")]
50pub enum OrganizeStrategy {
51    /// Group by semantic topic/content similarity
52    ByTopic,
53    /// Group by file type first, then by content
54    ByType,
55    /// Group by project/module structure
56    ByProject,
57    /// Custom grouping with specified categories
58    Custom { categories: Vec<String> },
59}
60
61/// A proposed semantic operation plan.
62#[derive(Debug, Clone, Serialize, Deserialize)]
63pub struct SemanticPlan {
64    /// Unique plan identifier
65    pub id: Uuid,
66    /// When the plan was created
67    pub created_at: DateTime<Utc>,
68    /// Type of operation
69    pub operation: PlanOperation,
70    /// Human-readable description
71    pub description: String,
72    /// Proposed file operations
73    pub actions: Vec<PlanAction>,
74    /// Status of the plan
75    pub status: PlanStatus,
76    /// Estimated impact (files affected)
77    pub impact: PlanImpact,
78}
79
80/// Type of semantic operation.
81#[derive(Debug, Clone, Serialize, Deserialize)]
82#[serde(rename_all = "snake_case")]
83pub enum PlanOperation {
84    /// Organize files into groups
85    Organize {
86        scope: PathBuf,
87        strategy: OrganizeStrategy,
88    },
89    /// Clean up files
90    Cleanup { scope: PathBuf },
91    /// Deduplicate files
92    Dedupe { scope: PathBuf },
93}
94
95/// A single action in a plan.
96#[derive(Debug, Clone, Serialize, Deserialize)]
97pub struct PlanAction {
98    /// Type of action
99    pub action: ActionType,
100    /// Confidence score (0.0-1.0)
101    pub confidence: f32,
102    /// Reason for this action
103    pub reason: String,
104}
105
106/// Type of file action.
107#[derive(Debug, Clone, Serialize, Deserialize)]
108#[serde(rename_all = "snake_case")]
109pub enum ActionType {
110    /// Move a file to a new location
111    Move { from: PathBuf, to: PathBuf },
112    /// Create a new directory
113    Mkdir { path: PathBuf },
114    /// Delete a file (will use soft delete)
115    Delete { path: PathBuf },
116    /// Create a symlink
117    Symlink { target: PathBuf, link: PathBuf },
118}
119
120/// Status of a plan.
121#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
122#[serde(rename_all = "snake_case")]
123pub enum PlanStatus {
124    /// Plan is pending review
125    Pending,
126    /// Plan was approved and is being executed
127    Approved,
128    /// Plan was rejected
129    Rejected,
130    /// Plan was executed successfully
131    Completed,
132    /// Plan execution failed
133    Failed { error: String },
134}
135
136/// Impact summary of a plan.
137#[derive(Debug, Clone, Default, Serialize, Deserialize)]
138pub struct PlanImpact {
139    /// Total files affected
140    pub files_affected: usize,
141    /// Directories created
142    pub dirs_created: usize,
143    /// Files moved
144    pub files_moved: usize,
145    /// Files deleted
146    pub files_deleted: usize,
147}
148
149/// Analysis of cleanup candidates.
150#[derive(Debug, Clone, Serialize, Deserialize)]
151pub struct CleanupAnalysis {
152    /// When the analysis was performed
153    pub analyzed_at: DateTime<Utc>,
154    /// Total files analyzed
155    pub total_files: usize,
156    /// Cleanup candidates
157    pub candidates: Vec<CleanupCandidate>,
158    /// Potential space savings in bytes
159    pub potential_savings_bytes: u64,
160}
161
162/// A file that could be cleaned up.
163#[derive(Debug, Clone, Serialize, Deserialize)]
164pub struct CleanupCandidate {
165    /// File path
166    pub path: PathBuf,
167    /// Reason for cleanup suggestion
168    pub reason: CleanupReason,
169    /// Confidence score (0.0-1.0)
170    pub confidence: f32,
171    /// File size in bytes
172    pub size_bytes: u64,
173}
174
175/// Reason a file is suggested for cleanup.
176#[derive(Debug, Clone, Serialize, Deserialize)]
177#[serde(rename_all = "snake_case")]
178pub enum CleanupReason {
179    /// File appears to be a duplicate
180    Duplicate {
181        similar_to: PathBuf,
182        similarity: f32,
183    },
184    /// File hasn't been accessed in a long time
185    Stale { last_accessed: DateTime<Utc> },
186    /// Temporary file pattern
187    Temporary,
188    /// Generated file that can be recreated
189    Generated { source: PathBuf },
190    /// Empty or near-empty file
191    Empty,
192}
193
194/// Groups of duplicate/similar files.
195#[derive(Debug, Clone, Serialize, Deserialize)]
196pub struct DuplicateGroups {
197    /// When the analysis was performed
198    pub analyzed_at: DateTime<Utc>,
199    /// Minimum similarity threshold used
200    pub threshold: f32,
201    /// Groups of similar files
202    pub groups: Vec<DuplicateGroup>,
203    /// Total potential savings if duplicates removed
204    pub potential_savings_bytes: u64,
205}
206
207/// A group of similar files.
208#[derive(Debug, Clone, Serialize, Deserialize)]
209pub struct DuplicateGroup {
210    /// Group identifier
211    pub id: Uuid,
212    /// Representative file (keep this one)
213    pub representative: PathBuf,
214    /// Similar files (candidates for removal)
215    pub duplicates: Vec<DuplicateEntry>,
216    /// Total size of duplicates
217    pub wasted_bytes: u64,
218}
219
220/// A duplicate file entry.
221#[derive(Debug, Clone, Serialize, Deserialize)]
222pub struct DuplicateEntry {
223    /// File path
224    pub path: PathBuf,
225    /// Similarity to representative (0.0-1.0)
226    pub similarity: f32,
227    /// File size
228    pub size_bytes: u64,
229}
230
231/// Result of finding similar files.
232#[derive(Debug, Clone, Serialize, Deserialize)]
233pub struct SimilarFilesResult {
234    /// Source file
235    pub source: PathBuf,
236    /// Similar files found
237    pub similar: Vec<SimilarFile>,
238}
239
240/// A similar file.
241#[derive(Debug, Clone, Serialize, Deserialize)]
242pub struct SimilarFile {
243    /// File path
244    pub path: PathBuf,
245    /// Similarity score (0.0-1.0)
246    pub similarity: f32,
247    /// Preview of content
248    #[serde(skip_serializing_if = "Option::is_none")]
249    pub preview: Option<String>,
250}
251
252/// Configuration for semantic operations.
253#[derive(Debug, Clone)]
254pub struct SemanticConfig {
255    /// Minimum similarity for duplicate detection
256    pub duplicate_threshold: f32,
257    /// Number of results for similar file search
258    pub similar_limit: usize,
259    /// Maximum plan retention (in hours)
260    pub plan_retention_hours: u32,
261    /// Base directory for persistence (plans, etc.)
262    pub data_dir: PathBuf,
263}
264
265impl Default for SemanticConfig {
266    fn default() -> Self {
267        let data_dir = dirs::data_local_dir()
268            .unwrap_or_else(|| PathBuf::from("."))
269            .join("ragfs");
270
271        Self {
272            duplicate_threshold: 0.95,
273            similar_limit: 10,
274            plan_retention_hours: 24,
275            data_dir,
276        }
277    }
278}
279
280/// Result of executing a plan action.
281#[derive(Debug, Clone, Serialize, Deserialize)]
282pub struct ActionResult {
283    /// Whether the action succeeded
284    pub success: bool,
285    /// ID for undoing this action (if reversible)
286    #[serde(skip_serializing_if = "Option::is_none")]
287    pub undo_id: Option<Uuid>,
288    /// Error message if failed
289    #[serde(skip_serializing_if = "Option::is_none")]
290    pub error: Option<String>,
291    /// When the action was executed
292    pub executed_at: DateTime<Utc>,
293}
294
295/// Semantic manager for intelligent file operations.
296pub struct SemanticManager {
297    /// Source directory root
298    source: PathBuf,
299    /// Vector store for similarity search
300    store: Option<Arc<dyn VectorStore>>,
301    /// Embedder for generating embeddings
302    embedder: Option<Arc<dyn Embedder>>,
303    /// Configuration
304    config: SemanticConfig,
305    /// Pending plans (`plan_id` -> plan)
306    pending_plans: Arc<RwLock<HashMap<Uuid, SemanticPlan>>>,
307    /// Last similar files result
308    last_similar_result: Arc<RwLock<Option<SimilarFilesResult>>>,
309    /// Cached cleanup analysis
310    cleanup_cache: Arc<RwLock<Option<CleanupAnalysis>>>,
311    /// Cached duplicate groups
312    dedupe_cache: Arc<RwLock<Option<DuplicateGroups>>>,
313    /// Directory for storing plans
314    plans_dir: PathBuf,
315    /// Operations manager for executing actions
316    ops_manager: Option<Arc<crate::ops::OpsManager>>,
317}
318
319impl SemanticManager {
320    /// Create a new semantic manager.
321    pub fn new(
322        source: PathBuf,
323        store: Option<Arc<dyn VectorStore>>,
324        embedder: Option<Arc<dyn Embedder>>,
325        config: Option<SemanticConfig>,
326    ) -> Self {
327        let config = config.unwrap_or_default();
328
329        // Create index hash for isolation (same pattern as SafetyManager)
330        let index_hash = blake3::hash(source.to_string_lossy().as_bytes())
331            .to_hex()
332            .chars()
333            .take(16)
334            .collect::<String>();
335
336        let plans_dir = config.data_dir.join("plans").join(&index_hash);
337
338        // Ensure plans directory exists
339        if let Err(e) = fs::create_dir_all(&plans_dir) {
340            warn!("Failed to create plans directory: {e}");
341        }
342
343        // Load existing plans from disk
344        let plans = Self::load_plans(&plans_dir);
345        info!("Loaded {} existing semantic plans", plans.len());
346
347        Self {
348            source,
349            store,
350            embedder,
351            config,
352            pending_plans: Arc::new(RwLock::new(plans)),
353            last_similar_result: Arc::new(RwLock::new(None)),
354            cleanup_cache: Arc::new(RwLock::new(None)),
355            dedupe_cache: Arc::new(RwLock::new(None)),
356            plans_dir,
357            ops_manager: None,
358        }
359    }
360
361    /// Create a semantic manager with an operations manager for plan execution.
362    pub fn with_ops(
363        source: PathBuf,
364        store: Option<Arc<dyn VectorStore>>,
365        embedder: Option<Arc<dyn Embedder>>,
366        config: Option<SemanticConfig>,
367        ops_manager: Arc<crate::ops::OpsManager>,
368    ) -> Self {
369        let mut manager = Self::new(source, store, embedder, config);
370        manager.ops_manager = Some(ops_manager);
371        manager
372    }
373
374    /// Set the operations manager.
375    pub fn set_ops_manager(&mut self, ops_manager: Arc<crate::ops::OpsManager>) {
376        self.ops_manager = Some(ops_manager);
377    }
378
379    /// Load all plans from disk.
380    fn load_plans(plans_dir: &PathBuf) -> HashMap<Uuid, SemanticPlan> {
381        let mut plans = HashMap::new();
382
383        if !plans_dir.exists() {
384            return plans;
385        }
386
387        let entries = match fs::read_dir(plans_dir) {
388            Ok(e) => e,
389            Err(e) => {
390                warn!("Failed to read plans directory: {e}");
391                return plans;
392            }
393        };
394
395        for entry in entries.flatten() {
396            let path = entry.path();
397            if path.extension().is_some_and(|e| e == "json")
398                && let Ok(content) = fs::read_to_string(&path)
399            {
400                match serde_json::from_str::<SemanticPlan>(&content) {
401                    Ok(plan) => {
402                        plans.insert(plan.id, plan);
403                    }
404                    Err(e) => {
405                        warn!("Failed to parse plan file {:?}: {e}", path);
406                    }
407                }
408            }
409        }
410
411        plans
412    }
413
414    /// Save a single plan to disk.
415    fn save_plan(&self, plan: &SemanticPlan) -> std::io::Result<()> {
416        let plan_path = self.plans_dir.join(format!("{}.json", plan.id));
417        let temp_path = self.plans_dir.join(format!("{}.json.tmp", plan.id));
418
419        // Write to temp file first for atomic operation
420        let content = serde_json::to_string_pretty(plan)?;
421        fs::write(&temp_path, content)?;
422
423        // Atomic rename
424        fs::rename(&temp_path, &plan_path)?;
425
426        Ok(())
427    }
428
429    /// Delete a plan file from disk.
430    fn delete_plan_file(&self, plan_id: Uuid) -> std::io::Result<()> {
431        let plan_path = self.plans_dir.join(format!("{plan_id}.json"));
432        if plan_path.exists() {
433            fs::remove_file(&plan_path)?;
434        }
435        Ok(())
436    }
437
438    /// Purge expired plans from memory and disk.
439    pub async fn purge_expired_plans(&self) -> usize {
440        let now = Utc::now();
441        let retention = chrono::Duration::hours(i64::from(self.config.plan_retention_hours));
442        let cutoff = now - retention;
443
444        let mut plans = self.pending_plans.write().await;
445        let expired: Vec<Uuid> = plans
446            .iter()
447            .filter(|(_, p)| {
448                // Purge completed/rejected/failed plans past retention
449                // Keep pending plans regardless of age
450                matches!(
451                    p.status,
452                    PlanStatus::Completed | PlanStatus::Rejected | PlanStatus::Failed { .. }
453                ) && p.created_at < cutoff
454            })
455            .map(|(id, _)| *id)
456            .collect();
457
458        let mut purged = 0;
459        for id in &expired {
460            plans.remove(id);
461            if let Err(e) = self.delete_plan_file(*id) {
462                warn!("Failed to delete expired plan file {}: {e}", id);
463            } else {
464                purged += 1;
465            }
466        }
467
468        if purged > 0 {
469            info!("Purged {} expired semantic plans", purged);
470        }
471
472        purged
473    }
474
475    /// Check if semantic operations are available.
476    #[must_use]
477    pub fn is_available(&self) -> bool {
478        self.store.is_some() && self.embedder.is_some()
479    }
480
481    /// Resolve a path and reject anything that escapes the source root.
482    fn resolve_path(&self, path: &Path) -> Result<PathBuf, String> {
483        ragfs_core::resolve_under_root(&self.source, path)
484    }
485
486    /// Canonicalize for comparison; keep the original path if it is gone.
487    fn canonical_or_owned(path: &Path) -> PathBuf {
488        path.canonicalize().unwrap_or_else(|_| path.to_path_buf())
489    }
490
491    /// Async canonicalize so FUSE-driven tokio workers are not blocked on stat.
492    async fn canonical_or_owned_async(path: PathBuf) -> PathBuf {
493        tokio::fs::canonicalize(&path).await.unwrap_or(path)
494    }
495
496    /// Find files similar to a given path.
497    pub async fn find_similar(&self, path: &PathBuf) -> Result<SimilarFilesResult, String> {
498        let full_path = self.resolve_path(path)?;
499        let store = self.store.as_ref().ok_or("Vector store not available")?;
500        let embedder = self.embedder.as_ref().ok_or("Embedder not available")?;
501
502        debug!("Finding files similar to: {}", full_path.display());
503
504        // Read the file content
505        let content =
506            std::fs::read_to_string(&full_path).map_err(|e| format!("Failed to read file: {e}"))?;
507
508        // Generate embedding for the content
509        let config = EmbeddingConfig::default();
510        let embedding_output = embedder
511            .embed_query(&content, &config)
512            .await
513            .map_err(|e| format!("Failed to generate embedding: {e}"))?;
514
515        // Search for similar files
516        let query = SearchQuery {
517            embedding: embedding_output.embedding,
518            text: None,
519            limit: self.config.similar_limit + 1, // +1 to exclude self
520            filters: Vec::new(),
521            metric: DistanceMetric::Cosine,
522            scope_prefix: None,
523        };
524        let results = store
525            .search(query)
526            .await
527            .map_err(|e| format!("Search failed: {e}"))?;
528
529        // Convert results, excluding the source file itself.
530        // Index paths may be non-canonical; compare after canonicalize.
531        let source_canonical = Self::canonical_or_owned_async(full_path.clone()).await;
532        let similar_limit = self.config.similar_limit;
533        let similar: Vec<SimilarFile> = tokio::task::spawn_blocking(move || {
534            results
535                .into_iter()
536                .filter(|r| Self::canonical_or_owned(&r.file_path) != source_canonical)
537                .take(similar_limit)
538                .map(|r| SimilarFile {
539                    path: r.file_path,
540                    similarity: r.score, // score is already similarity (higher = more similar)
541                    preview: Some(truncate_content(&r.content, 200)),
542                })
543                .collect()
544        })
545        .await
546        .map_err(|e| format!("Failed to canonicalize similar-file paths: {e}"))?;
547
548        let result = SimilarFilesResult {
549            source: full_path,
550            similar,
551        };
552
553        // Cache the result
554        *self.last_similar_result.write().await = Some(result.clone());
555
556        info!("Found {} similar files", result.similar.len());
557        Ok(result)
558    }
559
560    /// Get the last similar files result.
561    pub async fn get_last_similar_result(&self) -> Option<SimilarFilesResult> {
562        self.last_similar_result.read().await.clone()
563    }
564
565    /// Analyze files for cleanup candidates.
566    pub async fn analyze_cleanup(&self) -> Result<CleanupAnalysis, String> {
567        let store = self.store.as_ref().ok_or("Vector store not available")?;
568
569        debug!("Analyzing files for cleanup candidates");
570
571        // Get all file records from the store
572        let stats = store
573            .stats()
574            .await
575            .map_err(|e| format!("Failed to get stats: {e}"))?;
576
577        let mut candidates = Vec::new();
578        let mut potential_savings: u64 = 0;
579
580        // For now, we'll focus on duplicate detection as the primary cleanup criterion
581        // This could be expanded to include stale file detection, etc.
582
583        // Get duplicate groups and convert high-confidence duplicates to cleanup candidates
584        if let Ok(dupes) = self.find_duplicates().await {
585            for group in &dupes.groups {
586                for dup in &group.duplicates {
587                    if dup.similarity >= self.config.duplicate_threshold {
588                        candidates.push(CleanupCandidate {
589                            path: dup.path.clone(),
590                            reason: CleanupReason::Duplicate {
591                                similar_to: group.representative.clone(),
592                                similarity: dup.similarity,
593                            },
594                            confidence: dup.similarity,
595                            size_bytes: dup.size_bytes,
596                        });
597                        potential_savings += dup.size_bytes;
598                    }
599                }
600            }
601        }
602
603        let analysis = CleanupAnalysis {
604            analyzed_at: Utc::now(),
605            total_files: stats.total_files as usize,
606            candidates,
607            potential_savings_bytes: potential_savings,
608        };
609
610        // Cache the result
611        *self.cleanup_cache.write().await = Some(analysis.clone());
612
613        info!(
614            "Cleanup analysis: {} candidates, {} bytes potential savings",
615            analysis.candidates.len(),
616            analysis.potential_savings_bytes
617        );
618
619        Ok(analysis)
620    }
621
622    /// Get cached cleanup analysis.
623    pub async fn get_cleanup_analysis(&self) -> Option<CleanupAnalysis> {
624        self.cleanup_cache.read().await.clone()
625    }
626
627    /// Find duplicate file groups.
628    pub async fn find_duplicates(&self) -> Result<DuplicateGroups, String> {
629        let store = self.store.as_ref().ok_or("Vector store not available")?;
630        let _embedder = self.embedder.as_ref().ok_or("Embedder not available")?;
631
632        debug!("Finding duplicate files");
633
634        // Get all chunks and files from the store
635        let all_chunks = store
636            .get_all_chunks()
637            .await
638            .map_err(|e| format!("Failed to get chunks: {e}"))?;
639
640        let all_files = store
641            .get_all_files()
642            .await
643            .map_err(|e| format!("Failed to get files: {e}"))?;
644
645        if all_files.is_empty() {
646            return Ok(DuplicateGroups {
647                analyzed_at: Utc::now(),
648                threshold: self.config.duplicate_threshold,
649                groups: Vec::new(),
650                potential_savings_bytes: 0,
651            });
652        }
653
654        // Build a map of file_path -> chunks with embeddings
655        let mut file_chunks: HashMap<PathBuf, Vec<&Chunk>> = HashMap::new();
656        for chunk in &all_chunks {
657            if chunk.embedding.is_some() {
658                file_chunks
659                    .entry(chunk.file_path.clone())
660                    .or_default()
661                    .push(chunk);
662            }
663        }
664
665        // Build file info map
666        let file_info: HashMap<PathBuf, &FileRecord> =
667            all_files.iter().map(|f| (f.path.clone(), f)).collect();
668
669        // Calculate average embedding for each file
670        let file_embeddings: HashMap<PathBuf, Vec<f32>> = file_chunks
671            .iter()
672            .filter_map(|(path, chunks)| {
673                let embeddings: Vec<&Vec<f32>> =
674                    chunks.iter().filter_map(|c| c.embedding.as_ref()).collect();
675
676                if embeddings.is_empty() {
677                    return None;
678                }
679
680                // Average the embeddings
681                let dim = embeddings[0].len();
682                let mut avg = vec![0.0f32; dim];
683                for emb in &embeddings {
684                    for (i, &v) in emb.iter().enumerate() {
685                        avg[i] += v;
686                    }
687                }
688                let count = embeddings.len() as f32;
689                for v in &mut avg {
690                    *v /= count;
691                }
692
693                // Normalize the averaged embedding
694                let norm: f32 = avg.iter().map(|x| x * x).sum::<f32>().sqrt();
695                if norm > 0.0 {
696                    for v in &mut avg {
697                        *v /= norm;
698                    }
699                }
700
701                Some((path.clone(), avg))
702            })
703            .collect();
704
705        // Find similar file pairs using cosine similarity
706        let file_paths: Vec<&PathBuf> = file_embeddings.keys().collect();
707        let mut similarity_pairs: Vec<(PathBuf, PathBuf, f32)> = Vec::new();
708
709        for (i, path_a) in file_paths.iter().enumerate() {
710            let emb_a = &file_embeddings[*path_a];
711            for path_b in file_paths.iter().skip(i + 1) {
712                let emb_b = &file_embeddings[*path_b];
713                let similarity = cosine_similarity(emb_a, emb_b);
714
715                if similarity >= self.config.duplicate_threshold {
716                    similarity_pairs.push(((*path_a).clone(), (*path_b).clone(), similarity));
717                }
718            }
719        }
720
721        // Cluster similar files using Union-Find
722        let mut groups: Vec<DuplicateGroup> = Vec::new();
723        let mut processed: HashSet<PathBuf> = HashSet::new();
724
725        for (path_a, path_b, similarity) in similarity_pairs {
726            if processed.contains(&path_a) || processed.contains(&path_b) {
727                continue;
728            }
729
730            // Find or create group for path_a
731            let size_a = file_info.get(&path_a).map_or(0, |f| f.size_bytes);
732            let size_b = file_info.get(&path_b).map_or(0, |f| f.size_bytes);
733
734            // Use the larger file as representative
735            let (representative, duplicate, dup_similarity, dup_size) = if size_a >= size_b {
736                (path_a.clone(), path_b.clone(), similarity, size_b)
737            } else {
738                (path_b.clone(), path_a.clone(), similarity, size_a)
739            };
740
741            // Check if representative already has a group
742            if let Some(group) = groups
743                .iter_mut()
744                .find(|g| g.representative == representative)
745            {
746                group.duplicates.push(DuplicateEntry {
747                    path: duplicate.clone(),
748                    similarity: dup_similarity,
749                    size_bytes: dup_size,
750                });
751                group.wasted_bytes += dup_size;
752                processed.insert(duplicate);
753            } else {
754                // Create new group
755                groups.push(DuplicateGroup {
756                    id: Uuid::new_v4(),
757                    representative: representative.clone(),
758                    duplicates: vec![DuplicateEntry {
759                        path: duplicate.clone(),
760                        similarity: dup_similarity,
761                        size_bytes: dup_size,
762                    }],
763                    wasted_bytes: dup_size,
764                });
765                processed.insert(representative);
766                processed.insert(duplicate);
767            }
768        }
769
770        let potential_savings: u64 = groups.iter().map(|g| g.wasted_bytes).sum();
771
772        info!(
773            "Found {} duplicate groups with {} bytes potential savings",
774            groups.len(),
775            potential_savings
776        );
777
778        let result = DuplicateGroups {
779            analyzed_at: Utc::now(),
780            threshold: self.config.duplicate_threshold,
781            groups,
782            potential_savings_bytes: potential_savings,
783        };
784
785        // Cache the result
786        *self.dedupe_cache.write().await = Some(result.clone());
787
788        Ok(result)
789    }
790
791    /// Get cached duplicate groups.
792    pub async fn get_duplicate_groups(&self) -> Option<DuplicateGroups> {
793        self.dedupe_cache.read().await.clone()
794    }
795
796    /// Create an organization plan.
797    pub async fn create_organize_plan(
798        &self,
799        request: OrganizeRequest,
800    ) -> Result<SemanticPlan, String> {
801        let scope_path = Self::canonical_or_owned_async(self.resolve_path(&request.scope)?).await;
802        let store = self.store.as_ref().ok_or("Vector store not available")?;
803        let embedder = self.embedder.as_ref();
804
805        debug!(
806            "Creating organization plan for: {}",
807            request.scope.display()
808        );
809
810        // Get all chunks and files
811        let all_chunks = store
812            .get_all_chunks()
813            .await
814            .map_err(|e| format!("Failed to get chunks: {e}"))?;
815
816        let all_files = store
817            .get_all_files()
818            .await
819            .map_err(|e| format!("Failed to get files: {e}"))?;
820
821        let scope_for_cmp = scope_path.clone();
822        let file_paths: Vec<PathBuf> = all_files.iter().map(|f| f.path.clone()).collect();
823        let chunk_paths: Vec<PathBuf> = all_chunks.iter().map(|c| c.file_path.clone()).collect();
824        let (in_scope_files, in_scope_chunks) = tokio::task::spawn_blocking(move || {
825            let files: HashSet<PathBuf> = file_paths
826                .into_iter()
827                .filter(|p| Self::canonical_or_owned(p).starts_with(&scope_for_cmp))
828                .collect();
829            let chunks: HashSet<PathBuf> = chunk_paths
830                .into_iter()
831                .filter(|p| Self::canonical_or_owned(p).starts_with(&scope_for_cmp))
832                .collect();
833            (files, chunks)
834        })
835        .await
836        .map_err(|e| format!("Failed to canonicalize scoped paths: {e}"))?;
837
838        let scoped_files: Vec<&FileRecord> = all_files
839            .iter()
840            .filter(|f| in_scope_files.contains(&f.path))
841            .collect();
842
843        if scoped_files.is_empty() {
844            return Ok(SemanticPlan {
845                id: Uuid::new_v4(),
846                created_at: Utc::now(),
847                operation: PlanOperation::Organize {
848                    scope: request.scope.clone(),
849                    strategy: request.strategy.clone(),
850                },
851                description: format!("No files found in scope: {}", request.scope.display()),
852                actions: Vec::new(),
853                status: PlanStatus::Pending,
854                impact: PlanImpact::default(),
855            });
856        }
857
858        // Build file embeddings map
859        let mut file_chunks: HashMap<PathBuf, Vec<&Chunk>> = HashMap::new();
860        for chunk in &all_chunks {
861            if chunk.embedding.is_some() && in_scope_chunks.contains(&chunk.file_path) {
862                file_chunks
863                    .entry(chunk.file_path.clone())
864                    .or_default()
865                    .push(chunk);
866            }
867        }
868
869        // Calculate average embedding for each file
870        let file_embeddings: HashMap<PathBuf, Vec<f32>> = file_chunks
871            .iter()
872            .filter_map(|(path, chunks)| {
873                let embeddings: Vec<&Vec<f32>> =
874                    chunks.iter().filter_map(|c| c.embedding.as_ref()).collect();
875
876                if embeddings.is_empty() {
877                    return None;
878                }
879
880                let dim = embeddings[0].len();
881                let mut avg = vec![0.0f32; dim];
882                for emb in &embeddings {
883                    for (i, &v) in emb.iter().enumerate() {
884                        avg[i] += v;
885                    }
886                }
887                let count = embeddings.len() as f32;
888                for v in &mut avg {
889                    *v /= count;
890                }
891
892                // Normalize
893                let norm: f32 = avg.iter().map(|x| x * x).sum::<f32>().sqrt();
894                if norm > 0.0 {
895                    for v in &mut avg {
896                        *v /= norm;
897                    }
898                }
899
900                Some((path.clone(), avg))
901            })
902            .collect();
903
904        // Generate actions based on strategy
905        let (actions, description) = match &request.strategy {
906            OrganizeStrategy::ByTopic => self.plan_by_topic(
907                &file_embeddings,
908                &scope_path,
909                request.max_groups,
910                request.similarity_threshold,
911            ),
912            OrganizeStrategy::ByType => self.plan_by_type(&scoped_files, &scope_path),
913            OrganizeStrategy::ByProject => self.plan_by_project(&scoped_files, &scope_path),
914            OrganizeStrategy::Custom { categories } => {
915                self.plan_by_custom(&file_embeddings, &scope_path, categories, embedder)
916                    .await
917            }
918        };
919
920        let dirs_created = actions
921            .iter()
922            .filter(|a| matches!(a.action, ActionType::Mkdir { .. }))
923            .count();
924        let files_moved = actions
925            .iter()
926            .filter(|a| matches!(a.action, ActionType::Move { .. }))
927            .count();
928
929        let plan = SemanticPlan {
930            id: Uuid::new_v4(),
931            created_at: Utc::now(),
932            operation: PlanOperation::Organize {
933                scope: request.scope,
934                strategy: request.strategy,
935            },
936            description,
937            actions,
938            status: PlanStatus::Pending,
939            impact: PlanImpact {
940                files_affected: files_moved,
941                dirs_created,
942                files_moved,
943                files_deleted: 0,
944            },
945        };
946
947        // Store the plan in memory
948        self.pending_plans
949            .write()
950            .await
951            .insert(plan.id, plan.clone());
952
953        // Persist to disk
954        if let Err(e) = self.save_plan(&plan) {
955            warn!("Failed to persist plan {}: {e}", plan.id);
956        }
957
958        info!(
959            "Created organization plan: {} with {} actions",
960            plan.id,
961            plan.actions.len()
962        );
963        Ok(plan)
964    }
965
966    /// Plan organization by semantic topic using clustering.
967    fn plan_by_topic(
968        &self,
969        file_embeddings: &HashMap<PathBuf, Vec<f32>>,
970        scope_path: &PathBuf,
971        max_groups: usize,
972        similarity_threshold: f32,
973    ) -> (Vec<PlanAction>, String) {
974        if file_embeddings.is_empty() {
975            return (Vec::new(), "No files with embeddings found".to_string());
976        }
977
978        // Simple clustering: find centroids and group files
979        let file_paths: Vec<&PathBuf> = file_embeddings.keys().collect();
980        let num_files = file_paths.len();
981        let num_clusters = max_groups.min(num_files);
982
983        // Initialize clusters with k random files (here we use evenly spaced indices)
984        let step = if num_files > num_clusters {
985            num_files / num_clusters
986        } else {
987            1
988        };
989        let mut centroids: Vec<Vec<f32>> = (0..num_clusters)
990            .map(|i| file_embeddings[file_paths[i * step.min(num_files - 1)]].clone())
991            .collect();
992
993        // Simple k-means iterations
994        let mut cluster_assignments: HashMap<PathBuf, usize> = HashMap::new();
995
996        for _ in 0..5 {
997            // Assign each file to nearest centroid
998            cluster_assignments.clear();
999            for path in &file_paths {
1000                let emb = &file_embeddings[*path];
1001                let mut best_cluster = 0;
1002                let mut best_sim = -1.0f32;
1003
1004                for (cluster_idx, centroid) in centroids.iter().enumerate() {
1005                    let sim = cosine_similarity(emb, centroid);
1006                    if sim > best_sim {
1007                        best_sim = sim;
1008                        best_cluster = cluster_idx;
1009                    }
1010                }
1011
1012                cluster_assignments.insert((*path).clone(), best_cluster);
1013            }
1014
1015            // Update centroids
1016            for (cluster_idx, centroid) in centroids.iter_mut().enumerate() {
1017                let members: Vec<&PathBuf> = cluster_assignments
1018                    .iter()
1019                    .filter(|&(_, c)| *c == cluster_idx)
1020                    .map(|(p, _)| p)
1021                    .collect();
1022
1023                if members.is_empty() {
1024                    continue;
1025                }
1026
1027                let dim = centroid.len();
1028                let mut new_centroid = vec![0.0f32; dim];
1029
1030                for path in &members {
1031                    let emb = &file_embeddings[*path];
1032                    for (i, &v) in emb.iter().enumerate() {
1033                        new_centroid[i] += v;
1034                    }
1035                }
1036
1037                let count = members.len() as f32;
1038                for v in &mut new_centroid {
1039                    *v /= count;
1040                }
1041
1042                // Normalize
1043                let norm: f32 = new_centroid.iter().map(|x| x * x).sum::<f32>().sqrt();
1044                if norm > 0.0 {
1045                    for v in &mut new_centroid {
1046                        *v /= norm;
1047                    }
1048                }
1049
1050                *centroid = new_centroid;
1051            }
1052        }
1053
1054        // Generate actions
1055        let mut actions = Vec::new();
1056
1057        // Create topic directories
1058        for cluster_idx in 0..num_clusters {
1059            let topic_dir = scope_path.join(format!("topic_{}", cluster_idx + 1));
1060            actions.push(PlanAction {
1061                action: ActionType::Mkdir { path: topic_dir },
1062                confidence: 1.0,
1063                reason: format!("Create directory for topic cluster {}", cluster_idx + 1),
1064            });
1065        }
1066
1067        // Move files to their clusters
1068        for (path, &cluster_idx) in &cluster_assignments {
1069            let file_name = path.file_name().unwrap_or_default();
1070            let topic_dir = scope_path.join(format!("topic_{}", cluster_idx + 1));
1071            let new_path = topic_dir.join(file_name);
1072
1073            if new_path != *path {
1074                // Calculate confidence based on distance to centroid
1075                let emb = &file_embeddings[path];
1076                let centroid = &centroids[cluster_idx];
1077                let confidence = cosine_similarity(emb, centroid).max(similarity_threshold);
1078
1079                actions.push(PlanAction {
1080                    action: ActionType::Move {
1081                        from: path.clone(),
1082                        to: new_path,
1083                    },
1084                    confidence,
1085                    reason: format!(
1086                        "Move to topic cluster {} based on content similarity",
1087                        cluster_idx + 1
1088                    ),
1089                });
1090            }
1091        }
1092
1093        let description = format!(
1094            "Organize {} files into {} topic clusters",
1095            file_paths.len(),
1096            num_clusters
1097        );
1098
1099        (actions, description)
1100    }
1101
1102    /// Plan organization by file type.
1103    fn plan_by_type(
1104        &self,
1105        files: &[&FileRecord],
1106        scope_path: &PathBuf,
1107    ) -> (Vec<PlanAction>, String) {
1108        let mut actions = Vec::new();
1109        let mut type_dirs: HashSet<String> = HashSet::new();
1110
1111        for file in files {
1112            // Determine type directory based on extension or MIME type
1113            let type_dir = if let Some(ext) = file.path.extension() {
1114                ext.to_string_lossy().to_string()
1115            } else {
1116                // Use MIME type category
1117                file.mime_type
1118                    .split('/')
1119                    .next()
1120                    .unwrap_or("other")
1121                    .to_string()
1122            };
1123
1124            // Create type directory if needed
1125            if type_dirs.insert(type_dir.clone()) {
1126                actions.push(PlanAction {
1127                    action: ActionType::Mkdir {
1128                        path: scope_path.join(&type_dir),
1129                    },
1130                    confidence: 1.0,
1131                    reason: format!("Create directory for {type_dir} files"),
1132                });
1133            }
1134
1135            // Move file
1136            let file_name = file.path.file_name().unwrap_or_default();
1137            let new_path = scope_path.join(&type_dir).join(file_name);
1138
1139            if new_path != file.path {
1140                actions.push(PlanAction {
1141                    action: ActionType::Move {
1142                        from: file.path.clone(),
1143                        to: new_path,
1144                    },
1145                    confidence: 1.0,
1146                    reason: format!("Move to {type_dir} directory based on file type"),
1147                });
1148            }
1149        }
1150
1151        let description = format!(
1152            "Organize {} files into {} type-based directories",
1153            files.len(),
1154            type_dirs.len()
1155        );
1156
1157        (actions, description)
1158    }
1159
1160    /// Plan organization by project structure (based on imports/dependencies).
1161    fn plan_by_project(
1162        &self,
1163        files: &[&FileRecord],
1164        scope_path: &PathBuf,
1165    ) -> (Vec<PlanAction>, String) {
1166        // For project-based organization, we look at file paths to infer structure
1167        // This is a simplified implementation
1168        let mut actions = Vec::new();
1169        let mut project_dirs: HashSet<String> = HashSet::new();
1170
1171        for file in files {
1172            // Use the first directory component after scope as "project"
1173            let relative = file.path.strip_prefix(scope_path).unwrap_or(&file.path);
1174            let project = relative.components().next().map_or_else(
1175                || "root".to_string(),
1176                |c| c.as_os_str().to_string_lossy().to_string(),
1177            );
1178
1179            if project_dirs.insert(project.clone()) && !project.contains('.') {
1180                actions.push(PlanAction {
1181                    action: ActionType::Mkdir {
1182                        path: scope_path.join(&project),
1183                    },
1184                    confidence: 0.8,
1185                    reason: format!("Create project directory: {project}"),
1186                });
1187            }
1188        }
1189
1190        let description = format!(
1191            "Organize {} files into {} project directories",
1192            files.len(),
1193            project_dirs.len()
1194        );
1195
1196        (actions, description)
1197    }
1198
1199    /// Plan organization by custom categories.
1200    ///
1201    /// Generates embeddings for category names and assigns files to the
1202    /// best matching category using cosine similarity.
1203    async fn plan_by_custom(
1204        &self,
1205        file_embeddings: &HashMap<PathBuf, Vec<f32>>,
1206        scope_path: &PathBuf,
1207        categories: &[String],
1208        embedder: Option<&Arc<dyn Embedder>>,
1209    ) -> (Vec<PlanAction>, String) {
1210        let mut actions = Vec::new();
1211
1212        // Create category directories
1213        for category in categories {
1214            actions.push(PlanAction {
1215                action: ActionType::Mkdir {
1216                    path: scope_path.join(category),
1217                },
1218                confidence: 1.0,
1219                reason: format!("Create custom category directory: {category}"),
1220            });
1221        }
1222
1223        // Try to generate embeddings for categories and assign files automatically
1224        if let Some(embedder) = embedder {
1225            let category_texts: Vec<&str> = categories.iter().map(String::as_str).collect();
1226            let config = EmbeddingConfig::default();
1227
1228            match embedder.embed_text(&category_texts, &config).await {
1229                Ok(category_embeddings) if category_embeddings.len() == categories.len() => {
1230                    // Minimum similarity threshold for assignment
1231                    const MIN_SIMILARITY: f32 = 0.3;
1232                    let mut assigned_count = 0;
1233
1234                    for (file_path, file_emb) in file_embeddings {
1235                        // Find best matching category
1236                        let best = category_embeddings
1237                            .iter()
1238                            .zip(categories.iter())
1239                            .map(|(emb, cat)| (cat, cosine_similarity(file_emb, &emb.embedding)))
1240                            .max_by(|a, b| {
1241                                a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal)
1242                            });
1243
1244                        if let Some((category, score)) = best
1245                            && score >= MIN_SIMILARITY
1246                            && let Some(file_name) = file_path.file_name()
1247                        {
1248                            let new_path = scope_path.join(category).join(file_name);
1249                            if new_path != *file_path {
1250                                actions.push(PlanAction {
1251                                    action: ActionType::Move {
1252                                        from: file_path.clone(),
1253                                        to: new_path,
1254                                    },
1255                                    confidence: score,
1256                                    reason: format!(
1257                                        "Move to category '{category}' (similarity: {score:.2})"
1258                                    ),
1259                                });
1260                                assigned_count += 1;
1261                            }
1262                        }
1263                    }
1264
1265                    let description = format!(
1266                        "Organize {} files into {} custom categories ({} files assigned)",
1267                        file_embeddings.len(),
1268                        categories.len(),
1269                        assigned_count
1270                    );
1271                    return (actions, description);
1272                }
1273                Ok(_) => {
1274                    warn!(
1275                        "Category embedding count mismatch, falling back to directory creation only"
1276                    );
1277                }
1278                Err(e) => {
1279                    warn!(
1280                        "Failed to embed categories: {}, falling back to directory creation only",
1281                        e
1282                    );
1283                }
1284            }
1285        }
1286
1287        // Fallback: just create directories without automatic assignment
1288        let description = format!(
1289            "Created {} custom category directories for {} files (manual assignment needed)",
1290            categories.len(),
1291            file_embeddings.len()
1292        );
1293
1294        (actions, description)
1295    }
1296
1297    /// List all pending plans.
1298    pub async fn list_pending_plans(&self) -> Vec<SemanticPlan> {
1299        self.pending_plans
1300            .read()
1301            .await
1302            .values()
1303            .filter(|p| p.status == PlanStatus::Pending)
1304            .cloned()
1305            .collect()
1306    }
1307
1308    /// Get a specific plan.
1309    pub async fn get_plan(&self, plan_id: Uuid) -> Option<SemanticPlan> {
1310        self.pending_plans.read().await.get(&plan_id).cloned()
1311    }
1312
1313    /// Execute a single action via `OpsManager`.
1314    async fn execute_action(&self, action: &ActionType) -> Result<ActionResult, String> {
1315        let ops = self
1316            .ops_manager
1317            .as_ref()
1318            .ok_or("OpsManager not configured - cannot execute plan actions")?;
1319
1320        let result = match action {
1321            ActionType::Move { from, to } => ops.move_file(from, to).await,
1322            ActionType::Mkdir { path } => ops.mkdir(path).await,
1323            ActionType::Delete { path } => ops.delete(path).await,
1324            ActionType::Symlink { target, link } => ops.symlink(target, link).await,
1325        };
1326
1327        Ok(ActionResult {
1328            success: result.success,
1329            undo_id: result.undo_id,
1330            error: if result.success {
1331                None
1332            } else {
1333                Some(result.error.unwrap_or_else(|| "Unknown error".to_string()))
1334            },
1335            executed_at: Utc::now(),
1336        })
1337    }
1338
1339    /// Approve and execute a plan.
1340    pub async fn approve_plan(&self, plan_id: Uuid) -> Result<SemanticPlan, String> {
1341        // Verify OpsManager is available before starting
1342        if self.ops_manager.is_none() {
1343            return Err("OpsManager not configured - cannot execute plan actions".to_string());
1344        }
1345
1346        let mut plans = self.pending_plans.write().await;
1347        let plan = plans
1348            .get_mut(&plan_id)
1349            .ok_or_else(|| "Plan not found".to_string())?;
1350
1351        if plan.status != PlanStatus::Pending {
1352            return Err(format!("Plan is not pending: {:?}", plan.status));
1353        }
1354
1355        info!(
1356            "Approving plan: {} with {} actions",
1357            plan_id,
1358            plan.actions.len()
1359        );
1360        plan.status = PlanStatus::Approved;
1361
1362        // Execute actions sequentially, stopping on first failure
1363        let total_actions = plan.actions.len();
1364        let mut completed_actions = 0;
1365
1366        // Clone actions to avoid holding lock during execution
1367        let actions_to_execute: Vec<ActionType> =
1368            plan.actions.iter().map(|a| a.action.clone()).collect();
1369
1370        // Release write lock during execution to avoid deadlock
1371        drop(plans);
1372
1373        for (idx, action) in actions_to_execute.iter().enumerate() {
1374            debug!(
1375                "Executing action {}/{}: {:?}",
1376                idx + 1,
1377                total_actions,
1378                action
1379            );
1380
1381            match self.execute_action(action).await {
1382                Ok(result) if result.success => {
1383                    completed_actions += 1;
1384                    debug!(
1385                        "Action {}/{} succeeded (undo_id: {:?})",
1386                        idx + 1,
1387                        total_actions,
1388                        result.undo_id
1389                    );
1390                }
1391                Ok(result) => {
1392                    // Action failed
1393                    let error_msg = result.error.unwrap_or_else(|| "Unknown error".to_string());
1394                    warn!("Action {}/{} failed: {}", idx + 1, total_actions, error_msg);
1395
1396                    // Update plan status to failed
1397                    let mut plans = self.pending_plans.write().await;
1398                    if let Some(plan) = plans.get_mut(&plan_id) {
1399                        plan.status = PlanStatus::Failed {
1400                            error: format!(
1401                                "Action {} of {} failed: {}",
1402                                idx + 1,
1403                                total_actions,
1404                                error_msg
1405                            ),
1406                        };
1407
1408                        let result = plan.clone();
1409                        if let Err(e) = self.save_plan(&result) {
1410                            warn!("Failed to persist failed plan {}: {e}", plan_id);
1411                        }
1412                        return Ok(result);
1413                    }
1414                    return Err("Plan disappeared during execution".to_string());
1415                }
1416                Err(e) => {
1417                    // Execution error (OpsManager issue)
1418                    warn!(
1419                        "Failed to execute action {}/{}: {}",
1420                        idx + 1,
1421                        total_actions,
1422                        e
1423                    );
1424
1425                    let mut plans = self.pending_plans.write().await;
1426                    if let Some(plan) = plans.get_mut(&plan_id) {
1427                        plan.status = PlanStatus::Failed { error: e.clone() };
1428
1429                        let result = plan.clone();
1430                        if let Err(e) = self.save_plan(&result) {
1431                            warn!("Failed to persist failed plan {}: {e}", plan_id);
1432                        }
1433                        return Ok(result);
1434                    }
1435                    return Err("Plan disappeared during execution".to_string());
1436                }
1437            }
1438        }
1439
1440        // All actions completed successfully
1441        let mut plans = self.pending_plans.write().await;
1442        if let Some(plan) = plans.get_mut(&plan_id) {
1443            plan.status = PlanStatus::Completed;
1444            info!(
1445                "Plan {} completed successfully: {} actions executed",
1446                plan_id, completed_actions
1447            );
1448
1449            let result = plan.clone();
1450            if let Err(e) = self.save_plan(&result) {
1451                warn!("Failed to persist completed plan {}: {e}", plan_id);
1452            }
1453            return Ok(result);
1454        }
1455
1456        Err("Plan disappeared during execution".to_string())
1457    }
1458
1459    /// Reject a plan.
1460    pub async fn reject_plan(&self, plan_id: Uuid) -> Result<SemanticPlan, String> {
1461        let mut plans = self.pending_plans.write().await;
1462        let plan = plans
1463            .get_mut(&plan_id)
1464            .ok_or_else(|| "Plan not found".to_string())?;
1465
1466        if plan.status != PlanStatus::Pending {
1467            return Err(format!("Plan is not pending: {:?}", plan.status));
1468        }
1469
1470        info!("Rejecting plan: {}", plan_id);
1471        plan.status = PlanStatus::Rejected;
1472
1473        let result = plan.clone();
1474
1475        // Persist the status change
1476        if let Err(e) = self.save_plan(&result) {
1477            warn!("Failed to persist rejected plan {}: {e}", plan_id);
1478        }
1479
1480        Ok(result)
1481    }
1482
1483    /// Get cleanup analysis as JSON bytes (for FUSE read).
1484    pub async fn get_cleanup_json(&self) -> Vec<u8> {
1485        if let Some(analysis) = self.get_cleanup_analysis().await {
1486            serde_json::to_string_pretty(&analysis)
1487                .unwrap_or_else(|_| "{}".to_string())
1488                .into_bytes()
1489        } else {
1490            // Return a message indicating analysis hasn't been run
1491            let msg = serde_json::json!({
1492                "message": "No cleanup analysis available. Run analyze_cleanup first.",
1493                "hint": "Write any content to .semantic/.cleanup to trigger analysis"
1494            });
1495            serde_json::to_string_pretty(&msg)
1496                .unwrap_or_default()
1497                .into_bytes()
1498        }
1499    }
1500
1501    /// Get duplicate groups as JSON bytes (for FUSE read).
1502    pub async fn get_dedupe_json(&self) -> Vec<u8> {
1503        if let Some(groups) = self.get_duplicate_groups().await {
1504            serde_json::to_string_pretty(&groups)
1505                .unwrap_or_else(|_| "{}".to_string())
1506                .into_bytes()
1507        } else {
1508            let msg = serde_json::json!({
1509                "message": "No duplicate analysis available. Run find_duplicates first.",
1510                "hint": "Write any content to .semantic/.dedupe to trigger analysis"
1511            });
1512            serde_json::to_string_pretty(&msg)
1513                .unwrap_or_default()
1514                .into_bytes()
1515        }
1516    }
1517
1518    /// Get similar files result as JSON bytes (for FUSE read).
1519    pub async fn get_similar_json(&self) -> Vec<u8> {
1520        if let Some(result) = self.get_last_similar_result().await {
1521            serde_json::to_string_pretty(&result)
1522                .unwrap_or_else(|_| "{}".to_string())
1523                .into_bytes()
1524        } else {
1525            let msg = serde_json::json!({
1526                "message": "No similar files search performed yet.",
1527                "hint": "Write a file path to .semantic/.similar to find similar files"
1528            });
1529            serde_json::to_string_pretty(&msg)
1530                .unwrap_or_default()
1531                .into_bytes()
1532        }
1533    }
1534
1535    /// Get pending plans directory listing.
1536    pub async fn get_pending_plan_ids(&self) -> Vec<String> {
1537        self.pending_plans
1538            .read()
1539            .await
1540            .iter()
1541            .filter(|(_, p)| p.status == PlanStatus::Pending)
1542            .map(|(id, _)| id.to_string())
1543            .collect()
1544    }
1545
1546    /// Get a plan as JSON bytes (for FUSE read).
1547    pub async fn get_plan_json(&self, plan_id: &str) -> Vec<u8> {
1548        if let Ok(uuid) = Uuid::parse_str(plan_id)
1549            && let Some(plan) = self.get_plan(uuid).await
1550        {
1551            return serde_json::to_string_pretty(&plan)
1552                .unwrap_or_else(|_| "{}".to_string())
1553                .into_bytes();
1554        }
1555        let msg = serde_json::json!({
1556            "error": "Plan not found",
1557            "plan_id": plan_id
1558        });
1559        serde_json::to_string_pretty(&msg)
1560            .unwrap_or_default()
1561            .into_bytes()
1562    }
1563}
1564
1565/// Truncate content for preview.
1566fn truncate_content(content: &str, max_len: usize) -> String {
1567    if content.len() <= max_len {
1568        content.to_string()
1569    } else {
1570        format!("{}...", &content[..max_len])
1571    }
1572}
1573
1574/// Calculate cosine similarity between two embeddings.
1575fn cosine_similarity(a: &[f32], b: &[f32]) -> f32 {
1576    if a.len() != b.len() {
1577        return 0.0;
1578    }
1579
1580    let dot: f32 = a.iter().zip(b.iter()).map(|(x, y)| x * y).sum();
1581    let norm_a: f32 = a.iter().map(|x| x * x).sum::<f32>().sqrt();
1582    let norm_b: f32 = b.iter().map(|x| x * x).sum::<f32>().sqrt();
1583
1584    if norm_a == 0.0 || norm_b == 0.0 {
1585        return 0.0;
1586    }
1587
1588    (dot / (norm_a * norm_b)).clamp(-1.0, 1.0)
1589}
1590
1591#[cfg(test)]
1592mod tests {
1593    use super::*;
1594
1595    #[test]
1596    fn test_organize_request_serialization() {
1597        let request = OrganizeRequest {
1598            scope: PathBuf::from("docs/"),
1599            strategy: OrganizeStrategy::ByTopic,
1600            max_groups: 5,
1601            similarity_threshold: 0.8,
1602        };
1603
1604        let json = serde_json::to_string(&request).unwrap();
1605        let parsed: OrganizeRequest = serde_json::from_str(&json).unwrap();
1606
1607        assert_eq!(parsed.scope, request.scope);
1608        assert_eq!(parsed.max_groups, 5);
1609    }
1610
1611    #[test]
1612    fn test_organize_request_defaults() {
1613        let json = r#"{"scope":"src/","strategy":"by_topic"}"#;
1614        let request: OrganizeRequest = serde_json::from_str(json).unwrap();
1615
1616        assert_eq!(request.max_groups, 10);
1617        assert!((request.similarity_threshold - 0.7).abs() < f32::EPSILON);
1618    }
1619
1620    #[test]
1621    fn test_plan_status_serialization() {
1622        let status = PlanStatus::Failed {
1623            error: "test error".to_string(),
1624        };
1625        let json = serde_json::to_string(&status).unwrap();
1626        assert!(json.contains("failed"));
1627        assert!(json.contains("test error"));
1628    }
1629
1630    #[test]
1631    fn test_cleanup_reason_variants() {
1632        let duplicate = CleanupReason::Duplicate {
1633            similar_to: PathBuf::from("/original.txt"),
1634            similarity: 0.98,
1635        };
1636        let json = serde_json::to_string(&duplicate).unwrap();
1637        assert!(json.contains("duplicate"));
1638
1639        let stale = CleanupReason::Stale {
1640            last_accessed: Utc::now(),
1641        };
1642        let json = serde_json::to_string(&stale).unwrap();
1643        assert!(json.contains("stale"));
1644    }
1645
1646    #[test]
1647    fn test_semantic_config_default() {
1648        let config = SemanticConfig::default();
1649        assert!((config.duplicate_threshold - 0.95).abs() < f32::EPSILON);
1650        assert_eq!(config.similar_limit, 10);
1651        assert_eq!(config.plan_retention_hours, 24);
1652    }
1653
1654    #[test]
1655    fn test_truncate_content() {
1656        assert_eq!(truncate_content("short", 100), "short");
1657        assert_eq!(truncate_content("hello world", 5), "hello...");
1658    }
1659
1660    #[test]
1661    fn test_action_type_serialization() {
1662        let action = ActionType::Move {
1663            from: PathBuf::from("/old/path.txt"),
1664            to: PathBuf::from("/new/path.txt"),
1665        };
1666        let json = serde_json::to_string(&action).unwrap();
1667        assert!(json.contains("move"));
1668        assert!(json.contains("/old/path.txt"));
1669    }
1670
1671    #[test]
1672    fn test_similar_file_serialization() {
1673        let similar = SimilarFile {
1674            path: PathBuf::from("/doc.txt"),
1675            similarity: 0.85,
1676            preview: Some("This is a preview...".to_string()),
1677        };
1678        let json = serde_json::to_string(&similar).unwrap();
1679        assert!(json.contains("0.85"));
1680        assert!(json.contains("preview"));
1681    }
1682
1683    #[tokio::test]
1684    async fn test_semantic_manager_without_store() {
1685        let manager = SemanticManager::new(PathBuf::from("/tmp"), None, None, None);
1686        assert!(!manager.is_available());
1687    }
1688
1689    #[tokio::test]
1690    async fn test_pending_plans_empty() {
1691        let manager = SemanticManager::new(PathBuf::from("/tmp"), None, None, None);
1692        let plans = manager.list_pending_plans().await;
1693        assert!(plans.is_empty());
1694    }
1695
1696    #[tokio::test]
1697    async fn test_get_plan_not_found() {
1698        let manager = SemanticManager::new(PathBuf::from("/tmp"), None, None, None);
1699        let plan = manager.get_plan(Uuid::new_v4()).await;
1700        assert!(plan.is_none());
1701    }
1702
1703    #[tokio::test]
1704    async fn test_get_cleanup_json_empty() {
1705        let manager = SemanticManager::new(PathBuf::from("/tmp"), None, None, None);
1706        let json = manager.get_cleanup_json().await;
1707        let json_str = String::from_utf8(json).unwrap();
1708        assert!(json_str.contains("No cleanup analysis"));
1709    }
1710
1711    #[tokio::test]
1712    async fn test_get_dedupe_json_empty() {
1713        let manager = SemanticManager::new(PathBuf::from("/tmp"), None, None, None);
1714        let json = manager.get_dedupe_json().await;
1715        let json_str = String::from_utf8(json).unwrap();
1716        assert!(json_str.contains("No duplicate analysis"));
1717    }
1718
1719    #[tokio::test]
1720    async fn test_get_similar_json_empty() {
1721        let manager = SemanticManager::new(PathBuf::from("/tmp"), None, None, None);
1722        let json = manager.get_similar_json().await;
1723        let json_str = String::from_utf8(json).unwrap();
1724        assert!(json_str.contains("No similar files search"));
1725    }
1726
1727    #[tokio::test]
1728    async fn test_plan_by_custom_without_embedder() {
1729        // Test the fallback behavior when embedder is None
1730        let manager = SemanticManager::new(PathBuf::from("/tmp/test"), None, None, None);
1731        let mut file_embeddings = HashMap::new();
1732        file_embeddings.insert(PathBuf::from("/tmp/test/doc1.txt"), vec![0.1, 0.2, 0.3]);
1733        file_embeddings.insert(PathBuf::from("/tmp/test/doc2.txt"), vec![0.4, 0.5, 0.6]);
1734
1735        let scope_path = PathBuf::from("/tmp/test");
1736        let categories = vec!["code".to_string(), "docs".to_string()];
1737
1738        let (actions, description) = manager
1739            .plan_by_custom(&file_embeddings, &scope_path, &categories, None)
1740            .await;
1741
1742        // Should create 2 category directories but no file moves (no embedder)
1743        let mkdir_count = actions
1744            .iter()
1745            .filter(|a| matches!(a.action, ActionType::Mkdir { .. }))
1746            .count();
1747        let move_count = actions
1748            .iter()
1749            .filter(|a| matches!(a.action, ActionType::Move { .. }))
1750            .count();
1751
1752        assert_eq!(mkdir_count, 2, "Should create 2 category directories");
1753        assert_eq!(move_count, 0, "Should not move files without embedder");
1754        assert!(
1755            description.contains("manual assignment needed"),
1756            "Description should indicate manual assignment needed"
1757        );
1758    }
1759
1760    #[test]
1761    fn test_custom_categories_serialization() {
1762        let request = OrganizeRequest {
1763            scope: PathBuf::from("src/"),
1764            strategy: OrganizeStrategy::Custom {
1765                categories: vec!["code".to_string(), "docs".to_string(), "tests".to_string()],
1766            },
1767            max_groups: 10,
1768            similarity_threshold: 0.7,
1769        };
1770
1771        let json = serde_json::to_string(&request).unwrap();
1772        assert!(json.contains("custom"));
1773        assert!(json.contains("code"));
1774        assert!(json.contains("docs"));
1775        assert!(json.contains("tests"));
1776
1777        let parsed: OrganizeRequest = serde_json::from_str(&json).unwrap();
1778        if let OrganizeStrategy::Custom { categories } = parsed.strategy {
1779            assert_eq!(categories.len(), 3);
1780        } else {
1781            panic!("Expected Custom strategy");
1782        }
1783    }
1784
1785    #[tokio::test]
1786    async fn test_find_similar_rejects_escaped_path() {
1787        let temp = tempfile::TempDir::new().unwrap();
1788        let manager = SemanticManager::new(temp.path().to_path_buf(), None, None, None);
1789        let err = manager
1790            .find_similar(&PathBuf::from("../secret.txt"))
1791            .await
1792            .unwrap_err();
1793        assert!(err.contains("escapes"), "{err}");
1794    }
1795
1796    #[tokio::test]
1797    async fn test_find_similar_rejects_absolute_outside_root() {
1798        let temp = tempfile::TempDir::new().unwrap();
1799        let outside = tempfile::TempDir::new().unwrap();
1800        let manager = SemanticManager::new(temp.path().to_path_buf(), None, None, None);
1801        let err = manager
1802            .find_similar(&outside.path().join("other.txt"))
1803            .await
1804            .unwrap_err();
1805        assert!(err.contains("escapes"), "{err}");
1806    }
1807
1808    #[tokio::test]
1809    async fn test_create_organize_plan_rejects_escaped_scope() {
1810        let temp = tempfile::TempDir::new().unwrap();
1811        let manager = SemanticManager::new(temp.path().to_path_buf(), None, None, None);
1812        let request = OrganizeRequest {
1813            scope: PathBuf::from("../../etc"),
1814            strategy: OrganizeStrategy::ByTopic,
1815            max_groups: 5,
1816            similarity_threshold: 0.8,
1817        };
1818        let err = manager.create_organize_plan(request).await.unwrap_err();
1819        assert!(err.contains("escapes"), "{err}");
1820    }
1821
1822    #[test]
1823    fn test_canonical_or_owned_equalizes_dot_components() {
1824        let temp = tempfile::TempDir::new().unwrap();
1825        let file = temp.path().join("a.txt");
1826        std::fs::write(&file, "x").unwrap();
1827        let dotted = temp.path().join(".").join("a.txt");
1828        assert_eq!(
1829            SemanticManager::canonical_or_owned(&dotted),
1830            SemanticManager::canonical_or_owned(&file)
1831        );
1832        let scope = temp.path().join("docs");
1833        std::fs::create_dir(&scope).unwrap();
1834        let nested = scope.join("a.txt");
1835        std::fs::write(&nested, "x").unwrap();
1836        let dotted_nested = temp.path().join(".").join("docs").join("a.txt");
1837        assert!(
1838            SemanticManager::canonical_or_owned(&dotted_nested)
1839                .starts_with(SemanticManager::canonical_or_owned(&scope))
1840        );
1841    }
1842}