recall_echo/graph/
dedup.rs1use 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
31pub enum ResolvedEntity {
33 Created(Entity),
34 Merged(Entity),
35 Skipped,
36}
37
38pub async fn resolve_entity(
45 gm: &GraphMemory,
46 llm: &dyn LlmProvider,
47 candidate: &ExtractedEntity,
48 session_id: &str,
49) -> Result<ResolvedEntity, GraphError> {
50 let similar = gm.search(&candidate.abstract_text, 5).await?;
52
53 let relevant: Vec<_> = similar.iter().filter(|r| r.score > 0.7).collect();
55
56 if relevant.is_empty() {
57 let entity = gm.add_entity(candidate.to_new_entity(session_id)).await?;
59 return Ok(ResolvedEntity::Created(entity));
60 }
61
62 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 let target_entity = gm.get_entity(&target).await?;
81 let Some(target_entity) = target_entity else {
82 let entity = gm.add_entity(candidate.to_new_entity(session_id)).await?;
84 return Ok(ResolvedEntity::Created(entity));
85 };
86
87 if !target_entity.mutable {
89 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
100async 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
168pub 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 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}