Skip to main content

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}