Skip to main content

recall_echo/graph/
dedup.rs

1//! LLM-powered entity deduplication — skip, create, or merge decisions.
2
3use std::fmt::Write as _;
4
5use super::error::GraphError;
6use super::llm::LlmProvider;
7use super::types::*;
8use super::GraphMemory;
9
10const DEDUP_SYSTEM_PROMPT: &str = r#"You are a deduplication system for a knowledge graph. Given a candidate entity and existing similar entities, decide:
11
121. "skip" — The candidate is a duplicate. It adds no new information.
132. "create" — The candidate is genuinely new despite surface similarity.
143. "merge" — The candidate adds new information to an existing entity. Specify which one.
15
16Return EXACTLY this JSON (no markdown fencing, no explanation):
17
18{
19  "decision": "skip" | "create" | "merge",
20  "target": "Name of existing entity to merge into (only if merge)",
21  "reason": "Brief explanation"
22}
23
24Rules:
25- Same entity with minor name variations (e.g., "ElevenLabs" vs "Eleven Labs"): merge
26- Same concept but genuinely different instances: create
27- Candidate adds meaningful new detail to an existing entity: merge
28- Candidate is less detailed than existing: skip
29- When in doubt between create and merge: prefer create (avoid data loss)"#;
30
31/// Resolved entity after dedup — either newly created or existing (merged/skipped).
32pub enum ResolvedEntity {
33    Created(Entity),
34    Merged(Entity),
35    Skipped,
36}
37
38/// Run the full dedup pipeline for one extracted entity.
39///
40/// 1. Vector search for similar entities
41/// 2. If none similar: CREATE directly
42/// 3. If similar found: ask LLM for skip/create/merge decision
43/// 4. For merge on immutable types: fall back to CREATE
44pub async fn resolve_entity(
45    gm: &GraphMemory,
46    llm: &dyn LlmProvider,
47    candidate: &ExtractedEntity,
48    session_id: &str,
49) -> Result<ResolvedEntity, GraphError> {
50    // Search for similar entities
51    let similar = gm.search(&candidate.abstract_text, 5).await?;
52
53    // Filter to meaningful similarity (> 0.7 blended score)
54    let relevant: Vec<_> = similar.iter().filter(|r| r.score > 0.7).collect();
55
56    if relevant.is_empty() {
57        // No similar entities — create directly
58        let entity = gm.add_entity(candidate.to_new_entity(session_id)).await?;
59        return Ok(ResolvedEntity::Created(entity));
60    }
61
62    // Ask LLM for dedup decision
63    let user_message = build_dedup_message(candidate, &relevant);
64    let response = llm
65        .complete(DEDUP_SYSTEM_PROMPT, &user_message, 300)
66        .await?;
67
68    let decision = parse_dedup_response(&response)?;
69
70    match decision {
71        DedupDecision::Skip => Ok(ResolvedEntity::Skipped),
72
73        DedupDecision::Create => {
74            let entity = gm.add_entity(candidate.to_new_entity(session_id)).await?;
75            Ok(ResolvedEntity::Created(entity))
76        }
77
78        DedupDecision::Merge { target } => {
79            // Find the target entity
80            let target_entity = gm.get_entity(&target).await?;
81            let Some(target_entity) = target_entity else {
82                // Target not found — fall back to create
83                let entity = gm.add_entity(candidate.to_new_entity(session_id)).await?;
84                return Ok(ResolvedEntity::Created(entity));
85            };
86
87            // Check mutability
88            if !target_entity.mutable {
89                // Immutable — can't merge, create instead
90                let entity = gm.add_entity(candidate.to_new_entity(session_id)).await?;
91                return Ok(ResolvedEntity::Created(entity));
92            }
93
94            let merged = merge_entity(gm, &target_entity, candidate).await?;
95            Ok(ResolvedEntity::Merged(merged))
96        }
97    }
98}
99
100/// Merge candidate data into an existing entity.
101///
102/// Rules:
103/// - Abstract: use longer/more detailed version
104/// - Overview: concatenate if both exist
105/// - Content: append candidate content
106/// - Attributes: deep-merge (candidate wins on conflict)
107async fn merge_entity(
108    gm: &GraphMemory,
109    target: &Entity,
110    candidate: &ExtractedEntity,
111) -> Result<Entity, GraphError> {
112    let new_abstract = if candidate.abstract_text.len() > target.abstract_text.len() {
113        Some(candidate.abstract_text.clone())
114    } else {
115        None
116    };
117
118    let new_overview = candidate.overview.as_ref().map(|co| {
119        if target.overview.is_empty() {
120            co.clone()
121        } else {
122            format!("{}\n\n{}", target.overview, co)
123        }
124    });
125
126    let new_content = candidate.content.as_ref().map(|cc| match &target.content {
127        Some(tc) => format!("{tc}\n\n{cc}"),
128        None => cc.clone(),
129    });
130
131    let new_attributes = candidate
132        .attributes
133        .as_ref()
134        .map(|ca| match &target.attributes {
135            Some(ta) => merge_json_objects(ta, ca),
136            None => ca.clone(),
137        });
138
139    let updates = EntityUpdate {
140        abstract_text: new_abstract,
141        overview: new_overview,
142        content: new_content,
143        attributes: new_attributes,
144    };
145
146    gm.update_entity(&target.id_string(), updates).await
147}
148
149fn build_dedup_message(candidate: &ExtractedEntity, similar: &[&SearchResult]) -> String {
150    let mut msg = format!(
151        "CANDIDATE:\n  Name: {}\n  Type: {}\n  Abstract: {}\n\nEXISTING SIMILAR ENTITIES:\n",
152        candidate.name, candidate.entity_type, candidate.abstract_text
153    );
154    for (i, r) in similar.iter().enumerate() {
155        let _ = write!(
156            msg,
157            "\n{}. Name: {} (score: {:.3})\n   Type: {}\n   Abstract: {}\n",
158            i + 1,
159            r.entity.name,
160            r.score,
161            r.entity.entity_type,
162            r.entity.abstract_text
163        );
164    }
165    msg
166}
167
168/// Parse the LLM's dedup decision from JSON.
169pub fn parse_dedup_response(text: &str) -> Result<DedupDecision, GraphError> {
170    let cleaned = strip_markdown_fencing(text);
171
172    let v: serde_json::Value = serde_json::from_str(&cleaned).map_err(|e| {
173        // Try extracting JSON from surrounding text
174        if let Some(json_str) = extract_json_object(&cleaned) {
175            if let Ok(v) = serde_json::from_str::<serde_json::Value>(json_str) {
176                return parse_decision_value(&v)
177                    .err()
178                    .unwrap_or_else(|| GraphError::Parse(e.to_string()));
179            }
180        }
181        GraphError::Parse(format!("dedup response not valid JSON: {e}"))
182    })?;
183
184    parse_decision_value(&v)
185}
186
187fn parse_decision_value(v: &serde_json::Value) -> Result<DedupDecision, GraphError> {
188    let decision = v
189        .get("decision")
190        .and_then(|d| d.as_str())
191        .ok_or_else(|| GraphError::Parse("missing 'decision' field".into()))?;
192
193    match decision {
194        "skip" => Ok(DedupDecision::Skip),
195        "create" => Ok(DedupDecision::Create),
196        "merge" => {
197            let target = v
198                .get("target")
199                .and_then(|t| t.as_str())
200                .ok_or_else(|| GraphError::Parse("merge decision missing 'target' field".into()))?;
201            Ok(DedupDecision::Merge {
202                target: target.to_string(),
203            })
204        }
205        other => Err(GraphError::Parse(format!("unknown decision: {other}"))),
206    }
207}
208
209use super::util::{extract_json_object, merge_json_objects, strip_markdown_fencing};
210
211#[cfg(test)]
212mod tests {
213    use super::*;
214
215    #[test]
216    fn parse_skip_decision() {
217        let json = r#"{"decision": "skip", "reason": "duplicate"}"#;
218        let decision = parse_dedup_response(json).unwrap();
219        assert_eq!(decision, DedupDecision::Skip);
220    }
221
222    #[test]
223    fn parse_create_decision() {
224        let json = r#"{"decision": "create", "reason": "genuinely new"}"#;
225        let decision = parse_dedup_response(json).unwrap();
226        assert_eq!(decision, DedupDecision::Create);
227    }
228
229    #[test]
230    fn parse_merge_decision() {
231        let json = r#"{"decision": "merge", "target": "Rust", "reason": "same entity"}"#;
232        let decision = parse_dedup_response(json).unwrap();
233        assert_eq!(
234            decision,
235            DedupDecision::Merge {
236                target: "Rust".into()
237            }
238        );
239    }
240
241    #[test]
242    fn parse_with_fencing() {
243        let json = "```json\n{\"decision\": \"skip\", \"reason\": \"dup\"}\n```";
244        let decision = parse_dedup_response(json).unwrap();
245        assert_eq!(decision, DedupDecision::Skip);
246    }
247
248    #[test]
249    fn merge_json_objects_test() {
250        let base = serde_json::json!({"a": 1, "b": 2});
251        let overlay = serde_json::json!({"b": 3, "c": 4});
252        let merged = merge_json_objects(&base, &overlay);
253        assert_eq!(merged, serde_json::json!({"a": 1, "b": 3, "c": 4}));
254    }
255}