use std::collections::{BTreeMap, BTreeSet};
use rusqlite::{Connection, OptionalExtension};
use crate::namespace_census::{self, NamespaceCensus, NamespaceConstraint};
#[path = "namespace_move_knowledge.rs"]
mod knowledge_move;
#[path = "namespace_move_pack_tables.rs"]
mod pack_tables;
#[path = "namespace_move_routes.rs"]
mod subject_routes;
use knowledge_move::{move_knowledge_atoms, ordinary_atom_count};
use pack_tables::{move_task_audit, settle_pack_tables};
use subject_routes::routed_subjects;
#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub enum SubjectClass {
Note(String),
Entity(String),
Edge,
EdgeRelation(String),
Atom,
Domain,
}
impl SubjectClass {
pub fn parse(key: &str) -> Result<Self, MoveError> {
match key {
"edge" => return Ok(Self::Edge),
"atom" => return Ok(Self::Atom),
"domain" => return Ok(Self::Domain),
_ => {}
}
match key.split_once(':') {
Some(("note", kind)) if !kind.is_empty() => Ok(Self::Note(kind.to_string())),
Some(("entity", kind)) if !kind.is_empty() => Ok(Self::Entity(kind.to_string())),
Some(("edge", relation))
if khive_types::EdgeRelation::VALID_NAMES.contains(&relation) =>
{
Ok(Self::EdgeRelation(relation.to_string()))
}
_ => Err(MoveError::UnknownSubjectClass {
key: key.to_string(),
}),
}
}
pub fn render(&self) -> String {
match self {
Self::Note(kind) => format!("note:{kind}"),
Self::Entity(kind) => format!("entity:{kind}"),
Self::Edge => "edge".to_string(),
Self::EdgeRelation(relation) => format!("edge:{relation}"),
Self::Atom => "atom".to_string(),
Self::Domain => "domain".to_string(),
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct MoveRoute {
pub class: SubjectClass,
pub target: String,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct MoveRequest {
pub source: String,
pub routes: Vec<MoveRoute>,
pub now_micros: Option<i64>,
}
impl MoveRequest {
pub fn new(source: impl Into<String>, routes: Vec<MoveRoute>) -> Self {
Self {
source: source.into(),
routes,
now_micros: None,
}
}
pub fn at(mut self, now_micros: i64) -> Self {
self.now_micros = Some(now_micros);
self
}
fn route_for(&self, class: &SubjectClass) -> Option<&MoveRoute> {
self.routes
.iter()
.find(|route| &route.class == class)
.or_else(|| match class {
SubjectClass::EdgeRelation(_) => self
.routes
.iter()
.find(|route| route.class == SubjectClass::Edge),
_ => None,
})
}
fn single_target(&self) -> Option<&str> {
let mut targets = self.routes.iter().map(|r| r.target.as_str());
let first = targets.next()?;
targets.all(|t| t == first).then_some(first)
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct MoveCounts {
pub subjects: BTreeMap<String, u64>,
pub rows: BTreeMap<String, u64>,
pub left_behind: BTreeMap<String, u64>,
pub ann_log_appended: u64,
pub live_policies_left_behind: u64,
pub grants_in_force_left_behind: u64,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct StreamMember {
pub note_id: String,
pub stream: String,
pub seq: i64,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Collision {
pub table: String,
pub constraint: String,
pub target: String,
pub key: String,
}
#[derive(Debug)]
pub enum MoveError {
UnknownSubjectClass {
key: String,
},
DuplicateRoute {
class: String,
},
TargetIsSource {
class: String,
},
UnknownTable {
table: String,
rows: u64,
},
UnroutedClass {
class: String,
rows: u64,
},
StreamMembers {
notes: Vec<StreamMember>,
},
UnroutableVector {
subject_id: String,
table: String,
destinations: u64,
},
Collisions {
collisions: Vec<Collision>,
},
Sqlite(rusqlite::Error),
}
impl From<rusqlite::Error> for MoveError {
fn from(error: rusqlite::Error) -> Self {
Self::Sqlite(error)
}
}
impl std::fmt::Display for MoveError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::UnknownSubjectClass { key } => write!(
f,
"no subject class named {key:?}; expected note:<kind>, entity:<kind>, edge:<relation>, edge, atom or domain"
),
Self::DuplicateRoute { class } => {
write!(f, "{class} is routed more than once")
}
Self::TargetIsSource { class } => {
write!(f, "{class} is routed to the namespace it is already in")
}
Self::UnknownTable { table, rows } => write!(
f,
"{table} holds {rows} row(s) in the source namespace and this build has no rule \
for it; a table added by a later migration refuses a move rather than being \
left behind by one"
),
Self::UnroutedClass { class, rows } => write!(
f,
"{class} has {rows} row(s) in the source namespace and no route; \
a class routed with zero rows succeeds reporting zero, an unrouted one refuses"
),
Self::StreamMembers { notes } => {
write!(f, "{} note(s) belong to a stream and cannot change namespace: ", notes.len())?;
for (i, member) in notes.iter().enumerate() {
if i > 0 {
f.write_str(", ")?;
}
write!(f, "{} ({} seq {})", member.note_id, member.stream, member.seq)?;
}
Ok(())
}
Self::UnroutableVector { subject_id, table, destinations } => write!(
f,
"unroutable_vector: {table:?} subject {subject_id:?} has {destinations} source-subject destinations; a partitioned move requires exactly one"
),
Self::Collisions { collisions } => {
write!(f, "{} collision(s): ", collisions.len())?;
for (i, collision) in collisions.iter().enumerate() {
if i > 0 {
f.write_str(", ")?;
}
write!(
f,
"{}.{} in {} already holds {}",
collision.table, collision.constraint, collision.target, collision.key
)?;
}
Ok(())
}
Self::Sqlite(error) => write!(f, "{error}"),
}
}
}
impl std::error::Error for MoveError {}
const NAMESPACE_SCOPED_TABLES: &[&str] = &[
"brain_profile_snapshots",
"brain_event_log",
"proposals_open",
];
const SUBJECT_KEYED_TABLES: &[&str] = &["brain_implicit_mass", "brain_serve_ledger"];
const MEMORY_VISIBILITY_RECEIPTS: &str = "memory_visibility_receipts";
const MEMORY_VISIBILITY_FENCES: &str = "memory_visibility_fences";
const MEMORY_VISIBILITY_EPOCHS: &str = "memory_visibility_epochs";
const LEAVE_BEHIND_TABLES: &[&str] = &["comm_sender_transport", "sessions", "session_messages"];
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum TableDisposition {
Subject,
Derived { trigger_maintained: bool },
DerivedRowidMap,
Appended,
History,
ConsumerWatermark,
RefusedBySchema,
NamespaceScopedAggregate,
SubjectKeyed { subject_column: &'static str },
LeaveBehind,
}
pub fn disposition(table: &namespace_census::NamespaceTable) -> Option<TableDisposition> {
use TableDisposition::*;
Some(match table.name.as_str() {
"notes" | "entities" | "graph_edges" | "knowledge_atoms" | "knowledge_domains" => Subject,
"knowledge_sections" => SubjectKeyed {
subject_column: "atom_id",
},
"proposals_open" => NamespaceScopedAggregate,
"fts_knowledge" | "fts_sections" => Derived {
trigger_maintained: true,
},
"fts_notes" | "fts_entities" => Derived {
trigger_maintained: false,
},
"fts_notes_rowids" | "fts_entities_rowids" => DerivedRowidMap,
"vector_provenance" => Derived {
trigger_maintained: false,
},
"ann_write_log" => Appended,
"events" => History,
"ann_consumer_watermark" | "ann_consumer_pending" => ConsumerWatermark,
"note_streams" => RefusedBySchema,
"brain_profile_snapshots" | "brain_event_log" => NamespaceScopedAggregate,
"brain_implicit_mass" | "brain_serve_ledger" => SubjectKeyed {
subject_column: "target_id",
},
"memory_visibility_receipts" => SubjectKeyed {
subject_column: "note_id",
},
"memory_visibility_fences" => SubjectKeyed {
subject_column: "note_id",
},
"memory_visibility_epochs" => SubjectKeyed {
subject_column: "note_id",
},
"gtd_lifecycle_audit" => SubjectKeyed {
subject_column: "note_id",
},
"knowledge_eval_runs" => NamespaceScopedAggregate,
"exec_runs" | "exec_events" | "git_receipts" => NamespaceScopedAggregate,
"tool_policy" | "tool_grants" => LeaveBehind,
"retrieval_snapshots" => Derived {
trigger_maintained: false,
},
_ if LEAVE_BEHIND_TABLES.contains(&table.name.as_str()) => LeaveBehind,
_ if is_runtime_vector_table(table) => Derived {
trigger_maintained: false,
},
_ => return None,
})
}
fn is_runtime_vector_table(table: &namespace_census::NamespaceTable) -> bool {
table.name.starts_with("vec_") && table.virtual_table
}
#[derive(Debug, Default)]
struct SourceInventory {
note_kinds: BTreeMap<String, u64>,
entity_kinds: BTreeMap<String, u64>,
edge_relations: BTreeMap<String, u64>,
atoms: u64,
domains: u64,
unknown: Vec<(String, u64)>,
}
fn count_in_namespace(conn: &Connection, table: &str, namespace: &str) -> rusqlite::Result<u64> {
let sql = format!(
"SELECT COUNT(*) FROM {} WHERE namespace = ?1",
namespace_census::quote_ident(table)
);
conn.query_row(&sql, [namespace], |row| row.get::<_, i64>(0))
.map(|n| n as u64)
}
fn kinds_in_namespace(
conn: &Connection,
table: &str,
namespace: &str,
) -> rusqlite::Result<BTreeMap<String, u64>> {
let sql = format!(
"SELECT kind, COUNT(*) FROM {} WHERE namespace = ?1 GROUP BY kind",
namespace_census::quote_ident(table)
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt.query_map([namespace], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)? as u64))
})?;
let mut out = BTreeMap::new();
for row in rows {
let (kind, count) = row?;
out.insert(kind, count);
}
Ok(out)
}
fn edge_relations_in_namespace(
conn: &Connection,
namespace: &str,
) -> rusqlite::Result<BTreeMap<String, u64>> {
conn.prepare(
"SELECT relation, COUNT(*) FROM graph_edges WHERE namespace = ?1 GROUP BY relation",
)?
.query_map([namespace], |row| {
Ok((row.get(0)?, row.get::<_, i64>(1)? as u64))
})?
.collect()
}
fn read_source(
conn: &Connection,
census: &NamespaceCensus,
source: &str,
) -> Result<SourceInventory, MoveError> {
let mut inventory = SourceInventory {
note_kinds: kinds_in_namespace(conn, "notes", source)?,
entity_kinds: kinds_in_namespace(conn, "entities", source)?,
edge_relations: edge_relations_in_namespace(conn, source)?,
atoms: ordinary_atom_count(conn, source)?,
domains: count_in_namespace(conn, "knowledge_domains", source)?,
unknown: Vec::new(),
};
for table in &census.tables {
if disposition(table).is_some() {
continue;
}
let count = count_in_namespace(conn, &table.name, source)?;
if count > 0 {
inventory.unknown.push((table.name.clone(), count));
}
}
Ok(inventory)
}
fn stream_members(conn: &Connection, source: &str) -> rusqlite::Result<Vec<StreamMember>> {
let mut stmt = conn.prepare(
"SELECT note_id, stream, seq FROM note_streams \
WHERE namespace = ?1 ORDER BY stream, seq",
)?;
let rows = stmt.query_map([source], |row| {
Ok(StreamMember {
note_id: row.get(0)?,
stream: row.get(1)?,
seq: row.get(2)?,
})
})?;
rows.collect()
}
pub fn validate(
conn: &Connection,
census: &NamespaceCensus,
request: &MoveRequest,
) -> Result<(), MoveError> {
let mut seen = BTreeSet::new();
for route in &request.routes {
if let SubjectClass::EdgeRelation(relation) = &route.class {
if !khive_types::EdgeRelation::VALID_NAMES.contains(&relation.as_str()) {
return Err(MoveError::UnknownSubjectClass {
key: route.class.render(),
});
}
}
if !seen.insert(route.class.clone()) {
return Err(MoveError::DuplicateRoute {
class: route.class.render(),
});
}
if route.target == request.source {
return Err(MoveError::TargetIsSource {
class: route.class.render(),
});
}
}
let inventory = read_source(conn, census, &request.source)?;
if let Some((table, rows)) = inventory.unknown.first() {
return Err(MoveError::UnknownTable {
table: table.clone(),
rows: *rows,
});
}
let unrouted = |class: SubjectClass, rows: u64| -> Result<(), MoveError> {
if rows > 0 && request.route_for(&class).is_none() {
return Err(MoveError::UnroutedClass {
class: class.render(),
rows,
});
}
Ok(())
};
for (kind, rows) in &inventory.note_kinds {
unrouted(SubjectClass::Note(kind.clone()), *rows)?;
}
for (kind, rows) in &inventory.entity_kinds {
unrouted(SubjectClass::Entity(kind.clone()), *rows)?;
}
for (relation, rows) in &inventory.edge_relations {
unrouted(SubjectClass::EdgeRelation(relation.clone()), *rows)?;
}
unrouted(SubjectClass::Atom, inventory.atoms)?;
unrouted(SubjectClass::Domain, inventory.domains)?;
let pinned = stream_members(conn, &request.source)?;
if !pinned.is_empty() {
return Err(MoveError::StreamMembers { notes: pinned });
}
validate_partitioned_vector_moves(conn, census, request)?;
Ok(())
}
fn validate_partitioned_vector_moves(
conn: &Connection,
census: &NamespaceCensus,
request: &MoveRequest,
) -> Result<(), MoveError> {
if request.single_target().is_some() {
return Ok(());
}
let (selector, parameters) = routed_subjects(request);
for table in census
.tables
.iter()
.filter(|table| is_runtime_vector_table(table))
{
let unresolved = conn
.query_row(
&format!(
"WITH routed AS ({selector}) \
SELECT vector.subject_id, COUNT(DISTINCT routed.target) \
FROM {} AS vector LEFT JOIN routed ON routed.subject_id = vector.subject_id \
WHERE vector.namespace = ?1 GROUP BY vector.subject_id \
HAVING COUNT(DISTINCT routed.target) != 1 \
ORDER BY vector.subject_id LIMIT 1",
namespace_census::quote_ident(&table.name)
),
rusqlite::params_from_iter(¶meters),
|row| Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)? as u64)),
)
.optional()?;
if let Some((subject_id, destinations)) = unresolved {
return Err(MoveError::UnroutableVector {
subject_id,
table: table.name.clone(),
destinations,
});
}
}
Ok(())
}
fn edge_collision_target_predicate(
request: &MoveRequest,
target: &str,
parameters: &mut Vec<rusqlite::types::Value>,
) -> String {
let fallback = request
.route_for(&SubjectClass::Edge)
.is_some_and(|route| route.target == target);
let mut matching = Vec::new();
let mut specific = Vec::new();
for route in &request.routes {
if let SubjectClass::EdgeRelation(relation) = &route.class {
if route.target != target && !fallback {
continue;
}
parameters.push(relation.clone().into());
let parameter = format!("?{}", parameters.len());
if route.target == target {
matching.push(parameter.clone());
}
if fallback {
specific.push(parameter);
}
}
}
let mut alternatives = Vec::new();
if !matching.is_empty() {
alternatives.push(format!("source.relation IN ({})", matching.join(", ")));
}
if fallback {
alternatives.push(if specific.is_empty() {
"1".to_owned()
} else {
format!("source.relation NOT IN ({})", specific.join(", "))
});
}
if alternatives.is_empty() {
"0".to_owned()
} else {
alternatives.join(" OR ")
}
}
fn collisions_for(
conn: &Connection,
constraint: &NamespaceConstraint,
request: &MoveRequest,
target: &str,
) -> rusqlite::Result<Vec<Collision>> {
if constraint.partial || !constraint.columns_are_nameable() {
return Ok(Vec::new());
}
let names: Vec<&str> = constraint
.columns
.iter()
.filter_map(|c| c.as_deref())
.collect();
let others: Vec<&str> = names
.iter()
.copied()
.filter(|c| !c.eq_ignore_ascii_case("namespace"))
.collect();
if others.len() != names.len() - 1 {
return Ok(Vec::new());
}
if others.is_empty() {
return collisions_on_namespace_alone(conn, constraint, &request.source, target);
}
let table = namespace_census::quote_ident(&constraint.table);
let join = others
.iter()
.map(|c| {
let q = namespace_census::quote_ident(c);
format!("target.{q} IS source.{q}")
})
.collect::<Vec<_>>()
.join(" AND ");
let select = others
.iter()
.map(|c| format!("source.{}", namespace_census::quote_ident(c)))
.collect::<Vec<_>>()
.join(", ");
let mut parameters = vec![request.source.clone().into(), target.to_owned().into()];
let mut sql = format!(
"SELECT {select} FROM {table} AS source \
JOIN {table} AS target ON target.namespace = ?2 AND {join} \
WHERE source.namespace = ?1"
);
if constraint.table == "graph_edges" && constraint.index == "idx_graph_edges_unique_triple" {
let predicate = edge_collision_target_predicate(request, target, &mut parameters);
sql.push_str(&format!(" AND ({predicate})"));
}
let mut stmt = conn.prepare(&sql)?;
let column_count = others.len();
let rows = stmt.query_map(rusqlite::params_from_iter(parameters), move |row| {
let mut parts = Vec::with_capacity(column_count);
for i in 0..column_count {
parts.push(match row.get_ref(i)? {
rusqlite::types::ValueRef::Null => "NULL".to_string(),
rusqlite::types::ValueRef::Integer(v) => v.to_string(),
rusqlite::types::ValueRef::Real(v) => v.to_string(),
rusqlite::types::ValueRef::Text(v) => String::from_utf8_lossy(v).into_owned(),
rusqlite::types::ValueRef::Blob(_) => "<blob>".to_string(),
});
}
Ok(parts.join(", "))
})?;
let mut found = Vec::new();
for key in rows {
found.push(Collision {
table: constraint.table.clone(),
constraint: constraint.index.clone(),
target: target.to_string(),
key: key?,
});
}
Ok(found)
}
fn collisions_on_namespace_alone(
conn: &Connection,
constraint: &NamespaceConstraint,
source: &str,
target: &str,
) -> rusqlite::Result<Vec<Collision>> {
let table = namespace_census::quote_ident(&constraint.table);
let sql = format!(
"SELECT (SELECT COUNT(*) FROM {table} WHERE namespace = ?1) \
* (SELECT COUNT(*) FROM {table} WHERE namespace = ?2)"
);
let product: i64 = conn.query_row(&sql, [source, target], |row| row.get(0))?;
Ok(if product > 0 {
vec![Collision {
table: constraint.table.clone(),
constraint: constraint.index.clone(),
target: target.to_string(),
key: "(namespace alone)".to_string(),
}]
} else {
Vec::new()
})
}
struct KindedTables {
base: &'static str,
fts: &'static str,
rowids: &'static str,
}
const NOTE_TABLES: KindedTables = KindedTables {
base: "notes",
fts: "fts_notes",
rowids: "fts_notes_rowids",
};
const ENTITY_TABLES: KindedTables = KindedTables {
base: "entities",
fts: "fts_entities",
rowids: "fts_entities_rowids",
};
fn move_kinded_subject(
conn: &Connection,
tables: &KindedTables,
source: &str,
target: &str,
kind: &str,
rows: &mut BTreeMap<String, u64>,
) -> rusqlite::Result<u64> {
let KindedTables { base, fts, rowids } = *tables;
let statement = if base == "entities" {
"UPDATE entities SET namespace = ?2, version = version + 1 \
WHERE namespace = ?1 AND kind = ?3"
.to_owned()
} else {
format!(
"UPDATE {} SET namespace = ?2 WHERE namespace = ?1 AND kind = ?3",
namespace_census::quote_ident(base)
)
};
let moved = conn.execute(&statement, rusqlite::params![source, target, kind])? as u64;
*rows.entry(base.to_string()).or_default() += moved;
let selector = format!(
"SELECT id FROM {} WHERE namespace = ?2 AND kind = ?3",
namespace_census::quote_ident(base)
);
for derived in [fts, rowids] {
let n = conn.execute(
&format!(
"UPDATE {} SET namespace = ?2 \
WHERE namespace = ?1 AND subject_id IN ({selector})",
namespace_census::quote_ident(derived)
),
rusqlite::params![source, target, kind],
)? as u64;
*rows.entry(derived.to_string()).or_default() += n;
}
Ok(moved)
}
fn move_whole_table(
conn: &Connection,
table: &str,
source: &str,
target: &str,
rows: &mut BTreeMap<String, u64>,
) -> rusqlite::Result<u64> {
let moved = conn.execute(
&format!(
"UPDATE {} SET namespace = ?2 WHERE namespace = ?1",
namespace_census::quote_ident(table)
),
rusqlite::params![source, target],
)? as u64;
*rows.entry(table.to_string()).or_default() += moved;
Ok(moved)
}
fn move_edges(
conn: &Connection,
request: &MoveRequest,
route: &MoveRoute,
rows: &mut BTreeMap<String, u64>,
) -> rusqlite::Result<u64> {
let mut parameters: Vec<rusqlite::types::Value> =
vec![request.source.clone().into(), route.target.clone().into()];
let predicate = if let SubjectClass::EdgeRelation(relation) = &route.class {
parameters.push(relation.clone().into());
"relation = ?3".to_owned()
} else {
let mut excluded = Vec::new();
for specific in &request.routes {
if let SubjectClass::EdgeRelation(relation) = &specific.class {
parameters.push(relation.clone().into());
excluded.push(format!("?{}", parameters.len()));
}
}
if excluded.is_empty() {
return move_whole_table(conn, "graph_edges", &request.source, &route.target, rows);
}
format!("relation NOT IN ({})", excluded.join(", "))
};
let moved = conn.execute(
&format!("UPDATE graph_edges SET namespace = ?2 WHERE namespace = ?1 AND {predicate}"),
rusqlite::params_from_iter(parameters),
)? as u64;
*rows.entry("graph_edges".into()).or_default() += moved;
Ok(moved)
}
fn move_memory_visibility(
conn: &Connection,
source: &str,
target: &str,
rows: &mut BTreeMap<String, u64>,
) -> rusqlite::Result<()> {
let receipts = conn.execute(
"INSERT INTO memory_visibility_receipts (namespace, note_id, model_count) \
SELECT ?2, receipt.note_id, receipt.model_count \
FROM memory_visibility_receipts AS receipt \
JOIN notes AS note ON note.id = receipt.note_id \
WHERE receipt.namespace = ?1 AND note.namespace = ?2",
rusqlite::params![source, target],
)? as u64;
*rows.entry(MEMORY_VISIBILITY_RECEIPTS.into()).or_default() += receipts;
let fences = conn.execute(
"UPDATE memory_visibility_fences SET namespace = ?2 \
WHERE namespace = ?1 AND note_id IN (\
SELECT note_id FROM memory_visibility_receipts WHERE namespace = ?2)",
rusqlite::params![source, target],
)? as u64;
*rows.entry(MEMORY_VISIBILITY_FENCES.into()).or_default() += fences;
let removed = conn.execute(
"DELETE FROM memory_visibility_receipts \
WHERE namespace = ?1 AND note_id IN (\
SELECT note_id FROM memory_visibility_receipts WHERE namespace = ?2)",
rusqlite::params![source, target],
)? as u64;
debug_assert_eq!(removed, receipts);
Ok(())
}
fn move_vectors(
conn: &Connection,
table: &str,
source: &str,
target: Option<&str>,
) -> rusqlite::Result<VectorMove> {
let quoted = namespace_census::quote_ident(table);
let columns = "subject_id, namespace, kind, field, embedding_model, embedding";
conn.execute_batch("DROP TABLE IF EXISTS temp.namespace_move_vectors")?;
let staged = if let Some(target) = target {
conn.execute(
&format!(
"CREATE TEMP TABLE namespace_move_vectors AS \
SELECT {columns}, ?2 AS target FROM {quoted} WHERE namespace = ?1"
),
rusqlite::params![source, target],
)
} else {
conn.execute(
&format!(
"CREATE TEMP TABLE namespace_move_vectors AS \
SELECT vector.subject_id, vector.namespace, vector.kind, vector.field, \
vector.embedding_model, vector.embedding, routed.target \
FROM {quoted} AS vector JOIN (\
SELECT DISTINCT subject_id, target FROM temp.namespace_move_subject_targets\
) AS routed ON routed.subject_id = vector.subject_id \
WHERE vector.namespace = ?1"
),
[source],
)
};
staged?;
let staged_rows: i64 = conn.query_row(
"SELECT COUNT(*) FROM temp.namespace_move_vectors",
[],
|r| r.get(0),
)?;
conn.execute(
&format!("DELETE FROM {quoted} WHERE namespace = ?1"),
[source],
)?;
let inserted = conn.execute(
&format!(
"INSERT INTO {quoted} ({columns}) \
SELECT subject_id, target, kind, field, embedding_model, embedding \
FROM temp.namespace_move_vectors"
),
[],
)? as u64;
debug_assert_eq!(
inserted, staged_rows as u64,
"every staged vector is re-inserted or the move is losing embeddings"
);
let has_provenance: bool = conn.query_row(
"SELECT EXISTS(SELECT 1 FROM sqlite_master \
WHERE type = 'table' AND name = 'vector_provenance')",
[],
|row| row.get(0),
)?;
if has_provenance {
let model_key = table
.strip_prefix("vec_")
.expect("runtime vector tables use the vec_ prefix");
conn.execute(
"DELETE FROM vector_provenance \
WHERE model_key = ?1 \
AND subject_id IN (SELECT subject_id FROM temp.namespace_move_vectors)",
[model_key],
)?;
}
let deleted = conn.execute(
"INSERT INTO ann_write_log (namespace, embedding_model, kind, field, subject_id, op) \
SELECT ?1, embedding_model, kind, field, subject_id, 'delete' \
FROM temp.namespace_move_vectors",
[source],
)? as u64;
let upserts = conn
.prepare(
"INSERT INTO ann_write_log \
(namespace, embedding_model, kind, field, subject_id, op) \
SELECT target, embedding_model, kind, field, subject_id, 'upsert' \
FROM temp.namespace_move_vectors \
RETURNING seq, subject_id, embedding_model, kind, field",
)?
.query_map([], |row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, String>(3)?,
row.get::<_, String>(4)?,
))
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
for (seq, subject, model, kind, field) in &upserts {
if kind == "note" && field == "note.content" {
refresh_moved_memory_fence(conn, source, subject, model, *seq)?;
}
}
let appended = deleted + upserts.len() as u64;
conn.execute_batch("DROP TABLE temp.namespace_move_vectors")?;
Ok(VectorMove {
moved: inserted,
ann_appended: appended,
})
}
fn refresh_moved_memory_fence(
conn: &Connection,
source: &str,
subject: &str,
model: &str,
seq: i64,
) -> rusqlite::Result<()> {
conn.execute(
"UPDATE memory_visibility_fences SET ann_write_log_seq = ?4 \
WHERE namespace = ?1 AND note_id = ?2 AND model = ?3",
rusqlite::params![source, subject, model, seq],
)?;
Ok(())
}
struct VectorMove {
moved: u64,
ann_appended: u64,
}
pub fn move_namespace(conn: &Connection, request: &MoveRequest) -> Result<MoveCounts, MoveError> {
let census = namespace_census::census(conn)?;
validate(conn, &census, request)?;
let mut targets: BTreeSet<&str> = BTreeSet::new();
for route in &request.routes {
targets.insert(route.target.as_str());
}
let mut collisions = Vec::new();
for constraint in namespace_census::reachable_constraints(&census) {
for target in &targets {
collisions.extend(collisions_for(conn, constraint, request, target)?);
}
}
if !collisions.is_empty() {
return Err(MoveError::Collisions { collisions });
}
let mut counts = MoveCounts::default();
let source = request.source.as_str();
if request.single_target().is_none() {
let (selector, parameters) = routed_subjects(request);
conn.execute_batch("DROP TABLE IF EXISTS temp.namespace_move_subject_targets")?;
conn.execute(
&format!("CREATE TEMP TABLE namespace_move_subject_targets AS {selector}"),
rusqlite::params_from_iter(parameters),
)?;
}
move_task_audit(conn, &census, request, &mut counts.rows)?;
for route in &request.routes {
let target = route.target.as_str();
let moved = match &route.class {
SubjectClass::Note(kind) => {
move_kinded_subject(conn, &NOTE_TABLES, source, target, kind, &mut counts.rows)?
}
SubjectClass::Entity(kind) => {
move_kinded_subject(conn, &ENTITY_TABLES, source, target, kind, &mut counts.rows)?
}
SubjectClass::Edge | SubjectClass::EdgeRelation(_) => {
move_edges(conn, request, route, &mut counts.rows)?
}
SubjectClass::Atom => {
move_knowledge_atoms(conn, request, target, false, &mut counts.rows)?
}
SubjectClass::Domain => {
move_knowledge_atoms(conn, request, target, true, &mut counts.rows)?;
move_whole_table(conn, "knowledge_domains", source, target, &mut counts.rows)?
}
};
counts.subjects.insert(route.class.render(), moved);
}
if let Some(target) = request.single_target() {
move_whole_table(conn, "knowledge_sections", source, target, &mut counts.rows)?;
}
for table in &census.tables {
if !is_runtime_vector_table(table) {
continue;
}
let moved = move_vectors(conn, &table.name, source, request.single_target())?;
*counts.rows.entry(table.name.clone()).or_default() += moved.moved;
counts.ann_log_appended += moved.ann_appended;
}
if request.single_target().is_none() {
conn.execute_batch("DROP TABLE temp.namespace_move_subject_targets")?;
}
if census
.tables
.iter()
.any(|table| table.name == "vector_provenance")
{
let left = count_in_namespace(conn, "vector_provenance", source)?;
if left > 0 {
counts.left_behind.insert("vector_provenance".into(), left);
}
}
for table in SUBJECT_KEYED_TABLES {
for target in &targets {
let moved = conn.execute(
&format!(
"UPDATE {} SET namespace = ?2 WHERE namespace = ?1 AND target_id IN (\
SELECT id FROM notes WHERE namespace = ?2 \
UNION ALL SELECT id FROM entities WHERE namespace = ?2 \
UNION ALL SELECT id FROM knowledge_atoms WHERE namespace = ?2)",
namespace_census::quote_ident(table)
),
rusqlite::params![source, target],
)? as u64;
*counts.rows.entry((*table).to_string()).or_default() += moved;
}
let left = count_in_namespace(conn, table, source)?;
if left > 0 {
counts.left_behind.insert((*table).to_string(), left);
}
}
for target in &targets {
let moved = conn.execute(
"UPDATE memory_visibility_epochs SET namespace = ?2 \
WHERE namespace = ?1 AND note_id IN (SELECT id FROM notes WHERE namespace = ?2)",
rusqlite::params![source, target],
)? as u64;
*counts
.rows
.entry(MEMORY_VISIBILITY_EPOCHS.into())
.or_default() += moved;
}
let left = count_in_namespace(conn, MEMORY_VISIBILITY_EPOCHS, source)?;
if left > 0 {
counts
.left_behind
.insert(MEMORY_VISIBILITY_EPOCHS.into(), left);
}
for target in &targets {
move_memory_visibility(conn, source, target, &mut counts.rows)?;
}
for table in [MEMORY_VISIBILITY_RECEIPTS, MEMORY_VISIBILITY_FENCES] {
let left = count_in_namespace(conn, table, source)?;
if left > 0 {
counts.left_behind.insert((*table).to_string(), left);
}
}
for table in LEAVE_BEHIND_TABLES {
let left = count_in_namespace(conn, table, source)?;
if left > 0 {
counts.left_behind.insert((*table).to_string(), left);
}
}
for table in NAMESPACE_SCOPED_TABLES {
match request.single_target() {
Some(target) => {
move_whole_table(conn, table, source, target, &mut counts.rows)?;
}
None => {
let left = count_in_namespace(conn, table, source)?;
if left > 0 {
counts.left_behind.insert((*table).to_string(), left);
}
}
}
}
settle_pack_tables(conn, &census, request, &mut counts)?;
Ok(counts)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::migrations::run_migrations_for_test as run_migrations;
use rusqlite::Connection;
pub(super) fn migrated() -> Connection {
let mut conn = Connection::open_in_memory().expect("open");
run_migrations(&mut conn).expect("migrate");
conn
}
pub(super) fn seed_note(conn: &Connection, id: &str, namespace: &str, kind: &str) {
conn.execute(
"INSERT INTO notes (id, namespace, kind, name, content, created_at, updated_at) \
VALUES (?1, ?2, ?3, 'a name', 'some content', 1, 1)",
rusqlite::params![id, namespace, kind],
)
.expect("seed note");
}
pub(super) fn route(key: &str, target: &str) -> MoveRoute {
MoveRoute {
class: SubjectClass::parse(key).expect("route key"),
target: target.to_string(),
}
}
#[test]
fn a_routed_class_with_no_rows_succeeds_reporting_zero() {
let conn = migrated();
let request = MoveRequest::new(
"empty-source",
vec![route("note:observation", "target"), route("atom", "target")],
);
let counts = move_namespace(&conn, &request).expect("a backend with nothing routed here");
assert_eq!(counts.subjects.get("note:observation"), Some(&0));
assert_eq!(counts.subjects.get("atom"), Some(&0));
assert_eq!(counts.rows.get(MEMORY_VISIBILITY_RECEIPTS), Some(&0));
assert_eq!(counts.rows.get(MEMORY_VISIBILITY_FENCES), Some(&0));
}
#[test]
fn a_memory_visibility_receipt_follows_its_note_and_reports_count() {
let conn = migrated();
conn.pragma_update(None, "foreign_keys", "ON").unwrap();
seed_note(&conn, "n1", "source", "observation");
conn.execute(
"INSERT INTO memory_visibility_receipts (namespace, note_id, model_count) \
VALUES ('source', 'n1', 0)",
[],
)
.unwrap();
let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
let counts = move_namespace(&conn, &request).expect("the receipt follows its note");
assert_eq!(counts.subjects.get("note:observation"), Some(&1));
assert_eq!(counts.rows.get(MEMORY_VISIBILITY_RECEIPTS), Some(&1));
assert_eq!(counts.rows.get(MEMORY_VISIBILITY_FENCES), Some(&0));
let stored: (String, i64) = conn
.query_row(
"SELECT namespace, model_count FROM memory_visibility_receipts \
WHERE note_id = 'n1'",
[],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.unwrap();
assert_eq!(stored, ("target".into(), 0));
let source_rows: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memory_visibility_receipts WHERE namespace = 'source'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(source_rows, 0);
}
#[test]
fn memory_visibility_fences_follow_their_note_with_foreign_keys_enabled() {
let conn = migrated();
conn.pragma_update(None, "foreign_keys", "ON").unwrap();
seed_note(&conn, "n1", "source", "observation");
seed_note(&conn, "n2", "source", "decision");
conn.execute(
"INSERT INTO memory_visibility_receipts (namespace, note_id, model_count) \
VALUES ('source', 'n1', 2), ('source', 'n2', 0)",
[],
)
.unwrap();
conn.execute(
"INSERT INTO memory_visibility_fences \
(namespace, note_id, model, ann_write_log_seq) VALUES \
('source', 'n1', 'model-a', 10), ('source', 'n1', 'model-b', 11)",
[],
)
.unwrap();
let request = MoveRequest::new(
"source",
vec![
route("note:observation", "target-a"),
route("note:decision", "target-b"),
],
);
let counts = move_namespace(&conn, &request).expect("both receipts follow their notes");
assert_eq!(counts.rows.get(MEMORY_VISIBILITY_RECEIPTS), Some(&2));
assert_eq!(counts.rows.get(MEMORY_VISIBILITY_FENCES), Some(&2));
let receipt_places: Vec<(String, String)> = conn
.prepare("SELECT note_id, namespace FROM memory_visibility_receipts ORDER BY note_id")
.unwrap()
.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
.unwrap()
.collect::<rusqlite::Result<_>>()
.unwrap();
assert_eq!(
receipt_places,
vec![
("n1".into(), "target-a".into()),
("n2".into(), "target-b".into())
]
);
let fences: Vec<(String, String, i64)> = conn
.prepare(
"SELECT namespace, model, ann_write_log_seq \
FROM memory_visibility_fences ORDER BY model",
)
.unwrap()
.query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))
.unwrap()
.collect::<rusqlite::Result<_>>()
.unwrap();
assert_eq!(
fences,
vec![
("target-a".into(), "model-a".into(), 10),
("target-a".into(), "model-b".into(), 11)
]
);
let foreign_key_errors: i64 = conn
.query_row("SELECT COUNT(*) FROM pragma_foreign_key_check", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(foreign_key_errors, 0);
}
#[cfg(feature = "vectors")]
#[test]
fn moved_visibility_fences_use_each_exact_upsert_and_preserve_unrelated_receipts() {
crate::extension::ensure_extensions_loaded();
let mut conn = migrated();
conn.pragma_update(None, "foreign_keys", "ON").unwrap();
for (id, namespace) in [
("moved", "source"),
("zero", "source"),
("unfenced", "source"),
("existing", "target"),
] {
seed_note(&conn, id, namespace, "memory");
}
conn.execute_batch(
"INSERT INTO memory_visibility_receipts (namespace, note_id, model_count) VALUES \
('source', 'moved', 2), ('source', 'zero', 0), ('target', 'existing', 1); \
INSERT INTO memory_visibility_fences (namespace, note_id, model, ann_write_log_seq) VALUES \
('source', 'moved', 'model-a', 17), ('source', 'moved', 'model-b', 18), \
('target', 'existing', 'model-a', 99);",
).unwrap();
for (table, model) in [("vec_model_a", "model-a"), ("vec_model_b", "model-b")] {
conn.execute_batch(&format!(
"CREATE VIRTUAL TABLE {table} USING vec0(\
subject_id TEXT PRIMARY KEY, namespace TEXT NOT NULL, \
kind TEXT NOT NULL, field TEXT NOT NULL, embedding_model TEXT NOT NULL, \
embedding float[2] distance_metric=cosine)"
))
.unwrap();
conn.execute(&format!(
"INSERT INTO {table} (subject_id, namespace, kind, field, embedding_model, embedding) \
VALUES ('moved', 'source', 'note', 'note.content', ?1, '[0.1, 0.2]')"
), [model]).unwrap();
}
conn.execute_batch(
"INSERT INTO vec_model_a (subject_id, namespace, kind, field, embedding_model, embedding) VALUES \
('unfenced', 'source', 'note', 'note.content', 'model-a', '[0.1, 0.2]'), \
('existing', 'target', 'note', 'note.content', 'model-a', '[0.1, 0.2]'); \
INSERT INTO ann_write_log (seq, namespace, embedding_model, kind, field, subject_id, op) \
VALUES (100, 'target', 'model-a', 'note', 'note.content', 'existing', 'upsert');",
).unwrap();
let transaction = conn.transaction().unwrap();
let counts = move_namespace(
&transaction,
&MoveRequest::new("source", vec![route("note:memory", "target")]),
)
.expect("move memory receipts and both vector models");
assert_eq!(counts.ann_log_appended, 6);
for model in ["model-a", "model-b"] {
let fence: i64 = transaction
.query_row(
"SELECT ann_write_log_seq FROM memory_visibility_fences \
WHERE namespace = 'target' AND note_id = 'moved' AND model = ?1",
[model],
|row| row.get(0),
)
.unwrap();
let upsert: i64 = transaction
.query_row(
"SELECT seq FROM ann_write_log WHERE namespace = 'target' \
AND subject_id = 'moved' AND embedding_model = ?1 AND op = 'upsert'",
[model],
|row| row.get(0),
)
.unwrap();
assert!(upsert > 100);
assert_eq!(fence, upsert, "each model owes its own destination upsert");
}
let untouched: i64 = transaction
.query_row(
"SELECT ann_write_log_seq FROM memory_visibility_fences \
WHERE namespace = 'target' AND note_id = 'existing' AND model = 'model-a'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(untouched, 99);
let zero_count: i64 = transaction
.query_row(
"SELECT model_count FROM memory_visibility_receipts \
WHERE namespace = 'target' AND note_id = 'zero'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(zero_count, 0);
let invented: i64 = transaction.query_row(
"SELECT COUNT(*) FROM memory_visibility_fences WHERE note_id IN ('zero', 'unfenced')",
[], |row| row.get(0),
).unwrap();
assert_eq!(invented, 0);
transaction.commit().unwrap();
}
#[test]
fn a_class_with_rows_and_no_route_refuses_and_says_how_many() {
let conn = migrated();
seed_note(&conn, "n1", "source", "observation");
seed_note(&conn, "n2", "source", "decision");
let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
let error = move_namespace(&conn, &request).expect_err("decision notes are unrouted");
match error {
MoveError::UnroutedClass { class, rows } => {
assert_eq!(class, "note:decision");
assert_eq!(rows, 1);
}
other => panic!("expected an unrouted class, got {other}"),
}
let still_here: i64 = conn
.query_row(
"SELECT COUNT(*) FROM notes WHERE namespace = 'source'",
[],
|r| r.get(0),
)
.expect("count");
assert_eq!(still_here, 2, "a refusal writes nothing");
}
#[test]
fn a_namespace_table_this_build_has_no_rule_for_refuses_the_move() {
let conn = migrated();
seed_note(&conn, "n1", "source", "observation");
conn.execute_batch(
"CREATE TABLE later_migration_added_this (\
id TEXT PRIMARY KEY, namespace TEXT NOT NULL);\
INSERT INTO later_migration_added_this VALUES ('x', 'source');",
)
.expect("a migration lands");
let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
let error = move_namespace(&conn, &request).expect_err("an unknown table refuses");
match error {
MoveError::UnknownTable { table, rows } => {
assert_eq!(table, "later_migration_added_this");
assert_eq!(rows, 1);
}
other => panic!("expected an unknown table, got {other}"),
}
}
#[test]
fn a_table_named_like_a_vector_table_but_not_one_refuses_by_name() {
let conn = migrated();
seed_note(&conn, "n1", "source", "observation");
conn.execute_batch(
"CREATE TABLE vec_audit (\
id TEXT PRIMARY KEY, namespace TEXT NOT NULL);\
INSERT INTO vec_audit VALUES ('x', 'source');",
)
.expect("a migration lands a table whose name starts with the prefix");
let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
let error = move_namespace(&conn, &request).expect_err("the prefix is not enough");
match error {
MoveError::UnknownTable { table, rows } => {
assert_eq!(table, "vec_audit");
assert_eq!(rows, 1);
}
other => panic!("expected an unknown table, got {other}"),
}
}
#[test]
fn an_unknown_table_holding_nothing_here_does_not_refuse() {
let conn = migrated();
seed_note(&conn, "n1", "source", "observation");
conn.execute_batch(
"CREATE TABLE later_migration_added_this (\
id TEXT PRIMARY KEY, namespace TEXT NOT NULL);\
INSERT INTO later_migration_added_this VALUES ('x', 'somewhere-else');",
)
.expect("a migration lands");
let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
let counts = move_namespace(&conn, &request).expect("nothing of ours is in that table");
assert_eq!(counts.subjects.get("note:observation"), Some(&1));
}
#[test]
fn a_route_to_the_namespace_it_is_already_in_refuses() {
let conn = migrated();
let request = MoveRequest::new("source", vec![route("atom", "source")]);
match move_namespace(&conn, &request).expect_err("a no-op written as an instruction") {
MoveError::TargetIsSource { class } => assert_eq!(class, "atom"),
other => panic!("expected target-is-source, got {other}"),
}
}
#[test]
fn the_same_class_routed_twice_refuses_rather_than_picking_one() {
let conn = migrated();
let request = MoveRequest::new(
"source",
vec![route("atom", "one"), route("atom", "another")],
);
match move_namespace(&conn, &request).expect_err("two targets, no rule to choose") {
MoveError::DuplicateRoute { class } => assert_eq!(class, "atom"),
other => panic!("expected a duplicate route, got {other}"),
}
}
#[test]
fn an_unknown_route_key_names_what_it_was_given() {
match SubjectClass::parse("notes:observation").expect_err("plural is a typo") {
MoveError::UnknownSubjectClass { key } => assert_eq!(key, "notes:observation"),
other => panic!("expected an unknown class, got {other}"),
}
assert_eq!(
SubjectClass::parse("note:observation").expect("singular"),
SubjectClass::Note("observation".into())
);
}
#[test]
fn a_partitioning_move_reports_the_aggregates_it_leaves_behind() {
let conn = migrated();
seed_note(&conn, "n1", "source", "observation");
seed_note(&conn, "n2", "source", "decision");
conn.execute(
"INSERT INTO brain_profile_snapshots (profile_id, namespace, snapshot_json, updated_at) \
VALUES ('p', 'source', '{}', 1)",
[],
)
.expect("seed a snapshot");
let request = MoveRequest::new(
"source",
vec![
route("note:observation", "one"),
route("note:decision", "another"),
],
);
let counts = move_namespace(&conn, &request).expect("a partitioning move");
assert_eq!(counts.left_behind.get("brain_profile_snapshots"), Some(&1));
let stayed: i64 = conn
.query_row(
"SELECT COUNT(*) FROM brain_profile_snapshots WHERE namespace = 'source'",
[],
|r| r.get(0),
)
.expect("count");
assert_eq!(stayed, 1);
}
#[test]
fn a_total_single_target_move_carries_the_aggregates() {
let conn = migrated();
seed_note(&conn, "n1", "source", "observation");
conn.execute(
"INSERT INTO brain_profile_snapshots (profile_id, namespace, snapshot_json, updated_at) \
VALUES ('p', 'source', '{}', 1)",
[],
)
.expect("seed a snapshot");
let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
let counts = move_namespace(&conn, &request).expect("a total move");
assert!(counts.left_behind.is_empty(), "{:?}", counts.left_behind);
let moved: i64 = conn
.query_row(
"SELECT COUNT(*) FROM brain_profile_snapshots WHERE namespace = 'target'",
[],
|r| r.get(0),
)
.expect("count");
assert_eq!(moved, 1);
}
#[cfg(feature = "vectors")]
#[test]
fn a_vector_move_tells_the_source_side_to_drop_what_left() {
crate::extension::ensure_extensions_loaded();
let conn = migrated();
seed_note(&conn, "n1", "source", "observation");
conn.execute_batch(
"CREATE VIRTUAL TABLE vec_test_model USING vec0(\
subject_id TEXT PRIMARY KEY, \
namespace TEXT NOT NULL, \
kind TEXT NOT NULL, \
field TEXT NOT NULL, \
embedding_model TEXT NOT NULL, \
embedding float[4] distance_metric=cosine\
)",
)
.expect("the vector table an embedding model creates at runtime");
conn.execute(
"INSERT INTO vec_test_model \
(subject_id, namespace, kind, field, embedding_model, embedding) \
VALUES ('n1', 'source', 'observation', 'content', 'test-model', \
'[0.1, 0.2, 0.3, 0.4]')",
[],
)
.expect("seed a vector");
let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
let counts = move_namespace(&conn, &request).expect("a total move");
assert_eq!(
counts.rows.get("vec_test_model"),
Some(&1),
"the vector itself moved"
);
let left_in_source: i64 = conn
.query_row(
"SELECT COUNT(*) FROM vec_test_model WHERE namespace = 'source'",
[],
|r| r.get(0),
)
.expect("count");
assert_eq!(left_in_source, 0);
assert_eq!(counts.ann_log_appended, 2);
let dropped_from_source: i64 = conn
.query_row(
"SELECT COUNT(*) FROM ann_write_log \
WHERE namespace = 'source' AND op = 'delete' \
AND subject_id = 'n1' AND embedding_model = 'test-model'",
[],
|r| r.get(0),
)
.expect("count");
assert_eq!(
dropped_from_source, 1,
"the source consumer is never told to drop the vector, so its index \
keeps answering with a subject that has left the namespace"
);
let taken_by_target: i64 = conn
.query_row(
"SELECT COUNT(*) FROM ann_write_log \
WHERE namespace = 'target' AND op = 'upsert' \
AND subject_id = 'n1' AND embedding_model = 'test-model'",
[],
|r| r.get(0),
)
.expect("count");
assert_eq!(taken_by_target, 1);
}
#[cfg(feature = "vectors")]
#[tokio::test]
async fn vector_provenance_namespace_move_clears_sidecar() {
use std::sync::Arc;
use khive_storage::VectorStore;
use crate::pool::{ConnectionPool, PoolConfig};
use crate::stores::vectors::SqliteVecStore;
crate::extension::ensure_extensions_loaded();
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("namespace-move-provenance.db");
let mut conn = Connection::open(&path).expect("open database");
run_migrations(&mut conn).expect("migrate database");
let subject_id = uuid::Uuid::new_v4();
let subject = subject_id.to_string();
seed_note(&conn, &subject, "source", "observation");
conn.execute_batch(
"CREATE VIRTUAL TABLE vec_test_model USING vec0(\
subject_id TEXT PRIMARY KEY, namespace TEXT NOT NULL, \
kind TEXT NOT NULL, field TEXT NOT NULL, \
embedding_model TEXT NOT NULL, embedding float[2] distance_metric=cosine)",
)
.unwrap();
conn.execute(
"INSERT INTO vec_test_model \
(subject_id, namespace, kind, field, embedding_model, embedding) \
VALUES (?1, 'source', 'observation', 'content', 'test-model', '[0.1, 0.2]')",
[&subject],
)
.unwrap();
let before_blob: Vec<u8> = conn
.query_row(
"SELECT embedding FROM vec_test_model WHERE subject_id = ?1",
[&subject],
|row| row.get(0),
)
.unwrap();
let digest = blake3::hash(&before_blob).to_hex().to_string();
conn.execute(
"INSERT INTO vector_provenance \
(model_key, subject_id, namespace, embedding_digest, text_fingerprint, updated_at) \
VALUES ('test_model', ?1, 'source', ?2, ?3, '2026-09-25T12:34:56Z')",
rusqlite::params![&subject, &digest, "a".repeat(64)],
)
.unwrap();
let stale_target_id = uuid::Uuid::new_v4();
let stale_target_subject = stale_target_id.to_string();
seed_note(&conn, &stale_target_subject, "source", "observation");
conn.execute(
"INSERT INTO vec_test_model \
(subject_id, namespace, kind, field, embedding_model, embedding) \
VALUES (?1, 'source', 'observation', 'content', 'test-model', '[0.1, 0.2]')",
[&stale_target_subject],
)
.unwrap();
conn.execute(
"INSERT INTO vector_provenance \
(model_key, subject_id, namespace, embedding_digest, text_fingerprint, updated_at) \
VALUES ('test_model', ?1, 'target', ?2, ?3, '2026-09-25T12:34:56Z')",
rusqlite::params![&stale_target_subject, &digest, "b".repeat(64)],
)
.unwrap();
conn.execute_batch("BEGIN IMMEDIATE").unwrap();
move_namespace(
&conn,
&MoveRequest::new("source", vec![route("note:observation", "target")]),
)
.unwrap();
conn.execute_batch("COMMIT").unwrap();
let sidecars: i64 = conn
.query_row(
"SELECT COUNT(*) FROM vector_provenance \
WHERE model_key = 'test_model' AND subject_id IN (?1, ?2)",
rusqlite::params![&subject, &stale_target_subject],
|row| row.get(0),
)
.unwrap();
assert_eq!(
sidecars, 0,
"a move must invalidate even an equal-BLOB sidecar"
);
let after_blob: Vec<u8> = conn
.query_row(
"SELECT embedding FROM vec_test_model \
WHERE subject_id = ?1 AND namespace = 'target'",
[&subject],
|row| row.get(0),
)
.unwrap();
assert_eq!(after_blob, before_blob, "the vector itself must survive");
drop(conn);
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(path),
write_queue_enabled: Some(false),
..PoolConfig::for_test()
})
.expect("reopen moved database"),
);
let vectors = SqliteVecStore::new(
pool,
true,
"test_model".into(),
"test-model".into(),
2,
"target".into(),
)
.expect("open target vector store");
let observed = vectors
.provenance(subject_id)
.await
.expect("read moved vector")
.expect("moved vector remains present");
assert_eq!(observed.text_fingerprint, None);
assert_eq!(observed.updated_at, None);
let stale_target = vectors
.provenance(stale_target_id)
.await
.expect("read formerly stale target vector")
.expect("second moved vector remains present");
assert_eq!(stale_target.text_fingerprint, None);
assert_eq!(stale_target.updated_at, None);
}
#[cfg(feature = "vectors")]
#[test]
fn a_vector_the_target_already_held_is_not_in_the_instructions() {
crate::extension::ensure_extensions_loaded();
let conn = migrated();
seed_note(&conn, "n1", "source", "observation");
conn.execute_batch(
"CREATE VIRTUAL TABLE vec_test_model USING vec0(\
subject_id TEXT PRIMARY KEY, \
namespace TEXT NOT NULL, \
kind TEXT NOT NULL, \
field TEXT NOT NULL, \
embedding_model TEXT NOT NULL, \
embedding float[4] distance_metric=cosine\
)",
)
.expect("the vector table an embedding model creates at runtime");
conn.execute(
"INSERT INTO vec_test_model \
(subject_id, namespace, kind, field, embedding_model, embedding) \
VALUES ('n1', 'source', 'observation', 'content', 'test-model', \
'[0.1, 0.2, 0.3, 0.4]')",
[],
)
.expect("the vector that moves");
conn.execute(
"INSERT INTO vec_test_model \
(subject_id, namespace, kind, field, embedding_model, embedding) \
VALUES ('already-there', 'target', 'observation', 'content', 'test-model', \
'[0.5, 0.6, 0.7, 0.8]')",
[],
)
.expect("a vector the target already holds");
let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
let counts = move_namespace(&conn, &request).expect("a total move");
assert_eq!(
counts.ann_log_appended, 2,
"two instructions are owed for the one vector that moved"
);
let about_the_resident: i64 = conn
.query_row(
"SELECT COUNT(*) FROM ann_write_log WHERE subject_id = 'already-there'",
[],
|r| r.get(0),
)
.expect("count");
assert_eq!(
about_the_resident, 0,
"a vector that did not move is told nothing, and is certainly not \
dropped from a namespace it was never in"
);
let about_the_mover: i64 = conn
.query_row(
"SELECT COUNT(*) FROM ann_write_log WHERE subject_id = 'n1'",
[],
|r| r.get(0),
)
.expect("count");
assert_eq!(about_the_mover, 2);
}
#[test]
fn two_routes_to_one_target_report_a_shared_collision_once() {
let conn = migrated();
seed_note(&conn, "n1", "source", "observation");
seed_note(&conn, "n2", "source", "insight");
for (id, namespace) in [("a1", "source"), ("a2", "target")] {
conn.execute(
"INSERT INTO knowledge_atoms \
(id, namespace, slug, name, created_at, updated_at) \
VALUES (?1, ?2, 'shared-slug', 'an atom', 1, 1)",
rusqlite::params![id, namespace],
)
.expect("seed an atom on each side of the move");
}
let request = MoveRequest::new(
"source",
vec![
route("note:observation", "target"),
route("note:insight", "target"),
route("atom", "target"),
],
);
let error = move_namespace(&conn, &request).expect_err("the pre-flight refuses");
let MoveError::Collisions { collisions } = error else {
panic!("expected a named collision list, got {error:?}");
};
assert_eq!(
collisions.len(),
1,
"three routes share one target, so the one blocking row is reported \
once: {collisions:?}"
);
assert_eq!(collisions[0].table, "knowledge_atoms");
assert_eq!(collisions[0].constraint, "idx_knowledge_atoms_ns_slug");
assert_eq!(collisions[0].key, "shared-slug");
}
}
#[cfg(test)]
#[test]
fn issue2673_namespace_move_advances_entity_version_without_changing_timestamp() {
let mut conn = Connection::open_in_memory().unwrap();
crate::migrations::run_migrations(&mut conn).unwrap();
conn.execute("INSERT INTO entities(id,namespace,kind,name,created_at,updated_at) VALUES('versioned','source','concept','moved',7,7)", []).unwrap();
conn.execute(
"INSERT INTO notes(id,namespace,kind,content,created_at,updated_at) \
VALUES('note-moved','source','observation','moved',11,11), \
('note-control','unrelated','observation','unchanged',13,13)",
[],
)
.unwrap();
let request = MoveRequest::new(
"source",
vec![
MoveRoute {
class: SubjectClass::Entity("concept".into()),
target: "target".into(),
},
MoveRoute {
class: SubjectClass::Note("observation".into()),
target: "target".into(),
},
],
);
let tx = conn.transaction().unwrap();
let moved = move_namespace(&tx, &request).unwrap();
assert_eq!(moved.subjects.get("entity:concept"), Some(&1));
assert_eq!(moved.subjects.get("note:observation"), Some(&1));
tx.commit().unwrap();
let stored = conn
.query_row(
"SELECT namespace,updated_at,version FROM entities WHERE id='versioned'",
[],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, i64>(1)?,
row.get::<_, i64>(2)?,
))
},
)
.unwrap();
assert_eq!(stored, ("target".into(), 7, 2));
for (id, namespace, timestamp, version) in [
("note-moved", "target", 11_i64, 2_i64),
("note-control", "unrelated", 13_i64, 1_i64),
] {
let stored = conn
.query_row(
"SELECT namespace,updated_at,version FROM notes WHERE id=?1",
[id],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, i64>(1)?,
row.get::<_, i64>(2)?,
))
},
)
.unwrap();
assert_eq!(stored, (namespace.into(), timestamp, version));
}
let tx = conn.transaction().unwrap();
let repeated = move_namespace(&tx, &request).unwrap();
assert_eq!(repeated.subjects.get("entity:concept"), Some(&0));
assert_eq!(repeated.subjects.get("note:observation"), Some(&0));
tx.commit().unwrap();
for (sql, expected) in [
("SELECT version FROM entities WHERE id='versioned'", 2_i64),
("SELECT version FROM notes WHERE id='note-moved'", 2_i64),
("SELECT version FROM notes WHERE id='note-control'", 1_i64),
] {
assert_eq!(
conn.query_row(sql, [], |row| row.get::<_, i64>(0)).unwrap(),
expected
);
}
}
#[cfg(all(test, feature = "vectors"))]
#[path = "namespace_move_partition_tests.rs"]
mod partition_tests;
#[cfg(test)]
#[path = "namespace_move_edge_tests.rs"]
mod edge_tests;
#[cfg(test)]
#[path = "namespace_move_pack_tables_tests.rs"]
mod pack_table_tests;
#[cfg(all(test, feature = "vectors"))]
#[path = "namespace_move_orphan_section_tests.rs"]
mod orphan_section_tests;