1use rusqlite::{Connection, OptionalExtension, Transaction};
56use std::collections::BTreeMap;
57use std::sync::{Mutex, OnceLock};
58
59use crate::error::SqliteStoreError;
60
61const CREATE_LEDGER_SQL: &str = "CREATE TABLE IF NOT EXISTS main.meerkat_schema (
62 domain TEXT PRIMARY KEY,
63 version INTEGER NOT NULL
64)";
65
66const CUSTODY_SAVEPOINT_SQL: &str = "SAVEPOINT meerkat_migration_custody";
71const CUSTODY_RELEASE_SQL: &str = "RELEASE SAVEPOINT meerkat_migration_custody";
72
73#[derive(Debug)]
75pub struct Migration {
76 pub version: i64,
79 pub name: &'static str,
81 pub apply: fn(&Transaction<'_>) -> Result<(), rusqlite::Error>,
86}
87
88#[derive(Debug)]
90pub struct SchemaPredecessor {
91 pub version: i64,
92 pub verify: fn(&Connection) -> Result<(), String>,
93}
94
95#[derive(Debug, Clone, Copy, PartialEq, Eq)]
97pub enum SchemaObjectKind {
98 Table,
99 Index,
100 Trigger,
101 View,
102}
103
104impl SchemaObjectKind {
105 fn sqlite_name(self) -> &'static str {
106 match self {
107 Self::Table => "table",
108 Self::Index => "index",
109 Self::Trigger => "trigger",
110 Self::View => "view",
111 }
112 }
113}
114
115#[derive(Debug, Clone, Copy, PartialEq, Eq)]
122pub struct SchemaObject {
123 pub kind: SchemaObjectKind,
124 pub name: &'static str,
125}
126
127#[derive(Debug)]
129pub struct SchemaDomain {
130 pub name: &'static str,
132 pub migrations: &'static [Migration],
134 pub initialize_current: fn(&Transaction<'_>) -> Result<(), rusqlite::Error>,
142 pub allowed_existing_versions: &'static [i64],
146 pub released_predecessors: &'static [SchemaPredecessor],
148 pub owned_objects: &'static [SchemaObject],
151 pub retired_objects: &'static [SchemaObject],
155}
156
157impl SchemaDomain {
158 pub fn supported_version(&self) -> i64 {
160 self.migrations.last().map_or(0, |m| m.version)
161 }
162
163 fn validate(&self) -> Result<(), SqliteStoreError> {
164 for (idx, migration) in self.migrations.iter().enumerate() {
165 let expected = idx as i64 + 1;
166 if migration.version != expected {
167 return Err(SqliteStoreError::InvalidMigrationList {
168 domain: self.name.to_string(),
169 detail: format!(
170 "migration at position {idx} has version {}, expected {expected} \
171 (versions must be contiguous from 1)",
172 migration.version
173 ),
174 });
175 }
176 }
177 let supported = self.supported_version();
178 let mut previous = None;
179 for &version in self.allowed_existing_versions {
180 if version <= 0 || version > supported {
181 return Err(SqliteStoreError::InvalidMigrationList {
182 domain: self.name.to_string(),
183 detail: format!(
184 "allowed existing version {version} is outside 1..={supported}"
185 ),
186 });
187 }
188 if previous.is_some_and(|value| value >= version) {
189 return Err(SqliteStoreError::InvalidMigrationList {
190 domain: self.name.to_string(),
191 detail: "allowed existing versions must be strictly increasing".to_string(),
192 });
193 }
194 previous = Some(version);
195 }
196 if !self.allowed_existing_versions.contains(&supported) {
197 return Err(SqliteStoreError::InvalidMigrationList {
198 domain: self.name.to_string(),
199 detail: format!(
200 "allowed existing versions must explicitly include current version {supported}"
201 ),
202 });
203 }
204 for &version in self
205 .allowed_existing_versions
206 .iter()
207 .filter(|&&version| version < supported)
208 {
209 let matches = self
210 .released_predecessors
211 .iter()
212 .filter(|predecessor| predecessor.version == version)
213 .count();
214 if matches != 1 {
215 return Err(SqliteStoreError::InvalidMigrationList {
216 domain: self.name.to_string(),
217 detail: format!(
218 "allowed predecessor version {version} must have exactly one frozen \
219 verifier, found {matches}"
220 ),
221 });
222 }
223 }
224 for predecessor in self.released_predecessors {
225 if predecessor.version >= supported
226 || !self
227 .allowed_existing_versions
228 .contains(&predecessor.version)
229 {
230 return Err(SqliteStoreError::InvalidMigrationList {
231 domain: self.name.to_string(),
232 detail: format!(
233 "fingerprint verifier for version {} is not an allowed predecessor",
234 predecessor.version
235 ),
236 });
237 }
238 }
239 for (idx, object) in self
240 .owned_objects
241 .iter()
242 .chain(self.retired_objects)
243 .enumerate()
244 {
245 if object.name.is_empty() || object.name == "meerkat_schema" {
246 return Err(SqliteStoreError::InvalidMigrationList {
247 domain: self.name.to_string(),
248 detail: format!(
249 "owned object at position {idx} has reserved or empty name `{}`",
250 object.name
251 ),
252 });
253 }
254 if self
255 .owned_objects
256 .iter()
257 .chain(self.retired_objects)
258 .take(idx)
259 .any(|prior| prior.name == object.name)
260 {
261 return Err(SqliteStoreError::InvalidMigrationList {
262 domain: self.name.to_string(),
263 detail: format!("owned object name `{}` is duplicated", object.name),
264 });
265 }
266 }
267 Ok(())
268 }
269
270 fn accepts_existing_version(&self, version: i64) -> bool {
271 self.allowed_existing_versions.contains(&version)
272 }
273
274 fn verify_predecessor(&self, conn: &Connection, version: i64) -> Result<(), SqliteStoreError> {
275 if version == self.supported_version() {
276 return verify_current_schema_fingerprint(conn, self).map_err(|detail| {
277 SqliteStoreError::SchemaFingerprintMismatch {
278 domain: self.name.to_string(),
279 version,
280 detail,
281 }
282 });
283 }
284 let predecessor = self
285 .released_predecessors
286 .iter()
287 .find(|predecessor| predecessor.version == version)
288 .ok_or_else(|| unsupported_predecessor(self, version))?;
289 (predecessor.verify)(conn).map_err(|detail| SqliteStoreError::SchemaFingerprintMismatch {
290 domain: self.name.to_string(),
291 version,
292 detail,
293 })
294 }
295}
296
297#[derive(Debug, Clone, Copy, PartialEq, Eq)]
299pub struct LedgerReport {
300 pub from_version: i64,
302 pub to_version: i64,
304}
305
306impl LedgerReport {
307 pub fn migrated(&self) -> bool {
309 self.to_version > self.from_version
310 }
311}
312
313#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
315pub struct MaintenancePrepareReport {
316 pub changed: usize,
318}
319
320#[derive(Debug, Clone, Copy, PartialEq, Eq)]
322pub struct MaintenanceBridgeReport {
323 pub from_version: i64,
325 pub to_version: i64,
327 pub prepared: usize,
329}
330
331impl MaintenanceBridgeReport {
332 pub fn migrated(&self) -> bool {
334 self.to_version > self.from_version
335 }
336
337 pub fn changed(&self) -> bool {
339 self.migrated() || self.prepared > 0
340 }
341}
342
343pub fn domain_version(conn: &Connection, domain: &str) -> Result<Option<i64>, SqliteStoreError> {
351 if !ledger_table_exists(conn)? {
352 return Ok(None);
353 }
354 validate_ledger_shape(conn)?;
355 read_version(conn, domain)
356}
357
358pub fn preflight_schema_eligibility(
373 conn: &Connection,
374 domain: &SchemaDomain,
375) -> Result<(), SqliteStoreError> {
376 domain.validate()?;
377 let supported = domain.supported_version();
378 match domain_version(conn, domain.name)? {
379 Some(found) if found > supported => {
380 return Err(SqliteStoreError::SchemaFromTheFuture {
381 domain: domain.name.to_string(),
382 found,
383 supported,
384 });
385 }
386 Some(found) if !domain.accepts_existing_version(found) => {
387 return Err(unsupported_predecessor(domain, found));
388 }
389 Some(found) => domain.verify_predecessor(conn, found)?,
390 None => {
391 let objects = find_owned_objects(conn, domain)?;
392 if !objects.is_empty() {
393 return Err(SqliteStoreError::UnledgeredDomainObjects {
394 domain: domain.name.to_string(),
395 objects,
396 });
397 }
398 }
399 }
400 Ok(())
401}
402
403pub fn apply_domain_migrations(
410 conn: &mut Connection,
411 domain: &SchemaDomain,
412) -> Result<LedgerReport, SqliteStoreError> {
413 domain.validate()?;
414 let supported = domain.supported_version();
415
416 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
417 let current = if ledger_table_exists(&tx)? {
420 validate_ledger_shape(&tx)?;
421 read_version(&tx, domain.name)?
422 } else {
423 None
424 };
425 if let Some(found) = current {
426 if found > supported {
427 return Err(SqliteStoreError::SchemaFromTheFuture {
428 domain: domain.name.to_string(),
429 found,
430 supported,
431 });
432 }
433 if !domain.accepts_existing_version(found) {
434 return Err(unsupported_predecessor(domain, found));
435 }
436 domain.verify_predecessor(&tx, found)?;
437 } else {
438 let objects = find_owned_objects(&tx, domain)?;
439 if !objects.is_empty() {
440 return Err(SqliteStoreError::UnledgeredDomainObjects {
441 domain: domain.name.to_string(),
442 objects,
443 });
444 }
445 }
446 let current = current.unwrap_or(0);
447 if current == supported {
448 return Ok(LedgerReport {
449 from_version: current,
450 to_version: current,
451 });
452 }
453
454 if !ledger_table_exists(&tx)? {
457 tx.execute_batch(CREATE_LEDGER_SQL)?;
458 validate_ledger_shape(&tx)?;
459 }
460
461 if current == 0 {
462 tx.execute_batch(CUSTODY_SAVEPOINT_SQL)?;
463 (domain.initialize_current)(&tx).map_err(|source| SqliteStoreError::MigrationFailed {
464 domain: domain.name.to_string(),
465 version: supported,
466 name: "initialize-current".to_string(),
467 source,
468 })?;
469 if tx.is_autocommit() || tx.execute_batch(CUSTODY_RELEASE_SQL).is_err() {
470 return Err(SqliteStoreError::MigrationBrokeTransaction {
471 domain: domain.name.to_string(),
472 version: supported,
473 name: "initialize-current".to_string(),
474 });
475 }
476 } else {
477 for migration in domain.migrations.iter().filter(|m| m.version > current) {
478 tx.execute_batch(CUSTODY_SAVEPOINT_SQL)?;
479 (migration.apply)(&tx).map_err(|source| SqliteStoreError::MigrationFailed {
480 domain: domain.name.to_string(),
481 version: migration.version,
482 name: migration.name.to_string(),
483 source,
484 })?;
485 if tx.is_autocommit() || tx.execute_batch(CUSTODY_RELEASE_SQL).is_err() {
495 return Err(SqliteStoreError::MigrationBrokeTransaction {
496 domain: domain.name.to_string(),
497 version: migration.version,
498 name: migration.name.to_string(),
499 });
500 }
501 }
502 }
503 verify_current_schema_fingerprint(&tx, domain).map_err(|detail| {
504 SqliteStoreError::SchemaFingerprintMismatch {
505 domain: domain.name.to_string(),
506 version: supported,
507 detail,
508 }
509 })?;
510 tx.execute(
511 "INSERT INTO main.meerkat_schema (domain, version) VALUES (?1, ?2)
512 ON CONFLICT(domain) DO UPDATE SET version = excluded.version",
513 rusqlite::params![domain.name, supported],
514 )?;
515 verify_ledger_stamp(&tx, domain.name, supported)?;
516 tx.commit()?;
517
518 Ok(LedgerReport {
519 from_version: current,
520 to_version: supported,
521 })
522}
523
524pub fn bridge_unledgered_domain(
549 conn: &mut Connection,
550 domain: &SchemaDomain,
551 target_version: i64,
552 recoverable_source_versions: &[i64],
553 prepare: Option<fn(&Transaction<'_>) -> Result<MaintenancePrepareReport, rusqlite::Error>>,
554) -> Result<MaintenanceBridgeReport, SqliteStoreError> {
555 domain.validate()?;
556 let supported = domain.supported_version();
557 if target_version > supported {
558 return Err(SqliteStoreError::SchemaFromTheFuture {
559 domain: domain.name.to_string(),
560 found: target_version,
561 supported,
562 });
563 }
564 if !domain.accepts_existing_version(target_version) {
565 return Err(unsupported_predecessor(domain, target_version));
566 }
567 validate_recoverable_source_versions(domain, target_version, recoverable_source_versions)?;
568
569 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
570 let current = if ledger_table_exists(&tx)? {
571 validate_ledger_shape(&tx)?;
572 read_version(&tx, domain.name)?
573 } else {
574 None
575 };
576
577 let mut inferred_oracles = None;
578 let from_version = if let Some(found) = current {
579 if found > target_version {
580 return Err(SqliteStoreError::SchemaFromTheFuture {
581 domain: domain.name.to_string(),
582 found,
583 supported: target_version,
584 });
585 }
586 if !domain.accepts_existing_version(found) {
587 return Err(unsupported_predecessor(domain, found));
588 }
589 domain.verify_predecessor(&tx, found)?;
590 if found == target_version && prepare.is_none() {
591 return Ok(MaintenanceBridgeReport {
592 from_version: found,
593 to_version: found,
594 prepared: 0,
595 });
596 }
597 found
598 } else {
599 let objects = find_owned_objects(&tx, domain)?;
600 if objects.is_empty() {
601 return Ok(MaintenanceBridgeReport {
602 from_version: 0,
603 to_version: 0,
604 prepared: 0,
605 });
606 }
607
608 let oracles = build_migration_prefix_oracles(domain, target_version)?;
609 let actual = domain_catalog_fingerprint(&tx, domain).map_err(|detail| {
610 SqliteStoreError::UnledgeredSchemaNoMatch {
611 domain: domain.name.to_string(),
612 target_version,
613 objects: vec![detail],
614 }
615 })?;
616 let mut matches = oracles
617 .iter()
618 .filter_map(|(version, fingerprint)| {
619 (recoverable_source_versions.contains(version) && fingerprint == &actual)
620 .then_some(*version)
621 })
622 .collect::<Vec<_>>();
623 for predecessor in domain.released_predecessors.iter().filter(|predecessor| {
631 recoverable_source_versions.contains(&predecessor.version)
632 && predecessor.version <= target_version
633 }) {
634 if (predecessor.verify)(&tx).is_ok() && !matches.contains(&predecessor.version) {
635 matches.push(predecessor.version);
636 }
637 }
638 matches.sort_unstable();
639 let matched = match matches.as_slice() {
640 [version] => *version,
641 [] => {
642 return Err(SqliteStoreError::UnledgeredSchemaNoMatch {
643 domain: domain.name.to_string(),
644 target_version,
645 objects,
646 });
647 }
648 _ => {
649 return Err(SqliteStoreError::UnledgeredSchemaAmbiguous {
650 domain: domain.name.to_string(),
651 target_version,
652 matches,
653 });
654 }
655 };
656 inferred_oracles = Some(oracles);
657 matched
658 };
659
660 let oracles = match inferred_oracles {
661 Some(oracles) => oracles,
662 None => build_migration_prefix_oracles(domain, target_version)?,
663 };
664 validate_domain_trigger_isolation(&tx, domain, from_version)?;
665
666 let prepared = match prepare {
667 Some(prepare) => {
668 run_with_custody(&tx, domain, from_version, "maintenance-prepare", prepare)?.changed
669 }
670 None => 0,
671 };
672 for migration in domain
673 .migrations
674 .iter()
675 .filter(|migration| migration.version > from_version && migration.version <= target_version)
676 {
677 run_with_custody(
678 &tx,
679 domain,
680 migration.version,
681 migration.name,
682 migration.apply,
683 )?;
684 }
685
686 let target = oracles
687 .iter()
688 .find_map(|(version, fingerprint)| (*version == target_version).then_some(fingerprint))
689 .ok_or_else(|| SqliteStoreError::InvalidMigrationList {
690 domain: domain.name.to_string(),
691 detail: format!(
692 "migration-prefix oracle did not produce requested target version {target_version}"
693 ),
694 })?;
695 let converged = domain_catalog_fingerprint(&tx, domain).map_err(|detail| {
696 SqliteStoreError::SchemaFingerprintMismatch {
697 domain: domain.name.to_string(),
698 version: target_version,
699 detail,
700 }
701 })?;
702 if &converged != target {
703 return Err(SqliteStoreError::SchemaFingerprintMismatch {
704 domain: domain.name.to_string(),
705 version: target_version,
706 detail: format!(
707 "migration-prefix catalog differs: expected {target:?}, found {converged:?}"
708 ),
709 });
710 }
711 domain.verify_predecessor(&tx, target_version)?;
712
713 if current != Some(target_version) {
714 if !ledger_table_exists(&tx)? {
715 tx.execute_batch(CREATE_LEDGER_SQL)?;
716 validate_ledger_shape(&tx)?;
717 }
718 tx.execute(
719 "INSERT INTO main.meerkat_schema (domain, version) VALUES (?1, ?2)
720 ON CONFLICT(domain) DO UPDATE SET version = excluded.version",
721 rusqlite::params![domain.name, target_version],
722 )?;
723 }
724 verify_ledger_stamp(&tx, domain.name, target_version)?;
725 tx.commit()?;
726
727 Ok(MaintenanceBridgeReport {
728 from_version,
729 to_version: target_version,
730 prepared,
731 })
732}
733
734fn validate_recoverable_source_versions(
735 domain: &SchemaDomain,
736 target_version: i64,
737 versions: &[i64],
738) -> Result<(), SqliteStoreError> {
739 let mut previous = None;
740 for &version in versions {
741 if version <= 0 || version > target_version {
742 return Err(SqliteStoreError::InvalidMigrationList {
743 domain: domain.name.to_string(),
744 detail: format!(
745 "recoverable source version {version} is outside 1..={target_version}"
746 ),
747 });
748 }
749 if previous.is_some_and(|prior| prior >= version) {
750 return Err(SqliteStoreError::InvalidMigrationList {
751 domain: domain.name.to_string(),
752 detail: "recoverable source versions must be strictly increasing".to_string(),
753 });
754 }
755 previous = Some(version);
756 }
757 Ok(())
758}
759
760fn verify_ledger_stamp(
761 conn: &Connection,
762 domain: &str,
763 expected: i64,
764) -> Result<(), SqliteStoreError> {
765 let found = read_version(conn, domain)?;
766 if found != Some(expected) {
767 return Err(malformed(format!(
768 "domain `{domain}` stamp did not persist exact version {expected}; found {found:?}"
769 )));
770 }
771 Ok(())
772}
773
774fn validate_domain_trigger_isolation(
775 conn: &Connection,
776 domain: &SchemaDomain,
777 source_version: i64,
778) -> Result<(), SqliteStoreError> {
779 let all_objects = all_domain_objects(domain);
780 let allowed_triggers = all_objects
781 .iter()
782 .filter(|object| object.kind == SchemaObjectKind::Trigger)
783 .map(|object| object.name)
784 .collect::<Vec<_>>();
785 let mut statement = conn
786 .prepare(
787 "SELECT 'main', name FROM main.sqlite_schema
788 WHERE type = 'trigger' AND tbl_name = ?1 COLLATE NOCASE
789 UNION ALL
790 SELECT 'temp', name FROM temp.sqlite_schema
791 WHERE type = 'trigger' AND tbl_name = ?1 COLLATE NOCASE
792 ORDER BY 1, 2",
793 )
794 .map_err(SqliteStoreError::Sqlite)?;
795 let mut refused = Vec::new();
796 for target in all_objects.iter().filter(|object| {
797 matches!(
798 object.kind,
799 SchemaObjectKind::Table | SchemaObjectKind::View
800 )
801 }) {
802 let rows = statement
803 .query_map([target.name], |row| {
804 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
805 })
806 .map_err(SqliteStoreError::Sqlite)?;
807 for row in rows {
808 let (schema, trigger) = row.map_err(SqliteStoreError::Sqlite)?;
809 let declared = schema == "main"
810 && allowed_triggers
811 .iter()
812 .any(|allowed| allowed.eq_ignore_ascii_case(&trigger));
813 if !declared {
814 refused.push(format!("{schema}.{trigger} on {}", target.name));
815 }
816 }
817 }
818 refused.sort();
819 refused.dedup();
820 if !refused.is_empty() {
821 return Err(SqliteStoreError::SchemaFingerprintMismatch {
822 domain: domain.name.to_string(),
823 version: source_version,
824 detail: format!(
825 "undeclared or TEMP triggers can intercept maintenance writes: {refused:?}"
826 ),
827 });
828 }
829 Ok(())
830}
831
832fn run_with_custody<T>(
833 tx: &Transaction<'_>,
834 domain: &SchemaDomain,
835 version: i64,
836 name: &str,
837 body: fn(&Transaction<'_>) -> Result<T, rusqlite::Error>,
838) -> Result<T, SqliteStoreError> {
839 tx.execute_batch(CUSTODY_SAVEPOINT_SQL)?;
840 let body_result = body(tx);
841 if tx.is_autocommit() || tx.execute_batch(CUSTODY_RELEASE_SQL).is_err() {
842 return Err(SqliteStoreError::MigrationBrokeTransaction {
843 domain: domain.name.to_string(),
844 version,
845 name: name.to_string(),
846 });
847 }
848 body_result.map_err(|source| SqliteStoreError::MigrationFailed {
849 domain: domain.name.to_string(),
850 version,
851 name: name.to_string(),
852 source,
853 })
854}
855
856fn build_migration_prefix_oracles(
857 domain: &SchemaDomain,
858 target_version: i64,
859) -> Result<Vec<(i64, DomainCatalogFingerprint)>, SqliteStoreError> {
860 let mut expected = Connection::open_in_memory().map_err(SqliteStoreError::Sqlite)?;
861 let tx = expected
862 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
863 .map_err(SqliteStoreError::Sqlite)?;
864 let mut oracles = Vec::with_capacity(target_version as usize);
865 for migration in domain
866 .migrations
867 .iter()
868 .filter(|migration| migration.version <= target_version)
869 {
870 (migration.apply)(&tx).map_err(|source| SqliteStoreError::MigrationFailed {
871 domain: domain.name.to_string(),
872 version: migration.version,
873 name: format!("migration-prefix-oracle:{}", migration.name),
874 source,
875 })?;
876 let fingerprint = domain_catalog_fingerprint(&tx, domain).map_err(|detail| {
877 SqliteStoreError::InvalidMigrationList {
878 domain: domain.name.to_string(),
879 detail: format!(
880 "migration-prefix oracle version {} is inconsistent with ownership: {detail}",
881 migration.version
882 ),
883 }
884 })?;
885 oracles.push((migration.version, fingerprint));
886 }
887 Ok(oracles)
888}
889
890#[derive(Debug, PartialEq, Eq)]
891struct DomainCatalogFingerprint {
892 names: Vec<(String, String)>,
893 objects: Vec<CatalogObjectFingerprint>,
894}
895
896fn domain_catalog_fingerprint(
897 conn: &Connection,
898 domain: &SchemaDomain,
899) -> Result<DomainCatalogFingerprint, String> {
900 let all_objects = all_domain_objects(domain);
901 let names = catalog_names(conn, &all_objects).map_err(|error| error.to_string())?;
902 let declared = all_objects
903 .iter()
904 .map(|object| (object.name, object))
905 .collect::<BTreeMap<_, _>>();
906 let mut objects = Vec::with_capacity(names.len());
907 for (kind, name) in &names {
908 let Some(object) = declared.get(name.as_str()) else {
909 return Err(format!("undeclared owned object `{name}`"));
910 };
911 if kind != object.kind.sqlite_name() {
912 return Err(format!(
913 "owned object `{name}` has kind `{kind}`, expected `{}`",
914 object.kind.sqlite_name()
915 ));
916 }
917 objects.push(
918 catalog_fingerprint(conn, object)
919 .map_err(|error| format!("fingerprint object `{name}`: {error}"))?,
920 );
921 }
922 Ok(DomainCatalogFingerprint { names, objects })
923}
924
925static EXPECTED_CURRENT_CATALOGS: OnceLock<Mutex<BTreeMap<String, Result<String, String>>>> =
926 OnceLock::new();
927
928fn verify_current_schema_fingerprint(
932 actual: &Connection,
933 domain: &SchemaDomain,
934) -> Result<(), String> {
935 let expected = {
936 let cache = EXPECTED_CURRENT_CATALOGS.get_or_init(|| Mutex::new(BTreeMap::new()));
937 let key = current_catalog_cache_key(domain);
938 let cached = cache
939 .lock()
940 .map_err(|_| "current catalog cache lock is poisoned".to_string())?
941 .get(&key)
942 .cloned();
943 if let Some(cached) = cached {
944 cached?
945 } else {
946 let built = build_current_catalog_fingerprint(domain);
947 cache
948 .lock()
949 .map_err(|_| "current catalog cache lock is poisoned".to_string())?
950 .insert(key, built.clone());
951 built?
952 }
953 };
954 let actual = compact_catalog_fingerprint(actual, domain, domain.owned_objects)?;
955 if actual != expected {
956 return Err(format!(
957 "current owned catalog differs: expected {expected}, found {actual}"
958 ));
959 }
960 Ok(())
961}
962
963fn current_catalog_cache_key(domain: &SchemaDomain) -> String {
968 let mut key = format!(
969 "{}\u{1f}{}\u{1f}{:x}",
970 domain.name,
971 domain.supported_version(),
972 domain.initialize_current as usize
973 );
974 for object in domain.owned_objects {
975 key.push_str(&format!(
976 "\u{1e}current:{}:{}",
977 object.kind.sqlite_name(),
978 object.name
979 ));
980 }
981 for object in domain.retired_objects {
982 key.push_str(&format!(
983 "\u{1e}retired:{}:{}",
984 object.kind.sqlite_name(),
985 object.name
986 ));
987 }
988 key
989}
990
991fn build_current_catalog_fingerprint(domain: &SchemaDomain) -> Result<String, String> {
992 let mut expected =
993 Connection::open_in_memory().map_err(|error| format!("open current oracle: {error}"))?;
994 let tx = expected
995 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
996 .map_err(|error| format!("begin current oracle: {error}"))?;
997 (domain.initialize_current)(&tx).map_err(|error| format!("build current oracle: {error}"))?;
998 tx.commit()
999 .map_err(|error| format!("commit current oracle: {error}"))?;
1000 compact_catalog_fingerprint(&expected, domain, domain.owned_objects)
1001}
1002
1003fn compact_catalog_fingerprint(
1004 conn: &Connection,
1005 domain: &SchemaDomain,
1006 expected_objects: &[SchemaObject],
1007) -> Result<String, String> {
1008 let all_objects = all_domain_objects(domain);
1009 let owned_by_name = all_objects
1010 .iter()
1011 .map(|object| (object.name, object))
1012 .collect::<BTreeMap<_, _>>();
1013 let current_by_name = expected_objects
1014 .iter()
1015 .map(|object| (object.name, object))
1016 .collect::<BTreeMap<_, _>>();
1017 let mut actual_names = Vec::new();
1018 let mut entries = Vec::with_capacity(expected_objects.len());
1019 let mut statement = conn
1020 .prepare(
1021 "SELECT type, name, tbl_name, sql
1022 FROM main.sqlite_schema
1023 WHERE name NOT LIKE 'sqlite_%'
1024 ORDER BY type, name",
1025 )
1026 .map_err(|error| error.to_string())?;
1027 let rows = statement
1028 .query_map([], |row| {
1029 Ok((
1030 row.get::<_, String>(0)?,
1031 row.get::<_, String>(1)?,
1032 row.get::<_, String>(2)?,
1033 row.get::<_, Option<String>>(3)?,
1034 ))
1035 })
1036 .map_err(|error| error.to_string())?;
1037 for row in rows {
1038 let (kind, name, table_name, sql) = row.map_err(|error| error.to_string())?;
1039 if owned_by_name.contains_key(name.as_str()) {
1040 actual_names.push((kind.clone(), name.clone()));
1041 }
1042 if current_by_name.contains_key(name.as_str()) {
1043 entries.push(format!(
1044 "{kind}\u{1f}{name}\u{1f}{table_name}\u{1f}{}",
1045 sql.map(|sql| normalize_schema_sql(&sql))
1046 .unwrap_or_default()
1047 ));
1048 }
1049 }
1050 actual_names.sort();
1051 let mut expected_names = expected_objects
1052 .iter()
1053 .map(|object| {
1054 (
1055 object.kind.sqlite_name().to_string(),
1056 object.name.to_string(),
1057 )
1058 })
1059 .collect::<Vec<_>>();
1060 expected_names.sort();
1061 if actual_names != expected_names {
1062 return Err(format!(
1063 "owned object set differs: expected {expected_names:?}, found {actual_names:?}"
1064 ));
1065 }
1066
1067 entries.sort();
1068 Ok(entries.join("\u{1e}"))
1069}
1070
1071fn all_domain_objects(domain: &SchemaDomain) -> Vec<SchemaObject> {
1072 domain
1073 .owned_objects
1074 .iter()
1075 .chain(domain.retired_objects)
1076 .copied()
1077 .collect()
1078}
1079
1080pub fn verify_released_schema_fingerprint(
1090 actual: &Connection,
1091 domain: &SchemaDomain,
1092 released_objects: &[SchemaObject],
1093 build_released: fn(&Transaction<'_>) -> Result<(), rusqlite::Error>,
1094) -> Result<(), String> {
1095 let mut expected = Connection::open_in_memory()
1096 .map_err(|error| format!("open fingerprint oracle: {error}"))?;
1097 let tx = expected
1098 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1099 .map_err(|error| format!("begin fingerprint oracle: {error}"))?;
1100 build_released(&tx).map_err(|error| format!("build fingerprint oracle: {error}"))?;
1101 tx.commit()
1102 .map_err(|error| format!("commit fingerprint oracle: {error}"))?;
1103
1104 let expected_names = catalog_names(&expected, released_objects)
1105 .map_err(|error| format!("read fingerprint oracle: {error}"))?;
1106 let mut declared_expected = released_objects
1107 .iter()
1108 .map(|object| {
1109 (
1110 object.kind.sqlite_name().to_string(),
1111 object.name.to_string(),
1112 )
1113 })
1114 .collect::<Vec<_>>();
1115 declared_expected.sort();
1116 if expected_names != declared_expected {
1117 return Err(format!(
1118 "frozen builder produced {expected_names:?}, manifest declares {declared_expected:?}"
1119 ));
1120 }
1121
1122 let actual_names = catalog_names(actual, &all_domain_objects(domain))
1123 .map_err(|error| format!("read actual catalog: {error}"))?;
1124 if actual_names != declared_expected {
1125 return Err(format!(
1126 "owned object set differs: expected {declared_expected:?}, found {actual_names:?}"
1127 ));
1128 }
1129
1130 for object in released_objects {
1131 let wanted = catalog_fingerprint(&expected, object)
1132 .map_err(|error| format!("fingerprint oracle {}: {error}", object.name))?;
1133 let found = catalog_fingerprint(actual, object)
1134 .map_err(|error| format!("fingerprint actual {}: {error}", object.name))?;
1135 if found != wanted {
1136 return Err(format!(
1137 "object `{}` differs: expected {wanted:?}, found {found:?}",
1138 object.name
1139 ));
1140 }
1141 }
1142 Ok(())
1143}
1144
1145#[derive(Debug, PartialEq, Eq)]
1146struct CatalogObjectFingerprint {
1147 kind: String,
1148 name: String,
1149 table_name: String,
1150 normalized_sql: Option<String>,
1151 table_columns: Vec<TableColumnFingerprint>,
1152 foreign_keys: Vec<ForeignKeyFingerprint>,
1153 index: Option<IndexFingerprint>,
1154}
1155
1156#[derive(Debug, PartialEq, Eq)]
1157struct TableColumnFingerprint {
1158 cid: i64,
1159 name: String,
1160 declared_type: String,
1161 not_null: bool,
1162 default_value: Option<String>,
1163 primary_key_position: i64,
1164 hidden: i64,
1165}
1166
1167#[derive(Debug, PartialEq, Eq)]
1168struct ForeignKeyFingerprint {
1169 id: i64,
1170 sequence: i64,
1171 target_table: String,
1172 from_column: String,
1173 to_column: Option<String>,
1174 on_update: String,
1175 on_delete: String,
1176 match_clause: String,
1177}
1178
1179#[derive(Debug, PartialEq, Eq)]
1180struct IndexFingerprint {
1181 unique: bool,
1182 origin: String,
1183 partial: bool,
1184 columns: Vec<IndexColumnFingerprint>,
1185}
1186
1187#[derive(Debug, PartialEq, Eq)]
1188struct IndexColumnFingerprint {
1189 sequence: i64,
1190 column_id: i64,
1191 name: Option<String>,
1192 descending: bool,
1193 collation: Option<String>,
1194 key: bool,
1195}
1196
1197fn catalog_names(
1198 conn: &Connection,
1199 objects: &[SchemaObject],
1200) -> Result<Vec<(String, String)>, rusqlite::Error> {
1201 let mut found = Vec::new();
1202 let mut statement = conn.prepare(
1203 "SELECT type, name FROM main.sqlite_schema
1204 WHERE name = ?1 AND name NOT LIKE 'sqlite_%'
1205 ORDER BY type, name",
1206 )?;
1207 for object in objects {
1208 let rows = statement.query_map([object.name], |row| Ok((row.get(0)?, row.get(1)?)))?;
1209 found.extend(rows.collect::<Result<Vec<_>, _>>()?);
1210 }
1211 found.sort();
1212 found.dedup();
1213 Ok(found)
1214}
1215
1216fn catalog_fingerprint(
1217 conn: &Connection,
1218 object: &SchemaObject,
1219) -> Result<CatalogObjectFingerprint, rusqlite::Error> {
1220 let (kind, name, table_name, sql): (String, String, String, Option<String>) = conn.query_row(
1221 "SELECT type, name, tbl_name, sql FROM main.sqlite_schema WHERE name = ?1",
1222 [object.name],
1223 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
1224 )?;
1225 let (table_columns, foreign_keys) = if kind == "table" || kind == "view" {
1226 (
1227 table_columns(conn, object.name)?,
1228 foreign_keys(conn, object.name)?,
1229 )
1230 } else {
1231 (Vec::new(), Vec::new())
1232 };
1233 let index = if kind == "index" {
1234 Some(index_fingerprint(conn, &table_name, object.name)?)
1235 } else {
1236 None
1237 };
1238 Ok(CatalogObjectFingerprint {
1239 kind,
1240 name,
1241 table_name,
1242 normalized_sql: sql.map(|sql| normalize_schema_sql(&sql)),
1243 table_columns,
1244 foreign_keys,
1245 index,
1246 })
1247}
1248
1249fn table_columns(
1250 conn: &Connection,
1251 table: &str,
1252) -> Result<Vec<TableColumnFingerprint>, rusqlite::Error> {
1253 let mut statement = conn.prepare(
1254 "SELECT cid, name, type, \"notnull\", dflt_value, pk, hidden
1255 FROM pragma_table_xinfo(?1, 'main')
1256 ORDER BY cid",
1257 )?;
1258 let rows = statement.query_map([table], |row| {
1259 Ok(TableColumnFingerprint {
1260 cid: row.get(0)?,
1261 name: row.get(1)?,
1262 declared_type: row.get(2)?,
1263 not_null: row.get(3)?,
1264 default_value: row.get(4)?,
1265 primary_key_position: row.get(5)?,
1266 hidden: row.get(6)?,
1267 })
1268 })?;
1269 rows.collect()
1270}
1271
1272fn foreign_keys(
1273 conn: &Connection,
1274 table: &str,
1275) -> Result<Vec<ForeignKeyFingerprint>, rusqlite::Error> {
1276 let mut statement = conn.prepare(
1277 "SELECT id, seq, \"table\", \"from\", \"to\", on_update, on_delete, \"match\"
1278 FROM pragma_foreign_key_list(?1, 'main')
1279 ORDER BY id, seq",
1280 )?;
1281 let rows = statement.query_map([table], |row| {
1282 Ok(ForeignKeyFingerprint {
1283 id: row.get(0)?,
1284 sequence: row.get(1)?,
1285 target_table: row.get(2)?,
1286 from_column: row.get(3)?,
1287 to_column: row.get(4)?,
1288 on_update: row.get(5)?,
1289 on_delete: row.get(6)?,
1290 match_clause: row.get(7)?,
1291 })
1292 })?;
1293 rows.collect()
1294}
1295
1296fn index_fingerprint(
1297 conn: &Connection,
1298 table: &str,
1299 index: &str,
1300) -> Result<IndexFingerprint, rusqlite::Error> {
1301 let (unique, origin, partial): (bool, String, bool) = conn.query_row(
1302 "SELECT \"unique\", origin, partial
1303 FROM pragma_index_list(?1, 'main')
1304 WHERE name = ?2",
1305 rusqlite::params![table, index],
1306 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
1307 )?;
1308 let mut statement = conn.prepare(
1309 "SELECT seqno, cid, name, desc, coll, key
1310 FROM pragma_index_xinfo(?1, 'main')
1311 ORDER BY seqno",
1312 )?;
1313 let columns = statement
1314 .query_map([index], |row| {
1315 Ok(IndexColumnFingerprint {
1316 sequence: row.get(0)?,
1317 column_id: row.get(1)?,
1318 name: row.get(2)?,
1319 descending: row.get(3)?,
1320 collation: row.get(4)?,
1321 key: row.get(5)?,
1322 })
1323 })?
1324 .collect::<Result<Vec<_>, _>>()?;
1325 Ok(IndexFingerprint {
1326 unique,
1327 origin,
1328 partial,
1329 columns,
1330 })
1331}
1332
1333fn normalize_schema_sql(sql: &str) -> String {
1334 #[derive(Clone, Copy)]
1335 enum LexState {
1336 Normal,
1337 SingleQuoted,
1338 DoubleQuoted,
1339 BacktickQuoted,
1340 BracketQuoted,
1341 LineComment,
1342 BlockComment,
1343 }
1344
1345 let bytes = sql.as_bytes();
1346 let mut collapsed = Vec::with_capacity(bytes.len());
1347 let mut state = LexState::Normal;
1348 let mut index = 0;
1349 while index < bytes.len() {
1350 let byte = bytes[index];
1351 match state {
1352 LexState::Normal => {
1353 if byte.is_ascii_whitespace() {
1354 if collapsed.last().is_some_and(|last| *last != b' ') {
1355 collapsed.push(b' ');
1356 }
1357 } else {
1358 collapsed.push(byte);
1359 state = match byte {
1360 b'\'' => LexState::SingleQuoted,
1361 b'"' => LexState::DoubleQuoted,
1362 b'`' => LexState::BacktickQuoted,
1363 b'[' => LexState::BracketQuoted,
1364 b'-' if bytes.get(index + 1) == Some(&b'-') => LexState::LineComment,
1365 b'/' if bytes.get(index + 1) == Some(&b'*') => LexState::BlockComment,
1366 _ => LexState::Normal,
1367 };
1368 }
1369 }
1370 LexState::SingleQuoted | LexState::DoubleQuoted | LexState::BacktickQuoted => {
1371 collapsed.push(byte);
1372 let delimiter = match state {
1373 LexState::SingleQuoted => b'\'',
1374 LexState::DoubleQuoted => b'"',
1375 LexState::BacktickQuoted => b'`',
1376 _ => unreachable!(),
1377 };
1378 if byte == delimiter {
1379 if bytes.get(index + 1) == Some(&delimiter) {
1380 index += 1;
1381 collapsed.push(delimiter);
1382 } else {
1383 state = LexState::Normal;
1384 }
1385 }
1386 }
1387 LexState::BracketQuoted => {
1388 collapsed.push(byte);
1389 if byte == b']' {
1390 if bytes.get(index + 1) == Some(&b']') {
1391 index += 1;
1392 collapsed.push(b']');
1393 } else {
1394 state = LexState::Normal;
1395 }
1396 }
1397 }
1398 LexState::LineComment => {
1399 collapsed.push(byte);
1400 if byte == b'\n' || byte == b'\r' {
1401 state = LexState::Normal;
1402 }
1403 }
1404 LexState::BlockComment => {
1405 collapsed.push(byte);
1406 if byte == b'*' && bytes.get(index + 1) == Some(&b'/') {
1407 index += 1;
1408 collapsed.push(b'/');
1409 state = LexState::Normal;
1410 }
1411 }
1412 }
1413 index += 1;
1414 }
1415
1416 const IF_NOT_EXISTS: &[u8] = b"IF NOT EXISTS";
1417 let mut normalized = Vec::with_capacity(collapsed.len());
1418 let mut state = LexState::Normal;
1419 let mut index = 0;
1420 while index < collapsed.len() {
1421 let byte = collapsed[index];
1422 if matches!(state, LexState::Normal)
1423 && collapsed
1424 .get(index..index + IF_NOT_EXISTS.len())
1425 .is_some_and(|candidate| candidate.eq_ignore_ascii_case(IF_NOT_EXISTS))
1426 && (index == 0 || !is_sql_identifier_byte(collapsed[index - 1]))
1427 && collapsed
1428 .get(index + IF_NOT_EXISTS.len())
1429 .is_none_or(|after| !is_sql_identifier_byte(*after))
1430 {
1431 index += IF_NOT_EXISTS.len();
1432 if collapsed.get(index) == Some(&b' ') {
1433 index += 1;
1434 }
1435 continue;
1436 }
1437 normalized.push(byte);
1438 match state {
1439 LexState::Normal => {
1440 state = match byte {
1441 b'\'' => LexState::SingleQuoted,
1442 b'"' => LexState::DoubleQuoted,
1443 b'`' => LexState::BacktickQuoted,
1444 b'[' => LexState::BracketQuoted,
1445 b'-' if collapsed.get(index + 1) == Some(&b'-') => LexState::LineComment,
1446 b'/' if collapsed.get(index + 1) == Some(&b'*') => LexState::BlockComment,
1447 _ => LexState::Normal,
1448 };
1449 }
1450 LexState::SingleQuoted | LexState::DoubleQuoted | LexState::BacktickQuoted => {
1451 let delimiter = match state {
1452 LexState::SingleQuoted => b'\'',
1453 LexState::DoubleQuoted => b'"',
1454 LexState::BacktickQuoted => b'`',
1455 _ => unreachable!(),
1456 };
1457 if byte == delimiter {
1458 if collapsed.get(index + 1) == Some(&delimiter) {
1459 index += 1;
1460 normalized.push(delimiter);
1461 } else {
1462 state = LexState::Normal;
1463 }
1464 }
1465 }
1466 LexState::BracketQuoted => {
1467 if byte == b']' {
1468 if collapsed.get(index + 1) == Some(&b']') {
1469 index += 1;
1470 normalized.push(b']');
1471 } else {
1472 state = LexState::Normal;
1473 }
1474 }
1475 }
1476 LexState::LineComment => {
1477 if byte == b'\n' || byte == b'\r' {
1478 state = LexState::Normal;
1479 }
1480 }
1481 LexState::BlockComment => {
1482 if byte == b'*' && collapsed.get(index + 1) == Some(&b'/') {
1483 index += 1;
1484 normalized.push(b'/');
1485 state = LexState::Normal;
1486 }
1487 }
1488 }
1489 index += 1;
1490 }
1491 String::from_utf8(normalized).unwrap_or_else(|_| sql.to_string())
1492}
1493
1494fn is_sql_identifier_byte(byte: u8) -> bool {
1495 byte.is_ascii_alphanumeric() || byte == b'_'
1496}
1497
1498fn unsupported_predecessor(domain: &SchemaDomain, found: i64) -> SqliteStoreError {
1499 SqliteStoreError::UnsupportedSchemaPredecessor {
1500 domain: domain.name.to_string(),
1501 found,
1502 supported: domain.supported_version(),
1503 allowed: domain.allowed_existing_versions.to_vec(),
1504 }
1505}
1506
1507fn find_owned_objects(
1512 conn: &Connection,
1513 domain: &SchemaDomain,
1514) -> Result<Vec<String>, SqliteStoreError> {
1515 let mut found = Vec::new();
1516 let mut stmt = conn.prepare(
1517 "SELECT type, name FROM main.sqlite_schema
1518 WHERE name = ?1 AND name NOT LIKE 'sqlite_%'
1519 ORDER BY type, name",
1520 )?;
1521 for expected in domain.owned_objects.iter().chain(domain.retired_objects) {
1522 let rows = stmt.query_map([expected.name], |row| {
1523 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
1524 })?;
1525 for row in rows {
1526 let (actual_kind, name) = row?;
1527 found.push(format!(
1528 "{actual_kind}:{name} (expected {})",
1529 expected.kind.sqlite_name()
1530 ));
1531 }
1532 }
1533 found.sort();
1534 found.dedup();
1535 Ok(found)
1536}
1537
1538fn ledger_table_exists(conn: &Connection) -> Result<bool, SqliteStoreError> {
1539 let exists = conn
1540 .query_row(
1541 "SELECT 1 FROM main.sqlite_master WHERE type = 'table' AND name = 'meerkat_schema'",
1542 [],
1543 |_| Ok(()),
1544 )
1545 .optional()?
1546 .is_some();
1547 Ok(exists)
1548}
1549
1550fn malformed(detail: String) -> SqliteStoreError {
1551 SqliteStoreError::LedgerMalformed { detail }
1552}
1553
1554fn validate_ledger_shape(conn: &Connection) -> Result<(), SqliteStoreError> {
1563 let mut stmt = conn.prepare(
1564 "SELECT name, type, \"notnull\", pk FROM pragma_table_info('meerkat_schema', 'main')",
1565 )?;
1566 let mut rows = stmt.query([])?;
1567 let mut domain_ok = false;
1568 let mut version_ok = false;
1569 while let Some(row) = rows.next()? {
1570 let name: String = row.get(0)?;
1571 let decl_type: String = row.get(1)?;
1572 let notnull: bool = row.get(2)?;
1573 let pk: i64 = row.get(3)?;
1574 match name.as_str() {
1575 "domain" => {
1576 if !decl_type.eq_ignore_ascii_case("TEXT") || pk != 1 {
1577 return Err(malformed(format!(
1578 "column `domain` must be `TEXT PRIMARY KEY`, found type `{decl_type}` \
1579 with pk position {pk}"
1580 )));
1581 }
1582 domain_ok = true;
1583 }
1584 "version" => {
1585 if !decl_type.eq_ignore_ascii_case("INTEGER") || !notnull || pk != 0 {
1586 return Err(malformed(format!(
1587 "column `version` must be non-key `INTEGER NOT NULL`, found type \
1588 `{decl_type}` notnull={notnull} pk position {pk}"
1589 )));
1590 }
1591 version_ok = true;
1592 }
1593 other => {
1594 if pk != 0 {
1595 return Err(malformed(format!(
1596 "unexpected primary-key column `{other}`"
1597 )));
1598 }
1599 }
1600 }
1601 }
1602 if !domain_ok || !version_ok {
1603 return Err(malformed(
1604 "table lacks the pinned `domain`/`version` columns".to_string(),
1605 ));
1606 }
1607 let mut trigger_stmt = conn.prepare(
1608 "SELECT 'main', name FROM main.sqlite_schema
1609 WHERE type = 'trigger' AND tbl_name = 'meerkat_schema' COLLATE NOCASE
1610 UNION ALL
1611 SELECT 'temp', name FROM temp.sqlite_schema
1612 WHERE type = 'trigger' AND tbl_name = 'meerkat_schema' COLLATE NOCASE
1613 ORDER BY 1, 2",
1614 )?;
1615 let triggers = trigger_stmt
1616 .query_map([], |row| {
1617 Ok(format!(
1618 "{}.{}",
1619 row.get::<_, String>(0)?,
1620 row.get::<_, String>(1)?
1621 ))
1622 })?
1623 .collect::<Result<Vec<_>, _>>()?;
1624 if !triggers.is_empty() {
1625 return Err(malformed(format!(
1626 "table has attached triggers {triggers:?}; ledger writes must be isolated"
1627 )));
1628 }
1629 Ok(())
1630}
1631
1632fn read_version(conn: &Connection, domain: &str) -> Result<Option<i64>, SqliteStoreError> {
1639 let mut stmt = conn.prepare("SELECT version FROM main.meerkat_schema WHERE domain = ?1")?;
1640 let mut rows = stmt.query([domain])?;
1641 let Some(row) = rows.next()? else {
1642 return Ok(None);
1643 };
1644 let version: i64 = row.get(0)?;
1645 if rows.next()?.is_some() {
1646 return Err(malformed(format!(
1647 "multiple ledger rows for domain `{domain}`"
1648 )));
1649 }
1650 if version <= 0 {
1651 return Err(malformed(format!(
1652 "domain `{domain}` records non-positive version {version}"
1653 )));
1654 }
1655 Ok(Some(version))
1656}
1657
1658#[cfg(test)]
1659#[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
1660mod tests {
1661 use super::*;
1662 use crate::profile::{ConnectionProfile, open};
1663
1664 fn create_t1(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
1665 tx.execute_batch("CREATE TABLE IF NOT EXISTS t1 (x INTEGER)")
1666 }
1667
1668 fn add_column_guarded(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
1669 let has_column = tx
1670 .prepare("PRAGMA table_info(t1)")?
1671 .query_map([], |row| row.get::<_, String>(1))?
1672 .collect::<Result<Vec<_>, _>>()?
1673 .iter()
1674 .any(|name| name == "y");
1675 if !has_column {
1676 tx.execute_batch("ALTER TABLE t1 ADD COLUMN y TEXT")?;
1677 }
1678 Ok(())
1679 }
1680
1681 fn initialize_v2(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
1682 create_t1(tx)?;
1683 add_column_guarded(tx)
1684 }
1685
1686 fn initialize_v2_alt(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
1687 tx.execute_batch("CREATE TABLE t1 (x INTEGER, z BLOB)")
1688 }
1689
1690 const RELEASED_V1_OBJECTS: &[SchemaObject] = &[SchemaObject {
1691 kind: SchemaObjectKind::Table,
1692 name: "t1",
1693 }];
1694
1695 fn verify_v1(conn: &Connection) -> Result<(), String> {
1696 verify_released_schema_fingerprint(conn, &DOMAIN_V2, RELEASED_V1_OBJECTS, create_t1)
1697 }
1698
1699 const DOMAIN_V1: SchemaDomain = SchemaDomain {
1700 name: "test-domain",
1701 migrations: &[Migration {
1702 version: 1,
1703 name: "base",
1704 apply: create_t1,
1705 }],
1706 initialize_current: create_t1,
1707 allowed_existing_versions: &[1],
1708 released_predecessors: &[],
1709 owned_objects: &[SchemaObject {
1710 kind: SchemaObjectKind::Table,
1711 name: "t1",
1712 }],
1713 retired_objects: &[],
1714 };
1715
1716 const DOMAIN_V2: SchemaDomain = SchemaDomain {
1717 name: "test-domain",
1718 migrations: &[
1719 Migration {
1720 version: 1,
1721 name: "base",
1722 apply: create_t1,
1723 },
1724 Migration {
1725 version: 2,
1726 name: "add-y",
1727 apply: add_column_guarded,
1728 },
1729 ],
1730 initialize_current: initialize_v2,
1731 allowed_existing_versions: &[1, 2],
1732 released_predecessors: &[SchemaPredecessor {
1733 version: 1,
1734 verify: verify_v1,
1735 }],
1736 owned_objects: &[SchemaObject {
1737 kind: SchemaObjectKind::Table,
1738 name: "t1",
1739 }],
1740 retired_objects: &[],
1741 };
1742
1743 const DOMAIN_V2_ALT_INITIALIZER: SchemaDomain = SchemaDomain {
1744 name: "test-domain",
1745 migrations: &[
1746 Migration {
1747 version: 1,
1748 name: "base",
1749 apply: create_t1,
1750 },
1751 Migration {
1752 version: 2,
1753 name: "alt-current",
1754 apply: add_column_guarded,
1755 },
1756 ],
1757 initialize_current: initialize_v2_alt,
1758 allowed_existing_versions: &[2],
1759 released_predecessors: &[],
1760 owned_objects: &[SchemaObject {
1761 kind: SchemaObjectKind::Table,
1762 name: "t1",
1763 }],
1764 retired_objects: &[],
1765 };
1766
1767 fn temp_conn(dir: &tempfile::TempDir) -> Connection {
1768 open(&dir.path().join("db.sqlite3"), ConnectionProfile::PRIMARY).expect("open")
1769 }
1770
1771 #[test]
1772 fn fresh_file_initializes_current_and_stamps() {
1773 let dir = tempfile::tempdir().expect("tempdir");
1774 let mut conn = temp_conn(&dir);
1775 let report = apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("apply");
1776 assert_eq!(
1777 report,
1778 LedgerReport {
1779 from_version: 0,
1780 to_version: 2
1781 }
1782 );
1783 assert_eq!(domain_version(&conn, "test-domain").expect("read"), Some(2));
1784 conn.execute("INSERT INTO t1 (x, y) VALUES (1, 'a')", [])
1785 .expect("schema converged");
1786 }
1787
1788 #[test]
1789 fn second_open_is_current_noop() {
1790 let dir = tempfile::tempdir().expect("tempdir");
1791 let mut conn = temp_conn(&dir);
1792 apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("first");
1793 let report = apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("second");
1794 assert!(!report.migrated());
1795 }
1796
1797 #[test]
1798 fn current_oracle_cache_binds_initializer_and_manifest_not_only_name_version() {
1799 let first_dir = tempfile::tempdir().expect("first tempdir");
1800 let mut first = temp_conn(&first_dir);
1801 apply_domain_migrations(&mut first, &DOMAIN_V2).expect("first current");
1802
1803 let second_dir = tempfile::tempdir().expect("second tempdir");
1804 let mut second = temp_conn(&second_dir);
1805 apply_domain_migrations(&mut second, &DOMAIN_V2_ALT_INITIALIZER).expect("alt current");
1806 let columns: Vec<String> = second
1807 .prepare("PRAGMA table_info(t1)")
1808 .expect("prepare")
1809 .query_map([], |row| row.get(1))
1810 .expect("query")
1811 .collect::<Result<_, _>>()
1812 .expect("columns");
1813 assert_eq!(columns, vec!["x", "z"]);
1814 }
1815
1816 #[test]
1817 fn current_row_with_partial_catalog_is_refused_before_noop() {
1818 let dir = tempfile::tempdir().expect("tempdir");
1819 let mut conn = temp_conn(&dir);
1820 apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("current");
1821 conn.execute_batch("ALTER TABLE t1 ADD COLUMN candidate_partial TEXT")
1822 .expect("partial candidate mutation");
1823 let err = apply_domain_migrations(&mut conn, &DOMAIN_V2).expect_err("refuse current shape");
1824 assert!(matches!(
1825 err,
1826 SqliteStoreError::SchemaFingerprintMismatch { version: 2, .. }
1827 ));
1828 assert_eq!(
1829 domain_version(&conn, DOMAIN_V2.name).expect("ledger"),
1830 Some(2)
1831 );
1832 }
1833
1834 #[test]
1835 fn upgrade_applies_only_pending() {
1836 let dir = tempfile::tempdir().expect("tempdir");
1837 let mut conn = temp_conn(&dir);
1838 apply_domain_migrations(&mut conn, &DOMAIN_V1).expect("v1");
1839 let report = apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("v2");
1840 assert_eq!(
1841 report,
1842 LedgerReport {
1843 from_version: 1,
1844 to_version: 2
1845 }
1846 );
1847 }
1848
1849 #[test]
1850 fn allowed_version_with_wrong_catalog_is_refused_without_migration() {
1851 let dir = tempfile::tempdir().expect("tempdir");
1852 let mut conn = temp_conn(&dir);
1853 apply_domain_migrations(&mut conn, &DOMAIN_V1).expect("v1");
1854 conn.execute_batch("ALTER TABLE t1 ADD COLUMN candidate_only TEXT")
1855 .expect("candidate shape");
1856 let err = apply_domain_migrations(&mut conn, &DOMAIN_V2).expect_err("refuse fingerprint");
1857 assert!(matches!(
1858 err,
1859 SqliteStoreError::SchemaFingerprintMismatch { version: 1, .. }
1860 ));
1861 assert_eq!(
1862 domain_version(&conn, DOMAIN_V1.name).expect("ledger"),
1863 Some(1),
1864 "fingerprint refusal advanced the ledger"
1865 );
1866 let columns: Vec<String> = conn
1867 .prepare("PRAGMA table_info(t1)")
1868 .expect("prepare")
1869 .query_map([], |row| row.get(1))
1870 .expect("query")
1871 .collect::<Result<_, _>>()
1872 .expect("columns");
1873 assert_eq!(columns, vec!["x", "candidate_only"]);
1874 }
1875
1876 #[test]
1877 fn pre_floor_and_gap_versions_are_refused_without_mutation() {
1878 fn no_op(_tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
1879 Ok(())
1880 }
1881 fn verify_v2(_conn: &Connection) -> Result<(), String> {
1882 Ok(())
1883 }
1884 const DOMAIN_V3_FLOOR_2: SchemaDomain = SchemaDomain {
1885 name: "floor-domain",
1886 migrations: &[
1887 Migration {
1888 version: 1,
1889 name: "base",
1890 apply: no_op,
1891 },
1892 Migration {
1893 version: 2,
1894 name: "released-floor",
1895 apply: no_op,
1896 },
1897 Migration {
1898 version: 3,
1899 name: "current",
1900 apply: no_op,
1901 },
1902 ],
1903 initialize_current: no_op,
1904 allowed_existing_versions: &[2, 3],
1905 released_predecessors: &[SchemaPredecessor {
1906 version: 2,
1907 verify: verify_v2,
1908 }],
1909 owned_objects: &[],
1910 retired_objects: &[],
1911 };
1912 const DOMAIN_V4_GAP_3: SchemaDomain = SchemaDomain {
1913 name: "gap-domain",
1914 migrations: &[
1915 Migration {
1916 version: 1,
1917 name: "old",
1918 apply: no_op,
1919 },
1920 Migration {
1921 version: 2,
1922 name: "released-floor",
1923 apply: no_op,
1924 },
1925 Migration {
1926 version: 3,
1927 name: "unreleased-candidate",
1928 apply: no_op,
1929 },
1930 Migration {
1931 version: 4,
1932 name: "current",
1933 apply: no_op,
1934 },
1935 ],
1936 initialize_current: no_op,
1937 allowed_existing_versions: &[2, 4],
1938 released_predecessors: &[SchemaPredecessor {
1939 version: 2,
1940 verify: verify_v2,
1941 }],
1942 owned_objects: &[],
1943 retired_objects: &[],
1944 };
1945 let dir = tempfile::tempdir().expect("tempdir");
1946 let mut conn = temp_conn(&dir);
1947 conn.execute_batch(CREATE_LEDGER_SQL).expect("ledger");
1948 conn.execute(
1949 "INSERT INTO main.meerkat_schema (domain, version) VALUES (?1, 1)",
1950 [DOMAIN_V3_FLOOR_2.name],
1951 )
1952 .expect("pre-floor row");
1953 let err =
1954 apply_domain_migrations(&mut conn, &DOMAIN_V3_FLOOR_2).expect_err("refuse pre-floor");
1955 assert!(matches!(
1956 err,
1957 SqliteStoreError::UnsupportedSchemaPredecessor { found: 1, .. }
1958 ));
1959 assert_eq!(
1960 domain_version(&conn, DOMAIN_V3_FLOOR_2.name).expect("ledger"),
1961 Some(1)
1962 );
1963
1964 conn.execute(
1965 "INSERT INTO main.meerkat_schema (domain, version) VALUES (?1, 3)",
1966 [DOMAIN_V4_GAP_3.name],
1967 )
1968 .expect("gap row");
1969 let err = apply_domain_migrations(&mut conn, &DOMAIN_V4_GAP_3).expect_err("refuse gap");
1970 assert!(matches!(
1971 err,
1972 SqliteStoreError::UnsupportedSchemaPredecessor { found: 3, .. }
1973 ));
1974 assert_eq!(
1975 domain_version(&conn, DOMAIN_V4_GAP_3.name).expect("ledger"),
1976 Some(3)
1977 );
1978 }
1979
1980 #[test]
1981 fn unledgered_owned_objects_are_refused_without_mutation() {
1982 let dir = tempfile::tempdir().expect("tempdir");
1983 let mut conn = temp_conn(&dir);
1984 conn.execute_batch("CREATE TABLE t1 (x INTEGER)")
1985 .expect("unknown unledgered ddl");
1986 let err = apply_domain_migrations(&mut conn, &DOMAIN_V2).expect_err("refuse");
1987 assert!(matches!(
1988 err,
1989 SqliteStoreError::UnledgeredDomainObjects { .. }
1990 ));
1991 assert!(
1992 !ledger_table_exists(&conn).expect("ledger presence"),
1993 "eligibility refusal must not create the ledger"
1994 );
1995 let columns: Vec<String> = conn
1996 .prepare("PRAGMA table_info(t1)")
1997 .expect("prepare")
1998 .query_map([], |row| row.get(1))
1999 .expect("query")
2000 .collect::<Result<_, _>>()
2001 .expect("columns");
2002 assert_eq!(columns, vec!["x"], "refusal mutated unknown schema");
2003 }
2004
2005 #[test]
2006 fn maintenance_bridge_authenticates_exact_v1_and_migrates_to_v2() {
2007 let dir = tempfile::tempdir().expect("tempdir");
2008 let mut conn = temp_conn(&dir);
2009 conn.execute_batch("CREATE TABLE t1 (x INTEGER)")
2010 .expect("historical unledgered v1");
2011
2012 let report =
2013 bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], None).expect("bridge");
2014 assert_eq!(
2015 report,
2016 MaintenanceBridgeReport {
2017 from_version: 1,
2018 to_version: 2,
2019 prepared: 0,
2020 }
2021 );
2022 assert_eq!(
2023 domain_version(&conn, DOMAIN_V2.name).expect("ledger"),
2024 Some(2)
2025 );
2026 let columns = conn
2027 .prepare("PRAGMA table_info(t1)")
2028 .expect("prepare")
2029 .query_map([], |row| row.get::<_, String>(1))
2030 .expect("query")
2031 .collect::<Result<Vec<_>, _>>()
2032 .expect("columns");
2033 assert_eq!(columns, vec!["x", "y"]);
2034
2035 let second = bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], None)
2036 .expect("idempotent bridge");
2037 assert_eq!(
2038 second,
2039 MaintenanceBridgeReport {
2040 from_version: 2,
2041 to_version: 2,
2042 prepared: 0,
2043 }
2044 );
2045 }
2046
2047 #[test]
2048 fn maintenance_bridge_prepares_and_upgrades_existing_v1_row() {
2049 fn normalize_existing(
2050 tx: &Transaction<'_>,
2051 ) -> Result<MaintenancePrepareReport, rusqlite::Error> {
2052 let changed = tx.execute("UPDATE t1 SET x = x + 1", [])?;
2053 Ok(MaintenancePrepareReport { changed })
2054 }
2055
2056 let dir = tempfile::tempdir().expect("tempdir");
2057 let mut conn = temp_conn(&dir);
2058 apply_domain_migrations(&mut conn, &DOMAIN_V1).expect("ledgered v1");
2059 conn.execute("INSERT INTO t1 (x) VALUES (7)", [])
2060 .expect("historical data");
2061
2062 let report =
2063 bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], Some(normalize_existing))
2064 .expect("bridge ledgered predecessor");
2065 assert_eq!(
2066 report,
2067 MaintenanceBridgeReport {
2068 from_version: 1,
2069 to_version: 2,
2070 prepared: 1,
2071 }
2072 );
2073 assert_eq!(
2074 domain_version(&conn, DOMAIN_V2.name).expect("ledger"),
2075 Some(2)
2076 );
2077 assert_eq!(
2078 conn.query_row("SELECT x FROM t1", [], |row| row.get::<_, i64>(0))
2079 .expect("prepared row"),
2080 8
2081 );
2082 conn.execute("INSERT INTO t1 (x, y) VALUES (9, 'migrated')", [])
2083 .expect("v2 shape");
2084 }
2085
2086 #[test]
2087 fn maintenance_bridge_prepares_existing_target_and_reports_durable_changes() {
2088 fn normalize_target(
2089 tx: &Transaction<'_>,
2090 ) -> Result<MaintenancePrepareReport, rusqlite::Error> {
2091 let changed = tx.execute("UPDATE t1 SET x = x + 1", [])?;
2092 Ok(MaintenancePrepareReport { changed })
2093 }
2094 fn mutate_then_fail(
2095 tx: &Transaction<'_>,
2096 ) -> Result<MaintenancePrepareReport, rusqlite::Error> {
2097 tx.execute("UPDATE t1 SET x = 99", [])?;
2098 tx.execute_batch("THIS IS NOT SQL")?;
2099 Ok(MaintenancePrepareReport { changed: 1 })
2100 }
2101 fn rolls_back_then_begins_and_fails(
2102 tx: &Transaction<'_>,
2103 ) -> Result<MaintenancePrepareReport, rusqlite::Error> {
2104 tx.execute_batch("UPDATE t1 SET x = 99; ROLLBACK; BEGIN; THIS IS NOT SQL")?;
2105 Ok(MaintenancePrepareReport { changed: 1 })
2106 }
2107
2108 let dir = tempfile::tempdir().expect("tempdir");
2109 let mut conn = temp_conn(&dir);
2110 apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("ledgered target");
2111 conn.execute("INSERT INTO t1 (x, y) VALUES (7, 'target')", [])
2112 .expect("target data");
2113
2114 let report =
2115 bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], Some(normalize_target))
2116 .expect("prepare target");
2117 assert_eq!(
2118 report,
2119 MaintenanceBridgeReport {
2120 from_version: 2,
2121 to_version: 2,
2122 prepared: 1,
2123 }
2124 );
2125 assert!(!report.migrated());
2126 assert!(report.changed());
2127 assert_eq!(
2128 conn.query_row("SELECT x FROM t1", [], |row| row.get::<_, i64>(0))
2129 .expect("prepared target row"),
2130 8
2131 );
2132 assert_eq!(
2133 domain_version(&conn, DOMAIN_V2.name).expect("ledger"),
2134 Some(2)
2135 );
2136
2137 let err = bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], Some(mutate_then_fail))
2138 .expect_err("target prepare failure");
2139 assert!(matches!(
2140 err,
2141 SqliteStoreError::MigrationFailed { ref name, .. }
2142 if name == "maintenance-prepare"
2143 ));
2144 assert_eq!(
2145 conn.query_row("SELECT x FROM t1", [], |row| row.get::<_, i64>(0))
2146 .expect("row after rollback"),
2147 8
2148 );
2149
2150 let err = bridge_unledgered_domain(
2151 &mut conn,
2152 &DOMAIN_V2,
2153 2,
2154 &[1],
2155 Some(rolls_back_then_begins_and_fails),
2156 )
2157 .expect_err("target custody loss");
2158 assert!(matches!(
2159 err,
2160 SqliteStoreError::MigrationBrokeTransaction { ref name, .. }
2161 if name == "maintenance-prepare"
2162 ));
2163 assert_eq!(
2164 conn.query_row("SELECT x FROM t1", [], |row| row.get::<_, i64>(0))
2165 .expect("row after custody refusal"),
2166 8
2167 );
2168 assert_eq!(
2169 domain_version(&conn, DOMAIN_V2.name).expect("ledger"),
2170 Some(2)
2171 );
2172 }
2173
2174 #[test]
2175 fn maintenance_bridge_refuses_target_that_only_matches_migration_oracle() {
2176 let dir = tempfile::tempdir().expect("tempdir");
2177 let mut conn = temp_conn(&dir);
2178 conn.execute_batch("CREATE TABLE t1 (x INTEGER)")
2179 .expect("historical unledgered v1");
2180
2181 let err = bridge_unledgered_domain(&mut conn, &DOMAIN_V2_ALT_INITIALIZER, 2, &[1], None)
2182 .expect_err("ordinary target verifier must reject drift");
2183 assert!(matches!(
2184 err,
2185 SqliteStoreError::SchemaFingerprintMismatch { version: 2, .. }
2186 ));
2187 assert!(!ledger_table_exists(&conn).expect("ledger presence"));
2188 let columns = conn
2189 .prepare("PRAGMA table_info(t1)")
2190 .expect("prepare")
2191 .query_map([], |row| row.get::<_, String>(1))
2192 .expect("query")
2193 .collect::<Result<Vec<_>, _>>()
2194 .expect("columns");
2195 assert_eq!(
2196 columns,
2197 vec!["x"],
2198 "target-verifier refusal did not roll back"
2199 );
2200 }
2201
2202 #[test]
2203 fn maintenance_bridge_refuses_catalog_outside_source_allowlist_across_data_only_gap() {
2204 fn no_op(_tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2205 Ok(())
2206 }
2207 const DATA_ONLY_GAP: SchemaDomain = SchemaDomain {
2208 name: "data-only-gap-domain",
2209 migrations: &[
2210 Migration {
2211 version: 1,
2212 name: "base",
2213 apply: create_t1,
2214 },
2215 Migration {
2216 version: 2,
2217 name: "data-only",
2218 apply: no_op,
2219 },
2220 Migration {
2221 version: 3,
2222 name: "add-y",
2223 apply: add_column_guarded,
2224 },
2225 ],
2226 initialize_current: initialize_v2,
2227 allowed_existing_versions: &[3],
2228 released_predecessors: &[],
2229 owned_objects: &[SchemaObject {
2230 kind: SchemaObjectKind::Table,
2231 name: "t1",
2232 }],
2233 retired_objects: &[],
2234 };
2235
2236 let dir = tempfile::tempdir().expect("tempdir");
2237 let mut conn = temp_conn(&dir);
2238 conn.execute_batch("CREATE TABLE t1 (x INTEGER)")
2239 .expect("v1-or-v2 catalog");
2240 let err = bridge_unledgered_domain(&mut conn, &DATA_ONLY_GAP, 3, &[3], None)
2241 .expect_err("excluded historical prefixes must not be inferred");
2242 assert!(matches!(
2243 err,
2244 SqliteStoreError::UnledgeredSchemaNoMatch { .. }
2245 ));
2246 assert!(!ledger_table_exists(&conn).expect("ledger presence"));
2247 let columns = conn
2248 .prepare("PRAGMA table_info(t1)")
2249 .expect("prepare")
2250 .query_map([], |row| row.get::<_, String>(1))
2251 .expect("query")
2252 .collect::<Result<Vec<_>, _>>()
2253 .expect("columns");
2254 assert_eq!(columns, vec!["x"]);
2255
2256 let err = bridge_unledgered_domain(&mut conn, &DATA_ONLY_GAP, 3, &[2, 1], None)
2257 .expect_err("unordered source authority must be refused");
2258 assert!(matches!(err, SqliteStoreError::InvalidMigrationList { .. }));
2259 }
2260
2261 #[test]
2262 fn maintenance_bridge_leaves_fresh_domain_for_normal_initializer() {
2263 let dir = tempfile::tempdir().expect("tempdir");
2264 let mut conn = temp_conn(&dir);
2265 let report =
2266 bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], None).expect("fresh no-op");
2267 assert_eq!(
2268 report,
2269 MaintenanceBridgeReport {
2270 from_version: 0,
2271 to_version: 0,
2272 prepared: 0,
2273 }
2274 );
2275 assert!(!ledger_table_exists(&conn).expect("ledger presence"));
2276 assert!(
2277 find_owned_objects(&conn, &DOMAIN_V2)
2278 .expect("objects")
2279 .is_empty()
2280 );
2281 }
2282
2283 #[test]
2284 fn maintenance_bridge_refuses_malformed_historical_shape_without_mutation() {
2285 let dir = tempfile::tempdir().expect("tempdir");
2286 let mut conn = temp_conn(&dir);
2287 conn.execute_batch("CREATE TABLE t1 (x INTEGER, candidate_only BLOB)")
2288 .expect("candidate schema");
2289
2290 let err = bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], None)
2291 .expect_err("refuse unauthenticated shape");
2292 assert!(matches!(
2293 err,
2294 SqliteStoreError::UnledgeredSchemaNoMatch {
2295 target_version: 2,
2296 ..
2297 }
2298 ));
2299 assert!(!ledger_table_exists(&conn).expect("ledger presence"));
2300 let columns = conn
2301 .prepare("PRAGMA table_info(t1)")
2302 .expect("prepare")
2303 .query_map([], |row| row.get::<_, String>(1))
2304 .expect("query")
2305 .collect::<Result<Vec<_>, _>>()
2306 .expect("columns");
2307 assert_eq!(columns, vec!["x", "candidate_only"]);
2308 }
2309
2310 #[test]
2311 fn maintenance_bridge_fingerprint_preserves_sql_literal_whitespace() {
2312 fn create_exact(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2313 tx.execute_batch("CREATE TABLE literal_t (x TEXT CHECK(x <> 'a b'))")
2314 }
2315 const LITERAL_DOMAIN: SchemaDomain = SchemaDomain {
2316 name: "literal-fingerprint-domain",
2317 migrations: &[Migration {
2318 version: 1,
2319 name: "base",
2320 apply: create_exact,
2321 }],
2322 initialize_current: create_exact,
2323 allowed_existing_versions: &[1],
2324 released_predecessors: &[],
2325 owned_objects: &[SchemaObject {
2326 kind: SchemaObjectKind::Table,
2327 name: "literal_t",
2328 }],
2329 retired_objects: &[],
2330 };
2331
2332 let dir = tempfile::tempdir().expect("tempdir");
2333 let mut conn = temp_conn(&dir);
2334 conn.execute_batch("CREATE TABLE literal_t (x TEXT CHECK(x <> 'a b'))")
2335 .expect("semantic mismatch");
2336 let err = bridge_unledgered_domain(&mut conn, &LITERAL_DOMAIN, 1, &[1], None)
2337 .expect_err("literal whitespace must remain fingerprint-significant");
2338 assert!(matches!(
2339 err,
2340 SqliteStoreError::UnledgeredSchemaNoMatch { .. }
2341 ));
2342 assert!(!ledger_table_exists(&conn).expect("ledger presence"));
2343
2344 assert_eq!(
2345 normalize_schema_sql("CREATE TABLE IF NOT EXISTS t (x CHECK(x <> 'IF NOT EXISTS'))"),
2346 "CREATE TABLE t (x CHECK(x <> 'IF NOT EXISTS'))"
2347 );
2348 }
2349
2350 #[test]
2351 fn maintenance_bridge_rolls_back_prepare_failure() {
2352 fn mutate_then_fail(
2353 tx: &Transaction<'_>,
2354 ) -> Result<MaintenancePrepareReport, rusqlite::Error> {
2355 tx.execute("UPDATE t1 SET x = 99", [])?;
2356 tx.execute_batch("THIS IS NOT SQL")?;
2357 Ok(MaintenancePrepareReport { changed: 1 })
2358 }
2359
2360 let dir = tempfile::tempdir().expect("tempdir");
2361 let mut conn = temp_conn(&dir);
2362 conn.execute_batch("CREATE TABLE t1 (x INTEGER); INSERT INTO t1 VALUES (7)")
2363 .expect("historical unledgered v1");
2364
2365 let err = bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], Some(mutate_then_fail))
2366 .expect_err("prepare must fail");
2367 assert!(matches!(
2368 err,
2369 SqliteStoreError::MigrationFailed { ref name, .. }
2370 if name == "maintenance-prepare"
2371 ));
2372 assert_eq!(
2373 conn.query_row("SELECT x FROM t1", [], |row| row.get::<_, i64>(0))
2374 .expect("original row"),
2375 7
2376 );
2377 assert!(!ledger_table_exists(&conn).expect("ledger presence"));
2378 let columns = conn
2379 .prepare("PRAGMA table_info(t1)")
2380 .expect("prepare")
2381 .query_map([], |row| row.get::<_, String>(1))
2382 .expect("query")
2383 .collect::<Result<Vec<_>, _>>()
2384 .expect("columns");
2385 assert_eq!(columns, vec!["x"]);
2386 }
2387
2388 #[test]
2389 fn maintenance_bridge_reports_custody_loss_before_callback_error() {
2390 fn commits_then_fails(
2391 tx: &Transaction<'_>,
2392 ) -> Result<MaintenancePrepareReport, rusqlite::Error> {
2393 tx.execute_batch("UPDATE t1 SET x = 99; COMMIT; THIS IS NOT SQL")?;
2394 Ok(MaintenancePrepareReport { changed: 1 })
2395 }
2396
2397 let dir = tempfile::tempdir().expect("tempdir");
2398 let mut conn = temp_conn(&dir);
2399 conn.execute_batch("CREATE TABLE t1 (x INTEGER); INSERT INTO t1 VALUES (7)")
2400 .expect("historical unledgered v1");
2401 let err =
2402 bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], Some(commits_then_fails))
2403 .expect_err("custody must dominate callback error");
2404 assert!(matches!(
2405 err,
2406 SqliteStoreError::MigrationBrokeTransaction { ref name, .. }
2407 if name == "maintenance-prepare"
2408 ));
2409 assert_eq!(domain_version(&conn, DOMAIN_V2.name).expect("ledger"), None);
2410 }
2411
2412 #[test]
2413 fn maintenance_bridge_refuses_undeclared_trigger_on_owned_table_before_prepare() {
2414 fn prepare_successor(
2415 tx: &Transaction<'_>,
2416 ) -> Result<MaintenancePrepareReport, rusqlite::Error> {
2417 let changed = tx.execute("UPDATE t1 SET x = 8 WHERE x = 7", [])?;
2418 Ok(MaintenancePrepareReport { changed })
2419 }
2420
2421 let dir = tempfile::tempdir().expect("tempdir");
2422 let mut conn = temp_conn(&dir);
2423 conn.execute_batch(
2424 "CREATE TABLE t1 (x INTEGER);
2425 INSERT INTO t1 VALUES (7);
2426 CREATE TRIGGER replace_prepared_successor
2427 AFTER UPDATE ON t1
2428 BEGIN
2429 UPDATE t1 SET x = 99 WHERE rowid = NEW.rowid;
2430 END",
2431 )
2432 .expect("intercepting trigger");
2433
2434 let err = bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], Some(prepare_successor))
2435 .expect_err("undeclared trigger must be refused before prepare");
2436 assert!(matches!(
2437 err,
2438 SqliteStoreError::SchemaFingerprintMismatch { version: 1, .. }
2439 ));
2440 assert_eq!(
2441 conn.query_row("SELECT x FROM t1", [], |row| row.get::<_, i64>(0))
2442 .expect("original row"),
2443 7
2444 );
2445 assert!(!ledger_table_exists(&conn).expect("ledger presence"));
2446 }
2447
2448 #[test]
2449 fn maintenance_bridge_refuses_undeclared_instead_of_trigger_on_owned_view_before_prepare() {
2450 fn create_owned_view(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2451 tx.execute_batch(
2452 "CREATE TABLE view_base (x INTEGER);
2453 CREATE VIEW owned_view AS SELECT x FROM view_base",
2454 )
2455 }
2456 fn prepare_view(tx: &Transaction<'_>) -> Result<MaintenancePrepareReport, rusqlite::Error> {
2457 let changed = tx.execute("UPDATE owned_view SET x = 8 WHERE x = 7", [])?;
2458 Ok(MaintenancePrepareReport { changed })
2459 }
2460 const VIEW_DOMAIN: SchemaDomain = SchemaDomain {
2461 name: "view-bridge-domain",
2462 migrations: &[Migration {
2463 version: 1,
2464 name: "base",
2465 apply: create_owned_view,
2466 }],
2467 initialize_current: create_owned_view,
2468 allowed_existing_versions: &[1],
2469 released_predecessors: &[],
2470 owned_objects: &[
2471 SchemaObject {
2472 kind: SchemaObjectKind::Table,
2473 name: "view_base",
2474 },
2475 SchemaObject {
2476 kind: SchemaObjectKind::View,
2477 name: "owned_view",
2478 },
2479 ],
2480 retired_objects: &[],
2481 };
2482
2483 let dir = tempfile::tempdir().expect("tempdir");
2484 let mut conn = temp_conn(&dir);
2485 conn.execute_batch(
2486 "CREATE TABLE view_base (x INTEGER);
2487 CREATE VIEW owned_view AS SELECT x FROM view_base;
2488 INSERT INTO view_base VALUES (7);
2489 CREATE TRIGGER replace_view_update
2490 INSTEAD OF UPDATE ON owned_view
2491 BEGIN
2492 UPDATE view_base SET x = 99 WHERE x = OLD.x;
2493 END",
2494 )
2495 .expect("intercepting view trigger");
2496
2497 let err = bridge_unledgered_domain(&mut conn, &VIEW_DOMAIN, 1, &[1], Some(prepare_view))
2498 .expect_err("undeclared view trigger must be refused before prepare");
2499 assert!(matches!(
2500 err,
2501 SqliteStoreError::SchemaFingerprintMismatch { version: 1, .. }
2502 ));
2503 assert_eq!(
2504 conn.query_row("SELECT x FROM view_base", [], |row| row.get::<_, i64>(0))
2505 .expect("original row"),
2506 7
2507 );
2508 assert!(!ledger_table_exists(&conn).expect("ledger presence"));
2509 }
2510
2511 #[test]
2512 fn maintenance_bridge_refuses_ambiguous_prefix_without_mutation() {
2513 fn no_op(_tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2514 Ok(())
2515 }
2516 const AMBIGUOUS: SchemaDomain = SchemaDomain {
2517 name: "ambiguous-bridge-domain",
2518 migrations: &[
2519 Migration {
2520 version: 1,
2521 name: "base",
2522 apply: create_t1,
2523 },
2524 Migration {
2525 version: 2,
2526 name: "data-only",
2527 apply: no_op,
2528 },
2529 ],
2530 initialize_current: create_t1,
2531 allowed_existing_versions: &[2],
2532 released_predecessors: &[],
2533 owned_objects: &[SchemaObject {
2534 kind: SchemaObjectKind::Table,
2535 name: "t1",
2536 }],
2537 retired_objects: &[],
2538 };
2539
2540 let dir = tempfile::tempdir().expect("tempdir");
2541 let mut conn = temp_conn(&dir);
2542 conn.execute_batch("CREATE TABLE t1 (x INTEGER)")
2543 .expect("ambiguous schema");
2544 let err = bridge_unledgered_domain(&mut conn, &AMBIGUOUS, 2, &[1, 2], None)
2545 .expect_err("refuse ambiguity");
2546 assert!(matches!(
2547 err,
2548 SqliteStoreError::UnledgeredSchemaAmbiguous { ref matches, .. }
2549 if matches == &[1, 2]
2550 ));
2551 assert!(!ledger_table_exists(&conn).expect("ledger presence"));
2552 }
2553
2554 #[test]
2555 fn maintenance_bridge_refuses_malformed_ledger_before_owned_schema_contact() {
2556 let dir = tempfile::tempdir().expect("tempdir");
2557 let mut conn = temp_conn(&dir);
2558 conn.execute_batch(
2559 "CREATE TABLE t1 (x INTEGER);
2560 CREATE TABLE meerkat_schema (domain TEXT, version INTEGER)",
2561 )
2562 .expect("malformed ledger");
2563 let err = bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], None)
2564 .expect_err("refuse malformed ledger");
2565 assert!(matches!(err, SqliteStoreError::LedgerMalformed { .. }));
2566 let columns = conn
2567 .prepare("PRAGMA table_info(t1)")
2568 .expect("prepare")
2569 .query_map([], |row| row.get::<_, String>(1))
2570 .expect("query")
2571 .collect::<Result<Vec<_>, _>>()
2572 .expect("columns");
2573 assert_eq!(columns, vec!["x"]);
2574 }
2575
2576 #[test]
2577 fn maintenance_bridge_refuses_mixed_case_trigger_that_mutates_foreign_ledger_row() {
2578 let dir = tempfile::tempdir().expect("tempdir");
2579 let mut conn = temp_conn(&dir);
2580 conn.execute_batch(
2581 "CREATE TABLE t1 (x INTEGER);
2582 CREATE TABLE meerkat_schema (
2583 domain TEXT PRIMARY KEY,
2584 version INTEGER NOT NULL
2585 );
2586 INSERT INTO meerkat_schema (domain, version) VALUES ('foreign-domain', 7);
2587 CREATE TRIGGER mutate_foreign_schema_row
2588 AFTER INSERT ON MEERKAT_SCHEMA
2589 BEGIN
2590 UPDATE MEERKAT_SCHEMA
2591 SET version = 999
2592 WHERE domain = 'foreign-domain';
2593 END",
2594 )
2595 .expect("hostile ledger trigger");
2596
2597 let err = bridge_unledgered_domain(&mut conn, &DOMAIN_V2, 2, &[1], None)
2598 .expect_err("mixed-case ledger trigger must be refused");
2599 assert!(matches!(err, SqliteStoreError::LedgerMalformed { .. }));
2600 let row_count = conn
2601 .query_row(
2602 "SELECT COUNT(*) FROM main.meerkat_schema WHERE domain = ?1",
2603 [DOMAIN_V2.name],
2604 |row| row.get::<_, i64>(0),
2605 )
2606 .expect("raw ledger count");
2607 assert_eq!(row_count, 0);
2608 let foreign_version = conn
2609 .query_row(
2610 "SELECT version FROM main.meerkat_schema WHERE domain = 'foreign-domain'",
2611 [],
2612 |row| row.get::<_, i64>(0),
2613 )
2614 .expect("foreign ledger row");
2615 assert_eq!(foreign_version, 7);
2616 let columns = conn
2617 .prepare("PRAGMA table_info(t1)")
2618 .expect("prepare")
2619 .query_map([], |row| row.get::<_, String>(1))
2620 .expect("query")
2621 .collect::<Result<Vec<_>, _>>()
2622 .expect("columns");
2623 assert_eq!(
2624 columns,
2625 vec!["x"],
2626 "failed stamp did not roll back migration"
2627 );
2628 }
2629
2630 #[test]
2631 fn fresh_domain_ignores_foreign_cotenant_objects() {
2632 let dir = tempfile::tempdir().expect("tempdir");
2633 let mut conn = temp_conn(&dir);
2634 conn.execute_batch("CREATE TABLE foreign_table (value TEXT)")
2635 .expect("foreign ddl");
2636 let report = apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("fresh domain");
2637 assert_eq!(
2638 report,
2639 LedgerReport {
2640 from_version: 0,
2641 to_version: 2
2642 }
2643 );
2644 conn.execute("INSERT INTO foreign_table VALUES ('kept')", [])
2645 .expect("foreign object survives");
2646 }
2647
2648 #[test]
2649 fn future_version_is_refused_before_any_mutation() {
2650 let dir = tempfile::tempdir().expect("tempdir");
2651 let mut conn = temp_conn(&dir);
2652 apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("stamp v2");
2653 let err = apply_domain_migrations(&mut conn, &DOMAIN_V1).expect_err("refuse");
2655 match err {
2656 SqliteStoreError::SchemaFromTheFuture {
2657 domain,
2658 found,
2659 supported,
2660 } => {
2661 assert_eq!(domain, "test-domain");
2662 assert_eq!(found, 2);
2663 assert_eq!(supported, 1);
2664 }
2665 other => panic!("wrong error: {other}"),
2666 }
2667 assert_eq!(domain_version(&conn, "test-domain").expect("read"), Some(2));
2669 }
2670
2671 #[test]
2672 fn foreign_domain_rows_are_untouched() {
2673 let dir = tempfile::tempdir().expect("tempdir");
2674 let mut conn = temp_conn(&dir);
2675 apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("mine");
2676 conn.execute(
2677 "INSERT INTO meerkat_schema (domain, version) VALUES ('foreign-domain', 7)",
2678 [],
2679 )
2680 .expect("foreign row");
2681 apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("noop");
2682 let foreign: i64 = conn
2683 .query_row(
2684 "SELECT version FROM meerkat_schema WHERE domain = 'foreign-domain'",
2685 [],
2686 |r| r.get(0),
2687 )
2688 .expect("foreign row survives");
2689 assert_eq!(foreign, 7);
2690 }
2691
2692 #[test]
2693 fn invalid_migration_list_is_refused_without_touching_the_file() {
2694 const BAD: SchemaDomain = SchemaDomain {
2695 name: "bad-domain",
2696 migrations: &[Migration {
2697 version: 3,
2698 name: "gap",
2699 apply: create_t1,
2700 }],
2701 initialize_current: create_t1,
2702 allowed_existing_versions: &[3],
2703 released_predecessors: &[],
2704 owned_objects: &[],
2705 retired_objects: &[],
2706 };
2707 let dir = tempfile::tempdir().expect("tempdir");
2708 let mut conn = temp_conn(&dir);
2709 let err = apply_domain_migrations(&mut conn, &BAD).expect_err("refuse");
2710 assert!(matches!(err, SqliteStoreError::InvalidMigrationList { .. }));
2711 assert!(!ledger_table_exists(&conn).expect("check"));
2712 }
2713
2714 #[test]
2715 fn failed_migration_rolls_back_atomically() {
2716 fn fail(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2717 tx.execute_batch("CREATE TABLE half_done (x INTEGER)")?;
2718 tx.execute_batch("THIS IS NOT SQL")
2719 }
2720 fn initialize_failing(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2721 create_t1(tx)?;
2722 fail(tx)
2723 }
2724 const FAILING: SchemaDomain = SchemaDomain {
2725 name: "failing-domain",
2726 migrations: &[
2727 Migration {
2728 version: 1,
2729 name: "base",
2730 apply: create_t1,
2731 },
2732 Migration {
2733 version: 2,
2734 name: "explodes",
2735 apply: fail,
2736 },
2737 ],
2738 initialize_current: initialize_failing,
2739 allowed_existing_versions: &[1, 2],
2740 released_predecessors: &[SchemaPredecessor {
2741 version: 1,
2742 verify: verify_v1,
2743 }],
2744 owned_objects: &[
2745 SchemaObject {
2746 kind: SchemaObjectKind::Table,
2747 name: "t1",
2748 },
2749 SchemaObject {
2750 kind: SchemaObjectKind::Table,
2751 name: "half_done",
2752 },
2753 ],
2754 retired_objects: &[],
2755 };
2756 let dir = tempfile::tempdir().expect("tempdir");
2757 let mut conn = temp_conn(&dir);
2758 let err = apply_domain_migrations(&mut conn, &FAILING).expect_err("must fail");
2759 assert!(matches!(
2760 err,
2761 SqliteStoreError::MigrationFailed { version: 2, .. }
2762 ));
2763 assert_eq!(domain_version(&conn, "failing-domain").expect("read"), None);
2766 let tables: Vec<String> = conn
2767 .prepare(
2768 "SELECT name FROM sqlite_master WHERE type='table' AND name IN ('t1','half_done')",
2769 )
2770 .expect("prepare")
2771 .query_map([], |r| r.get(0))
2772 .expect("query")
2773 .collect::<Result<_, _>>()
2774 .expect("rows");
2775 assert!(tables.is_empty(), "rollback left tables behind: {tables:?}");
2776 }
2777
2778 #[test]
2779 fn malformed_ledger_shape_is_refused_not_healed() {
2780 let dir = tempfile::tempdir().expect("tempdir");
2781 let mut conn = temp_conn(&dir);
2782 conn.execute_batch("CREATE TABLE meerkat_schema (x INTEGER)")
2784 .expect("foreign ddl");
2785 let err = domain_version(&conn, "test-domain").expect_err("refuse read");
2786 assert!(matches!(err, SqliteStoreError::LedgerMalformed { .. }));
2787 let err = apply_domain_migrations(&mut conn, &DOMAIN_V1).expect_err("refuse migrate");
2788 assert!(matches!(err, SqliteStoreError::LedgerMalformed { .. }));
2789 let count: i64 = conn
2791 .query_row("SELECT COUNT(*) FROM meerkat_schema", [], |r| r.get(0))
2792 .expect("foreign table survives");
2793 assert_eq!(count, 0);
2794 }
2795
2796 #[test]
2797 fn non_positive_versions_are_refused_not_healed() {
2798 for bad_version in [0i64, -3] {
2799 let dir = tempfile::tempdir().expect("tempdir");
2800 let mut conn = temp_conn(&dir);
2801 conn.execute_batch(
2802 "CREATE TABLE meerkat_schema (domain TEXT PRIMARY KEY, version INTEGER NOT NULL)",
2803 )
2804 .expect("ledger ddl");
2805 conn.execute(
2806 "INSERT INTO meerkat_schema (domain, version) VALUES ('test-domain', ?1)",
2807 [bad_version],
2808 )
2809 .expect("seed bad version");
2810 let err = domain_version(&conn, "test-domain").expect_err("refuse read");
2811 assert!(matches!(err, SqliteStoreError::LedgerMalformed { .. }));
2812 let err = apply_domain_migrations(&mut conn, &DOMAIN_V1).expect_err("refuse migrate");
2813 assert!(matches!(err, SqliteStoreError::LedgerMalformed { .. }));
2814 let stored: i64 = conn
2816 .query_row(
2817 "SELECT version FROM meerkat_schema WHERE domain = 'test-domain'",
2818 [],
2819 |r| r.get(0),
2820 )
2821 .expect("row survives");
2822 assert_eq!(stored, bad_version);
2823 }
2824 }
2825
2826 #[test]
2827 fn duplicate_domain_rows_are_refused() {
2828 let dir = tempfile::tempdir().expect("tempdir");
2829 let conn = temp_conn(&dir);
2830 conn.execute_batch(
2834 "CREATE TABLE meerkat_schema (domain TEXT, version INTEGER NOT NULL);
2835 INSERT INTO meerkat_schema VALUES ('dup-domain', 1);
2836 INSERT INTO meerkat_schema VALUES ('dup-domain', 2);",
2837 )
2838 .expect("seed duplicates");
2839 let err = read_version(&conn, "dup-domain").expect_err("refuse duplicates");
2840 match err {
2841 SqliteStoreError::LedgerMalformed { detail } => {
2842 assert!(detail.contains("multiple ledger rows"), "{detail}");
2843 }
2844 other => panic!("wrong error: {other}"),
2845 }
2846 }
2847
2848 #[test]
2849 fn temp_shadowing_cannot_hijack_the_ledger() {
2850 let dir = tempfile::tempdir().expect("tempdir");
2851 let mut conn = temp_conn(&dir);
2852 apply_domain_migrations(&mut conn, &DOMAIN_V1).expect("stamp v1");
2853 conn.execute_batch(
2856 "CREATE TEMP TABLE meerkat_schema (domain TEXT PRIMARY KEY, version INTEGER NOT NULL);
2857 INSERT INTO temp.meerkat_schema VALUES ('test-domain', 999);",
2858 )
2859 .expect("temp shadow");
2860 assert_eq!(domain_version(&conn, "test-domain").expect("read"), Some(1));
2861 let report = apply_domain_migrations(&mut conn, &DOMAIN_V1).expect("noop against main");
2862 assert!(!report.migrated());
2863 }
2864
2865 #[test]
2866 fn migration_that_ends_the_transaction_is_refused_unstamped() {
2867 fn no_op(_tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2868 Ok(())
2869 }
2870 fn verify_empty_predecessor(_conn: &Connection) -> Result<(), String> {
2871 Ok(())
2872 }
2873 fn commits_underneath(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2874 tx.execute_batch("CREATE TABLE escaped_commit (x INTEGER); COMMIT")
2875 }
2876 fn rolls_back_underneath(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2877 tx.execute_batch("ROLLBACK")
2878 }
2879 fn commits_then_begins(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2883 tx.execute_batch("CREATE TABLE escaped_commit_begin (x INTEGER); COMMIT; BEGIN")
2884 }
2885 fn rolls_back_then_begins(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2886 tx.execute_batch("ROLLBACK; BEGIN")
2887 }
2888 const COMMITS: SchemaDomain = SchemaDomain {
2889 name: "custody-commit",
2890 migrations: &[
2891 Migration {
2892 version: 1,
2893 name: "base",
2894 apply: no_op,
2895 },
2896 Migration {
2897 version: 2,
2898 name: "commits-underneath",
2899 apply: commits_underneath,
2900 },
2901 ],
2902 initialize_current: no_op,
2903 allowed_existing_versions: &[1, 2],
2904 released_predecessors: &[SchemaPredecessor {
2905 version: 1,
2906 verify: verify_empty_predecessor,
2907 }],
2908 owned_objects: &[],
2909 retired_objects: &[],
2910 };
2911 const ROLLS_BACK: SchemaDomain = SchemaDomain {
2912 name: "custody-rollback",
2913 migrations: &[
2914 Migration {
2915 version: 1,
2916 name: "base",
2917 apply: no_op,
2918 },
2919 Migration {
2920 version: 2,
2921 name: "rolls-back-underneath",
2922 apply: rolls_back_underneath,
2923 },
2924 ],
2925 initialize_current: no_op,
2926 allowed_existing_versions: &[1, 2],
2927 released_predecessors: &[SchemaPredecessor {
2928 version: 1,
2929 verify: verify_empty_predecessor,
2930 }],
2931 owned_objects: &[],
2932 retired_objects: &[],
2933 };
2934 const COMMITS_THEN_BEGINS: SchemaDomain = SchemaDomain {
2935 name: "custody-commit-begin",
2936 migrations: &[
2937 Migration {
2938 version: 1,
2939 name: "base",
2940 apply: no_op,
2941 },
2942 Migration {
2943 version: 2,
2944 name: "commits-then-begins",
2945 apply: commits_then_begins,
2946 },
2947 ],
2948 initialize_current: no_op,
2949 allowed_existing_versions: &[1, 2],
2950 released_predecessors: &[SchemaPredecessor {
2951 version: 1,
2952 verify: verify_empty_predecessor,
2953 }],
2954 owned_objects: &[],
2955 retired_objects: &[],
2956 };
2957 const ROLLS_BACK_THEN_BEGINS: SchemaDomain = SchemaDomain {
2958 name: "custody-rollback-begin",
2959 migrations: &[
2960 Migration {
2961 version: 1,
2962 name: "base",
2963 apply: no_op,
2964 },
2965 Migration {
2966 version: 2,
2967 name: "rolls-back-then-begins",
2968 apply: rolls_back_then_begins,
2969 },
2970 ],
2971 initialize_current: no_op,
2972 allowed_existing_versions: &[1, 2],
2973 released_predecessors: &[SchemaPredecessor {
2974 version: 1,
2975 verify: verify_empty_predecessor,
2976 }],
2977 owned_objects: &[],
2978 retired_objects: &[],
2979 };
2980 for (domain, expected_name) in [
2981 (&COMMITS, "commits-underneath"),
2982 (&ROLLS_BACK, "rolls-back-underneath"),
2983 (&COMMITS_THEN_BEGINS, "commits-then-begins"),
2984 (&ROLLS_BACK_THEN_BEGINS, "rolls-back-then-begins"),
2985 ] {
2986 let dir = tempfile::tempdir().expect("tempdir");
2987 let mut conn = temp_conn(&dir);
2988 conn.execute_batch(CREATE_LEDGER_SQL).expect("ledger");
2989 conn.execute(
2990 "INSERT INTO main.meerkat_schema (domain, version) VALUES (?1, 1)",
2991 [domain.name],
2992 )
2993 .expect("released predecessor");
2994 let err = apply_domain_migrations(&mut conn, domain).expect_err("custody violation");
2995 match err {
2996 SqliteStoreError::MigrationBrokeTransaction {
2997 domain: err_domain,
2998 version,
2999 name,
3000 } => {
3001 assert_eq!(err_domain, domain.name);
3002 assert_eq!(version, 2);
3003 assert_eq!(name, expected_name);
3004 }
3005 other => panic!("wrong error: {other}"),
3006 }
3007 assert_eq!(domain_version(&conn, domain.name).expect("read"), Some(1));
3010 }
3011 }
3012
3013 #[test]
3014 fn initializer_that_ends_the_transaction_is_refused_unstamped() {
3015 fn commits_underneath(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
3016 tx.execute_batch("CREATE TABLE escaped_initializer (x INTEGER); COMMIT")
3017 }
3018 const COMMITS: SchemaDomain = SchemaDomain {
3019 name: "initializer-custody-commit",
3020 migrations: &[Migration {
3021 version: 1,
3022 name: "base",
3023 apply: commits_underneath,
3024 }],
3025 initialize_current: commits_underneath,
3026 allowed_existing_versions: &[1],
3027 released_predecessors: &[],
3028 owned_objects: &[SchemaObject {
3029 kind: SchemaObjectKind::Table,
3030 name: "escaped_initializer",
3031 }],
3032 retired_objects: &[],
3033 };
3034 let dir = tempfile::tempdir().expect("tempdir");
3035 let mut conn = temp_conn(&dir);
3036 let err = apply_domain_migrations(&mut conn, &COMMITS).expect_err("custody violation");
3037 match err {
3038 SqliteStoreError::MigrationBrokeTransaction {
3039 domain,
3040 version,
3041 name,
3042 } => {
3043 assert_eq!(domain, COMMITS.name);
3044 assert_eq!(version, 1);
3045 assert_eq!(name, "initialize-current");
3046 }
3047 other => panic!("wrong error: {other}"),
3048 }
3049 assert_eq!(domain_version(&conn, COMMITS.name).expect("read"), None);
3050 }
3051
3052 #[test]
3053 fn initializer_using_its_own_savepoints_keeps_custody() {
3054 fn nests_savepoints(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
3058 tx.execute_batch(
3059 "SAVEPOINT body_sp;
3060 CREATE TABLE sp_t (x INTEGER);
3061 RELEASE SAVEPOINT body_sp",
3062 )
3063 }
3064 const NESTED: SchemaDomain = SchemaDomain {
3065 name: "custody-nested-savepoint",
3066 migrations: &[Migration {
3067 version: 1,
3068 name: "nests-savepoints",
3069 apply: nests_savepoints,
3070 }],
3071 initialize_current: nests_savepoints,
3072 allowed_existing_versions: &[1],
3073 released_predecessors: &[],
3074 owned_objects: &[SchemaObject {
3075 kind: SchemaObjectKind::Table,
3076 name: "sp_t",
3077 }],
3078 retired_objects: &[],
3079 };
3080 let dir = tempfile::tempdir().expect("tempdir");
3081 let mut conn = temp_conn(&dir);
3082 let report = apply_domain_migrations(&mut conn, &NESTED).expect("apply");
3083 assert_eq!(report.to_version, 1);
3084 assert_eq!(domain_version(&conn, NESTED.name).expect("read"), Some(1));
3085 }
3086
3087 #[test]
3088 fn schema_preflight_passes_fresh_and_current_refuses_future() {
3089 let dir = tempfile::tempdir().expect("tempdir");
3090 let mut conn = temp_conn(&dir);
3091 preflight_schema_eligibility(&conn, &DOMAIN_V1).expect("no ledger yet");
3092 apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("stamp v2");
3093 preflight_schema_eligibility(&conn, &DOMAIN_V2).expect("current");
3094 let err =
3095 preflight_schema_eligibility(&conn, &DOMAIN_V1).expect_err("future for old binary");
3096 assert!(matches!(
3097 err,
3098 SqliteStoreError::SchemaFromTheFuture {
3099 found: 2,
3100 supported: 1,
3101 ..
3102 }
3103 ));
3104 }
3105
3106 #[test]
3107 fn concurrent_opens_race_safely() {
3108 let dir = tempfile::tempdir().expect("tempdir");
3109 let path = dir.path().join("db.sqlite3");
3110 let mut handles = Vec::new();
3111 for _ in 0..8 {
3112 let path = path.clone();
3113 handles.push(std::thread::spawn(move || {
3114 let mut conn = open(&path, ConnectionProfile::PRIMARY).expect("open");
3115 apply_domain_migrations(&mut conn, &DOMAIN_V2).expect("apply")
3116 }));
3117 }
3118 let mut migrated = 0;
3119 for handle in handles {
3120 let report = handle.join().expect("thread");
3121 assert_eq!(report.to_version, 2);
3122 if report.migrated() {
3123 migrated += 1;
3124 }
3125 }
3126 assert!(migrated >= 1, "someone must have migrated");
3127 let conn = open(&path, ConnectionProfile::ReadOnly).expect("reopen");
3128 assert_eq!(domain_version(&conn, "test-domain").expect("read"), Some(2));
3129 }
3130}