1use std::collections::{BTreeMap, BTreeSet};
47
48use rusqlite::{Connection, OptionalExtension};
49
50use crate::namespace_census::{self, NamespaceCensus, NamespaceConstraint};
51
52#[path = "namespace_move_knowledge.rs"]
53mod knowledge_move;
54#[path = "namespace_move_pack_tables.rs"]
55mod pack_tables;
56#[path = "namespace_move_routes.rs"]
57mod subject_routes;
58use knowledge_move::{move_knowledge_atoms, ordinary_atom_count};
59use pack_tables::{move_task_audit, settle_pack_tables};
60use subject_routes::routed_subjects;
61
62#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
70pub enum SubjectClass {
71 Note(String),
72 Entity(String),
73 Edge,
74 EdgeRelation(String),
75 Atom,
76 Domain,
77}
78
79impl SubjectClass {
80 pub fn parse(key: &str) -> Result<Self, MoveError> {
83 match key {
84 "edge" => return Ok(Self::Edge),
85 "atom" => return Ok(Self::Atom),
86 "domain" => return Ok(Self::Domain),
87 _ => {}
88 }
89 match key.split_once(':') {
90 Some(("note", kind)) if !kind.is_empty() => Ok(Self::Note(kind.to_string())),
91 Some(("entity", kind)) if !kind.is_empty() => Ok(Self::Entity(kind.to_string())),
92 Some(("edge", relation))
93 if khive_types::EdgeRelation::VALID_NAMES.contains(&relation) =>
94 {
95 Ok(Self::EdgeRelation(relation.to_string()))
96 }
97 _ => Err(MoveError::UnknownSubjectClass {
98 key: key.to_string(),
99 }),
100 }
101 }
102
103 pub fn render(&self) -> String {
104 match self {
105 Self::Note(kind) => format!("note:{kind}"),
106 Self::Entity(kind) => format!("entity:{kind}"),
107 Self::Edge => "edge".to_string(),
108 Self::EdgeRelation(relation) => format!("edge:{relation}"),
109 Self::Atom => "atom".to_string(),
110 Self::Domain => "domain".to_string(),
111 }
112 }
113}
114
115#[derive(Clone, Debug, PartialEq, Eq)]
117pub struct MoveRoute {
118 pub class: SubjectClass,
119 pub target: String,
120}
121
122#[derive(Clone, Debug, PartialEq, Eq)]
124pub struct MoveRequest {
125 pub source: String,
126 pub routes: Vec<MoveRoute>,
127 pub now_micros: Option<i64>,
132}
133
134impl MoveRequest {
135 pub fn new(source: impl Into<String>, routes: Vec<MoveRoute>) -> Self {
139 Self {
140 source: source.into(),
141 routes,
142 now_micros: None,
143 }
144 }
145
146 pub fn at(mut self, now_micros: i64) -> Self {
148 self.now_micros = Some(now_micros);
149 self
150 }
151
152 fn route_for(&self, class: &SubjectClass) -> Option<&MoveRoute> {
153 self.routes
154 .iter()
155 .find(|route| &route.class == class)
156 .or_else(|| match class {
157 SubjectClass::EdgeRelation(_) => self
158 .routes
159 .iter()
160 .find(|route| route.class == SubjectClass::Edge),
161 _ => None,
162 })
163 }
164
165 fn single_target(&self) -> Option<&str> {
170 let mut targets = self.routes.iter().map(|r| r.target.as_str());
171 let first = targets.next()?;
172 targets.all(|t| t == first).then_some(first)
173 }
174}
175
176#[derive(Clone, Debug, Default, PartialEq, Eq)]
183pub struct MoveCounts {
184 pub subjects: BTreeMap<String, u64>,
189 pub rows: BTreeMap<String, u64>,
191 pub left_behind: BTreeMap<String, u64>,
194 pub ann_log_appended: u64,
196 pub live_policies_left_behind: u64,
200 pub grants_in_force_left_behind: u64,
204}
205
206#[derive(Clone, Debug, PartialEq, Eq)]
208pub struct StreamMember {
209 pub note_id: String,
210 pub stream: String,
211 pub seq: i64,
212}
213
214#[derive(Clone, Debug, PartialEq, Eq)]
216pub struct Collision {
217 pub table: String,
218 pub constraint: String,
219 pub target: String,
220 pub key: String,
226}
227
228#[derive(Debug)]
230pub enum MoveError {
231 UnknownSubjectClass {
233 key: String,
234 },
235 DuplicateRoute {
238 class: String,
239 },
240 TargetIsSource {
243 class: String,
244 },
245 UnknownTable {
253 table: String,
254 rows: u64,
255 },
256 UnroutedClass {
262 class: String,
263 rows: u64,
264 },
265 StreamMembers {
273 notes: Vec<StreamMember>,
274 },
275 UnroutableVector {
278 subject_id: String,
279 table: String,
280 destinations: u64,
281 },
282 Collisions {
284 collisions: Vec<Collision>,
285 },
286 Sqlite(rusqlite::Error),
287}
288
289impl From<rusqlite::Error> for MoveError {
290 fn from(error: rusqlite::Error) -> Self {
291 Self::Sqlite(error)
292 }
293}
294
295impl std::fmt::Display for MoveError {
296 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
297 match self {
298 Self::UnknownSubjectClass { key } => write!(
299 f,
300 "no subject class named {key:?}; expected note:<kind>, entity:<kind>, edge:<relation>, edge, atom or domain"
301 ),
302 Self::DuplicateRoute { class } => {
303 write!(f, "{class} is routed more than once")
304 }
305 Self::TargetIsSource { class } => {
306 write!(f, "{class} is routed to the namespace it is already in")
307 }
308 Self::UnknownTable { table, rows } => write!(
309 f,
310 "{table} holds {rows} row(s) in the source namespace and this build has no rule \
311 for it; a table added by a later migration refuses a move rather than being \
312 left behind by one"
313 ),
314 Self::UnroutedClass { class, rows } => write!(
315 f,
316 "{class} has {rows} row(s) in the source namespace and no route; \
317 a class routed with zero rows succeeds reporting zero, an unrouted one refuses"
318 ),
319 Self::StreamMembers { notes } => {
320 write!(f, "{} note(s) belong to a stream and cannot change namespace: ", notes.len())?;
321 for (i, member) in notes.iter().enumerate() {
322 if i > 0 {
323 f.write_str(", ")?;
324 }
325 write!(f, "{} ({} seq {})", member.note_id, member.stream, member.seq)?;
326 }
327 Ok(())
328 }
329 Self::UnroutableVector { subject_id, table, destinations } => write!(
330 f,
331 "unroutable_vector: {table:?} subject {subject_id:?} has {destinations} source-subject destinations; a partitioned move requires exactly one"
332 ),
333 Self::Collisions { collisions } => {
334 write!(f, "{} collision(s): ", collisions.len())?;
335 for (i, collision) in collisions.iter().enumerate() {
336 if i > 0 {
337 f.write_str(", ")?;
338 }
339 write!(
340 f,
341 "{}.{} in {} already holds {}",
342 collision.table, collision.constraint, collision.target, collision.key
343 )?;
344 }
345 Ok(())
346 }
347 Self::Sqlite(error) => write!(f, "{error}"),
348 }
349 }
350}
351
352impl std::error::Error for MoveError {}
353
354const NAMESPACE_SCOPED_TABLES: &[&str] = &[
362 "brain_profile_snapshots",
363 "brain_event_log",
364 "proposals_open",
365];
366
367const SUBJECT_KEYED_TABLES: &[&str] = &["brain_implicit_mass", "brain_serve_ledger"];
375
376const MEMORY_VISIBILITY_RECEIPTS: &str = "memory_visibility_receipts";
379const MEMORY_VISIBILITY_FENCES: &str = "memory_visibility_fences";
380const MEMORY_VISIBILITY_EPOCHS: &str = "memory_visibility_epochs";
381
382const LEAVE_BEHIND_TABLES: &[&str] = &["comm_sender_transport", "sessions", "session_messages"];
387
388#[derive(Clone, Copy, Debug, PartialEq, Eq)]
398pub enum TableDisposition {
399 Subject,
401 Derived { trigger_maintained: bool },
418 DerivedRowidMap,
420 Appended,
424 History,
427 ConsumerWatermark,
429 RefusedBySchema,
432 NamespaceScopedAggregate,
436 SubjectKeyed { subject_column: &'static str },
438 LeaveBehind,
441}
442
443pub fn disposition(table: &namespace_census::NamespaceTable) -> Option<TableDisposition> {
450 use TableDisposition::*;
451 Some(match table.name.as_str() {
452 "notes" | "entities" | "graph_edges" | "knowledge_atoms" | "knowledge_domains" => Subject,
453
454 "knowledge_sections" => SubjectKeyed {
457 subject_column: "atom_id",
458 },
459
460 "proposals_open" => NamespaceScopedAggregate,
466
467 "fts_knowledge" | "fts_sections" => Derived {
468 trigger_maintained: true,
469 },
470 "fts_notes" | "fts_entities" => Derived {
471 trigger_maintained: false,
472 },
473 "fts_notes_rowids" | "fts_entities_rowids" => DerivedRowidMap,
474 "vector_provenance" => Derived {
476 trigger_maintained: false,
477 },
478
479 "ann_write_log" => Appended,
480 "events" => History,
481 "ann_consumer_watermark" | "ann_consumer_pending" => ConsumerWatermark,
482 "note_streams" => RefusedBySchema,
483
484 "brain_profile_snapshots" | "brain_event_log" => NamespaceScopedAggregate,
485 "brain_implicit_mass" | "brain_serve_ledger" => SubjectKeyed {
486 subject_column: "target_id",
487 },
488 "memory_visibility_receipts" => SubjectKeyed {
489 subject_column: "note_id",
490 },
491 "memory_visibility_fences" => SubjectKeyed {
492 subject_column: "note_id",
493 },
494 "memory_visibility_epochs" => SubjectKeyed {
495 subject_column: "note_id",
496 },
497
498 "gtd_lifecycle_audit" => SubjectKeyed {
505 subject_column: "note_id",
506 },
507 "knowledge_eval_runs" => NamespaceScopedAggregate,
510 "exec_runs" | "exec_events" | "git_receipts" => NamespaceScopedAggregate,
516 "tool_policy" | "tool_grants" => LeaveBehind,
520 "retrieval_snapshots" => Derived {
524 trigger_maintained: false,
525 },
526 _ if LEAVE_BEHIND_TABLES.contains(&table.name.as_str()) => LeaveBehind,
527
528 _ if is_runtime_vector_table(table) => Derived {
540 trigger_maintained: false,
541 },
542
543 _ => return None,
544 })
545}
546
547fn is_runtime_vector_table(table: &namespace_census::NamespaceTable) -> bool {
553 table.name.starts_with("vec_") && table.virtual_table
554}
555
556#[derive(Debug, Default)]
559struct SourceInventory {
560 note_kinds: BTreeMap<String, u64>,
562 entity_kinds: BTreeMap<String, u64>,
563 edge_relations: BTreeMap<String, u64>,
564 atoms: u64,
565 domains: u64,
566 unknown: Vec<(String, u64)>,
569}
570
571fn count_in_namespace(conn: &Connection, table: &str, namespace: &str) -> rusqlite::Result<u64> {
572 let sql = format!(
573 "SELECT COUNT(*) FROM {} WHERE namespace = ?1",
574 namespace_census::quote_ident(table)
575 );
576 conn.query_row(&sql, [namespace], |row| row.get::<_, i64>(0))
577 .map(|n| n as u64)
578}
579
580fn kinds_in_namespace(
581 conn: &Connection,
582 table: &str,
583 namespace: &str,
584) -> rusqlite::Result<BTreeMap<String, u64>> {
585 let sql = format!(
586 "SELECT kind, COUNT(*) FROM {} WHERE namespace = ?1 GROUP BY kind",
587 namespace_census::quote_ident(table)
588 );
589 let mut stmt = conn.prepare(&sql)?;
590 let rows = stmt.query_map([namespace], |row| {
591 Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)? as u64))
592 })?;
593 let mut out = BTreeMap::new();
594 for row in rows {
595 let (kind, count) = row?;
596 out.insert(kind, count);
597 }
598 Ok(out)
599}
600
601fn edge_relations_in_namespace(
602 conn: &Connection,
603 namespace: &str,
604) -> rusqlite::Result<BTreeMap<String, u64>> {
605 conn.prepare(
606 "SELECT relation, COUNT(*) FROM graph_edges WHERE namespace = ?1 GROUP BY relation",
607 )?
608 .query_map([namespace], |row| {
609 Ok((row.get(0)?, row.get::<_, i64>(1)? as u64))
610 })?
611 .collect()
612}
613
614fn read_source(
622 conn: &Connection,
623 census: &NamespaceCensus,
624 source: &str,
625) -> Result<SourceInventory, MoveError> {
626 let mut inventory = SourceInventory {
627 note_kinds: kinds_in_namespace(conn, "notes", source)?,
628 entity_kinds: kinds_in_namespace(conn, "entities", source)?,
629 edge_relations: edge_relations_in_namespace(conn, source)?,
630 atoms: ordinary_atom_count(conn, source)?,
631 domains: count_in_namespace(conn, "knowledge_domains", source)?,
632 unknown: Vec::new(),
633 };
634
635 for table in &census.tables {
636 if disposition(table).is_some() {
637 continue;
638 }
639 let count = count_in_namespace(conn, &table.name, source)?;
640 if count > 0 {
641 inventory.unknown.push((table.name.clone(), count));
642 }
643 }
644 Ok(inventory)
645}
646
647fn stream_members(conn: &Connection, source: &str) -> rusqlite::Result<Vec<StreamMember>> {
654 let mut stmt = conn.prepare(
655 "SELECT note_id, stream, seq FROM note_streams \
656 WHERE namespace = ?1 ORDER BY stream, seq",
657 )?;
658 let rows = stmt.query_map([source], |row| {
659 Ok(StreamMember {
660 note_id: row.get(0)?,
661 stream: row.get(1)?,
662 seq: row.get(2)?,
663 })
664 })?;
665 rows.collect()
666}
667
668pub fn validate(
673 conn: &Connection,
674 census: &NamespaceCensus,
675 request: &MoveRequest,
676) -> Result<(), MoveError> {
677 let mut seen = BTreeSet::new();
678 for route in &request.routes {
679 if let SubjectClass::EdgeRelation(relation) = &route.class {
680 if !khive_types::EdgeRelation::VALID_NAMES.contains(&relation.as_str()) {
681 return Err(MoveError::UnknownSubjectClass {
682 key: route.class.render(),
683 });
684 }
685 }
686 if !seen.insert(route.class.clone()) {
687 return Err(MoveError::DuplicateRoute {
688 class: route.class.render(),
689 });
690 }
691 if route.target == request.source {
692 return Err(MoveError::TargetIsSource {
693 class: route.class.render(),
694 });
695 }
696 }
697
698 let inventory = read_source(conn, census, &request.source)?;
699
700 if let Some((table, rows)) = inventory.unknown.first() {
701 return Err(MoveError::UnknownTable {
702 table: table.clone(),
703 rows: *rows,
704 });
705 }
706
707 let unrouted = |class: SubjectClass, rows: u64| -> Result<(), MoveError> {
708 if rows > 0 && request.route_for(&class).is_none() {
709 return Err(MoveError::UnroutedClass {
710 class: class.render(),
711 rows,
712 });
713 }
714 Ok(())
715 };
716 for (kind, rows) in &inventory.note_kinds {
717 unrouted(SubjectClass::Note(kind.clone()), *rows)?;
718 }
719 for (kind, rows) in &inventory.entity_kinds {
720 unrouted(SubjectClass::Entity(kind.clone()), *rows)?;
721 }
722 for (relation, rows) in &inventory.edge_relations {
723 unrouted(SubjectClass::EdgeRelation(relation.clone()), *rows)?;
724 }
725 unrouted(SubjectClass::Atom, inventory.atoms)?;
726 unrouted(SubjectClass::Domain, inventory.domains)?;
727
728 let pinned = stream_members(conn, &request.source)?;
729 if !pinned.is_empty() {
730 return Err(MoveError::StreamMembers { notes: pinned });
731 }
732
733 validate_partitioned_vector_moves(conn, census, request)?;
734 Ok(())
735}
736
737fn validate_partitioned_vector_moves(
738 conn: &Connection,
739 census: &NamespaceCensus,
740 request: &MoveRequest,
741) -> Result<(), MoveError> {
742 if request.single_target().is_some() {
743 return Ok(());
744 }
745 let (selector, parameters) = routed_subjects(request);
746 for table in census
747 .tables
748 .iter()
749 .filter(|table| is_runtime_vector_table(table))
750 {
751 let unresolved = conn
752 .query_row(
753 &format!(
754 "WITH routed AS ({selector}) \
755 SELECT vector.subject_id, COUNT(DISTINCT routed.target) \
756 FROM {} AS vector LEFT JOIN routed ON routed.subject_id = vector.subject_id \
757 WHERE vector.namespace = ?1 GROUP BY vector.subject_id \
758 HAVING COUNT(DISTINCT routed.target) != 1 \
759 ORDER BY vector.subject_id LIMIT 1",
760 namespace_census::quote_ident(&table.name)
761 ),
762 rusqlite::params_from_iter(¶meters),
763 |row| Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)? as u64)),
764 )
765 .optional()?;
766 if let Some((subject_id, destinations)) = unresolved {
767 return Err(MoveError::UnroutableVector {
768 subject_id,
769 table: table.name.clone(),
770 destinations,
771 });
772 }
773 }
774 Ok(())
775}
776
777fn edge_collision_target_predicate(
778 request: &MoveRequest,
779 target: &str,
780 parameters: &mut Vec<rusqlite::types::Value>,
781) -> String {
782 let fallback = request
783 .route_for(&SubjectClass::Edge)
784 .is_some_and(|route| route.target == target);
785 let mut matching = Vec::new();
786 let mut specific = Vec::new();
787 for route in &request.routes {
788 if let SubjectClass::EdgeRelation(relation) = &route.class {
789 if route.target != target && !fallback {
790 continue;
791 }
792 parameters.push(relation.clone().into());
793 let parameter = format!("?{}", parameters.len());
794 if route.target == target {
795 matching.push(parameter.clone());
796 }
797 if fallback {
798 specific.push(parameter);
799 }
800 }
801 }
802 let mut alternatives = Vec::new();
803 if !matching.is_empty() {
804 alternatives.push(format!("source.relation IN ({})", matching.join(", ")));
805 }
806 if fallback {
807 alternatives.push(if specific.is_empty() {
808 "1".to_owned()
809 } else {
810 format!("source.relation NOT IN ({})", specific.join(", "))
811 });
812 }
813 if alternatives.is_empty() {
814 "0".to_owned()
815 } else {
816 alternatives.join(" OR ")
817 }
818}
819
820fn collisions_for(
873 conn: &Connection,
874 constraint: &NamespaceConstraint,
875 request: &MoveRequest,
876 target: &str,
877) -> rusqlite::Result<Vec<Collision>> {
878 if constraint.partial || !constraint.columns_are_nameable() {
879 return Ok(Vec::new());
880 }
881 let names: Vec<&str> = constraint
882 .columns
883 .iter()
884 .filter_map(|c| c.as_deref())
885 .collect();
886 let others: Vec<&str> = names
887 .iter()
888 .copied()
889 .filter(|c| !c.eq_ignore_ascii_case("namespace"))
890 .collect();
891 if others.len() != names.len() - 1 {
892 return Ok(Vec::new());
894 }
895 if others.is_empty() {
896 return collisions_on_namespace_alone(conn, constraint, &request.source, target);
900 }
901
902 let table = namespace_census::quote_ident(&constraint.table);
903 let join = others
904 .iter()
905 .map(|c| {
906 let q = namespace_census::quote_ident(c);
907 format!("target.{q} IS source.{q}")
911 })
912 .collect::<Vec<_>>()
913 .join(" AND ");
914 let select = others
915 .iter()
916 .map(|c| format!("source.{}", namespace_census::quote_ident(c)))
917 .collect::<Vec<_>>()
918 .join(", ");
919 let mut parameters = vec![request.source.clone().into(), target.to_owned().into()];
920 let mut sql = format!(
921 "SELECT {select} FROM {table} AS source \
922 JOIN {table} AS target ON target.namespace = ?2 AND {join} \
923 WHERE source.namespace = ?1"
924 );
925 if constraint.table == "graph_edges" && constraint.index == "idx_graph_edges_unique_triple" {
926 let predicate = edge_collision_target_predicate(request, target, &mut parameters);
927 sql.push_str(&format!(" AND ({predicate})"));
928 }
929
930 let mut stmt = conn.prepare(&sql)?;
931 let column_count = others.len();
932 let rows = stmt.query_map(rusqlite::params_from_iter(parameters), move |row| {
933 let mut parts = Vec::with_capacity(column_count);
934 for i in 0..column_count {
935 parts.push(match row.get_ref(i)? {
936 rusqlite::types::ValueRef::Null => "NULL".to_string(),
937 rusqlite::types::ValueRef::Integer(v) => v.to_string(),
938 rusqlite::types::ValueRef::Real(v) => v.to_string(),
939 rusqlite::types::ValueRef::Text(v) => String::from_utf8_lossy(v).into_owned(),
940 rusqlite::types::ValueRef::Blob(_) => "<blob>".to_string(),
941 });
942 }
943 Ok(parts.join(", "))
944 })?;
945
946 let mut found = Vec::new();
947 for key in rows {
948 found.push(Collision {
949 table: constraint.table.clone(),
950 constraint: constraint.index.clone(),
951 target: target.to_string(),
952 key: key?,
953 });
954 }
955 Ok(found)
956}
957
958fn collisions_on_namespace_alone(
959 conn: &Connection,
960 constraint: &NamespaceConstraint,
961 source: &str,
962 target: &str,
963) -> rusqlite::Result<Vec<Collision>> {
964 let table = namespace_census::quote_ident(&constraint.table);
965 let sql = format!(
966 "SELECT (SELECT COUNT(*) FROM {table} WHERE namespace = ?1) \
967 * (SELECT COUNT(*) FROM {table} WHERE namespace = ?2)"
968 );
969 let product: i64 = conn.query_row(&sql, [source, target], |row| row.get(0))?;
970 Ok(if product > 0 {
971 vec![Collision {
972 table: constraint.table.clone(),
973 constraint: constraint.index.clone(),
974 target: target.to_string(),
975 key: "(namespace alone)".to_string(),
976 }]
977 } else {
978 Vec::new()
979 })
980}
981
982struct KindedTables {
994 base: &'static str,
995 fts: &'static str,
996 rowids: &'static str,
997}
998
999const NOTE_TABLES: KindedTables = KindedTables {
1000 base: "notes",
1001 fts: "fts_notes",
1002 rowids: "fts_notes_rowids",
1003};
1004
1005const ENTITY_TABLES: KindedTables = KindedTables {
1006 base: "entities",
1007 fts: "fts_entities",
1008 rowids: "fts_entities_rowids",
1009};
1010
1011fn move_kinded_subject(
1012 conn: &Connection,
1013 tables: &KindedTables,
1014 source: &str,
1015 target: &str,
1016 kind: &str,
1017 rows: &mut BTreeMap<String, u64>,
1018) -> rusqlite::Result<u64> {
1019 let KindedTables { base, fts, rowids } = *tables;
1020 let statement = if base == "entities" {
1024 "UPDATE entities SET namespace = ?2, version = version + 1 \
1025 WHERE namespace = ?1 AND kind = ?3"
1026 .to_owned()
1027 } else {
1028 format!(
1029 "UPDATE {} SET namespace = ?2 WHERE namespace = ?1 AND kind = ?3",
1030 namespace_census::quote_ident(base)
1031 )
1032 };
1033 let moved = conn.execute(&statement, rusqlite::params![source, target, kind])? as u64;
1034 *rows.entry(base.to_string()).or_default() += moved;
1035
1036 let selector = format!(
1039 "SELECT id FROM {} WHERE namespace = ?2 AND kind = ?3",
1040 namespace_census::quote_ident(base)
1041 );
1042 for derived in [fts, rowids] {
1043 let n = conn.execute(
1044 &format!(
1045 "UPDATE {} SET namespace = ?2 \
1046 WHERE namespace = ?1 AND subject_id IN ({selector})",
1047 namespace_census::quote_ident(derived)
1048 ),
1049 rusqlite::params![source, target, kind],
1050 )? as u64;
1051 *rows.entry(derived.to_string()).or_default() += n;
1052 }
1053 Ok(moved)
1054}
1055
1056fn move_whole_table(
1058 conn: &Connection,
1059 table: &str,
1060 source: &str,
1061 target: &str,
1062 rows: &mut BTreeMap<String, u64>,
1063) -> rusqlite::Result<u64> {
1064 let moved = conn.execute(
1065 &format!(
1066 "UPDATE {} SET namespace = ?2 WHERE namespace = ?1",
1067 namespace_census::quote_ident(table)
1068 ),
1069 rusqlite::params![source, target],
1070 )? as u64;
1071 *rows.entry(table.to_string()).or_default() += moved;
1072 Ok(moved)
1073}
1074
1075fn move_edges(
1076 conn: &Connection,
1077 request: &MoveRequest,
1078 route: &MoveRoute,
1079 rows: &mut BTreeMap<String, u64>,
1080) -> rusqlite::Result<u64> {
1081 let mut parameters: Vec<rusqlite::types::Value> =
1082 vec![request.source.clone().into(), route.target.clone().into()];
1083 let predicate = if let SubjectClass::EdgeRelation(relation) = &route.class {
1084 parameters.push(relation.clone().into());
1085 "relation = ?3".to_owned()
1086 } else {
1087 let mut excluded = Vec::new();
1088 for specific in &request.routes {
1089 if let SubjectClass::EdgeRelation(relation) = &specific.class {
1090 parameters.push(relation.clone().into());
1091 excluded.push(format!("?{}", parameters.len()));
1092 }
1093 }
1094 if excluded.is_empty() {
1095 return move_whole_table(conn, "graph_edges", &request.source, &route.target, rows);
1096 }
1097 format!("relation NOT IN ({})", excluded.join(", "))
1098 };
1099 let moved = conn.execute(
1100 &format!("UPDATE graph_edges SET namespace = ?2 WHERE namespace = ?1 AND {predicate}"),
1101 rusqlite::params_from_iter(parameters),
1102 )? as u64;
1103 *rows.entry("graph_edges".into()).or_default() += moved;
1104 Ok(moved)
1105}
1106
1107fn move_memory_visibility(
1112 conn: &Connection,
1113 source: &str,
1114 target: &str,
1115 rows: &mut BTreeMap<String, u64>,
1116) -> rusqlite::Result<()> {
1117 let receipts = conn.execute(
1118 "INSERT INTO memory_visibility_receipts (namespace, note_id, model_count) \
1119 SELECT ?2, receipt.note_id, receipt.model_count \
1120 FROM memory_visibility_receipts AS receipt \
1121 JOIN notes AS note ON note.id = receipt.note_id \
1122 WHERE receipt.namespace = ?1 AND note.namespace = ?2",
1123 rusqlite::params![source, target],
1124 )? as u64;
1125 *rows.entry(MEMORY_VISIBILITY_RECEIPTS.into()).or_default() += receipts;
1126
1127 let fences = conn.execute(
1128 "UPDATE memory_visibility_fences SET namespace = ?2 \
1129 WHERE namespace = ?1 AND note_id IN (\
1130 SELECT note_id FROM memory_visibility_receipts WHERE namespace = ?2)",
1131 rusqlite::params![source, target],
1132 )? as u64;
1133 *rows.entry(MEMORY_VISIBILITY_FENCES.into()).or_default() += fences;
1134
1135 let removed = conn.execute(
1136 "DELETE FROM memory_visibility_receipts \
1137 WHERE namespace = ?1 AND note_id IN (\
1138 SELECT note_id FROM memory_visibility_receipts WHERE namespace = ?2)",
1139 rusqlite::params![source, target],
1140 )? as u64;
1141 debug_assert_eq!(removed, receipts);
1142 Ok(())
1143}
1144
1145fn move_vectors(
1158 conn: &Connection,
1159 table: &str,
1160 source: &str,
1161 target: Option<&str>,
1162) -> rusqlite::Result<VectorMove> {
1163 let quoted = namespace_census::quote_ident(table);
1164 let columns = "subject_id, namespace, kind, field, embedding_model, embedding";
1165
1166 conn.execute_batch("DROP TABLE IF EXISTS temp.namespace_move_vectors")?;
1167 let staged = if let Some(target) = target {
1168 conn.execute(
1169 &format!(
1170 "CREATE TEMP TABLE namespace_move_vectors AS \
1171 SELECT {columns}, ?2 AS target FROM {quoted} WHERE namespace = ?1"
1172 ),
1173 rusqlite::params![source, target],
1174 )
1175 } else {
1176 conn.execute(
1177 &format!(
1178 "CREATE TEMP TABLE namespace_move_vectors AS \
1179 SELECT vector.subject_id, vector.namespace, vector.kind, vector.field, \
1180 vector.embedding_model, vector.embedding, routed.target \
1181 FROM {quoted} AS vector JOIN (\
1182 SELECT DISTINCT subject_id, target FROM temp.namespace_move_subject_targets\
1183 ) AS routed ON routed.subject_id = vector.subject_id \
1184 WHERE vector.namespace = ?1"
1185 ),
1186 [source],
1187 )
1188 };
1189 staged?;
1192 let staged_rows: i64 = conn.query_row(
1193 "SELECT COUNT(*) FROM temp.namespace_move_vectors",
1194 [],
1195 |r| r.get(0),
1196 )?;
1197
1198 conn.execute(
1199 &format!("DELETE FROM {quoted} WHERE namespace = ?1"),
1200 [source],
1201 )?;
1202 let inserted = conn.execute(
1203 &format!(
1204 "INSERT INTO {quoted} ({columns}) \
1205 SELECT subject_id, target, kind, field, embedding_model, embedding \
1206 FROM temp.namespace_move_vectors"
1207 ),
1208 [],
1209 )? as u64;
1210 debug_assert_eq!(
1211 inserted, staged_rows as u64,
1212 "every staged vector is re-inserted or the move is losing embeddings"
1213 );
1214
1215 let has_provenance: bool = conn.query_row(
1221 "SELECT EXISTS(SELECT 1 FROM sqlite_master \
1222 WHERE type = 'table' AND name = 'vector_provenance')",
1223 [],
1224 |row| row.get(0),
1225 )?;
1226 if has_provenance {
1227 let model_key = table
1228 .strip_prefix("vec_")
1229 .expect("runtime vector tables use the vec_ prefix");
1230 conn.execute(
1231 "DELETE FROM vector_provenance \
1232 WHERE model_key = ?1 \
1233 AND subject_id IN (SELECT subject_id FROM temp.namespace_move_vectors)",
1234 [model_key],
1235 )?;
1236 }
1237
1238 let deleted = conn.execute(
1247 "INSERT INTO ann_write_log (namespace, embedding_model, kind, field, subject_id, op) \
1248 SELECT ?1, embedding_model, kind, field, subject_id, 'delete' \
1249 FROM temp.namespace_move_vectors",
1250 [source],
1251 )? as u64;
1252 let upserts = conn
1253 .prepare(
1254 "INSERT INTO ann_write_log \
1255 (namespace, embedding_model, kind, field, subject_id, op) \
1256 SELECT target, embedding_model, kind, field, subject_id, 'upsert' \
1257 FROM temp.namespace_move_vectors \
1258 RETURNING seq, subject_id, embedding_model, kind, field",
1259 )?
1260 .query_map([], |row| {
1261 Ok((
1262 row.get::<_, i64>(0)?,
1263 row.get::<_, String>(1)?,
1264 row.get::<_, String>(2)?,
1265 row.get::<_, String>(3)?,
1266 row.get::<_, String>(4)?,
1267 ))
1268 })?
1269 .collect::<rusqlite::Result<Vec<_>>>()?;
1270 for (seq, subject, model, kind, field) in &upserts {
1271 if kind == "note" && field == "note.content" {
1272 refresh_moved_memory_fence(conn, source, subject, model, *seq)?;
1273 }
1274 }
1275 let appended = deleted + upserts.len() as u64;
1276
1277 conn.execute_batch("DROP TABLE temp.namespace_move_vectors")?;
1278
1279 Ok(VectorMove {
1280 moved: inserted,
1281 ann_appended: appended,
1282 })
1283}
1284
1285fn refresh_moved_memory_fence(
1288 conn: &Connection,
1289 source: &str,
1290 subject: &str,
1291 model: &str,
1292 seq: i64,
1293) -> rusqlite::Result<()> {
1294 conn.execute(
1295 "UPDATE memory_visibility_fences SET ann_write_log_seq = ?4 \
1296 WHERE namespace = ?1 AND note_id = ?2 AND model = ?3",
1297 rusqlite::params![source, subject, model, seq],
1298 )?;
1299 Ok(())
1300}
1301
1302struct VectorMove {
1317 moved: u64,
1318 ann_appended: u64,
1319}
1320
1321pub fn move_namespace(conn: &Connection, request: &MoveRequest) -> Result<MoveCounts, MoveError> {
1329 let census = namespace_census::census(conn)?;
1330 validate(conn, &census, request)?;
1331
1332 let mut targets: BTreeSet<&str> = BTreeSet::new();
1335 for route in &request.routes {
1336 targets.insert(route.target.as_str());
1337 }
1338 let mut collisions = Vec::new();
1339 for constraint in namespace_census::reachable_constraints(&census) {
1340 for target in &targets {
1341 collisions.extend(collisions_for(conn, constraint, request, target)?);
1342 }
1343 }
1344 if !collisions.is_empty() {
1345 return Err(MoveError::Collisions { collisions });
1346 }
1347
1348 let mut counts = MoveCounts::default();
1349 let source = request.source.as_str();
1350
1351 if request.single_target().is_none() {
1354 let (selector, parameters) = routed_subjects(request);
1355 conn.execute_batch("DROP TABLE IF EXISTS temp.namespace_move_subject_targets")?;
1356 conn.execute(
1357 &format!("CREATE TEMP TABLE namespace_move_subject_targets AS {selector}"),
1358 rusqlite::params_from_iter(parameters),
1359 )?;
1360 }
1361
1362 move_task_audit(conn, &census, request, &mut counts.rows)?;
1365
1366 for route in &request.routes {
1367 let target = route.target.as_str();
1368 let moved = match &route.class {
1369 SubjectClass::Note(kind) => {
1370 move_kinded_subject(conn, &NOTE_TABLES, source, target, kind, &mut counts.rows)?
1371 }
1372 SubjectClass::Entity(kind) => {
1373 move_kinded_subject(conn, &ENTITY_TABLES, source, target, kind, &mut counts.rows)?
1374 }
1375 SubjectClass::Edge | SubjectClass::EdgeRelation(_) => {
1376 move_edges(conn, request, route, &mut counts.rows)?
1377 }
1378 SubjectClass::Atom => {
1379 move_knowledge_atoms(conn, request, target, false, &mut counts.rows)?
1380 }
1381 SubjectClass::Domain => {
1382 move_knowledge_atoms(conn, request, target, true, &mut counts.rows)?;
1384 move_whole_table(conn, "knowledge_domains", source, target, &mut counts.rows)?
1385 }
1386 };
1387 counts.subjects.insert(route.class.render(), moved);
1388 }
1389
1390 if let Some(target) = request.single_target() {
1393 move_whole_table(conn, "knowledge_sections", source, target, &mut counts.rows)?;
1394 }
1395
1396 for table in &census.tables {
1400 if !is_runtime_vector_table(table) {
1401 continue;
1402 }
1403 let moved = move_vectors(conn, &table.name, source, request.single_target())?;
1404 *counts.rows.entry(table.name.clone()).or_default() += moved.moved;
1405 counts.ann_log_appended += moved.ann_appended;
1406 }
1407 if request.single_target().is_none() {
1408 conn.execute_batch("DROP TABLE temp.namespace_move_subject_targets")?;
1409 }
1410 if census
1411 .tables
1412 .iter()
1413 .any(|table| table.name == "vector_provenance")
1414 {
1415 let left = count_in_namespace(conn, "vector_provenance", source)?;
1416 if left > 0 {
1417 counts.left_behind.insert("vector_provenance".into(), left);
1418 }
1419 }
1420
1421 for table in SUBJECT_KEYED_TABLES {
1425 for target in &targets {
1426 let moved = conn.execute(
1427 &format!(
1428 "UPDATE {} SET namespace = ?2 WHERE namespace = ?1 AND target_id IN (\
1429 SELECT id FROM notes WHERE namespace = ?2 \
1430 UNION ALL SELECT id FROM entities WHERE namespace = ?2 \
1431 UNION ALL SELECT id FROM knowledge_atoms WHERE namespace = ?2)",
1432 namespace_census::quote_ident(table)
1433 ),
1434 rusqlite::params![source, target],
1435 )? as u64;
1436 *counts.rows.entry((*table).to_string()).or_default() += moved;
1437 }
1438 let left = count_in_namespace(conn, table, source)?;
1439 if left > 0 {
1440 counts.left_behind.insert((*table).to_string(), left);
1441 }
1442 }
1443
1444 for target in &targets {
1446 let moved = conn.execute(
1447 "UPDATE memory_visibility_epochs SET namespace = ?2 \
1448 WHERE namespace = ?1 AND note_id IN (SELECT id FROM notes WHERE namespace = ?2)",
1449 rusqlite::params![source, target],
1450 )? as u64;
1451 *counts
1452 .rows
1453 .entry(MEMORY_VISIBILITY_EPOCHS.into())
1454 .or_default() += moved;
1455 }
1456 let left = count_in_namespace(conn, MEMORY_VISIBILITY_EPOCHS, source)?;
1457 if left > 0 {
1458 counts
1459 .left_behind
1460 .insert(MEMORY_VISIBILITY_EPOCHS.into(), left);
1461 }
1462
1463 for target in &targets {
1467 move_memory_visibility(conn, source, target, &mut counts.rows)?;
1468 }
1469 for table in [MEMORY_VISIBILITY_RECEIPTS, MEMORY_VISIBILITY_FENCES] {
1470 let left = count_in_namespace(conn, table, source)?;
1471 if left > 0 {
1472 counts.left_behind.insert((*table).to_string(), left);
1473 }
1474 }
1475
1476 for table in LEAVE_BEHIND_TABLES {
1477 let left = count_in_namespace(conn, table, source)?;
1478 if left > 0 {
1479 counts.left_behind.insert((*table).to_string(), left);
1480 }
1481 }
1482
1483 for table in NAMESPACE_SCOPED_TABLES {
1487 match request.single_target() {
1488 Some(target) => {
1489 move_whole_table(conn, table, source, target, &mut counts.rows)?;
1490 }
1491 None => {
1492 let left = count_in_namespace(conn, table, source)?;
1493 if left > 0 {
1494 counts.left_behind.insert((*table).to_string(), left);
1495 }
1496 }
1497 }
1498 }
1499
1500 settle_pack_tables(conn, &census, request, &mut counts)?;
1501
1502 Ok(counts)
1503}
1504
1505#[cfg(test)]
1506mod tests {
1507 use super::*;
1508 use crate::migrations::run_migrations_for_test as run_migrations;
1509 use rusqlite::Connection;
1510
1511 pub(super) fn migrated() -> Connection {
1512 let mut conn = Connection::open_in_memory().expect("open");
1513 run_migrations(&mut conn).expect("migrate");
1514 conn
1515 }
1516
1517 pub(super) fn seed_note(conn: &Connection, id: &str, namespace: &str, kind: &str) {
1527 conn.execute(
1528 "INSERT INTO notes (id, namespace, kind, name, content, created_at, updated_at) \
1529 VALUES (?1, ?2, ?3, 'a name', 'some content', 1, 1)",
1530 rusqlite::params![id, namespace, kind],
1531 )
1532 .expect("seed note");
1533 }
1534
1535 pub(super) fn route(key: &str, target: &str) -> MoveRoute {
1536 MoveRoute {
1537 class: SubjectClass::parse(key).expect("route key"),
1538 target: target.to_string(),
1539 }
1540 }
1541
1542 #[test]
1547 fn a_routed_class_with_no_rows_succeeds_reporting_zero() {
1548 let conn = migrated();
1549 let request = MoveRequest::new(
1550 "empty-source",
1551 vec![route("note:observation", "target"), route("atom", "target")],
1552 );
1553 let counts = move_namespace(&conn, &request).expect("a backend with nothing routed here");
1554 assert_eq!(counts.subjects.get("note:observation"), Some(&0));
1555 assert_eq!(counts.subjects.get("atom"), Some(&0));
1556 assert_eq!(counts.rows.get(MEMORY_VISIBILITY_RECEIPTS), Some(&0));
1557 assert_eq!(counts.rows.get(MEMORY_VISIBILITY_FENCES), Some(&0));
1558 }
1559
1560 #[test]
1561 fn a_memory_visibility_receipt_follows_its_note_and_reports_count() {
1562 let conn = migrated();
1563 conn.pragma_update(None, "foreign_keys", "ON").unwrap();
1564 seed_note(&conn, "n1", "source", "observation");
1565 conn.execute(
1566 "INSERT INTO memory_visibility_receipts (namespace, note_id, model_count) \
1567 VALUES ('source', 'n1', 0)",
1568 [],
1569 )
1570 .unwrap();
1571
1572 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1573 let counts = move_namespace(&conn, &request).expect("the receipt follows its note");
1574 assert_eq!(counts.subjects.get("note:observation"), Some(&1));
1575 assert_eq!(counts.rows.get(MEMORY_VISIBILITY_RECEIPTS), Some(&1));
1576 assert_eq!(counts.rows.get(MEMORY_VISIBILITY_FENCES), Some(&0));
1577 let stored: (String, i64) = conn
1578 .query_row(
1579 "SELECT namespace, model_count FROM memory_visibility_receipts \
1580 WHERE note_id = 'n1'",
1581 [],
1582 |row| Ok((row.get(0)?, row.get(1)?)),
1583 )
1584 .unwrap();
1585 assert_eq!(stored, ("target".into(), 0));
1586 let source_rows: i64 = conn
1587 .query_row(
1588 "SELECT COUNT(*) FROM memory_visibility_receipts WHERE namespace = 'source'",
1589 [],
1590 |row| row.get(0),
1591 )
1592 .unwrap();
1593 assert_eq!(source_rows, 0);
1594 }
1595
1596 #[test]
1597 fn memory_visibility_fences_follow_their_note_with_foreign_keys_enabled() {
1598 let conn = migrated();
1599 conn.pragma_update(None, "foreign_keys", "ON").unwrap();
1600 seed_note(&conn, "n1", "source", "observation");
1601 seed_note(&conn, "n2", "source", "decision");
1602 conn.execute(
1603 "INSERT INTO memory_visibility_receipts (namespace, note_id, model_count) \
1604 VALUES ('source', 'n1', 2), ('source', 'n2', 0)",
1605 [],
1606 )
1607 .unwrap();
1608 conn.execute(
1609 "INSERT INTO memory_visibility_fences \
1610 (namespace, note_id, model, ann_write_log_seq) VALUES \
1611 ('source', 'n1', 'model-a', 10), ('source', 'n1', 'model-b', 11)",
1612 [],
1613 )
1614 .unwrap();
1615
1616 let request = MoveRequest::new(
1617 "source",
1618 vec![
1619 route("note:observation", "target-a"),
1620 route("note:decision", "target-b"),
1621 ],
1622 );
1623 let counts = move_namespace(&conn, &request).expect("both receipts follow their notes");
1624 assert_eq!(counts.rows.get(MEMORY_VISIBILITY_RECEIPTS), Some(&2));
1625 assert_eq!(counts.rows.get(MEMORY_VISIBILITY_FENCES), Some(&2));
1626 let receipt_places: Vec<(String, String)> = conn
1627 .prepare("SELECT note_id, namespace FROM memory_visibility_receipts ORDER BY note_id")
1628 .unwrap()
1629 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1630 .unwrap()
1631 .collect::<rusqlite::Result<_>>()
1632 .unwrap();
1633 assert_eq!(
1634 receipt_places,
1635 vec![
1636 ("n1".into(), "target-a".into()),
1637 ("n2".into(), "target-b".into())
1638 ]
1639 );
1640 let fences: Vec<(String, String, i64)> = conn
1641 .prepare(
1642 "SELECT namespace, model, ann_write_log_seq \
1643 FROM memory_visibility_fences ORDER BY model",
1644 )
1645 .unwrap()
1646 .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))
1647 .unwrap()
1648 .collect::<rusqlite::Result<_>>()
1649 .unwrap();
1650 assert_eq!(
1651 fences,
1652 vec![
1653 ("target-a".into(), "model-a".into(), 10),
1654 ("target-a".into(), "model-b".into(), 11)
1655 ]
1656 );
1657 let foreign_key_errors: i64 = conn
1658 .query_row("SELECT COUNT(*) FROM pragma_foreign_key_check", [], |row| {
1659 row.get(0)
1660 })
1661 .unwrap();
1662 assert_eq!(foreign_key_errors, 0);
1663 }
1664
1665 #[cfg(feature = "vectors")]
1666 #[test]
1667 fn moved_visibility_fences_use_each_exact_upsert_and_preserve_unrelated_receipts() {
1668 crate::extension::ensure_extensions_loaded();
1669 let mut conn = migrated();
1670 conn.pragma_update(None, "foreign_keys", "ON").unwrap();
1671 for (id, namespace) in [
1672 ("moved", "source"),
1673 ("zero", "source"),
1674 ("unfenced", "source"),
1675 ("existing", "target"),
1676 ] {
1677 seed_note(&conn, id, namespace, "memory");
1678 }
1679 conn.execute_batch(
1680 "INSERT INTO memory_visibility_receipts (namespace, note_id, model_count) VALUES \
1681 ('source', 'moved', 2), ('source', 'zero', 0), ('target', 'existing', 1); \
1682 INSERT INTO memory_visibility_fences (namespace, note_id, model, ann_write_log_seq) VALUES \
1683 ('source', 'moved', 'model-a', 17), ('source', 'moved', 'model-b', 18), \
1684 ('target', 'existing', 'model-a', 99);",
1685 ).unwrap();
1686 for (table, model) in [("vec_model_a", "model-a"), ("vec_model_b", "model-b")] {
1687 conn.execute_batch(&format!(
1688 "CREATE VIRTUAL TABLE {table} USING vec0(\
1689 subject_id TEXT PRIMARY KEY, namespace TEXT NOT NULL, \
1690 kind TEXT NOT NULL, field TEXT NOT NULL, embedding_model TEXT NOT NULL, \
1691 embedding float[2] distance_metric=cosine)"
1692 ))
1693 .unwrap();
1694 conn.execute(&format!(
1695 "INSERT INTO {table} (subject_id, namespace, kind, field, embedding_model, embedding) \
1696 VALUES ('moved', 'source', 'note', 'note.content', ?1, '[0.1, 0.2]')"
1697 ), [model]).unwrap();
1698 }
1699 conn.execute_batch(
1700 "INSERT INTO vec_model_a (subject_id, namespace, kind, field, embedding_model, embedding) VALUES \
1701 ('unfenced', 'source', 'note', 'note.content', 'model-a', '[0.1, 0.2]'), \
1702 ('existing', 'target', 'note', 'note.content', 'model-a', '[0.1, 0.2]'); \
1703 INSERT INTO ann_write_log (seq, namespace, embedding_model, kind, field, subject_id, op) \
1704 VALUES (100, 'target', 'model-a', 'note', 'note.content', 'existing', 'upsert');",
1705 ).unwrap();
1706 let transaction = conn.transaction().unwrap();
1707 let counts = move_namespace(
1708 &transaction,
1709 &MoveRequest::new("source", vec![route("note:memory", "target")]),
1710 )
1711 .expect("move memory receipts and both vector models");
1712 assert_eq!(counts.ann_log_appended, 6);
1713 for model in ["model-a", "model-b"] {
1714 let fence: i64 = transaction
1715 .query_row(
1716 "SELECT ann_write_log_seq FROM memory_visibility_fences \
1717 WHERE namespace = 'target' AND note_id = 'moved' AND model = ?1",
1718 [model],
1719 |row| row.get(0),
1720 )
1721 .unwrap();
1722 let upsert: i64 = transaction
1723 .query_row(
1724 "SELECT seq FROM ann_write_log WHERE namespace = 'target' \
1725 AND subject_id = 'moved' AND embedding_model = ?1 AND op = 'upsert'",
1726 [model],
1727 |row| row.get(0),
1728 )
1729 .unwrap();
1730 assert!(upsert > 100);
1731 assert_eq!(fence, upsert, "each model owes its own destination upsert");
1732 }
1733 let untouched: i64 = transaction
1734 .query_row(
1735 "SELECT ann_write_log_seq FROM memory_visibility_fences \
1736 WHERE namespace = 'target' AND note_id = 'existing' AND model = 'model-a'",
1737 [],
1738 |row| row.get(0),
1739 )
1740 .unwrap();
1741 assert_eq!(untouched, 99);
1742 let zero_count: i64 = transaction
1743 .query_row(
1744 "SELECT model_count FROM memory_visibility_receipts \
1745 WHERE namespace = 'target' AND note_id = 'zero'",
1746 [],
1747 |row| row.get(0),
1748 )
1749 .unwrap();
1750 assert_eq!(zero_count, 0);
1751 let invented: i64 = transaction.query_row(
1752 "SELECT COUNT(*) FROM memory_visibility_fences WHERE note_id IN ('zero', 'unfenced')",
1753 [], |row| row.get(0),
1754 ).unwrap();
1755 assert_eq!(invented, 0);
1756 transaction.commit().unwrap();
1757 }
1758
1759 #[test]
1763 fn a_class_with_rows_and_no_route_refuses_and_says_how_many() {
1764 let conn = migrated();
1765 seed_note(&conn, "n1", "source", "observation");
1766 seed_note(&conn, "n2", "source", "decision");
1767
1768 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1769 let error = move_namespace(&conn, &request).expect_err("decision notes are unrouted");
1770 match error {
1771 MoveError::UnroutedClass { class, rows } => {
1772 assert_eq!(class, "note:decision");
1773 assert_eq!(rows, 1);
1774 }
1775 other => panic!("expected an unrouted class, got {other}"),
1776 }
1777
1778 let still_here: i64 = conn
1779 .query_row(
1780 "SELECT COUNT(*) FROM notes WHERE namespace = 'source'",
1781 [],
1782 |r| r.get(0),
1783 )
1784 .expect("count");
1785 assert_eq!(still_here, 2, "a refusal writes nothing");
1786 }
1787
1788 #[test]
1793 fn a_namespace_table_this_build_has_no_rule_for_refuses_the_move() {
1794 let conn = migrated();
1795 seed_note(&conn, "n1", "source", "observation");
1796 conn.execute_batch(
1797 "CREATE TABLE later_migration_added_this (\
1798 id TEXT PRIMARY KEY, namespace TEXT NOT NULL);\
1799 INSERT INTO later_migration_added_this VALUES ('x', 'source');",
1800 )
1801 .expect("a migration lands");
1802
1803 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1804 let error = move_namespace(&conn, &request).expect_err("an unknown table refuses");
1805 match error {
1806 MoveError::UnknownTable { table, rows } => {
1807 assert_eq!(table, "later_migration_added_this");
1808 assert_eq!(rows, 1);
1809 }
1810 other => panic!("expected an unknown table, got {other}"),
1811 }
1812 }
1813
1814 #[test]
1826 fn a_table_named_like_a_vector_table_but_not_one_refuses_by_name() {
1827 let conn = migrated();
1828 seed_note(&conn, "n1", "source", "observation");
1829 conn.execute_batch(
1830 "CREATE TABLE vec_audit (\
1831 id TEXT PRIMARY KEY, namespace TEXT NOT NULL);\
1832 INSERT INTO vec_audit VALUES ('x', 'source');",
1833 )
1834 .expect("a migration lands a table whose name starts with the prefix");
1835
1836 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1837 let error = move_namespace(&conn, &request).expect_err("the prefix is not enough");
1838 match error {
1839 MoveError::UnknownTable { table, rows } => {
1840 assert_eq!(table, "vec_audit");
1841 assert_eq!(rows, 1);
1842 }
1843 other => panic!("expected an unknown table, got {other}"),
1844 }
1845 }
1846
1847 #[test]
1851 fn an_unknown_table_holding_nothing_here_does_not_refuse() {
1852 let conn = migrated();
1853 seed_note(&conn, "n1", "source", "observation");
1854 conn.execute_batch(
1855 "CREATE TABLE later_migration_added_this (\
1856 id TEXT PRIMARY KEY, namespace TEXT NOT NULL);\
1857 INSERT INTO later_migration_added_this VALUES ('x', 'somewhere-else');",
1858 )
1859 .expect("a migration lands");
1860
1861 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1862 let counts = move_namespace(&conn, &request).expect("nothing of ours is in that table");
1863 assert_eq!(counts.subjects.get("note:observation"), Some(&1));
1864 }
1865
1866 #[test]
1867 fn a_route_to_the_namespace_it_is_already_in_refuses() {
1868 let conn = migrated();
1869 let request = MoveRequest::new("source", vec![route("atom", "source")]);
1870 match move_namespace(&conn, &request).expect_err("a no-op written as an instruction") {
1871 MoveError::TargetIsSource { class } => assert_eq!(class, "atom"),
1872 other => panic!("expected target-is-source, got {other}"),
1873 }
1874 }
1875
1876 #[test]
1877 fn the_same_class_routed_twice_refuses_rather_than_picking_one() {
1878 let conn = migrated();
1879 let request = MoveRequest::new(
1880 "source",
1881 vec![route("atom", "one"), route("atom", "another")],
1882 );
1883 match move_namespace(&conn, &request).expect_err("two targets, no rule to choose") {
1884 MoveError::DuplicateRoute { class } => assert_eq!(class, "atom"),
1885 other => panic!("expected a duplicate route, got {other}"),
1886 }
1887 }
1888
1889 #[test]
1890 fn an_unknown_route_key_names_what_it_was_given() {
1891 match SubjectClass::parse("notes:observation").expect_err("plural is a typo") {
1892 MoveError::UnknownSubjectClass { key } => assert_eq!(key, "notes:observation"),
1893 other => panic!("expected an unknown class, got {other}"),
1894 }
1895 assert_eq!(
1896 SubjectClass::parse("note:observation").expect("singular"),
1897 SubjectClass::Note("observation".into())
1898 );
1899 }
1900
1901 #[test]
1906 fn a_partitioning_move_reports_the_aggregates_it_leaves_behind() {
1907 let conn = migrated();
1908 seed_note(&conn, "n1", "source", "observation");
1909 seed_note(&conn, "n2", "source", "decision");
1910 conn.execute(
1911 "INSERT INTO brain_profile_snapshots (profile_id, namespace, snapshot_json, updated_at) \
1912 VALUES ('p', 'source', '{}', 1)",
1913 [],
1914 )
1915 .expect("seed a snapshot");
1916
1917 let request = MoveRequest::new(
1918 "source",
1919 vec![
1920 route("note:observation", "one"),
1921 route("note:decision", "another"),
1922 ],
1923 );
1924 let counts = move_namespace(&conn, &request).expect("a partitioning move");
1925 assert_eq!(counts.left_behind.get("brain_profile_snapshots"), Some(&1));
1926
1927 let stayed: i64 = conn
1928 .query_row(
1929 "SELECT COUNT(*) FROM brain_profile_snapshots WHERE namespace = 'source'",
1930 [],
1931 |r| r.get(0),
1932 )
1933 .expect("count");
1934 assert_eq!(stayed, 1);
1935 }
1936
1937 #[test]
1940 fn a_total_single_target_move_carries_the_aggregates() {
1941 let conn = migrated();
1942 seed_note(&conn, "n1", "source", "observation");
1943 conn.execute(
1944 "INSERT INTO brain_profile_snapshots (profile_id, namespace, snapshot_json, updated_at) \
1945 VALUES ('p', 'source', '{}', 1)",
1946 [],
1947 )
1948 .expect("seed a snapshot");
1949
1950 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1951 let counts = move_namespace(&conn, &request).expect("a total move");
1952 assert!(counts.left_behind.is_empty(), "{:?}", counts.left_behind);
1953
1954 let moved: i64 = conn
1955 .query_row(
1956 "SELECT COUNT(*) FROM brain_profile_snapshots WHERE namespace = 'target'",
1957 [],
1958 |r| r.get(0),
1959 )
1960 .expect("count");
1961 assert_eq!(moved, 1);
1962 }
1963 #[cfg(feature = "vectors")]
1976 #[test]
1977 fn a_vector_move_tells_the_source_side_to_drop_what_left() {
1978 crate::extension::ensure_extensions_loaded();
1981 let conn = migrated();
1982 seed_note(&conn, "n1", "source", "observation");
1983 conn.execute_batch(
1984 "CREATE VIRTUAL TABLE vec_test_model USING vec0(\
1985 subject_id TEXT PRIMARY KEY, \
1986 namespace TEXT NOT NULL, \
1987 kind TEXT NOT NULL, \
1988 field TEXT NOT NULL, \
1989 embedding_model TEXT NOT NULL, \
1990 embedding float[4] distance_metric=cosine\
1991 )",
1992 )
1993 .expect("the vector table an embedding model creates at runtime");
1994 conn.execute(
1995 "INSERT INTO vec_test_model \
1996 (subject_id, namespace, kind, field, embedding_model, embedding) \
1997 VALUES ('n1', 'source', 'observation', 'content', 'test-model', \
1998 '[0.1, 0.2, 0.3, 0.4]')",
1999 [],
2000 )
2001 .expect("seed a vector");
2002
2003 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
2004 let counts = move_namespace(&conn, &request).expect("a total move");
2005
2006 assert_eq!(
2007 counts.rows.get("vec_test_model"),
2008 Some(&1),
2009 "the vector itself moved"
2010 );
2011 let left_in_source: i64 = conn
2012 .query_row(
2013 "SELECT COUNT(*) FROM vec_test_model WHERE namespace = 'source'",
2014 [],
2015 |r| r.get(0),
2016 )
2017 .expect("count");
2018 assert_eq!(left_in_source, 0);
2019
2020 assert_eq!(counts.ann_log_appended, 2);
2023 let dropped_from_source: i64 = conn
2024 .query_row(
2025 "SELECT COUNT(*) FROM ann_write_log \
2026 WHERE namespace = 'source' AND op = 'delete' \
2027 AND subject_id = 'n1' AND embedding_model = 'test-model'",
2028 [],
2029 |r| r.get(0),
2030 )
2031 .expect("count");
2032 assert_eq!(
2033 dropped_from_source, 1,
2034 "the source consumer is never told to drop the vector, so its index \
2035 keeps answering with a subject that has left the namespace"
2036 );
2037 let taken_by_target: i64 = conn
2038 .query_row(
2039 "SELECT COUNT(*) FROM ann_write_log \
2040 WHERE namespace = 'target' AND op = 'upsert' \
2041 AND subject_id = 'n1' AND embedding_model = 'test-model'",
2042 [],
2043 |r| r.get(0),
2044 )
2045 .expect("count");
2046 assert_eq!(taken_by_target, 1);
2047 }
2048
2049 #[cfg(feature = "vectors")]
2050 #[tokio::test]
2051 async fn vector_provenance_namespace_move_clears_sidecar() {
2052 use std::sync::Arc;
2053
2054 use khive_storage::VectorStore;
2055
2056 use crate::pool::{ConnectionPool, PoolConfig};
2057 use crate::stores::vectors::SqliteVecStore;
2058
2059 crate::extension::ensure_extensions_loaded();
2060 let dir = tempfile::tempdir().expect("tempdir");
2061 let path = dir.path().join("namespace-move-provenance.db");
2062 let mut conn = Connection::open(&path).expect("open database");
2063 run_migrations(&mut conn).expect("migrate database");
2064 let subject_id = uuid::Uuid::new_v4();
2065 let subject = subject_id.to_string();
2066 seed_note(&conn, &subject, "source", "observation");
2067 conn.execute_batch(
2068 "CREATE VIRTUAL TABLE vec_test_model USING vec0(\
2069 subject_id TEXT PRIMARY KEY, namespace TEXT NOT NULL, \
2070 kind TEXT NOT NULL, field TEXT NOT NULL, \
2071 embedding_model TEXT NOT NULL, embedding float[2] distance_metric=cosine)",
2072 )
2073 .unwrap();
2074 conn.execute(
2075 "INSERT INTO vec_test_model \
2076 (subject_id, namespace, kind, field, embedding_model, embedding) \
2077 VALUES (?1, 'source', 'observation', 'content', 'test-model', '[0.1, 0.2]')",
2078 [&subject],
2079 )
2080 .unwrap();
2081 let before_blob: Vec<u8> = conn
2082 .query_row(
2083 "SELECT embedding FROM vec_test_model WHERE subject_id = ?1",
2084 [&subject],
2085 |row| row.get(0),
2086 )
2087 .unwrap();
2088 let digest = blake3::hash(&before_blob).to_hex().to_string();
2089 conn.execute(
2090 "INSERT INTO vector_provenance \
2091 (model_key, subject_id, namespace, embedding_digest, text_fingerprint, updated_at) \
2092 VALUES ('test_model', ?1, 'source', ?2, ?3, '2026-09-25T12:34:56Z')",
2093 rusqlite::params![&subject, &digest, "a".repeat(64)],
2094 )
2095 .unwrap();
2096
2097 let stale_target_id = uuid::Uuid::new_v4();
2100 let stale_target_subject = stale_target_id.to_string();
2101 seed_note(&conn, &stale_target_subject, "source", "observation");
2102 conn.execute(
2103 "INSERT INTO vec_test_model \
2104 (subject_id, namespace, kind, field, embedding_model, embedding) \
2105 VALUES (?1, 'source', 'observation', 'content', 'test-model', '[0.1, 0.2]')",
2106 [&stale_target_subject],
2107 )
2108 .unwrap();
2109 conn.execute(
2110 "INSERT INTO vector_provenance \
2111 (model_key, subject_id, namespace, embedding_digest, text_fingerprint, updated_at) \
2112 VALUES ('test_model', ?1, 'target', ?2, ?3, '2026-09-25T12:34:56Z')",
2113 rusqlite::params![&stale_target_subject, &digest, "b".repeat(64)],
2114 )
2115 .unwrap();
2116
2117 conn.execute_batch("BEGIN IMMEDIATE").unwrap();
2118 move_namespace(
2119 &conn,
2120 &MoveRequest::new("source", vec![route("note:observation", "target")]),
2121 )
2122 .unwrap();
2123 conn.execute_batch("COMMIT").unwrap();
2124 let sidecars: i64 = conn
2125 .query_row(
2126 "SELECT COUNT(*) FROM vector_provenance \
2127 WHERE model_key = 'test_model' AND subject_id IN (?1, ?2)",
2128 rusqlite::params![&subject, &stale_target_subject],
2129 |row| row.get(0),
2130 )
2131 .unwrap();
2132 assert_eq!(
2133 sidecars, 0,
2134 "a move must invalidate even an equal-BLOB sidecar"
2135 );
2136 let after_blob: Vec<u8> = conn
2137 .query_row(
2138 "SELECT embedding FROM vec_test_model \
2139 WHERE subject_id = ?1 AND namespace = 'target'",
2140 [&subject],
2141 |row| row.get(0),
2142 )
2143 .unwrap();
2144 assert_eq!(after_blob, before_blob, "the vector itself must survive");
2145 drop(conn);
2146
2147 let pool = Arc::new(
2148 ConnectionPool::new(PoolConfig {
2149 path: Some(path),
2150 write_queue_enabled: Some(false),
2151 ..PoolConfig::for_test()
2152 })
2153 .expect("reopen moved database"),
2154 );
2155 let vectors = SqliteVecStore::new(
2156 pool,
2157 true,
2158 "test_model".into(),
2159 "test-model".into(),
2160 2,
2161 "target".into(),
2162 )
2163 .expect("open target vector store");
2164 let observed = vectors
2165 .provenance(subject_id)
2166 .await
2167 .expect("read moved vector")
2168 .expect("moved vector remains present");
2169 assert_eq!(observed.text_fingerprint, None);
2170 assert_eq!(observed.updated_at, None);
2171 let stale_target = vectors
2172 .provenance(stale_target_id)
2173 .await
2174 .expect("read formerly stale target vector")
2175 .expect("second moved vector remains present");
2176 assert_eq!(stale_target.text_fingerprint, None);
2177 assert_eq!(stale_target.updated_at, None);
2178 }
2179
2180 #[cfg(feature = "vectors")]
2191 #[test]
2192 fn a_vector_the_target_already_held_is_not_in_the_instructions() {
2193 crate::extension::ensure_extensions_loaded();
2194 let conn = migrated();
2195 seed_note(&conn, "n1", "source", "observation");
2196 conn.execute_batch(
2197 "CREATE VIRTUAL TABLE vec_test_model USING vec0(\
2198 subject_id TEXT PRIMARY KEY, \
2199 namespace TEXT NOT NULL, \
2200 kind TEXT NOT NULL, \
2201 field TEXT NOT NULL, \
2202 embedding_model TEXT NOT NULL, \
2203 embedding float[4] distance_metric=cosine\
2204 )",
2205 )
2206 .expect("the vector table an embedding model creates at runtime");
2207 conn.execute(
2208 "INSERT INTO vec_test_model \
2209 (subject_id, namespace, kind, field, embedding_model, embedding) \
2210 VALUES ('n1', 'source', 'observation', 'content', 'test-model', \
2211 '[0.1, 0.2, 0.3, 0.4]')",
2212 [],
2213 )
2214 .expect("the vector that moves");
2215 conn.execute(
2216 "INSERT INTO vec_test_model \
2217 (subject_id, namespace, kind, field, embedding_model, embedding) \
2218 VALUES ('already-there', 'target', 'observation', 'content', 'test-model', \
2219 '[0.5, 0.6, 0.7, 0.8]')",
2220 [],
2221 )
2222 .expect("a vector the target already holds");
2223
2224 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
2225 let counts = move_namespace(&conn, &request).expect("a total move");
2226
2227 assert_eq!(
2228 counts.ann_log_appended, 2,
2229 "two instructions are owed for the one vector that moved"
2230 );
2231 let about_the_resident: i64 = conn
2232 .query_row(
2233 "SELECT COUNT(*) FROM ann_write_log WHERE subject_id = 'already-there'",
2234 [],
2235 |r| r.get(0),
2236 )
2237 .expect("count");
2238 assert_eq!(
2239 about_the_resident, 0,
2240 "a vector that did not move is told nothing, and is certainly not \
2241 dropped from a namespace it was never in"
2242 );
2243 let about_the_mover: i64 = conn
2246 .query_row(
2247 "SELECT COUNT(*) FROM ann_write_log WHERE subject_id = 'n1'",
2248 [],
2249 |r| r.get(0),
2250 )
2251 .expect("count");
2252 assert_eq!(about_the_mover, 2);
2253 }
2254 #[test]
2266 fn two_routes_to_one_target_report_a_shared_collision_once() {
2267 let conn = migrated();
2268 seed_note(&conn, "n1", "source", "observation");
2273 seed_note(&conn, "n2", "source", "insight");
2274 for (id, namespace) in [("a1", "source"), ("a2", "target")] {
2275 conn.execute(
2276 "INSERT INTO knowledge_atoms \
2277 (id, namespace, slug, name, created_at, updated_at) \
2278 VALUES (?1, ?2, 'shared-slug', 'an atom', 1, 1)",
2279 rusqlite::params![id, namespace],
2280 )
2281 .expect("seed an atom on each side of the move");
2282 }
2283
2284 let request = MoveRequest::new(
2285 "source",
2286 vec![
2287 route("note:observation", "target"),
2288 route("note:insight", "target"),
2289 route("atom", "target"),
2290 ],
2291 );
2292 let error = move_namespace(&conn, &request).expect_err("the pre-flight refuses");
2293 let MoveError::Collisions { collisions } = error else {
2294 panic!("expected a named collision list, got {error:?}");
2295 };
2296
2297 assert_eq!(
2298 collisions.len(),
2299 1,
2300 "three routes share one target, so the one blocking row is reported \
2301 once: {collisions:?}"
2302 );
2303 assert_eq!(collisions[0].table, "knowledge_atoms");
2304 assert_eq!(collisions[0].constraint, "idx_knowledge_atoms_ns_slug");
2305 assert_eq!(collisions[0].key, "shared-slug");
2306 }
2307}
2308
2309#[cfg(test)]
2310#[test]
2311fn issue2673_namespace_move_advances_entity_version_without_changing_timestamp() {
2312 let mut conn = Connection::open_in_memory().unwrap();
2313 crate::migrations::run_migrations(&mut conn).unwrap();
2314 conn.execute("INSERT INTO entities(id,namespace,kind,name,created_at,updated_at) VALUES('versioned','source','concept','moved',7,7)", []).unwrap();
2315 conn.execute(
2316 "INSERT INTO notes(id,namespace,kind,content,created_at,updated_at) \
2317 VALUES('note-moved','source','observation','moved',11,11), \
2318 ('note-control','unrelated','observation','unchanged',13,13)",
2319 [],
2320 )
2321 .unwrap();
2322 let request = MoveRequest::new(
2323 "source",
2324 vec![
2325 MoveRoute {
2326 class: SubjectClass::Entity("concept".into()),
2327 target: "target".into(),
2328 },
2329 MoveRoute {
2330 class: SubjectClass::Note("observation".into()),
2331 target: "target".into(),
2332 },
2333 ],
2334 );
2335 let tx = conn.transaction().unwrap();
2336 let moved = move_namespace(&tx, &request).unwrap();
2337 assert_eq!(moved.subjects.get("entity:concept"), Some(&1));
2338 assert_eq!(moved.subjects.get("note:observation"), Some(&1));
2339 tx.commit().unwrap();
2340 let stored = conn
2341 .query_row(
2342 "SELECT namespace,updated_at,version FROM entities WHERE id='versioned'",
2343 [],
2344 |row| {
2345 Ok((
2346 row.get::<_, String>(0)?,
2347 row.get::<_, i64>(1)?,
2348 row.get::<_, i64>(2)?,
2349 ))
2350 },
2351 )
2352 .unwrap();
2353 assert_eq!(stored, ("target".into(), 7, 2));
2354 for (id, namespace, timestamp, version) in [
2355 ("note-moved", "target", 11_i64, 2_i64),
2356 ("note-control", "unrelated", 13_i64, 1_i64),
2357 ] {
2358 let stored = conn
2359 .query_row(
2360 "SELECT namespace,updated_at,version FROM notes WHERE id=?1",
2361 [id],
2362 |row| {
2363 Ok((
2364 row.get::<_, String>(0)?,
2365 row.get::<_, i64>(1)?,
2366 row.get::<_, i64>(2)?,
2367 ))
2368 },
2369 )
2370 .unwrap();
2371 assert_eq!(stored, (namespace.into(), timestamp, version));
2372 }
2373
2374 let tx = conn.transaction().unwrap();
2375 let repeated = move_namespace(&tx, &request).unwrap();
2376 assert_eq!(repeated.subjects.get("entity:concept"), Some(&0));
2377 assert_eq!(repeated.subjects.get("note:observation"), Some(&0));
2378 tx.commit().unwrap();
2379 for (sql, expected) in [
2380 ("SELECT version FROM entities WHERE id='versioned'", 2_i64),
2381 ("SELECT version FROM notes WHERE id='note-moved'", 2_i64),
2382 ("SELECT version FROM notes WHERE id='note-control'", 1_i64),
2383 ] {
2384 assert_eq!(
2385 conn.query_row(sql, [], |row| row.get::<_, i64>(0)).unwrap(),
2386 expected
2387 );
2388 }
2389}
2390
2391#[cfg(all(test, feature = "vectors"))]
2392#[path = "namespace_move_partition_tests.rs"]
2393mod partition_tests;
2394
2395#[cfg(test)]
2396#[path = "namespace_move_edge_tests.rs"]
2397mod edge_tests;
2398
2399#[cfg(test)]
2400#[path = "namespace_move_pack_tables_tests.rs"]
2401mod pack_table_tests;
2402
2403#[cfg(all(test, feature = "vectors"))]
2404#[path = "namespace_move_orphan_section_tests.rs"]
2405mod orphan_section_tests;