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 }