use std::collections::{BTreeMap, BTreeSet};
use rusqlite::Connection;
use crate::namespace_census::{self, NamespaceCensus, NamespaceConstraint};
#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub enum SubjectClass {
Note(String),
Entity(String),
Edge,
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())),
_ => 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::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>,
}
impl MoveRequest {
pub fn new(source: impl Into<String>, routes: Vec<MoveRoute>) -> Self {
Self {
source: source.into(),
routes,
}
}
fn route_for(&self, class: &SubjectClass) -> Option<&MoveRoute> {
self.routes.iter().find(|route| &route.class == class)
}
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,
}
#[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>,
},
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, 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::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"];
#[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 },
}
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,
"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",
},
_ 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>,
edges: 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 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)?,
edges: count_in_namespace(conn, "graph_edges", source)?,
atoms: count_in_namespace(conn, "knowledge_atoms", 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 !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)?;
}
unrouted(SubjectClass::Edge, inventory.edges)?;
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 });
}
Ok(())
}
fn collisions_for(
conn: &Connection,
constraint: &NamespaceConstraint,
source: &str,
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, 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 sql = format!(
"SELECT {select} FROM {table} AS source \
JOIN {table} AS target ON target.namespace = ?2 AND {join} \
WHERE source.namespace = ?1"
);
let mut stmt = conn.prepare(&sql)?;
let column_count = others.len();
let rows = stmt.query_map([source, target], 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_vectors(
conn: &Connection,
table: &str,
source: &str,
target: &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 = conn.execute(
&format!(
"CREATE TEMP TABLE namespace_move_vectors AS \
SELECT {columns} FROM {quoted} WHERE 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, ?1, kind, field, embedding_model, embedding \
FROM temp.namespace_move_vectors"
),
[target],
)? as u64;
debug_assert_eq!(
inserted, staged_rows as u64,
"every staged vector is re-inserted or the move is losing embeddings"
);
let appended = 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
+ conn.execute(
"INSERT INTO ann_write_log \
(namespace, embedding_model, kind, field, subject_id, op) \
SELECT ?1, embedding_model, kind, field, subject_id, 'upsert' \
FROM temp.namespace_move_vectors",
[target],
)? as u64;
conn.execute_batch("DROP TABLE temp.namespace_move_vectors")?;
Ok(VectorMove {
moved: inserted,
ann_appended: appended,
})
}
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.source, target)?);
}
}
if !collisions.is_empty() {
return Err(MoveError::Collisions { collisions });
}
let mut counts = MoveCounts::default();
let source = request.source.as_str();
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 => {
move_whole_table(conn, "graph_edges", source, target, &mut counts.rows)?
}
SubjectClass::Atom => {
let moved =
move_whole_table(conn, "knowledge_atoms", source, target, &mut counts.rows)?;
let sections = conn.execute(
"UPDATE knowledge_sections SET namespace = ?2 \
WHERE namespace = ?1 \
AND atom_id IN (SELECT id FROM knowledge_atoms WHERE namespace = ?2)",
rusqlite::params![source, target],
)? as u64;
*counts.rows.entry("knowledge_sections".into()).or_default() += sections;
moved
}
SubjectClass::Domain => {
move_whole_table(conn, "knowledge_domains", source, target, &mut counts.rows)?
}
};
counts.subjects.insert(route.class.render(), moved);
}
for table in &census.tables {
if !is_runtime_vector_table(table) {
continue;
}
if let Some(target) = request.single_target() {
let moved = move_vectors(conn, &table.name, source, target)?;
*counts.rows.entry(table.name.clone()).or_default() += moved.moved;
counts.ann_log_appended += moved.ann_appended;
} else {
let left = count_in_namespace(conn, &table.name, source)?;
if left > 0 {
*counts.left_behind.entry(table.name.clone()).or_default() += 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 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);
}
}
}
}
Ok(counts)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::migrations::run_migrations;
use rusqlite::Connection;
fn migrated() -> Connection {
let mut conn = Connection::open_in_memory().expect("open");
run_migrations(&mut conn).expect("migrate");
conn
}
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");
}
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));
}
#[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")]
#[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
);
}
}