1use 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 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#[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 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 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
562fn 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
578fn 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 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 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 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 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 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 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 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 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 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 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 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 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 assert_eq!((count, revision), (1, 3), "survivor count and bumps");
1381 }
1382}