khive_db/namespace_move.rs
1//! Moving records between namespaces (ADR-189).
2//!
3//! Namespace is attribution-only and an open string, so moving a record between
4//! namespaces is sound in principle. It is not a `UPDATE ... SET namespace`
5//! sweep, for reasons that are properties of this schema rather than matters of
6//! taste: six of the affected tables are fts5 virtual tables that accept no
7//! `UPDATE` of an indexed column, the vector tables are created at runtime and
8//! appear in no static list, and the ANN bookkeeping has ordering semantics an
9//! `UPDATE` violates silently.
10//!
11//! Three scoping facts decide the shape of everything below.
12//!
13//! **One backend.** A pack can be assigned its own backend, and then its records
14//! live in a different SQLite file; the live configuration on a development
15//! machine puts comm's notes and the knowledge atoms in two files beside the
16//! main one. SQLite has no transaction across unattached databases, so this
17//! operates on the connection it is given and a store with three backends is
18//! the same route map applied three times. That composes because a class routed
19//! with no rows here succeeds reporting zero. Atomicity does not compose, and
20//! this module does not pretend otherwise.
21//!
22//! **Re-runnable, so a partial application is a resume point.** A backend whose
23//! move did not run holds exactly the state that existed before anyone asked:
24//! records under the source namespace. That is not a new failure mode, it is the
25//! one the move was called to fix, still present for the unmoved subset. Because
26//! a routed class with no rows succeeds reporting zero, a second run over an
27//! already-moved backend is a no-op and a second run over the failed one is a
28//! first run. So what a multi-backend caller owes is not all-or-nothing across
29//! files, which SQLite cannot give it, but per-backend atomicity, a per-backend
30//! report so a partial outcome is known rather than silent, and the willingness
31//! to run the same request again.
32//!
33//! **No transaction of its own.** `WriterTaskHandle::send` hands its closure a
34//! connection already inside the `BEGIN IMMEDIATE` it opened and owns the commit
35//! or rollback, and a nested bare `BEGIN IMMEDIATE` is an error. So the entry
36//! point here is DML-only, and the enumerating `SELECT`s run inside the caller's
37//! transaction with the writes they feed — the same TOCTOU reason
38//! `Fts5TextSearch::rename_namespace` records.
39//!
40//! **Refuse rather than resolve.** A collision means two rows the caller wrote
41//! claim one identity in the target namespace. There is no conflict policy: any
42//! `ON CONFLICT` form picks a winner over a caller's data while satisfying a
43//! counts-in-equals-counts-out assertion, which is the worst available pairing —
44//! the destructive outcome and the reassuring receipt arrive together.
45
46use std::collections::{BTreeMap, BTreeSet};
47
48use rusqlite::Connection;
49
50use crate::namespace_census::{self, NamespaceCensus, NamespaceConstraint};
51
52/// A routable subject class: a record that exists in its own right.
53///
54/// Qualified by kind for the two classes that carry one. Nothing in the schema
55/// stops a note kind and an entity kind sharing a spelling, so an unqualified
56/// key would route both on a store that has them, and a refusal naming
57/// `note:observation` is one a caller can act on where `observation` is not.
58#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
59pub enum SubjectClass {
60 Note(String),
61 Entity(String),
62 Edge,
63 Atom,
64 Domain,
65}
66
67impl SubjectClass {
68 /// Parse a route key. The error names what was given, because a typo here
69 /// is the most likely caller mistake and the least likely to be obvious.
70 pub fn parse(key: &str) -> Result<Self, MoveError> {
71 match key {
72 "edge" => return Ok(Self::Edge),
73 "atom" => return Ok(Self::Atom),
74 "domain" => return Ok(Self::Domain),
75 _ => {}
76 }
77 match key.split_once(':') {
78 Some(("note", kind)) if !kind.is_empty() => Ok(Self::Note(kind.to_string())),
79 Some(("entity", kind)) if !kind.is_empty() => Ok(Self::Entity(kind.to_string())),
80 _ => Err(MoveError::UnknownSubjectClass {
81 key: key.to_string(),
82 }),
83 }
84 }
85
86 pub fn render(&self) -> String {
87 match self {
88 Self::Note(kind) => format!("note:{kind}"),
89 Self::Entity(kind) => format!("entity:{kind}"),
90 Self::Edge => "edge".to_string(),
91 Self::Atom => "atom".to_string(),
92 Self::Domain => "domain".to_string(),
93 }
94 }
95}
96
97/// One `(subject class, target namespace)` pair.
98#[derive(Clone, Debug, PartialEq, Eq)]
99pub struct MoveRoute {
100 pub class: SubjectClass,
101 pub target: String,
102}
103
104/// A whole move, validated before anything is written.
105#[derive(Clone, Debug, PartialEq, Eq)]
106pub struct MoveRequest {
107 pub source: String,
108 pub routes: Vec<MoveRoute>,
109}
110
111impl MoveRequest {
112 /// The route map as data, validated whole before the first write. A list can
113 /// be checked, logged, replayed and tested against a fixture; a callback
114 /// can do none of those.
115 pub fn new(source: impl Into<String>, routes: Vec<MoveRoute>) -> Self {
116 Self {
117 source: source.into(),
118 routes,
119 }
120 }
121
122 fn route_for(&self, class: &SubjectClass) -> Option<&MoveRoute> {
123 self.routes.iter().find(|route| &route.class == class)
124 }
125
126 /// Every route names the same target, and every class present in the source
127 /// is routed. Only then can a per-namespace aggregate with no subject be
128 /// carried anywhere.
129 fn single_target(&self) -> Option<&str> {
130 let mut targets = self.routes.iter().map(|r| r.target.as_str());
131 let first = targets.next()?;
132 targets.all(|t| t == first).then_some(first)
133 }
134}
135
136/// What a move did, per route and per table.
137///
138/// `left_behind` is not an error column. A per-namespace aggregate with no
139/// subject cannot be split across a partitioning move, so it stays, and the
140/// caller is told rather than left to discover it.
141#[derive(Clone, Debug, Default, PartialEq, Eq)]
142pub struct MoveCounts {
143 /// Subjects moved, keyed by the rendered route key. A routed class with no
144 /// rows appears here with zero — that is what makes the map verifiable by
145 /// its caller, and what makes the same map re-runnable against each backend
146 /// of a split store.
147 pub subjects: BTreeMap<String, u64>,
148 /// Rows written, keyed by table. Derived rows appear here and in no route.
149 pub rows: BTreeMap<String, u64>,
150 /// Rows left where they were, keyed by table, with the reason in the docs
151 /// above rather than in the data.
152 pub left_behind: BTreeMap<String, u64>,
153 /// Entries appended to `ann_write_log` under the target namespaces.
154 pub ann_log_appended: u64,
155}
156
157/// A note that cannot move, and its position in the stream that pins it.
158#[derive(Clone, Debug, PartialEq, Eq)]
159pub struct StreamMember {
160 pub note_id: String,
161 pub stream: String,
162 pub seq: i64,
163}
164
165/// A row that would claim an identity already taken in the target namespace.
166#[derive(Clone, Debug, PartialEq, Eq)]
167pub struct Collision {
168 pub table: String,
169 pub constraint: String,
170 pub target: String,
171 /// The key values that clash, rendered in the constraint's own column order.
172 ///
173 /// This identifies both rows and needs no second field: the row already in
174 /// the target holds this key there, and the row that would have moved holds
175 /// the same key in the source namespace.
176 pub key: String,
177}
178
179/// Why a move refused, or how it failed.
180#[derive(Debug)]
181pub enum MoveError {
182 /// A route key that names no subject class.
183 UnknownSubjectClass {
184 key: String,
185 },
186 /// The same class routed twice. Ambiguous rather than redundant: the two
187 /// targets may differ, and picking one would be a guess.
188 DuplicateRoute {
189 class: String,
190 },
191 /// A route whose target is the source. A no-op written as an instruction is
192 /// more likely a mistake than an intent.
193 TargetIsSource {
194 class: String,
195 },
196 /// A namespace-bearing table holding rows here that this code has no rule
197 /// for.
198 ///
199 /// The census finds tables; only a reader of the code can say what a move
200 /// does with one. So a table arriving in a migration after this was written
201 /// refuses the move by name, rather than letting the subjects around it move
202 /// and leaving the new table pointing at a namespace nothing else is in.
203 UnknownTable {
204 table: String,
205 rows: u64,
206 },
207 /// A subject class with rows in the source namespace and no route.
208 ///
209 /// Distinct from a routed class with no rows, which succeeds reporting zero.
210 /// A host that binds a pack writing nothing needs "routed, nothing there" to
211 /// be distinguishable from "you forgot this one".
212 UnroutedClass {
213 class: String,
214 rows: u64,
215 },
216 /// Notes the stream schema pins to their namespace.
217 ///
218 /// Four triggers make this absolute: an `UPDATE` of a member note naming
219 /// `namespace` aborts, the delete aborts, and both writes to the ledger row
220 /// abort. So the refusal comes from a read taken before any write, and names
221 /// the notes, rather than from a trigger's abort string from somewhere in
222 /// the middle of the transaction.
223 StreamMembers {
224 notes: Vec<StreamMember>,
225 },
226 /// Rows that would collide in a target namespace.
227 Collisions {
228 collisions: Vec<Collision>,
229 },
230 Sqlite(rusqlite::Error),
231}
232
233impl From<rusqlite::Error> for MoveError {
234 fn from(error: rusqlite::Error) -> Self {
235 Self::Sqlite(error)
236 }
237}
238
239impl std::fmt::Display for MoveError {
240 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
241 match self {
242 Self::UnknownSubjectClass { key } => write!(
243 f,
244 "no subject class named {key:?}; expected note:<kind>, entity:<kind>, edge, atom or domain"
245 ),
246 Self::DuplicateRoute { class } => {
247 write!(f, "{class} is routed more than once")
248 }
249 Self::TargetIsSource { class } => {
250 write!(f, "{class} is routed to the namespace it is already in")
251 }
252 Self::UnknownTable { table, rows } => write!(
253 f,
254 "{table} holds {rows} row(s) in the source namespace and this build has no rule \
255 for it; a table added by a later migration refuses a move rather than being \
256 left behind by one"
257 ),
258 Self::UnroutedClass { class, rows } => write!(
259 f,
260 "{class} has {rows} row(s) in the source namespace and no route; \
261 a class routed with zero rows succeeds reporting zero, an unrouted one refuses"
262 ),
263 Self::StreamMembers { notes } => {
264 write!(f, "{} note(s) belong to a stream and cannot change namespace: ", notes.len())?;
265 for (i, member) in notes.iter().enumerate() {
266 if i > 0 {
267 f.write_str(", ")?;
268 }
269 write!(f, "{} ({} seq {})", member.note_id, member.stream, member.seq)?;
270 }
271 Ok(())
272 }
273 Self::Collisions { collisions } => {
274 write!(f, "{} collision(s): ", collisions.len())?;
275 for (i, collision) in collisions.iter().enumerate() {
276 if i > 0 {
277 f.write_str(", ")?;
278 }
279 write!(
280 f,
281 "{}.{} in {} already holds {}",
282 collision.table, collision.constraint, collision.target, collision.key
283 )?;
284 }
285 Ok(())
286 }
287 Self::Sqlite(error) => write!(f, "{error}"),
288 }
289 }
290}
291
292impl std::error::Error for MoveError {}
293
294/// Tables holding rows keyed by namespace with no subject to be carried by.
295///
296/// `brain_profile_snapshots` is `(profile_id, namespace)` and `brain_event_log`
297/// likewise: one row per profile per namespace, describing an aggregate. A move
298/// routing five classes to five namespaces has no target to carry a single
299/// profile snapshot to, and splitting it would invent numbers. They move only
300/// when the move is total and single-target.
301const NAMESPACE_SCOPED_TABLES: &[&str] = &[
302 "brain_profile_snapshots",
303 "brain_event_log",
304 "proposals_open",
305];
306
307/// Tables keyed by a subject the caller routes, which therefore follow it.
308///
309/// Both are `(namespace, target_id, …)` shapes in the brain pack, and `target_id`
310/// names a note, an entity or an atom. That is why the follow query below is a
311/// union over the three subject tables rather than a per-route join: the column
312/// does not say which kind of subject it points at, and a wrong guess would
313/// leave learned state under a namespace its subject has left.
314const SUBJECT_KEYED_TABLES: &[&str] = &["brain_implicit_mass", "brain_serve_ledger"];
315
316/// What a move does with one namespace-bearing table.
317///
318/// This is the half of the design that cannot be derived, because it is a
319/// statement about meaning rather than about schema. The census answers which
320/// tables carry the column; only a reader of the code can say whether a row in
321/// one of them is a subject that moves, a projection that follows, an append-only
322/// record of something that already happened, or an aggregate that cannot be
323/// split. So the list is written down — and the guard is that a table the census
324/// finds and this function does not name is a REFUSAL, not a default.
325#[derive(Clone, Copy, Debug, PartialEq, Eq)]
326pub enum TableDisposition {
327 /// Rows the caller routes by subject class.
328 Subject,
329 /// A projection rebuilt from its subject, never routed on its own.
330 ///
331 /// Split two ways on purpose. `fts_knowledge` and `fts_sections` are
332 /// maintained by triggers that fire on `UPDATE OF ... namespace`
333 /// (`sql/026-knowledge-fts-repair.sql:49` and
334 /// `sql/002-narrow-fts-sections-update-trigger.sql:10`), so writing the base
335 /// row carries them. Both live declarations are column-scoped, and
336 /// `fts_sections_au` reached that shape by being narrowed: `sql/schema.sql`
337 /// declares it `AFTER UPDATE` unconditioned and V2 drops and recreates it
338 /// over a named column list, to stop reindex-only updates from paying an
339 /// fts5 delete-and-reinsert. So this disposition depends on `namespace`
340 /// staying in that list, which a later narrowing could shorten without
341 /// touching anything here. `fts_notes` and `fts_entities` have no triggers at all —
342 /// their contents are written from Rust — so the move writes them itself, and
343 /// because fts5 refuses an `UPDATE` of an indexed column that write is a
344 /// delete followed by an insert.
345 Derived { trigger_maintained: bool },
346 /// The rowid map beside an fts5 table, keyed `(namespace, subject_id)`.
347 DerivedRowidMap,
348 /// Appended to, never rewritten: new entries land under the target namespace
349 /// and the entries already there stay where they are, because they record
350 /// writes that happened under the old name.
351 Appended,
352 /// A record of what happened. Rewriting it would make the audit trail
353 /// describe a past that did not occur.
354 History,
355 /// Consumer bookkeeping whose ordering semantics an `UPDATE` violates.
356 ConsumerWatermark,
357 /// The schema itself refuses: `sql/029-note-streams.sql` installs four
358 /// triggers that abort any namespace change to a member note or its ledger.
359 RefusedBySchema,
360 /// Keyed by `(profile_id, namespace)` with no subject, so a partitioning
361 /// move has no target to carry it to. Moves only when the request is total
362 /// and single-target; otherwise reported as left behind.
363 NamespaceScopedAggregate,
364 /// Keyed by a subject the caller routes, so it follows that subject.
365 SubjectKeyed { subject_column: &'static str },
366}
367
368/// The disposition of a table, or `None` if this code has never seen it.
369///
370/// `None` is the whole point. A migration that adds a namespace-bearing table
371/// after this was written lands in no branch here, and a move that finds rows in
372/// it refuses by name rather than moving the subjects around it and leaving the
373/// new table pointing at a namespace nothing else is in.
374pub fn disposition(table: &namespace_census::NamespaceTable) -> Option<TableDisposition> {
375 use TableDisposition::*;
376 Some(match table.name.as_str() {
377 "notes" | "entities" | "graph_edges" | "knowledge_atoms" | "knowledge_domains" => Subject,
378
379 // A section is not a subject: it carries `atom_id REFERENCES
380 // knowledge_atoms(id)` and its uniqueness is `(atom_id, content_hash)`,
381 // with no namespace in it (`sql/schema.sql:166-182`). It follows its
382 // atom, and a caller cannot route it away from one.
383 "knowledge_sections" => SubjectKeyed {
384 subject_column: "atom_id",
385 },
386
387 // An open proposal is keyed by `proposal_id` and references no subject,
388 // so it belongs to the namespace rather than to anything inside it. In a
389 // move routing several classes to several targets there is no namespace
390 // for it to belong to afterwards, which is the same problem the two brain
391 // aggregates have and takes the same answer.
392 "proposals_open" => NamespaceScopedAggregate,
393
394 "fts_knowledge" | "fts_sections" => Derived {
395 trigger_maintained: true,
396 },
397 "fts_notes" | "fts_entities" => Derived {
398 trigger_maintained: false,
399 },
400 "fts_notes_rowids" | "fts_entities_rowids" => DerivedRowidMap,
401
402 "ann_write_log" => Appended,
403 "events" => History,
404 "ann_consumer_watermark" | "ann_consumer_pending" => ConsumerWatermark,
405 "note_streams" => RefusedBySchema,
406
407 "brain_profile_snapshots" | "brain_event_log" => NamespaceScopedAggregate,
408 "brain_implicit_mass" | "brain_serve_ledger" => SubjectKeyed {
409 subject_column: "target_id",
410 },
411
412 // Created at runtime, one per embedding model, and in no source file, so
413 // the live store is the only place they can be identified from.
414 //
415 // The NAME is not enough to identify one, and treating it as enough put a
416 // hole straight through the refusal above: a migration adding an ordinary
417 // namespace-bearing table called `vec_audit` would be classed here, handed
418 // to `move_vectors`, and die on `no such column: embedding` in the middle
419 // of the caller's transaction -- a bare SQLite error in place of the
420 // refusal by name that every other unnamed table gets. A vector table is a
421 // `CREATE VIRTUAL TABLE`, which an ordinary migration's table is not, so
422 // the census's own reading of that is the second half of the test.
423 _ if is_runtime_vector_table(table) => Derived {
424 trigger_maintained: false,
425 },
426
427 _ => return None,
428 })
429}
430
431/// A vector table created at runtime by an embedding model.
432///
433/// One predicate rather than two spellings of `starts_with("vec_")`: the
434/// disposition and the loop that moves them have to agree, or a table one of
435/// them admits reaches code the other never cleared.
436fn is_runtime_vector_table(table: &namespace_census::NamespaceTable) -> bool {
437 table.name.starts_with("vec_") && table.virtual_table
438}
439
440/// What the source namespace actually holds, read inside the caller's
441/// transaction with the writes it feeds.
442#[derive(Debug, Default)]
443struct SourceInventory {
444 /// `notes` and `entities` rows per `kind`, which is what a route key names.
445 note_kinds: BTreeMap<String, u64>,
446 entity_kinds: BTreeMap<String, u64>,
447 edges: u64,
448 atoms: u64,
449 domains: u64,
450 /// Namespace-bearing tables holding rows here that [`disposition`] has no
451 /// rule for. Non-empty means refuse.
452 unknown: Vec<(String, u64)>,
453}
454
455fn count_in_namespace(conn: &Connection, table: &str, namespace: &str) -> rusqlite::Result<u64> {
456 let sql = format!(
457 "SELECT COUNT(*) FROM {} WHERE namespace = ?1",
458 namespace_census::quote_ident(table)
459 );
460 conn.query_row(&sql, [namespace], |row| row.get::<_, i64>(0))
461 .map(|n| n as u64)
462}
463
464fn kinds_in_namespace(
465 conn: &Connection,
466 table: &str,
467 namespace: &str,
468) -> rusqlite::Result<BTreeMap<String, u64>> {
469 let sql = format!(
470 "SELECT kind, COUNT(*) FROM {} WHERE namespace = ?1 GROUP BY kind",
471 namespace_census::quote_ident(table)
472 );
473 let mut stmt = conn.prepare(&sql)?;
474 let rows = stmt.query_map([namespace], |row| {
475 Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)? as u64))
476 })?;
477 let mut out = BTreeMap::new();
478 for row in rows {
479 let (kind, count) = row?;
480 out.insert(kind, count);
481 }
482 Ok(out)
483}
484
485/// Soft-deleted rows are counted and moved with the rest, deliberately.
486///
487/// They are rows, they carry the namespace, and several of the partial unique
488/// indexes exclude them, so they cannot collide. Leaving them behind would
489/// strand a record's tombstone in a namespace its subject no longer occupies,
490/// which is the state the delete path reads when it decides whether a later
491/// write is a resurrection.
492fn read_source(
493 conn: &Connection,
494 census: &NamespaceCensus,
495 source: &str,
496) -> Result<SourceInventory, MoveError> {
497 let mut inventory = SourceInventory {
498 note_kinds: kinds_in_namespace(conn, "notes", source)?,
499 entity_kinds: kinds_in_namespace(conn, "entities", source)?,
500 edges: count_in_namespace(conn, "graph_edges", source)?,
501 atoms: count_in_namespace(conn, "knowledge_atoms", source)?,
502 domains: count_in_namespace(conn, "knowledge_domains", source)?,
503 unknown: Vec::new(),
504 };
505
506 for table in &census.tables {
507 if disposition(table).is_some() {
508 continue;
509 }
510 let count = count_in_namespace(conn, &table.name, source)?;
511 if count > 0 {
512 inventory.unknown.push((table.name.clone(), count));
513 }
514 }
515 Ok(inventory)
516}
517
518/// Notes the stream schema pins where they are.
519///
520/// Read before any write, so the refusal names the notes rather than arriving as
521/// a trigger's abort string from the middle of the transaction. Four triggers in
522/// `sql/029-note-streams.sql` make it absolute: a namespace-naming `UPDATE` of a
523/// member note aborts, the delete aborts, and both writes to the ledger abort.
524fn stream_members(conn: &Connection, source: &str) -> rusqlite::Result<Vec<StreamMember>> {
525 let mut stmt = conn.prepare(
526 "SELECT note_id, stream, seq FROM note_streams \
527 WHERE namespace = ?1 ORDER BY stream, seq",
528 )?;
529 let rows = stmt.query_map([source], |row| {
530 Ok(StreamMember {
531 note_id: row.get(0)?,
532 stream: row.get(1)?,
533 seq: row.get(2)?,
534 })
535 })?;
536 rows.collect()
537}
538
539/// Everything that can refuse before a single row is written.
540///
541/// Ordered by how much the caller can do about it: a malformed request first,
542/// then a request the store refuses, then a request the store would corrupt.
543pub fn validate(
544 conn: &Connection,
545 census: &NamespaceCensus,
546 request: &MoveRequest,
547) -> Result<(), MoveError> {
548 let mut seen = BTreeSet::new();
549 for route in &request.routes {
550 if !seen.insert(route.class.clone()) {
551 return Err(MoveError::DuplicateRoute {
552 class: route.class.render(),
553 });
554 }
555 if route.target == request.source {
556 return Err(MoveError::TargetIsSource {
557 class: route.class.render(),
558 });
559 }
560 }
561
562 let inventory = read_source(conn, census, &request.source)?;
563
564 if let Some((table, rows)) = inventory.unknown.first() {
565 return Err(MoveError::UnknownTable {
566 table: table.clone(),
567 rows: *rows,
568 });
569 }
570
571 let unrouted = |class: SubjectClass, rows: u64| -> Result<(), MoveError> {
572 if rows > 0 && request.route_for(&class).is_none() {
573 return Err(MoveError::UnroutedClass {
574 class: class.render(),
575 rows,
576 });
577 }
578 Ok(())
579 };
580 for (kind, rows) in &inventory.note_kinds {
581 unrouted(SubjectClass::Note(kind.clone()), *rows)?;
582 }
583 for (kind, rows) in &inventory.entity_kinds {
584 unrouted(SubjectClass::Entity(kind.clone()), *rows)?;
585 }
586 unrouted(SubjectClass::Edge, inventory.edges)?;
587 unrouted(SubjectClass::Atom, inventory.atoms)?;
588 unrouted(SubjectClass::Domain, inventory.domains)?;
589
590 let pinned = stream_members(conn, &request.source)?;
591 if !pinned.is_empty() {
592 return Err(MoveError::StreamMembers { notes: pinned });
593 }
594
595 Ok(())
596}
597
598/// Collisions the pre-flight can enumerate exactly, for one constraint.
599///
600/// Exactness is the admission criterion, not coverage. A constraint qualifies
601/// only when every key column has a name and the index is not partial, because
602/// those are the two cases where the query below means precisely what the
603/// constraint means:
604///
605/// - an expression key column is reported by `PRAGMA index_xinfo` with a null
606/// name, so there is nothing to join on;
607/// - a partial index carries a `WHERE` clause that `index_xinfo` does not
608/// report, so a join ignoring it reports clashes the constraint would not have
609/// raised.
610///
611/// The second rule covers two cases that are not alike, and saying so here keeps
612/// the comment from presenting one reason for both.
613/// `idx_comm_message_external_id` (`sql/005-unique-comm-external-id.sql:33`) is
614/// unreachable either way: its third key column is `json_extract(properties,
615/// '$.external_id')` and its predicate calls the same function twice, so neither
616/// the key nor the filter can be expressed without evaluating it on both sides.
617/// `idx_notes_namespace_kind_key` (`sql/028-notes-key.sql:4`) is not like that at
618/// all: three plain column names and `WHERE key IS NOT NULL AND deleted_at IS
619/// NULL`, which a source/target join CAN express exactly. It is excluded only
620/// because `index_xinfo` does not hand over the `WHERE`, and it is the collision
621/// a consolidation of two namespaces is most likely to hit, since two notes
622/// sharing a key under one kind is the ordinary case rather than the exotic one.
623/// Admitting it means deciding which predicate shapes a parse may accept, and a
624/// predicate read permissively would refuse moves SQLite allows, so it is left
625/// out until that rule exists rather than guessed at here.
626///
627/// What the excluded constraints get instead is the constraint itself: the move
628/// issues plain statements, they error, and the caller's transaction rolls back.
629/// So this function decides the QUALITY of a refusal, never whether one happens.
630/// A collision it cannot enumerate still aborts the move.
631///
632/// That refusal is thinner than it sounds, and it is worth writing down because
633/// it is the reason the exclusion above costs something. SQLite names the
634/// COLUMNS, not the index: a caller whose note-key move is refused receives
635/// `UNIQUE constraint failed: notes.namespace, notes.kind, notes.key` and gets
636/// no index name to look up and no rows. Measured.
637///
638/// The reverse also happens, and it does not show up here at all: a constraint
639/// this function DOES enumerate can be unreachable because of one the census
640/// never reported. `graph_edges` is `PRIMARY KEY (namespace, id)`, which names
641/// `namespace` and so arrives here as a live key — but
642/// `sql/014-graph-edges-id-unique.sql:25` puts a UNIQUE index on `id` alone,
643/// globally, so no two rows in the database can share an `id` and the clash this
644/// arm looks for cannot exist in any store that reached V13. That index names no
645/// namespace, so a census keyed on the column cannot see it, and nothing in the
646/// enumerated set says the arm is dead. `idx_graph_edges_unique_triple`
647/// (`namespace, source_id, target_id, relation`) is the constraint that actually
648/// refuses a graph edge move, and it is enumerated here.
649fn collisions_for(
650 conn: &Connection,
651 constraint: &NamespaceConstraint,
652 source: &str,
653 target: &str,
654) -> rusqlite::Result<Vec<Collision>> {
655 if constraint.partial || !constraint.columns_are_nameable() {
656 return Ok(Vec::new());
657 }
658 let names: Vec<&str> = constraint
659 .columns
660 .iter()
661 .filter_map(|c| c.as_deref())
662 .collect();
663 let others: Vec<&str> = names
664 .iter()
665 .copied()
666 .filter(|c| !c.eq_ignore_ascii_case("namespace"))
667 .collect();
668 if others.len() != names.len() - 1 {
669 // `namespace` is not in this constraint, so moving cannot collide on it.
670 return Ok(Vec::new());
671 }
672 if others.is_empty() {
673 // A uniqueness constraint on `namespace` alone: one row per namespace,
674 // and a second one arriving is a collision whatever its other columns.
675 // Handled by the same query with an empty key rendering.
676 return collisions_on_namespace_alone(conn, constraint, source, target);
677 }
678
679 let table = namespace_census::quote_ident(&constraint.table);
680 let join = others
681 .iter()
682 .map(|c| {
683 let q = namespace_census::quote_ident(c);
684 // `IS` rather than `=` so two NULLs in a nullable key column compare
685 // equal, which is what a UNIQUE index does NOT do. This over-reports
686 // in exactly one direction and the direction is the safe one.
687 format!("target.{q} IS source.{q}")
688 })
689 .collect::<Vec<_>>()
690 .join(" AND ");
691 let select = others
692 .iter()
693 .map(|c| format!("source.{}", namespace_census::quote_ident(c)))
694 .collect::<Vec<_>>()
695 .join(", ");
696 let sql = format!(
697 "SELECT {select} FROM {table} AS source \
698 JOIN {table} AS target ON target.namespace = ?2 AND {join} \
699 WHERE source.namespace = ?1"
700 );
701
702 let mut stmt = conn.prepare(&sql)?;
703 let column_count = others.len();
704 let rows = stmt.query_map([source, target], move |row| {
705 let mut parts = Vec::with_capacity(column_count);
706 for i in 0..column_count {
707 parts.push(match row.get_ref(i)? {
708 rusqlite::types::ValueRef::Null => "NULL".to_string(),
709 rusqlite::types::ValueRef::Integer(v) => v.to_string(),
710 rusqlite::types::ValueRef::Real(v) => v.to_string(),
711 rusqlite::types::ValueRef::Text(v) => String::from_utf8_lossy(v).into_owned(),
712 rusqlite::types::ValueRef::Blob(_) => "<blob>".to_string(),
713 });
714 }
715 Ok(parts.join(", "))
716 })?;
717
718 let mut found = Vec::new();
719 for key in rows {
720 found.push(Collision {
721 table: constraint.table.clone(),
722 constraint: constraint.index.clone(),
723 target: target.to_string(),
724 key: key?,
725 });
726 }
727 Ok(found)
728}
729
730fn collisions_on_namespace_alone(
731 conn: &Connection,
732 constraint: &NamespaceConstraint,
733 source: &str,
734 target: &str,
735) -> rusqlite::Result<Vec<Collision>> {
736 let table = namespace_census::quote_ident(&constraint.table);
737 let sql = format!(
738 "SELECT (SELECT COUNT(*) FROM {table} WHERE namespace = ?1) \
739 * (SELECT COUNT(*) FROM {table} WHERE namespace = ?2)"
740 );
741 let product: i64 = conn.query_row(&sql, [source, target], |row| row.get(0))?;
742 Ok(if product > 0 {
743 vec![Collision {
744 table: constraint.table.clone(),
745 constraint: constraint.index.clone(),
746 target: target.to_string(),
747 key: "(namespace alone)".to_string(),
748 }]
749 } else {
750 Vec::new()
751 })
752}
753
754/// Move one note or entity kind, and the rows derived from it.
755///
756/// The derived writes run AFTER the base update and select through it, so no id
757/// list is ever held in memory and a large namespace costs the same as a small
758/// one. They are still exact: the `WHERE namespace = :source` on the derived
759/// table excludes rows that were already in the target before this ran.
760/// The three tables a note or entity kind is spread across.
761///
762/// Grouped rather than passed as three strings because they are one fact: the
763/// map is keyed on the rowid the fts table holds, so naming them apart invites
764/// a call site that pairs a base with the wrong map.
765struct KindedTables {
766 base: &'static str,
767 fts: &'static str,
768 rowids: &'static str,
769}
770
771const NOTE_TABLES: KindedTables = KindedTables {
772 base: "notes",
773 fts: "fts_notes",
774 rowids: "fts_notes_rowids",
775};
776
777const ENTITY_TABLES: KindedTables = KindedTables {
778 base: "entities",
779 fts: "fts_entities",
780 rowids: "fts_entities_rowids",
781};
782
783fn move_kinded_subject(
784 conn: &Connection,
785 tables: &KindedTables,
786 source: &str,
787 target: &str,
788 kind: &str,
789 rows: &mut BTreeMap<String, u64>,
790) -> rusqlite::Result<u64> {
791 let KindedTables { base, fts, rowids } = *tables;
792 // Entity writers advance revisions explicitly; note revisions belong to
793 // their update trigger. Keep the entity target and both assignment lists
794 // literal so the two writer contracts can be inspected independently.
795 let statement = if base == "entities" {
796 "UPDATE entities SET namespace = ?2, version = version + 1 \
797 WHERE namespace = ?1 AND kind = ?3"
798 .to_owned()
799 } else {
800 format!(
801 "UPDATE {} SET namespace = ?2 WHERE namespace = ?1 AND kind = ?3",
802 namespace_census::quote_ident(base)
803 )
804 };
805 let moved = conn.execute(&statement, rusqlite::params![source, target, kind])? as u64;
806 *rows.entry(base.to_string()).or_default() += moved;
807
808 // An ordinary fts5 table accepts this and preserves the rowid, which is what
809 // the map below is keyed on. Measured; the arm lives in the tests.
810 let selector = format!(
811 "SELECT id FROM {} WHERE namespace = ?2 AND kind = ?3",
812 namespace_census::quote_ident(base)
813 );
814 for derived in [fts, rowids] {
815 let n = conn.execute(
816 &format!(
817 "UPDATE {} SET namespace = ?2 \
818 WHERE namespace = ?1 AND subject_id IN ({selector})",
819 namespace_census::quote_ident(derived)
820 ),
821 rusqlite::params![source, target, kind],
822 )? as u64;
823 *rows.entry(derived.to_string()).or_default() += n;
824 }
825 Ok(moved)
826}
827
828/// Move a whole table's rows out of the source namespace.
829fn move_whole_table(
830 conn: &Connection,
831 table: &str,
832 source: &str,
833 target: &str,
834 rows: &mut BTreeMap<String, u64>,
835) -> rusqlite::Result<u64> {
836 let moved = conn.execute(
837 &format!(
838 "UPDATE {} SET namespace = ?2 WHERE namespace = ?1",
839 namespace_census::quote_ident(table)
840 ),
841 rusqlite::params![source, target],
842 )? as u64;
843 *rows.entry(table.to_string()).or_default() += moved;
844 Ok(moved)
845}
846
847/// A vec0 row moves by delete and re-insert, carrying the stored embedding.
848///
849/// No `UPDATE` against vec0 exists anywhere in the tree, and re-embedding would
850/// be wasted work rather than a safer alternative: namespace is not an input to
851/// an embedding, so a re-index pass recomputes byte-identical vectors.
852///
853/// The order is forced by the key. `subject_id` is declared `PRIMARY KEY` and
854/// does not change, so an insert issued before the delete would collide with the
855/// row it is replacing. The moving rows therefore land in a temporary table
856/// first. That table is `TEMP`, so it is invisible to the store and dropped with
857/// the connection, and it is created and dropped inside the caller's
858/// transaction with everything else.
859fn move_vectors(
860 conn: &Connection,
861 table: &str,
862 source: &str,
863 target: &str,
864) -> rusqlite::Result<VectorMove> {
865 let quoted = namespace_census::quote_ident(table);
866 let columns = "subject_id, namespace, kind, field, embedding_model, embedding";
867
868 conn.execute_batch("DROP TABLE IF EXISTS temp.namespace_move_vectors")?;
869 let staged = conn.execute(
870 &format!(
871 "CREATE TEMP TABLE namespace_move_vectors AS \
872 SELECT {columns} FROM {quoted} WHERE namespace = ?1"
873 ),
874 [source],
875 );
876 // `CREATE TABLE ... AS SELECT` reports no row count, so the count comes from
877 // the staging table itself rather than from the statement.
878 staged?;
879 let staged_rows: i64 = conn.query_row(
880 "SELECT COUNT(*) FROM temp.namespace_move_vectors",
881 [],
882 |r| r.get(0),
883 )?;
884
885 conn.execute(
886 &format!("DELETE FROM {quoted} WHERE namespace = ?1"),
887 [source],
888 )?;
889 let inserted = conn.execute(
890 &format!(
891 "INSERT INTO {quoted} ({columns}) \
892 SELECT subject_id, ?1, kind, field, embedding_model, embedding \
893 FROM temp.namespace_move_vectors"
894 ),
895 [target],
896 )? as u64;
897 debug_assert_eq!(
898 inserted, staged_rows as u64,
899 "every staged vector is re-inserted or the move is losing embeddings"
900 );
901
902 // The staging table is still here because THIS is what the write log has to
903 // be built from. It holds the moved rows and nothing else; the live table now
904 // holds them beside whatever the target already had, and a log built by
905 // reading the target back cannot tell the two apart. Measured against the
906 // read-back form: a target already holding one other subject's vector
907 // produced a `delete` under the source for a subject the source never held,
908 // a second `upsert` for a vector that never moved, and an appended count of
909 // four where two vectors' worth of instructions were owed.
910 let appended = conn.execute(
911 "INSERT INTO ann_write_log (namespace, embedding_model, kind, field, subject_id, op) \
912 SELECT ?1, embedding_model, kind, field, subject_id, 'delete' \
913 FROM temp.namespace_move_vectors",
914 [source],
915 )? as u64
916 + conn.execute(
917 "INSERT INTO ann_write_log \
918 (namespace, embedding_model, kind, field, subject_id, op) \
919 SELECT ?1, embedding_model, kind, field, subject_id, 'upsert' \
920 FROM temp.namespace_move_vectors",
921 [target],
922 )? as u64;
923
924 conn.execute_batch("DROP TABLE temp.namespace_move_vectors")?;
925
926 Ok(VectorMove {
927 moved: inserted,
928 ann_appended: appended,
929 })
930}
931
932/// What one vector table's move did: the rows carried, and the instructions
933/// appended for the ANN consumers on BOTH sides.
934///
935/// The write log is appended to and never rewritten: its existing entries record
936/// writes that happened under the old name and are true. What a move adds is two
937/// entries per moved vector, not one.
938///
939/// One entry is not enough and the asymmetry is easy to miss. A consumer builds
940/// its index per `(namespace, embedding_model)` and advances a watermark over
941/// this log. An `upsert` under the target tells the target's consumer to take
942/// the vector. Nothing tells the SOURCE's consumer to drop it, so without the
943/// paired `delete` that index keeps answering searches with a subject that is no
944/// longer in its namespace — the same silent outcome as doing nothing to the
945/// vectors at all, moved one layer out.
946struct VectorMove {
947 moved: u64,
948 ann_appended: u64,
949}
950
951/// Move records out of one namespace, per the route map, inside the caller's
952/// transaction.
953///
954/// DML only. The caller opened `BEGIN IMMEDIATE` and owns the commit or the
955/// rollback; a nested one here is a SQLite error. Every refusal happens before
956/// the first write, except the ones only a constraint can raise, and those abort
957/// the caller's transaction whole.
958pub fn move_namespace(conn: &Connection, request: &MoveRequest) -> Result<MoveCounts, MoveError> {
959 let census = namespace_census::census(conn)?;
960 validate(conn, &census, request)?;
961
962 // Over DISTINCT targets, not over routes. `collisions_for` is a function of
963 // the constraint and the two namespaces and does not read the route's class,
964 // so two routes sharing a target ask the same question twice and the answers
965 // are byte-identical. A `Collision` carries no route, so the repeats are not
966 // a second fact about a second class, they are the same row printed again.
967 // Found by the fixture, which planted three note kinds bound for one target
968 // and read back the same collision three times.
969 let mut targets: BTreeSet<&str> = BTreeSet::new();
970 for route in &request.routes {
971 targets.insert(route.target.as_str());
972 }
973 let mut collisions = Vec::new();
974 for constraint in namespace_census::reachable_constraints(&census) {
975 for target in &targets {
976 collisions.extend(collisions_for(conn, constraint, &request.source, target)?);
977 }
978 }
979 if !collisions.is_empty() {
980 return Err(MoveError::Collisions { collisions });
981 }
982
983 let mut counts = MoveCounts::default();
984 let source = request.source.as_str();
985
986 for route in &request.routes {
987 let target = route.target.as_str();
988 let moved = match &route.class {
989 SubjectClass::Note(kind) => {
990 move_kinded_subject(conn, &NOTE_TABLES, source, target, kind, &mut counts.rows)?
991 }
992 SubjectClass::Entity(kind) => {
993 move_kinded_subject(conn, &ENTITY_TABLES, source, target, kind, &mut counts.rows)?
994 }
995 SubjectClass::Edge => {
996 move_whole_table(conn, "graph_edges", source, target, &mut counts.rows)?
997 }
998 SubjectClass::Atom => {
999 let moved =
1000 move_whole_table(conn, "knowledge_atoms", source, target, &mut counts.rows)?;
1001 // Sections follow their atom by `atom_id`, and `fts_knowledge`
1002 // and `fts_sections` follow both by trigger. Writing either
1003 // virtual table here would be a no-op that reports success.
1004 let sections = conn.execute(
1005 "UPDATE knowledge_sections SET namespace = ?2 \
1006 WHERE namespace = ?1 \
1007 AND atom_id IN (SELECT id FROM knowledge_atoms WHERE namespace = ?2)",
1008 rusqlite::params![source, target],
1009 )? as u64;
1010 *counts.rows.entry("knowledge_sections".into()).or_default() += sections;
1011 moved
1012 }
1013 SubjectClass::Domain => {
1014 move_whole_table(conn, "knowledge_domains", source, target, &mut counts.rows)?
1015 }
1016 };
1017 counts.subjects.insert(route.class.render(), moved);
1018 }
1019
1020 // Vectors are enumerated from the live store, never from a constant list: a
1021 // store using a model this build was never compiled against still has its
1022 // `vec_*` table found here.
1023 for table in &census.tables {
1024 if !is_runtime_vector_table(table) {
1025 continue;
1026 }
1027 if let Some(target) = request.single_target() {
1028 let moved = move_vectors(conn, &table.name, source, target)?;
1029 *counts.rows.entry(table.name.clone()).or_default() += moved.moved;
1030 counts.ann_log_appended += moved.ann_appended;
1031 } else {
1032 // A partitioning move cannot send one vector table to several
1033 // targets in one statement, and splitting it needs the subject each
1034 // row belongs to, which is the next thing this grows.
1035 let left = count_in_namespace(conn, &table.name, source)?;
1036 if left > 0 {
1037 *counts.left_behind.entry(table.name.clone()).or_default() += left;
1038 }
1039 }
1040 }
1041
1042 // Learned state follows the subject it is about. Run per distinct target, so
1043 // a partitioning move sends each row after the subject it names. The set is
1044 // the one the pre-flight above already built, for the same reason.
1045 for table in SUBJECT_KEYED_TABLES {
1046 for target in &targets {
1047 let moved = conn.execute(
1048 &format!(
1049 "UPDATE {} SET namespace = ?2 WHERE namespace = ?1 AND target_id IN (\
1050 SELECT id FROM notes WHERE namespace = ?2 \
1051 UNION ALL SELECT id FROM entities WHERE namespace = ?2 \
1052 UNION ALL SELECT id FROM knowledge_atoms WHERE namespace = ?2)",
1053 namespace_census::quote_ident(table)
1054 ),
1055 rusqlite::params![source, target],
1056 )? as u64;
1057 *counts.rows.entry((*table).to_string()).or_default() += moved;
1058 }
1059 let left = count_in_namespace(conn, table, source)?;
1060 if left > 0 {
1061 counts.left_behind.insert((*table).to_string(), left);
1062 }
1063 }
1064
1065 // Per-namespace aggregates with no subject. A partitioning move has no
1066 // target to carry them to, so they stay and the caller is told, rather than
1067 // being left to find out.
1068 for table in NAMESPACE_SCOPED_TABLES {
1069 match request.single_target() {
1070 Some(target) => {
1071 move_whole_table(conn, table, source, target, &mut counts.rows)?;
1072 }
1073 None => {
1074 let left = count_in_namespace(conn, table, source)?;
1075 if left > 0 {
1076 counts.left_behind.insert((*table).to_string(), left);
1077 }
1078 }
1079 }
1080 }
1081
1082 Ok(counts)
1083}
1084
1085#[cfg(test)]
1086mod tests {
1087 use super::*;
1088 use crate::migrations::run_migrations;
1089 use rusqlite::Connection;
1090
1091 fn migrated() -> Connection {
1092 let mut conn = Connection::open_in_memory().expect("open");
1093 run_migrations(&mut conn).expect("migrate");
1094 conn
1095 }
1096
1097 /// Seeds through raw SQL, which is enough for every arm below and is NOT
1098 /// enough for the ones that are deliberately absent.
1099 ///
1100 /// `fts_notes` and `fts_entities` have no triggers: their contents are
1101 /// written from Rust, so a SQL seed leaves them empty and an arm asserting
1102 /// the move carried them would pass against a store where there was nothing
1103 /// to carry. Those arms need a fixture built through the store's own
1104 /// writers. `fts_knowledge` and `fts_sections` ARE trigger-maintained, so
1105 /// they are reachable from here and are exercised.
1106 fn seed_note(conn: &Connection, id: &str, namespace: &str, kind: &str) {
1107 conn.execute(
1108 "INSERT INTO notes (id, namespace, kind, name, content, created_at, updated_at) \
1109 VALUES (?1, ?2, ?3, 'a name', 'some content', 1, 1)",
1110 rusqlite::params![id, namespace, kind],
1111 )
1112 .expect("seed note");
1113 }
1114
1115 fn route(key: &str, target: &str) -> MoveRoute {
1116 MoveRoute {
1117 class: SubjectClass::parse(key).expect("route key"),
1118 target: target.to_string(),
1119 }
1120 }
1121
1122 /// The property the whole multi-backend story rests on. A store split across
1123 /// several SQLite files runs the same request once per backend, and that
1124 /// composes only because a backend holding none of a routed class is a
1125 /// SUCCESS reporting zero rather than a refusal.
1126 #[test]
1127 fn a_routed_class_with_no_rows_succeeds_reporting_zero() {
1128 let conn = migrated();
1129 let request = MoveRequest::new(
1130 "empty-source",
1131 vec![route("note:observation", "target"), route("atom", "target")],
1132 );
1133 let counts = move_namespace(&conn, &request).expect("a backend with nothing routed here");
1134 assert_eq!(counts.subjects.get("note:observation"), Some(&0));
1135 assert_eq!(counts.subjects.get("atom"), Some(&0));
1136 }
1137
1138 /// And the case it must stay distinguishable from. A host binding a pack
1139 /// that writes nothing needs "routed, nothing there" to read differently
1140 /// from "you forgot this one".
1141 #[test]
1142 fn a_class_with_rows_and_no_route_refuses_and_says_how_many() {
1143 let conn = migrated();
1144 seed_note(&conn, "n1", "source", "observation");
1145 seed_note(&conn, "n2", "source", "decision");
1146
1147 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1148 let error = move_namespace(&conn, &request).expect_err("decision notes are unrouted");
1149 match error {
1150 MoveError::UnroutedClass { class, rows } => {
1151 assert_eq!(class, "note:decision");
1152 assert_eq!(rows, 1);
1153 }
1154 other => panic!("expected an unrouted class, got {other}"),
1155 }
1156
1157 let still_here: i64 = conn
1158 .query_row(
1159 "SELECT COUNT(*) FROM notes WHERE namespace = 'source'",
1160 [],
1161 |r| r.get(0),
1162 )
1163 .expect("count");
1164 assert_eq!(still_here, 2, "a refusal writes nothing");
1165 }
1166
1167 /// The drift guard, and the reason `disposition` returning `None` is a
1168 /// refusal rather than a default. This arm mints a namespace-bearing table
1169 /// the way a migration would, so it passes only if the refusal is derived
1170 /// from the census rather than from a list somebody remembered to edit.
1171 #[test]
1172 fn a_namespace_table_this_build_has_no_rule_for_refuses_the_move() {
1173 let conn = migrated();
1174 seed_note(&conn, "n1", "source", "observation");
1175 conn.execute_batch(
1176 "CREATE TABLE later_migration_added_this (\
1177 id TEXT PRIMARY KEY, namespace TEXT NOT NULL);\
1178 INSERT INTO later_migration_added_this VALUES ('x', 'source');",
1179 )
1180 .expect("a migration lands");
1181
1182 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1183 let error = move_namespace(&conn, &request).expect_err("an unknown table refuses");
1184 match error {
1185 MoveError::UnknownTable { table, rows } => {
1186 assert_eq!(table, "later_migration_added_this");
1187 assert_eq!(rows, 1);
1188 }
1189 other => panic!("expected an unknown table, got {other}"),
1190 }
1191 }
1192
1193 /// The same refusal, for a table whose NAME says it is a vector table.
1194 ///
1195 /// The vector tables are created at runtime by embedding models and appear in
1196 /// no source file, so they are recognised from the live store. Recognising
1197 /// them by name alone puts a hole through the refusal above: a migration
1198 /// adding an ordinary table called `vec_audit` would be classed as a vector
1199 /// table, handed to the vector move, and die on `no such column: embedding`
1200 /// somewhere inside the caller's transaction. That is the one outcome the
1201 /// refusal exists to prevent -- a bare SQLite error in place of a named
1202 /// refusal. A real vector table is a `CREATE VIRTUAL TABLE` and this one is
1203 /// not, which is what separates them here.
1204 #[test]
1205 fn a_table_named_like_a_vector_table_but_not_one_refuses_by_name() {
1206 let conn = migrated();
1207 seed_note(&conn, "n1", "source", "observation");
1208 conn.execute_batch(
1209 "CREATE TABLE vec_audit (\
1210 id TEXT PRIMARY KEY, namespace TEXT NOT NULL);\
1211 INSERT INTO vec_audit VALUES ('x', 'source');",
1212 )
1213 .expect("a migration lands a table whose name starts with the prefix");
1214
1215 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1216 let error = move_namespace(&conn, &request).expect_err("the prefix is not enough");
1217 match error {
1218 MoveError::UnknownTable { table, rows } => {
1219 assert_eq!(table, "vec_audit");
1220 assert_eq!(rows, 1);
1221 }
1222 other => panic!("expected an unknown table, got {other}"),
1223 }
1224 }
1225
1226 /// Control for the arm above: the same table with no rows in the source
1227 /// namespace does NOT refuse, because a move that never touches it has
1228 /// nothing to be wrong about.
1229 #[test]
1230 fn an_unknown_table_holding_nothing_here_does_not_refuse() {
1231 let conn = migrated();
1232 seed_note(&conn, "n1", "source", "observation");
1233 conn.execute_batch(
1234 "CREATE TABLE later_migration_added_this (\
1235 id TEXT PRIMARY KEY, namespace TEXT NOT NULL);\
1236 INSERT INTO later_migration_added_this VALUES ('x', 'somewhere-else');",
1237 )
1238 .expect("a migration lands");
1239
1240 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1241 let counts = move_namespace(&conn, &request).expect("nothing of ours is in that table");
1242 assert_eq!(counts.subjects.get("note:observation"), Some(&1));
1243 }
1244
1245 #[test]
1246 fn a_route_to_the_namespace_it_is_already_in_refuses() {
1247 let conn = migrated();
1248 let request = MoveRequest::new("source", vec![route("atom", "source")]);
1249 match move_namespace(&conn, &request).expect_err("a no-op written as an instruction") {
1250 MoveError::TargetIsSource { class } => assert_eq!(class, "atom"),
1251 other => panic!("expected target-is-source, got {other}"),
1252 }
1253 }
1254
1255 #[test]
1256 fn the_same_class_routed_twice_refuses_rather_than_picking_one() {
1257 let conn = migrated();
1258 let request = MoveRequest::new(
1259 "source",
1260 vec![route("atom", "one"), route("atom", "another")],
1261 );
1262 match move_namespace(&conn, &request).expect_err("two targets, no rule to choose") {
1263 MoveError::DuplicateRoute { class } => assert_eq!(class, "atom"),
1264 other => panic!("expected a duplicate route, got {other}"),
1265 }
1266 }
1267
1268 #[test]
1269 fn an_unknown_route_key_names_what_it_was_given() {
1270 match SubjectClass::parse("notes:observation").expect_err("plural is a typo") {
1271 MoveError::UnknownSubjectClass { key } => assert_eq!(key, "notes:observation"),
1272 other => panic!("expected an unknown class, got {other}"),
1273 }
1274 assert_eq!(
1275 SubjectClass::parse("note:observation").expect("singular"),
1276 SubjectClass::Note("observation".into())
1277 );
1278 }
1279
1280 /// A partitioning move has no target to carry a per-namespace aggregate to,
1281 /// so it stays AND is reported. Silence here would be the failure: the seat
1282 /// that eventually binds brain would discover the rows instead of reading a
1283 /// line about them.
1284 #[test]
1285 fn a_partitioning_move_reports_the_aggregates_it_leaves_behind() {
1286 let conn = migrated();
1287 seed_note(&conn, "n1", "source", "observation");
1288 seed_note(&conn, "n2", "source", "decision");
1289 conn.execute(
1290 "INSERT INTO brain_profile_snapshots (profile_id, namespace, snapshot_json, updated_at) \
1291 VALUES ('p', 'source', '{}', 1)",
1292 [],
1293 )
1294 .expect("seed a snapshot");
1295
1296 let request = MoveRequest::new(
1297 "source",
1298 vec![
1299 route("note:observation", "one"),
1300 route("note:decision", "another"),
1301 ],
1302 );
1303 let counts = move_namespace(&conn, &request).expect("a partitioning move");
1304 assert_eq!(counts.left_behind.get("brain_profile_snapshots"), Some(&1));
1305
1306 let stayed: i64 = conn
1307 .query_row(
1308 "SELECT COUNT(*) FROM brain_profile_snapshots WHERE namespace = 'source'",
1309 [],
1310 |r| r.get(0),
1311 )
1312 .expect("count");
1313 assert_eq!(stayed, 1);
1314 }
1315
1316 /// The same aggregate DOES move when every route names one target, because
1317 /// then there is somewhere for it to belong.
1318 #[test]
1319 fn a_total_single_target_move_carries_the_aggregates() {
1320 let conn = migrated();
1321 seed_note(&conn, "n1", "source", "observation");
1322 conn.execute(
1323 "INSERT INTO brain_profile_snapshots (profile_id, namespace, snapshot_json, updated_at) \
1324 VALUES ('p', 'source', '{}', 1)",
1325 [],
1326 )
1327 .expect("seed a snapshot");
1328
1329 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1330 let counts = move_namespace(&conn, &request).expect("a total move");
1331 assert!(counts.left_behind.is_empty(), "{:?}", counts.left_behind);
1332
1333 let moved: i64 = conn
1334 .query_row(
1335 "SELECT COUNT(*) FROM brain_profile_snapshots WHERE namespace = 'target'",
1336 [],
1337 |r| r.get(0),
1338 )
1339 .expect("count");
1340 assert_eq!(moved, 1);
1341 }
1342 /// The half of a vector move that is invisible from the side it leaves.
1343 ///
1344 /// An ANN consumer builds its index per `(namespace, embedding_model)` by
1345 /// advancing a watermark over `ann_write_log`. Moving the row in the `vec_*`
1346 /// table and appending only the target's `upsert` leaves the source's index
1347 /// intact and still answering searches with a subject that is no longer in
1348 /// its namespace, which is the same observable as never having touched the
1349 /// vectors at all. This arm fails if nothing tells the source side to drop
1350 /// what left.
1351 ///
1352 /// The consumer itself lives above this crate, so what is asserted here is
1353 /// the instruction it reads, not the index it builds from it.
1354 #[cfg(feature = "vectors")]
1355 #[test]
1356 fn a_vector_move_tells_the_source_side_to_drop_what_left() {
1357 // Registration is an auto-extension, so it only reaches connections
1358 // opened after it. This has to come before `migrated`.
1359 crate::extension::ensure_extensions_loaded();
1360 let conn = migrated();
1361 seed_note(&conn, "n1", "source", "observation");
1362 conn.execute_batch(
1363 "CREATE VIRTUAL TABLE vec_test_model USING vec0(\
1364 subject_id TEXT PRIMARY KEY, \
1365 namespace TEXT NOT NULL, \
1366 kind TEXT NOT NULL, \
1367 field TEXT NOT NULL, \
1368 embedding_model TEXT NOT NULL, \
1369 embedding float[4] distance_metric=cosine\
1370 )",
1371 )
1372 .expect("the vector table an embedding model creates at runtime");
1373 conn.execute(
1374 "INSERT INTO vec_test_model \
1375 (subject_id, namespace, kind, field, embedding_model, embedding) \
1376 VALUES ('n1', 'source', 'observation', 'content', 'test-model', \
1377 '[0.1, 0.2, 0.3, 0.4]')",
1378 [],
1379 )
1380 .expect("seed a vector");
1381
1382 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1383 let counts = move_namespace(&conn, &request).expect("a total move");
1384
1385 assert_eq!(
1386 counts.rows.get("vec_test_model"),
1387 Some(&1),
1388 "the vector itself moved"
1389 );
1390 let left_in_source: i64 = conn
1391 .query_row(
1392 "SELECT COUNT(*) FROM vec_test_model WHERE namespace = 'source'",
1393 [],
1394 |r| r.get(0),
1395 )
1396 .expect("count");
1397 assert_eq!(left_in_source, 0);
1398
1399 // Two entries per moved vector, and the one that matters here is the
1400 // first: without it the source's index is never told anything.
1401 assert_eq!(counts.ann_log_appended, 2);
1402 let dropped_from_source: i64 = conn
1403 .query_row(
1404 "SELECT COUNT(*) FROM ann_write_log \
1405 WHERE namespace = 'source' AND op = 'delete' \
1406 AND subject_id = 'n1' AND embedding_model = 'test-model'",
1407 [],
1408 |r| r.get(0),
1409 )
1410 .expect("count");
1411 assert_eq!(
1412 dropped_from_source, 1,
1413 "the source consumer is never told to drop the vector, so its index \
1414 keeps answering with a subject that has left the namespace"
1415 );
1416 let taken_by_target: i64 = conn
1417 .query_row(
1418 "SELECT COUNT(*) FROM ann_write_log \
1419 WHERE namespace = 'target' AND op = 'upsert' \
1420 AND subject_id = 'n1' AND embedding_model = 'test-model'",
1421 [],
1422 |r| r.get(0),
1423 )
1424 .expect("count");
1425 assert_eq!(taken_by_target, 1);
1426 }
1427
1428 /// The instructions are about the vectors that MOVED, and a target is
1429 /// allowed to have vectors of its own already.
1430 ///
1431 /// Built by reading the live table back after the insert, the source side is
1432 /// "everything now under the target", which is the moved rows plus whatever
1433 /// was already there. That tells the source's consumer to drop a subject the
1434 /// source never held, re-upserts a vector that did not move, and reports an
1435 /// appended count of four where two instructions were owed. The staged rows
1436 /// are the only reading of "what moved" that survives the insert, which is
1437 /// why the log is built before they are dropped.
1438 #[cfg(feature = "vectors")]
1439 #[test]
1440 fn a_vector_the_target_already_held_is_not_in_the_instructions() {
1441 crate::extension::ensure_extensions_loaded();
1442 let conn = migrated();
1443 seed_note(&conn, "n1", "source", "observation");
1444 conn.execute_batch(
1445 "CREATE VIRTUAL TABLE vec_test_model USING vec0(\
1446 subject_id TEXT PRIMARY KEY, \
1447 namespace TEXT NOT NULL, \
1448 kind TEXT NOT NULL, \
1449 field TEXT NOT NULL, \
1450 embedding_model TEXT NOT NULL, \
1451 embedding float[4] distance_metric=cosine\
1452 )",
1453 )
1454 .expect("the vector table an embedding model creates at runtime");
1455 conn.execute(
1456 "INSERT INTO vec_test_model \
1457 (subject_id, namespace, kind, field, embedding_model, embedding) \
1458 VALUES ('n1', 'source', 'observation', 'content', 'test-model', \
1459 '[0.1, 0.2, 0.3, 0.4]')",
1460 [],
1461 )
1462 .expect("the vector that moves");
1463 conn.execute(
1464 "INSERT INTO vec_test_model \
1465 (subject_id, namespace, kind, field, embedding_model, embedding) \
1466 VALUES ('already-there', 'target', 'observation', 'content', 'test-model', \
1467 '[0.5, 0.6, 0.7, 0.8]')",
1468 [],
1469 )
1470 .expect("a vector the target already holds");
1471
1472 let request = MoveRequest::new("source", vec![route("note:observation", "target")]);
1473 let counts = move_namespace(&conn, &request).expect("a total move");
1474
1475 assert_eq!(
1476 counts.ann_log_appended, 2,
1477 "two instructions are owed for the one vector that moved"
1478 );
1479 let about_the_resident: i64 = conn
1480 .query_row(
1481 "SELECT COUNT(*) FROM ann_write_log WHERE subject_id = 'already-there'",
1482 [],
1483 |r| r.get(0),
1484 )
1485 .expect("count");
1486 assert_eq!(
1487 about_the_resident, 0,
1488 "a vector that did not move is told nothing, and is certainly not \
1489 dropped from a namespace it was never in"
1490 );
1491 // The control, so the arm cannot pass on a move that logged nothing at
1492 // all: the vector that did move still has both of its instructions.
1493 let about_the_mover: i64 = conn
1494 .query_row(
1495 "SELECT COUNT(*) FROM ann_write_log WHERE subject_id = 'n1'",
1496 [],
1497 |r| r.get(0),
1498 )
1499 .expect("count");
1500 assert_eq!(about_the_mover, 2);
1501 }
1502 /// Two routes bound for one target report a shared collision once, not twice.
1503 ///
1504 /// The pre-flight reads a constraint and two namespaces and never reads the
1505 /// route's class, so iterating routes asked the same question once per route
1506 /// and pushed byte-identical rows. A `Collision` carries no route, so a
1507 /// repeat says nothing a reader can act on: it inflates the list in
1508 /// proportion to how finely the caller partitioned its request, which is the
1509 /// one thing the refusal should be independent of.
1510 ///
1511 /// The arm fails on the unfixed code by reporting the same collision three
1512 /// times, once per route. Restoring the route-keyed loop is the control.
1513 #[test]
1514 fn two_routes_to_one_target_report_a_shared_collision_once() {
1515 let conn = migrated();
1516 // The clash is on the atom slug, which is a plain two-column unique
1517 // index and the one collision a SQL seed can plant honestly. Every class
1518 // present in the source must be routed or `validate` refuses first, so
1519 // the source holds exactly what these three routes name.
1520 seed_note(&conn, "n1", "source", "observation");
1521 seed_note(&conn, "n2", "source", "insight");
1522 for (id, namespace) in [("a1", "source"), ("a2", "target")] {
1523 conn.execute(
1524 "INSERT INTO knowledge_atoms \
1525 (id, namespace, slug, name, created_at, updated_at) \
1526 VALUES (?1, ?2, 'shared-slug', 'an atom', 1, 1)",
1527 rusqlite::params![id, namespace],
1528 )
1529 .expect("seed an atom on each side of the move");
1530 }
1531
1532 let request = MoveRequest::new(
1533 "source",
1534 vec![
1535 route("note:observation", "target"),
1536 route("note:insight", "target"),
1537 route("atom", "target"),
1538 ],
1539 );
1540 let error = move_namespace(&conn, &request).expect_err("the pre-flight refuses");
1541 let MoveError::Collisions { collisions } = error else {
1542 panic!("expected a named collision list, got {error:?}");
1543 };
1544
1545 assert_eq!(
1546 collisions.len(),
1547 1,
1548 "three routes share one target, so the one blocking row is reported \
1549 once: {collisions:?}"
1550 );
1551 assert_eq!(collisions[0].table, "knowledge_atoms");
1552 assert_eq!(collisions[0].constraint, "idx_knowledge_atoms_ns_slug");
1553 assert_eq!(collisions[0].key, "shared-slug");
1554 }
1555}
1556
1557#[cfg(test)]
1558#[test]
1559fn issue2673_namespace_move_advances_entity_version_without_changing_timestamp() {
1560 let mut conn = Connection::open_in_memory().unwrap();
1561 crate::migrations::run_migrations(&mut conn).unwrap();
1562 conn.execute("INSERT INTO entities(id,namespace,kind,name,created_at,updated_at) VALUES('versioned','source','concept','moved',7,7)", []).unwrap();
1563 conn.execute(
1564 "INSERT INTO notes(id,namespace,kind,content,created_at,updated_at) \
1565 VALUES('note-moved','source','observation','moved',11,11), \
1566 ('note-control','unrelated','observation','unchanged',13,13)",
1567 [],
1568 )
1569 .unwrap();
1570 let request = MoveRequest::new(
1571 "source",
1572 vec![
1573 MoveRoute {
1574 class: SubjectClass::Entity("concept".into()),
1575 target: "target".into(),
1576 },
1577 MoveRoute {
1578 class: SubjectClass::Note("observation".into()),
1579 target: "target".into(),
1580 },
1581 ],
1582 );
1583 let tx = conn.transaction().unwrap();
1584 let moved = move_namespace(&tx, &request).unwrap();
1585 assert_eq!(moved.subjects.get("entity:concept"), Some(&1));
1586 assert_eq!(moved.subjects.get("note:observation"), Some(&1));
1587 tx.commit().unwrap();
1588 let stored = conn
1589 .query_row(
1590 "SELECT namespace,updated_at,version FROM entities WHERE id='versioned'",
1591 [],
1592 |row| {
1593 Ok((
1594 row.get::<_, String>(0)?,
1595 row.get::<_, i64>(1)?,
1596 row.get::<_, i64>(2)?,
1597 ))
1598 },
1599 )
1600 .unwrap();
1601 assert_eq!(stored, ("target".into(), 7, 2));
1602 for (id, namespace, timestamp, version) in [
1603 ("note-moved", "target", 11_i64, 2_i64),
1604 ("note-control", "unrelated", 13_i64, 1_i64),
1605 ] {
1606 let stored = conn
1607 .query_row(
1608 "SELECT namespace,updated_at,version FROM notes WHERE id=?1",
1609 [id],
1610 |row| {
1611 Ok((
1612 row.get::<_, String>(0)?,
1613 row.get::<_, i64>(1)?,
1614 row.get::<_, i64>(2)?,
1615 ))
1616 },
1617 )
1618 .unwrap();
1619 assert_eq!(stored, (namespace.into(), timestamp, version));
1620 }
1621
1622 let tx = conn.transaction().unwrap();
1623 let repeated = move_namespace(&tx, &request).unwrap();
1624 assert_eq!(repeated.subjects.get("entity:concept"), Some(&0));
1625 assert_eq!(repeated.subjects.get("note:observation"), Some(&0));
1626 tx.commit().unwrap();
1627 for (sql, expected) in [
1628 ("SELECT version FROM entities WHERE id='versioned'", 2_i64),
1629 ("SELECT version FROM notes WHERE id='note-moved'", 2_i64),
1630 ("SELECT version FROM notes WHERE id='note-control'", 1_i64),
1631 ] {
1632 assert_eq!(
1633 conn.query_row(sql, [], |row| row.get::<_, i64>(0)).unwrap(),
1634 expected
1635 );
1636 }
1637}