1use 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
31const 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 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 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 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 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 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 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 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
205fn 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
237fn 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 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 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 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 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"); assert!(sub.analyze(&ContradictionSweep::new(), 10_000).is_empty());
380 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"); assert!(sub.analyze(&ContradictionSweep::new(), 10_000).is_empty());
395 }
396}