use rusqlite::{params, Connection};
use serde::{Deserialize, Serialize};
use std::collections::{HashMap, HashSet};
use crate::error::Result;
pub const FORMULA_VERSION: i32 = 1;
pub const DEPENDENCE_WEIGHT_SOURCE: f64 = 0.5;
pub const DEPENDENCE_WEIGHT_PIPELINE: f64 = 0.3;
pub const DEPENDENCE_WEIGHT_SELF_GEN: f64 = 0.7;
const MAX_MODALITIES: f64 = 6.0;
pub mod state_status {
pub const FRESH: &str = "fresh";
pub const RECOMPUTING: &str = "recomputing";
pub const FAILED: &str = "failed";
pub const STALE_FORMULA: &str = "stale_formula";
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct MobilityState {
pub proposition_id: String,
pub regime: String,
pub snapshot_ts: f64,
pub support_mass: Option<f64>, pub attack_mass: Option<f64>, pub source_diversity: Option<f64>, pub effective_independence: Option<f64>, pub temporal_coherence: Option<f64>, pub transportability: Option<f64>, pub mutability: Option<f64>, pub load_bearingness: Option<f64>, pub modality_consilience: Option<f64>, pub self_gen_local: Option<f64>, pub self_gen_ancestral: Option<f64>, pub contamination_risk: Option<f64>, pub novelty_isolation: Option<f64>, pub tier_write_components: Vec<String>,
pub tier_read_components: Vec<String>,
pub tier_bg_components: Vec<String>,
pub formula_version: i32,
pub content_hash: String,
pub live_claim_count: i64,
pub state_status: String,
pub computed_at: i64,
}
#[derive(Debug, Clone)]
struct ClaimRow {
claim_id: String,
polarity: i32,
weight: f64,
extractor: String,
source_lineage: Vec<String>, self_generated: bool,
modality_signal: String,
source_memory_rid: Option<String>,
namespace: String,
valid_from: Option<f64>,
valid_to: Option<f64>,
}
impl crate::engine::YantrikDB {
pub fn compute_write_tier_mobility(
&self,
proposition_id: &str,
regime: &str,
) -> Result<MobilityState> {
let conn = self.conn.lock();
compute_write_tier_mobility_conn(&conn, proposition_id, regime)
}
pub fn get_mobility_state(
&self,
proposition_id: &str,
regime: &str,
) -> Result<Option<MobilityState>> {
let conn = self.conn.lock();
read_latest_state(&conn, proposition_id, regime)
}
pub fn upsert_mobility_state(&self, state: &MobilityState) -> Result<()> {
let conn = self.conn.lock();
upsert_mobility_state_inner(&conn, state)
}
pub fn compute_background_mobility(
&self,
proposition_id: &str,
regime: &str,
) -> Result<Option<MobilityState>> {
let conn = self.conn.lock();
let Some(mut state) = read_latest_state(&conn, proposition_id, regime)? else {
return Ok(None);
};
let tau = compute_temporal_coherence(&conn, proposition_id, regime)?;
let lambda = compute_load_bearingness(&conn, proposition_id)?;
let psi_a = compute_self_gen_ancestral(&conn, proposition_id, 2)?;
state.temporal_coherence = Some(tau);
state.load_bearingness = Some(lambda);
state.self_gen_ancestral = Some(psi_a);
for name in ["temporal_coherence", "load_bearingness", "self_gen_ancestral"] {
if !state.tier_bg_components.iter().any(|c| c == name) {
state.tier_bg_components.push(name.to_string());
}
}
upsert_mobility_state_inner(&conn, &state)?;
Ok(Some(state))
}
pub fn recompute_background_mobility_batch(&self, limit: usize) -> Result<usize> {
let pending = self.list_background_pending(limit)?;
let mut count = 0;
for (prop_id, regime) in pending {
if self.compute_background_mobility(&prop_id, ®ime)?.is_some() {
count += 1;
}
}
Ok(count)
}
fn list_background_pending(&self, limit: usize) -> Result<Vec<(String, String)>> {
let conn = self.conn.lock();
let mut stmt = conn.prepare(
"SELECT proposition_id, regime FROM mobility_state \
WHERE temporal_coherence IS NULL \
ORDER BY computed_at ASC LIMIT ?1",
)?;
let rows = stmt
.query_map(params![limit as i64], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
})?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(rows)
}
}
fn read_latest_state(
conn: &Connection,
proposition_id: &str,
regime: &str,
) -> Result<Option<MobilityState>> {
let mut stmt = conn.prepare(
"SELECT proposition_id, regime, snapshot_ts, \
support_mass, attack_mass, source_diversity, effective_independence, \
temporal_coherence, transportability, mutability, load_bearingness, \
modality_consilience, self_gen_local, self_gen_ancestral, \
contamination_risk, novelty_isolation, \
tier_write_components, tier_read_components, tier_bg_components, \
formula_version, content_hash, live_claim_count, state_status, computed_at \
FROM mobility_state \
WHERE proposition_id = ?1 AND regime = ?2 \
ORDER BY snapshot_ts DESC LIMIT 1",
)?;
let result = stmt.query_row(params![proposition_id, regime], row_to_mobility_state);
match result {
Ok(state) => Ok(Some(state)),
Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
Err(e) => Err(e.into()),
}
}
fn upsert_mobility_state_inner(conn: &Connection, state: &MobilityState) -> Result<()> {
let tier_write = serde_json::to_string(&state.tier_write_components)?;
let tier_read = serde_json::to_string(&state.tier_read_components)?;
let tier_bg = serde_json::to_string(&state.tier_bg_components)?;
conn.execute(
"INSERT INTO mobility_state (\
proposition_id, regime, snapshot_ts, \
support_mass, attack_mass, source_diversity, effective_independence, \
temporal_coherence, transportability, mutability, load_bearingness, \
modality_consilience, self_gen_local, self_gen_ancestral, \
contamination_risk, novelty_isolation, \
tier_write_components, tier_read_components, tier_bg_components, \
formula_version, content_hash, live_claim_count, state_status, computed_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, \
?17, ?18, ?19, ?20, ?21, ?22, ?23, ?24) \
ON CONFLICT(proposition_id, regime, snapshot_ts) DO UPDATE SET \
support_mass = excluded.support_mass, \
attack_mass = excluded.attack_mass, \
source_diversity = excluded.source_diversity, \
effective_independence = excluded.effective_independence, \
temporal_coherence = excluded.temporal_coherence, \
transportability = excluded.transportability, \
mutability = excluded.mutability, \
load_bearingness = excluded.load_bearingness, \
modality_consilience = excluded.modality_consilience, \
self_gen_local = excluded.self_gen_local, \
self_gen_ancestral = excluded.self_gen_ancestral, \
contamination_risk = excluded.contamination_risk, \
novelty_isolation = excluded.novelty_isolation, \
tier_write_components = excluded.tier_write_components, \
tier_read_components = excluded.tier_read_components, \
tier_bg_components = excluded.tier_bg_components, \
formula_version = excluded.formula_version, \
content_hash = excluded.content_hash, \
live_claim_count = excluded.live_claim_count, \
state_status = excluded.state_status, \
computed_at = excluded.computed_at",
params![
state.proposition_id, state.regime, state.snapshot_ts,
state.support_mass, state.attack_mass, state.source_diversity,
state.effective_independence, state.temporal_coherence,
state.transportability, state.mutability, state.load_bearingness,
state.modality_consilience, state.self_gen_local,
state.self_gen_ancestral, state.contamination_risk,
state.novelty_isolation,
tier_write, tier_read, tier_bg,
state.formula_version, state.content_hash, state.live_claim_count,
state.state_status, state.computed_at,
],
)?;
Ok(())
}
pub(super) fn compute_write_tier_mobility_conn(
conn: &Connection,
proposition_id: &str,
regime: &str,
) -> Result<MobilityState> {
let claims = fetch_claims_for_mobility(conn, proposition_id, regime)?;
let hash = content_hash(&claims);
if let Some(existing) = read_latest_state(conn, proposition_id, regime)? {
if existing.formula_version == FORMULA_VERSION
&& existing.content_hash == hash
&& existing.state_status == state_status::FRESH
{
return Ok(existing);
}
}
let support: Vec<&ClaimRow> = claims.iter().filter(|c| c.polarity == 1).collect();
let attack: Vec<&ClaimRow> = claims.iter().filter(|c| c.polarity == -1).collect();
let support_mass = accumulate_mass(&support);
let attack_mass = accumulate_mass(&attack);
let chi = compute_modality_consilience(&support);
let psi_l = compute_self_gen_local(&support);
let state = MobilityState {
proposition_id: proposition_id.to_string(),
regime: regime.to_string(),
snapshot_ts: crate::engine::now(),
support_mass: Some(support_mass),
attack_mass: Some(attack_mass),
modality_consilience: Some(chi),
self_gen_local: Some(psi_l),
tier_write_components: vec![
"support_mass".into(),
"attack_mass".into(),
"modality_consilience".into(),
"self_gen_local".into(),
],
formula_version: FORMULA_VERSION,
content_hash: hash,
live_claim_count: claims.len() as i64,
state_status: state_status::FRESH.to_string(),
computed_at: unix_seconds(),
..Default::default()
};
upsert_mobility_state_inner(conn, &state)?;
Ok(state)
}
fn fetch_claims_for_mobility(
conn: &Connection,
proposition_id: &str,
regime: &str,
) -> Result<Vec<ClaimRow>> {
let mut stmt = conn.prepare(
"SELECT claim_id, polarity, weight, extractor, source_lineage, self_generated, \
modality_signal, source_memory_rid, namespace, valid_from, valid_to \
FROM claims \
WHERE proposition_id = ?1 AND regime_tag = ?2 AND tombstoned = 0 \
ORDER BY claim_id ASC",
)?;
let rows = stmt
.query_map(params![proposition_id, regime], |row| {
let lineage_json: String = row.get(4)?;
let lineage = normalize_lineage(&lineage_json);
let self_gen_int: i64 = row.get(5)?;
Ok(ClaimRow {
claim_id: row.get(0)?,
polarity: row.get(1)?,
weight: row.get(2)?,
extractor: row.get(3)?,
source_lineage: lineage,
self_generated: self_gen_int != 0,
modality_signal: row.get(6)?,
source_memory_rid: row.get(7)?,
namespace: row.get(8)?,
valid_from: row.get(9)?,
valid_to: row.get(10)?,
})
})?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(rows)
}
fn normalize_lineage(lineage_json: &str) -> Vec<String> {
let parsed: Vec<String> = serde_json::from_str(lineage_json).unwrap_or_default();
let mut set: Vec<String> = parsed.into_iter().collect::<HashSet<_>>().into_iter().collect();
set.sort();
set
}
fn content_hash(claims: &[ClaimRow]) -> String {
let mut hasher = blake3::Hasher::new();
hasher.update(b"yantrikdb.warrant.v");
hasher.update(&FORMULA_VERSION.to_le_bytes());
hasher.update(b"\x00claims\x00");
let mut indexed: Vec<(&String, &ClaimRow)> = claims.iter().map(|c| (&c.claim_id, c)).collect();
indexed.sort_by(|a, b| a.0.cmp(b.0));
for (cid, c) in indexed {
hasher.update(cid.as_bytes());
hasher.update(b"|");
hasher.update(&c.polarity.to_le_bytes());
hasher.update(b"|");
hasher.update(&[c.self_generated as u8]);
hasher.update(b"|");
hasher.update(c.extractor.as_bytes());
hasher.update(b"|");
hasher.update(c.modality_signal.as_bytes());
hasher.update(b"|");
let w_scaled = (c.weight * 1000.0).round() as i64;
hasher.update(&w_scaled.to_le_bytes());
hasher.update(b"|");
for src in &c.source_lineage {
hasher.update(src.as_bytes());
hasher.update(b",");
}
hasher.update(b"\x00");
}
hex::encode(hasher.finalize().as_bytes())
}
fn unix_seconds() -> i64 {
use std::time::{SystemTime, UNIX_EPOCH};
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
pub(super) fn accumulate_mass(claims: &[&ClaimRow]) -> f64 {
compute_omegas(claims).iter().enumerate().map(|(k, w)| w * claims[k].weight).sum()
}
fn compute_omegas(claims: &[&ClaimRow]) -> Vec<f64> {
if claims.is_empty() {
return Vec::new();
}
let mut freq: HashMap<&str, usize> = HashMap::new();
for c in claims {
for e in &c.source_lineage {
*freq.entry(e.as_str()).or_insert(0) += 1;
}
}
let total_distinct = freq.len();
let mut omegas = Vec::with_capacity(claims.len());
for (k, claim) in claims.iter().enumerate() {
let d_k = leave_one_out_jaccard(&claim.source_lineage, &freq, total_distinct);
let p_k = pipeline_overlap_ratio(claim, claims, k);
let s_k = self_gen_overlap_binary(claim, claims, k);
let discount = 1.0
+ DEPENDENCE_WEIGHT_SOURCE * d_k
+ DEPENDENCE_WEIGHT_PIPELINE * p_k
+ DEPENDENCE_WEIGHT_SELF_GEN * s_k;
omegas.push(1.0 / discount);
}
omegas
}
fn leave_one_out_jaccard(
claim_set: &[String],
freq: &HashMap<&str, usize>,
total_distinct: usize,
) -> f64 {
if claim_set.is_empty() || total_distinct == 0 {
return 0.0;
}
let rest_union_is_empty = {
let claim_set_lookup: HashSet<&str> = claim_set.iter().map(|s| s.as_str()).collect();
freq.iter().all(|(e, c)| *c == 1 && claim_set_lookup.contains(e))
};
if rest_union_is_empty {
return 0.0;
}
let intersection = claim_set
.iter()
.filter(|e| freq.get(e.as_str()).copied().unwrap_or(0) > 1)
.count();
intersection as f64 / total_distinct as f64
}
fn pipeline_overlap_ratio(claim: &ClaimRow, claims: &[&ClaimRow], own_index: usize) -> f64 {
if claims.len() <= 1 {
return 0.0;
}
let matching = claims
.iter()
.enumerate()
.filter(|(i, c)| *i != own_index && c.extractor == claim.extractor)
.count();
matching as f64 / (claims.len() - 1) as f64
}
fn self_gen_overlap_binary(claim: &ClaimRow, claims: &[&ClaimRow], own_index: usize) -> f64 {
if !claim.self_generated {
return 0.0;
}
let has_other = claims
.iter()
.enumerate()
.any(|(i, c)| i != own_index && c.self_generated);
if has_other {
1.0
} else {
0.0
}
}
fn compute_modality_consilience(claims: &[&ClaimRow]) -> f64 {
if claims.is_empty() {
return 0.0;
}
let distinct: HashSet<&String> = claims.iter().map(|c| &c.modality_signal).collect();
(distinct.len() as f64 / MAX_MODALITIES).min(1.0)
}
fn compute_self_gen_local(claims: &[&ClaimRow]) -> f64 {
if claims.is_empty() {
return 0.0;
}
let self_gen_count = claims.iter().filter(|c| c.self_generated).count();
self_gen_count as f64 / claims.len() as f64
}
fn compute_temporal_coherence(
conn: &Connection,
proposition_id: &str,
regime: &str,
) -> Result<f64> {
let mut stmt = conn.prepare(
"SELECT polarity FROM claims \
WHERE proposition_id = ?1 AND regime_tag = ?2 \
ORDER BY created_at ASC, claim_id ASC",
)?;
let polarities: Vec<i32> = stmt
.query_map(params![proposition_id, regime], |row| row.get::<_, i32>(0))?
.collect::<std::result::Result<Vec<_>, _>>()?;
if polarities.len() < 2 {
return Ok(1.0);
}
let mut flips = 0;
for w in polarities.windows(2) {
if w[0] != w[1] {
flips += 1;
}
}
let max_flips = (polarities.len() - 1) as f64;
Ok(1.0 - (flips as f64 / max_flips))
}
fn compute_load_bearingness(conn: &Connection, proposition_id: &str) -> Result<f64> {
let count: i64 = conn.query_row(
"SELECT COUNT(DISTINCT mie.move_id) \
FROM move_input_edge mie \
INNER JOIN claims c ON mie.claim_id = c.claim_id \
WHERE c.proposition_id = ?1",
params![proposition_id],
|row| row.get(0),
)?;
Ok(count as f64)
}
fn compute_self_gen_ancestral(
conn: &Connection,
proposition_id: &str,
max_depth: u32,
) -> Result<f64> {
use std::collections::HashSet;
let seed_claims = fetch_proposition_live_claim_ids(conn, proposition_id)?;
if seed_claims.is_empty() {
return Ok(0.0);
}
let mut ancestry: HashSet<String> = HashSet::new();
let mut frontier: Vec<String> = seed_claims;
for _depth in 0..max_depth {
if frontier.is_empty() {
break;
}
let mut next: Vec<String> = Vec::new();
for claim_id in &frontier {
let moves = fetch_producing_moves(conn, claim_id)?;
for mv in &moves {
let inputs = fetch_move_input_claim_ids(conn, mv)?;
for inp in inputs {
if ancestry.insert(inp.clone()) {
next.push(inp);
}
}
}
}
frontier = next;
}
if ancestry.is_empty() {
return Ok(0.0);
}
let placeholders: String = (0..ancestry.len())
.map(|i| format!("?{}", i + 1))
.collect::<Vec<_>>()
.join(",");
let sql = format!(
"SELECT COUNT(*) FROM claims \
WHERE claim_id IN ({}) AND self_generated = 1",
placeholders
);
let params_vec: Vec<Box<dyn rusqlite::types::ToSql>> = ancestry
.iter()
.map(|c| Box::new(c.clone()) as Box<dyn rusqlite::types::ToSql>)
.collect();
let params_ref: Vec<&dyn rusqlite::types::ToSql> =
params_vec.iter().map(|b| b.as_ref()).collect();
let self_gen: i64 = conn.query_row(&sql, params_ref.as_slice(), |row| row.get(0))?;
Ok(self_gen as f64 / ancestry.len() as f64)
}
fn fetch_proposition_live_claim_ids(conn: &Connection, proposition_id: &str) -> Result<Vec<String>> {
let mut stmt = conn.prepare(
"SELECT claim_id FROM claims \
WHERE proposition_id = ?1 AND tombstoned = 0",
)?;
let rows = stmt
.query_map(params![proposition_id], |row| row.get::<_, String>(0))?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(rows)
}
fn fetch_producing_moves(conn: &Connection, claim_id: &str) -> Result<Vec<String>> {
let mut stmt = conn.prepare("SELECT move_id FROM move_output_edge WHERE claim_id = ?1")?;
let rows = stmt
.query_map(params![claim_id], |row| row.get::<_, String>(0))?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(rows)
}
fn fetch_move_input_claim_ids(conn: &Connection, move_id: &str) -> Result<Vec<String>> {
let mut stmt = conn.prepare("SELECT claim_id FROM move_input_edge WHERE move_id = ?1")?;
let rows = stmt
.query_map(params![move_id], |row| row.get::<_, String>(0))?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(rows)
}
pub const CONTEST_DERIVATION_VERSION: i32 = 1;
pub mod contest_flags {
pub const DUPLICATION_RISK: u64 = 1 << 0;
pub const SAME_SOURCE_CONFLICT: u64 = 1 << 1;
pub const REFERENT_HETEROGENEITY_PRESENT: u64 = 1 << 2;
pub const SAME_ARTIFACT_EXTRACTOR_CONFLICT: u64 = 1 << 3;
pub const PRESENT_TENSE_CONFLICT: u64 = 1 << 4;
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct ContestState {
pub proposition_id: String,
pub regime: String,
pub support_mass: f64,
pub attack_mass: f64,
pub support_effective_independence: f64,
pub attack_effective_independence: f64,
pub support_distinct_source_count: i64,
pub attack_distinct_source_count: i64,
pub same_source_opposite_polarity_count: i64,
pub same_artifact_extractor_polarity_conflict_count: i64,
pub temporal_overlap_conflict_count: i64,
pub temporal_separable_opposition_count: i64,
pub referent_schema_heterogeneity_count: i64,
pub heuristic_flags: u64,
pub derivation_version: i32,
pub content_hash: String,
pub live_claim_count: i64,
pub state_status: String,
pub computed_at: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct ContestConflictReport {
pub proposition_id: String,
pub regime: String,
pub heuristic_flags: u64,
pub same_source_opposite_polarity_pairs: Vec<(String, String)>,
pub same_artifact_extractor_conflict_pairs: Vec<(String, String)>,
pub temporal_overlap_conflict_pairs: Vec<(String, String)>,
}
impl crate::engine::YantrikDB {
pub fn compute_contest_state(
&self,
proposition_id: &str,
regime: &str,
) -> Result<ContestState> {
let conn = self.conn.lock();
compute_contest_state_conn(&conn, proposition_id, regime)
}
pub fn get_contest_state(
&self,
proposition_id: &str,
regime: &str,
) -> Result<Option<ContestState>> {
let conn = self.conn.lock();
read_contest_state(&conn, proposition_id, regime)
}
pub fn list_flagged_propositions(
&self,
flag_mask: u64,
limit: usize,
) -> Result<Vec<ContestState>> {
if flag_mask == 0 {
return Ok(Vec::new());
}
let conn = self.conn.lock();
let mut stmt = conn.prepare(
"SELECT proposition_id, regime, support_mass, attack_mass, \
support_effective_independence, attack_effective_independence, \
support_distinct_source_count, attack_distinct_source_count, \
same_source_opposite_polarity_count, \
same_artifact_extractor_polarity_conflict_count, \
temporal_overlap_conflict_count, temporal_separable_opposition_count, \
referent_schema_heterogeneity_count, heuristic_flags, \
derivation_version, content_hash, live_claim_count, state_status, computed_at \
FROM contest_state \
WHERE (heuristic_flags & ?1) != 0 \
ORDER BY computed_at DESC \
LIMIT ?2",
)?;
let rows = stmt
.query_map(params![flag_mask as i64, limit as i64], |row| {
let flags_int: i64 = row.get(13)?;
Ok(ContestState {
proposition_id: row.get(0)?,
regime: row.get(1)?,
support_mass: row.get(2)?,
attack_mass: row.get(3)?,
support_effective_independence: row.get(4)?,
attack_effective_independence: row.get(5)?,
support_distinct_source_count: row.get(6)?,
attack_distinct_source_count: row.get(7)?,
same_source_opposite_polarity_count: row.get(8)?,
same_artifact_extractor_polarity_conflict_count: row.get(9)?,
temporal_overlap_conflict_count: row.get(10)?,
temporal_separable_opposition_count: row.get(11)?,
referent_schema_heterogeneity_count: row.get(12)?,
heuristic_flags: flags_int as u64,
derivation_version: row.get(14)?,
content_hash: row.get(15)?,
live_claim_count: row.get(16)?,
state_status: row.get(17)?,
computed_at: row.get(18)?,
})
})?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(rows)
}
pub fn inspect_contest_conflicts(
&self,
proposition_id: &str,
regime: &str,
) -> Result<Option<ContestConflictReport>> {
let conn = self.conn.lock();
let Some(state) = read_contest_state(&conn, proposition_id, regime)? else {
return Ok(None);
};
let claims = fetch_claims_for_mobility(&conn, proposition_id, regime)?;
drop(conn);
let supports: Vec<&ClaimRow> = claims.iter().filter(|c| c.polarity == 1).collect();
let attacks: Vec<&ClaimRow> = claims.iter().filter(|c| c.polarity == -1).collect();
Ok(Some(ContestConflictReport {
proposition_id: proposition_id.to_string(),
regime: regime.to_string(),
heuristic_flags: state.heuristic_flags,
same_source_opposite_polarity_pairs: same_source_pairs(&supports, &attacks),
same_artifact_extractor_conflict_pairs: same_artifact_pairs(&supports, &attacks),
temporal_overlap_conflict_pairs: temporal_overlap_pairs(&supports, &attacks),
}))
}
}
fn same_source_pairs(supports: &[&ClaimRow], attacks: &[&ClaimRow]) -> Vec<(String, String)> {
let mut out = Vec::new();
for s in supports {
if s.source_lineage.is_empty() {
continue;
}
let s_key = lineage_key(&s.source_lineage);
for a in attacks {
if a.source_lineage.is_empty() {
continue;
}
if lineage_key(&a.source_lineage) == s_key {
out.push((s.claim_id.clone(), a.claim_id.clone()));
}
}
}
out
}
fn same_artifact_pairs(supports: &[&ClaimRow], attacks: &[&ClaimRow]) -> Vec<(String, String)> {
let mut out = Vec::new();
for s in supports {
let Some(s_rid) = &s.source_memory_rid else { continue };
for a in attacks {
if let Some(a_rid) = &a.source_memory_rid {
if a_rid == s_rid && a.extractor != s.extractor {
out.push((s.claim_id.clone(), a.claim_id.clone()));
}
}
}
}
out
}
fn temporal_overlap_pairs(supports: &[&ClaimRow], attacks: &[&ClaimRow]) -> Vec<(String, String)> {
let mut out = Vec::new();
for s in supports {
for a in attacks {
if !intervals_disjoint(s, a) {
out.push((s.claim_id.clone(), a.claim_id.clone()));
}
}
}
out
}
pub(super) fn compute_contest_state_conn(
conn: &Connection,
proposition_id: &str,
regime: &str,
) -> Result<ContestState> {
let claims = fetch_claims_for_mobility(conn, proposition_id, regime)?;
let hash = contest_content_hash(&claims);
if let Some(existing) = read_contest_state(conn, proposition_id, regime)? {
if existing.derivation_version == CONTEST_DERIVATION_VERSION
&& existing.content_hash == hash
&& existing.state_status == state_status::FRESH
{
return Ok(existing);
}
}
let prior_flags = read_contest_state(conn, proposition_id, regime)?
.map(|s| s.heuristic_flags)
.unwrap_or(0);
let state = derive_contest_state(proposition_id, regime, &claims, hash);
upsert_contest_state(conn, &state)?;
let newly_set = state.heuristic_flags & !prior_flags;
if newly_set & 0b0010 != 0 {
auto_file_adversarial_for_proposition(conn, proposition_id, "contradiction")?;
}
if newly_set & 0b1000 != 0 {
auto_file_adversarial_for_proposition(conn, proposition_id, "contradiction")?;
}
Ok(state)
}
fn auto_file_adversarial_for_proposition(
conn: &Connection,
proposition_id: &str,
discovered_via: &str,
) -> Result<()> {
let now_ts = crate::engine::now();
let mut stmt = conn.prepare(
"SELECT DISTINCT moe.move_id FROM move_output_edge moe \
INNER JOIN claims c ON moe.claim_id = c.claim_id \
WHERE c.proposition_id = ?1",
)?;
let moves: Vec<String> = stmt
.query_map(params![proposition_id], |row| row.get::<_, String>(0))?
.collect::<std::result::Result<Vec<_>, _>>()?;
drop(stmt);
for move_id in moves {
let already_exists: bool = conn
.query_row(
"SELECT 1 FROM move_adversarial_instance \
WHERE move_id = ?1 AND discovered_via = ?2 LIMIT 1",
params![move_id, discovered_via],
|_| Ok(true),
)
.unwrap_or(false);
if already_exists {
continue;
}
let instance_id = crate::id::new_id();
let root_cause = format!(
"auto-generated: proposition {} contest flag transitioned for move {}",
proposition_id, move_id
);
conn.execute(
"INSERT INTO move_adversarial_instance (\
instance_id, move_id, status, discovered_via, traced_root_cause, \
discovered_at, created_at) \
VALUES (?1, ?2, 'candidate', ?3, ?4, ?5, ?5)",
params![instance_id, move_id, discovered_via, root_cause, now_ts],
)?;
}
Ok(())
}
fn derive_contest_state(
proposition_id: &str,
regime: &str,
claims: &[ClaimRow],
content_hash: String,
) -> ContestState {
let supports: Vec<&ClaimRow> = claims.iter().filter(|c| c.polarity == 1).collect();
let attacks: Vec<&ClaimRow> = claims.iter().filter(|c| c.polarity == -1).collect();
let support_omegas = compute_omegas(&supports);
let attack_omegas = compute_omegas(&attacks);
let support_mass: f64 = support_omegas
.iter()
.enumerate()
.map(|(k, w)| w * supports[k].weight)
.sum();
let attack_mass: f64 = attack_omegas
.iter()
.enumerate()
.map(|(k, w)| w * attacks[k].weight)
.sum();
let support_effective_independence: f64 = support_omegas.iter().sum();
let attack_effective_independence: f64 = attack_omegas.iter().sum();
let support_distinct_source_count = distinct_source_count(&supports);
let attack_distinct_source_count = distinct_source_count(&attacks);
let same_source_opposite_polarity_count =
count_same_source_opposite_polarity(&supports, &attacks);
let same_artifact_extractor_polarity_conflict_count =
count_same_artifact_extractor_conflict(&supports, &attacks);
let (temporal_overlap_conflict_count, temporal_separable_opposition_count) =
count_temporal_conflicts(&supports, &attacks);
let referent_schema_heterogeneity_count = count_referent_heterogeneity(claims);
let mut flags: u64 = 0;
if support_mass > 2.0 && support_effective_independence < 2.0 {
flags |= contest_flags::DUPLICATION_RISK;
}
if same_source_opposite_polarity_count > 0 {
flags |= contest_flags::SAME_SOURCE_CONFLICT;
}
if referent_schema_heterogeneity_count > 0 {
flags |= contest_flags::REFERENT_HETEROGENEITY_PRESENT;
}
if same_artifact_extractor_polarity_conflict_count > 0 {
flags |= contest_flags::SAME_ARTIFACT_EXTRACTOR_CONFLICT;
}
if temporal_overlap_conflict_count > 0 {
flags |= contest_flags::PRESENT_TENSE_CONFLICT;
}
ContestState {
proposition_id: proposition_id.to_string(),
regime: regime.to_string(),
support_mass,
attack_mass,
support_effective_independence,
attack_effective_independence,
support_distinct_source_count,
attack_distinct_source_count,
same_source_opposite_polarity_count,
same_artifact_extractor_polarity_conflict_count,
temporal_overlap_conflict_count,
temporal_separable_opposition_count,
referent_schema_heterogeneity_count,
heuristic_flags: flags,
derivation_version: CONTEST_DERIVATION_VERSION,
content_hash,
live_claim_count: claims.len() as i64,
state_status: state_status::FRESH.to_string(),
computed_at: unix_seconds(),
}
}
fn distinct_source_count(claims: &[&ClaimRow]) -> i64 {
let mut set: HashSet<&str> = HashSet::new();
for c in claims {
for e in &c.source_lineage {
set.insert(e.as_str());
}
}
set.len() as i64
}
fn count_same_source_opposite_polarity(supports: &[&ClaimRow], attacks: &[&ClaimRow]) -> i64 {
if supports.is_empty() || attacks.is_empty() {
return 0;
}
let mut support_groups: HashMap<String, usize> = HashMap::new();
for s in supports {
if s.source_lineage.is_empty() {
continue; }
*support_groups.entry(lineage_key(&s.source_lineage)).or_insert(0) += 1;
}
let mut count: i64 = 0;
for a in attacks {
if a.source_lineage.is_empty() {
continue;
}
if let Some(n_support) = support_groups.get(&lineage_key(&a.source_lineage)) {
count += *n_support as i64;
}
}
count
}
fn count_same_artifact_extractor_conflict(
supports: &[&ClaimRow],
attacks: &[&ClaimRow],
) -> i64 {
if supports.is_empty() || attacks.is_empty() {
return 0;
}
let mut by_artifact: HashMap<&str, Vec<&str>> = HashMap::new();
for s in supports {
if let Some(rid) = &s.source_memory_rid {
by_artifact.entry(rid.as_str()).or_default().push(s.extractor.as_str());
}
}
let mut count: i64 = 0;
for a in attacks {
if let Some(rid) = &a.source_memory_rid {
if let Some(support_extractors) = by_artifact.get(rid.as_str()) {
for se in support_extractors {
if *se != a.extractor {
count += 1;
}
}
}
}
}
count
}
fn count_temporal_conflicts(supports: &[&ClaimRow], attacks: &[&ClaimRow]) -> (i64, i64) {
if supports.is_empty() || attacks.is_empty() {
return (0, 0);
}
let mut overlap = 0i64;
let mut separable = 0i64;
for s in supports {
for a in attacks {
if intervals_disjoint(s, a) {
separable += 1;
} else {
overlap += 1;
}
}
}
(overlap, separable)
}
fn intervals_disjoint(a: &ClaimRow, b: &ClaimRow) -> bool {
let a_fully_open = a.valid_from.is_none() && a.valid_to.is_none();
let b_fully_open = b.valid_from.is_none() && b.valid_to.is_none();
if a_fully_open && b_fully_open {
return false;
}
let a_from = a.valid_from.unwrap_or(f64::NEG_INFINITY);
let a_to = a.valid_to.unwrap_or(f64::INFINITY);
let b_from = b.valid_from.unwrap_or(f64::NEG_INFINITY);
let b_to = b.valid_to.unwrap_or(f64::INFINITY);
a_to < b_from || b_to < a_from
}
fn count_referent_heterogeneity(claims: &[ClaimRow]) -> i64 {
let distinct: HashSet<&str> = claims.iter().map(|c| c.namespace.as_str()).collect();
if distinct.len() > 1 {
distinct.len() as i64
} else {
0
}
}
fn lineage_key(lineage: &[String]) -> String {
lineage.join("\x01")
}
fn contest_content_hash(claims: &[ClaimRow]) -> String {
let mut hasher = blake3::Hasher::new();
hasher.update(b"yantrikdb.contest.v");
hasher.update(&CONTEST_DERIVATION_VERSION.to_le_bytes());
hasher.update(b"\x00claims\x00");
let mut indexed: Vec<(&String, &ClaimRow)> = claims.iter().map(|c| (&c.claim_id, c)).collect();
indexed.sort_by(|a, b| a.0.cmp(b.0));
for (cid, c) in indexed {
hasher.update(cid.as_bytes());
hasher.update(b"|");
hasher.update(&c.polarity.to_le_bytes());
hasher.update(b"|");
hasher.update(c.extractor.as_bytes());
hasher.update(b"|");
hasher.update(c.namespace.as_bytes());
hasher.update(b"|");
hasher.update(c.source_memory_rid.as_deref().unwrap_or("").as_bytes());
hasher.update(b"|");
let w_scaled = (c.weight * 1000.0).round() as i64;
hasher.update(&w_scaled.to_le_bytes());
hasher.update(b"|");
for src in &c.source_lineage {
hasher.update(src.as_bytes());
hasher.update(b",");
}
hasher.update(b"|");
let vf = c.valid_from.map(|t| (t * 1000.0).round() as i64).unwrap_or(i64::MIN);
let vt = c.valid_to.map(|t| (t * 1000.0).round() as i64).unwrap_or(i64::MAX);
hasher.update(&vf.to_le_bytes());
hasher.update(&vt.to_le_bytes());
hasher.update(b"\x00");
}
hex::encode(hasher.finalize().as_bytes())
}
fn read_contest_state(
conn: &Connection,
proposition_id: &str,
regime: &str,
) -> Result<Option<ContestState>> {
let mut stmt = conn.prepare(
"SELECT proposition_id, regime, support_mass, attack_mass, \
support_effective_independence, attack_effective_independence, \
support_distinct_source_count, attack_distinct_source_count, \
same_source_opposite_polarity_count, \
same_artifact_extractor_polarity_conflict_count, \
temporal_overlap_conflict_count, temporal_separable_opposition_count, \
referent_schema_heterogeneity_count, heuristic_flags, \
derivation_version, content_hash, live_claim_count, state_status, computed_at \
FROM contest_state \
WHERE proposition_id = ?1 AND regime = ?2",
)?;
let result = stmt.query_row(params![proposition_id, regime], |row| {
let flags_int: i64 = row.get(13)?;
Ok(ContestState {
proposition_id: row.get(0)?,
regime: row.get(1)?,
support_mass: row.get(2)?,
attack_mass: row.get(3)?,
support_effective_independence: row.get(4)?,
attack_effective_independence: row.get(5)?,
support_distinct_source_count: row.get(6)?,
attack_distinct_source_count: row.get(7)?,
same_source_opposite_polarity_count: row.get(8)?,
same_artifact_extractor_polarity_conflict_count: row.get(9)?,
temporal_overlap_conflict_count: row.get(10)?,
temporal_separable_opposition_count: row.get(11)?,
referent_schema_heterogeneity_count: row.get(12)?,
heuristic_flags: flags_int as u64,
derivation_version: row.get(14)?,
content_hash: row.get(15)?,
live_claim_count: row.get(16)?,
state_status: row.get(17)?,
computed_at: row.get(18)?,
})
});
match result {
Ok(state) => Ok(Some(state)),
Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
Err(e) => Err(e.into()),
}
}
fn upsert_contest_state(conn: &Connection, s: &ContestState) -> Result<()> {
conn.execute(
"INSERT INTO contest_state (\
proposition_id, regime, support_mass, attack_mass, \
support_effective_independence, attack_effective_independence, \
support_distinct_source_count, attack_distinct_source_count, \
same_source_opposite_polarity_count, \
same_artifact_extractor_polarity_conflict_count, \
temporal_overlap_conflict_count, temporal_separable_opposition_count, \
referent_schema_heterogeneity_count, heuristic_flags, \
derivation_version, content_hash, live_claim_count, state_status, computed_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19) \
ON CONFLICT(proposition_id, regime) DO UPDATE SET \
support_mass = excluded.support_mass, \
attack_mass = excluded.attack_mass, \
support_effective_independence = excluded.support_effective_independence, \
attack_effective_independence = excluded.attack_effective_independence, \
support_distinct_source_count = excluded.support_distinct_source_count, \
attack_distinct_source_count = excluded.attack_distinct_source_count, \
same_source_opposite_polarity_count = excluded.same_source_opposite_polarity_count, \
same_artifact_extractor_polarity_conflict_count = excluded.same_artifact_extractor_polarity_conflict_count, \
temporal_overlap_conflict_count = excluded.temporal_overlap_conflict_count, \
temporal_separable_opposition_count = excluded.temporal_separable_opposition_count, \
referent_schema_heterogeneity_count = excluded.referent_schema_heterogeneity_count, \
heuristic_flags = excluded.heuristic_flags, \
derivation_version = excluded.derivation_version, \
content_hash = excluded.content_hash, \
live_claim_count = excluded.live_claim_count, \
state_status = excluded.state_status, \
computed_at = excluded.computed_at",
params![
s.proposition_id, s.regime, s.support_mass, s.attack_mass,
s.support_effective_independence, s.attack_effective_independence,
s.support_distinct_source_count, s.attack_distinct_source_count,
s.same_source_opposite_polarity_count,
s.same_artifact_extractor_polarity_conflict_count,
s.temporal_overlap_conflict_count, s.temporal_separable_opposition_count,
s.referent_schema_heterogeneity_count,
s.heuristic_flags as i64,
s.derivation_version, s.content_hash, s.live_claim_count,
s.state_status, s.computed_at,
],
)?;
Ok(())
}
fn row_to_mobility_state(row: &rusqlite::Row) -> rusqlite::Result<MobilityState> {
let tier_write_json: String = row.get(16)?;
let tier_read_json: String = row.get(17)?;
let tier_bg_json: String = row.get(18)?;
Ok(MobilityState {
proposition_id: row.get(0)?,
regime: row.get(1)?,
snapshot_ts: row.get(2)?,
support_mass: row.get(3)?,
attack_mass: row.get(4)?,
source_diversity: row.get(5)?,
effective_independence: row.get(6)?,
temporal_coherence: row.get(7)?,
transportability: row.get(8)?,
mutability: row.get(9)?,
load_bearingness: row.get(10)?,
modality_consilience: row.get(11)?,
self_gen_local: row.get(12)?,
self_gen_ancestral: row.get(13)?,
contamination_risk: row.get(14)?,
novelty_isolation: row.get(15)?,
tier_write_components: serde_json::from_str(&tier_write_json).unwrap_or_default(),
tier_read_components: serde_json::from_str(&tier_read_json).unwrap_or_default(),
tier_bg_components: serde_json::from_str(&tier_bg_json).unwrap_or_default(),
formula_version: row.get(19)?,
content_hash: row.get(20)?,
live_claim_count: row.get(21)?,
state_status: row.get(22)?,
computed_at: row.get(23)?,
})
}
#[cfg(test)]
mod tests {
use super::*;
fn mk_claim(
claim_id: &str,
extractor: &str,
lineage: &[&str],
self_gen: bool,
modality: &str,
weight: f64,
) -> ClaimRow {
let mut normalized: Vec<String> = lineage.iter().map(|s| s.to_string()).collect();
normalized.sort();
normalized.dedup();
ClaimRow {
claim_id: claim_id.to_string(),
polarity: 1,
weight,
extractor: extractor.to_string(),
source_lineage: normalized,
self_generated: self_gen,
modality_signal: modality.to_string(),
source_memory_rid: None,
namespace: "default".to_string(),
valid_from: None,
valid_to: None,
}
}
#[test]
fn accumulate_single_claim_returns_weight() {
let c = mk_claim("c1", "ext_a", &["src_1"], false, "text", 1.0);
let refs = vec![&c];
let total = accumulate_mass(&refs);
assert!((total - 1.0).abs() < 1e-9);
}
#[test]
fn accumulate_independent_claims_scales_linearly() {
let c1 = mk_claim("c1", "ext_a", &["src_a1", "src_a2"], false, "text", 1.0);
let c2 = mk_claim("c2", "ext_b", &["src_b1"], false, "image", 1.0);
let c3 = mk_claim("c3", "ext_c", &["src_c1"], false, "numeric", 1.0);
let total = accumulate_mass(&[&c1, &c2, &c3]);
assert!(
(total - 3.0).abs() < 1e-9,
"independent claims should sum to raw weights, got {}",
total
);
}
#[test]
fn accumulate_duplicate_lineage_discounts() {
let c1 = mk_claim("c1", "ext_a", &["src_shared"], false, "text", 1.0);
let c2 = mk_claim("c2", "ext_a", &["src_shared"], false, "text", 1.0);
let c3 = mk_claim("c3", "ext_a", &["src_shared"], false, "text", 1.0);
let total = accumulate_mass(&[&c1, &c2, &c3]);
assert!(total < 2.0, "duplicate lineage should discount, got {}", total);
assert!(total > 1.5, "discount shouldn't be excessive, got {}", total);
}
#[test]
fn accumulate_self_generated_is_strongly_discounted() {
let c1 = mk_claim("c1", "self_reasoning", &["self_1"], true, "text", 1.0);
let c2 = mk_claim("c2", "self_reasoning", &["self_1"], true, "text", 1.0);
let total = accumulate_mass(&[&c1, &c2]);
assert!(total < 1.0, "self-generated duplicates should collapse; got {}", total);
}
#[test]
fn accumulate_is_order_invariant() {
let c1 = mk_claim("c1", "ext_a", &["src_a"], false, "text", 1.0);
let c2 = mk_claim("c2", "ext_b", &["src_b"], false, "image", 1.0);
let c3 = mk_claim("c3", "ext_c", &["src_c"], false, "numeric", 1.0);
let forward = accumulate_mass(&[&c1, &c2, &c3]);
let reverse = accumulate_mass(&[&c3, &c2, &c1]);
let shuffled = accumulate_mass(&[&c2, &c3, &c1]);
assert!((forward - reverse).abs() < 1e-9);
assert!((forward - shuffled).abs() < 1e-9);
}
#[test]
fn modality_consilience_tracks_distinct_modalities() {
let c1 = mk_claim("c1", "ext_a", &["src_1"], false, "text", 1.0);
let c2 = mk_claim("c2", "ext_a", &["src_2"], false, "text", 1.0);
let mono = compute_modality_consilience(&[&c1, &c2]);
let c3 = mk_claim("c3", "ext_a", &["src_3"], false, "image", 1.0);
let c4 = mk_claim("c4", "ext_a", &["src_4"], false, "numeric", 1.0);
let multi = compute_modality_consilience(&[&c1, &c3, &c4]);
assert!(multi > mono);
}
#[test]
fn self_gen_local_is_correct_ratio() {
let c1 = mk_claim("c1", "ext_a", &["src_1"], true, "text", 1.0);
let c2 = mk_claim("c2", "ext_a", &["src_2"], false, "text", 1.0);
let c3 = mk_claim("c3", "ext_a", &["src_3"], true, "text", 1.0);
let r = compute_self_gen_local(&[&c1, &c2, &c3]);
assert!((r - 2.0 / 3.0).abs() < 1e-9);
}
#[test]
fn content_hash_is_order_invariant() {
let c1 = mk_claim("c1", "ext_a", &["src_a"], false, "text", 1.0);
let c2 = mk_claim("c2", "ext_b", &["src_b"], false, "image", 1.0);
let c3 = mk_claim("c3", "ext_c", &["src_c"], false, "numeric", 1.0);
let h_forward = content_hash(&[c1.clone(), c2.clone(), c3.clone()]);
let h_reverse = content_hash(&[c3.clone(), c2.clone(), c1.clone()]);
assert_eq!(h_forward, h_reverse);
}
#[test]
fn content_hash_discriminates_on_lineage_change() {
let c1 = mk_claim("c1", "ext_a", &["src_a"], false, "text", 1.0);
let c2 = mk_claim("c1", "ext_a", &["src_b"], false, "text", 1.0);
assert_ne!(content_hash(&[c1]), content_hash(&[c2]));
}
#[test]
fn jaccard_empty_inputs_are_zero() {
let freq: HashMap<&str, usize> = HashMap::new();
assert_eq!(leave_one_out_jaccard(&[], &freq, 0), 0.0);
let claim = vec!["a".to_string()];
assert_eq!(leave_one_out_jaccard(&claim, &freq, 0), 0.0);
}
#[test]
fn jaccard_disjoint_is_zero() {
let mut freq: HashMap<&str, usize> = HashMap::new();
freq.insert("a", 1);
freq.insert("b", 1);
let j = leave_one_out_jaccard(&["a".to_string()], &freq, 2);
assert_eq!(j, 0.0);
}
#[test]
fn jaccard_fully_shared_is_one() {
let mut freq: HashMap<&str, usize> = HashMap::new();
freq.insert("a", 2);
let j = leave_one_out_jaccard(&["a".to_string()], &freq, 1);
assert_eq!(j, 1.0);
}
}