Skip to main content

areev_loop/analyzers/
contradiction_sweep.rs

1//! Contradiction sweep (T0). Flags subjects holding two or more live objects
2//! under a *functional* relation (one that should be single-valued). Ships with
3//! a seeded functional-relation list so it fires on day one; the from-file
4//! learner (single-valued for ≥80% of subjects) is deferred. Resolving a
5//! contradiction is a judgment call, so it never auto-applies.
6//!
7//! **With a decision backend** (`Engine::with_decider`, proposal row E2): for
8//! every (namespace, subject, relation) OUTSIDE the functional set holding
9//! two or more distinct live objects, each pair of distinct values is asked
10//! "can both of these be true at the same time?" — batched, at most
11//! `pair_cap` pairs per run. When a CALIBRATED `1 − p ≥ DECIDE_MIN_P` for any
12//! pair, the same supersede-older draft is proposed over the values in the
13//! conflicting pairs (the newest of them wins), under the
14//! `contradiction.judged` summary that says the relation was not seeded and
15//! with the probabilities on `judged_by`. It carries no recurrence metric: a
16//! relation nobody declared single-valued may legitimately gain values later,
17//! and counting them would read as a regression. Uncalibrated → no proposals;
18//! a failed request → none from that batch.
19
20use crate::analyzer::{AnalyzeCtx, Analyzer};
21use crate::decide::{Ask, DECIDE_MIN_P, QUESTIONS_PER_REQUEST};
22use crate::analyzers::bound_evidence;
23use crate::cal;
24use crate::error::Result;
25use crate::manifest::*;
26use crate::model::{normalize_ident, ActionKind, GrainRecord, Severity};
27use crate::recommendation::{MetricSnapshot, Proposal, RecDraft, Summary};
28use serde_json::{json, Map, Value};
29use std::collections::BTreeMap;
30
31/// Relations that are single-valued by convention (a subset of the built-in
32/// `mg:` vocabulary plus common agent relations).
33const SEEDED_FUNCTIONAL: &[&str] = &[
34    "deploy_target",
35    "lives_in",
36    "reports_to",
37    "status",
38    "tier",
39    "owner",
40    "region",
41    "assigned_to",
42    "primary_email",
43    "current_plan",
44];
45
46pub struct ContradictionSweep {
47    manifest: AnalyzerManifest,
48}
49
50impl ContradictionSweep {
51    pub fn new() -> Self {
52        ContradictionSweep {
53            manifest: AnalyzerManifest {
54                id: "loop.contradiction_sweep/1".into(),
55                title: "Contradiction sweep".into(),
56                description: "Flags conflicting live values under functional relations.".into(),
57                tier: Tier::T0,
58                cadence: CadenceClass::Fast,
59                requires: vec![],
60                target_classes: vec![TargetClass::Memory],
61                auto_apply: AutoApplyClass::Never,
62                trust_class: TrustClass::Builtin,
63                params: vec![ParamSpec::Str {
64                    name: "extra_relations".into(),
65                    default: String::new(),
66                    max_len: 2000,
67                    description: "Additional functional (single-valued) relations to check, \
68                                  comma-separated — e.g. a healthcare deployment adds \
69                                  \"insurance_plan,prior_auth,next_appt\"."
70                        .into(),
71                }],
72                default_on: true,
73            },
74        }
75    }
76}
77
78impl Default for ContradictionSweep {
79    fn default() -> Self {
80        Self::new()
81    }
82}
83
84impl Analyzer for ContradictionSweep {
85    fn manifest(&self) -> &AnalyzerManifest {
86        &self.manifest
87    }
88
89    fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
90        // Seeded functional relations + any host-supplied domain relations.
91        let mut functional: std::collections::BTreeSet<String> =
92            SEEDED_FUNCTIONAL.iter().map(|s| s.to_string()).collect();
93        for r in ctx.params().get_str("extra_relations").split(',') {
94            let r = normalize_ident(r);
95            if !r.is_empty() {
96                functional.insert(r);
97            }
98        }
99
100        let facts = ctx.facts()?;
101        // (ns, subject, relation) → live facts, only for functional relations.
102        // The rest are kept aside for the decision backend (E2), if any.
103        let mut groups: BTreeMap<(String, String, String), Vec<GrainRecord>> = BTreeMap::new();
104        let mut other: BTreeMap<(String, String, String), Vec<GrainRecord>> = BTreeMap::new();
105        for f in facts {
106            let (Some(s), Some(r)) = (f.fact_subject(), f.fact_relation()) else {
107                continue;
108            };
109            let key = (
110                normalize_ident(&f.namespace),
111                normalize_ident(s),
112                normalize_ident(r),
113            );
114            if !functional.contains(&key.2) {
115                other.entry(key).or_default().push(f);
116                continue;
117            }
118            groups.entry(key).or_default().push(f);
119        }
120
121        let mut drafts = Vec::new();
122        for ((ns, subject, relation), mut members) in groups {
123            // Distinct live objects?
124            let distinct: std::collections::BTreeSet<String> = members
125                .iter()
126                .filter_map(|m| m.fact_object().map(normalize_ident))
127                .collect();
128            if distinct.len() < 2 {
129                continue;
130            }
131            // Resolve-to-latest: keep the newest, supersede the older values.
132            members.sort_by(|a, b| {
133                a.created_at_ms
134                    .cmp(&b.created_at_ms)
135                    .then(a.hash.cmp(&b.hash))
136            });
137            let latest = members.last().unwrap().clone();
138            let mut latest_fields = Map::new();
139            latest_fields.insert("subject".into(), json!(latest.fact_subject().unwrap_or("")));
140            latest_fields.insert(
141                "relation".into(),
142                json!(latest.fact_relation().unwrap_or("")),
143            );
144            latest_fields.insert("object".into(), json!(latest.fact_object().unwrap_or("")));
145            // The resolution supersedes older values with a NEW grain built
146            // from these fields — carry the namespace or the winning value
147            // would migrate to the store default namespace.
148            if !latest.namespace.is_empty() {
149                latest_fields.insert("namespace".into(), json!(latest.namespace));
150            }
151
152            let mut statements = Vec::new();
153            for older in &members[..members.len() - 1] {
154                statements.push(cal::supersede(&older.hash, "fact", &latest_fields));
155            }
156            let evidence = bound_evidence(members.iter().map(|m| m.hash.clone()).collect());
157
158            let mut args = Map::new();
159            args.insert("subject".into(), json!(subject));
160            args.insert("relation".into(), json!(relation));
161            args.insert("count".into(), json!(distinct.len()));
162
163            drafts.push(
164                RecDraft::new(
165                    format!("entity:{ns}/{subject}"),
166                    ActionKind::FlagContradiction,
167                    Summary::new("contradiction.functional", args),
168                    Proposal::Cal {
169                        cal: cal::batch(&statements),
170                    },
171                )
172                .severity(Severity::Medium)
173                .evidence(evidence)
174                .metric(MetricSnapshot {
175                    // After resolving to the latest value, does the subject
176                    // again hold ≥2 live values under this functional
177                    // relation? Baseline 0 = one live value; any excess at a
178                    // checkpoint is a regression → outcome review proposes a
179                    // revert for human judgment.
180                    metric: "contradiction_recurrence".into(),
181                    baseline: 0.0,
182                    unit: "count".into(),
183                    n: members.len() as u64,
184                    window: "live".into(),
185                    subject: Some(subject.clone()),
186                    namespace: (!ns.is_empty()).then(|| ns.clone()),
187                    relation: Some(relation.clone()),
188                    query: format!(
189                        "RECALL facts WHERE subject = \"{subject}\" AND relation = \"{relation}\" | COUNT DISTINCT object > 1"
190                    ),
191                    review_after_ms: 86_400_000,
192                    horizons_ms: vec![86_400_000, 7 * 86_400_000, 30 * 86_400_000],
193                    checkpoints: Vec::new(),
194                    // A count of excess live values: fewer is better.
195                    higher_is_better: false,
196                }),
197            );
198        }
199        drafts.extend(judged(ctx, other));
200        drafts.sort_by(|a, b| a.target_ref.cmp(&b.target_ref));
201        Ok(drafts)
202    }
203}
204
205/// The supersede-older CAL: every member but `latest` is superseded with
206/// `latest`'s value (namespace carried, or the winner would migrate to the
207/// store default namespace).
208fn supersede_older(members: &[GrainRecord], latest: &GrainRecord) -> String {
209    let mut latest_fields = Map::new();
210    latest_fields.insert("subject".into(), json!(latest.fact_subject().unwrap_or("")));
211    latest_fields.insert("relation".into(), json!(latest.fact_relation().unwrap_or("")));
212    latest_fields.insert("object".into(), json!(latest.fact_object().unwrap_or("")));
213    if !latest.namespace.is_empty() {
214        latest_fields.insert("namespace".into(), json!(latest.namespace));
215    }
216    let statements: Vec<String> = members
217        .iter()
218        .filter(|m| m.hash != latest.hash)
219        .map(|older| cal::supersede(&older.hash, "fact", &latest_fields))
220        .collect();
221    cal::batch(&statements)
222}
223
224fn fact_text(f: &GrainRecord) -> String {
225    format!(
226        "{} {} {}",
227        f.fact_subject().unwrap_or(""),
228        f.fact_relation().unwrap_or(""),
229        f.fact_object().unwrap_or("")
230    )
231}
232
233fn oldest_first(a: &GrainRecord, b: &GrainRecord) -> std::cmp::Ordering {
234    a.created_at_ms.cmp(&b.created_at_ms).then(a.hash.cmp(&b.hash))
235}
236
237/// E2: pairs of distinct values under relations nobody declared functional,
238/// asked of a CALIBRATED decision backend.
239fn judged(
240    ctx: &AnalyzeCtx,
241    other: BTreeMap<(String, String, String), Vec<GrainRecord>>,
242) -> Vec<RecDraft> {
243    let Some(d) = ctx.decider() else {
244        return Vec::new();
245    };
246    if !d.calibrated() {
247        return Vec::new();
248    }
249    // One representative grain per distinct normalized object (the newest),
250    // then every unordered pair of them, in deterministic order.
251    struct Pair {
252        group: usize,
253        a: GrainRecord,
254        b: GrainRecord,
255    }
256    let groups: Vec<((String, String, String), Vec<GrainRecord>)> = other.into_iter().collect();
257    let mut pairs: Vec<Pair> = Vec::new();
258    'groups: for (g, (_, members)) in groups.iter().enumerate() {
259        let mut by_object: BTreeMap<String, GrainRecord> = BTreeMap::new();
260        for m in members {
261            let Some(o) = m.fact_object() else { continue };
262            let slot = by_object.entry(normalize_ident(o)).or_insert_with(|| m.clone());
263            if oldest_first(slot, m).is_lt() {
264                *slot = m.clone();
265            }
266        }
267        if by_object.len() < 2 {
268            continue;
269        }
270        let reps: Vec<&GrainRecord> = by_object.values().collect();
271        for i in 0..reps.len() {
272            for j in (i + 1)..reps.len() {
273                if pairs.len() >= d.pair_cap() {
274                    break 'groups;
275                }
276                pairs.push(Pair { group: g, a: reps[i].clone(), b: reps[j].clone() });
277            }
278        }
279    }
280    // group → (conflicting grains by hash, p(cannot both be true) per pair,
281    // the provenance of the first answer that found one).
282    let backend = d.describe();
283    type Conflict = (BTreeMap<String, GrainRecord>, BTreeMap<String, f64>, crate::decide::Answered);
284    let mut conflicts: BTreeMap<usize, Conflict> = BTreeMap::new();
285    for chunk in pairs.chunks(QUESTIONS_PER_REQUEST) {
286        let mut state = Map::new();
287        let mut asks = Vec::new();
288        for (n, p) in chunk.iter().enumerate() {
289            let id = format!("p{n}");
290            state.insert(id.clone(), json!({"a": fact_text(&p.a), "b": fact_text(&p.b)}));
291            asks.push(Ask::Noul {
292                instructions: format!(
293                    "Can statements \"a\" and \"b\" of pair \"{id}\" (in state.pairs) both be true at the same time?"
294                ),
295                id,
296            });
297        }
298        // Fail-soft: a failed or uncalibrated answer contributes nothing.
299        let Ok(a) = d.ask(json!({ "pairs": Value::Object(state) }), &asks) else {
300            continue;
301        };
302        if !a.calibrated {
303            continue;
304        }
305        for (n, p) in chunk.iter().enumerate() {
306            let Some(&both) = a.noul.get(&format!("p{n}")) else { continue };
307            let p_no = 1.0 - both;
308            if p_no < DECIDE_MIN_P {
309                continue;
310            }
311            let entry = conflicts
312                .entry(p.group)
313                .or_insert_with(|| (BTreeMap::new(), BTreeMap::new(), a.clone()));
314            entry.0.insert(p.a.hash.clone(), p.a.clone());
315            entry.0.insert(p.b.hash.clone(), p.b.clone());
316            // Keyed by the pair itself: "not_both:<a>|<b>" → p(cannot both be true).
317            entry.1.insert(
318                format!(
319                    "not_both:{}|{}",
320                    p.a.fact_object().unwrap_or(""),
321                    p.b.fact_object().unwrap_or("")
322                ),
323                p_no,
324            );
325        }
326    }
327
328    let mut drafts = Vec::new();
329    for (g, (involved, ps, answered)) in conflicts {
330        let ((ns, subject, relation), _) = &groups[g];
331        let mut members: Vec<GrainRecord> = involved.into_values().collect();
332        members.sort_by(oldest_first);
333        let latest = members.last().expect("a conflicting pair has two members").clone();
334        let evidence = bound_evidence(members.iter().map(|m| m.hash.clone()).collect());
335        let p_max = ps.values().copied().fold(0.0_f64, f64::max);
336        let mut args = Map::new();
337        args.insert("subject".into(), json!(subject));
338        args.insert("relation".into(), json!(relation));
339        args.insert("count".into(), json!(members.len()));
340        args.insert("p".into(), json!((p_max * 1000.0).round() / 1000.0));
341        drafts.push(
342            RecDraft::new(
343                format!("entity:{ns}/{subject}"),
344                ActionKind::FlagContradiction,
345                Summary::new("contradiction.judged", args),
346                Proposal::Cal {
347                    cal: supersede_older(&members, &latest),
348                },
349            )
350            .severity(Severity::Medium)
351            .evidence(evidence)
352            .judged_by(answered.judged_by(&backend, "contradiction", ps)),
353        );
354    }
355    drafts
356}
357
358#[cfg(test)]
359mod tests {
360    use super::*;
361    use crate::testkit::TestSubstrate;
362
363    #[test]
364    fn flags_two_live_deploy_targets() {
365        let mut sub = TestSubstrate::new();
366        sub.add_fact("acme", "deploy_target", "us-east-1");
367        sub.add_fact("acme", "deploy_target", "eu-west-1");
368        let drafts = sub.analyze(&ContradictionSweep::new(), 10_000);
369        assert_eq!(drafts.len(), 1);
370        assert_eq!(drafts[0].action_kind, ActionKind::FlagContradiction);
371    }
372
373    #[test]
374    fn extra_relations_extend_the_functional_set() {
375        let mut sub = TestSubstrate::new();
376        sub.add_fact("bob", "insurance_plan", "aetna");
377        sub.add_fact("bob", "insurance_plan", "cigna"); // not in the seeded list
378        // Without the param, insurance_plan isn't treated as functional.
379        assert!(sub.analyze(&ContradictionSweep::new(), 10_000).is_empty());
380        // A healthcare deployment adds it.
381        let drafts = sub.analyze_with(
382            &ContradictionSweep::new(),
383            10_000,
384            &[("extra_relations", serde_json::json!("insurance_plan,prior_auth"))],
385        );
386        assert_eq!(drafts.len(), 1, "the custom functional relation is now checked");
387    }
388
389    #[test]
390    fn ignores_non_functional_relations() {
391        let mut sub = TestSubstrate::new();
392        sub.add_fact("acme", "likes", "pizza");
393        sub.add_fact("acme", "likes", "sushi"); // multi-valued relation — fine
394        assert!(sub.analyze(&ContradictionSweep::new(), 10_000).is_empty());
395    }
396}