initialize.rs (73313B)
1 //! Create-new initialization for one service-owned SQLite database. 2 3 use core::{fmt, future::Future}; 4 use std::{ 5 error::Error, 6 pin::Pin, 7 sync::{ 8 Arc, 9 atomic::{AtomicBool, Ordering}, 10 }, 11 }; 12 13 use futures::{future::BoxFuture, stream::BoxStream}; 14 use sqlx::{ 15 Either, Execute, Executor, SqlStr, Sqlite, SqliteConnection, 16 sqlite::{SqliteQueryResult, SqliteRow, SqliteStatement, SqliteTypeInfo}, 17 }; 18 19 use crate::{ 20 OpenMode, SchemaCatalog, ServiceDatabaseMetadata, ServiceSqliteError, ServiceSqliteErrorKind, 21 ServiceSqlitePaths, WriterAuthority, 22 }; 23 24 /// A sealed initialization executor that never exposes its SQLite connection. 25 /// 26 /// Service repositories may create their product schema with ordinary typed 27 /// SQLx queries through a mutable borrow. The service-SQLite runner exclusively 28 /// owns transaction begin, commit, rollback, shared metadata, and the migration 29 /// ledger. 30 /// 31 /// ``` 32 /// use radroots_service_sqlite::ServiceSqliteInitializer; 33 /// 34 /// async fn create_product_schema( 35 /// initializer: &mut ServiceSqliteInitializer<'_>, 36 /// ) -> Result<(), sqlx::Error> { 37 /// sqlx::query(concat!("CREATE ", "TABLE product_items (id INTEGER PRIMARY KEY)")) 38 /// .execute(initializer) 39 /// .await?; 40 /// Ok(()) 41 /// } 42 /// ``` 43 /// 44 /// Transaction control and the raw connection remain inaccessible: 45 /// 46 /// ```compile_fail 47 /// use radroots_service_sqlite::ServiceSqliteInitializer; 48 /// 49 /// async fn bypass(initializer: ServiceSqliteInitializer<'_>) { 50 /// initializer.commit().await.unwrap(); 51 /// } 52 /// ``` 53 pub struct ServiceSqliteInitializer<'connection> { 54 connection: &'connection mut SqliteConnection, 55 statement_control_rejected: Arc<AtomicBool>, 56 } 57 58 struct RestrictedInitializationExecute<Q> { 59 query: Q, 60 statement_control_rejected: Arc<AtomicBool>, 61 } 62 63 impl<'query, Q> Execute<'query, Sqlite> for RestrictedInitializationExecute<Q> 64 where 65 Q: Execute<'query, Sqlite>, 66 { 67 fn sql(self) -> SqlStr { 68 restricted_initialization_sql(self.query.sql(), &self.statement_control_rejected) 69 } 70 71 fn statement(&self) -> Option<&SqliteStatement> { 72 None 73 } 74 75 fn take_arguments( 76 &mut self, 77 ) -> Result<Option<<Sqlite as sqlx::Database>::Arguments>, sqlx::error::BoxDynError> { 78 self.query.take_arguments() 79 } 80 81 fn persistent(&self) -> bool { 82 self.query.persistent() 83 } 84 } 85 86 fn restricted_initialization_sql(sql: SqlStr, statement_control_rejected: &AtomicBool) -> SqlStr { 87 if crate::statement_policy::contains_forbidden_statement_control(sql.as_str()) { 88 statement_control_rejected.store(true, Ordering::Release); 89 SqlStr::from_static("RADROOTS_FORBIDDEN_INITIALIZATION_STATEMENT_CONTROL") 90 } else { 91 sql 92 } 93 } 94 95 impl fmt::Debug for ServiceSqliteInitializer<'_> { 96 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 97 formatter.write_str("ServiceSqliteInitializer([redacted])") 98 } 99 } 100 101 impl<'executor, 'connection> Executor<'executor> 102 for &'executor mut ServiceSqliteInitializer<'connection> 103 where 104 'connection: 'executor, 105 { 106 type Database = Sqlite; 107 108 fn fetch_many<'e, 'q: 'e, Q>( 109 self, 110 query: Q, 111 ) -> BoxStream<'e, Result<Either<SqliteQueryResult, SqliteRow>, sqlx::Error>> 112 where 113 'executor: 'e, 114 Q: 'q + Execute<'q, Self::Database>, 115 { 116 (&mut *self.connection).fetch_many(RestrictedInitializationExecute { 117 query, 118 statement_control_rejected: Arc::clone(&self.statement_control_rejected), 119 }) 120 } 121 122 fn fetch_optional<'e, 'q: 'e, Q>( 123 self, 124 query: Q, 125 ) -> BoxFuture<'e, Result<Option<SqliteRow>, sqlx::Error>> 126 where 127 'executor: 'e, 128 Q: 'q + Execute<'q, Self::Database>, 129 { 130 (&mut *self.connection).fetch_optional(RestrictedInitializationExecute { 131 query, 132 statement_control_rejected: Arc::clone(&self.statement_control_rejected), 133 }) 134 } 135 136 fn prepare_with<'e>( 137 self, 138 sql: SqlStr, 139 parameters: &'e [SqliteTypeInfo], 140 ) -> BoxFuture<'e, Result<SqliteStatement, sqlx::Error>> 141 where 142 'executor: 'e, 143 { 144 (&mut *self.connection).prepare_with( 145 restricted_initialization_sql(sql, &self.statement_control_rejected), 146 parameters, 147 ) 148 } 149 } 150 151 /// Boxed callback future tied to the borrowed initialization executor. 152 pub type ServiceSqliteInitializerFuture<'a, E> = 153 Pin<Box<dyn Future<Output = Result<(), E>> + Send + 'a>>; 154 155 #[cfg(any(target_os = "linux", target_os = "macos"))] 156 #[derive(Debug)] 157 pub(crate) enum InitializeDatabaseOutcome { 158 Initialized(WriterAuthority), 159 Existing(WriterAuthority), 160 } 161 162 /// Creates and initializes a missing service database while holding sole writer authority. 163 /// 164 /// The callback receives only a mutable borrow of the sealed initialization 165 /// executor. Product schema, shared metadata, the empty v1 migration ledger, 166 /// and catalog verification occur in one runner-owned transaction. Cancellation 167 /// or failure quarantines the one-shot connection and removes only the exact 168 /// inode reserved by this call. 169 pub async fn initialize_database<F, E>( 170 paths: &ServiceSqlitePaths, 171 mode: OpenMode, 172 metadata: &ServiceDatabaseMetadata, 173 schema_catalog: &SchemaCatalog, 174 initialize_schema: F, 175 ) -> Result<WriterAuthority, ServiceSqliteError> 176 where 177 F: for<'a> FnOnce( 178 &'a mut ServiceSqliteInitializer<'_>, 179 ) -> ServiceSqliteInitializerFuture<'a, E>, 180 E: Error + Send + Sync + 'static, 181 { 182 if mode != OpenMode::Initialize { 183 return Err(initialization_error(InitializationCause::new( 184 InitializationFailureKind::UnsupportedMode, 185 ))); 186 } 187 if !metadata.matches_paths(paths) { 188 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata)); 189 } 190 191 let authority = WriterAuthority::acquire(paths, mode)?.ok_or_else(|| { 192 initialization_error(InitializationCause::new( 193 InitializationFailureKind::UnsupportedMode, 194 )) 195 })?; 196 197 #[cfg(any(target_os = "linux", target_os = "macos"))] 198 { 199 authority.validate_for(paths)?; 200 let recovery = crate::restore::refuse_unresolved_recovery(authority.directory()); 201 authority.validate_for(paths)?; 202 recovery?; 203 } 204 205 #[cfg(any(target_os = "linux", target_os = "macos"))] 206 { 207 let failpoints = crate::failpoint::DurabilityFailpoints::default(); 208 match initialize_with_ops( 209 paths, 210 authority, 211 metadata, 212 schema_catalog, 213 initialize_schema, 214 &SystemInitializationOperations, 215 &failpoints, 216 ) 217 .await? 218 { 219 InitializeDatabaseOutcome::Initialized(authority) => Ok(authority), 220 InitializeDatabaseOutcome::Existing(_authority) => Err(initialization_error( 221 InitializationCause::new(InitializationFailureKind::StateAlreadyExists), 222 )), 223 } 224 } 225 226 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 227 { 228 drop((authority, metadata, schema_catalog, initialize_schema)); 229 Err(initialization_error(InitializationCause::new( 230 InitializationFailureKind::CreateUnavailable, 231 ))) 232 } 233 } 234 235 #[cfg(any(target_os = "linux", target_os = "macos"))] 236 pub(crate) async fn initialize_or_existing_database<F, E>( 237 paths: &ServiceSqlitePaths, 238 metadata: &ServiceDatabaseMetadata, 239 schema_catalog: &SchemaCatalog, 240 initialize_schema: F, 241 ) -> Result<InitializeDatabaseOutcome, ServiceSqliteError> 242 where 243 F: for<'a> FnOnce( 244 &'a mut ServiceSqliteInitializer<'_>, 245 ) -> ServiceSqliteInitializerFuture<'a, E>, 246 E: Error + Send + Sync + 'static, 247 { 248 if !metadata.matches_paths(paths) { 249 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata)); 250 } 251 let authority = WriterAuthority::acquire(paths, OpenMode::Initialize)?.ok_or_else(|| { 252 initialization_error(InitializationCause::new( 253 InitializationFailureKind::UnsupportedMode, 254 )) 255 })?; 256 let failpoints = crate::failpoint::DurabilityFailpoints::default(); 257 initialize_with_ops( 258 paths, 259 authority, 260 metadata, 261 schema_catalog, 262 initialize_schema, 263 &SystemInitializationOperations, 264 &failpoints, 265 ) 266 .await 267 } 268 269 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 270 enum InitializationFailureKind { 271 UnsupportedMode, 272 #[cfg(any(target_os = "linux", target_os = "macos"))] 273 StateAlreadyExists, 274 CreateUnavailable, 275 #[cfg(any(target_os = "linux", target_os = "macos"))] 276 InvalidDatabase, 277 #[cfg(any(target_os = "linux", target_os = "macos"))] 278 SchemaInitializationFailed, 279 #[cfg(any(target_os = "linux", target_os = "macos"))] 280 DatabaseSyncFailed, 281 #[cfg(any(target_os = "linux", target_os = "macos"))] 282 DatabaseReplaced, 283 #[cfg(any(target_os = "linux", target_os = "macos"))] 284 DirectorySyncFailed, 285 #[cfg(any(target_os = "linux", target_os = "macos"))] 286 CleanupFailed, 287 #[cfg(any(target_os = "linux", target_os = "macos"))] 288 InjectedFailure, 289 } 290 291 struct InitializationCause { 292 kind: InitializationFailureKind, 293 source: Option<Box<dyn Error + Send + Sync + 'static>>, 294 } 295 296 impl InitializationCause { 297 const fn new(kind: InitializationFailureKind) -> Self { 298 Self { kind, source: None } 299 } 300 301 #[cfg(any(target_os = "linux", target_os = "macos"))] 302 fn with_source( 303 kind: InitializationFailureKind, 304 source: impl Error + Send + Sync + 'static, 305 ) -> Self { 306 Self { 307 kind, 308 source: Some(Box::new(source)), 309 } 310 } 311 } 312 313 impl fmt::Debug for InitializationCause { 314 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 315 formatter 316 .debug_struct("InitializationCause") 317 .field("kind", &self.kind) 318 .field("source", &self.source.as_ref().map(|_| "[redacted]")) 319 .finish() 320 } 321 } 322 323 impl fmt::Display for InitializationCause { 324 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 325 formatter.write_str(match self.kind { 326 InitializationFailureKind::UnsupportedMode => { 327 "SQLite initialization requires initialize mode" 328 } 329 #[cfg(any(target_os = "linux", target_os = "macos"))] 330 InitializationFailureKind::StateAlreadyExists => "SQLite state already exists", 331 InitializationFailureKind::CreateUnavailable => "SQLite state could not be reserved", 332 #[cfg(any(target_os = "linux", target_os = "macos"))] 333 InitializationFailureKind::InvalidDatabase => "SQLite state file has invalid metadata", 334 #[cfg(any(target_os = "linux", target_os = "macos"))] 335 InitializationFailureKind::SchemaInitializationFailed => { 336 "SQLite schema initialization failed" 337 } 338 #[cfg(any(target_os = "linux", target_os = "macos"))] 339 InitializationFailureKind::DatabaseSyncFailed => { 340 "SQLite state could not be synchronized" 341 } 342 #[cfg(any(target_os = "linux", target_os = "macos"))] 343 InitializationFailureKind::DatabaseReplaced => { 344 "SQLite state identity changed during initialization" 345 } 346 #[cfg(any(target_os = "linux", target_os = "macos"))] 347 InitializationFailureKind::DirectorySyncFailed => { 348 "SQLite state directory could not be synchronized" 349 } 350 #[cfg(any(target_os = "linux", target_os = "macos"))] 351 InitializationFailureKind::CleanupFailed => "SQLite initialization cleanup failed", 352 #[cfg(any(target_os = "linux", target_os = "macos"))] 353 InitializationFailureKind::InjectedFailure => { 354 "SQLite initialization durability boundary failed" 355 } 356 }) 357 } 358 } 359 360 impl Error for InitializationCause { 361 fn source(&self) -> Option<&(dyn Error + 'static)> { 362 self.source 363 .as_deref() 364 .map(|source| source as &(dyn Error + 'static)) 365 } 366 } 367 368 fn initialization_error(cause: InitializationCause) -> ServiceSqliteError { 369 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Create, cause) 370 } 371 372 #[cfg(any(test, target_os = "linux", target_os = "macos"))] 373 fn require_initialization_condition( 374 condition: bool, 375 kind: InitializationFailureKind, 376 ) -> Result<(), InitializationCause> { 377 condition 378 .then_some(()) 379 .ok_or_else(|| InitializationCause::new(kind)) 380 } 381 382 #[cfg(test)] 383 mod failure_tests { 384 385 use super::*; 386 387 #[test] 388 fn initialization_failure_inventory_is_complete_and_source_aware() { 389 let cases = [ 390 ( 391 InitializationFailureKind::UnsupportedMode, 392 "SQLite initialization requires initialize mode", 393 ), 394 ( 395 InitializationFailureKind::CreateUnavailable, 396 "SQLite state could not be reserved", 397 ), 398 ] 399 .into_iter(); 400 #[cfg(any(target_os = "linux", target_os = "macos"))] 401 let cases = cases.chain([ 402 ( 403 InitializationFailureKind::StateAlreadyExists, 404 "SQLite state already exists", 405 ), 406 ( 407 InitializationFailureKind::InvalidDatabase, 408 "SQLite state file has invalid metadata", 409 ), 410 ( 411 InitializationFailureKind::SchemaInitializationFailed, 412 "SQLite schema initialization failed", 413 ), 414 ( 415 InitializationFailureKind::DatabaseSyncFailed, 416 "SQLite state could not be synchronized", 417 ), 418 ( 419 InitializationFailureKind::DatabaseReplaced, 420 "SQLite state identity changed during initialization", 421 ), 422 ( 423 InitializationFailureKind::DirectorySyncFailed, 424 "SQLite state directory could not be synchronized", 425 ), 426 ( 427 InitializationFailureKind::CleanupFailed, 428 "SQLite initialization cleanup failed", 429 ), 430 ( 431 InitializationFailureKind::InjectedFailure, 432 "SQLite initialization durability boundary failed", 433 ), 434 ]); 435 436 for (kind, message) in cases { 437 let plain = InitializationCause::new(kind); 438 assert_eq!(plain.to_string(), message); 439 assert!(plain.source().is_none()); 440 assert!(format!("{plain:?}").contains("source: None")); 441 assert!(require_initialization_condition(true, kind).is_ok()); 442 assert_eq!( 443 require_initialization_condition(false, kind) 444 .expect_err("false condition") 445 .kind, 446 kind 447 ); 448 449 #[cfg(any(target_os = "linux", target_os = "macos"))] 450 { 451 let sourced = 452 InitializationCause::with_source(kind, std::io::Error::other("private-cause")); 453 assert_eq!(sourced.to_string(), message); 454 assert!(sourced.source().is_some()); 455 let debug = format!("{sourced:?}"); 456 assert!(debug.contains("[redacted]")); 457 assert!(!debug.contains("private-cause")); 458 } 459 } 460 } 461 } 462 463 #[cfg(any(target_os = "linux", target_os = "macos"))] 464 mod supported { 465 use std::{fs::File, os::fd::AsRawFd}; 466 467 use rustix::{ 468 fs::{AtFlags, FileType, Mode, OFlags, fchmod, fstat, lstat, openat, statat, unlinkat}, 469 process::geteuid, 470 }; 471 472 use super::*; 473 474 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 475 struct FileIdentity { 476 device: u64, 477 inode: u64, 478 } 479 480 pub(super) trait InitializationOperations { 481 fn sync_database(&self, database: &File) -> Result<(), InitializationCause>; 482 fn sync_directory(&self, directory: &File) -> Result<(), InitializationCause>; 483 fn unlink_database(&self, directory: &File) -> Result<(), InitializationCause>; 484 } 485 486 pub(super) struct SystemInitializationOperations; 487 488 impl InitializationOperations for SystemInitializationOperations { 489 #[cfg_attr(coverage_nightly, coverage(off))] 490 fn sync_database(&self, database: &File) -> Result<(), InitializationCause> { 491 database.sync_all().map_err(|_| { 492 InitializationCause::new(InitializationFailureKind::DatabaseSyncFailed) 493 }) 494 } 495 496 #[cfg_attr(coverage_nightly, coverage(off))] 497 fn sync_directory(&self, directory: &File) -> Result<(), InitializationCause> { 498 directory.sync_all().map_err(|_| { 499 InitializationCause::new(InitializationFailureKind::DirectorySyncFailed) 500 }) 501 } 502 503 #[cfg_attr(coverage_nightly, coverage(off))] 504 fn unlink_database(&self, directory: &File) -> Result<(), InitializationCause> { 505 unlinkat( 506 directory, 507 radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME, 508 AtFlags::empty(), 509 ) 510 .map_err(|_| InitializationCause::new(InitializationFailureKind::CleanupFailed)) 511 } 512 } 513 514 struct PendingDatabase<'a, O: InitializationOperations> { 515 directory: &'a File, 516 database: File, 517 identity: FileIdentity, 518 operations: &'a O, 519 committed: bool, 520 } 521 522 impl<'a, O: InitializationOperations> PendingDatabase<'a, O> { 523 fn create(directory: &'a File, operations: &'a O) -> Result<Self, InitializationCause> { 524 let descriptor = openat( 525 directory, 526 radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME, 527 OFlags::RDWR | OFlags::CREATE | OFlags::EXCL | OFlags::NOFOLLOW | OFlags::CLOEXEC, 528 Mode::RUSR | Mode::WUSR, 529 ) 530 .map_err(|error| { 531 if error == rustix::io::Errno::EXIST { 532 InitializationCause::new(InitializationFailureKind::StateAlreadyExists) 533 } else { 534 InitializationCause::new(InitializationFailureKind::CreateUnavailable) 535 } 536 })?; 537 let database = File::from(descriptor); 538 let identity = descriptor_identity(&database)?; 539 let pending = Self { 540 directory, 541 database, 542 identity, 543 operations, 544 committed: false, 545 }; 546 fchmod(&pending.database, Mode::RUSR | Mode::WUSR).map_err(|_| { 547 InitializationCause::new(InitializationFailureKind::InvalidDatabase) 548 })?; 549 pending.validate()?; 550 Ok(pending) 551 } 552 553 fn validate(&self) -> Result<(), InitializationCause> { 554 let descriptor_identity = validate_descriptor(&self.database)?; 555 require_initialization_condition( 556 descriptor_identity == self.identity, 557 InitializationFailureKind::InvalidDatabase, 558 )?; 559 self.validate_entry() 560 } 561 562 fn sqlite_descriptor_path(&self) -> String { 563 let descriptor = self.database.as_raw_fd(); 564 #[cfg(target_os = "linux")] 565 let path = format!("/proc/self/fd/{descriptor}"); 566 #[cfg(target_os = "macos")] 567 let path = format!("/dev/fd/{descriptor}"); 568 path 569 } 570 571 fn validate_entry(&self) -> Result<(), InitializationCause> { 572 let status = statat( 573 self.directory, 574 radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME, 575 AtFlags::SYMLINK_NOFOLLOW, 576 ) 577 .map_err(|_| InitializationCause::new(InitializationFailureKind::DatabaseReplaced))?; 578 let device = crate::native_metadata::device(status.st_dev).map_err(|_| { 579 InitializationCause::new(InitializationFailureKind::InvalidDatabase) 580 })?; 581 let current = validate_status( 582 FileType::from_raw_mode(status.st_mode).is_file(), 583 crate::native_metadata::link_count(status.st_nlink), 584 status.st_uid, 585 crate::native_metadata::mode(status.st_mode), 586 device, 587 status.st_ino, 588 )?; 589 require_initialization_condition( 590 current == self.identity, 591 InitializationFailureKind::DatabaseReplaced, 592 )?; 593 Ok(()) 594 } 595 596 fn current_entry_identity(&self) -> Result<FileIdentity, InitializationCause> { 597 let status = statat( 598 self.directory, 599 radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME, 600 AtFlags::SYMLINK_NOFOLLOW, 601 ) 602 .map_err(|_| InitializationCause::new(InitializationFailureKind::DatabaseReplaced))?; 603 let device = crate::native_metadata::device(status.st_dev).map_err(|_| { 604 InitializationCause::new(InitializationFailureKind::DatabaseReplaced) 605 })?; 606 Ok(FileIdentity { 607 device, 608 inode: status.st_ino, 609 }) 610 } 611 612 fn validate_canonical_path( 613 &self, 614 path: &std::path::Path, 615 ) -> Result<(), InitializationCause> { 616 let status = lstat(path).map_err(|_| { 617 InitializationCause::new(InitializationFailureKind::DatabaseReplaced) 618 })?; 619 let device = crate::native_metadata::device(status.st_dev).map_err(|_| { 620 InitializationCause::new(InitializationFailureKind::InvalidDatabase) 621 })?; 622 let current = validate_status( 623 FileType::from_raw_mode(status.st_mode).is_file(), 624 crate::native_metadata::link_count(status.st_nlink), 625 status.st_uid, 626 crate::native_metadata::mode(status.st_mode), 627 device, 628 status.st_ino, 629 )?; 630 require_initialization_condition( 631 current == self.identity, 632 InitializationFailureKind::DatabaseReplaced, 633 )?; 634 Ok(()) 635 } 636 637 fn commit( 638 &mut self, 639 canonical_path: &std::path::Path, 640 failpoints: &crate::failpoint::DurabilityFailpoints, 641 ) -> Result<(), InitializationCause> { 642 hit( 643 failpoints, 644 crate::failpoint::DurabilityFailpoint::InitializeBeforeFileSync, 645 )?; 646 self.operations.sync_database(&self.database)?; 647 hit( 648 failpoints, 649 crate::failpoint::DurabilityFailpoint::InitializeAfterFileSync, 650 )?; 651 self.validate()?; 652 self.validate_canonical_path(canonical_path)?; 653 hit( 654 failpoints, 655 crate::failpoint::DurabilityFailpoint::InitializeBeforeCommitDirectorySync, 656 )?; 657 self.operations.sync_directory(self.directory)?; 658 hit( 659 failpoints, 660 crate::failpoint::DurabilityFailpoint::InitializeAfterCommitDirectorySync, 661 )?; 662 self.committed = true; 663 Ok(()) 664 } 665 666 fn rollback(&mut self) -> Result<(), InitializationCause> { 667 if self.committed { 668 return Ok(()); 669 } 670 require_initialization_condition( 671 self.current_entry_identity()? == self.identity, 672 InitializationFailureKind::DatabaseReplaced, 673 )?; 674 self.operations.unlink_database(self.directory)?; 675 self.operations.sync_directory(self.directory)?; 676 self.committed = true; 677 Ok(()) 678 } 679 } 680 681 impl<O: InitializationOperations> Drop for PendingDatabase<'_, O> { 682 fn drop(&mut self) { 683 let _ = self.rollback(); 684 } 685 } 686 687 fn validate_descriptor( 688 descriptor: &impl std::os::fd::AsFd, 689 ) -> Result<FileIdentity, InitializationCause> { 690 let status = fstat(descriptor) 691 .map_err(|_| InitializationCause::new(InitializationFailureKind::InvalidDatabase))?; 692 let device = crate::native_metadata::device(status.st_dev) 693 .map_err(|_| InitializationCause::new(InitializationFailureKind::InvalidDatabase))?; 694 validate_status( 695 FileType::from_raw_mode(status.st_mode).is_file(), 696 crate::native_metadata::link_count(status.st_nlink), 697 status.st_uid, 698 crate::native_metadata::mode(status.st_mode), 699 device, 700 status.st_ino, 701 ) 702 } 703 704 fn descriptor_identity( 705 descriptor: &impl std::os::fd::AsFd, 706 ) -> Result<FileIdentity, InitializationCause> { 707 let status = fstat(descriptor) 708 .map_err(|_| InitializationCause::new(InitializationFailureKind::InvalidDatabase))?; 709 let device = crate::native_metadata::device(status.st_dev) 710 .map_err(|_| InitializationCause::new(InitializationFailureKind::InvalidDatabase))?; 711 Ok(FileIdentity { 712 device, 713 inode: status.st_ino, 714 }) 715 } 716 717 fn validate_status( 718 is_regular_file: bool, 719 link_count: u64, 720 actual_uid: u32, 721 mode: u32, 722 device: u64, 723 inode: u64, 724 ) -> Result<FileIdentity, InitializationCause> { 725 require_initialization_condition( 726 crate::native_metadata::exact_regular_file( 727 is_regular_file, 728 link_count, 729 actual_uid, 730 geteuid().as_raw(), 731 mode, 732 ), 733 InitializationFailureKind::InvalidDatabase, 734 )?; 735 Ok(FileIdentity { device, inode }) 736 } 737 738 fn rollback_failure( 739 primary: InitializationCause, 740 cleanup: Result<(), InitializationCause>, 741 ) -> ServiceSqliteError { 742 match cleanup { 743 Ok(()) => initialization_error(primary), 744 Err(_cleanup) if primary.kind == InitializationFailureKind::DatabaseReplaced => { 745 initialization_error(primary) 746 } 747 Err(cleanup) => { 748 initialization_error(InitializationCause::with_source(cleanup.kind, primary)) 749 } 750 } 751 } 752 753 async fn fail_with_rollback<O: InitializationOperations>( 754 mut pending: PendingDatabase<'_, O>, 755 primary: InitializationCause, 756 ) -> Result<InitializeDatabaseOutcome, ServiceSqliteError> { 757 Err(rollback_failure(primary, pending.rollback())) 758 } 759 760 async fn fail_metadata_with_rollback<O: InitializationOperations>( 761 mut pending: PendingDatabase<'_, O>, 762 primary: ServiceSqliteError, 763 ) -> Result<InitializeDatabaseOutcome, ServiceSqliteError> { 764 match pending.rollback() { 765 Ok(()) => Err(primary), 766 Err(cleanup) => Err(initialization_error(InitializationCause::with_source( 767 cleanup.kind, 768 primary, 769 ))), 770 } 771 } 772 773 pub(super) async fn initialize_with_ops<F, E, O>( 774 paths: &ServiceSqlitePaths, 775 authority: WriterAuthority, 776 metadata: &ServiceDatabaseMetadata, 777 schema_catalog: &SchemaCatalog, 778 initialize_schema: F, 779 operations: &O, 780 failpoints: &crate::failpoint::DurabilityFailpoints, 781 ) -> Result<InitializeDatabaseOutcome, ServiceSqliteError> 782 where 783 F: for<'a> FnOnce( 784 &'a mut ServiceSqliteInitializer<'_>, 785 ) -> ServiceSqliteInitializerFuture<'a, E>, 786 E: Error + Send + Sync + 'static, 787 O: InitializationOperations, 788 { 789 hit( 790 failpoints, 791 crate::failpoint::DurabilityFailpoint::InitializeBeforeCreate, 792 ) 793 .map_err(initialization_error)?; 794 let initialization_directory = authority.directory().try_clone().map_err(|source| { 795 initialization_error(InitializationCause::with_source( 796 InitializationFailureKind::CreateUnavailable, 797 source, 798 )) 799 })?; 800 let mut pending = match PendingDatabase::create(&initialization_directory, operations) { 801 Ok(pending) => pending, 802 Err(error) if error.kind == InitializationFailureKind::StateAlreadyExists => { 803 return Ok(InitializeDatabaseOutcome::Existing(authority)); 804 } 805 Err(error) => return Err(initialization_error(error)), 806 }; 807 if let Err(error) = hit( 808 failpoints, 809 crate::failpoint::DurabilityFailpoint::InitializeAfterCreate, 810 ) { 811 return fail_with_rollback(pending, error).await; 812 } 813 if let Err(error) = hit( 814 failpoints, 815 crate::failpoint::DurabilityFailpoint::InitializeBeforeReservationDirectorySync, 816 ) { 817 return fail_with_rollback(pending, error).await; 818 } 819 if let Err(error) = operations.sync_directory(authority.directory()) { 820 return fail_with_rollback(pending, error).await; 821 } 822 if let Err(error) = hit( 823 failpoints, 824 crate::failpoint::DurabilityFailpoint::InitializeAfterReservationDirectorySync, 825 ) { 826 return fail_with_rollback(pending, error).await; 827 } 828 authority.validate_for(paths)?; 829 let recovery = crate::restore::refuse_unresolved_recovery(authority.directory()); 830 authority.validate_for(paths)?; 831 if let Err(error) = recovery { 832 return fail_metadata_with_rollback(pending, error).await; 833 } 834 if let Err(error) = pending.validate_canonical_path(paths.state_database()) { 835 return fail_with_rollback(pending, error).await; 836 } 837 let metadata_result = initialize_transaction( 838 paths, 839 &authority, 840 &pending, 841 metadata, 842 schema_catalog, 843 initialize_schema, 844 ) 845 .await; 846 authority.validate_for(paths)?; 847 if let Err(error) = pending.validate() { 848 return fail_with_rollback(pending, error).await; 849 } 850 if let Err(error) = pending.validate_canonical_path(paths.state_database()) { 851 return fail_with_rollback(pending, error).await; 852 } 853 if let Err(error) = metadata_result { 854 return fail_metadata_with_rollback(pending, error).await; 855 } 856 if let Err(error) = pending.validate() { 857 return fail_with_rollback(pending, error).await; 858 } 859 if let Err(error) = pending.validate_canonical_path(paths.state_database()) { 860 return fail_with_rollback(pending, error).await; 861 } 862 if let Err(error) = pending.commit(paths.state_database(), failpoints) { 863 return fail_with_rollback(pending, error).await; 864 } 865 drop(pending); 866 Ok(InitializeDatabaseOutcome::Initialized(authority)) 867 } 868 869 async fn initialize_transaction<F, E, O>( 870 paths: &ServiceSqlitePaths, 871 authority: &WriterAuthority, 872 pending: &PendingDatabase<'_, O>, 873 metadata: &ServiceDatabaseMetadata, 874 schema_catalog: &SchemaCatalog, 875 initialize_schema: F, 876 ) -> Result<(), ServiceSqliteError> 877 where 878 F: for<'a> FnOnce( 879 &'a mut ServiceSqliteInitializer<'_>, 880 ) -> ServiceSqliteInitializerFuture<'a, E>, 881 E: Error + Send + Sync + 'static, 882 O: InitializationOperations, 883 { 884 use sqlx::{ 885 ConnectOptions, Connection, 886 sqlite::{SqliteConnectOptions, SqliteJournalMode}, 887 }; 888 889 authority.validate_for(paths)?; 890 pending 891 .validate() 892 .and_then(|()| pending.validate_canonical_path(paths.state_database())) 893 .map_err(initialization_error)?; 894 let options = SqliteConnectOptions::new() 895 .filename(pending.sqlite_descriptor_path()) 896 .create_if_missing(false) 897 .foreign_keys(true) 898 .journal_mode(SqliteJournalMode::Memory) 899 .pragma("trusted_schema", "OFF") 900 .disable_statement_logging(); 901 let connected = SqliteConnection::connect_with(&options).await; 902 authority.validate_for(paths)?; 903 pending 904 .validate() 905 .and_then(|()| pending.validate_canonical_path(paths.state_database())) 906 .map_err(initialization_error)?; 907 let mut connection = connected.map_err(|source| { 908 initialization_error(InitializationCause::with_source( 909 InitializationFailureKind::SchemaInitializationFailed, 910 source, 911 )) 912 })?; 913 let installed = 914 crate::transaction_control::TransactionControlGate::install(&mut connection).await; 915 authority.validate_for(paths)?; 916 pending 917 .validate() 918 .and_then(|()| pending.validate_canonical_path(paths.state_database())) 919 .map_err(initialization_error)?; 920 let gate = installed.map_err(|source| { 921 initialization_error(InitializationCause::with_source( 922 InitializationFailureKind::SchemaInitializationFailed, 923 source, 924 )) 925 })?; 926 let begun = connection.begin_with("BEGIN IMMEDIATE").await; 927 authority.validate_for(paths)?; 928 pending 929 .validate() 930 .and_then(|()| pending.validate_canonical_path(paths.state_database())) 931 .map_err(initialization_error)?; 932 let mut transaction = begun.map_err(|source| { 933 initialization_error(InitializationCause::with_source( 934 InitializationFailureKind::SchemaInitializationFailed, 935 source, 936 )) 937 })?; 938 939 let statement_control_rejected = Arc::new(AtomicBool::new(false)); 940 let callback = { 941 let mut initializer = ServiceSqliteInitializer { 942 connection: &mut transaction, 943 statement_control_rejected: Arc::clone(&statement_control_rejected), 944 }; 945 initialize_schema(&mut initializer).await 946 }; 947 authority.validate_for(paths)?; 948 pending 949 .validate() 950 .and_then(|()| pending.validate_canonical_path(paths.state_database())) 951 .map_err(initialization_error)?; 952 let governed_after_callback = 953 crate::migration::assert_governed_transaction(&mut transaction) 954 .await 955 .is_ok(); 956 authority.validate_for(paths)?; 957 958 let operation = match callback { 959 Ok(()) 960 if !statement_control_rejected.load(Ordering::Acquire) 961 && !gate.control_violation_observed() 962 && governed_after_callback => 963 { 964 crate::metadata::write_database_metadata_in_transaction( 965 &mut transaction, 966 metadata, 967 schema_catalog, 968 ) 969 .await 970 } 971 Ok(()) => Err(initialization_error(InitializationCause::new( 972 InitializationFailureKind::SchemaInitializationFailed, 973 ))), 974 Err(source) => Err(initialization_error(InitializationCause::with_source( 975 InitializationFailureKind::SchemaInitializationFailed, 976 source, 977 ))), 978 }; 979 authority.validate_for(paths)?; 980 pending 981 .validate() 982 .and_then(|()| pending.validate_canonical_path(paths.state_database())) 983 .map_err(initialization_error)?; 984 let governed_before_commit = 985 crate::migration::assert_governed_transaction(&mut transaction) 986 .await 987 .is_ok(); 988 authority.validate_for(paths)?; 989 990 if let Err(primary) = operation { 991 let permit = gate.permit_runner_rollback(); 992 let rollback = transaction.rollback().await; 993 drop(permit); 994 authority.validate_for(paths)?; 995 pending 996 .validate() 997 .and_then(|()| pending.validate_canonical_path(paths.state_database())) 998 .map_err(initialization_error)?; 999 let remove = gate.remove(&mut connection).await; 1000 authority.validate_for(paths)?; 1001 let close = connection.close().await; 1002 authority.validate_for(paths)?; 1003 if let Err(source) = rollback.or(remove).or(close) { 1004 return Err(initialization_error(InitializationCause::with_source( 1005 InitializationFailureKind::SchemaInitializationFailed, 1006 source, 1007 ))); 1008 } 1009 return Err(primary); 1010 } 1011 1012 if gate.control_violation_observed() || !governed_before_commit { 1013 let permit = gate.permit_runner_rollback(); 1014 let rollback = transaction.rollback().await; 1015 drop(permit); 1016 let remove = gate.remove(&mut connection).await; 1017 let close = connection.close().await; 1018 authority.validate_for(paths)?; 1019 rollback.or(remove).or(close).map_err(|source| { 1020 initialization_error(InitializationCause::with_source( 1021 InitializationFailureKind::SchemaInitializationFailed, 1022 source, 1023 )) 1024 })?; 1025 return Err(initialization_error(InitializationCause::new( 1026 InitializationFailureKind::SchemaInitializationFailed, 1027 ))); 1028 } 1029 1030 let permit = gate.permit_outer_commit(); 1031 let committed = transaction.commit().await; 1032 drop(permit); 1033 authority.validate_for(paths)?; 1034 pending 1035 .validate() 1036 .and_then(|()| pending.validate_canonical_path(paths.state_database())) 1037 .map_err(initialization_error)?; 1038 let removed = gate.remove(&mut connection).await; 1039 authority.validate_for(paths)?; 1040 let closed = connection.close().await; 1041 authority.validate_for(paths)?; 1042 committed.or(removed).or(closed).map_err(|source| { 1043 initialization_error(InitializationCause::with_source( 1044 InitializationFailureKind::SchemaInitializationFailed, 1045 source, 1046 )) 1047 }) 1048 } 1049 1050 fn hit( 1051 failpoints: &crate::failpoint::DurabilityFailpoints, 1052 point: crate::failpoint::DurabilityFailpoint, 1053 ) -> Result<(), InitializationCause> { 1054 failpoints.hit(point).map_err(|source| { 1055 InitializationCause::with_source(InitializationFailureKind::InjectedFailure, source) 1056 }) 1057 } 1058 1059 #[cfg(test)] 1060 mod tests { 1061 use std::{ 1062 cell::{Cell, RefCell}, 1063 convert::Infallible, 1064 fs, 1065 future::{pending, ready}, 1066 io, 1067 num::NonZeroU32, 1068 os::unix::fs::{MetadataExt, PermissionsExt, symlink}, 1069 path::Path, 1070 sync::Arc, 1071 }; 1072 1073 use radroots_runtime_paths::{ 1074 InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, 1075 RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, 1076 ServiceId, 1077 }; 1078 use radroots_storage::event::SourceGeneration; 1079 use tokio::sync::Notify; 1080 1081 use super::*; 1082 1083 #[derive(Debug)] 1084 struct CallbackFailure; 1085 1086 impl fmt::Display for CallbackFailure { 1087 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1088 formatter.write_str("secret callback path=/private/state.sqlite") 1089 } 1090 } 1091 1092 impl Error for CallbackFailure {} 1093 1094 #[derive(Default)] 1095 struct RecordingOperations { 1096 events: RefCell<Vec<&'static str>>, 1097 fail_database_sync: Cell<bool>, 1098 fail_directory_sync_on_call: Cell<Option<usize>>, 1099 directory_sync_calls: Cell<usize>, 1100 fail_unlink: Cell<bool>, 1101 } 1102 1103 impl InitializationOperations for RecordingOperations { 1104 fn sync_database(&self, _database: &File) -> Result<(), InitializationCause> { 1105 self.events.borrow_mut().push("sync_database"); 1106 if self.fail_database_sync.replace(false) { 1107 Err(InitializationCause::new( 1108 InitializationFailureKind::DatabaseSyncFailed, 1109 )) 1110 } else { 1111 Ok(()) 1112 } 1113 } 1114 1115 fn sync_directory(&self, _directory: &File) -> Result<(), InitializationCause> { 1116 self.events.borrow_mut().push("sync_directory"); 1117 let call = self.directory_sync_calls.get() + 1; 1118 self.directory_sync_calls.set(call); 1119 if self.fail_directory_sync_on_call.get() == Some(call) { 1120 Err(InitializationCause::new( 1121 InitializationFailureKind::DirectorySyncFailed, 1122 )) 1123 } else { 1124 Ok(()) 1125 } 1126 } 1127 1128 fn unlink_database(&self, directory: &File) -> Result<(), InitializationCause> { 1129 self.events.borrow_mut().push("unlink_database"); 1130 if self.fail_unlink.get() { 1131 return Err(InitializationCause::new( 1132 InitializationFailureKind::CleanupFailed, 1133 )); 1134 } 1135 SystemInitializationOperations.unlink_database(directory) 1136 } 1137 } 1138 1139 fn paths(root: &Path, instance: &str) -> ServiceSqlitePaths { 1140 let context = RuntimeContext::resolve( 1141 &RadrootsPathResolver::new( 1142 RadrootsPlatform::Linux, 1143 RadrootsHostEnvironment::default(), 1144 ), 1145 RuntimeContextBootstrap::new( 1146 RadrootsPathProfile::RepoLocal, 1147 Some(root.to_path_buf()), 1148 RuntimeContextSource::BootstrapCli, 1149 RuntimeContextSource::BootstrapCli, 1150 ) 1151 .expect("bootstrap"), 1152 ServiceId::new("myc").expect("service"), 1153 InstanceId::new(instance).expect("instance"), 1154 ) 1155 .expect("runtime context"); 1156 ServiceSqlitePaths::from_runtime_context(&context).expect("SQLite paths") 1157 } 1158 1159 fn prepare(paths: &ServiceSqlitePaths) { 1160 fs::create_dir_all(paths.state_database().parent().expect("state directory")) 1161 .expect("create state directory"); 1162 } 1163 1164 fn metadata(paths: &ServiceSqlitePaths) -> ServiceDatabaseMetadata { 1165 ServiceDatabaseMetadata::new( 1166 paths, 1167 SourceGeneration::new([7; 32]).expect("source generation"), 1168 NonZeroU32::new(1).expect("schema version"), 1169 1_700_000_000_000, 1170 crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), 1171 ) 1172 .expect("database metadata") 1173 } 1174 1175 fn schema_catalog(service_objects: Vec<crate::SchemaObject>) -> crate::SchemaCatalog { 1176 let migrations = crate::MigrationCatalog::new([]).expect("empty migrations"); 1177 let digest = 1178 crate::SchemaVersionCatalog::computed_digest(1, service_objects.iter().cloned()) 1179 .expect("schema digest"); 1180 let version = crate::SchemaVersionCatalog::new(1, service_objects, digest) 1181 .expect("schema version"); 1182 crate::SchemaCatalog::new(&migrations, [version]).expect("schema catalog") 1183 } 1184 1185 fn base_schema_catalog() -> crate::SchemaCatalog { 1186 schema_catalog(Vec::new()) 1187 } 1188 1189 fn service_schema_catalog() -> crate::SchemaCatalog { 1190 const SQL: &str = "CREATE TABLE service_schema (id INTEGER PRIMARY KEY)"; 1191 let object = crate::SchemaObject::new( 1192 crate::SchemaObjectKind::Table, 1193 "service_schema", 1194 "service_schema", 1195 SQL, 1196 crate::SchemaObject::computed_digest( 1197 crate::SchemaObjectKind::Table, 1198 "service_schema", 1199 "service_schema", 1200 SQL, 1201 ) 1202 .expect("service schema digest"), 1203 ) 1204 .expect("service schema object"); 1205 schema_catalog(vec![object]) 1206 } 1207 1208 #[tokio::test(flavor = "current_thread")] 1209 async fn successful_initialization_is_create_new_mode_0600_and_retains_authority() { 1210 let root = tempfile::tempdir().expect("root"); 1211 let paths = paths(root.path(), "success"); 1212 prepare(&paths); 1213 let metadata = metadata(&paths); 1214 let schema_catalog = service_schema_catalog(); 1215 let mut authority = initialize_database( 1216 &paths, 1217 OpenMode::Initialize, 1218 &metadata, 1219 &schema_catalog, 1220 |initializer| { 1221 Box::pin(async move { 1222 assert_eq!( 1223 format!("{initializer:?}"), 1224 "ServiceSqliteInitializer([redacted])" 1225 ); 1226 sqlx::query("CREATE TABLE service_schema (id INTEGER PRIMARY KEY)") 1227 .execute(initializer) 1228 .await?; 1229 Ok::<(), sqlx::Error>(()) 1230 }) 1231 }, 1232 ) 1233 .await 1234 .expect("initialize"); 1235 1236 let filesystem_metadata = fs::metadata(paths.state_database()).unwrap(); 1237 assert_eq!(filesystem_metadata.permissions().mode() & 0o777, 0o600); 1238 assert_eq!(filesystem_metadata.nlink(), 1); 1239 assert!(authority.is_held()); 1240 assert!(WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting).is_err()); 1241 authority.release().expect("release"); 1242 assert!( 1243 WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 1244 .expect("reacquire") 1245 .is_some() 1246 ); 1247 } 1248 1249 #[tokio::test(flavor = "current_thread")] 1250 async fn existing_regular_symlink_directory_and_hardlink_are_never_mutated() { 1251 for shape in ["regular", "symlink", "directory", "hardlink"] { 1252 let root = tempfile::tempdir().expect("root"); 1253 let paths = paths(root.path(), shape); 1254 prepare(&paths); 1255 let target = root.path().join("existing-target"); 1256 fs::write(&target, b"preserve-me").expect("target"); 1257 match shape { 1258 "regular" => fs::write(paths.state_database(), b"existing").unwrap(), 1259 "symlink" => symlink(&target, paths.state_database()).unwrap(), 1260 "directory" => fs::create_dir(paths.state_database()).unwrap(), 1261 "hardlink" => fs::hard_link(&target, paths.state_database()).unwrap(), 1262 _ => unreachable!(), 1263 } 1264 let called = Cell::new(false); 1265 let metadata = metadata(&paths); 1266 let schema_catalog = base_schema_catalog(); 1267 let error = initialize_database( 1268 &paths, 1269 OpenMode::Initialize, 1270 &metadata, 1271 &schema_catalog, 1272 |_| { 1273 called.set(true); 1274 Box::pin(ready(Ok::<(), CallbackFailure>(()))) 1275 }, 1276 ) 1277 .await 1278 .expect_err("existing state must fail"); 1279 assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); 1280 assert!(!called.get()); 1281 assert_eq!(fs::read(&target).unwrap(), b"preserve-me"); 1282 } 1283 } 1284 1285 #[tokio::test(flavor = "current_thread")] 1286 async fn non_initialize_modes_have_zero_callback_and_filesystem_effects() { 1287 for mode in [OpenMode::ReadWriteExisting, OpenMode::ReadOnlyInspection] { 1288 let root = tempfile::tempdir().expect("root"); 1289 let paths = paths(root.path(), "wrong-mode"); 1290 let called = Cell::new(false); 1291 let metadata = metadata(&paths); 1292 let schema_catalog = base_schema_catalog(); 1293 let error = initialize_database(&paths, mode, &metadata, &schema_catalog, |_| { 1294 called.set(true); 1295 Box::pin(ready(Ok::<(), CallbackFailure>(()))) 1296 }) 1297 .await 1298 .expect_err("mode must reject"); 1299 assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); 1300 assert!(!called.get()); 1301 assert!(!paths.state_database().exists()); 1302 assert!(!paths.state_lock().exists()); 1303 assert!(!paths.state_database().parent().unwrap().exists()); 1304 } 1305 } 1306 1307 #[tokio::test(flavor = "current_thread")] 1308 async fn mismatched_metadata_paths_have_zero_callback_and_filesystem_effects() { 1309 let root = tempfile::tempdir().expect("root"); 1310 let expected_paths = paths(root.path(), "expected"); 1311 let other = paths(root.path(), "other"); 1312 let called = Cell::new(false); 1313 let other_metadata = metadata(&other); 1314 let error = initialize_database( 1315 &expected_paths, 1316 OpenMode::Initialize, 1317 &other_metadata, 1318 &base_schema_catalog(), 1319 |_| { 1320 called.set(true); 1321 Box::pin(ready(Ok::<(), CallbackFailure>(()))) 1322 }, 1323 ) 1324 .await 1325 .expect_err("metadata paths must reject"); 1326 assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata); 1327 assert!(!called.get()); 1328 assert!(!expected_paths.state_database().exists()); 1329 assert!(!expected_paths.state_lock().exists()); 1330 let error = initialize_or_existing_database( 1331 &expected_paths, 1332 &other_metadata, 1333 &base_schema_catalog(), 1334 |_| { 1335 called.set(true); 1336 Box::pin(ready(Ok::<(), CallbackFailure>(()))) 1337 }, 1338 ) 1339 .await 1340 .expect_err("existing-or-new initialization must reject mismatched identity"); 1341 assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata); 1342 assert!(!called.get()); 1343 assert!(!expected_paths.state_database().exists()); 1344 assert!(!expected_paths.state_lock().exists()); 1345 } 1346 1347 #[tokio::test(flavor = "current_thread")] 1348 async fn callback_failure_cleans_up_releases_authority_and_preserves_trusted_cause() { 1349 let root = tempfile::tempdir().expect("root"); 1350 let paths = paths(root.path(), "callback-failure"); 1351 prepare(&paths); 1352 let metadata = metadata(&paths); 1353 let schema_catalog = base_schema_catalog(); 1354 let error = initialize_database( 1355 &paths, 1356 OpenMode::Initialize, 1357 &metadata, 1358 &schema_catalog, 1359 |_| Box::pin(ready(Err::<(), _>(CallbackFailure))), 1360 ) 1361 .await 1362 .expect_err("callback failure"); 1363 assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); 1364 assert!(!paths.state_database().exists()); 1365 assert!( 1366 WriterAuthority::acquire(&paths, OpenMode::Initialize) 1367 .unwrap() 1368 .is_some() 1369 ); 1370 1371 let display = error.to_string(); 1372 let debug = format!("{error:?}"); 1373 assert!(!display.contains("secret")); 1374 assert!(!debug.contains("secret")); 1375 let first = error.source().expect("initialization cause"); 1376 assert_eq!(first.to_string(), "SQLite schema initialization failed"); 1377 assert_eq!( 1378 first.source().map(ToString::to_string).as_deref(), 1379 Some("secret callback path=/private/state.sqlite") 1380 ); 1381 1382 let retry = initialize_database( 1383 &paths, 1384 OpenMode::Initialize, 1385 &metadata, 1386 &schema_catalog, 1387 |_| Box::pin(ready(Ok::<(), CallbackFailure>(()))), 1388 ) 1389 .await 1390 .expect("retry after cleanup"); 1391 assert!(retry.is_held()); 1392 assert!(paths.state_database().exists()); 1393 } 1394 1395 #[tokio::test(flavor = "current_thread")] 1396 async fn transaction_control_and_attachments_are_rejected_before_sqlite_compilation() { 1397 for scenario in ["commit", "rollback_begin", "attach_detach"] { 1398 let root = tempfile::tempdir().expect("root"); 1399 let paths = paths(root.path(), scenario); 1400 prepare(&paths); 1401 let external = root.path().join("must-not-exist.sqlite"); 1402 let external_for_callback = external.clone(); 1403 let metadata = metadata(&paths); 1404 let schema_catalog = base_schema_catalog(); 1405 let error = initialize_database( 1406 &paths, 1407 OpenMode::Initialize, 1408 &metadata, 1409 &schema_catalog, 1410 move |initializer| { 1411 Box::pin(async move { 1412 let sql = match scenario { 1413 "commit" => "COMMIT".to_owned(), 1414 "rollback_begin" => { 1415 "ROLLBACK; BEGIN DEFERRED; CREATE TABLE escaped(value INTEGER)" 1416 .to_owned() 1417 } 1418 "attach_detach" => format!( 1419 "ATTACH DATABASE '{}' AS extra; DETACH DATABASE extra", 1420 external_for_callback.display() 1421 ), 1422 _ => unreachable!(), 1423 }; 1424 let _ = sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 1425 .execute(initializer) 1426 .await; 1427 Ok::<(), Infallible>(()) 1428 }) 1429 }, 1430 ) 1431 .await 1432 .expect_err("forbidden statement control must fail initialization"); 1433 assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); 1434 for rendered in [error.to_string(), format!("{error:?}")] { 1435 assert!(!rendered.contains("must-not-exist")); 1436 assert!(!rendered.contains("ATTACH")); 1437 assert!(!rendered.contains("ROLLBACK")); 1438 } 1439 assert!(!paths.state_database().exists()); 1440 assert!(!external.exists()); 1441 assert!( 1442 WriterAuthority::acquire(&paths, OpenMode::Initialize) 1443 .expect("reacquire") 1444 .is_some() 1445 ); 1446 } 1447 } 1448 1449 #[tokio::test(flavor = "current_thread")] 1450 async fn metadata_or_migration_ledger_conflict_cleans_the_exact_reserved_database() { 1451 let root = tempfile::tempdir().expect("root"); 1452 for (instance, statement, expected_kind) in [ 1453 ( 1454 "metadata-failure", 1455 "CREATE TABLE radroots_service_metadata (value TEXT)", 1456 ServiceSqliteErrorKind::Metadata, 1457 ), 1458 ( 1459 "migration-ledger-failure", 1460 "CREATE TABLE schema_migrations (value TEXT)", 1461 ServiceSqliteErrorKind::Migration, 1462 ), 1463 ] { 1464 let paths = paths(root.path(), instance); 1465 prepare(&paths); 1466 let metadata = metadata(&paths); 1467 let schema_catalog = base_schema_catalog(); 1468 let error = initialize_database( 1469 &paths, 1470 OpenMode::Initialize, 1471 &metadata, 1472 &schema_catalog, 1473 move |initializer| { 1474 Box::pin(async move { 1475 sqlx::query(statement).execute(initializer).await?; 1476 Ok::<(), sqlx::Error>(()) 1477 }) 1478 }, 1479 ) 1480 .await 1481 .expect_err("conflicting governed table must fail"); 1482 1483 assert_eq!(error.kind(), expected_kind); 1484 assert!(!paths.state_database().exists()); 1485 assert!( 1486 WriterAuthority::acquire(&paths, OpenMode::Initialize) 1487 .expect("reacquire") 1488 .is_some() 1489 ); 1490 } 1491 } 1492 1493 #[tokio::test(flavor = "current_thread")] 1494 async fn schema_catalog_mismatch_cleans_the_exact_reserved_database() { 1495 let root = tempfile::tempdir().expect("root"); 1496 let paths = paths(root.path(), "schema-mismatch"); 1497 prepare(&paths); 1498 let metadata = metadata(&paths); 1499 let schema_catalog = base_schema_catalog(); 1500 let error = initialize_database( 1501 &paths, 1502 OpenMode::Initialize, 1503 &metadata, 1504 &schema_catalog, 1505 |initializer| { 1506 Box::pin(async move { 1507 sqlx::query("CREATE TABLE unexpected (value INTEGER)") 1508 .execute(initializer) 1509 .await?; 1510 Ok::<(), sqlx::Error>(()) 1511 }) 1512 }, 1513 ) 1514 .await 1515 .expect_err("unlisted schema object must fail initialization"); 1516 assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); 1517 assert!(!paths.state_database().exists()); 1518 assert!( 1519 WriterAuthority::acquire(&paths, OpenMode::Initialize) 1520 .unwrap() 1521 .is_some() 1522 ); 1523 } 1524 1525 #[tokio::test(flavor = "current_thread")] 1526 async fn cancellation_rolls_back_and_releases_authority() { 1527 let root = tempfile::tempdir().expect("root"); 1528 let paths = paths(root.path(), "cancelled"); 1529 prepare(&paths); 1530 let metadata = metadata(&paths); 1531 let schema_catalog = base_schema_catalog(); 1532 let task_paths = paths.clone(); 1533 let reached = Arc::new(Notify::new()); 1534 let reached_from_callback = Arc::clone(&reached); 1535 let task = tokio::spawn(async move { 1536 initialize_database( 1537 &task_paths, 1538 OpenMode::Initialize, 1539 &metadata, 1540 &schema_catalog, 1541 move |initializer| { 1542 Box::pin(async move { 1543 sqlx::query("CREATE TABLE cancelled_probe (value INTEGER)") 1544 .execute(initializer) 1545 .await?; 1546 reached_from_callback.notify_one(); 1547 pending::<Result<(), sqlx::Error>>().await 1548 }) 1549 }, 1550 ) 1551 .await 1552 }); 1553 reached.notified().await; 1554 assert!(paths.state_database().exists()); 1555 assert!(WriterAuthority::acquire(&paths, OpenMode::Initialize).is_err()); 1556 task.abort(); 1557 task.await.expect_err("initialization task is cancelled"); 1558 assert!(!paths.state_database().exists()); 1559 assert!( 1560 WriterAuthority::acquire(&paths, OpenMode::Initialize) 1561 .unwrap() 1562 .is_some() 1563 ); 1564 } 1565 1566 #[tokio::test(flavor = "current_thread")] 1567 async fn replacement_is_detected_and_never_deleted() { 1568 let root = tempfile::tempdir().expect("root"); 1569 let paths = paths(root.path(), "replacement"); 1570 prepare(&paths); 1571 let metadata = metadata(&paths); 1572 let schema_catalog = base_schema_catalog(); 1573 let replacement_path = paths.state_database().to_path_buf(); 1574 let error = initialize_database( 1575 &paths, 1576 OpenMode::Initialize, 1577 &metadata, 1578 &schema_catalog, 1579 move |_| { 1580 Box::pin(async move { 1581 fs::remove_file(&replacement_path)?; 1582 fs::write(&replacement_path, b"replacement")?; 1583 Ok::<(), io::Error>(()) 1584 }) 1585 }, 1586 ) 1587 .await 1588 .expect_err("replacement must fail"); 1589 assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); 1590 assert_eq!(fs::read(paths.state_database()).unwrap(), b"replacement"); 1591 } 1592 1593 #[tokio::test(flavor = "current_thread")] 1594 async fn parent_directory_replacement_cannot_rebind_the_callback_path() { 1595 let root = tempfile::tempdir().expect("root"); 1596 let paths = paths(root.path(), "parent-replacement"); 1597 prepare(&paths); 1598 let metadata = metadata(&paths); 1599 let schema_catalog = base_schema_catalog(); 1600 let state_directory = paths.state_database().parent().unwrap().to_path_buf(); 1601 let displaced_directory = state_directory.with_file_name("parent-replacement-old"); 1602 let displaced_for_callback = displaced_directory.clone(); 1603 let replacement_path = paths.state_database().to_path_buf(); 1604 1605 let error = initialize_database( 1606 &paths, 1607 OpenMode::Initialize, 1608 &metadata, 1609 &schema_catalog, 1610 move |_| { 1611 Box::pin(async move { 1612 fs::rename(&state_directory, &displaced_for_callback)?; 1613 fs::create_dir(&state_directory)?; 1614 fs::write(&replacement_path, b"replacement")?; 1615 fs::set_permissions(&replacement_path, fs::Permissions::from_mode(0o600))?; 1616 Ok::<(), io::Error>(()) 1617 }) 1618 }, 1619 ) 1620 .await 1621 .expect_err("canonical path replacement must fail"); 1622 1623 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 1624 assert_eq!(fs::read(paths.state_database()).unwrap(), b"replacement"); 1625 assert!(!displaced_directory.join("state.sqlite").exists()); 1626 } 1627 1628 #[tokio::test(flavor = "current_thread")] 1629 async fn injected_sync_and_cleanup_failures_preserve_exact_ordering() { 1630 let scenarios = [ 1631 "reservation-sync", 1632 "database-sync", 1633 "directory-sync", 1634 "cleanup", 1635 ]; 1636 for scenario in scenarios { 1637 let root = tempfile::tempdir().expect("root"); 1638 let paths = paths(root.path(), scenario); 1639 prepare(&paths); 1640 let metadata = metadata(&paths); 1641 let schema_catalog = base_schema_catalog(); 1642 let authority = WriterAuthority::acquire(&paths, OpenMode::Initialize) 1643 .unwrap() 1644 .unwrap(); 1645 let operations = RecordingOperations::default(); 1646 match scenario { 1647 "reservation-sync" => { 1648 operations.fail_directory_sync_on_call.set(Some(1)); 1649 } 1650 "database-sync" => operations.fail_database_sync.set(true), 1651 "directory-sync" => { 1652 operations.fail_directory_sync_on_call.set(Some(2)); 1653 } 1654 "cleanup" => operations.fail_unlink.set(true), 1655 _ => unreachable!(), 1656 } 1657 let result = if scenario == "cleanup" { 1658 initialize_with_ops( 1659 &paths, 1660 authority, 1661 &metadata, 1662 &schema_catalog, 1663 |_| Box::pin(ready(Err::<(), _>(CallbackFailure))), 1664 &operations, 1665 &crate::failpoint::DurabilityFailpoints::default(), 1666 ) 1667 .await 1668 } else { 1669 initialize_with_ops( 1670 &paths, 1671 authority, 1672 &metadata, 1673 &schema_catalog, 1674 |_| Box::pin(ready(Ok::<(), CallbackFailure>(()))), 1675 &operations, 1676 &crate::failpoint::DurabilityFailpoints::default(), 1677 ) 1678 .await 1679 }; 1680 assert_eq!( 1681 result.expect_err("injected failure").kind(), 1682 ServiceSqliteErrorKind::Create 1683 ); 1684 let events = operations.events.borrow(); 1685 match scenario { 1686 "reservation-sync" => assert_eq!( 1687 events.as_slice(), 1688 ["sync_directory", "unlink_database", "sync_directory"] 1689 ), 1690 "database-sync" => assert_eq!( 1691 events.as_slice(), 1692 [ 1693 "sync_directory", 1694 "sync_database", 1695 "unlink_database", 1696 "sync_directory" 1697 ] 1698 ), 1699 "directory-sync" => assert_eq!( 1700 events.as_slice(), 1701 [ 1702 "sync_directory", 1703 "sync_database", 1704 "sync_directory", 1705 "unlink_database", 1706 "sync_directory" 1707 ] 1708 ), 1709 "cleanup" => { 1710 assert_eq!( 1711 events.as_slice(), 1712 ["sync_directory", "unlink_database", "unlink_database"] 1713 ) 1714 } 1715 _ => unreachable!(), 1716 } 1717 } 1718 } 1719 1720 #[test] 1721 fn rollback_failure_classifier_preserves_primary_and_cleanup_precedence() { 1722 let ordinary = 1723 || InitializationCause::new(InitializationFailureKind::SchemaInitializationFailed); 1724 let replaced = || InitializationCause::new(InitializationFailureKind::DatabaseReplaced); 1725 let cleanup = || InitializationCause::new(InitializationFailureKind::CleanupFailed); 1726 1727 let ordinary_success = rollback_failure(ordinary(), Ok(())); 1728 assert_eq!(ordinary_success.kind(), ServiceSqliteErrorKind::Create); 1729 assert_eq!( 1730 ordinary_success 1731 .source() 1732 .map(ToString::to_string) 1733 .as_deref(), 1734 Some("SQLite schema initialization failed") 1735 ); 1736 1737 let replaced_failure = rollback_failure(replaced(), Err(cleanup())); 1738 assert_eq!(replaced_failure.kind(), ServiceSqliteErrorKind::Create); 1739 assert_eq!( 1740 replaced_failure 1741 .source() 1742 .map(ToString::to_string) 1743 .as_deref(), 1744 Some("SQLite state identity changed during initialization") 1745 ); 1746 1747 let cleanup_failure = rollback_failure(ordinary(), Err(cleanup())); 1748 assert_eq!(cleanup_failure.kind(), ServiceSqliteErrorKind::Create); 1749 let outer = cleanup_failure.source().expect("cleanup source"); 1750 assert_eq!(outer.to_string(), "SQLite initialization cleanup failed"); 1751 assert_eq!( 1752 outer.source().map(ToString::to_string).as_deref(), 1753 Some("SQLite schema initialization failed") 1754 ); 1755 } 1756 1757 #[tokio::test(flavor = "current_thread")] 1758 async fn every_initialization_durability_edge_fails_once_and_rolls_back() { 1759 use crate::failpoint::{DurabilityFailpoint, DurabilityFailpoints}; 1760 1761 for (index, (point, expected_events)) in [ 1762 (DurabilityFailpoint::InitializeBeforeCreate, &[][..]), 1763 ( 1764 DurabilityFailpoint::InitializeAfterCreate, 1765 &["unlink_database", "sync_directory"][..], 1766 ), 1767 ( 1768 DurabilityFailpoint::InitializeBeforeReservationDirectorySync, 1769 &["unlink_database", "sync_directory"][..], 1770 ), 1771 ( 1772 DurabilityFailpoint::InitializeAfterReservationDirectorySync, 1773 &["sync_directory", "unlink_database", "sync_directory"][..], 1774 ), 1775 ( 1776 DurabilityFailpoint::InitializeBeforeFileSync, 1777 &["sync_directory", "unlink_database", "sync_directory"][..], 1778 ), 1779 ( 1780 DurabilityFailpoint::InitializeAfterFileSync, 1781 &[ 1782 "sync_directory", 1783 "sync_database", 1784 "unlink_database", 1785 "sync_directory", 1786 ][..], 1787 ), 1788 ( 1789 DurabilityFailpoint::InitializeBeforeCommitDirectorySync, 1790 &[ 1791 "sync_directory", 1792 "sync_database", 1793 "unlink_database", 1794 "sync_directory", 1795 ][..], 1796 ), 1797 ( 1798 DurabilityFailpoint::InitializeAfterCommitDirectorySync, 1799 &[ 1800 "sync_directory", 1801 "sync_database", 1802 "sync_directory", 1803 "unlink_database", 1804 "sync_directory", 1805 ][..], 1806 ), 1807 ] 1808 .into_iter() 1809 .enumerate() 1810 { 1811 let root = tempfile::tempdir().expect("root"); 1812 let paths = paths(root.path(), &format!("failpoint-{index}")); 1813 prepare(&paths); 1814 let metadata = metadata(&paths); 1815 let schema_catalog = base_schema_catalog(); 1816 let authority = WriterAuthority::acquire(&paths, OpenMode::Initialize) 1817 .expect("authority acquisition") 1818 .expect("initialize authority"); 1819 let failpoints = DurabilityFailpoints::armed(point); 1820 let operations = RecordingOperations::default(); 1821 let error = initialize_with_ops( 1822 &paths, 1823 authority, 1824 &metadata, 1825 &schema_catalog, 1826 |_| Box::pin(ready(Ok::<(), CallbackFailure>(()))), 1827 &operations, 1828 &failpoints, 1829 ) 1830 .await 1831 .expect_err("durability edge must fail"); 1832 assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); 1833 assert!(failpoints.fired()); 1834 assert_eq!( 1835 failpoints.reached().last(), 1836 Some(&point), 1837 "named edge must be the last reached boundary" 1838 ); 1839 assert_eq!( 1840 operations.events.borrow().as_slice(), 1841 expected_events, 1842 "named before/after edge must bracket the expected sync operation" 1843 ); 1844 assert!(!paths.state_database().exists()); 1845 let mut recovered = WriterAuthority::acquire(&paths, OpenMode::Initialize) 1846 .expect("reacquire after rollback") 1847 .expect("initialize authority"); 1848 recovered.release().expect("release recovered authority"); 1849 } 1850 } 1851 } 1852 } 1853 1854 #[cfg(any(target_os = "linux", target_os = "macos"))] 1855 use supported::{SystemInitializationOperations, initialize_with_ops};