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//!
7//! **With a decision backend** (`Engine::with_decider`, proposal row E2):
8//! observation pairs the Jaccard rule cannot decide — similarity in
9//! `[JUDGED_JACCARD_FLOOR, jaccard)`, same namespace, neither already
10//! clustered — are asked "do these two state the same claim?" in batches, at
11//! most `pair_cap` pairs per run. A CALIBRATED `p ≥ DECIDE_MIN_P` proposes
12//! the same supersede draft the Jaccard path does, with the probability on
13//! its `judged_by` record; an uncalibrated backend proposes nothing (a rank
14//! means nothing here), and a failed request contributes nothing.
15
16use crate::analyzer::{AnalyzeCtx, Analyzer};
17use crate::decide::{Ask, DECIDE_MIN_P, QUESTIONS_PER_REQUEST};
18use crate::analyzers::bound_evidence;
19use crate::cal;
20use crate::error::Result;
21use crate::manifest::*;
22use crate::model::{normalize_ident, ActionKind, GrainRecord, Severity};
23use crate::recommendation::{Proposal, RecDraft, Summary};
24use serde_json::{json, Map, Value};
25use std::collections::BTreeMap;
26
27/// The lowest token-set Jaccard at which a pair is worth asking a decision
28/// backend about. Below it two observations share too little wording for a
29/// "same claim" to be the likely reading, and asking would spend the pair cap
30/// on noise.
31pub const JUDGED_JACCARD_FLOOR: f64 = 0.5;
32
33pub struct DuplicateSweep {
34    manifest: AnalyzerManifest,
35}
36
37impl DuplicateSweep {
38    pub fn new() -> Self {
39        DuplicateSweep {
40            manifest: AnalyzerManifest {
41                id: "loop.duplicate_sweep/1".into(),
42                title: "Duplicate sweep".into(),
43                description: "Consolidates exact-duplicate facts and near-duplicate observations."
44                    .into(),
45                tier: Tier::T1,
46                cadence: CadenceClass::Batch,
47                requires: vec![],
48                target_classes: vec![TargetClass::Memory],
49                auto_apply: AutoApplyClass::StructuralCuration,
50                trust_class: TrustClass::Builtin,
51                params: vec![ParamSpec::Float {
52                    name: "jaccard".into(),
53                    default: 0.9,
54                    min: 0.5,
55                    max: 1.0,
56                    description: "Near-duplicate token-set Jaccard threshold.".into(),
57                }],
58                default_on: true,
59            },
60        }
61    }
62}
63
64impl Default for DuplicateSweep {
65    fn default() -> Self {
66        Self::new()
67    }
68}
69
70impl Analyzer for DuplicateSweep {
71    fn manifest(&self) -> &AnalyzerManifest {
72        &self.manifest
73    }
74
75    fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
76        let mut drafts = self.exact_facts(ctx)?;
77        drafts.extend(self.near_observations(ctx)?);
78        drafts.sort_by(|a, b| a.target_ref.cmp(&b.target_ref));
79        Ok(drafts)
80    }
81}
82
83impl DuplicateSweep {
84    fn exact_facts(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
85        let facts = ctx.facts()?;
86        // key = (ns, subject, relation, object) normalized.
87        let mut groups: BTreeMap<(String, String, String, String), Vec<GrainRecord>> =
88            BTreeMap::new();
89        for f in facts {
90            let (Some(s), Some(r), Some(o)) =
91                (f.fact_subject(), f.fact_relation(), f.fact_object())
92            else {
93                continue;
94            };
95            let key = (
96                normalize_ident(&f.namespace),
97                normalize_ident(s),
98                normalize_ident(r),
99                normalize_ident(o),
100            );
101            groups.entry(key).or_default().push(f);
102        }
103
104        let mut drafts = Vec::new();
105        for ((_, subject, _, _), mut members) in groups {
106            if members.len() < 2 {
107                continue;
108            }
109            // Canonical = earliest; supersede the rest.
110            members.sort_by(|a, b| {
111                a.created_at_ms
112                    .cmp(&b.created_at_ms)
113                    .then(a.hash.cmp(&b.hash))
114            });
115            let canonical = members[0].clone();
116            let mut canonical_fields = Map::new();
117            canonical_fields.insert(
118                "subject".into(),
119                json!(canonical.fact_subject().unwrap_or("")),
120            );
121            canonical_fields.insert(
122                "relation".into(),
123                json!(canonical.fact_relation().unwrap_or("")),
124            );
125            canonical_fields.insert(
126                "object".into(),
127                json!(canonical.fact_object().unwrap_or("")),
128            );
129            // The replacement must carry the original's namespace — a
130            // supersession builds a NEW grain from exactly these fields, so an
131            // absent namespace would silently move the fact to the store
132            // default namespace and out of every ns-scoped recall.
133            if !canonical.namespace.is_empty() {
134                canonical_fields.insert("namespace".into(), json!(canonical.namespace));
135            }
136
137            let mut statements = Vec::new();
138            for extra in &members[1..] {
139                statements.push(cal::supersede(&extra.hash, "fact", &canonical_fields));
140            }
141            let evidence =
142                bound_evidence(members.iter().map(|m| m.hash.clone()).collect::<Vec<_>>());
143
144            let mut args = Map::new();
145            args.insert("count".into(), json!(members.len()));
146            args.insert("subject".into(), json!(subject));
147
148            // No recurrence metric here (unlike contradiction_sweep): a
149            // supersession creates a NEW replacement grain, so post-apply the
150            // canonical + its copy both stay live and a live-grain count
151            // never drops — a grain-count metric would read "regressed" the
152            // moment it was applied. Head-based recall (`latest`) already
153            // returns one value; measuring duplicate recurrence honestly
154            // needs a supersede-by-existing substrate primitive first.
155            drafts.push(
156                RecDraft::new(
157                    format!(
158                        "entity:{}/{}",
159                        normalize_ident(&canonical.namespace),
160                        subject
161                    ),
162                    ActionKind::Consolidate,
163                    Summary::new("duplicate.exact", args),
164                    Proposal::Cal {
165                        cal: cal::batch(&statements),
166                    },
167                )
168                .severity(Severity::Low)
169                .evidence(evidence),
170            );
171        }
172        Ok(drafts)
173    }
174
175    fn near_observations(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
176        let threshold = ctx.params().get_float("jaccard");
177        let obs = ctx.observations()?;
178
179        // Greedy clustering within a namespace by token-set Jaccard.
180        let mut tokenized: Vec<(GrainRecord, std::collections::BTreeSet<String>)> = obs
181            .into_iter()
182            .filter_map(|o| {
183                let tokens = tokenize(obs_text(&o)?);
184                Some((o, tokens))
185            })
186            .collect();
187        tokenized.sort_by(|a, b| {
188            a.0.created_at_ms
189                .cmp(&b.0.created_at_ms)
190                .then(a.0.hash.cmp(&b.0.hash))
191        });
192
193        let mut used = vec![false; tokenized.len()];
194        let mut drafts = Vec::new();
195        for i in 0..tokenized.len() {
196            if used[i] {
197                continue;
198            }
199            let mut cluster = vec![i];
200            for j in (i + 1)..tokenized.len() {
201                if used[j] || tokenized[i].0.namespace != tokenized[j].0.namespace {
202                    continue;
203                }
204                if jaccard(&tokenized[i].1, &tokenized[j].1) >= threshold {
205                    used[j] = true;
206                    cluster.push(j);
207                }
208            }
209            if cluster.len() < 2 {
210                continue;
211            }
212            used[i] = true;
213            let mut args = Map::new();
214            args.insert("count".into(), json!(cluster.len()));
215            args.insert("threshold".into(), json!(threshold));
216            drafts.push(consolidate_draft(&tokenized, &cluster, Summary::new("duplicate.near", args)));
217        }
218        drafts.extend(judged_pairs(ctx, &tokenized, &mut used, threshold));
219        Ok(drafts)
220    }
221}
222
223/// The supersede-into-the-earliest draft over one cluster (indices into
224/// `tokenized`, earliest first) — shared by the Jaccard and judged paths so
225/// both propose exactly the same change.
226fn consolidate_draft(
227    tokenized: &[(GrainRecord, std::collections::BTreeSet<String>)],
228    cluster: &[usize],
229    summary: Summary,
230) -> RecDraft {
231    let canonical = &tokenized[cluster[0]].0;
232    let mut canonical_fields = Map::new();
233    canonical_fields.insert("body".into(), json!(obs_text(canonical).unwrap_or("")));
234    // Keep the cluster's namespace on the replacement (clusters never cross
235    // namespaces — both paths filter on it).
236    if !canonical.namespace.is_empty() {
237        canonical_fields.insert("namespace".into(), json!(canonical.namespace));
238    }
239    let mut statements = Vec::new();
240    for &k in &cluster[1..] {
241        statements.push(cal::supersede(&tokenized[k].0.hash, "observation", &canonical_fields));
242    }
243    let evidence = bound_evidence(cluster.iter().map(|&k| tokenized[k].0.hash.clone()).collect());
244    RecDraft::new(
245        format!("grain:{}", canonical.hash),
246        ActionKind::Consolidate,
247        summary,
248        Proposal::Cal {
249            cal: cal::batch(&statements),
250        },
251    )
252    .severity(Severity::Info)
253    .evidence(evidence)
254}
255
256/// E2: the pairs the Jaccard rule left undecided, asked of a CALIBRATED
257/// decision backend. Pairs are taken in (earlier, later) creation order, so
258/// the cap cuts deterministically. A judged pair whose members a previous
259/// judged pair already consumed is skipped (greedy, like the Jaccard path).
260fn judged_pairs(
261    ctx: &AnalyzeCtx,
262    tokenized: &[(GrainRecord, std::collections::BTreeSet<String>)],
263    used: &mut [bool],
264    threshold: f64,
265) -> Vec<RecDraft> {
266    let Some(d) = ctx.decider() else {
267        return Vec::new();
268    };
269    if !d.calibrated() {
270        return Vec::new();
271    }
272    let mut pairs: Vec<(usize, usize, f64)> = Vec::new();
273    'outer: for i in 0..tokenized.len() {
274        if used[i] {
275            continue;
276        }
277        for j in (i + 1)..tokenized.len() {
278            if used[j] || tokenized[i].0.namespace != tokenized[j].0.namespace {
279                continue;
280            }
281            let sim = jaccard(&tokenized[i].1, &tokenized[j].1);
282            if (JUDGED_JACCARD_FLOOR..threshold).contains(&sim) {
283                if pairs.len() >= d.pair_cap() {
284                    break 'outer;
285                }
286                pairs.push((i, j, sim));
287            }
288        }
289    }
290    let backend = d.describe();
291    let mut drafts = Vec::new();
292    for chunk in pairs.chunks(QUESTIONS_PER_REQUEST) {
293        let mut state = Map::new();
294        let mut asks = Vec::new();
295        for (n, &(i, j, _)) in chunk.iter().enumerate() {
296            let id = format!("p{n}");
297            state.insert(
298                id.clone(),
299                json!({"a": obs_text(&tokenized[i].0).unwrap_or(""), "b": obs_text(&tokenized[j].0).unwrap_or("")}),
300            );
301            asks.push(Ask::Noul {
302                instructions: format!(
303                    "Do texts \"a\" and \"b\" of pair \"{id}\" (in state.pairs) state the same claim?"
304                ),
305                id,
306            });
307        }
308        // Fail-soft: a failed or uncalibrated answer contributes nothing.
309        let Ok(a) = d.ask(json!({ "pairs": Value::Object(state) }), &asks) else {
310            continue;
311        };
312        if !a.calibrated {
313            continue;
314        }
315        for (n, &(i, j, sim)) in chunk.iter().enumerate() {
316            let id = format!("p{n}");
317            let Some(&p) = a.noul.get(&id) else { continue };
318            if p < DECIDE_MIN_P || used[i] || used[j] {
319                continue;
320            }
321            used[i] = true;
322            used[j] = true;
323            let mut args = Map::new();
324            args.insert("count".into(), json!(2));
325            args.insert("p".into(), json!(round3(p)));
326            args.insert("similarity".into(), json!(round3(sim)));
327            let judged = a.judged_by(&backend, "duplicate", BTreeMap::from([("same_claim".to_string(), p)]));
328            drafts.push(
329                consolidate_draft(tokenized, &[i, j], Summary::new("duplicate.judged", args))
330                    .judged_by(judged),
331            );
332        }
333    }
334    drafts
335}
336
337fn round3(x: f64) -> f64 {
338    (x * 1000.0).round() / 1000.0
339}
340
341fn obs_text(o: &GrainRecord) -> Option<&str> {
342    o.str_field("body")
343        .or_else(|| o.str_field("content"))
344        .or_else(|| o.str_field("text"))
345}
346
347pub(crate) fn tokenize(text: &str) -> std::collections::BTreeSet<String> {
348    text.to_lowercase()
349        .split(|c: char| !c.is_alphanumeric())
350        .filter(|t| !t.is_empty())
351        .map(|t| t.to_string())
352        .collect()
353}
354
355pub(crate) fn jaccard(a: &std::collections::BTreeSet<String>, b: &std::collections::BTreeSet<String>) -> f64 {
356    if a.is_empty() && b.is_empty() {
357        return 1.0;
358    }
359    let inter = a.intersection(b).count() as f64;
360    let union = a.union(b).count() as f64;
361    if union == 0.0 {
362        0.0
363    } else {
364        inter / union
365    }
366}
367
368#[cfg(test)]
369mod tests {
370    use super::*;
371    use crate::testkit::TestSubstrate;
372
373    #[test]
374    fn consolidates_exact_duplicate_facts() {
375        let mut sub = TestSubstrate::new();
376        sub.add_fact("caller", "tier", "Enterprise");
377        sub.add_fact("caller", "tier", "enterprise"); // case variant → same
378        sub.add_fact("caller", "tier", "Enterprise");
379        let drafts = sub.analyze(&DuplicateSweep::new(), 10_000);
380        assert_eq!(drafts.len(), 1);
381        assert_eq!(drafts[0].action_kind, ActionKind::Consolidate);
382        assert_eq!(drafts[0].evidence.len(), 3);
383    }
384
385    #[test]
386    fn near_duplicate_observations_cluster() {
387        let mut sub = TestSubstrate::new();
388        sub.add_observation(
389            "caller",
390            "user asked about pricing tiers refunds billing invoices today",
391        );
392        sub.add_observation(
393            "caller",
394            "user asked about pricing tiers refunds billing invoices today please",
395        );
396        // Superset differs by one token of eleven → Jaccard ≈ 0.91 ≥ 0.9.
397        let drafts = sub.analyze(&DuplicateSweep::new(), 10_000);
398        assert_eq!(drafts.len(), 1, "the two near-dup observations cluster");
399    }
400
401    #[test]
402    fn distinct_facts_are_left_alone() {
403        let mut sub = TestSubstrate::new();
404        sub.add_fact("caller", "tier", "Enterprise");
405        sub.add_fact("caller", "tier", "Free");
406        assert!(sub.analyze(&DuplicateSweep::new(), 10_000).is_empty());
407    }
408}