lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

migration.rs (129474B)


      1 //! Deterministic migration identity, ledger, and governed execution mechanics.
      2 
      3 use core::fmt;
      4 use std::{collections::BTreeSet, error::Error, future::Future, pin::Pin};
      5 
      6 #[cfg(any(target_os = "linux", target_os = "macos"))]
      7 use std::collections::BTreeMap;
      8 
      9 use sha2::{Digest, Sha256};
     10 
     11 use sqlx::SqliteConnection;
     12 #[cfg(any(target_os = "linux", target_os = "macos"))]
     13 use sqlx::{Connection, Row};
     14 
     15 use crate::ServiceSqliteError;
     16 #[cfg(any(target_os = "linux", target_os = "macos"))]
     17 use crate::{SchemaCatalog, ServiceSqliteErrorKind};
     18 
     19 const MIGRATION_CONTENT_DOMAIN: &[u8] = b"radroots.service_sqlite.migration_content.v1\0";
     20 const MIGRATION_CATALOG_DOMAIN: &[u8] = b"radroots.service_sqlite.migration_catalog.v1\0";
     21 const BASE_SCHEMA_VERSION: u32 = 1;
     22 const MAX_MIGRATION_NAME_UTF8_BYTES: usize = 128;
     23 const MAX_MIGRATION_CONTENT_BYTES: usize = 1024 * 1024;
     24 const MAX_MIGRATION_COUNT: usize = 4096;
     25 const MAX_MIGRATION_BUILD_ID_UTF8_BYTES: usize = 128;
     26 
     27 /// A bounded stable lower-snake migration name.
     28 #[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
     29 pub struct MigrationName(&'static str);
     30 
     31 impl MigrationName {
     32     /// Validates an embedded migration name.
     33     pub fn new(value: &'static str) -> Result<Self, MigrationContractError> {
     34         if !valid_name(value) {
     35             return Err(MigrationContractError::InvalidName);
     36         }
     37         Ok(Self(value))
     38     }
     39 
     40     /// Returns the validated stable name.
     41     #[must_use]
     42     pub const fn as_str(self) -> &'static str {
     43         self.0
     44     }
     45 }
     46 
     47 impl fmt::Debug for MigrationName {
     48     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     49         formatter
     50             .debug_tuple("MigrationName")
     51             .field(&self.0)
     52             .finish()
     53     }
     54 }
     55 
     56 /// The execution kind whose canonical content is bound by a checksum.
     57 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
     58 pub enum MigrationKind {
     59     Sql,
     60     Callback,
     61 }
     62 
     63 impl MigrationKind {
     64     const fn tag(self) -> u8 {
     65         match self {
     66             Self::Sql => 0,
     67             Self::Callback => 1,
     68         }
     69     }
     70 }
     71 
     72 /// A SHA-256 digest over one migration body or one ordered catalog.
     73 #[derive(Clone, Copy, PartialEq, Eq, Hash)]
     74 pub struct MigrationChecksum([u8; 32]);
     75 
     76 impl MigrationChecksum {
     77     /// Constructs an independently pinned checksum from exact reviewed bytes.
     78     #[must_use]
     79     pub const fn from_bytes(bytes: [u8; 32]) -> Self {
     80         Self(bytes)
     81     }
     82 
     83     /// Computes the frozen SQL-content checksum.
     84     #[must_use]
     85     pub fn for_sql(sql: &str) -> Self {
     86         Self::for_content(MigrationKind::Sql, sql.as_bytes())
     87     }
     88 
     89     /// Computes the frozen callback-definition checksum.
     90     #[must_use]
     91     pub fn for_callback(callback_definition: &[u8]) -> Self {
     92         Self::for_content(MigrationKind::Callback, callback_definition)
     93     }
     94 
     95     /// Returns the exact digest bytes.
     96     #[must_use]
     97     pub const fn as_bytes(&self) -> &[u8; 32] {
     98         &self.0
     99     }
    100 
    101     fn for_content(kind: MigrationKind, content: &[u8]) -> Self {
    102         let mut hasher = Sha256::new();
    103         hasher.update(MIGRATION_CONTENT_DOMAIN);
    104         hasher.update([kind.tag()]);
    105         let content_len = u64::try_from(content.len()).expect("migration content bound fits u64");
    106         hasher.update(content_len.to_be_bytes());
    107         hasher.update(content);
    108         Self(hasher.finalize().into())
    109     }
    110 }
    111 
    112 impl fmt::Debug for MigrationChecksum {
    113     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    114         formatter.write_str("MigrationChecksum([redacted])")
    115     }
    116 }
    117 
    118 /// One immutable future-schema migration identity.
    119 #[derive(Clone, PartialEq, Eq)]
    120 pub struct MigrationDescriptor {
    121     target_version: u32,
    122     name: MigrationName,
    123     kind: MigrationKind,
    124     checksum: MigrationChecksum,
    125     content: &'static [u8],
    126 }
    127 
    128 impl MigrationDescriptor {
    129     /// Defines an embedded SQL migration after verifying its expected checksum.
    130     pub fn sql(
    131         target_version: u32,
    132         name: &'static str,
    133         sql: &'static str,
    134         expected_checksum: MigrationChecksum,
    135     ) -> Result<Self, MigrationContractError> {
    136         Self::new(
    137             target_version,
    138             name,
    139             MigrationKind::Sql,
    140             sql.as_bytes(),
    141             expected_checksum,
    142         )
    143     }
    144 
    145     /// Defines an execution-free callback identity from canonical embedded bytes.
    146     pub fn callback(
    147         target_version: u32,
    148         name: &'static str,
    149         callback_definition: &'static [u8],
    150         expected_checksum: MigrationChecksum,
    151     ) -> Result<Self, MigrationContractError> {
    152         Self::new(
    153             target_version,
    154             name,
    155             MigrationKind::Callback,
    156             callback_definition,
    157             expected_checksum,
    158         )
    159     }
    160 
    161     fn new(
    162         target_version: u32,
    163         name: &'static str,
    164         kind: MigrationKind,
    165         content: &'static [u8],
    166         expected_checksum: MigrationChecksum,
    167     ) -> Result<Self, MigrationContractError> {
    168         if target_version <= BASE_SCHEMA_VERSION {
    169             return Err(MigrationContractError::InvalidTargetVersion);
    170         }
    171         if content.is_empty() {
    172             return Err(MigrationContractError::EmptyContent);
    173         }
    174         if content.len() > MAX_MIGRATION_CONTENT_BYTES {
    175             return Err(MigrationContractError::ContentTooLarge);
    176         }
    177         let name = MigrationName::new(name)?;
    178         let actual_checksum = MigrationChecksum::for_content(kind, content);
    179         if actual_checksum != expected_checksum {
    180             return Err(MigrationContractError::ChecksumMismatch);
    181         }
    182         Ok(Self {
    183             target_version,
    184             name,
    185             kind,
    186             checksum: actual_checksum,
    187             content,
    188         })
    189     }
    190 
    191     /// Returns the schema version produced by this migration.
    192     #[must_use]
    193     pub const fn target_version(&self) -> u32 {
    194         self.target_version
    195     }
    196 
    197     /// Returns the stable migration name.
    198     #[must_use]
    199     pub const fn name(&self) -> MigrationName {
    200         self.name
    201     }
    202 
    203     /// Returns whether this descriptor binds SQL or callback-definition bytes.
    204     #[must_use]
    205     pub const fn kind(&self) -> MigrationKind {
    206         self.kind
    207     }
    208 
    209     /// Returns the verified content checksum.
    210     #[must_use]
    211     pub const fn checksum(&self) -> MigrationChecksum {
    212         self.checksum
    213     }
    214 }
    215 
    216 impl fmt::Debug for MigrationDescriptor {
    217     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    218         formatter
    219             .debug_struct("MigrationDescriptor")
    220             .field("target_version", &self.target_version)
    221             .field("name", &self.name)
    222             .field("kind", &self.kind)
    223             .field("checksum", &self.checksum)
    224             .field("content", &"[redacted]")
    225             .finish()
    226     }
    227 }
    228 
    229 /// One validated ordered migration catalog whose baseline is schema v1.
    230 #[derive(Clone, PartialEq, Eq)]
    231 pub struct MigrationCatalog {
    232     descriptors: Box<[MigrationDescriptor]>,
    233     current_version: u32,
    234     digest: MigrationChecksum,
    235 }
    236 
    237 impl MigrationCatalog {
    238     /// Validates and owns at most 4096 future migrations in exact caller order.
    239     pub fn new<I>(descriptors: I) -> Result<Self, MigrationContractError>
    240     where
    241         I: IntoIterator<Item = MigrationDescriptor>,
    242     {
    243         let descriptors: Vec<_> = descriptors
    244             .into_iter()
    245             .take(MAX_MIGRATION_COUNT + 1)
    246             .collect();
    247         if descriptors.len() > MAX_MIGRATION_COUNT {
    248             return Err(MigrationContractError::TooManyMigrations);
    249         }
    250         validate_catalog(&descriptors)?;
    251         let current_version = descriptors
    252             .last()
    253             .map_or(BASE_SCHEMA_VERSION, MigrationDescriptor::target_version);
    254         let digest = catalog_digest(&descriptors);
    255         Ok(Self {
    256             descriptors: descriptors.into_boxed_slice(),
    257             current_version,
    258             digest,
    259         })
    260     }
    261 
    262     /// Returns the immutable descriptors in exact execution order.
    263     #[must_use]
    264     pub fn descriptors(&self) -> &[MigrationDescriptor] {
    265         &self.descriptors
    266     }
    267 
    268     /// Returns schema v1 for an empty catalog or the last target version.
    269     #[must_use]
    270     pub const fn current_version(&self) -> u32 {
    271         self.current_version
    272     }
    273 
    274     /// Returns the deterministic digest of the ordered catalog identity.
    275     #[must_use]
    276     pub const fn digest(&self) -> MigrationChecksum {
    277         self.digest
    278     }
    279 }
    280 
    281 impl fmt::Debug for MigrationCatalog {
    282     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    283         formatter
    284             .debug_struct("MigrationCatalog")
    285             .field("descriptor_count", &self.descriptors.len())
    286             .field("current_version", &self.current_version)
    287             .field("digest", &self.digest)
    288             .finish()
    289     }
    290 }
    291 
    292 /// Injected Unix timestamp recorded for one applied migration.
    293 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
    294 pub struct MigrationAppliedAtUnixSeconds(u64);
    295 
    296 impl MigrationAppliedAtUnixSeconds {
    297     /// Validates a timestamp representable by SQLite's signed integer storage.
    298     pub const fn new(value: u64) -> Result<Self, MigrationEvidenceError> {
    299         if value > i64::MAX as u64 {
    300             return Err(MigrationEvidenceError::InvalidAppliedTime);
    301         }
    302         Ok(Self(value))
    303     }
    304 
    305     /// Returns the injected Unix timestamp in seconds.
    306     #[must_use]
    307     pub const fn get(self) -> u64 {
    308         self.0
    309     }
    310 }
    311 
    312 /// Complete deterministic application-build identity recorded with a migration.
    313 #[derive(Clone, PartialEq, Eq, Hash)]
    314 pub struct MigrationBuildIdentity {
    315     service_version: String,
    316     service_commit: String,
    317     lib_revision: String,
    318     rust_version: String,
    319     target: String,
    320     feature_profile: String,
    321     config_contract_version: u32,
    322     state_contract_version: u32,
    323     admin_contract_version: u32,
    324     status_contract_version: u32,
    325     provider_contract_version: u32,
    326 }
    327 
    328 impl MigrationBuildIdentity {
    329     /// Validates the complete timestamp-free build identity used by service hosts.
    330     #[allow(clippy::too_many_arguments)]
    331     pub fn new(
    332         service_version: impl AsRef<str>,
    333         service_commit: impl AsRef<str>,
    334         lib_revision: impl AsRef<str>,
    335         rust_version: impl AsRef<str>,
    336         target: impl AsRef<str>,
    337         feature_profile: impl AsRef<str>,
    338         config_contract_version: u32,
    339         state_contract_version: u32,
    340         admin_contract_version: u32,
    341         status_contract_version: u32,
    342         provider_contract_version: u32,
    343     ) -> Result<Self, MigrationEvidenceError> {
    344         let service_version = service_version.as_ref();
    345         let service_commit = service_commit.as_ref();
    346         let lib_revision = lib_revision.as_ref();
    347         let rust_version = rust_version.as_ref();
    348         let target = target.as_ref();
    349         let feature_profile = feature_profile.as_ref();
    350         if !valid_build_text(service_version)
    351             || !valid_revision(service_commit)
    352             || !valid_revision(lib_revision)
    353             || !valid_build_text(rust_version)
    354             || !valid_build_text(target)
    355             || !valid_build_text(feature_profile)
    356             || [
    357                 config_contract_version,
    358                 state_contract_version,
    359                 admin_contract_version,
    360                 status_contract_version,
    361                 provider_contract_version,
    362             ]
    363             .contains(&0)
    364         {
    365             return Err(MigrationEvidenceError::InvalidBuildIdentity);
    366         }
    367         Ok(Self {
    368             service_version: service_version.to_owned(),
    369             service_commit: service_commit.to_owned(),
    370             lib_revision: lib_revision.to_owned(),
    371             rust_version: rust_version.to_owned(),
    372             target: target.to_owned(),
    373             feature_profile: feature_profile.to_owned(),
    374             config_contract_version,
    375             state_contract_version,
    376             admin_contract_version,
    377             status_contract_version,
    378             provider_contract_version,
    379         })
    380     }
    381 
    382     #[must_use]
    383     pub fn service_version(&self) -> &str {
    384         &self.service_version
    385     }
    386 
    387     #[must_use]
    388     pub fn service_commit(&self) -> &str {
    389         &self.service_commit
    390     }
    391 
    392     #[must_use]
    393     pub fn lib_revision(&self) -> &str {
    394         &self.lib_revision
    395     }
    396 
    397     #[must_use]
    398     pub fn rust_version(&self) -> &str {
    399         &self.rust_version
    400     }
    401 
    402     #[must_use]
    403     pub fn target(&self) -> &str {
    404         &self.target
    405     }
    406 
    407     #[must_use]
    408     pub fn feature_profile(&self) -> &str {
    409         &self.feature_profile
    410     }
    411 
    412     #[must_use]
    413     pub const fn config_contract_version(&self) -> u32 {
    414         self.config_contract_version
    415     }
    416 
    417     #[must_use]
    418     pub const fn state_contract_version(&self) -> u32 {
    419         self.state_contract_version
    420     }
    421 
    422     #[must_use]
    423     pub const fn admin_contract_version(&self) -> u32 {
    424         self.admin_contract_version
    425     }
    426 
    427     #[must_use]
    428     pub const fn status_contract_version(&self) -> u32 {
    429         self.status_contract_version
    430     }
    431 
    432     #[must_use]
    433     pub const fn provider_contract_version(&self) -> u32 {
    434         self.provider_contract_version
    435     }
    436 }
    437 
    438 impl fmt::Debug for MigrationBuildIdentity {
    439     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    440         formatter
    441             .debug_struct("MigrationBuildIdentity")
    442             .field("text", &"[redacted]")
    443             .field("revisions", &"[redacted]")
    444             .field(
    445                 "contract_versions",
    446                 &[
    447                     self.config_contract_version,
    448                     self.state_contract_version,
    449                     self.admin_contract_version,
    450                     self.status_contract_version,
    451                     self.provider_contract_version,
    452                 ],
    453             )
    454             .finish()
    455     }
    456 }
    457 
    458 /// Invalid injected migration ledger evidence.
    459 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    460 pub enum MigrationEvidenceError {
    461     InvalidAppliedTime,
    462     InvalidBuildIdentity,
    463 }
    464 
    465 impl fmt::Display for MigrationEvidenceError {
    466     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    467         formatter.write_str(match self {
    468             Self::InvalidAppliedTime => "migration application time is invalid",
    469             Self::InvalidBuildIdentity => "migration application build identity is invalid",
    470         })
    471     }
    472 }
    473 
    474 impl Error for MigrationEvidenceError {}
    475 
    476 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    477 pub struct MigrationApplicationOutcome {
    478     initial_version: u32,
    479     final_version: u32,
    480     applied_count: u32,
    481 }
    482 
    483 impl MigrationApplicationOutcome {
    484     /// Returns the schema version observed under the migration transaction.
    485     #[must_use]
    486     pub const fn initial_version(self) -> u32 {
    487         self.initial_version
    488     }
    489 
    490     /// Returns the schema version committed or already present.
    491     #[must_use]
    492     pub const fn final_version(self) -> u32 {
    493         self.final_version
    494     }
    495 
    496     /// Returns the number of newly committed migration rows.
    497     #[must_use]
    498     pub const fn applied_count(self) -> u32 {
    499         self.applied_count
    500     }
    501 }
    502 
    503 /// Stable, content-free migration contract failures.
    504 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    505 pub enum MigrationContractError {
    506     InvalidName,
    507     InvalidTargetVersion,
    508     EmptyContent,
    509     ContentTooLarge,
    510     ChecksumMismatch,
    511     TooManyMigrations,
    512     DuplicateVersion,
    513     DuplicateName,
    514     OutOfOrder,
    515     VersionGap,
    516 }
    517 
    518 impl fmt::Display for MigrationContractError {
    519     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    520         formatter.write_str(match self {
    521             Self::InvalidName => "migration name is invalid",
    522             Self::InvalidTargetVersion => "migration target version is invalid",
    523             Self::EmptyContent => "migration content is empty",
    524             Self::ContentTooLarge => "migration content exceeds the limit",
    525             Self::ChecksumMismatch => "migration checksum does not match",
    526             Self::TooManyMigrations => "migration catalog exceeds the limit",
    527             Self::DuplicateVersion => "migration target version is duplicated",
    528             Self::DuplicateName => "migration name is duplicated",
    529             Self::OutOfOrder => "migration catalog is out of order",
    530             Self::VersionGap => "migration catalog contains a version gap",
    531         })
    532     }
    533 }
    534 
    535 impl Error for MigrationContractError {}
    536 
    537 fn valid_name(value: &str) -> bool {
    538     let bytes = value.as_bytes();
    539     if bytes.is_empty() || bytes.len() > MAX_MIGRATION_NAME_UTF8_BYTES {
    540         return false;
    541     }
    542     let is_boundary = |byte: u8| byte.is_ascii_lowercase() || byte.is_ascii_digit();
    543     if !is_boundary(bytes[0]) || !is_boundary(bytes[bytes.len() - 1]) {
    544         return false;
    545     }
    546     let mut previous_underscore = false;
    547     for byte in bytes {
    548         if *byte == b'_' {
    549             if previous_underscore {
    550                 return false;
    551             }
    552             previous_underscore = true;
    553         } else if is_boundary(*byte) {
    554             previous_underscore = false;
    555         } else {
    556             return false;
    557         }
    558     }
    559     true
    560 }
    561 
    562 fn valid_build_text(value: &str) -> bool {
    563     let mut bytes = value.bytes();
    564     let Some(first) = bytes.next() else {
    565         return false;
    566     };
    567     value.len() <= MAX_MIGRATION_BUILD_ID_UTF8_BYTES
    568         && first.is_ascii_alphanumeric()
    569         && bytes
    570             .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b':' | b'-'))
    571 }
    572 
    573 fn valid_revision(value: &str) -> bool {
    574     value.len() == 40
    575         && value
    576             .bytes()
    577             .all(|byte| byte.is_ascii_digit() || matches!(byte, b'a'..=b'f'))
    578 }
    579 
    580 fn validate_catalog(descriptors: &[MigrationDescriptor]) -> Result<(), MigrationContractError> {
    581     let mut versions = BTreeSet::new();
    582     let mut names = BTreeSet::new();
    583     for descriptor in descriptors {
    584         if !versions.insert(descriptor.target_version) {
    585             return Err(MigrationContractError::DuplicateVersion);
    586         }
    587         if !names.insert(descriptor.name) {
    588             return Err(MigrationContractError::DuplicateName);
    589         }
    590     }
    591     if descriptors
    592         .windows(2)
    593         .any(|pair| pair[0].target_version > pair[1].target_version)
    594     {
    595         return Err(MigrationContractError::OutOfOrder);
    596     }
    597     for (index, descriptor) in descriptors.iter().enumerate() {
    598         let expected =
    599             BASE_SCHEMA_VERSION + u32::try_from(index).expect("catalog bound fits u32") + 1;
    600         if descriptor.target_version != expected {
    601             return Err(MigrationContractError::VersionGap);
    602         }
    603     }
    604     Ok(())
    605 }
    606 
    607 fn catalog_digest(descriptors: &[MigrationDescriptor]) -> MigrationChecksum {
    608     let mut hasher = Sha256::new();
    609     hasher.update(MIGRATION_CATALOG_DOMAIN);
    610     let descriptor_count =
    611         u32::try_from(descriptors.len()).expect("migration catalog bound fits u32");
    612     hasher.update(descriptor_count.to_be_bytes());
    613     for descriptor in descriptors {
    614         hasher.update(descriptor.target_version.to_be_bytes());
    615         let name = descriptor.name.as_str().as_bytes();
    616         let name_len = u64::try_from(name.len()).expect("migration name bound fits u64");
    617         hasher.update(name_len.to_be_bytes());
    618         hasher.update(name);
    619         hasher.update([descriptor.kind.tag()]);
    620         hasher.update(descriptor.checksum.as_bytes());
    621     }
    622     MigrationChecksum(hasher.finalize().into())
    623 }
    624 
    625 pub type MigrationCallbackFuture<'a> =
    626     Pin<Box<dyn Future<Output = Result<(), ServiceSqliteError>> + Send + 'a>>;
    627 
    628 pub type MigrationCallback =
    629     for<'a> fn(&'a mut MigrationTransactionExecutor<'_>) -> MigrationCallbackFuture<'a>;
    630 
    631 pub struct MigrationTransactionExecutor<'a> {
    632     connection: &'a mut SqliteConnection,
    633     statement_control_rejected: bool,
    634 }
    635 
    636 impl MigrationTransactionExecutor<'_> {
    637     pub async fn execute(&mut self, sql: &'static str) -> Result<(), ServiceSqliteError> {
    638         if crate::statement_policy::contains_forbidden_statement_control(sql) {
    639             self.statement_control_rejected = true;
    640             return Err(ServiceSqliteError::new(
    641                 crate::ServiceSqliteErrorKind::Migration,
    642             ));
    643         }
    644         #[cfg(any(target_os = "linux", target_os = "macos"))]
    645         {
    646             assert_governed_transaction(self.connection).await?;
    647             let execution = sqlx::raw_sql(sql).execute(&mut *self.connection).await;
    648             let transaction = assert_governed_transaction(self.connection).await;
    649             transaction?;
    650             execution
    651                 .map(|_| ())
    652                 .map_err(|source| migration_source(MigrationFailureKind::Execution, source))
    653         }
    654         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    655         {
    656             let _ = (&mut self.connection, sql);
    657             Err(ServiceSqliteError::new(
    658                 crate::ServiceSqliteErrorKind::Migration,
    659             ))
    660         }
    661     }
    662 }
    663 
    664 #[cfg(any(target_os = "linux", target_os = "macos"))]
    665 #[derive(PartialEq, Eq)]
    666 pub(crate) struct MigrationConnectionPolicy {
    667     application_id: i64,
    668     journal_mode: String,
    669     synchronous: i64,
    670     foreign_keys: i64,
    671     trusted_schema: i64,
    672     busy_timeout: i64,
    673     query_only: i64,
    674 }
    675 
    676 #[derive(Clone, Copy)]
    677 pub struct MigrationCallbackBinding {
    678     #[cfg(any(target_os = "linux", target_os = "macos"))]
    679     target_version: u32,
    680     #[cfg(any(target_os = "linux", target_os = "macos"))]
    681     name: MigrationName,
    682     #[cfg(any(target_os = "linux", target_os = "macos"))]
    683     checksum: MigrationChecksum,
    684     #[cfg(any(target_os = "linux", target_os = "macos"))]
    685     callback: MigrationCallback,
    686 }
    687 
    688 impl MigrationCallbackBinding {
    689     pub const fn new(
    690         target_version: u32,
    691         name: MigrationName,
    692         checksum: MigrationChecksum,
    693         callback: MigrationCallback,
    694     ) -> Self {
    695         #[cfg(any(target_os = "linux", target_os = "macos"))]
    696         {
    697             Self {
    698                 target_version,
    699                 name,
    700                 checksum,
    701                 callback,
    702             }
    703         }
    704         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    705         {
    706             let _ = (target_version, name, checksum, callback);
    707             Self {}
    708         }
    709     }
    710 }
    711 
    712 #[cfg(any(target_os = "linux", target_os = "macos"))]
    713 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    714 enum MigrationFailureKind {
    715     CatalogMismatch,
    716     HistoryCorrupt,
    717     CallbackBinding,
    718     Execution,
    719     LedgerWrite,
    720     MetadataAdvance,
    721     Commit,
    722 }
    723 
    724 #[cfg(any(target_os = "linux", target_os = "macos"))]
    725 #[derive(Debug)]
    726 struct MigrationFailure(MigrationFailureKind);
    727 
    728 #[cfg(any(target_os = "linux", target_os = "macos"))]
    729 impl fmt::Display for MigrationFailure {
    730     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    731         formatter.write_str(match self.0 {
    732             MigrationFailureKind::CatalogMismatch => "migration catalog does not match state",
    733             MigrationFailureKind::HistoryCorrupt => "migration history is corrupt",
    734             MigrationFailureKind::CallbackBinding => "migration callback binding is invalid",
    735             MigrationFailureKind::Execution => "migration execution failed",
    736             MigrationFailureKind::LedgerWrite => "migration ledger write failed",
    737             MigrationFailureKind::MetadataAdvance => "migration metadata advance failed",
    738             MigrationFailureKind::Commit => "migration commit outcome is unavailable",
    739         })
    740     }
    741 }
    742 
    743 #[cfg(any(target_os = "linux", target_os = "macos"))]
    744 impl Error for MigrationFailure {}
    745 
    746 #[cfg(any(target_os = "linux", target_os = "macos"))]
    747 struct MigrationSource {
    748     kind: MigrationFailureKind,
    749     source: Box<dyn Error + Send + Sync + 'static>,
    750 }
    751 
    752 #[cfg(any(target_os = "linux", target_os = "macos"))]
    753 impl fmt::Debug for MigrationSource {
    754     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    755         formatter
    756             .debug_struct("MigrationSource")
    757             .field("kind", &self.kind)
    758             .field("source", &"[redacted]")
    759             .finish()
    760     }
    761 }
    762 
    763 #[cfg(any(target_os = "linux", target_os = "macos"))]
    764 impl fmt::Display for MigrationSource {
    765     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    766         MigrationFailure(self.kind).fmt(formatter)
    767     }
    768 }
    769 
    770 #[cfg(any(target_os = "linux", target_os = "macos"))]
    771 impl Error for MigrationSource {
    772     fn source(&self) -> Option<&(dyn Error + 'static)> {
    773         Some(self.source.as_ref())
    774     }
    775 }
    776 
    777 #[cfg(any(target_os = "linux", target_os = "macos"))]
    778 fn migration_error(kind: MigrationFailureKind) -> ServiceSqliteError {
    779     ServiceSqliteError::with_source(ServiceSqliteErrorKind::Migration, MigrationFailure(kind))
    780 }
    781 
    782 #[cfg(any(target_os = "linux", target_os = "macos"))]
    783 fn migration_source(
    784     kind: MigrationFailureKind,
    785     source: impl Error + Send + Sync + 'static,
    786 ) -> ServiceSqliteError {
    787     ServiceSqliteError::with_source(
    788         ServiceSqliteErrorKind::Migration,
    789         MigrationSource {
    790             kind,
    791             source: Box::new(source),
    792         },
    793     )
    794 }
    795 
    796 #[cfg(any(target_os = "linux", target_os = "macos"))]
    797 fn require_migration_condition(
    798     condition: bool,
    799     kind: MigrationFailureKind,
    800 ) -> Result<(), ServiceSqliteError> {
    801     condition.then_some(()).ok_or_else(|| migration_error(kind))
    802 }
    803 
    804 #[cfg(any(target_os = "linux", target_os = "macos"))]
    805 #[derive(Clone, PartialEq, Eq)]
    806 struct AppliedMigration {
    807     version: u32,
    808     name: String,
    809     checksum: MigrationChecksum,
    810     applied_at: MigrationAppliedAtUnixSeconds,
    811     build: MigrationBuildIdentity,
    812 }
    813 
    814 #[cfg(any(target_os = "linux", target_os = "macos"))]
    815 pub(crate) async fn verify_migration_history(
    816     connection: &mut SqliteConnection,
    817     catalog: &MigrationCatalog,
    818     schema_catalog: &SchemaCatalog,
    819     require_current: bool,
    820 ) -> Result<u32, ServiceSqliteError> {
    821     schema_catalog
    822         .matches_migrations(catalog)
    823         .then_some(())
    824         .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
    825     let mut transaction = connection
    826         .begin()
    827         .await
    828         .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?;
    829     let result = verify_migration_history_snapshot(
    830         &mut transaction,
    831         catalog,
    832         schema_catalog,
    833         require_current,
    834     )
    835     .await;
    836     let rollback = transaction.rollback().await;
    837     match (result, rollback) {
    838         (Err(error), _) => Err(error),
    839         (Ok(_), Err(source)) => Err(migration_source(
    840             MigrationFailureKind::HistoryCorrupt,
    841             source,
    842         )),
    843         (Ok(version), Ok(())) => Ok(version),
    844     }
    845 }
    846 
    847 #[cfg(any(target_os = "linux", target_os = "macos"))]
    848 pub(crate) async fn verify_migration_history_snapshot(
    849     connection: &mut SqliteConnection,
    850     catalog: &MigrationCatalog,
    851     schema_catalog: &SchemaCatalog,
    852     require_current: bool,
    853 ) -> Result<u32, ServiceSqliteError> {
    854     let version = read_state_schema_version(connection).await?;
    855     let history = read_migration_history(connection).await?;
    856     validate_migration_prefix(catalog, version, &history)?;
    857     crate::integrity::verify_schema_catalog(connection, schema_catalog, version).await?;
    858     if require_current && version != catalog.current_version() {
    859         return Err(migration_error(MigrationFailureKind::CatalogMismatch));
    860     }
    861     Ok(version)
    862 }
    863 
    864 #[cfg(any(target_os = "linux", target_os = "macos"))]
    865 pub(crate) async fn apply_governed_migrations<V>(
    866     connection: &mut SqliteConnection,
    867     catalog: &MigrationCatalog,
    868     schema_catalog: &SchemaCatalog,
    869     applied_at: MigrationAppliedAtUnixSeconds,
    870     build: &MigrationBuildIdentity,
    871     callback_bindings: &[MigrationCallbackBinding],
    872     validate_authority: &mut V,
    873 ) -> Result<MigrationApplicationOutcome, ServiceSqliteError>
    874 where
    875     V: FnMut() -> Result<(), ServiceSqliteError>,
    876 {
    877     let mut after_commit = || Ok(());
    878     apply_governed_migrations_with_observer(
    879         connection,
    880         catalog,
    881         schema_catalog,
    882         applied_at,
    883         build,
    884         callback_bindings,
    885         validate_authority,
    886         &mut after_commit,
    887     )
    888     .await
    889 }
    890 
    891 #[cfg(any(target_os = "linux", target_os = "macos"))]
    892 #[allow(clippy::too_many_arguments)]
    893 async fn apply_governed_migrations_with_observer<V, O>(
    894     connection: &mut SqliteConnection,
    895     catalog: &MigrationCatalog,
    896     schema_catalog: &SchemaCatalog,
    897     applied_at: MigrationAppliedAtUnixSeconds,
    898     build: &MigrationBuildIdentity,
    899     callback_bindings: &[MigrationCallbackBinding],
    900     validate_authority: &mut V,
    901     after_commit: &mut O,
    902 ) -> Result<MigrationApplicationOutcome, ServiceSqliteError>
    903 where
    904     V: FnMut() -> Result<(), ServiceSqliteError>,
    905     O: FnMut() -> Result<(), ServiceSqliteError>,
    906 {
    907     schema_catalog
    908         .matches_migrations(catalog)
    909         .then_some(())
    910         .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
    911     let callbacks = validate_callback_bindings(catalog, callback_bindings)?;
    912     let mut initial_version = None;
    913     let mut applied_count = 0_u32;
    914     validate_authority()?;
    915     let initial_policy_result = read_connection_policy(connection).await;
    916     validate_authority()?;
    917     let initial_policy = initial_policy_result?;
    918 
    919     loop {
    920         validate_authority()?;
    921         let gate_result = crate::transaction_control::TransactionControlGate::install(connection)
    922             .await
    923             .map_err(|source| migration_source(MigrationFailureKind::Execution, source));
    924         validate_authority()?;
    925         let commit_gate = gate_result?;
    926         let transaction_result = connection.begin_with("BEGIN IMMEDIATE").await;
    927         validate_authority()?;
    928         let mut transaction = transaction_result
    929             .map_err(|source| migration_source(MigrationFailureKind::Execution, source))?;
    930         let transactional_result =
    931             verify_migration_history_snapshot(&mut transaction, catalog, schema_catalog, false)
    932                 .await;
    933         validate_authority()?;
    934         let current = transactional_result?;
    935         initial_version.get_or_insert(current);
    936         if current == catalog.current_version() {
    937             let rollback_permit = commit_gate.permit_runner_rollback();
    938             let rollback_result = transaction.rollback().await;
    939             drop(rollback_permit);
    940             validate_authority()?;
    941             rollback_result
    942                 .map_err(|source| migration_source(MigrationFailureKind::Commit, source))?;
    943             let remove_result = commit_gate
    944                 .remove(connection)
    945                 .await
    946                 .map_err(|source| migration_source(MigrationFailureKind::Commit, source));
    947             validate_authority()?;
    948             remove_result?;
    949             break;
    950         }
    951         let descriptor_index = usize::try_from(current.saturating_sub(BASE_SCHEMA_VERSION))
    952             .map_err(|_| migration_error(MigrationFailureKind::CatalogMismatch))?;
    953         let descriptor = catalog
    954             .descriptors()
    955             .get(descriptor_index)
    956             .ok_or_else(|| migration_error(MigrationFailureKind::CatalogMismatch))?;
    957         validate_authority()?;
    958         let execution_result = execute_descriptor(&mut transaction, descriptor, &callbacks).await;
    959         validate_authority()?;
    960         execution_result?;
    961         require_migration_condition(
    962             !commit_gate.control_violation_observed(),
    963             MigrationFailureKind::Execution,
    964         )?;
    965         let transaction_result = assert_governed_transaction(&mut transaction).await;
    966         validate_authority()?;
    967         transaction_result?;
    968         let policy_result = read_connection_policy(&mut transaction).await;
    969         validate_authority()?;
    970         require_migration_condition(
    971             policy_result? == initial_policy,
    972             MigrationFailureKind::Execution,
    973         )?;
    974         let transaction_result = assert_governed_transaction(&mut transaction).await;
    975         validate_authority()?;
    976         transaction_result?;
    977         let schema_result = crate::integrity::verify_schema_catalog(
    978             &mut transaction,
    979             schema_catalog,
    980             descriptor.target_version(),
    981         )
    982         .await;
    983         validate_authority()?;
    984         schema_result?;
    985         let insert_result =
    986             insert_migration_row(&mut transaction, descriptor, applied_at, build).await;
    987         validate_authority()?;
    988         insert_result?;
    989         let transaction_result = assert_governed_transaction(&mut transaction).await;
    990         validate_authority()?;
    991         transaction_result?;
    992         let advance_result =
    993             advance_schema_version(&mut transaction, current, descriptor.target_version()).await;
    994         validate_authority()?;
    995         advance_result?;
    996         let transaction_result = assert_governed_transaction(&mut transaction).await;
    997         validate_authority()?;
    998         transaction_result?;
    999         require_migration_condition(
   1000             !commit_gate.control_violation_observed(),
   1001             MigrationFailureKind::Execution,
   1002         )?;
   1003         let policy_result = read_connection_policy(&mut transaction).await;
   1004         validate_authority()?;
   1005         require_migration_condition(
   1006             policy_result? == initial_policy,
   1007             MigrationFailureKind::Execution,
   1008         )?;
   1009         let transaction_result = assert_governed_transaction(&mut transaction).await;
   1010         validate_authority()?;
   1011         transaction_result?;
   1012         require_migration_condition(
   1013             !commit_gate.control_violation_observed(),
   1014             MigrationFailureKind::Execution,
   1015         )?;
   1016         let permit = commit_gate.permit_outer_commit();
   1017         let commit_result = transaction.commit().await;
   1018         drop(permit);
   1019         validate_authority()?;
   1020         commit_result.map_err(|source| migration_source(MigrationFailureKind::Commit, source))?;
   1021         let remove_result = commit_gate
   1022             .remove(connection)
   1023             .await
   1024             .map_err(|source| migration_source(MigrationFailureKind::Commit, source));
   1025         validate_authority()?;
   1026         remove_result?;
   1027         let observed = after_commit();
   1028         validate_authority()?;
   1029         observed?;
   1030         applied_count = applied_count.saturating_add(1);
   1031     }
   1032 
   1033     validate_authority()?;
   1034     let final_result = verify_migration_history(connection, catalog, schema_catalog, true).await;
   1035     validate_authority()?;
   1036     let final_version = final_result?;
   1037     Ok(MigrationApplicationOutcome {
   1038         initial_version: initial_version.unwrap_or(final_version),
   1039         final_version,
   1040         applied_count,
   1041     })
   1042 }
   1043 
   1044 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1045 fn validate_callback_bindings(
   1046     catalog: &MigrationCatalog,
   1047     bindings: &[MigrationCallbackBinding],
   1048 ) -> Result<BTreeMap<u32, MigrationCallback>, ServiceSqliteError> {
   1049     let expected = catalog
   1050         .descriptors()
   1051         .iter()
   1052         .filter(|descriptor| descriptor.kind() == MigrationKind::Callback)
   1053         .count();
   1054     if bindings.len() != expected {
   1055         return Err(migration_error(MigrationFailureKind::CallbackBinding));
   1056     }
   1057     let mut callbacks = BTreeMap::new();
   1058     for binding in bindings {
   1059         let descriptor = catalog
   1060             .descriptors()
   1061             .iter()
   1062             .find(|descriptor| descriptor.target_version() == binding.target_version)
   1063             .ok_or_else(|| migration_error(MigrationFailureKind::CallbackBinding))?;
   1064         let unique = callbacks
   1065             .insert(binding.target_version, binding.callback)
   1066             .is_none();
   1067         if !crate::all_constraints([
   1068             descriptor.kind() == MigrationKind::Callback,
   1069             descriptor.name() == binding.name,
   1070             descriptor.checksum() == binding.checksum,
   1071             unique,
   1072         ]) {
   1073             return Err(migration_error(MigrationFailureKind::CallbackBinding));
   1074         }
   1075     }
   1076     Ok(callbacks)
   1077 }
   1078 
   1079 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1080 async fn execute_descriptor(
   1081     connection: &mut SqliteConnection,
   1082     descriptor: &MigrationDescriptor,
   1083     callbacks: &BTreeMap<u32, MigrationCallback>,
   1084 ) -> Result<(), ServiceSqliteError> {
   1085     let mut executor = MigrationTransactionExecutor {
   1086         connection,
   1087         statement_control_rejected: false,
   1088     };
   1089     match descriptor.kind() {
   1090         MigrationKind::Sql => {
   1091             let sql = core::str::from_utf8(descriptor.content)
   1092                 .map_err(|source| migration_source(MigrationFailureKind::Execution, source))?;
   1093             executor.execute(sql).await?;
   1094         }
   1095         MigrationKind::Callback => {
   1096             let callback = callbacks
   1097                 .get(&descriptor.target_version())
   1098                 .ok_or_else(|| migration_error(MigrationFailureKind::CallbackBinding))?;
   1099             callback(&mut executor)
   1100                 .await
   1101                 .map_err(|source| migration_source(MigrationFailureKind::Execution, source))?;
   1102             assert_governed_transaction(executor.connection).await?;
   1103         }
   1104     }
   1105     if executor.statement_control_rejected {
   1106         return Err(migration_error(MigrationFailureKind::Execution));
   1107     }
   1108     Ok(())
   1109 }
   1110 
   1111 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1112 pub(crate) async fn assert_governed_transaction(
   1113     connection: &mut SqliteConnection,
   1114 ) -> Result<(), ServiceSqliteError> {
   1115     sqlx::raw_sql(
   1116         "SAVEPOINT radroots_migration_transaction_probe;
   1117          RELEASE SAVEPOINT radroots_migration_transaction_probe;",
   1118     )
   1119     .execute(connection)
   1120     .await
   1121     .map(|_| ())
   1122     .map_err(|source| migration_source(MigrationFailureKind::Execution, source))
   1123 }
   1124 
   1125 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1126 pub(crate) async fn read_connection_policy(
   1127     connection: &mut SqliteConnection,
   1128 ) -> Result<MigrationConnectionPolicy, ServiceSqliteError> {
   1129     let text = |source| migration_source(MigrationFailureKind::Execution, source);
   1130     let application_id = sqlx::query_scalar::<_, i64>("PRAGMA application_id")
   1131         .fetch_one(&mut *connection)
   1132         .await
   1133         .map_err(text)?;
   1134     let journal_mode = sqlx::query_scalar::<_, String>("PRAGMA journal_mode")
   1135         .fetch_one(&mut *connection)
   1136         .await
   1137         .map_err(text)?;
   1138     require_migration_condition(journal_mode.len() <= 16, MigrationFailureKind::Execution)?;
   1139     let synchronous = sqlx::query_scalar::<_, i64>("PRAGMA synchronous")
   1140         .fetch_one(&mut *connection)
   1141         .await
   1142         .map_err(text)?;
   1143     let foreign_keys = sqlx::query_scalar::<_, i64>("PRAGMA foreign_keys")
   1144         .fetch_one(&mut *connection)
   1145         .await
   1146         .map_err(text)?;
   1147     let trusted_schema = sqlx::query_scalar::<_, i64>("PRAGMA trusted_schema")
   1148         .fetch_one(&mut *connection)
   1149         .await
   1150         .map_err(text)?;
   1151     let busy_timeout = sqlx::query_scalar::<_, i64>("PRAGMA busy_timeout")
   1152         .fetch_one(&mut *connection)
   1153         .await
   1154         .map_err(text)?;
   1155     let query_only = sqlx::query_scalar::<_, i64>("PRAGMA query_only")
   1156         .fetch_one(&mut *connection)
   1157         .await
   1158         .map_err(text)?;
   1159     let databases = sqlx::query_scalar::<_, String>(
   1160         "SELECT CASE
   1161              WHEN typeof(name) = 'text' AND length(CAST(name AS BLOB)) <= 4 THEN name
   1162              ELSE ''
   1163          END
   1164          FROM pragma_database_list
   1165          ORDER BY seq
   1166          LIMIT 2",
   1167     )
   1168     .fetch_all(connection)
   1169     .await
   1170     .map_err(text)?;
   1171     require_migration_condition(databases == ["main"], MigrationFailureKind::Execution)?;
   1172     Ok(MigrationConnectionPolicy {
   1173         application_id,
   1174         journal_mode,
   1175         synchronous,
   1176         foreign_keys,
   1177         trusted_schema,
   1178         busy_timeout,
   1179         query_only,
   1180     })
   1181 }
   1182 
   1183 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1184 async fn insert_migration_row(
   1185     connection: &mut SqliteConnection,
   1186     descriptor: &MigrationDescriptor,
   1187     applied_at: MigrationAppliedAtUnixSeconds,
   1188     build: &MigrationBuildIdentity,
   1189 ) -> Result<(), ServiceSqliteError> {
   1190     let result = sqlx::query(
   1191         "INSERT INTO schema_migrations (
   1192             version, name, checksum, applied_at_unix_s,
   1193             service_version, service_commit, lib_revision, rust_version, target, feature_profile,
   1194             config_contract_version, state_contract_version, admin_contract_version,
   1195             status_contract_version, provider_contract_version
   1196          ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
   1197     )
   1198     .bind(i64::from(descriptor.target_version()))
   1199     .bind(descriptor.name().as_str())
   1200     .bind(descriptor.checksum().as_bytes().as_slice())
   1201     .bind(
   1202         i64::try_from(applied_at.get())
   1203             .map_err(|_| migration_error(MigrationFailureKind::LedgerWrite))?,
   1204     )
   1205     .bind(build.service_version())
   1206     .bind(build.service_commit())
   1207     .bind(build.lib_revision())
   1208     .bind(build.rust_version())
   1209     .bind(build.target())
   1210     .bind(build.feature_profile())
   1211     .bind(i64::from(build.config_contract_version()))
   1212     .bind(i64::from(build.state_contract_version()))
   1213     .bind(i64::from(build.admin_contract_version()))
   1214     .bind(i64::from(build.status_contract_version()))
   1215     .bind(i64::from(build.provider_contract_version()))
   1216     .execute(connection)
   1217     .await
   1218     .map_err(|source| migration_source(MigrationFailureKind::LedgerWrite, source))?;
   1219     require_migration_condition(
   1220         result.rows_affected() == 1,
   1221         MigrationFailureKind::LedgerWrite,
   1222     )?;
   1223     Ok(())
   1224 }
   1225 
   1226 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1227 async fn advance_schema_version(
   1228     connection: &mut SqliteConnection,
   1229     current: u32,
   1230     target: u32,
   1231 ) -> Result<(), ServiceSqliteError> {
   1232     let result = sqlx::query(
   1233         "UPDATE radroots_service_metadata
   1234          SET state_schema_version = ?
   1235          WHERE singleton = 1 AND state_schema_version = ?",
   1236     )
   1237     .bind(i64::from(target))
   1238     .bind(i64::from(current))
   1239     .execute(connection)
   1240     .await
   1241     .map_err(|source| migration_source(MigrationFailureKind::MetadataAdvance, source))?;
   1242     require_migration_condition(
   1243         result.rows_affected() == 1,
   1244         MigrationFailureKind::MetadataAdvance,
   1245     )?;
   1246     Ok(())
   1247 }
   1248 
   1249 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1250 async fn read_state_schema_version(
   1251     connection: &mut SqliteConnection,
   1252 ) -> Result<u32, ServiceSqliteError> {
   1253     let rows = sqlx::query(
   1254         "SELECT state_schema_version, typeof(state_schema_version) AS version_type
   1255          FROM radroots_service_metadata
   1256          WHERE singleton = 1
   1257          LIMIT 2",
   1258     )
   1259     .fetch_all(connection)
   1260     .await
   1261     .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?;
   1262     let [row] = rows.as_slice() else {
   1263         return Err(migration_error(MigrationFailureKind::HistoryCorrupt));
   1264     };
   1265     if row
   1266         .try_get::<String, _>("version_type")
   1267         .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?
   1268         != "integer"
   1269     {
   1270         return Err(migration_error(MigrationFailureKind::HistoryCorrupt));
   1271     }
   1272     u32::try_from(
   1273         row.try_get::<i64, _>("state_schema_version")
   1274             .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?,
   1275     )
   1276     .map_err(|_| migration_error(MigrationFailureKind::HistoryCorrupt))
   1277 }
   1278 
   1279 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1280 async fn read_migration_history(
   1281     connection: &mut SqliteConnection,
   1282 ) -> Result<Vec<AppliedMigration>, ServiceSqliteError> {
   1283     let rows = sqlx::query(
   1284         "SELECT
   1285             version,
   1286             applied_at_unix_s,
   1287             config_contract_version, state_contract_version, admin_contract_version,
   1288             status_contract_version, provider_contract_version,
   1289             typeof(version) = 'integer' AS version_type_ok,
   1290             typeof(name) = 'text' AS name_type_ok,
   1291             length(CAST(name AS BLOB)) AS name_length,
   1292             substr(CAST(name AS BLOB), 1, 129) AS name_prefix,
   1293             typeof(checksum) = 'blob' AS checksum_type_ok,
   1294             length(checksum) AS checksum_length,
   1295             substr(checksum, 1, 33) AS checksum_prefix,
   1296             typeof(applied_at_unix_s) = 'integer' AS applied_at_type_ok,
   1297             typeof(service_version) = 'text' AS service_version_type_ok,
   1298             length(CAST(service_version AS BLOB)) AS service_version_length,
   1299             substr(CAST(service_version AS BLOB), 1, 129) AS service_version_prefix,
   1300             typeof(service_commit) = 'text' AS service_commit_type_ok,
   1301             length(CAST(service_commit AS BLOB)) AS service_commit_length,
   1302             substr(CAST(service_commit AS BLOB), 1, 41) AS service_commit_prefix,
   1303             typeof(lib_revision) = 'text' AS lib_revision_type_ok,
   1304             length(CAST(lib_revision AS BLOB)) AS lib_revision_length,
   1305             substr(CAST(lib_revision AS BLOB), 1, 41) AS lib_revision_prefix,
   1306             typeof(rust_version) = 'text' AS rust_version_type_ok,
   1307             length(CAST(rust_version AS BLOB)) AS rust_version_length,
   1308             substr(CAST(rust_version AS BLOB), 1, 129) AS rust_version_prefix,
   1309             typeof(target) = 'text' AS target_type_ok,
   1310             length(CAST(target AS BLOB)) AS target_length,
   1311             substr(CAST(target AS BLOB), 1, 129) AS target_prefix,
   1312             typeof(feature_profile) = 'text' AS feature_profile_type_ok,
   1313             length(CAST(feature_profile AS BLOB)) AS feature_profile_length,
   1314             substr(CAST(feature_profile AS BLOB), 1, 129) AS feature_profile_prefix,
   1315             typeof(config_contract_version) = 'integer' AS config_contract_version_type_ok,
   1316             typeof(state_contract_version) = 'integer' AS state_contract_version_type_ok,
   1317             typeof(admin_contract_version) = 'integer' AS admin_contract_version_type_ok,
   1318             typeof(status_contract_version) = 'integer' AS status_contract_version_type_ok,
   1319             typeof(provider_contract_version) = 'integer' AS provider_contract_version_type_ok
   1320          FROM schema_migrations
   1321          ORDER BY version
   1322          LIMIT 4097",
   1323     )
   1324     .fetch_all(connection)
   1325     .await
   1326     .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?;
   1327     require_migration_condition(
   1328         rows.len() <= MAX_MIGRATION_COUNT,
   1329         MigrationFailureKind::HistoryCorrupt,
   1330     )?;
   1331     rows.iter().map(parse_applied_migration).collect()
   1332 }
   1333 
   1334 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1335 fn parse_applied_migration(
   1336     row: &sqlx::sqlite::SqliteRow,
   1337 ) -> Result<AppliedMigration, ServiceSqliteError> {
   1338     for column in [
   1339         "version_type_ok",
   1340         "applied_at_type_ok",
   1341         "config_contract_version_type_ok",
   1342         "state_contract_version_type_ok",
   1343         "admin_contract_version_type_ok",
   1344         "status_contract_version_type_ok",
   1345         "provider_contract_version_type_ok",
   1346     ] {
   1347         require_migration_condition(
   1348             row.try_get::<i64, _>(column)
   1349                 .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?
   1350                 == 1,
   1351             MigrationFailureKind::HistoryCorrupt,
   1352         )?;
   1353     }
   1354     let version = u32::try_from(
   1355         row.try_get::<i64, _>("version")
   1356             .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?,
   1357     )
   1358     .map_err(|_| migration_error(MigrationFailureKind::HistoryCorrupt))?;
   1359     let name = crate::persisted_value::bounded_utf8(
   1360         row,
   1361         "name_type_ok",
   1362         "name_length",
   1363         "name_prefix",
   1364         1,
   1365         MAX_MIGRATION_NAME_UTF8_BYTES,
   1366     )
   1367     .ok_or_else(|| migration_error(MigrationFailureKind::HistoryCorrupt))?;
   1368     if !valid_name(name) {
   1369         return Err(migration_error(MigrationFailureKind::HistoryCorrupt));
   1370     }
   1371     let checksum: [u8; 32] = crate::persisted_value::bounded_bytes(
   1372         row,
   1373         "checksum_type_ok",
   1374         "checksum_length",
   1375         "checksum_prefix",
   1376         32,
   1377         32,
   1378     )
   1379     .ok_or_else(|| migration_error(MigrationFailureKind::HistoryCorrupt))?
   1380     .try_into()
   1381     .map_err(|_| migration_error(MigrationFailureKind::HistoryCorrupt))?;
   1382     let applied_at = MigrationAppliedAtUnixSeconds::new(
   1383         u64::try_from(
   1384             row.try_get::<i64, _>("applied_at_unix_s")
   1385                 .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?,
   1386         )
   1387         .map_err(|_| migration_error(MigrationFailureKind::HistoryCorrupt))?,
   1388     )
   1389     .map_err(|_| migration_error(MigrationFailureKind::HistoryCorrupt))?;
   1390     let text = |type_column, length_column, prefix_column, minimum, maximum| {
   1391         crate::persisted_value::bounded_utf8(
   1392             row,
   1393             type_column,
   1394             length_column,
   1395             prefix_column,
   1396             minimum,
   1397             maximum,
   1398         )
   1399         .ok_or_else(|| migration_error(MigrationFailureKind::HistoryCorrupt))
   1400     };
   1401     let version_field = |column| {
   1402         u32::try_from(
   1403             row.try_get::<i64, _>(column)
   1404                 .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?,
   1405         )
   1406         .map_err(|_| migration_error(MigrationFailureKind::HistoryCorrupt))
   1407     };
   1408     let build = MigrationBuildIdentity::new(
   1409         text(
   1410             "service_version_type_ok",
   1411             "service_version_length",
   1412             "service_version_prefix",
   1413             1,
   1414             MAX_MIGRATION_BUILD_ID_UTF8_BYTES,
   1415         )?,
   1416         text(
   1417             "service_commit_type_ok",
   1418             "service_commit_length",
   1419             "service_commit_prefix",
   1420             40,
   1421             40,
   1422         )?,
   1423         text(
   1424             "lib_revision_type_ok",
   1425             "lib_revision_length",
   1426             "lib_revision_prefix",
   1427             40,
   1428             40,
   1429         )?,
   1430         text(
   1431             "rust_version_type_ok",
   1432             "rust_version_length",
   1433             "rust_version_prefix",
   1434             1,
   1435             MAX_MIGRATION_BUILD_ID_UTF8_BYTES,
   1436         )?,
   1437         text(
   1438             "target_type_ok",
   1439             "target_length",
   1440             "target_prefix",
   1441             1,
   1442             MAX_MIGRATION_BUILD_ID_UTF8_BYTES,
   1443         )?,
   1444         text(
   1445             "feature_profile_type_ok",
   1446             "feature_profile_length",
   1447             "feature_profile_prefix",
   1448             1,
   1449             MAX_MIGRATION_BUILD_ID_UTF8_BYTES,
   1450         )?,
   1451         version_field("config_contract_version")?,
   1452         version_field("state_contract_version")?,
   1453         version_field("admin_contract_version")?,
   1454         version_field("status_contract_version")?,
   1455         version_field("provider_contract_version")?,
   1456     )
   1457     .map_err(|_| migration_error(MigrationFailureKind::HistoryCorrupt))?;
   1458     Ok(AppliedMigration {
   1459         version,
   1460         name: name.to_owned(),
   1461         checksum: MigrationChecksum::from_bytes(checksum),
   1462         applied_at,
   1463         build,
   1464     })
   1465 }
   1466 
   1467 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1468 fn validate_migration_prefix(
   1469     catalog: &MigrationCatalog,
   1470     version: u32,
   1471     history: &[AppliedMigration],
   1472 ) -> Result<(), ServiceSqliteError> {
   1473     require_migration_condition(
   1474         version >= BASE_SCHEMA_VERSION && version <= catalog.current_version(),
   1475         MigrationFailureKind::CatalogMismatch,
   1476     )?;
   1477     let expected_len = usize::try_from(version - BASE_SCHEMA_VERSION)
   1478         .map_err(|_| migration_error(MigrationFailureKind::HistoryCorrupt))?;
   1479     require_migration_condition(
   1480         history.len() == expected_len,
   1481         MigrationFailureKind::CatalogMismatch,
   1482     )?;
   1483     for (applied, descriptor) in history.iter().zip(catalog.descriptors()) {
   1484         if !crate::all_constraints([
   1485             applied.version == descriptor.target_version(),
   1486             applied.name == descriptor.name().as_str(),
   1487             applied.checksum == descriptor.checksum(),
   1488         ]) {
   1489             return Err(migration_error(MigrationFailureKind::CatalogMismatch));
   1490         }
   1491         let _ = (applied.applied_at, &applied.build);
   1492     }
   1493     Ok(())
   1494 }
   1495 
   1496 #[cfg(test)]
   1497 mod tests {
   1498     use super::*;
   1499 
   1500     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1501     #[tokio::test(flavor = "current_thread")]
   1502     async fn schema_version_reads_require_one_integer_metadata_row() {
   1503         let mut connection = SqliteConnection::connect("sqlite::memory:").await.unwrap();
   1504         sqlx::query("CREATE TABLE radroots_service_metadata (singleton, state_schema_version)")
   1505             .execute(&mut connection)
   1506             .await
   1507             .unwrap();
   1508         for statement in [
   1509             "DELETE FROM radroots_service_metadata",
   1510             "INSERT INTO radroots_service_metadata VALUES (1, 1), (1, 1)",
   1511             "DELETE FROM radroots_service_metadata; INSERT INTO radroots_service_metadata VALUES (1, 'invalid')",
   1512             "DELETE FROM radroots_service_metadata; INSERT INTO radroots_service_metadata VALUES (1, NULL)",
   1513         ] {
   1514             sqlx::raw_sql(sqlx::AssertSqlSafe(statement))
   1515                 .execute(&mut connection)
   1516                 .await
   1517                 .unwrap();
   1518             assert_eq!(
   1519                 read_state_schema_version(&mut connection)
   1520                     .await
   1521                     .unwrap_err()
   1522                     .kind(),
   1523                 ServiceSqliteErrorKind::Migration
   1524             );
   1525         }
   1526         sqlx::raw_sql("DELETE FROM radroots_service_metadata; INSERT INTO radroots_service_metadata VALUES (1, 1)")
   1527             .execute(&mut connection).await.unwrap();
   1528         assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 1);
   1529         connection.close().await.unwrap();
   1530     }
   1531 
   1532     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1533     #[test]
   1534     fn migration_failure_inventory_is_complete_and_source_aware() {
   1535         use std::error::Error as _;
   1536 
   1537         let cases = [
   1538             (
   1539                 MigrationFailureKind::CatalogMismatch,
   1540                 "migration catalog does not match state",
   1541             ),
   1542             (
   1543                 MigrationFailureKind::HistoryCorrupt,
   1544                 "migration history is corrupt",
   1545             ),
   1546             (
   1547                 MigrationFailureKind::CallbackBinding,
   1548                 "migration callback binding is invalid",
   1549             ),
   1550             (
   1551                 MigrationFailureKind::Execution,
   1552                 "migration execution failed",
   1553             ),
   1554             (
   1555                 MigrationFailureKind::LedgerWrite,
   1556                 "migration ledger write failed",
   1557             ),
   1558             (
   1559                 MigrationFailureKind::MetadataAdvance,
   1560                 "migration metadata advance failed",
   1561             ),
   1562             (
   1563                 MigrationFailureKind::Commit,
   1564                 "migration commit outcome is unavailable",
   1565             ),
   1566         ];
   1567         for (kind, message) in cases {
   1568             let plain = MigrationFailure(kind);
   1569             assert_eq!(plain.to_string(), message);
   1570             assert!(plain.source().is_none());
   1571 
   1572             let sourced = MigrationSource {
   1573                 kind,
   1574                 source: Box::new(std::io::Error::other("private-cause")),
   1575             };
   1576             assert_eq!(sourced.to_string(), message);
   1577             assert!(sourced.source().is_some());
   1578             let debug = format!("{sourced:?}");
   1579             assert!(debug.contains("[redacted]"));
   1580             assert!(!debug.contains("private-cause"));
   1581         }
   1582     }
   1583 
   1584     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1585     #[test]
   1586     fn migration_condition_classifier_preserves_every_stable_kind() {
   1587         for kind in [
   1588             MigrationFailureKind::CatalogMismatch,
   1589             MigrationFailureKind::CallbackBinding,
   1590             MigrationFailureKind::HistoryCorrupt,
   1591             MigrationFailureKind::Execution,
   1592             MigrationFailureKind::LedgerWrite,
   1593             MigrationFailureKind::MetadataAdvance,
   1594             MigrationFailureKind::Commit,
   1595         ] {
   1596             assert!(require_migration_condition(true, kind).is_ok());
   1597             let error = require_migration_condition(false, kind).expect_err("failure");
   1598             assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   1599         }
   1600     }
   1601 
   1602     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1603     use std::{
   1604         num::NonZeroU32,
   1605         path::{Path, PathBuf},
   1606         sync::{
   1607             Mutex,
   1608             atomic::{AtomicUsize, Ordering as AtomicOrdering},
   1609         },
   1610     };
   1611 
   1612     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1613     use radroots_runtime_paths::{
   1614         InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver,
   1615         RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, ServiceId,
   1616     };
   1617     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1618     use radroots_storage::event::SourceGeneration;
   1619     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1620     use sqlx::sqlite::SqliteConnectOptions;
   1621 
   1622     const SQL_TWO: &str = "CREATE TABLE alpha (id INTEGER PRIMARY KEY);";
   1623     const CALLBACK_THREE: &[u8] = b"callback:rebuild_projection:v1";
   1624     const SQL_TWO_CHECKSUM: MigrationChecksum = MigrationChecksum::from_bytes([
   1625         0xd9, 0xa8, 0x5f, 0x7a, 0x59, 0x04, 0x0b, 0x3b, 0x25, 0x86, 0x56, 0x48, 0x02, 0x44, 0x10,
   1626         0x93, 0x07, 0xaa, 0x3d, 0x1a, 0x5d, 0xec, 0x04, 0x06, 0xa7, 0x50, 0x99, 0x4f, 0x17, 0xe8,
   1627         0x91, 0x13,
   1628     ]);
   1629     const CALLBACK_SQL_BYTES_CHECKSUM: MigrationChecksum = MigrationChecksum::from_bytes([
   1630         0x7a, 0x6e, 0x62, 0xf7, 0xf7, 0xa4, 0xf6, 0x1a, 0xb9, 0x14, 0x84, 0xbf, 0xe6, 0xa1, 0x2f,
   1631         0xf5, 0x0d, 0x62, 0x3d, 0x8d, 0xa2, 0x74, 0x8a, 0x16, 0xf9, 0x18, 0xd9, 0x9a, 0x52, 0xae,
   1632         0xf6, 0x15,
   1633     ]);
   1634     const CALLBACK_THREE_CHECKSUM: MigrationChecksum = MigrationChecksum::from_bytes([
   1635         0x7d, 0xca, 0x22, 0x77, 0x1b, 0x17, 0xa9, 0xf2, 0xc8, 0x04, 0x4b, 0xdc, 0xf6, 0xa6, 0xfa,
   1636         0xea, 0x41, 0x46, 0xc3, 0x56, 0xb2, 0x20, 0x17, 0xe1, 0x91, 0xd1, 0xe5, 0x42, 0xbb, 0x69,
   1637         0x47, 0x66,
   1638     ]);
   1639 
   1640     fn sql(version: u32, name: &'static str, source: &'static str) -> MigrationDescriptor {
   1641         MigrationDescriptor::sql(version, name, source, MigrationChecksum::for_sql(source))
   1642             .expect("valid SQL descriptor")
   1643     }
   1644 
   1645     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1646     fn table_object(name: &'static str, sql: &'static str) -> crate::SchemaObject {
   1647         crate::SchemaObject::new(
   1648             crate::SchemaObjectKind::Table,
   1649             name,
   1650             name,
   1651             sql,
   1652             crate::SchemaObject::computed_digest(crate::SchemaObjectKind::Table, name, name, sql)
   1653                 .expect("schema table digest"),
   1654         )
   1655         .expect("schema table")
   1656     }
   1657 
   1658     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1659     fn schema_catalog(
   1660         migrations: &MigrationCatalog,
   1661         versions: Vec<Vec<crate::SchemaObject>>,
   1662     ) -> crate::SchemaCatalog {
   1663         let versions = versions
   1664             .into_iter()
   1665             .enumerate()
   1666             .map(|(index, objects)| {
   1667                 let version = u32::try_from(index + 1).expect("schema version");
   1668                 let digest =
   1669                     crate::SchemaVersionCatalog::computed_digest(version, objects.iter().cloned())
   1670                         .expect("schema digest");
   1671                 crate::SchemaVersionCatalog::new(version, objects, digest)
   1672                     .expect("schema version catalog")
   1673             })
   1674             .collect::<Vec<_>>();
   1675         crate::SchemaCatalog::new(migrations, versions).expect("schema catalog")
   1676     }
   1677 
   1678     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1679     fn unchanged_schema_catalog(migrations: &MigrationCatalog) -> crate::SchemaCatalog {
   1680         schema_catalog(
   1681             migrations,
   1682             (0..migrations.current_version())
   1683                 .map(|_| Vec::new())
   1684                 .collect(),
   1685         )
   1686     }
   1687 
   1688     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1689     fn alpha_schema_catalog(migrations: &MigrationCatalog) -> crate::SchemaCatalog {
   1690         const ALPHA_SQL: &str = "CREATE TABLE alpha (id INTEGER PRIMARY KEY)";
   1691         let alpha = table_object("alpha", ALPHA_SQL);
   1692         let mut versions = vec![Vec::new()];
   1693         versions.extend((1..migrations.current_version()).map(|_| vec![alpha.clone()]));
   1694         schema_catalog(migrations, versions)
   1695     }
   1696 
   1697     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1698     fn alpha_beta_schema_catalog(
   1699         migrations: &MigrationCatalog,
   1700         beta_sql: &'static str,
   1701     ) -> crate::SchemaCatalog {
   1702         const ALPHA_SQL: &str = "CREATE TABLE alpha (id INTEGER PRIMARY KEY)";
   1703         let alpha = table_object("alpha", ALPHA_SQL);
   1704         let beta = table_object("beta", beta_sql);
   1705         schema_catalog(
   1706             migrations,
   1707             vec![Vec::new(), vec![alpha.clone()], vec![alpha, beta]],
   1708         )
   1709     }
   1710 
   1711     fn build_identity() -> MigrationBuildIdentity {
   1712         MigrationBuildIdentity::new(
   1713             "0.1.0-alpha",
   1714             "0123456789abcdef0123456789abcdef01234567",
   1715             "89abcdef0123456789abcdef0123456789abcdef",
   1716             "1.97.1",
   1717             "x86_64-unknown-linux-gnu",
   1718             "service-host",
   1719             1,
   1720             2,
   1721             3,
   1722             4,
   1723             5,
   1724         )
   1725         .expect("valid build identity")
   1726     }
   1727 
   1728     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1729     async fn initialized_memory_database() -> SqliteConnection {
   1730         initialized_database(SqliteConnectOptions::new().filename(":memory:")).await
   1731     }
   1732 
   1733     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1734     async fn initialized_file_database(path: &Path) -> SqliteConnection {
   1735         initialized_database(
   1736             SqliteConnectOptions::new()
   1737                 .filename(path)
   1738                 .create_if_missing(true),
   1739         )
   1740         .await
   1741     }
   1742 
   1743     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1744     async fn initialized_database(options: SqliteConnectOptions) -> SqliteConnection {
   1745         let context = RuntimeContext::resolve(
   1746             &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()),
   1747             RuntimeContextBootstrap::new(
   1748                 RadrootsPathProfile::RepoLocal,
   1749                 Some(PathBuf::from("/isolated/migration-tests")),
   1750                 RuntimeContextSource::BootstrapCli,
   1751                 RuntimeContextSource::BootstrapCli,
   1752             )
   1753             .expect("runtime bootstrap"),
   1754             ServiceId::new("myc").expect("service"),
   1755             InstanceId::new("primary").expect("instance"),
   1756         )
   1757         .expect("runtime context");
   1758         let paths =
   1759             crate::ServiceSqlitePaths::from_runtime_context(&context).expect("SQLite paths");
   1760         let metadata = crate::ServiceDatabaseMetadata::new(
   1761             &paths,
   1762             SourceGeneration::new([7; 32]).expect("generation"),
   1763             NonZeroU32::new(1).expect("schema"),
   1764             1_700_000_000_000,
   1765             crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"),
   1766         )
   1767         .expect("metadata");
   1768         let mut connection = SqliteConnection::connect_with(&options)
   1769             .await
   1770             .expect("test SQLite");
   1771         let migrations = MigrationCatalog::new([]).expect("empty migration catalog");
   1772         let schema_catalog = unchanged_schema_catalog(&migrations);
   1773         crate::metadata::write_database_metadata(&mut connection, &metadata, &schema_catalog)
   1774             .await
   1775             .expect("initialize metadata and ledger");
   1776         connection
   1777     }
   1778 
   1779     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1780     fn insert_projection_callback<'a>(
   1781         executor: &'a mut MigrationTransactionExecutor<'_>,
   1782     ) -> MigrationCallbackFuture<'a> {
   1783         Box::pin(async move { executor.execute("INSERT INTO alpha (id) VALUES (41)").await })
   1784     }
   1785 
   1786     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1787     fn pending_projection_callback<'a>(
   1788         executor: &'a mut MigrationTransactionExecutor<'_>,
   1789     ) -> MigrationCallbackFuture<'a> {
   1790         Box::pin(async move {
   1791             executor
   1792                 .execute("INSERT INTO alpha (id) VALUES (99)")
   1793                 .await?;
   1794             PENDING_CALLBACK_COUNT.fetch_add(1, AtomicOrdering::SeqCst);
   1795             core::future::pending::<Result<(), ServiceSqliteError>>().await
   1796         })
   1797     }
   1798 
   1799     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1800     fn rollback_escape_callback<'a>(
   1801         executor: &'a mut MigrationTransactionExecutor<'_>,
   1802     ) -> MigrationCallbackFuture<'a> {
   1803         Box::pin(async move {
   1804             let _ = executor
   1805                 .execute(
   1806                     "CREATE TABLE callback_rolled_back (id INTEGER PRIMARY KEY);
   1807                      ROLLBACK;
   1808                      BEGIN DEFERRED;
   1809                      CREATE TABLE callback_leaked (id INTEGER PRIMARY KEY);",
   1810                 )
   1811                 .await;
   1812             Ok(())
   1813         })
   1814     }
   1815 
   1816     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1817     fn ignored_statement_control_callback<'a>(
   1818         executor: &'a mut MigrationTransactionExecutor<'_>,
   1819     ) -> MigrationCallbackFuture<'a> {
   1820         let sql = IGNORED_STATEMENT_CONTROL_SQL
   1821             .lock()
   1822             .expect("statement-control SQL mutex")
   1823             .expect("statement-control SQL is installed");
   1824         Box::pin(async move {
   1825             let _ = executor.execute(sql).await;
   1826             Ok(())
   1827         })
   1828     }
   1829 
   1830     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1831     static PENDING_CALLBACK_COUNT: AtomicUsize = AtomicUsize::new(0);
   1832 
   1833     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1834     static IGNORED_STATEMENT_CONTROL_SQL: Mutex<Option<&'static str>> = Mutex::new(None);
   1835 
   1836     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1837     async fn replace_with_permissive_ledger(connection: &mut SqliteConnection) {
   1838         sqlx::raw_sql(
   1839             "DROP TRIGGER schema_migrations_no_update;
   1840              DROP TRIGGER schema_migrations_no_delete;
   1841              DROP TABLE schema_migrations;
   1842              CREATE TABLE schema_migrations (
   1843                 version, name, checksum, applied_at_unix_s,
   1844                 service_version, service_commit, lib_revision, rust_version, target,
   1845                 feature_profile, config_contract_version, state_contract_version,
   1846                 admin_contract_version, status_contract_version, provider_contract_version
   1847              );",
   1848         )
   1849         .execute(connection)
   1850         .await
   1851         .expect("replace ledger for corrupt-state test");
   1852     }
   1853 
   1854     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1855     async fn insert_permissive_history_row(
   1856         connection: &mut SqliteConnection,
   1857         version: i64,
   1858         name: &str,
   1859         checksum: &[u8],
   1860         applied_at: i64,
   1861         service_version: &str,
   1862     ) {
   1863         let build = build_identity();
   1864         sqlx::query(
   1865             "INSERT INTO schema_migrations (
   1866                 version, name, checksum, applied_at_unix_s,
   1867                 service_version, service_commit, lib_revision, rust_version, target,
   1868                 feature_profile, config_contract_version, state_contract_version,
   1869                 admin_contract_version, status_contract_version, provider_contract_version
   1870              ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, 2, 3, 4, 5)",
   1871         )
   1872         .bind(version)
   1873         .bind(name)
   1874         .bind(checksum)
   1875         .bind(applied_at)
   1876         .bind(service_version)
   1877         .bind(build.service_commit())
   1878         .bind(build.lib_revision())
   1879         .bind(build.rust_version())
   1880         .bind(build.target())
   1881         .bind(build.feature_profile())
   1882         .execute(connection)
   1883         .await
   1884         .expect("insert corrupt-state row");
   1885     }
   1886 
   1887     #[test]
   1888     fn applied_time_and_complete_build_identity_are_exact_and_bounded() {
   1889         assert_eq!(MigrationAppliedAtUnixSeconds::new(0).unwrap().get(), 0);
   1890         assert_eq!(
   1891             MigrationAppliedAtUnixSeconds::new(i64::MAX as u64)
   1892                 .unwrap()
   1893                 .get(),
   1894             i64::MAX as u64
   1895         );
   1896         assert_eq!(
   1897             MigrationAppliedAtUnixSeconds::new(i64::MAX as u64 + 1),
   1898             Err(MigrationEvidenceError::InvalidAppliedTime)
   1899         );
   1900 
   1901         let build = build_identity();
   1902         assert_eq!(build.service_version(), "0.1.0-alpha");
   1903         assert_eq!(
   1904             build.service_commit(),
   1905             "0123456789abcdef0123456789abcdef01234567"
   1906         );
   1907         assert_eq!(
   1908             build.lib_revision(),
   1909             "89abcdef0123456789abcdef0123456789abcdef"
   1910         );
   1911         assert_eq!(build.rust_version(), "1.97.1");
   1912         assert_eq!(build.target(), "x86_64-unknown-linux-gnu");
   1913         assert_eq!(build.feature_profile(), "service-host");
   1914         assert_eq!(
   1915             [
   1916                 build.config_contract_version(),
   1917                 build.state_contract_version(),
   1918                 build.admin_contract_version(),
   1919                 build.status_contract_version(),
   1920                 build.provider_contract_version(),
   1921             ],
   1922             [1, 2, 3, 4, 5]
   1923         );
   1924         let debug = format!("{build:?}");
   1925         for hidden in [
   1926             build.service_version(),
   1927             build.service_commit(),
   1928             build.lib_revision(),
   1929             build.rust_version(),
   1930             build.target(),
   1931             build.feature_profile(),
   1932         ] {
   1933             assert!(!debug.contains(hidden));
   1934         }
   1935 
   1936         for invalid in ["", ".bad", "bad value", "bad/value", "é"] {
   1937             assert_eq!(
   1938                 MigrationBuildIdentity::new(
   1939                     invalid,
   1940                     "0123456789abcdef0123456789abcdef01234567",
   1941                     "89abcdef0123456789abcdef0123456789abcdef",
   1942                     "1.97.1",
   1943                     "x86_64-unknown-linux-gnu",
   1944                     "service-host",
   1945                     1,
   1946                     2,
   1947                     3,
   1948                     4,
   1949                     5,
   1950                 ),
   1951                 Err(MigrationEvidenceError::InvalidBuildIdentity)
   1952             );
   1953         }
   1954         for invalid_revision in [
   1955             "0123456789abcdef0123456789abcdef0123456",
   1956             "0123456789ABCDEF0123456789abcdef01234567",
   1957             "g123456789abcdef0123456789abcdef01234567",
   1958         ] {
   1959             assert_eq!(
   1960                 MigrationBuildIdentity::new(
   1961                     "0.1.0-alpha",
   1962                     invalid_revision,
   1963                     "89abcdef0123456789abcdef0123456789abcdef",
   1964                     "1.97.1",
   1965                     "x86_64-unknown-linux-gnu",
   1966                     "service-host",
   1967                     1,
   1968                     2,
   1969                     3,
   1970                     4,
   1971                     5,
   1972                 ),
   1973                 Err(MigrationEvidenceError::InvalidBuildIdentity)
   1974             );
   1975         }
   1976         let maximum = "a".repeat(MAX_MIGRATION_BUILD_ID_UTF8_BYTES);
   1977         assert!(
   1978             MigrationBuildIdentity::new(
   1979                 &maximum,
   1980                 "0123456789abcdef0123456789abcdef01234567",
   1981                 "89abcdef0123456789abcdef0123456789abcdef",
   1982                 &maximum,
   1983                 &maximum,
   1984                 &maximum,
   1985                 1,
   1986                 2,
   1987                 3,
   1988                 4,
   1989                 5,
   1990             )
   1991             .is_ok()
   1992         );
   1993         let maximum_plus_one = "a".repeat(MAX_MIGRATION_BUILD_ID_UTF8_BYTES + 1);
   1994         for field in [0, 3, 4, 5] {
   1995             let mut values = [
   1996                 "0.1.0-alpha",
   1997                 "0123456789abcdef0123456789abcdef01234567",
   1998                 "89abcdef0123456789abcdef0123456789abcdef",
   1999                 "1.97.1",
   2000                 "x86_64-unknown-linux-gnu",
   2001                 "service-host",
   2002             ];
   2003             values[field] = &maximum_plus_one;
   2004             assert_eq!(
   2005                 MigrationBuildIdentity::new(
   2006                     values[0], values[1], values[2], values[3], values[4], values[5], 1, 2, 3, 4,
   2007                     5,
   2008                 ),
   2009                 Err(MigrationEvidenceError::InvalidBuildIdentity)
   2010             );
   2011         }
   2012         let very_large = "a".repeat(4 * 1024 * 1024);
   2013         for field in 0..6 {
   2014             let mut values = [
   2015                 "0.1.0-alpha",
   2016                 "0123456789abcdef0123456789abcdef01234567",
   2017                 "89abcdef0123456789abcdef0123456789abcdef",
   2018                 "1.97.1",
   2019                 "x86_64-unknown-linux-gnu",
   2020                 "service-host",
   2021             ];
   2022             values[field] = &very_large;
   2023             assert_eq!(
   2024                 MigrationBuildIdentity::new(
   2025                     values[0], values[1], values[2], values[3], values[4], values[5], 1, 2, 3, 4,
   2026                     5,
   2027                 ),
   2028                 Err(MigrationEvidenceError::InvalidBuildIdentity),
   2029                 "field {field} allocated before validation"
   2030             );
   2031         }
   2032         assert_eq!(
   2033             MigrationBuildIdentity::new(
   2034                 "0.1.0-alpha",
   2035                 "0123456789abcdef0123456789abcdef01234567",
   2036                 "89abcdef0123456789abcdef0123456789abcdef",
   2037                 "1.97.1",
   2038                 "x86_64-unknown-linux-gnu",
   2039                 "service-host",
   2040                 1,
   2041                 2,
   2042                 3,
   2043                 4,
   2044                 0,
   2045             ),
   2046             Err(MigrationEvidenceError::InvalidBuildIdentity)
   2047         );
   2048     }
   2049 
   2050     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2051     #[tokio::test(flavor = "current_thread")]
   2052     async fn sql_and_callback_migrations_commit_exact_restart_safe_ledger() {
   2053         let sql_descriptor =
   2054             MigrationDescriptor::sql(2, "create_alpha", SQL_TWO, SQL_TWO_CHECKSUM).unwrap();
   2055         let callback_descriptor = MigrationDescriptor::callback(
   2056             3,
   2057             "rebuild_projection",
   2058             CALLBACK_THREE,
   2059             CALLBACK_THREE_CHECKSUM,
   2060         )
   2061         .unwrap();
   2062         let callback = MigrationCallbackBinding::new(
   2063             callback_descriptor.target_version(),
   2064             callback_descriptor.name(),
   2065             callback_descriptor.checksum(),
   2066             insert_projection_callback,
   2067         );
   2068         let catalog =
   2069             MigrationCatalog::new([sql_descriptor, callback_descriptor]).expect("catalog");
   2070         let schema_catalog = alpha_schema_catalog(&catalog);
   2071         let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap();
   2072         let build = build_identity();
   2073         let mut connection = initialized_memory_database().await;
   2074         let mut validate = || Ok(());
   2075 
   2076         let outcome = apply_governed_migrations(
   2077             &mut connection,
   2078             &catalog,
   2079             &schema_catalog,
   2080             applied_at,
   2081             &build,
   2082             &[callback],
   2083             &mut validate,
   2084         )
   2085         .await
   2086         .expect("apply catalog");
   2087         assert_eq!(outcome.initial_version(), 1);
   2088         assert_eq!(outcome.final_version(), 3);
   2089         assert_eq!(outcome.applied_count(), 2);
   2090         assert_eq!(
   2091             sqlx::query_scalar::<_, i64>("SELECT id FROM alpha")
   2092                 .fetch_one(&mut connection)
   2093                 .await
   2094                 .unwrap(),
   2095             41
   2096         );
   2097         let rows = sqlx::query(
   2098             "SELECT version, name, checksum, applied_at_unix_s,
   2099                     service_version, service_commit, lib_revision, rust_version,
   2100                     target, feature_profile, config_contract_version,
   2101                     state_contract_version, admin_contract_version,
   2102                     status_contract_version, provider_contract_version
   2103              FROM schema_migrations ORDER BY version",
   2104         )
   2105         .fetch_all(&mut connection)
   2106         .await
   2107         .unwrap();
   2108         assert_eq!(rows.len(), 2);
   2109         assert_eq!(rows[0].try_get::<i64, _>("version").unwrap(), 2);
   2110         assert_eq!(
   2111             rows[0].try_get::<String, _>("name").unwrap(),
   2112             "create_alpha"
   2113         );
   2114         assert_eq!(
   2115             rows[0].try_get::<Vec<u8>, _>("checksum").unwrap(),
   2116             SQL_TWO_CHECKSUM.as_bytes().as_slice()
   2117         );
   2118         assert_eq!(rows[1].try_get::<i64, _>("version").unwrap(), 3);
   2119         assert_eq!(
   2120             rows[1].try_get::<String, _>("name").unwrap(),
   2121             "rebuild_projection"
   2122         );
   2123         for row in &rows {
   2124             assert_eq!(
   2125                 row.try_get::<i64, _>("applied_at_unix_s").unwrap(),
   2126                 1_800_000_000
   2127             );
   2128             assert_eq!(
   2129                 row.try_get::<String, _>("service_version").unwrap(),
   2130                 build.service_version()
   2131             );
   2132             assert_eq!(
   2133                 row.try_get::<String, _>("service_commit").unwrap(),
   2134                 build.service_commit()
   2135             );
   2136             assert_eq!(
   2137                 row.try_get::<String, _>("lib_revision").unwrap(),
   2138                 build.lib_revision()
   2139             );
   2140             assert_eq!(
   2141                 row.try_get::<String, _>("rust_version").unwrap(),
   2142                 build.rust_version()
   2143             );
   2144             assert_eq!(row.try_get::<String, _>("target").unwrap(), build.target());
   2145             assert_eq!(
   2146                 row.try_get::<String, _>("feature_profile").unwrap(),
   2147                 build.feature_profile()
   2148             );
   2149             assert_eq!(row.try_get::<i64, _>("config_contract_version").unwrap(), 1);
   2150             assert_eq!(row.try_get::<i64, _>("state_contract_version").unwrap(), 2);
   2151             assert_eq!(row.try_get::<i64, _>("admin_contract_version").unwrap(), 3);
   2152             assert_eq!(row.try_get::<i64, _>("status_contract_version").unwrap(), 4);
   2153             assert_eq!(
   2154                 row.try_get::<i64, _>("provider_contract_version").unwrap(),
   2155                 5
   2156             );
   2157         }
   2158         for statement in [
   2159             "UPDATE schema_migrations SET name = 'changed' WHERE version = 2",
   2160             "DELETE FROM schema_migrations WHERE version = 2",
   2161         ] {
   2162             assert!(
   2163                 sqlx::query(statement)
   2164                     .execute(&mut connection)
   2165                     .await
   2166                     .is_err(),
   2167                 "append-only ledger accepted `{statement}`"
   2168             );
   2169         }
   2170         assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 3);
   2171 
   2172         let reopened = apply_governed_migrations(
   2173             &mut connection,
   2174             &catalog,
   2175             &schema_catalog,
   2176             MigrationAppliedAtUnixSeconds::new(1_900_000_000).unwrap(),
   2177             &build,
   2178             &[callback],
   2179             &mut validate,
   2180         )
   2181         .await
   2182         .expect("lost response converges on exact committed history");
   2183         assert_eq!(reopened.initial_version(), 3);
   2184         assert_eq!(reopened.final_version(), 3);
   2185         assert_eq!(reopened.applied_count(), 0);
   2186         assert_eq!(
   2187             sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM alpha")
   2188                 .fetch_one(&mut connection)
   2189                 .await
   2190                 .unwrap(),
   2191             1
   2192         );
   2193     }
   2194 
   2195     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2196     #[tokio::test(flavor = "current_thread")]
   2197     async fn failing_step_rolls_back_only_that_step_and_exact_prefix_resumes() {
   2198         const INVALID_SQL: &str =
   2199             "CREATE TABLE broken (id INTEGER PRIMARY KEY); SELECT no_such_function();";
   2200         const RECOVERY_SQL: &str = "CREATE TABLE beta (id INTEGER PRIMARY KEY);";
   2201         let first = sql(2, "create_alpha", SQL_TWO);
   2202         let invalid = sql(3, "create_beta", INVALID_SQL);
   2203         let invalid_catalog = MigrationCatalog::new([first.clone(), invalid]).unwrap();
   2204         let invalid_schema_catalog = alpha_schema_catalog(&invalid_catalog);
   2205         let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap();
   2206         let build = build_identity();
   2207         let mut connection = initialized_memory_database().await;
   2208         let mut validate = || Ok(());
   2209 
   2210         let error = apply_governed_migrations(
   2211             &mut connection,
   2212             &invalid_catalog,
   2213             &invalid_schema_catalog,
   2214             applied_at,
   2215             &build,
   2216             &[],
   2217             &mut validate,
   2218         )
   2219         .await
   2220         .expect_err("invalid second step must roll back");
   2221         assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   2222         assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 2);
   2223         assert_eq!(
   2224             read_migration_history(&mut connection).await.unwrap().len(),
   2225             1
   2226         );
   2227         assert_eq!(
   2228             sqlx::query_scalar::<_, i64>(
   2229                 "SELECT COUNT(*) FROM sqlite_schema WHERE type = 'table' AND name = 'broken'",
   2230             )
   2231             .fetch_one(&mut connection)
   2232             .await
   2233             .unwrap(),
   2234             0
   2235         );
   2236 
   2237         let recovered_catalog =
   2238             MigrationCatalog::new([first, sql(3, "create_beta", RECOVERY_SQL)]).unwrap();
   2239         let recovered_schema_catalog = alpha_beta_schema_catalog(
   2240             &recovered_catalog,
   2241             "CREATE TABLE beta (id INTEGER PRIMARY KEY)",
   2242         );
   2243         let recovered = apply_governed_migrations(
   2244             &mut connection,
   2245             &recovered_catalog,
   2246             &recovered_schema_catalog,
   2247             applied_at,
   2248             &build,
   2249             &[],
   2250             &mut validate,
   2251         )
   2252         .await
   2253         .expect("resume exact prefix");
   2254         assert_eq!(recovered.initial_version(), 2);
   2255         assert_eq!(recovered.final_version(), 3);
   2256         assert_eq!(recovered.applied_count(), 1);
   2257     }
   2258 
   2259     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2260     #[tokio::test(flavor = "current_thread")]
   2261     async fn schema_mismatch_before_or_after_execution_never_commits_a_step() {
   2262         let catalog = MigrationCatalog::new([sql(2, "create_alpha", SQL_TWO)]).unwrap();
   2263         let wrong_target_catalog = unchanged_schema_catalog(&catalog);
   2264         let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap();
   2265         let build = build_identity();
   2266         let mut validate = || Ok(());
   2267 
   2268         let mut after_execution = initialized_memory_database().await;
   2269         let error = apply_governed_migrations(
   2270             &mut after_execution,
   2271             &catalog,
   2272             &wrong_target_catalog,
   2273             applied_at,
   2274             &build,
   2275             &[],
   2276             &mut validate,
   2277         )
   2278         .await
   2279         .expect_err("target schema mismatch must roll back");
   2280         assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity);
   2281         assert_eq!(
   2282             read_state_schema_version(&mut after_execution)
   2283                 .await
   2284                 .unwrap(),
   2285             1
   2286         );
   2287         assert!(
   2288             read_migration_history(&mut after_execution)
   2289                 .await
   2290                 .unwrap()
   2291                 .is_empty()
   2292         );
   2293         assert_eq!(
   2294             sqlx::query_scalar::<_, i64>(
   2295                 "SELECT COUNT(*) FROM sqlite_schema WHERE type = 'table' AND name = 'alpha'",
   2296             )
   2297             .fetch_one(&mut after_execution)
   2298             .await
   2299             .unwrap(),
   2300             0
   2301         );
   2302 
   2303         let mut before_execution = initialized_memory_database().await;
   2304         sqlx::query("CREATE TABLE unexpected (value INTEGER)")
   2305             .execute(&mut before_execution)
   2306             .await
   2307             .unwrap();
   2308         let expected_target = alpha_schema_catalog(&catalog);
   2309         let error = apply_governed_migrations(
   2310             &mut before_execution,
   2311             &catalog,
   2312             &expected_target,
   2313             applied_at,
   2314             &build,
   2315             &[],
   2316             &mut validate,
   2317         )
   2318         .await
   2319         .expect_err("current schema mismatch must fail before execution");
   2320         assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity);
   2321         assert_eq!(
   2322             read_state_schema_version(&mut before_execution)
   2323                 .await
   2324                 .unwrap(),
   2325             1
   2326         );
   2327         assert!(
   2328             read_migration_history(&mut before_execution)
   2329                 .await
   2330                 .unwrap()
   2331                 .is_empty()
   2332         );
   2333         assert_eq!(
   2334             sqlx::query_scalar::<_, i64>(
   2335                 "SELECT COUNT(*) FROM sqlite_schema WHERE type = 'table' AND name = 'alpha'",
   2336             )
   2337             .fetch_one(&mut before_execution)
   2338             .await
   2339             .unwrap(),
   2340             0
   2341         );
   2342     }
   2343 
   2344     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2345     #[tokio::test(flavor = "current_thread")]
   2346     async fn transaction_control_cannot_escape_schema_ledger_metadata_atomicity() {
   2347         const COMMIT_ESCAPE_SQL: &str =
   2348             "CREATE TABLE sql_leaked (id INTEGER PRIMARY KEY); COMMIT; SELECT no_such_function();";
   2349         const REPLACEMENT_ESCAPE_SQL: &str =
   2350             "CREATE TABLE sql_rolled_back (id INTEGER PRIMARY KEY);
   2351              ROLLBACK;
   2352              BEGIN DEFERRED;
   2353              CREATE TABLE sql_replacement_leaked (id INTEGER PRIMARY KEY);";
   2354         const ROLLBACK_CALLBACK_DEFINITION: &[u8] = b"callback:rollback_escape:v1";
   2355         let directory = tempfile::tempdir().unwrap();
   2356         let build = build_identity();
   2357         let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap();
   2358 
   2359         let sql_path = directory.path().join("sql-escape.sqlite");
   2360         let mut sql_connection = initialized_file_database(&sql_path).await;
   2361         let sql_catalog =
   2362             MigrationCatalog::new([sql(2, "attempt_commit_escape", COMMIT_ESCAPE_SQL)]).unwrap();
   2363         let sql_schema_catalog = unchanged_schema_catalog(&sql_catalog);
   2364         let mut validate = || Ok(());
   2365         let sql_error = apply_governed_migrations(
   2366             &mut sql_connection,
   2367             &sql_catalog,
   2368             &sql_schema_catalog,
   2369             applied_at,
   2370             &build,
   2371             &[],
   2372             &mut validate,
   2373         )
   2374         .await
   2375         .expect_err("embedded COMMIT must be rejected");
   2376         assert_eq!(sql_error.kind(), ServiceSqliteErrorKind::Migration);
   2377         drop(sql_connection);
   2378         assert_fresh_connection_has_no_migration_effect(&sql_path, "sql_leaked").await;
   2379 
   2380         let replacement_path = directory.path().join("sql-replacement-escape.sqlite");
   2381         let mut replacement_connection = initialized_file_database(&replacement_path).await;
   2382         let replacement_catalog = MigrationCatalog::new([sql(
   2383             2,
   2384             "attempt_transaction_replacement",
   2385             REPLACEMENT_ESCAPE_SQL,
   2386         )])
   2387         .unwrap();
   2388         let replacement_schema_catalog = unchanged_schema_catalog(&replacement_catalog);
   2389         let replacement_error = apply_governed_migrations(
   2390             &mut replacement_connection,
   2391             &replacement_catalog,
   2392             &replacement_schema_catalog,
   2393             applied_at,
   2394             &build,
   2395             &[],
   2396             &mut validate,
   2397         )
   2398         .await
   2399         .expect_err("replacement transaction must not inherit the governed commit permit");
   2400         assert_eq!(replacement_error.kind(), ServiceSqliteErrorKind::Migration);
   2401         drop(replacement_connection);
   2402         assert_fresh_connection_has_no_migration_effect(
   2403             &replacement_path,
   2404             "sql_replacement_leaked",
   2405         )
   2406         .await;
   2407 
   2408         let callback_path = directory.path().join("callback-escape.sqlite");
   2409         let mut callback_connection = initialized_file_database(&callback_path).await;
   2410         let callback_descriptor = MigrationDescriptor::callback(
   2411             2,
   2412             "attempt_rollback_escape",
   2413             ROLLBACK_CALLBACK_DEFINITION,
   2414             MigrationChecksum::for_callback(ROLLBACK_CALLBACK_DEFINITION),
   2415         )
   2416         .unwrap();
   2417         let callback = MigrationCallbackBinding::new(
   2418             callback_descriptor.target_version(),
   2419             callback_descriptor.name(),
   2420             callback_descriptor.checksum(),
   2421             rollback_escape_callback,
   2422         );
   2423         let callback_catalog = MigrationCatalog::new([callback_descriptor]).unwrap();
   2424         let callback_schema_catalog = unchanged_schema_catalog(&callback_catalog);
   2425         let callback_error = apply_governed_migrations(
   2426             &mut callback_connection,
   2427             &callback_catalog,
   2428             &callback_schema_catalog,
   2429             applied_at,
   2430             &build,
   2431             &[callback],
   2432             &mut validate,
   2433         )
   2434         .await
   2435         .expect_err("callback ROLLBACK must be rejected even when its error is ignored");
   2436         assert_eq!(callback_error.kind(), ServiceSqliteErrorKind::Migration);
   2437         drop(callback_connection);
   2438         assert_fresh_connection_has_no_migration_effect(&callback_path, "callback_leaked").await;
   2439 
   2440         const POLICY_ESCAPE_SQL: &str = "CREATE TABLE policy_leaked (id INTEGER PRIMARY KEY);
   2441              PRAGMA busy_timeout = 1;
   2442              ATTACH ':memory:' AS escaped;";
   2443         let policy_path = directory.path().join("policy-escape.sqlite");
   2444         let mut policy_connection = initialized_file_database(&policy_path).await;
   2445         let policy_catalog =
   2446             MigrationCatalog::new([sql(2, "attempt_policy_escape", POLICY_ESCAPE_SQL)]).unwrap();
   2447         let policy_schema_catalog = unchanged_schema_catalog(&policy_catalog);
   2448         let policy_error = apply_governed_migrations(
   2449             &mut policy_connection,
   2450             &policy_catalog,
   2451             &policy_schema_catalog,
   2452             applied_at,
   2453             &build,
   2454             &[],
   2455             &mut validate,
   2456         )
   2457         .await
   2458         .expect_err("connection-policy or attachment drift must block commit");
   2459         assert_eq!(policy_error.kind(), ServiceSqliteErrorKind::Migration);
   2460         drop(policy_connection);
   2461         assert_fresh_connection_has_no_migration_effect(&policy_path, "policy_leaked").await;
   2462     }
   2463 
   2464     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2465     #[tokio::test(flavor = "current_thread")]
   2466     async fn migration_sql_rejects_complete_statement_control_inventory_before_execution() {
   2467         for (name, statement) in [
   2468             (
   2469                 "reject_pragma",
   2470                 "CREATE TABLE leaked (id INTEGER); /* policy */ PrAgMa\ntrusted_schema=ON",
   2471             ),
   2472             (
   2473                 "reject_attach",
   2474                 "CREATE TABLE leaked (id INTEGER); ATTACH ':memory:' AS escaped",
   2475             ),
   2476             (
   2477                 "reject_detach",
   2478                 "CREATE TABLE leaked (id INTEGER); DETACH DATABASE escaped",
   2479             ),
   2480             (
   2481                 "reject_begin",
   2482                 "CREATE TABLE leaked (id INTEGER); BEGIN DEFERRED",
   2483             ),
   2484             ("reject_commit", "CREATE TABLE leaked (id INTEGER); COMMIT"),
   2485             (
   2486                 "reject_end",
   2487                 "CREATE TABLE leaked (id INTEGER); END TRANSACTION",
   2488             ),
   2489             (
   2490                 "reject_rollback",
   2491                 "CREATE TABLE leaked (id INTEGER); ROLLBACK",
   2492             ),
   2493             (
   2494                 "reject_savepoint",
   2495                 "CREATE TABLE leaked (id INTEGER); SAVEPOINT escaped",
   2496             ),
   2497             (
   2498                 "reject_release",
   2499                 "CREATE TABLE leaked (id INTEGER); RELEASE SAVEPOINT escaped",
   2500             ),
   2501         ] {
   2502             let directory = tempfile::tempdir().unwrap();
   2503             let database_path = directory.path().join("statement-control.sqlite");
   2504             let mut connection = initialized_file_database(&database_path).await;
   2505             let catalog = MigrationCatalog::new([sql(2, name, statement)]).unwrap();
   2506             let schema_catalog = unchanged_schema_catalog(&catalog);
   2507             let mut validate = || Ok(());
   2508             let error = apply_governed_migrations(
   2509                 &mut connection,
   2510                 &catalog,
   2511                 &schema_catalog,
   2512                 MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(),
   2513                 &build_identity(),
   2514                 &[],
   2515                 &mut validate,
   2516             )
   2517             .await
   2518             .expect_err("statement-control migration must be rejected");
   2519             assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   2520             assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 1);
   2521             assert!(
   2522                 read_migration_history(&mut connection)
   2523                     .await
   2524                     .unwrap()
   2525                     .is_empty()
   2526             );
   2527             assert_eq!(
   2528                 sqlx::query_scalar::<_, i64>(
   2529                     "SELECT COUNT(*) FROM sqlite_schema WHERE type = 'table' AND name = 'leaked'",
   2530                 )
   2531                 .fetch_one(&mut connection)
   2532                 .await
   2533                 .unwrap(),
   2534                 0
   2535             );
   2536         }
   2537     }
   2538 
   2539     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2540     #[tokio::test(flavor = "current_thread")]
   2541     async fn migration_executor_rejects_transient_attachment_before_file_creation() {
   2542         let directory = tempfile::tempdir().unwrap();
   2543         let database_path = directory.path().join("main.sqlite");
   2544         let external_path = directory.path().join("external.sqlite");
   2545         let migration_sql = Box::leak(
   2546             format!(
   2547                 "ATTACH DATABASE '{}' AS extra; DETACH DATABASE extra",
   2548                 external_path.display()
   2549             )
   2550             .into_boxed_str(),
   2551         );
   2552         let mut connection = initialized_file_database(&database_path).await;
   2553         let catalog =
   2554             MigrationCatalog::new([sql(2, "reject_transient_attachment", migration_sql)]).unwrap();
   2555         let schema_catalog = unchanged_schema_catalog(&catalog);
   2556         let mut validate = || Ok(());
   2557         let error = apply_governed_migrations(
   2558             &mut connection,
   2559             &catalog,
   2560             &schema_catalog,
   2561             MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(),
   2562             &build_identity(),
   2563             &[],
   2564             &mut validate,
   2565         )
   2566         .await
   2567         .expect_err("migration ATTACH/DETACH must fail before SQLite compilation");
   2568         assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   2569         assert!(!external_path.exists());
   2570         assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 1);
   2571         assert!(
   2572             read_migration_history(&mut connection)
   2573                 .await
   2574                 .unwrap()
   2575                 .is_empty()
   2576         );
   2577 
   2578         const CALLBACK_DEFINITION: &[u8] = b"callback:reject_transient_attachment:v1";
   2579         let callback_database_path = directory.path().join("callback-main.sqlite");
   2580         let callback_external_path = directory.path().join("callback-external.sqlite");
   2581         let callback_sql = Box::leak(
   2582             format!(
   2583                 "ATTACH DATABASE '{}' AS extra; DETACH DATABASE extra",
   2584                 callback_external_path.display()
   2585             )
   2586             .into_boxed_str(),
   2587         );
   2588         *IGNORED_STATEMENT_CONTROL_SQL
   2589             .lock()
   2590             .expect("statement-control SQL mutex") = Some(callback_sql);
   2591         let mut callback_connection = initialized_file_database(&callback_database_path).await;
   2592         let callback_descriptor = MigrationDescriptor::callback(
   2593             2,
   2594             "reject_callback_attachment",
   2595             CALLBACK_DEFINITION,
   2596             MigrationChecksum::for_callback(CALLBACK_DEFINITION),
   2597         )
   2598         .unwrap();
   2599         let callback = MigrationCallbackBinding::new(
   2600             callback_descriptor.target_version(),
   2601             callback_descriptor.name(),
   2602             callback_descriptor.checksum(),
   2603             ignored_statement_control_callback,
   2604         );
   2605         let callback_catalog = MigrationCatalog::new([callback_descriptor]).unwrap();
   2606         let callback_schema_catalog = unchanged_schema_catalog(&callback_catalog);
   2607         let callback_error = apply_governed_migrations(
   2608             &mut callback_connection,
   2609             &callback_catalog,
   2610             &callback_schema_catalog,
   2611             MigrationAppliedAtUnixSeconds::new(1_800_000_001).unwrap(),
   2612             &build_identity(),
   2613             &[callback],
   2614             &mut validate,
   2615         )
   2616         .await
   2617         .expect_err("ignored callback ATTACH/DETACH must still fail the migration");
   2618         *IGNORED_STATEMENT_CONTROL_SQL
   2619             .lock()
   2620             .expect("statement-control SQL mutex") = None;
   2621         assert_eq!(callback_error.kind(), ServiceSqliteErrorKind::Migration);
   2622         assert!(!callback_external_path.exists());
   2623         assert_eq!(
   2624             read_state_schema_version(&mut callback_connection)
   2625                 .await
   2626                 .unwrap(),
   2627             1
   2628         );
   2629         assert!(
   2630             read_migration_history(&mut callback_connection)
   2631                 .await
   2632                 .unwrap()
   2633                 .is_empty()
   2634         );
   2635     }
   2636 
   2637     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2638     async fn assert_fresh_connection_has_no_migration_effect(path: &Path, table: &str) {
   2639         let mut connection = SqliteConnection::connect_with(
   2640             &SqliteConnectOptions::new()
   2641                 .filename(path)
   2642                 .create_if_missing(false),
   2643         )
   2644         .await
   2645         .expect("fresh verification connection");
   2646         assert_eq!(
   2647             sqlx::query_scalar::<_, i64>(
   2648                 "SELECT COUNT(*) FROM sqlite_schema WHERE type = 'table' AND name = ?",
   2649             )
   2650             .bind(table)
   2651             .fetch_one(&mut connection)
   2652             .await
   2653             .unwrap(),
   2654             0
   2655         );
   2656         assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 1);
   2657         assert!(
   2658             read_migration_history(&mut connection)
   2659                 .await
   2660                 .unwrap()
   2661                 .is_empty()
   2662         );
   2663     }
   2664 
   2665     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2666     #[tokio::test(flavor = "current_thread")]
   2667     async fn cancellation_before_commit_leaves_an_exact_resumable_prefix() {
   2668         PENDING_CALLBACK_COUNT.store(0, AtomicOrdering::SeqCst);
   2669         let sql_descriptor =
   2670             MigrationDescriptor::sql(2, "create_alpha", SQL_TWO, SQL_TWO_CHECKSUM).unwrap();
   2671         let callback_descriptor = MigrationDescriptor::callback(
   2672             3,
   2673             "rebuild_projection",
   2674             CALLBACK_THREE,
   2675             CALLBACK_THREE_CHECKSUM,
   2676         )
   2677         .unwrap();
   2678         let catalog = MigrationCatalog::new([sql_descriptor, callback_descriptor.clone()]).unwrap();
   2679         let schema_catalog = alpha_schema_catalog(&catalog);
   2680         let pending_binding = MigrationCallbackBinding::new(
   2681             3,
   2682             callback_descriptor.name(),
   2683             callback_descriptor.checksum(),
   2684             pending_projection_callback,
   2685         );
   2686         let working_binding = MigrationCallbackBinding::new(
   2687             3,
   2688             callback_descriptor.name(),
   2689             callback_descriptor.checksum(),
   2690             insert_projection_callback,
   2691         );
   2692         let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap();
   2693         let build = build_identity();
   2694         let mut connection = initialized_memory_database().await;
   2695         let mut validate = || Ok(());
   2696         let pending_bindings = [pending_binding];
   2697 
   2698         let mut application = Box::pin(apply_governed_migrations(
   2699             &mut connection,
   2700             &catalog,
   2701             &schema_catalog,
   2702             applied_at,
   2703             &build,
   2704             &pending_bindings,
   2705             &mut validate,
   2706         ));
   2707         tokio::select! {
   2708             outcome = &mut application => panic!("pending callback completed: {outcome:?}"),
   2709             () = async {
   2710                 while PENDING_CALLBACK_COUNT.load(AtomicOrdering::SeqCst) == 0 {
   2711                     tokio::task::yield_now().await;
   2712                 }
   2713             } => {}
   2714         }
   2715         drop(application);
   2716 
   2717         assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 2);
   2718         assert_eq!(
   2719             read_migration_history(&mut connection).await.unwrap().len(),
   2720             1
   2721         );
   2722         assert_eq!(
   2723             sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM alpha")
   2724                 .fetch_one(&mut connection)
   2725                 .await
   2726                 .unwrap(),
   2727             0,
   2728             "callback write survived cancellation before commit"
   2729         );
   2730         let recovered = apply_governed_migrations(
   2731             &mut connection,
   2732             &catalog,
   2733             &schema_catalog,
   2734             applied_at,
   2735             &build,
   2736             &[working_binding],
   2737             &mut validate,
   2738         )
   2739         .await
   2740         .expect("resume cancelled callback");
   2741         assert_eq!(recovered.initial_version(), 2);
   2742         assert_eq!(recovered.final_version(), 3);
   2743         assert_eq!(recovered.applied_count(), 1);
   2744         assert_eq!(
   2745             sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM alpha")
   2746                 .fetch_one(&mut connection)
   2747                 .await
   2748                 .unwrap(),
   2749             1
   2750         );
   2751     }
   2752 
   2753     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2754     #[tokio::test(flavor = "current_thread")]
   2755     async fn commit_response_loss_is_resolved_from_history_without_replay() {
   2756         let first = sql(2, "create_alpha", SQL_TWO);
   2757         let second = sql(3, "create_beta", "CREATE TABLE beta (id INTEGER);");
   2758         let catalog = MigrationCatalog::new([first, second]).unwrap();
   2759         let schema_catalog = alpha_beta_schema_catalog(&catalog, "CREATE TABLE beta (id INTEGER)");
   2760         let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap();
   2761         let build = build_identity();
   2762         let mut connection = initialized_memory_database().await;
   2763         let mut validate = || Ok(());
   2764         let mut lose_first_response = || Err(migration_error(MigrationFailureKind::Commit));
   2765 
   2766         let error = apply_governed_migrations_with_observer(
   2767             &mut connection,
   2768             &catalog,
   2769             &schema_catalog,
   2770             applied_at,
   2771             &build,
   2772             &[],
   2773             &mut validate,
   2774             &mut lose_first_response,
   2775         )
   2776         .await
   2777         .expect_err("first committed response is lost");
   2778         assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   2779         assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 2);
   2780         assert_eq!(
   2781             read_migration_history(&mut connection).await.unwrap().len(),
   2782             1
   2783         );
   2784 
   2785         let resumed = apply_governed_migrations(
   2786             &mut connection,
   2787             &catalog,
   2788             &schema_catalog,
   2789             applied_at,
   2790             &build,
   2791             &[],
   2792             &mut validate,
   2793         )
   2794         .await
   2795         .expect("resolve committed prefix and continue");
   2796         assert_eq!(resumed.initial_version(), 2);
   2797         assert_eq!(resumed.final_version(), 3);
   2798         assert_eq!(resumed.applied_count(), 1);
   2799         assert_eq!(
   2800             sqlx::query_scalar::<_, i64>(
   2801                 "SELECT COUNT(*) FROM sqlite_schema WHERE type = 'table' AND name = 'alpha'",
   2802             )
   2803             .fetch_one(&mut connection)
   2804             .await
   2805             .unwrap(),
   2806             1
   2807         );
   2808     }
   2809 
   2810     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2811     #[test]
   2812     fn callback_binding_validation_rejects_each_independent_identity_drift() {
   2813         let callback = MigrationDescriptor::callback(
   2814             2,
   2815             "rebuild_projection",
   2816             CALLBACK_THREE,
   2817             CALLBACK_THREE_CHECKSUM,
   2818         )
   2819         .expect("callback descriptor");
   2820         let callback_catalog = MigrationCatalog::new([callback.clone()]).expect("catalog");
   2821 
   2822         for binding in [
   2823             MigrationCallbackBinding::new(
   2824                 2,
   2825                 MigrationName::new("wrong_projection").expect("name"),
   2826                 callback.checksum(),
   2827                 insert_projection_callback,
   2828             ),
   2829             MigrationCallbackBinding::new(
   2830                 2,
   2831                 callback.name(),
   2832                 MigrationChecksum::from_bytes([0x55; 32]),
   2833                 insert_projection_callback,
   2834             ),
   2835         ] {
   2836             assert_eq!(
   2837                 validate_callback_bindings(&callback_catalog, &[binding])
   2838                     .expect_err("binding drift must fail")
   2839                     .kind(),
   2840                 ServiceSqliteErrorKind::Migration
   2841             );
   2842         }
   2843 
   2844         let sql = MigrationDescriptor::sql(2, "create_alpha", SQL_TWO, SQL_TWO_CHECKSUM)
   2845             .expect("SQL descriptor");
   2846         let callback_three = MigrationDescriptor::callback(
   2847             3,
   2848             "rebuild_projection",
   2849             CALLBACK_THREE,
   2850             CALLBACK_THREE_CHECKSUM,
   2851         )
   2852         .expect("callback descriptor");
   2853         let mixed_catalog = MigrationCatalog::new([sql.clone(), callback_three]).expect("catalog");
   2854         let wrong_kind = MigrationCallbackBinding::new(
   2855             2,
   2856             sql.name(),
   2857             sql.checksum(),
   2858             insert_projection_callback,
   2859         );
   2860         assert_eq!(
   2861             validate_callback_bindings(&mixed_catalog, &[wrong_kind])
   2862                 .expect_err("SQL descriptor cannot bind a callback")
   2863                 .kind(),
   2864             ServiceSqliteErrorKind::Migration
   2865         );
   2866 
   2867         const CALLBACK_FOUR: &[u8] = b"callback:rebuild_secondary_projection:v1";
   2868         let callback_two = MigrationDescriptor::callback(
   2869             2,
   2870             "rebuild_projection",
   2871             CALLBACK_THREE,
   2872             CALLBACK_THREE_CHECKSUM,
   2873         )
   2874         .expect("callback descriptor");
   2875         let callback_four = MigrationDescriptor::callback(
   2876             3,
   2877             "rebuild_secondary_projection",
   2878             CALLBACK_FOUR,
   2879             MigrationChecksum::for_callback(CALLBACK_FOUR),
   2880         )
   2881         .expect("callback descriptor");
   2882         let duplicate_target_catalog =
   2883             MigrationCatalog::new([callback_two.clone(), callback_four]).expect("catalog");
   2884         let duplicate = MigrationCallbackBinding::new(
   2885             2,
   2886             callback_two.name(),
   2887             callback_two.checksum(),
   2888             insert_projection_callback,
   2889         );
   2890         assert_eq!(
   2891             validate_callback_bindings(&duplicate_target_catalog, &[duplicate, duplicate])
   2892                 .expect_err("duplicate callback target must fail")
   2893                 .kind(),
   2894             ServiceSqliteErrorKind::Migration
   2895         );
   2896     }
   2897 
   2898     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2899     #[tokio::test(flavor = "current_thread")]
   2900     async fn callback_bindings_and_history_mismatches_fail_before_replay() {
   2901         let callback_descriptor = MigrationDescriptor::callback(
   2902             2,
   2903             "rebuild_projection",
   2904             CALLBACK_THREE,
   2905             CALLBACK_THREE_CHECKSUM,
   2906         )
   2907         .unwrap();
   2908         let catalog = MigrationCatalog::new([callback_descriptor.clone()]).unwrap();
   2909         let schema_catalog = unchanged_schema_catalog(&catalog);
   2910         let correct = MigrationCallbackBinding::new(
   2911             2,
   2912             callback_descriptor.name(),
   2913             callback_descriptor.checksum(),
   2914             insert_projection_callback,
   2915         );
   2916         let wrong = MigrationCallbackBinding::new(
   2917             3,
   2918             callback_descriptor.name(),
   2919             callback_descriptor.checksum(),
   2920             insert_projection_callback,
   2921         );
   2922         let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap();
   2923         let build = build_identity();
   2924 
   2925         for bindings in [Vec::new(), vec![wrong], vec![correct, correct]] {
   2926             let mut connection = initialized_memory_database().await;
   2927             let mut validate = || Ok(());
   2928             let error = apply_governed_migrations(
   2929                 &mut connection,
   2930                 &catalog,
   2931                 &schema_catalog,
   2932                 applied_at,
   2933                 &build,
   2934                 &bindings,
   2935                 &mut validate,
   2936             )
   2937             .await
   2938             .expect_err("callback registry mismatch");
   2939             assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   2940             assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 1);
   2941             assert!(
   2942                 read_migration_history(&mut connection)
   2943                     .await
   2944                     .unwrap()
   2945                     .is_empty()
   2946             );
   2947         }
   2948 
   2949         let mut connection = initialized_memory_database().await;
   2950         sqlx::query(
   2951             "INSERT INTO schema_migrations (
   2952                 version, name, checksum, applied_at_unix_s,
   2953                 service_version, service_commit, lib_revision, rust_version, target,
   2954                 feature_profile, config_contract_version, state_contract_version,
   2955                 admin_contract_version, status_contract_version, provider_contract_version
   2956              ) VALUES (2, 'wrong_name', zeroblob(32), 0, ?, ?, ?, ?, ?, ?, 1, 2, 3, 4, 5)",
   2957         )
   2958         .bind(build.service_version())
   2959         .bind(build.service_commit())
   2960         .bind(build.lib_revision())
   2961         .bind(build.rust_version())
   2962         .bind(build.target())
   2963         .bind(build.feature_profile())
   2964         .execute(&mut connection)
   2965         .await
   2966         .unwrap();
   2967         sqlx::query(
   2968             "UPDATE radroots_service_metadata SET state_schema_version = 2 WHERE singleton = 1",
   2969         )
   2970         .execute(&mut connection)
   2971         .await
   2972         .unwrap();
   2973         let error = verify_migration_history(&mut connection, &catalog, &schema_catalog, true)
   2974             .await
   2975             .expect_err("mismatched history");
   2976         assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   2977     }
   2978 
   2979     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2980     #[tokio::test(flavor = "current_thread")]
   2981     async fn missing_extra_reordered_newer_and_corrupt_history_fail_closed() {
   2982         let first = sql(2, "create_alpha", SQL_TWO);
   2983         let second = sql(3, "create_beta", "CREATE TABLE beta (id INTEGER);");
   2984         let catalog = MigrationCatalog::new([first.clone(), second.clone()]).unwrap();
   2985         let schema_catalog = unchanged_schema_catalog(&catalog);
   2986 
   2987         let mut missing = initialized_memory_database().await;
   2988         sqlx::query(
   2989             "UPDATE radroots_service_metadata SET state_schema_version = 2 WHERE singleton = 1",
   2990         )
   2991         .execute(&mut missing)
   2992         .await
   2993         .unwrap();
   2994         assert_eq!(
   2995             verify_migration_history(&mut missing, &catalog, &schema_catalog, false)
   2996                 .await
   2997                 .expect_err("missing row")
   2998                 .kind(),
   2999             ServiceSqliteErrorKind::Migration
   3000         );
   3001 
   3002         let mut extra = initialized_memory_database().await;
   3003         replace_with_permissive_ledger(&mut extra).await;
   3004         insert_permissive_history_row(
   3005             &mut extra,
   3006             2,
   3007             first.name().as_str(),
   3008             first.checksum().as_bytes(),
   3009             0,
   3010             "0.1.0-alpha",
   3011         )
   3012         .await;
   3013         assert_eq!(
   3014             verify_migration_history(&mut extra, &catalog, &schema_catalog, false)
   3015                 .await
   3016                 .expect_err("extra row")
   3017                 .kind(),
   3018             ServiceSqliteErrorKind::Migration
   3019         );
   3020 
   3021         let mut reordered = initialized_memory_database().await;
   3022         replace_with_permissive_ledger(&mut reordered).await;
   3023         insert_permissive_history_row(
   3024             &mut reordered,
   3025             2,
   3026             second.name().as_str(),
   3027             second.checksum().as_bytes(),
   3028             0,
   3029             "0.1.0-alpha",
   3030         )
   3031         .await;
   3032         insert_permissive_history_row(
   3033             &mut reordered,
   3034             3,
   3035             first.name().as_str(),
   3036             first.checksum().as_bytes(),
   3037             0,
   3038             "0.1.0-alpha",
   3039         )
   3040         .await;
   3041         sqlx::query(
   3042             "UPDATE radroots_service_metadata SET state_schema_version = 3 WHERE singleton = 1",
   3043         )
   3044         .execute(&mut reordered)
   3045         .await
   3046         .unwrap();
   3047         assert_eq!(
   3048             verify_migration_history(&mut reordered, &catalog, &schema_catalog, true)
   3049                 .await
   3050                 .expect_err("reordered names and checksums")
   3051                 .kind(),
   3052             ServiceSqliteErrorKind::Migration
   3053         );
   3054 
   3055         let mut newer = initialized_memory_database().await;
   3056         sqlx::query(
   3057             "UPDATE radroots_service_metadata SET state_schema_version = 4 WHERE singleton = 1",
   3058         )
   3059         .execute(&mut newer)
   3060         .await
   3061         .unwrap();
   3062         assert_eq!(
   3063             verify_migration_history(&mut newer, &catalog, &schema_catalog, false)
   3064                 .await
   3065                 .expect_err("newer schema")
   3066                 .kind(),
   3067             ServiceSqliteErrorKind::Migration
   3068         );
   3069 
   3070         for (applied_at, service_version) in [(0_i64, "bad value"), (-1, "0.1.0-alpha")] {
   3071             let mut corrupt = initialized_memory_database().await;
   3072             replace_with_permissive_ledger(&mut corrupt).await;
   3073             insert_permissive_history_row(
   3074                 &mut corrupt,
   3075                 2,
   3076                 first.name().as_str(),
   3077                 first.checksum().as_bytes(),
   3078                 applied_at,
   3079                 service_version,
   3080             )
   3081             .await;
   3082             sqlx::query(
   3083                 "UPDATE radroots_service_metadata SET state_schema_version = 2 WHERE singleton = 1",
   3084             )
   3085             .execute(&mut corrupt)
   3086             .await
   3087             .unwrap();
   3088             assert_eq!(
   3089                 verify_migration_history(&mut corrupt, &catalog, &schema_catalog, false)
   3090                     .await
   3091                     .expect_err("corrupt time or build")
   3092                     .kind(),
   3093                 ServiceSqliteErrorKind::Migration
   3094             );
   3095         }
   3096     }
   3097 
   3098     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3099     #[tokio::test(flavor = "current_thread")]
   3100     async fn every_migration_ledger_projection_rejects_wrong_storage_and_values() {
   3101         let descriptor = sql(2, "create_alpha", SQL_TWO);
   3102         let catalog = MigrationCatalog::new([descriptor.clone()]).unwrap();
   3103         let schema_catalog = unchanged_schema_catalog(&catalog);
   3104         let wrong_storage_updates = [
   3105             "UPDATE schema_migrations SET version = '2' WHERE rowid = 1",
   3106             "UPDATE schema_migrations SET name = 2 WHERE rowid = 1",
   3107             "UPDATE schema_migrations SET checksum = 'checksum' WHERE rowid = 1",
   3108             "UPDATE schema_migrations SET applied_at_unix_s = '0' WHERE rowid = 1",
   3109             "UPDATE schema_migrations SET service_version = 2 WHERE rowid = 1",
   3110             "UPDATE schema_migrations SET service_commit = 2 WHERE rowid = 1",
   3111             "UPDATE schema_migrations SET lib_revision = 2 WHERE rowid = 1",
   3112             "UPDATE schema_migrations SET rust_version = 2 WHERE rowid = 1",
   3113             "UPDATE schema_migrations SET target = 2 WHERE rowid = 1",
   3114             "UPDATE schema_migrations SET feature_profile = 2 WHERE rowid = 1",
   3115             "UPDATE schema_migrations SET config_contract_version = '1' WHERE rowid = 1",
   3116             "UPDATE schema_migrations SET state_contract_version = '2' WHERE rowid = 1",
   3117             "UPDATE schema_migrations SET admin_contract_version = '3' WHERE rowid = 1",
   3118             "UPDATE schema_migrations SET status_contract_version = '4' WHERE rowid = 1",
   3119             "UPDATE schema_migrations SET provider_contract_version = '5' WHERE rowid = 1",
   3120         ];
   3121         let invalid_value_updates = [
   3122             "UPDATE schema_migrations SET version = -1 WHERE rowid = 1",
   3123             "UPDATE schema_migrations SET version = 4294967296 WHERE rowid = 1",
   3124             "UPDATE schema_migrations SET name = 'BadName' WHERE rowid = 1",
   3125             "UPDATE schema_migrations SET checksum = zeroblob(31) WHERE rowid = 1",
   3126             "UPDATE schema_migrations SET applied_at_unix_s = -1 WHERE rowid = 1",
   3127             "UPDATE schema_migrations SET service_version = 'bad value' WHERE rowid = 1",
   3128             "UPDATE schema_migrations SET service_commit = 'bad' WHERE rowid = 1",
   3129             "UPDATE schema_migrations SET lib_revision = 'bad' WHERE rowid = 1",
   3130             "UPDATE schema_migrations SET rust_version = 'bad value' WHERE rowid = 1",
   3131             "UPDATE schema_migrations SET target = 'bad value' WHERE rowid = 1",
   3132             "UPDATE schema_migrations SET feature_profile = 'bad value' WHERE rowid = 1",
   3133             "UPDATE schema_migrations SET config_contract_version = 0 WHERE rowid = 1",
   3134             "UPDATE schema_migrations SET state_contract_version = -1 WHERE rowid = 1",
   3135             "UPDATE schema_migrations SET admin_contract_version = 4294967296 WHERE rowid = 1",
   3136             "UPDATE schema_migrations SET status_contract_version = 0 WHERE rowid = 1",
   3137             "UPDATE schema_migrations SET provider_contract_version = 0 WHERE rowid = 1",
   3138         ];
   3139 
   3140         for update in wrong_storage_updates
   3141             .into_iter()
   3142             .chain(invalid_value_updates)
   3143         {
   3144             let mut connection = initialized_memory_database().await;
   3145             replace_with_permissive_ledger(&mut connection).await;
   3146             insert_permissive_history_row(
   3147                 &mut connection,
   3148                 2,
   3149                 descriptor.name().as_str(),
   3150                 descriptor.checksum().as_bytes(),
   3151                 0,
   3152                 "0.1.0-alpha",
   3153             )
   3154             .await;
   3155             sqlx::raw_sql(update)
   3156                 .execute(&mut connection)
   3157                 .await
   3158                 .expect("corrupt ledger projection");
   3159             sqlx::query(
   3160                 "UPDATE radroots_service_metadata SET state_schema_version = 2 WHERE singleton = 1",
   3161             )
   3162             .execute(&mut connection)
   3163             .await
   3164             .unwrap();
   3165             assert_eq!(
   3166                 verify_migration_history(&mut connection, &catalog, &schema_catalog, false)
   3167                     .await
   3168                     .expect_err("corrupt ledger projection")
   3169                     .kind(),
   3170                 ServiceSqliteErrorKind::Migration,
   3171                 "accepted corrupt projection update `{update}`"
   3172             );
   3173         }
   3174     }
   3175 
   3176     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3177     #[tokio::test(flavor = "current_thread")]
   3178     async fn oversized_corrupt_history_is_bounded_before_decode() {
   3179         let descriptor = sql(2, "create_alpha", SQL_TWO);
   3180         let catalog = MigrationCatalog::new([descriptor.clone()]).unwrap();
   3181         let schema_catalog = unchanged_schema_catalog(&catalog);
   3182         let oversized_text = "a".repeat(4 * 1024 * 1024);
   3183         for (column, update) in [
   3184             (
   3185                 "name",
   3186                 "UPDATE schema_migrations SET name = ? WHERE version = 2",
   3187             ),
   3188             (
   3189                 "service_version",
   3190                 "UPDATE schema_migrations SET service_version = ? WHERE version = 2",
   3191             ),
   3192             (
   3193                 "service_commit",
   3194                 "UPDATE schema_migrations SET service_commit = ? WHERE version = 2",
   3195             ),
   3196             (
   3197                 "lib_revision",
   3198                 "UPDATE schema_migrations SET lib_revision = ? WHERE version = 2",
   3199             ),
   3200             (
   3201                 "rust_version",
   3202                 "UPDATE schema_migrations SET rust_version = ? WHERE version = 2",
   3203             ),
   3204             (
   3205                 "target",
   3206                 "UPDATE schema_migrations SET target = ? WHERE version = 2",
   3207             ),
   3208             (
   3209                 "feature_profile",
   3210                 "UPDATE schema_migrations SET feature_profile = ? WHERE version = 2",
   3211             ),
   3212         ] {
   3213             let mut connection = initialized_memory_database().await;
   3214             replace_with_permissive_ledger(&mut connection).await;
   3215             insert_permissive_history_row(
   3216                 &mut connection,
   3217                 2,
   3218                 descriptor.name().as_str(),
   3219                 descriptor.checksum().as_bytes(),
   3220                 0,
   3221                 "0.1.0-alpha",
   3222             )
   3223             .await;
   3224             sqlx::query(update)
   3225                 .bind(&oversized_text)
   3226                 .execute(&mut connection)
   3227                 .await
   3228                 .unwrap();
   3229             sqlx::query(
   3230                 "UPDATE radroots_service_metadata SET state_schema_version = 2 WHERE singleton = 1",
   3231             )
   3232             .execute(&mut connection)
   3233             .await
   3234             .unwrap();
   3235             assert_eq!(
   3236                 verify_migration_history(&mut connection, &catalog, &schema_catalog, true)
   3237                     .await
   3238                     .expect_err("oversized text must fail before decode")
   3239                     .kind(),
   3240                 ServiceSqliteErrorKind::Migration,
   3241                 "column {column}"
   3242             );
   3243         }
   3244 
   3245         let mut checksum = initialized_memory_database().await;
   3246         replace_with_permissive_ledger(&mut checksum).await;
   3247         insert_permissive_history_row(
   3248             &mut checksum,
   3249             2,
   3250             descriptor.name().as_str(),
   3251             &vec![0_u8; 4 * 1024 * 1024],
   3252             0,
   3253             "0.1.0-alpha",
   3254         )
   3255         .await;
   3256         sqlx::query(
   3257             "UPDATE radroots_service_metadata SET state_schema_version = 2 WHERE singleton = 1",
   3258         )
   3259         .execute(&mut checksum)
   3260         .await
   3261         .unwrap();
   3262         assert_eq!(
   3263             verify_migration_history(&mut checksum, &catalog, &schema_catalog, true)
   3264                 .await
   3265                 .expect_err("oversized checksum must fail before decode")
   3266                 .kind(),
   3267             ServiceSqliteErrorKind::Migration
   3268         );
   3269     }
   3270 
   3271     #[test]
   3272     fn exact_content_checksums_are_deterministic_and_kind_separated() {
   3273         let sql = MigrationChecksum::for_sql(SQL_TWO);
   3274         let callback = MigrationChecksum::for_callback(SQL_TWO.as_bytes());
   3275         assert_eq!(sql, SQL_TWO_CHECKSUM);
   3276         assert_eq!(callback, CALLBACK_SQL_BYTES_CHECKSUM);
   3277         assert_eq!(sql, MigrationChecksum::for_sql(SQL_TWO));
   3278         assert_ne!(sql, callback);
   3279         assert_ne!(
   3280             sql,
   3281             MigrationChecksum::for_sql(" CREATE TABLE alpha (id INTEGER PRIMARY KEY);")
   3282         );
   3283         assert_ne!(
   3284             sql,
   3285             MigrationChecksum::for_sql("CREATE TABLE alpha (id INTEGER PRIMARY KEY);\n")
   3286         );
   3287         assert_eq!(
   3288             MigrationDescriptor::sql(2, "create_alpha", SQL_TWO, SQL_TWO_CHECKSUM)
   3289                 .unwrap()
   3290                 .checksum(),
   3291             SQL_TWO_CHECKSUM
   3292         );
   3293         assert_eq!(
   3294             MigrationDescriptor::callback(
   3295                 2,
   3296                 "callback_alpha",
   3297                 SQL_TWO.as_bytes(),
   3298                 CALLBACK_SQL_BYTES_CHECKSUM
   3299             )
   3300             .unwrap()
   3301             .checksum(),
   3302             CALLBACK_SQL_BYTES_CHECKSUM
   3303         );
   3304     }
   3305 
   3306     #[test]
   3307     fn descriptor_checksum_name_version_and_content_bounds_fail_closed() {
   3308         assert_eq!(
   3309             MigrationDescriptor::sql(2, "alpha", SQL_TWO, MigrationChecksum::for_sql("other")),
   3310             Err(MigrationContractError::ChecksumMismatch)
   3311         );
   3312         assert_eq!(
   3313             MigrationDescriptor::callback(
   3314                 2,
   3315                 "alpha",
   3316                 CALLBACK_THREE,
   3317                 MigrationChecksum::for_callback(b"other")
   3318             ),
   3319             Err(MigrationContractError::ChecksumMismatch)
   3320         );
   3321         for invalid in ["", "_alpha", "alpha_", "Alpha", "alpha-beta", "alpha__beta"] {
   3322             assert_eq!(
   3323                 MigrationName::new(invalid),
   3324                 Err(MigrationContractError::InvalidName)
   3325             );
   3326         }
   3327         let max_name = Box::leak("a".repeat(MAX_MIGRATION_NAME_UTF8_BYTES).into_boxed_str());
   3328         assert_eq!(MigrationName::new(max_name).unwrap().as_str(), max_name);
   3329         let long_name = Box::leak(
   3330             "a".repeat(MAX_MIGRATION_NAME_UTF8_BYTES + 1)
   3331                 .into_boxed_str(),
   3332         );
   3333         assert_eq!(
   3334             MigrationName::new(long_name),
   3335             Err(MigrationContractError::InvalidName)
   3336         );
   3337         for version in [0, 1] {
   3338             assert_eq!(
   3339                 MigrationDescriptor::sql(
   3340                     version,
   3341                     "alpha",
   3342                     SQL_TWO,
   3343                     MigrationChecksum::for_sql(SQL_TWO)
   3344                 ),
   3345                 Err(MigrationContractError::InvalidTargetVersion)
   3346             );
   3347         }
   3348         assert_eq!(
   3349             MigrationDescriptor::sql(2, "alpha", "", MigrationChecksum::for_sql("")),
   3350             Err(MigrationContractError::EmptyContent)
   3351         );
   3352         let max_content = Box::leak(vec![b'x'; MAX_MIGRATION_CONTENT_BYTES].into_boxed_slice());
   3353         assert!(
   3354             MigrationDescriptor::callback(
   3355                 2,
   3356                 "alpha",
   3357                 max_content,
   3358                 MigrationChecksum::for_callback(max_content)
   3359             )
   3360             .is_ok()
   3361         );
   3362         let oversized = Box::leak(vec![b'x'; MAX_MIGRATION_CONTENT_BYTES + 1].into_boxed_slice());
   3363         assert_eq!(
   3364             MigrationDescriptor::callback(
   3365                 2,
   3366                 "alpha",
   3367                 oversized,
   3368                 MigrationChecksum::for_callback(oversized)
   3369             ),
   3370             Err(MigrationContractError::ContentTooLarge)
   3371         );
   3372     }
   3373 
   3374     #[test]
   3375     fn empty_v1_and_ordered_catalog_digest_are_exact() {
   3376         let empty = MigrationCatalog::new([]).expect("empty v1 catalog");
   3377         assert!(empty.descriptors().is_empty());
   3378         assert_eq!(empty.current_version(), 1);
   3379         assert_eq!(
   3380             hex(empty.digest()),
   3381             "ec89dc8f7b6c2a11b967e33808e4031e29b3970ffee4959bff9bad352877ee9b"
   3382         );
   3383 
   3384         let catalog = MigrationCatalog::new([
   3385             MigrationDescriptor::sql(2, "create_alpha", SQL_TWO, SQL_TWO_CHECKSUM).unwrap(),
   3386             MigrationDescriptor::callback(
   3387                 3,
   3388                 "rebuild_projection",
   3389                 CALLBACK_THREE,
   3390                 CALLBACK_THREE_CHECKSUM,
   3391             )
   3392             .unwrap(),
   3393         ])
   3394         .expect("ordered catalog");
   3395         assert_eq!(catalog.current_version(), 3);
   3396         assert_eq!(
   3397             catalog
   3398                 .descriptors()
   3399                 .iter()
   3400                 .map(MigrationDescriptor::target_version)
   3401                 .collect::<Vec<_>>(),
   3402             vec![2, 3]
   3403         );
   3404         assert_eq!(
   3405             hex(catalog.digest()),
   3406             "318e8b0143859e58ffe995b7d97d1cc2488097d307367979666cb47b98665838"
   3407         );
   3408     }
   3409 
   3410     #[test]
   3411     fn duplicate_gap_and_ordering_failures_are_distinct() {
   3412         assert_eq!(
   3413             MigrationCatalog::new([sql(2, "alpha", "alpha"), sql(2, "beta", "beta")]),
   3414             Err(MigrationContractError::DuplicateVersion)
   3415         );
   3416         assert_eq!(
   3417             MigrationCatalog::new([sql(2, "alpha", "alpha"), sql(3, "alpha", "beta")]),
   3418             Err(MigrationContractError::DuplicateName)
   3419         );
   3420         assert_eq!(
   3421             MigrationCatalog::new([sql(3, "alpha", "alpha")]),
   3422             Err(MigrationContractError::VersionGap)
   3423         );
   3424         assert_eq!(
   3425             MigrationCatalog::new([sql(2, "alpha", "alpha"), sql(4, "beta", "beta")]),
   3426             Err(MigrationContractError::VersionGap)
   3427         );
   3428         assert_eq!(
   3429             MigrationCatalog::new([sql(3, "alpha", "alpha"), sql(2, "beta", "beta")]),
   3430             Err(MigrationContractError::OutOfOrder)
   3431         );
   3432         assert_eq!(
   3433             MigrationCatalog::new([sql(u32::MAX, "alpha", "alpha")]),
   3434             Err(MigrationContractError::VersionGap)
   3435         );
   3436     }
   3437 
   3438     #[test]
   3439     fn catalog_count_is_bounded_during_ingestion() {
   3440         let maximum = (0..MAX_MIGRATION_COUNT).map(|index| {
   3441             let version = u32::try_from(index).unwrap() + 2;
   3442             let name = Box::leak(format!("migration_{version}").into_boxed_str());
   3443             sql(version, name, "SELECT 1;")
   3444         });
   3445         assert_eq!(
   3446             MigrationCatalog::new(maximum).unwrap().descriptors().len(),
   3447             MAX_MIGRATION_COUNT
   3448         );
   3449 
   3450         let excessive = (0..=MAX_MIGRATION_COUNT).map(|index| {
   3451             let version = u32::try_from(index).unwrap() + 2;
   3452             let name = Box::leak(format!("migration_{version}").into_boxed_str());
   3453             sql(version, name, "SELECT 1;")
   3454         });
   3455         assert_eq!(
   3456             MigrationCatalog::new(excessive),
   3457             Err(MigrationContractError::TooManyMigrations)
   3458         );
   3459 
   3460         let infinite = (2_u32..).map(|version| {
   3461             let name = Box::leak(format!("migration_{version}").into_boxed_str());
   3462             sql(version, name, "SELECT 1;")
   3463         });
   3464         assert_eq!(
   3465             MigrationCatalog::new(infinite),
   3466             Err(MigrationContractError::TooManyMigrations)
   3467         );
   3468     }
   3469 
   3470     #[test]
   3471     fn debug_and_errors_never_expose_migration_content() {
   3472         const SECRET_SQL: &str = "SELECT 'migration-secret';";
   3473         let descriptor = sql(2, "safe_name", SECRET_SQL);
   3474         let catalog = MigrationCatalog::new([descriptor.clone()]).unwrap();
   3475         for rendered in [format!("{descriptor:?}"), format!("{catalog:?}")] {
   3476             assert!(!rendered.contains(SECRET_SQL));
   3477             assert!(!rendered.contains("migration-secret"));
   3478         }
   3479         for error in [
   3480             MigrationContractError::InvalidName,
   3481             MigrationContractError::InvalidTargetVersion,
   3482             MigrationContractError::EmptyContent,
   3483             MigrationContractError::ContentTooLarge,
   3484             MigrationContractError::ChecksumMismatch,
   3485             MigrationContractError::TooManyMigrations,
   3486             MigrationContractError::DuplicateVersion,
   3487             MigrationContractError::DuplicateName,
   3488             MigrationContractError::OutOfOrder,
   3489             MigrationContractError::VersionGap,
   3490         ] {
   3491             assert!(!error.to_string().contains("secret"));
   3492             assert!(error.source().is_none());
   3493         }
   3494     }
   3495 
   3496     fn hex(checksum: MigrationChecksum) -> String {
   3497         checksum
   3498             .as_bytes()
   3499             .iter()
   3500             .map(|byte| format!("{byte:02x}"))
   3501             .collect()
   3502     }
   3503 }