1use crate::analyzer::{AnalyzeCtx, Analyzer};
8use crate::analyzers::bound_evidence;
9use crate::cal;
10use crate::error::Result;
11use crate::manifest::*;
12use crate::model::{normalize_ident, ActionKind, GrainRecord, Severity};
13use crate::recommendation::{Proposal, RecDraft, Summary};
14use serde_json::{json, Map};
15use std::collections::BTreeMap;
16
17pub struct DuplicateSweep {
18 manifest: AnalyzerManifest,
19}
20
21impl DuplicateSweep {
22 pub fn new() -> Self {
23 DuplicateSweep {
24 manifest: AnalyzerManifest {
25 id: "loop.duplicate_sweep/1".into(),
26 title: "Duplicate sweep".into(),
27 description: "Consolidates exact-duplicate facts and near-duplicate observations."
28 .into(),
29 tier: Tier::T1,
30 cadence: CadenceClass::Batch,
31 requires: vec![],
32 target_classes: vec![TargetClass::Memory],
33 auto_apply: AutoApplyClass::StructuralCuration,
34 trust_class: TrustClass::Builtin,
35 params: vec![ParamSpec::Float {
36 name: "jaccard".into(),
37 default: 0.9,
38 min: 0.5,
39 max: 1.0,
40 description: "Near-duplicate token-set Jaccard threshold.".into(),
41 }],
42 default_on: true,
43 },
44 }
45 }
46}
47
48impl Default for DuplicateSweep {
49 fn default() -> Self {
50 Self::new()
51 }
52}
53
54impl Analyzer for DuplicateSweep {
55 fn manifest(&self) -> &AnalyzerManifest {
56 &self.manifest
57 }
58
59 fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
60 let mut drafts = self.exact_facts(ctx)?;
61 drafts.extend(self.near_observations(ctx)?);
62 drafts.sort_by(|a, b| a.target_ref.cmp(&b.target_ref));
63 Ok(drafts)
64 }
65}
66
67impl DuplicateSweep {
68 fn exact_facts(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
69 let facts = ctx.facts()?;
70 let mut groups: BTreeMap<(String, String, String, String), Vec<GrainRecord>> =
72 BTreeMap::new();
73 for f in facts {
74 let (Some(s), Some(r), Some(o)) =
75 (f.fact_subject(), f.fact_relation(), f.fact_object())
76 else {
77 continue;
78 };
79 let key = (
80 normalize_ident(&f.namespace),
81 normalize_ident(s),
82 normalize_ident(r),
83 normalize_ident(o),
84 );
85 groups.entry(key).or_default().push(f);
86 }
87
88 let mut drafts = Vec::new();
89 for ((_, subject, _, _), mut members) in groups {
90 if members.len() < 2 {
91 continue;
92 }
93 members.sort_by(|a, b| {
95 a.created_at_ms
96 .cmp(&b.created_at_ms)
97 .then(a.hash.cmp(&b.hash))
98 });
99 let canonical = members[0].clone();
100 let mut canonical_fields = Map::new();
101 canonical_fields.insert(
102 "subject".into(),
103 json!(canonical.fact_subject().unwrap_or("")),
104 );
105 canonical_fields.insert(
106 "relation".into(),
107 json!(canonical.fact_relation().unwrap_or("")),
108 );
109 canonical_fields.insert(
110 "object".into(),
111 json!(canonical.fact_object().unwrap_or("")),
112 );
113 if !canonical.namespace.is_empty() {
118 canonical_fields.insert("namespace".into(), json!(canonical.namespace));
119 }
120
121 let mut statements = Vec::new();
122 for extra in &members[1..] {
123 statements.push(cal::supersede(&extra.hash, "fact", &canonical_fields));
124 }
125 let evidence =
126 bound_evidence(members.iter().map(|m| m.hash.clone()).collect::<Vec<_>>());
127
128 let mut args = Map::new();
129 args.insert("count".into(), json!(members.len()));
130 args.insert("subject".into(), json!(subject));
131
132 drafts.push(
140 RecDraft::new(
141 format!(
142 "entity:{}/{}",
143 normalize_ident(&canonical.namespace),
144 subject
145 ),
146 ActionKind::Consolidate,
147 Summary::new("duplicate.exact", args),
148 Proposal::Cal {
149 cal: cal::batch(&statements),
150 },
151 )
152 .severity(Severity::Low)
153 .evidence(evidence),
154 );
155 }
156 Ok(drafts)
157 }
158
159 fn near_observations(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
160 let threshold = ctx.params().get_float("jaccard");
161 let obs = ctx.observations()?;
162
163 let mut tokenized: Vec<(GrainRecord, std::collections::BTreeSet<String>)> = obs
165 .into_iter()
166 .filter_map(|o| {
167 let tokens = tokenize(obs_text(&o)?);
168 Some((o, tokens))
169 })
170 .collect();
171 tokenized.sort_by(|a, b| {
172 a.0.created_at_ms
173 .cmp(&b.0.created_at_ms)
174 .then(a.0.hash.cmp(&b.0.hash))
175 });
176
177 let mut used = vec![false; tokenized.len()];
178 let mut drafts = Vec::new();
179 for i in 0..tokenized.len() {
180 if used[i] {
181 continue;
182 }
183 let mut cluster = vec![i];
184 for j in (i + 1)..tokenized.len() {
185 if used[j] || tokenized[i].0.namespace != tokenized[j].0.namespace {
186 continue;
187 }
188 if jaccard(&tokenized[i].1, &tokenized[j].1) >= threshold {
189 used[j] = true;
190 cluster.push(j);
191 }
192 }
193 if cluster.len() < 2 {
194 continue;
195 }
196 used[i] = true;
197 let canonical = &tokenized[cluster[0]].0;
198 let mut canonical_fields = Map::new();
199 canonical_fields.insert("body".into(), json!(obs_text(canonical).unwrap_or("")));
200 if !canonical.namespace.is_empty() {
203 canonical_fields.insert("namespace".into(), json!(canonical.namespace));
204 }
205
206 let mut statements = Vec::new();
207 for &k in &cluster[1..] {
208 statements.push(cal::supersede(
209 &tokenized[k].0.hash,
210 "observation",
211 &canonical_fields,
212 ));
213 }
214 let evidence = bound_evidence(
215 cluster
216 .iter()
217 .map(|&k| tokenized[k].0.hash.clone())
218 .collect(),
219 );
220
221 let mut args = Map::new();
222 args.insert("count".into(), json!(cluster.len()));
223 args.insert("threshold".into(), json!(threshold));
224
225 drafts.push(
226 RecDraft::new(
227 format!("grain:{}", canonical.hash),
228 ActionKind::Consolidate,
229 Summary::new("duplicate.near", args),
230 Proposal::Cal {
231 cal: cal::batch(&statements),
232 },
233 )
234 .severity(Severity::Info)
235 .evidence(evidence),
236 );
237 }
238 Ok(drafts)
239 }
240}
241
242fn obs_text(o: &GrainRecord) -> Option<&str> {
243 o.str_field("body")
244 .or_else(|| o.str_field("content"))
245 .or_else(|| o.str_field("text"))
246}
247
248pub(crate) fn tokenize(text: &str) -> std::collections::BTreeSet<String> {
249 text.to_lowercase()
250 .split(|c: char| !c.is_alphanumeric())
251 .filter(|t| !t.is_empty())
252 .map(|t| t.to_string())
253 .collect()
254}
255
256pub(crate) fn jaccard(a: &std::collections::BTreeSet<String>, b: &std::collections::BTreeSet<String>) -> f64 {
257 if a.is_empty() && b.is_empty() {
258 return 1.0;
259 }
260 let inter = a.intersection(b).count() as f64;
261 let union = a.union(b).count() as f64;
262 if union == 0.0 {
263 0.0
264 } else {
265 inter / union
266 }
267}
268
269#[cfg(test)]
270mod tests {
271 use super::*;
272 use crate::testkit::TestSubstrate;
273
274 #[test]
275 fn consolidates_exact_duplicate_facts() {
276 let mut sub = TestSubstrate::new();
277 sub.add_fact("caller", "tier", "Enterprise");
278 sub.add_fact("caller", "tier", "enterprise"); sub.add_fact("caller", "tier", "Enterprise");
280 let drafts = sub.analyze(&DuplicateSweep::new(), 10_000);
281 assert_eq!(drafts.len(), 1);
282 assert_eq!(drafts[0].action_kind, ActionKind::Consolidate);
283 assert_eq!(drafts[0].evidence.len(), 3);
284 }
285
286 #[test]
287 fn near_duplicate_observations_cluster() {
288 let mut sub = TestSubstrate::new();
289 sub.add_observation(
290 "caller",
291 "user asked about pricing tiers refunds billing invoices today",
292 );
293 sub.add_observation(
294 "caller",
295 "user asked about pricing tiers refunds billing invoices today please",
296 );
297 let drafts = sub.analyze(&DuplicateSweep::new(), 10_000);
299 assert_eq!(drafts.len(), 1, "the two near-dup observations cluster");
300 }
301
302 #[test]
303 fn distinct_facts_are_left_alone() {
304 let mut sub = TestSubstrate::new();
305 sub.add_fact("caller", "tier", "Enterprise");
306 sub.add_fact("caller", "tier", "Free");
307 assert!(sub.analyze(&DuplicateSweep::new(), 10_000).is_empty());
308 }
309}