Skip to main content

areev_loop/analyzers/
duplicate_sweep.rs

1//! Duplicate sweep (T0/T1). Exact triple duplicates (NFC + case-fold) among
2//! Facts, and near-duplicate Observations by token-set Jaccard. Consolidation
3//! keeps the earliest member canonical and supersedes the rest — structural,
4//! non-destructive. (Exact duplicates are auto-apply *eligible*; near-dups fail
5//! the engine's exact-equality shape check and stay pending — §6.3.)
6
7use 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        // key = (ns, subject, relation, object) normalized.
71        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            // Canonical = earliest; supersede the rest.
94            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            // The replacement must carry the original's namespace — a
114            // supersession builds a NEW grain from exactly these fields, so an
115            // absent namespace would silently move the fact to the store
116            // default namespace and out of every ns-scoped recall.
117            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            // No recurrence metric here (unlike contradiction_sweep): a
133            // supersession creates a NEW replacement grain, so post-apply the
134            // canonical + its copy both stay live and a live-grain count
135            // never drops — a grain-count metric would read "regressed" the
136            // moment it was applied. Head-based recall (`latest`) already
137            // returns one value; measuring duplicate recurrence honestly
138            // needs a supersede-by-existing substrate primitive first.
139            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        // Greedy clustering within a namespace by token-set Jaccard.
164        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            // Keep the cluster's namespace on the replacement (clusters never
201            // cross namespaces — see the filter above).
202            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"); // case variant → same
279        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        // Superset differs by one token of eleven → Jaccard ≈ 0.91 ≥ 0.9.
298        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}