Skip to main content

mcpmem_core/
mutation.rs

1//! The single graph write boundary. Snapshots, graph writes and derived
2//! counters share the writer transaction; no change is published before commit.
3use std::collections::{BTreeMap, BTreeSet};
4
5use rusqlite::{Connection, OptionalExtension, params};
6use serde::{Deserialize, Serialize};
7use uuid::Uuid;
8
9use crate::errors::{MCSError, Result};
10use crate::graph::{GraphHandle, TxGuard, name_hash};
11use crate::types::{Entity, EntityInput, Observation, ObservationInput, Relation};
12
13pub type MutationError = MCSError;
14
15#[derive(Clone, Debug, Serialize, Deserialize)]
16#[serde(rename_all = "camelCase", deny_unknown_fields)]
17pub struct MutationContext {
18    pub actor: String,
19    pub origin: String,
20    pub correlation_id: Uuid,
21    pub causation_id: Option<Uuid>,
22    pub hop_count: u8,
23    pub idempotency_key: Option<String>,
24}
25
26impl MutationContext {
27    /// Trusted in-process legacy ingress. Network callers must supply a
28    /// context derived from their authenticated principal instead.
29    pub fn local() -> Self {
30        Self {
31            actor: "local".into(),
32            origin: "mcp".into(),
33            correlation_id: Uuid::new_v4(),
34            causation_id: None,
35            hop_count: 0,
36            idempotency_key: None,
37        }
38    }
39
40    pub fn validate(self) -> Result<Self> {
41        if self.actor.trim().is_empty()
42            || self.actor.len() > 256
43            || self.origin.trim().is_empty()
44            || self.origin.len() > 256
45            || self.actor.chars().any(char::is_control)
46            || self.origin.chars().any(char::is_control)
47            || self.correlation_id.is_nil()
48            || self.hop_count > 15
49            || self.causation_id.is_some_and(|id| id.is_nil())
50            || (self.hop_count > 0) != self.causation_id.is_some()
51            || self.idempotency_key.as_ref().is_some_and(|key| {
52                key.is_empty() || key.len() > 128 || key.chars().any(char::is_control)
53            })
54        {
55            return Err(MCSError::InvalidParams(
56                "Invalid mutation provenance".into(),
57            ));
58        }
59        Ok(self)
60    }
61}
62
63#[derive(Clone, Debug, Serialize, Deserialize)]
64#[serde(rename_all = "camelCase", deny_unknown_fields)]
65pub struct ObservationUpdate {
66    pub entity_name: String,
67    pub contents: Vec<ObservationInput>,
68}
69
70#[derive(Clone, Debug, Serialize, Deserialize)]
71#[serde(tag = "operation", rename_all = "snake_case", deny_unknown_fields)]
72pub enum MutationRequest {
73    CreateEntities {
74        entities: Vec<EntityInput>,
75    },
76    UpsertEntities {
77        entities: Vec<EntityInput>,
78    },
79    DeleteEntities {
80        names: Vec<String>,
81    },
82    CreateRelations {
83        relations: Vec<Relation>,
84    },
85    DeleteRelations {
86        relations: Vec<Relation>,
87    },
88    AddObservations {
89        observations: Vec<ObservationUpdate>,
90    },
91    DeleteObservations {
92        observations: Vec<ObservationUpdate>,
93    },
94    MergeEntities {
95        source: String,
96        target: String,
97    },
98    RenameEntity {
99        old_name: String,
100        new_name: String,
101    },
102    PurgeDefinedEntities {
103        name: String,
104    },
105    Compact,
106    Wipe,
107}
108
109#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
110#[serde(rename_all = "camelCase")]
111pub struct EntitySnapshot {
112    pub entity_id: i64,
113    pub name: String,
114    pub entity_type: String,
115    pub observations: Vec<Observation>,
116}
117
118impl EntitySnapshot {
119    pub fn entity(&self) -> Entity {
120        Entity {
121            name: self.name.clone(),
122            entity_type: self.entity_type.clone(),
123            observations: self.observations.clone(),
124        }
125    }
126}
127
128#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
129#[serde(rename_all = "lowercase")]
130pub enum ChangeOperation {
131    Create,
132    Update,
133    Delete,
134    Rename,
135}
136
137#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
138pub struct RelationDelta {
139    pub added: Vec<Relation>,
140    pub removed: Vec<Relation>,
141}
142
143#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
144#[serde(rename_all = "camelCase")]
145pub struct EntityChange {
146    pub operation: ChangeOperation,
147    pub before: Option<EntitySnapshot>,
148    pub after: Option<EntitySnapshot>,
149    pub relation_delta: Option<RelationDelta>,
150    #[serde(default, skip_serializing_if = "Option::is_none")]
151    pub old_name: Option<String>,
152    #[serde(default, skip_serializing_if = "Option::is_none")]
153    pub new_name: Option<String>,
154}
155
156#[derive(Clone, Debug, Serialize, Deserialize)]
157#[serde(rename_all = "camelCase")]
158pub struct CommittedChangeSet {
159    pub transaction_id: Uuid,
160    pub changes: Vec<EntityChange>,
161}
162
163#[derive(Debug, Serialize, Deserialize)]
164#[serde(rename_all = "camelCase")]
165pub struct ObservationResult {
166    pub entity_name: String,
167    pub added_observations: Vec<Observation>,
168}
169
170/// Legacy response data is captured inside the same transaction, preventing
171/// an adapter from returning a concurrent writer's later state.
172#[derive(Debug, Serialize, Deserialize)]
173pub enum MutationResult {
174    Entities(Vec<Entity>),
175    Relations(Vec<Relation>),
176    Observations(Vec<ObservationResult>),
177    Entity(Entity),
178    Count(usize),
179    Unit,
180}
181
182#[derive(Debug, Serialize, Deserialize)]
183pub struct MutationOutcome {
184    pub changes: CommittedChangeSet,
185    pub result: MutationResult,
186    pub replayed: bool,
187}
188
189pub struct MutationService<'a> {
190    graph: &'a GraphHandle,
191}
192
193impl<'a> MutationService<'a> {
194    pub const fn new(graph: &'a GraphHandle) -> Self {
195        Self { graph }
196    }
197
198    pub fn apply(
199        &self,
200        request: MutationRequest,
201        context: MutationContext,
202    ) -> Result<CommittedChangeSet> {
203        self.apply_with_result(request, context)
204            .map(|(changes, _)| changes)
205    }
206
207    pub fn apply_with_result(
208        &self,
209        request: MutationRequest,
210        context: MutationContext,
211    ) -> Result<(CommittedChangeSet, MutationResult)> {
212        if context.idempotency_key.is_some() {
213            return Err(MCSError::InvalidParams(
214                "idempotent ingress requires a raw request fingerprint".into(),
215            ));
216        }
217        self.apply_inner(request, context, None)
218            .map(|outcome| (outcome.changes, outcome.result))
219    }
220
221    pub fn apply_idempotent(
222        &self,
223        request: MutationRequest,
224        context: MutationContext,
225        fingerprint: &str,
226    ) -> Result<MutationOutcome> {
227        if context.idempotency_key.is_none()
228            || fingerprint.len() != 64
229            || !fingerprint
230                .bytes()
231                .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
232        {
233            return Err(MCSError::InvalidParams(
234                "idempotent ingress requires a key and SHA-256 request fingerprint".into(),
235            ));
236        }
237        self.apply_inner(request, context, Some(fingerprint))
238    }
239
240    fn apply_inner(
241        &self,
242        request: MutationRequest,
243        context: MutationContext,
244        fingerprint: Option<&str>,
245    ) -> Result<MutationOutcome> {
246        let context = context.validate()?;
247        let conn = self.graph.writer.lock();
248        let tx = TxGuard::begin(&conn)?;
249        if let (Some(key), Some(fingerprint)) = (&context.idempotency_key, fingerprint) {
250            let prior: Option<(String,String)> = conn.query_row("SELECT request_fingerprint,response FROM idempotency_record WHERE principal_id=?1 AND idempotency_key=?2", params![context.actor,key], |r| Ok((r.get(0)?,r.get(1)?))).optional().map_err(sql_error)?;
251            if let Some((saved_fingerprint, response)) = prior {
252                if saved_fingerprint != fingerprint {
253                    return Err(MCSError::InvalidParams("idempotency_conflict".into()));
254                }
255                let mut outcome: MutationOutcome = serde_json::from_str(&response)?;
256                outcome.replayed = true;
257                tx.commit()?;
258                return Ok(outcome);
259            }
260        }
261        if let Some(parent_id) = context.causation_id {
262            let parent = crate::events::EventRepository::new(&conn)
263                .get(parent_id)?
264                .ok_or_else(|| MCSError::InvalidParams("unknown causation event".into()))?;
265            if parent.provenance.correlation_id != context.correlation_id
266                || parent.provenance.hop_count.checked_add(1) != Some(context.hop_count)
267            {
268                return Err(MCSError::InvalidParams("invalid causation chain".into()));
269            }
270        }
271        self.graph.refresh_seqs(&conn)?;
272        let rename = match &request {
273            MutationRequest::RenameEntity { old_name, new_name } => {
274                Some((old_name.clone(), new_name.clone()))
275            }
276            _ => None,
277        };
278        let names = affected_names(&conn, &request)?;
279        let before = capture(&conn, &names)?;
280        let result = execute(self.graph, &conn, request)?;
281        let after = capture(&conn, &names)?;
282        let changes = match rename {
283            Some((old_name, new_name)) if old_name != new_name => {
284                rename_changes(&before, &after, &old_name, &new_name)
285            }
286            _ => effective_changes(&before, &after),
287        };
288        update_counters(&conn, &before, &after, &changes)?;
289        self.graph.sync_seqs(&conn)?;
290        let committed = CommittedChangeSet {
291            transaction_id: Uuid::new_v4(),
292            changes,
293        };
294        crate::events::persist_changes(&conn, &committed, &context)?;
295        let outcome = MutationOutcome {
296            changes: committed,
297            result,
298            replayed: false,
299        };
300        if let (Some(key), Some(fingerprint)) = (&context.idempotency_key, fingerprint) {
301            conn.execute(
302                "INSERT INTO idempotency_record VALUES(?1,?2,?3,?4,?5)",
303                params![
304                    context.actor,
305                    key,
306                    fingerprint,
307                    serde_json::to_string(&outcome)?,
308                    now_us()
309                ],
310            )
311            .map_err(sql_error)?;
312        }
313        tx.commit()?;
314        Ok(outcome)
315    }
316}
317
318fn sql_error(error: rusqlite::Error) -> MCSError {
319    MCSError::IoError(std::io::Error::other(error))
320}
321
322fn now_us() -> i64 {
323    std::time::SystemTime::now()
324        .duration_since(std::time::UNIX_EPOCH)
325        .unwrap_or_default()
326        .as_micros() as i64
327}
328
329pub(crate) fn read_entity(conn: &Connection, name: &str) -> Result<Option<EntitySnapshot>> {
330    let row = conn.query_row(
331        "SELECT e.id, e.name, t.name FROM entity e JOIN type_dict t ON t.id=e.type_id WHERE e.name_hash=?1 AND e.name=?2 AND e.flags=0",
332        params![name_hash(name), name],
333        |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?, row.get::<_, String>(2)?)),
334    ).optional().map_err(sql_error)?;
335    row.map(|(entity_id, name, entity_type)| {
336        let mut stmt = conn
337            .prepare_cached("SELECT body,created_us,occurred_us,origin_entity_name FROM observation WHERE entity_id=?1 ORDER BY idx, id")
338            .map_err(sql_error)?;
339        let observations = stmt
340            .query_map([entity_id], |row| Ok(Observation { body: row.get(0)?, created_at_us: Some(row.get(1)?), occurred_at_us: row.get(2)?, origin_entity_name: row.get(3)? }))
341            .map_err(sql_error)?
342            .collect::<rusqlite::Result<Vec<Observation>>>()
343            .map_err(sql_error)?;
344        Ok(EntitySnapshot {
345            entity_id,
346            name,
347            entity_type,
348            observations,
349        })
350    })
351    .transpose()
352}
353
354fn require_entity(conn: &Connection, name: &str) -> Result<EntitySnapshot> {
355    read_entity(conn, name)?
356        .ok_or_else(|| MCSError::InvalidParams(format!("Entity '{name}' not found")))
357}
358
359pub(crate) fn relations_for(conn: &Connection, name: &str) -> Result<Vec<Relation>> {
360    let mut stmt = conn.prepare_cached(
361        "SELECT f.name, t.name, d.name FROM relation r JOIN entity f ON f.id=r.from_id JOIN entity t ON t.id=r.to_id JOIN type_dict d ON d.id=r.type_id WHERE f.flags=0 AND t.flags=0 AND (r.from_id IN (SELECT id FROM entity WHERE name_hash=?1 AND name=?2 AND flags=0) OR r.to_id IN (SELECT id FROM entity WHERE name_hash=?1 AND name=?2 AND flags=0)) ORDER BY f.name, t.name, d.name"
362    ).map_err(sql_error)?;
363    stmt.query_map(params![name_hash(name), name], |row| {
364        Ok(Relation {
365            from: row.get(0)?,
366            to: row.get(1)?,
367            relation_type: row.get(2)?,
368        })
369    })
370    .map_err(sql_error)?
371    .collect::<rusqlite::Result<Vec<_>>>()
372    .map_err(sql_error)
373}
374
375fn defined_names(conn: &Connection, name: &str) -> Result<Vec<String>> {
376    let mut names: Vec<String> = relations_for(conn, name)?
377        .into_iter()
378        .filter(|r| r.from == name && r.relation_type == "defines")
379        .map(|r| r.to)
380        .collect();
381    names.push(name.into());
382    names.sort();
383    names.dedup();
384    Ok(names)
385}
386
387fn affected_names(conn: &Connection, request: &MutationRequest) -> Result<BTreeSet<String>> {
388    let mut names: BTreeSet<String> = match request {
389        MutationRequest::CreateEntities { entities }
390        | MutationRequest::UpsertEntities { entities } => {
391            entities.iter().map(|e| e.name.clone()).collect()
392        }
393        MutationRequest::DeleteEntities { names } => names.iter().cloned().collect(),
394        MutationRequest::CreateRelations { relations }
395        | MutationRequest::DeleteRelations { relations } => relations
396            .iter()
397            .flat_map(|r| [r.from.clone(), r.to.clone()])
398            .collect(),
399        MutationRequest::AddObservations { observations }
400        | MutationRequest::DeleteObservations { observations } => {
401            observations.iter().map(|o| o.entity_name.clone()).collect()
402        }
403        MutationRequest::MergeEntities { source, target } => {
404            [source.clone(), target.clone()].into()
405        }
406        MutationRequest::RenameEntity { old_name, new_name } => {
407            [old_name.clone(), new_name.clone()].into()
408        }
409        MutationRequest::PurgeDefinedEntities { name } => {
410            defined_names(conn, name)?.into_iter().collect()
411        }
412        MutationRequest::Compact => BTreeSet::new(),
413        MutationRequest::Wipe => {
414            let mut stmt = conn
415                .prepare("SELECT name FROM entity WHERE flags=0")
416                .map_err(sql_error)?;
417            stmt.query_map([], |row| row.get(0))
418                .map_err(sql_error)?
419                .collect::<rusqlite::Result<_>>()
420                .map_err(sql_error)?
421        }
422    };
423    // Deletes/merges change surviving neighbours too. Resolve these before
424    // deleting any rows so their before snapshots and relation deltas survive.
425    if matches!(
426        request,
427        MutationRequest::DeleteEntities { .. }
428            | MutationRequest::MergeEntities { .. }
429            | MutationRequest::RenameEntity { .. }
430            | MutationRequest::PurgeDefinedEntities { .. }
431    ) {
432        let neighbours = names
433            .iter()
434            .map(|name| relations_for(conn, name))
435            .collect::<Result<Vec<_>>>()?
436            .into_iter()
437            .flatten()
438            .flat_map(|r| [r.from, r.to])
439            .collect::<Vec<_>>();
440        names.extend(neighbours);
441    }
442    Ok(names)
443}
444
445#[derive(Default)]
446struct Snapshot {
447    entities: BTreeMap<String, EntitySnapshot>,
448    relations: BTreeSet<Relation>,
449    relation_rows: BTreeMap<Relation, i64>,
450}
451
452fn capture(conn: &Connection, names: &BTreeSet<String>) -> Result<Snapshot> {
453    let mut snapshot = Snapshot::default();
454    for name in names {
455        if let Some(entity) = read_entity(conn, name)? {
456            snapshot.entities.insert(name.clone(), entity);
457        }
458        let mut relation_rows = BTreeMap::new();
459        for relation in relations_for(conn, name)? {
460            *relation_rows.entry(relation).or_default() += 1;
461        }
462        // Both endpoint queries return every physical row of the same relation.
463        // Replace the count rather than adding it twice; keep set semantics for
464        // committed deltas independently of legacy duplicate storage rows.
465        snapshot.relations.extend(relation_rows.keys().cloned());
466        snapshot.relation_rows.extend(relation_rows);
467    }
468    Ok(snapshot)
469}
470
471fn effective_changes(before: &Snapshot, after: &Snapshot) -> Vec<EntityChange> {
472    let mut deltas: BTreeMap<&str, RelationDelta> = BTreeMap::new();
473    for (added, relations) in [
474        (true, after.relations.difference(&before.relations)),
475        (false, before.relations.difference(&after.relations)),
476    ] {
477        for relation in relations {
478            for name in [&relation.from, &relation.to]
479                .into_iter()
480                .collect::<BTreeSet<_>>()
481            {
482                let delta = deltas.entry(name).or_default();
483                if added {
484                    delta.added.push(relation.clone());
485                } else {
486                    delta.removed.push(relation.clone());
487                }
488            }
489        }
490    }
491    before
492        .entities
493        .keys()
494        .chain(after.entities.keys())
495        .collect::<BTreeSet<_>>()
496        .into_iter()
497        .filter_map(|name| {
498            let old = before.entities.get(name);
499            let new = after.entities.get(name);
500            let delta = deltas.remove(name.as_str()).unwrap_or_default();
501            let has_delta = !delta.added.is_empty() || !delta.removed.is_empty();
502            if old == new && !has_delta {
503                return None;
504            }
505            let operation = match (old, new) {
506                (None, Some(_)) => ChangeOperation::Create,
507                (Some(_), None) => ChangeOperation::Delete,
508                _ => ChangeOperation::Update,
509            };
510            Some(EntityChange {
511                operation,
512                before: old.cloned(),
513                after: new.cloned(),
514                relation_delta: has_delta.then_some(delta),
515                old_name: None,
516                new_name: None,
517            })
518        })
519        .collect()
520}
521
522fn rename_changes(
523    before: &Snapshot,
524    after: &Snapshot,
525    old_name: &str,
526    new_name: &str,
527) -> Vec<EntityChange> {
528    let (Some(before), Some(after)) = (before.entities.get(old_name), after.entities.get(new_name))
529    else {
530        return Vec::new();
531    };
532    vec![EntityChange {
533        operation: ChangeOperation::Rename,
534        before: Some(before.clone()),
535        after: Some(after.clone()),
536        relation_delta: None,
537        old_name: Some(old_name.into()),
538        new_name: Some(new_name.into()),
539    }]
540}
541
542fn type_id(conn: &Connection, name: &str, kind: i64) -> Result<i64> {
543    if let Some(id) = conn
544        .query_row(
545            "SELECT id FROM type_dict WHERE kind=?1 AND name=?2",
546            params![kind, name],
547            |r| r.get(0),
548        )
549        .optional()
550        .map_err(sql_error)?
551    {
552        return Ok(id);
553    }
554    conn.execute(
555        "INSERT INTO type_dict(kind,name,count) VALUES(?1,?2,0)",
556        params![kind, name],
557    )
558    .map_err(sql_error)?;
559    Ok(conn.last_insert_rowid())
560}
561
562/// Queue one taxonomy subject for every serving profile, using the same
563/// lookup as the entity job path. Without a serving profile the job is
564/// explicitly held by the queue function itself.
565fn enqueue_taxonomy_jobs(
566    conn: &Connection,
567    kind: i64,
568    id: i64,
569    revision: i64,
570    operation: crate::jobs::IndexOperation,
571) -> Result<()> {
572    for profile_id in crate::jobs::serving_profile_ids(conn)? {
573        crate::jobs::enqueue_taxonomy(conn, kind, id, revision, operation, profile_id)?;
574    }
575    Ok(())
576}
577
578/// Tombstone the taxonomy mirror of one deleted relation triple and enqueue
579/// the delete. A missing mirror is a legacy row: insert it tombstoned.
580fn tombstone_relation_mirror(
581    conn: &Connection,
582    from_id: i64,
583    to_id: i64,
584    type_id: i64,
585) -> Result<()> {
586    let (id, revision): (i64, i64) = conn
587        .query_row(
588            "INSERT INTO taxonomy_relation(from_id,to_id,type_id,revision,deleted) VALUES(?1,?2,?3,1,1) \
589             ON CONFLICT(from_id,to_id,type_id) DO UPDATE SET revision=revision+1,deleted=1 \
590             RETURNING id,revision",
591            params![from_id, to_id, type_id],
592            |row| Ok((row.get(0)?, row.get(1)?)),
593        )
594        .map_err(sql_error)?;
595    crate::jobs::enqueue_chunk_change(conn, crate::jobs::OwnerKind::Relation, id, revision, true)
596}
597
598fn insert_observations(
599    graph: &GraphHandle,
600    conn: &Connection,
601    id: i64,
602    contents: &[ObservationInput],
603) -> Result<Vec<Observation>> {
604    let idx: i64 = conn
605        .query_row(
606            "SELECT COALESCE(MAX(idx),-1) FROM observation WHERE entity_id=?1",
607            [id],
608            |r| r.get(0),
609        )
610        .map_err(sql_error)?;
611    let mut stmt = conn
612        .prepare_cached(
613            "INSERT INTO observation(id,entity_id,idx,body,created_us,occurred_us) VALUES(?1,?2,?3,?4,?5,?6)",
614        )
615        .map_err(sql_error)?;
616    let mut inserted = Vec::with_capacity(contents.len());
617    for (offset, observation) in contents.iter().enumerate() {
618        if observation.occurred_at_us.is_some_and(|time| time < 0) {
619            return Err(MCSError::InvalidParams(
620                "occurredAtUs must be non-negative".into(),
621            ));
622        }
623        let created_at_us = now_us();
624        stmt.execute(params![
625            graph.next_obs_id(),
626            id,
627            idx + offset as i64 + 1,
628            observation.body,
629            created_at_us,
630            observation.occurred_at_us
631        ])
632        .map_err(sql_error)?;
633        inserted.push(Observation {
634            body: observation.body.clone(),
635            created_at_us: Some(created_at_us),
636            occurred_at_us: observation.occurred_at_us,
637            origin_entity_name: None,
638        });
639    }
640    Ok(inserted)
641}
642
643fn create_entity(graph: &GraphHandle, conn: &Connection, entity: &EntityInput) -> Result<bool> {
644    if entity.name.is_empty() || read_entity(conn, &entity.name)?.is_some() {
645        return Ok(false);
646    }
647    let id = graph.next_entity_id();
648    let kind = type_id(conn, &entity.entity_type, 0)?;
649    conn.execute("INSERT INTO entity(id,name_hash,name,type_id,obs_count,out_deg,in_deg,created_us,updated_us,flags) VALUES(?1,?2,?3,?4,0,0,0,?5,?5,0)", params![id,name_hash(&entity.name),entity.name,kind,now_us()]).map_err(sql_error)?;
650    insert_observations(graph, conn, id, &entity.observations)?;
651    conn.execute(
652        "INSERT INTO name_fts(rowid,name) VALUES(?1,?2)",
653        params![id, entity.name],
654    )
655    .map_err(sql_error)?;
656    Ok(true)
657}
658
659fn create_relation(conn: &Connection, relation: &Relation) -> Result<bool> {
660    let (Some(from), Some(to)) = (
661        read_entity(conn, &relation.from)?,
662        read_entity(conn, &relation.to)?,
663    ) else {
664        return Ok(false);
665    };
666    let kind = type_id(conn, &relation.relation_type, 1)?;
667    let changed = conn.execute("INSERT INTO relation(from_id,to_id,type_id,created_us) SELECT ?1,?2,?3,?4 WHERE NOT EXISTS(SELECT 1 FROM relation WHERE from_id=?1 AND to_id=?2 AND type_id=?3)", params![from.entity_id,to.entity_id,kind,now_us()]).map_err(sql_error)?;
668    if changed > 0 {
669        // Mirror the triple for the chunk worker. The mirror id is its own
670        // autoincrement, never the source rowid: SQLite reuses a freed rowid
671        // for a later row, and an explicit-id insert would collide with the
672        // tombstoned mirror of a different triple. A recreated triple keeps
673        // the id its mirror already owns.
674        let mirror_id: i64 = conn
675            .query_row(
676                "INSERT INTO taxonomy_relation(from_id,to_id,type_id,revision,deleted) VALUES(?1,?2,?3,1,0) \
677                 ON CONFLICT(from_id,to_id,type_id) DO UPDATE SET revision=1,deleted=0 \
678                 RETURNING id",
679                params![from.entity_id, to.entity_id, kind],
680                |row| row.get(0),
681            )
682            .map_err(sql_error)?;
683        crate::jobs::enqueue_chunk_change(
684            conn,
685            crate::jobs::OwnerKind::Relation,
686            mirror_id,
687            1,
688            false,
689        )?;
690    }
691    Ok(changed > 0)
692}
693
694fn delete_entities(conn: &Connection, names: &[String]) -> Result<()> {
695    for name in names.iter().collect::<BTreeSet<_>>() {
696        if let Some(entity) = read_entity(conn, name)? {
697            let triples = conn
698                .prepare_cached(
699                    "SELECT from_id, to_id, type_id FROM relation WHERE from_id=?1 OR to_id=?1",
700                )
701                .map_err(sql_error)?
702                .query_map([entity.entity_id], |row| {
703                    Ok((
704                        row.get::<_, i64>(0)?,
705                        row.get::<_, i64>(1)?,
706                        row.get::<_, i64>(2)?,
707                    ))
708                })
709                .map_err(sql_error)?
710                .collect::<rusqlite::Result<Vec<(i64, i64, i64)>>>()
711                .map_err(sql_error)?;
712            conn.execute(
713                "DELETE FROM observation WHERE entity_id=?1",
714                [entity.entity_id],
715            )
716            .map_err(sql_error)?;
717            conn.execute(
718                "DELETE FROM relation WHERE from_id=?1 OR to_id=?1",
719                [entity.entity_id],
720            )
721            .map_err(sql_error)?;
722            conn.execute(
723                "INSERT INTO name_fts(name_fts,rowid,name) VALUES('delete',?1,?2)",
724                params![entity.entity_id, entity.name],
725            )
726            .map_err(sql_error)?;
727            conn.execute("DELETE FROM entity WHERE id=?1", [entity.entity_id])
728                .map_err(sql_error)?;
729            for (from_id, to_id, type_id) in triples {
730                tombstone_relation_mirror(conn, from_id, to_id, type_id)?;
731            }
732        }
733    }
734    Ok(())
735}
736
737fn execute(
738    graph: &GraphHandle,
739    conn: &Connection,
740    request: MutationRequest,
741) -> Result<MutationResult> {
742    match request {
743        MutationRequest::CreateEntities { entities } => {
744            let mut created = Vec::new();
745            for entity in entities {
746                if create_entity(graph, conn, &entity)? {
747                    created.push(require_entity(conn, &entity.name)?.entity());
748                }
749            }
750            Ok(MutationResult::Entities(created))
751        }
752        MutationRequest::UpsertEntities { entities } => {
753            let mut result = Vec::new();
754            for entity in entities {
755                if let Some(existing) = read_entity(conn, &entity.name)? {
756                    if existing.entity_type != entity.entity_type {
757                        conn.execute(
758                            "UPDATE entity SET type_id=?1 WHERE id=?2",
759                            params![type_id(conn, &entity.entity_type, 0)?, existing.entity_id],
760                        )
761                        .map_err(sql_error)?;
762                    }
763                    let mut seen: BTreeSet<&str> = existing
764                        .observations
765                        .iter()
766                        .map(|o| o.body.as_str())
767                        .collect();
768                    let added: Vec<ObservationInput> = entity
769                        .observations
770                        .iter()
771                        .filter(|o| seen.insert(o.body.as_str()))
772                        .cloned()
773                        .collect();
774                    insert_observations(graph, conn, existing.entity_id, &added)?;
775                    result.push(require_entity(conn, &entity.name)?.entity());
776                } else if create_entity(graph, conn, &entity)? {
777                    result.push(require_entity(conn, &entity.name)?.entity());
778                }
779            }
780            Ok(MutationResult::Entities(result))
781        }
782        MutationRequest::DeleteEntities { names } => {
783            delete_entities(conn, &names)?;
784            Ok(MutationResult::Unit)
785        }
786        MutationRequest::CreateRelations { relations } => {
787            let mut created = Vec::new();
788            for relation in relations {
789                if create_relation(conn, &relation)? {
790                    created.push(relation);
791                }
792            }
793            Ok(MutationResult::Relations(created))
794        }
795        MutationRequest::DeleteRelations { relations } => {
796            for relation in relations {
797                let triples = conn.prepare_cached(
798                    "SELECT rowid, from_id, to_id, type_id FROM relation \
799                     WHERE from_id IN (SELECT id FROM entity WHERE name_hash=?1 AND name=?2 AND flags=0) \
800                     AND to_id IN (SELECT id FROM entity WHERE name_hash=?3 AND name=?4 AND flags=0) \
801                     AND type_id IN (SELECT id FROM type_dict WHERE kind=1 AND name=?5)",
802                )
803                .map_err(sql_error)?
804                .query_map(
805                    params![name_hash(&relation.from), relation.from, name_hash(&relation.to), relation.to, relation.relation_type],
806                    |row| {
807                        Ok((
808                            row.get::<_, i64>(1)?,
809                            row.get::<_, i64>(2)?,
810                            row.get::<_, i64>(3)?,
811                        ))
812                    },
813                )
814                .map_err(sql_error)?
815                .collect::<rusqlite::Result<Vec<(i64, i64, i64)>>>()
816                .map_err(sql_error)?;
817                conn.execute("DELETE FROM relation WHERE from_id IN (SELECT id FROM entity WHERE name_hash=?1 AND name=?2 AND flags=0) AND to_id IN (SELECT id FROM entity WHERE name_hash=?3 AND name=?4 AND flags=0) AND type_id IN (SELECT id FROM type_dict WHERE kind=1 AND name=?5)", params![name_hash(&relation.from),relation.from,name_hash(&relation.to),relation.to,relation.relation_type]).map_err(sql_error)?;
818                for (from_id, to_id, type_id) in triples {
819                    tombstone_relation_mirror(conn, from_id, to_id, type_id)?;
820                }
821            }
822            Ok(MutationResult::Unit)
823        }
824        MutationRequest::AddObservations { observations } => {
825            let mut result = Vec::new();
826            for update in observations {
827                let entity = require_entity(conn, &update.entity_name)?;
828                let inserted =
829                    insert_observations(graph, conn, entity.entity_id, &update.contents)?;
830                result.push(ObservationResult {
831                    entity_name: update.entity_name,
832                    added_observations: inserted,
833                });
834            }
835            Ok(MutationResult::Observations(result))
836        }
837        MutationRequest::DeleteObservations { observations } => {
838            for update in observations {
839                if update.contents.is_empty() {
840                    continue;
841                }
842                let entity = require_entity(conn, &update.entity_name)?;
843                for body in &update.contents {
844                    conn.execute(
845                        "DELETE FROM observation WHERE entity_id=?1 AND body=?2",
846                        params![entity.entity_id, body.body],
847                    )
848                    .map_err(sql_error)?;
849                }
850            }
851            Ok(MutationResult::Unit)
852        }
853        MutationRequest::MergeEntities { source, target } => {
854            let old = require_entity(conn, &source)?;
855            let into = require_entity(conn, &target)?;
856            if source != target {
857                // Body remains the observation identity. Equal target bodies keep
858                // their metadata; newly copied rows retain the original fact/write
859                // times and record this merge's immediate source as audit origin.
860                let mut seen: BTreeSet<&str> =
861                    into.observations.iter().map(|o| o.body.as_str()).collect();
862                let mut idx: i64 = conn
863                    .query_row(
864                        "SELECT COALESCE(MAX(idx),-1) FROM observation WHERE entity_id=?1",
865                        [into.entity_id],
866                        |r| r.get(0),
867                    )
868                    .map_err(sql_error)?;
869                for observation in &old.observations {
870                    if seen.insert(&observation.body) {
871                        idx += 1;
872                        conn.execute("INSERT INTO observation(id,entity_id,idx,body,created_us,occurred_us,origin_entity_id,origin_entity_name) VALUES(?1,?2,?3,?4,?5,?6,?7,?8)", params![graph.next_obs_id(),into.entity_id,idx,observation.body,observation.created_at_us,observation.occurred_at_us,old.entity_id,old.name]).map_err(sql_error)?;
873                    }
874                }
875                let relations = relations_for(conn, &source)?;
876                for mut relation in relations {
877                    if relation.from == source {
878                        relation.from = target.clone();
879                    }
880                    if relation.to == source {
881                        relation.to = target.clone();
882                    }
883                    create_relation(conn, &relation)?;
884                }
885                delete_entities(conn, std::slice::from_ref(&source))?;
886            }
887            Ok(MutationResult::Entity(
888                require_entity(conn, &target)?.entity(),
889            ))
890        }
891        MutationRequest::RenameEntity { old_name, new_name } => {
892            let entity = require_entity(conn, &old_name)?;
893            if old_name == new_name {
894                return Ok(MutationResult::Entity(entity.entity()));
895            }
896            if read_entity(conn, &new_name)?.is_some() {
897                return Err(MCSError::InvalidParams(format!(
898                    "Entity '{new_name}' already exists"
899                )));
900            }
901            conn.execute(
902                "UPDATE entity SET name_hash=?1,name=?2 WHERE id=?3",
903                params![name_hash(&new_name), new_name, entity.entity_id],
904            )
905            .map_err(sql_error)?;
906            conn.execute(
907                "INSERT INTO name_fts(name_fts,rowid,name) VALUES('delete',?1,?2)",
908                params![entity.entity_id, old_name],
909            )
910            .map_err(sql_error)?;
911            conn.execute(
912                "INSERT INTO name_fts(rowid,name) VALUES(?1,?2)",
913                params![entity.entity_id, new_name],
914            )
915            .map_err(sql_error)?;
916            Ok(MutationResult::Entity(
917                require_entity(conn, &new_name)?.entity(),
918            ))
919        }
920        MutationRequest::PurgeDefinedEntities { name } => {
921            let names = defined_names(conn, &name)?;
922            delete_entities(conn, &names)?;
923            Ok(MutationResult::Count(names.len()))
924        }
925        MutationRequest::Compact => {
926            conn.execute_batch("PRAGMA incremental_vacuum;")
927                .map_err(sql_error)?;
928            Ok(MutationResult::Unit)
929        }
930        MutationRequest::Wipe => {
931            let mut stmt = conn
932                .prepare("SELECT name FROM entity WHERE flags=0")
933                .map_err(sql_error)?;
934            let names = stmt
935                .query_map([], |row| row.get(0))
936                .map_err(sql_error)?
937                .collect::<rusqlite::Result<Vec<String>>>()
938                .map_err(sql_error)?;
939            delete_entities(conn, &names)?;
940            // External-content indexes may contain orphan postings left by
941            // legacy deletions. Reset the indexes inside this transaction too.
942            conn.execute_batch(
943                "INSERT INTO name_fts(name_fts) VALUES('delete-all');
944                 INSERT INTO obs_fts(obs_fts) VALUES('delete-all');",
945            )
946            .map_err(sql_error)?;
947            Ok(MutationResult::Unit)
948        }
949    }
950}
951
952fn update_counters(
953    conn: &Connection,
954    before: &Snapshot,
955    after: &Snapshot,
956    changes: &[EntityChange],
957) -> Result<()> {
958    let mut type_deltas: BTreeMap<(i64, &str), i64> = BTreeMap::new();
959    for old in before.entities.values() {
960        *type_deltas.entry((0, &old.entity_type)).or_default() -= 1;
961    }
962    for new in after.entities.values() {
963        *type_deltas.entry((0, &new.entity_type)).or_default() += 1;
964    }
965    for (old, count) in &before.relation_rows {
966        *type_deltas.entry((1, &old.relation_type)).or_default() -= count;
967    }
968    for (new, count) in &after.relation_rows {
969        *type_deltas.entry((1, &new.relation_type)).or_default() += count;
970    }
971    // Affected types are the union of the net-delta keys and the types of
972    // every entity change. Each affected type gets ONE count/revision update
973    // and ONE enqueue, never one per change.
974    let mut affected_types: BTreeSet<(i64, &str)> = type_deltas.keys().copied().collect();
975    for change in changes {
976        if let Some(before) = &change.before {
977            affected_types.insert((0, before.entity_type.as_str()));
978        }
979        if let Some(after) = &change.after {
980            affected_types.insert((0, after.entity_type.as_str()));
981        }
982        if let Some(delta) = &change.relation_delta {
983            for relation in delta.added.iter().chain(delta.removed.iter()) {
984                affected_types.insert((1, relation.relation_type.as_str()));
985            }
986        }
987    }
988    for (kind, name) in affected_types {
989        let delta = type_deltas.get(&(kind, name)).copied().unwrap_or(0);
990        let revision: i64 = conn
991            .query_row(
992                "UPDATE type_dict SET count=count+?1, revision=revision+1 WHERE kind=?2 AND name=?3 RETURNING revision",
993                params![delta, kind, name],
994                |row| row.get(0),
995            )
996            .map_err(sql_error)?;
997        enqueue_taxonomy_jobs(
998            conn,
999            kind,
1000            type_id(conn, name, kind)?,
1001            revision,
1002            crate::jobs::IndexOperation::Upsert,
1003        )?;
1004    }
1005    let observations = |snapshot: &Snapshot| {
1006        snapshot
1007            .entities
1008            .values()
1009            .map(|e| e.observations.len() as i64)
1010            .sum::<i64>()
1011    };
1012    for (key, delta) in [
1013        (
1014            "entities",
1015            after.entities.len() as i64 - before.entities.len() as i64,
1016        ),
1017        (
1018            "relations",
1019            after.relation_rows.values().sum::<i64>() - before.relation_rows.values().sum::<i64>(),
1020        ),
1021        ("observations", observations(after) - observations(before)),
1022    ] {
1023        if delta != 0 {
1024            conn.execute(
1025                "UPDATE graph_stat SET value=value+?1 WHERE key=?2",
1026                params![delta, key],
1027            )
1028            .map_err(sql_error)?;
1029        }
1030    }
1031    let mut degrees: BTreeMap<&str, (i64, i64)> = BTreeMap::new();
1032    for (relation, count) in &after.relation_rows {
1033        degrees.entry(&relation.from).or_default().0 += count;
1034        degrees.entry(&relation.to).or_default().1 += count;
1035    }
1036    for entity in changes
1037        .iter()
1038        .filter(|change| change.operation != ChangeOperation::Rename)
1039        .filter_map(|change| change.after.as_ref())
1040    {
1041        let (outgoing, incoming) = degrees
1042            .get(entity.name.as_str())
1043            .copied()
1044            .unwrap_or_default();
1045        conn.execute(
1046            "UPDATE entity SET obs_count=?1,out_deg=?2,in_deg=?3,updated_us=?4 WHERE id=?5",
1047            params![
1048                entity.observations.len() as i64,
1049                outgoing,
1050                incoming,
1051                now_us(),
1052                entity.entity_id
1053            ],
1054        )
1055        .map_err(sql_error)?;
1056    }
1057    Ok(())
1058}
1059
1060#[cfg(test)]
1061mod tests {
1062    use super::*;
1063    use crate::graph::GraphHandle;
1064    use crate::storage::{Durability, SqliteTuning};
1065    use crate::types::EntityInput as Entity;
1066    use std::num::NonZeroUsize;
1067    use std::ops::Deref;
1068    use std::path::PathBuf;
1069
1070    struct TestKg(GraphHandle, PathBuf);
1071
1072    impl Deref for TestKg {
1073        type Target = GraphHandle;
1074        fn deref(&self) -> &GraphHandle {
1075            &self.0
1076        }
1077    }
1078
1079    impl Drop for TestKg {
1080        fn drop(&mut self) {
1081            let _ = std::fs::remove_file(&self.1);
1082            let _ = std::fs::remove_file(self.1.with_extension("db-wal"));
1083            let _ = std::fs::remove_file(self.1.with_extension("db-shm"));
1084        }
1085    }
1086
1087    fn new_kg() -> TestKg {
1088        use std::sync::atomic::AtomicU64;
1089        use std::sync::atomic::Ordering;
1090        static COUNTER: AtomicU64 = AtomicU64::new(200_000);
1091        let n = COUNTER.fetch_add(1, Ordering::SeqCst);
1092        let path =
1093            std::env::temp_dir().join(format!("kg_mutation_{}_{}.db", std::process::id(), n));
1094        let _ = std::fs::remove_file(&path);
1095        let _ = std::fs::remove_file(path.with_extension("db-wal"));
1096        let _ = std::fs::remove_file(path.with_extension("db-shm"));
1097        let kg = GraphHandle::new(
1098            &path,
1099            Durability::Async,
1100            SqliteTuning::default(),
1101            NonZeroUsize::new(10000).unwrap(),
1102            4,
1103        )
1104        .expect("create test kg");
1105        TestKg(kg, path)
1106    }
1107
1108    /// One managed serving profile, so taxonomy jobs land pending instead of
1109    /// held.
1110    fn serving_profile(kg: &GraphHandle) -> Uuid {
1111        let profile = Uuid::new_v4();
1112        let conn = kg.writer.lock();
1113        conn.execute(
1114            "UPDATE index_profile_registry SET state='Active', serving_profile=?1 WHERE store_key='default'",
1115            [profile.to_string()],
1116        )
1117        .expect("activate serving profile");
1118        profile
1119    }
1120
1121    fn entity(name: &str, entity_type: &str) -> Entity {
1122        Entity {
1123            name: name.into(),
1124            entity_type: entity_type.into(),
1125            observations: vec![],
1126        }
1127    }
1128
1129    fn relation(from: &str, to: &str, relation_type: &str) -> Relation {
1130        Relation {
1131            from: from.into(),
1132            to: to.into(),
1133            relation_type: relation_type.into(),
1134        }
1135    }
1136
1137    /// (owner_id, owner_revision, operation, state) of the chunk jobs for
1138    /// relation owners. The relation funnels enqueue these instead of the
1139    /// retired taxonomy kind-2 rows.
1140    fn relation_chunk_jobs(kg: &GraphHandle) -> Vec<(i64, i64, String, String)> {
1141        let conn = kg.writer.lock();
1142        let mut stmt = conn
1143            .prepare(
1144                "SELECT owner_id, owner_revision, operation, state FROM chunk_index_job
1145                 WHERE owner_kind='relation' ORDER BY owner_id",
1146            )
1147            .unwrap();
1148        stmt.query_map([], |row| {
1149            Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?))
1150        })
1151        .unwrap()
1152        .collect::<rusqlite::Result<_>>()
1153        .unwrap()
1154    }
1155
1156    /// (subject_kind, subject_id, subject_revision, operation, state)
1157    fn taxonomy_jobs(kg: &GraphHandle) -> Vec<(i64, i64, i64, String, String)> {
1158        let conn = kg.writer.lock();
1159        let mut stmt = conn
1160            .prepare("SELECT subject_kind, subject_id, subject_revision, operation, state FROM taxonomy_job ORDER BY subject_kind, subject_id")
1161            .unwrap();
1162        stmt.query_map([], |row| {
1163            Ok((
1164                row.get(0)?,
1165                row.get(1)?,
1166                row.get(2)?,
1167                row.get(3)?,
1168                row.get(4)?,
1169            ))
1170        })
1171        .unwrap()
1172        .collect::<rusqlite::Result<_>>()
1173        .unwrap()
1174    }
1175
1176    /// (count, revision) of one type_dict row.
1177    fn type_row(kg: &GraphHandle, kind: i64, name: &str) -> (i64, i64) {
1178        let conn = kg.writer.lock();
1179        conn.query_row(
1180            "SELECT count, revision FROM type_dict WHERE kind=?1 AND name=?2",
1181            params![kind, name],
1182            |row| Ok((row.get(0)?, row.get(1)?)),
1183        )
1184        .unwrap()
1185    }
1186
1187    /// (id, revision, deleted) of the mirror row for one relation triple.
1188    fn mirror_row(kg: &GraphHandle, from: &str, to: &str, relation_type: &str) -> (i64, i64, i64) {
1189        let conn = kg.writer.lock();
1190        conn.query_row(
1191            "SELECT m.id, m.revision, m.deleted
1192             FROM taxonomy_relation m
1193             JOIN entity f ON f.id = m.from_id
1194             JOIN entity t ON t.id = m.to_id
1195             JOIN type_dict d ON d.id = m.type_id
1196             WHERE f.name=?1 AND t.name=?2 AND d.name=?3 AND d.kind=1",
1197            params![from, to, relation_type],
1198            |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
1199        )
1200        .unwrap()
1201    }
1202
1203    #[test]
1204    fn create_entity_enqueues_one_type_job_at_revision_one() {
1205        let kg = new_kg();
1206        serving_profile(&kg);
1207        kg.create_entities(&[entity("ada", "person")]).unwrap();
1208
1209        let jobs = taxonomy_jobs(&kg);
1210        assert_eq!(jobs.len(), 1, "exactly one taxonomy job expected");
1211        assert_eq!(jobs[0].0, 0, "entity type job expected");
1212        assert_eq!(jobs[0].2, 1, "first revision expected");
1213        assert_eq!(jobs[0].3, "upsert");
1214        assert_eq!(jobs[0].4, "pending");
1215        let (count, revision) = type_row(&kg, 0, "person");
1216        assert_eq!((count, revision), (1, 1));
1217    }
1218
1219    #[test]
1220    fn second_entity_of_same_type_upserts_job_to_revision_two() {
1221        let kg = new_kg();
1222        serving_profile(&kg);
1223        kg.create_entities(&[entity("ada", "person")]).unwrap();
1224        kg.create_entities(&[entity("bob", "person")]).unwrap();
1225
1226        let jobs = taxonomy_jobs(&kg);
1227        assert_eq!(jobs.len(), 1, "one job row per affected type expected");
1228        assert_eq!((jobs[0].0, jobs[0].2, jobs[0].3.as_str()), (0, 2, "upsert"));
1229        let (count, revision) = type_row(&kg, 0, "person");
1230        assert_eq!((count, revision), (2, 2));
1231    }
1232
1233    #[test]
1234    fn rename_bumps_type_revision_without_count_change() {
1235        let kg = new_kg();
1236        serving_profile(&kg);
1237        kg.create_entities(&[entity("ada", "person")]).unwrap();
1238        kg.rename_entity("ada", "ada lovelace").unwrap();
1239
1240        let (count, revision) = type_row(&kg, 0, "person");
1241        assert_eq!((count, revision), (1, 2), "count stable, revision bumped");
1242        let jobs = taxonomy_jobs(&kg);
1243        assert_eq!(jobs.len(), 1);
1244        assert_eq!((jobs[0].0, jobs[0].2, jobs[0].3.as_str()), (0, 2, "upsert"));
1245    }
1246
1247    #[test]
1248    fn create_relation_writes_mirror() {
1249        let kg = new_kg();
1250        serving_profile(&kg);
1251        kg.create_entities(&[entity("ada", "person"), entity("bob", "person")])
1252            .unwrap();
1253        kg.create_relations(&[relation("ada", "bob", "knows")])
1254            .unwrap();
1255
1256        let (mirror_id, mirror_revision, deleted) = mirror_row(&kg, "ada", "bob", "knows");
1257        assert_eq!((mirror_revision, deleted), (1, 0), "fresh mirror expected");
1258
1259        let jobs = relation_chunk_jobs(&kg);
1260        assert_eq!(jobs.len(), 1, "the mirror enqueues one relation chunk job");
1261        assert_eq!(
1262            (jobs[0].0, jobs[0].1, jobs[0].2.as_str(), jobs[0].3.as_str()),
1263            (mirror_id, 1, "upsert", "pending")
1264        );
1265        let (count, revision) = type_row(&kg, 1, "knows");
1266        assert_eq!((count, revision), (1, 1));
1267    }
1268
1269    #[test]
1270    fn delete_relation_tombstones_mirror_and_enqueues_delete() {
1271        let kg = new_kg();
1272        serving_profile(&kg);
1273        kg.create_entities(&[entity("ada", "person"), entity("bob", "person")])
1274            .unwrap();
1275        kg.create_relations(&[relation("ada", "bob", "knows")])
1276            .unwrap();
1277        let (mirror_id, _, _) = mirror_row(&kg, "ada", "bob", "knows");
1278
1279        kg.delete_relations(&[relation("ada", "bob", "knows")])
1280            .unwrap();
1281
1282        let (mirror_id_after, revision, deleted) = mirror_row(&kg, "ada", "bob", "knows");
1283        assert_eq!(
1284            (mirror_id_after, revision, deleted),
1285            (mirror_id, 2, 1),
1286            "tombstone expected"
1287        );
1288        let conn = kg.writer.lock();
1289        let remaining: i64 = conn
1290            .query_row("SELECT COUNT(*) FROM relation", [], |r| r.get(0))
1291            .unwrap();
1292        drop(conn);
1293        assert_eq!(remaining, 0, "physical triple deleted");
1294        let jobs = relation_chunk_jobs(&kg);
1295        assert_eq!(jobs.len(), 1);
1296        assert_eq!(
1297            (jobs[0].0, jobs[0].1, jobs[0].2.as_str()),
1298            (mirror_id, 2, "delete")
1299        );
1300    }
1301
1302    #[test]
1303    fn recreated_relation_reuses_no_freed_mirror_id() {
1304        // The mirror id must not track the physical rowid: SQLite reuses a
1305        // freed rowid for a later row, and an explicit-id mirror insert would
1306        // collide with the tombstoned mirror of a different triple.
1307        let kg = new_kg();
1308        serving_profile(&kg);
1309        kg.create_entities(&[
1310            entity("ada", "person"),
1311            entity("bob", "person"),
1312            entity("carol", "person"),
1313        ])
1314        .unwrap();
1315        kg.create_relations(&[relation("ada", "bob", "knows")])
1316            .unwrap();
1317        let (first_id, _, _) = mirror_row(&kg, "ada", "bob", "knows");
1318        kg.delete_relations(&[relation("ada", "bob", "knows")])
1319            .unwrap();
1320        // The new triple reuses the freed relation rowid; its mirror must get
1321        // a fresh id instead of colliding with the tombstoned mirror above.
1322        kg.create_relations(&[relation("ada", "carol", "knows")])
1323            .unwrap();
1324        let (second_id, revision, deleted) = mirror_row(&kg, "ada", "carol", "knows");
1325        assert_ne!(first_id, second_id);
1326        assert_eq!((revision, deleted), (1, 0));
1327        let jobs = relation_chunk_jobs(&kg);
1328        assert_eq!(jobs.len(), 2);
1329        assert_eq!(
1330            (jobs[1].0, jobs[1].1, jobs[1].2.as_str()),
1331            (second_id, 1, "upsert")
1332        );
1333    }
1334
1335    #[test]
1336    fn one_mutation_with_many_changes_bumps_each_type_once() {
1337        let kg = new_kg();
1338        serving_profile(&kg);
1339        // Two creations of one new type in a single call. The union ruling
1340        // demands one revision bump and one job row for the affected type.
1341        kg.create_entities(&[entity("ada", "person"), entity("bob", "person")])
1342            .unwrap();
1343
1344        let (count, revision) = type_row(&kg, 0, "person");
1345        assert_eq!((count, revision), (2, 1), "one bump, not one per change");
1346        let jobs = taxonomy_jobs(&kg);
1347        assert_eq!(jobs.len(), 1);
1348        assert_eq!((jobs[0].0, jobs[0].2), (0, 1));
1349    }
1350
1351    #[test]
1352    fn delete_entity_tombstones_its_relation_mirrors() {
1353        let kg = new_kg();
1354        serving_profile(&kg);
1355        kg.create_entities(&[entity("ada", "person"), entity("bob", "person")])
1356            .unwrap();
1357        kg.create_relations(&[relation("ada", "bob", "knows")])
1358            .unwrap();
1359        let (mirror_id, _, _) = mirror_row(&kg, "ada", "bob", "knows");
1360
1361        kg.delete_entities(&["ada".into()]).unwrap();
1362
1363        let conn = kg.writer.lock();
1364        let (revision, deleted): (i64, i64) = conn
1365            .query_row(
1366                "SELECT revision, deleted FROM taxonomy_relation WHERE id=?1",
1367                [mirror_id],
1368                |row| Ok((row.get(0)?, row.get(1)?)),
1369            )
1370            .unwrap();
1371        drop(conn);
1372        assert_eq!((revision, deleted), (2, 1), "cascade tombstone expected");
1373        let jobs = relation_chunk_jobs(&kg);
1374        assert_eq!(jobs.len(), 1);
1375        assert_eq!((jobs[0].1, jobs[0].2.as_str()), (2, "delete"));
1376        let (count, revision) = type_row(&kg, 0, "person");
1377        // Bump once at entity creation, once at relation creation (the
1378        // relation delta makes an entity change for each endpoint), once at
1379        // the delete. Count drops to the one survivor.
1380        assert_eq!((count, revision), (1, 3), "survivor count and bumps");
1381    }
1382}