open.rs (137670B)
1 //! Instance-bound SQLite paths and declarative open modes. 2 3 use core::fmt; 4 use std::{ 5 error::Error, 6 path::{Path, PathBuf}, 7 }; 8 9 #[cfg(any(target_os = "linux", target_os = "macos"))] 10 use std::{ 11 fs::File, 12 sync::{ 13 Arc, Mutex, 14 atomic::{AtomicBool, Ordering}, 15 }, 16 }; 17 18 #[cfg(any(target_os = "linux", target_os = "macos"))] 19 use fs2::FileExt; 20 #[cfg(any(target_os = "linux", target_os = "macos"))] 21 use futures::future::BoxFuture; 22 use radroots_runtime_paths::{ 23 InstanceId, RuntimeContext, ServiceId, default_service_instance_artifacts, 24 }; 25 use serde::Serialize; 26 #[cfg(any(target_os = "linux", target_os = "macos"))] 27 use sqlx::{ 28 ConnectOptions, Connection, Sqlite, SqliteConnection, SqlitePool, 29 pool::PoolConnection, 30 sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions, SqliteSynchronous}, 31 }; 32 33 #[cfg(any(target_os = "linux", target_os = "macos"))] 34 use crate::{ 35 ExistingServiceDatabaseIntent, MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, 36 MigrationCatalog, SchemaCatalog, ServiceDatabaseIdentity, ServiceDatabaseMetadata, 37 ServiceSqliteConnectionOptions, ServiceSqliteError, ServiceSqliteErrorKind, WriterAuthority, 38 }; 39 40 #[cfg(any(target_os = "linux", target_os = "macos"))] 41 const STATEMENT_CACHE_CAPACITY: usize = 100; 42 #[cfg(any(target_os = "linux", target_os = "macos"))] 43 const COMMAND_BUFFER_CAPACITY: usize = 50; 44 #[cfg(any(target_os = "linux", target_os = "macos"))] 45 const ROW_BUFFER_CAPACITY: usize = 50; 46 #[cfg(any(target_os = "linux", target_os = "macos"))] 47 const WAL_FILE_NAME: &str = "state.sqlite-wal"; 48 #[cfg(any(target_os = "linux", target_os = "macos"))] 49 const SHARED_MEMORY_FILE_NAME: &str = "state.sqlite-shm"; 50 51 #[cfg(any(target_os = "linux", target_os = "macos"))] 52 #[derive(Clone)] 53 struct PoolConnectionValidation { 54 binding: DirectoryBinding, 55 paths: ServiceSqlitePaths, 56 identity: ServiceDatabaseIdentity, 57 catalog: MigrationCatalog, 58 schema_catalog: SchemaCatalog, 59 mode: OpenMode, 60 policy: ServiceSqliteConnectionOptions, 61 } 62 63 #[cfg(any(target_os = "linux", target_os = "macos"))] 64 #[derive(Clone, Copy)] 65 enum ServiceDatabaseExpectation<'a> { 66 Exact(&'a ServiceDatabaseIdentity), 67 Existing(&'a ExistingServiceDatabaseIntent), 68 } 69 70 #[cfg(any(target_os = "linux", target_os = "macos"))] 71 impl ServiceDatabaseExpectation<'_> { 72 fn matches_paths(self, paths: &ServiceSqlitePaths) -> bool { 73 match self { 74 Self::Exact(identity) => identity.matches_paths(paths), 75 Self::Existing(intent) => intent.matches_paths(paths), 76 } 77 } 78 79 fn supported_state_schema_version(self) -> core::num::NonZeroU32 { 80 match self { 81 Self::Exact(identity) => identity.supported_state_schema_version(), 82 Self::Existing(intent) => intent.supported_state_schema_version(), 83 } 84 } 85 86 async fn verify_metadata( 87 self, 88 connection: &mut SqliteConnection, 89 ) -> Result<ServiceDatabaseMetadata, ServiceSqliteError> { 90 match self { 91 Self::Exact(identity) => { 92 crate::metadata::verify_database_metadata(connection, identity).await 93 } 94 Self::Existing(intent) => { 95 crate::metadata::verify_existing_database_intent(connection, intent).await 96 } 97 } 98 } 99 100 fn exact_identity(self, metadata: &ServiceDatabaseMetadata) -> ServiceDatabaseIdentity { 101 match self { 102 Self::Exact(identity) => identity.clone(), 103 Self::Existing(intent) => intent.identity_for(metadata), 104 } 105 } 106 } 107 108 #[cfg(any(target_os = "linux", target_os = "macos"))] 109 enum PoolConnectionValidationFailure { 110 Authority, 111 Pragma(sqlx::Error), 112 PolicyMismatch, 113 Metadata, 114 Migration, 115 Integrity, 116 } 117 118 #[cfg(any(target_os = "linux", target_os = "macos"))] 119 #[derive(Clone)] 120 struct PoolConnectionFailureFlags { 121 authority: Arc<AtomicBool>, 122 metadata: Arc<AtomicBool>, 123 migration: Arc<AtomicBool>, 124 integrity: Arc<AtomicBool>, 125 pragma: Arc<AtomicBool>, 126 } 127 128 #[cfg(any(target_os = "linux", target_os = "macos"))] 129 impl PoolConnectionFailureFlags { 130 fn record(&self, failure: &PoolConnectionValidationFailure) { 131 match failure { 132 PoolConnectionValidationFailure::Authority => &self.authority, 133 PoolConnectionValidationFailure::Metadata => &self.metadata, 134 PoolConnectionValidationFailure::Migration => &self.migration, 135 PoolConnectionValidationFailure::Integrity => &self.integrity, 136 PoolConnectionValidationFailure::Pragma(_) 137 | PoolConnectionValidationFailure::PolicyMismatch => &self.pragma, 138 } 139 .store(true, Ordering::Release); 140 } 141 142 fn kind(&self) -> ServiceSqliteErrorKind { 143 connection_failure_kind( 144 self.authority.load(Ordering::Acquire), 145 self.metadata.load(Ordering::Acquire), 146 self.migration.load(Ordering::Acquire), 147 self.integrity.load(Ordering::Acquire), 148 self.pragma.load(Ordering::Acquire), 149 ) 150 } 151 } 152 153 #[cfg(any(target_os = "linux", target_os = "macos"))] 154 impl PoolConnectionValidationFailure { 155 fn into_sqlx(self) -> sqlx::Error { 156 match self { 157 Self::Pragma(source) => source, 158 Self::Authority => { 159 sqlx::Error::Protocol("SQLite connection authority mismatch".to_owned()) 160 } 161 Self::PolicyMismatch => { 162 sqlx::Error::Protocol("SQLite connection policy mismatch".to_owned()) 163 } 164 Self::Metadata => { 165 sqlx::Error::Protocol("SQLite connection metadata mismatch".to_owned()) 166 } 167 Self::Migration | Self::Integrity => { 168 sqlx::Error::Protocol("SQLite migration history mismatch".to_owned()) 169 } 170 } 171 } 172 } 173 174 /// Canonical database and writer-lock paths for one validated service instance. 175 /// 176 /// Callers cannot forge paths or rebind the service and instance independently: 177 /// 178 /// ```compile_fail 179 /// use std::path::PathBuf; 180 /// use radroots_runtime_paths::{InstanceId, ServiceId}; 181 /// use radroots_service_sqlite::ServiceSqlitePaths; 182 /// 183 /// let _ = ServiceSqlitePaths { 184 /// service: ServiceId::new("example").unwrap(), 185 /// instance: InstanceId::new("primary").unwrap(), 186 /// state_database: PathBuf::from("/tmp/alternate.sqlite"), 187 /// state_lock: PathBuf::from("/tmp/alternate.lock"), 188 /// }; 189 /// ``` 190 #[derive(Clone, PartialEq, Eq)] 191 pub struct ServiceSqlitePaths { 192 service: ServiceId, 193 instance: InstanceId, 194 state_database: PathBuf, 195 state_lock: PathBuf, 196 } 197 198 impl ServiceSqlitePaths { 199 /// Derives the fixed SQLite artifacts from one immutable runtime context. 200 pub fn from_runtime_context(context: &RuntimeContext) -> Result<Self, ServiceSqlitePathError> { 201 validate_state_directory(context.paths().state())?; 202 let artifacts = default_service_instance_artifacts(context.paths()); 203 Ok(Self { 204 service: context.service().clone(), 205 instance: context.instance().clone(), 206 state_database: artifacts.state_database().to_path_buf(), 207 state_lock: artifacts.state_lock().to_path_buf(), 208 }) 209 } 210 211 /// Returns the validated service identity bound to these paths. 212 #[must_use] 213 pub fn service(&self) -> &ServiceId { 214 &self.service 215 } 216 217 /// Returns the validated instance identity bound to these paths. 218 #[must_use] 219 pub fn instance(&self) -> &InstanceId { 220 &self.instance 221 } 222 223 /// Returns the canonical `state.sqlite` path. 224 #[must_use] 225 pub fn state_database(&self) -> &Path { 226 &self.state_database 227 } 228 229 /// Returns the canonical retained `state.lock` path. 230 #[must_use] 231 pub fn state_lock(&self) -> &Path { 232 &self.state_lock 233 } 234 } 235 236 impl fmt::Debug for ServiceSqlitePaths { 237 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 238 formatter 239 .debug_struct("ServiceSqlitePaths") 240 .field("service", &self.service) 241 .field("instance", &self.instance) 242 .field("state_database", &"[redacted]") 243 .field("state_lock", &"[redacted]") 244 .finish() 245 } 246 } 247 248 /// Path-shape failure detected before any filesystem or SQLite operation. 249 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 250 pub enum ServiceSqlitePathError { 251 RelativeStateDirectory, 252 MissingStateDirectoryParent, 253 } 254 255 impl fmt::Display for ServiceSqlitePathError { 256 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 257 match self { 258 Self::RelativeStateDirectory => { 259 formatter.write_str("SQLite state directory must be absolute") 260 } 261 Self::MissingStateDirectoryParent => { 262 formatter.write_str("SQLite state directory must have a parent") 263 } 264 } 265 } 266 } 267 268 impl Error for ServiceSqlitePathError {} 269 270 fn validate_state_directory(path: &Path) -> Result<(), ServiceSqlitePathError> { 271 if !path.is_absolute() { 272 return Err(ServiceSqlitePathError::RelativeStateDirectory); 273 } 274 if path.parent().is_none() { 275 return Err(ServiceSqlitePathError::MissingStateDirectoryParent); 276 } 277 Ok(()) 278 } 279 280 /// Declarative behavior for opening one service-owned SQLite database. 281 #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)] 282 #[serde(rename_all = "snake_case")] 283 pub enum OpenMode { 284 Initialize, 285 ReadWriteExisting, 286 ReadOnlyInspection, 287 } 288 289 impl OpenMode { 290 /// Returns whether this mode permits creating missing state. 291 #[must_use] 292 pub const fn may_create(self) -> bool { 293 matches!(self, Self::Initialize) 294 } 295 296 /// Returns whether state must already exist before opening. 297 #[must_use] 298 pub const fn requires_existing(self) -> bool { 299 !matches!(self, Self::Initialize) 300 } 301 302 /// Returns whether exclusive writer authority is required. 303 #[must_use] 304 pub const fn requires_writer_authority(self) -> bool { 305 !matches!(self, Self::ReadOnlyInspection) 306 } 307 } 308 309 #[cfg(any(target_os = "linux", target_os = "macos"))] 310 pub(crate) struct PrivateConnectionPool { 311 pool: SqlitePool, 312 binding: DirectoryBinding, 313 paths: ServiceSqlitePaths, 314 identity: ServiceDatabaseIdentity, 315 catalog: MigrationCatalog, 316 schema_catalog: SchemaCatalog, 317 mode: OpenMode, 318 policy: ServiceSqliteConnectionOptions, 319 resources: Arc<Mutex<PrivateConnectionResources>>, 320 close_driver: tokio::sync::Mutex<PrivateCloseDriver>, 321 #[cfg(test)] 322 close_phase: std::sync::atomic::AtomicU8, 323 authority_failure: Arc<AtomicBool>, 324 metadata_failure: Arc<AtomicBool>, 325 migration_failure: Arc<AtomicBool>, 326 integrity_failure: Arc<AtomicBool>, 327 pragma_failure: Arc<AtomicBool>, 328 } 329 330 #[cfg(any(target_os = "linux", target_os = "macos"))] 331 struct PrivateConnectionResources { 332 authority: Option<WriterAuthority>, 333 inspection_guard: Option<ReadOnlyInspectionGuard>, 334 } 335 336 #[cfg(any(target_os = "linux", target_os = "macos"))] 337 #[derive(Clone)] 338 pub(crate) struct BackupSourceValidator { 339 binding: DirectoryBinding, 340 paths: ServiceSqlitePaths, 341 resources: Arc<Mutex<PrivateConnectionResources>>, 342 } 343 344 #[cfg(any(target_os = "linux", target_os = "macos"))] 345 impl BackupSourceValidator { 346 pub(crate) fn validate(&self) -> Result<(), ServiceSqliteError> { 347 let resources = self.resources.lock().map_err(|_| { 348 connection_error( 349 ServiceSqliteErrorKind::Authority, 350 ConnectionFailureKind::AuthorityMismatch, 351 ) 352 })?; 353 resources 354 .authority 355 .as_ref() 356 .ok_or_else(|| { 357 connection_error( 358 ServiceSqliteErrorKind::Authority, 359 ConnectionFailureKind::AuthorityMismatch, 360 ) 361 })? 362 .validate_for(&self.paths)?; 363 self.binding.validate(&self.paths) 364 } 365 } 366 367 #[cfg(any(target_os = "linux", target_os = "macos"))] 368 enum PrivateCloseDriver { 369 Pending, 370 Connecting(BoxFuture<'static, Result<SqliteConnection, ServiceSqliteError>>), 371 Connected(SqliteConnection), 372 Closing { 373 future: BoxFuture<'static, Result<(), ServiceSqliteError>>, 374 authority_error: Option<ServiceSqliteError>, 375 checkpoint_error: Option<ServiceSqliteError>, 376 connection_close_error: Option<ServiceSqliteError>, 377 }, 378 Complete(Option<ServiceSqliteErrorKind>), 379 } 380 381 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 382 pub(crate) const TEST_CLOSE_PHASE_CHECKPOINT: u8 = 4; 383 384 #[cfg(any(target_os = "linux", target_os = "macos"))] 385 impl PrivateConnectionPool { 386 fn connection_failure_kind(&self) -> ServiceSqliteErrorKind { 387 connection_failure_kind( 388 self.authority_failure.load(Ordering::Acquire), 389 self.metadata_failure.load(Ordering::Acquire), 390 self.migration_failure.load(Ordering::Acquire), 391 self.integrity_failure.load(Ordering::Acquire), 392 self.pragma_failure.load(Ordering::Acquire), 393 ) 394 } 395 396 pub(crate) fn validate(&self) -> Result<(), ServiceSqliteError> { 397 let resources = self.resources.lock().map_err(|_| { 398 connection_error( 399 ServiceSqliteErrorKind::Authority, 400 ConnectionFailureKind::AuthorityMismatch, 401 ) 402 })?; 403 match self.mode { 404 OpenMode::Initialize | OpenMode::ReadWriteExisting => resources 405 .authority 406 .as_ref() 407 .ok_or_else(|| { 408 connection_error( 409 ServiceSqliteErrorKind::Authority, 410 ConnectionFailureKind::AuthorityMismatch, 411 ) 412 })? 413 .validate_for(&self.paths)?, 414 OpenMode::ReadOnlyInspection => resources 415 .inspection_guard 416 .as_ref() 417 .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))? 418 .validate_for(&self.paths)?, 419 } 420 self.binding.validate(&self.paths) 421 } 422 423 pub(crate) const fn mode(&self) -> OpenMode { 424 self.mode 425 } 426 427 pub(crate) fn identity(&self) -> &ServiceDatabaseIdentity { 428 &self.identity 429 } 430 431 pub(crate) async fn database_metadata( 432 &self, 433 ) -> Result<ServiceDatabaseMetadata, ServiceSqliteError> { 434 self.validate()?; 435 let mut connection = self.acquire().await?; 436 let result = 437 crate::metadata::verify_database_metadata(&mut connection, &self.identity).await; 438 self.validate()?; 439 result 440 } 441 442 pub(crate) fn backup_source_validator(&self) -> BackupSourceValidator { 443 BackupSourceValidator { 444 binding: self.binding.clone(), 445 paths: self.paths.clone(), 446 resources: Arc::clone(&self.resources), 447 } 448 } 449 450 pub(crate) fn catalog(&self) -> &MigrationCatalog { 451 &self.catalog 452 } 453 454 pub(crate) fn schema_catalog(&self) -> &SchemaCatalog { 455 &self.schema_catalog 456 } 457 458 #[cfg(test)] 459 pub(crate) fn close_phase(&self) -> u8 { 460 self.close_phase.load(Ordering::Acquire) 461 } 462 463 pub(crate) async fn acquire(&self) -> Result<PoolConnection<Sqlite>, ServiceSqliteError> { 464 self.validate()?; 465 let result = self.pool.acquire().await; 466 self.validate()?; 467 let mut connection = 468 result.map_err(|source| connection_source(self.connection_failure_kind(), source))?; 469 let history = crate::migration::verify_migration_history( 470 &mut connection, 471 &self.catalog, 472 &self.schema_catalog, 473 true, 474 ) 475 .await; 476 self.validate()?; 477 history?; 478 Ok(connection) 479 } 480 481 pub(crate) async fn apply_migrations( 482 &self, 483 applied_at: MigrationAppliedAtUnixSeconds, 484 build: &MigrationBuildIdentity, 485 callbacks: &[crate::migration::MigrationCallbackBinding], 486 ) -> Result<crate::migration::MigrationApplicationOutcome, ServiceSqliteError> { 487 self.validate_writer_authority()?; 488 let acquired = self.pool.acquire().await; 489 self.validate_writer_authority()?; 490 let mut connection = 491 acquired.map_err(|source| connection_source(self.connection_failure_kind(), source))?; 492 // Migration execution installs connection-local fail-closed guards. Always 493 // discard this one-time connection so cancellation cannot return a guarded 494 // or callback-altered handle to the pool. 495 connection.close_on_drop(); 496 let mut validate_authority = || self.validate_writer_authority(); 497 let result = crate::migration::apply_governed_migrations( 498 &mut connection, 499 &self.catalog, 500 &self.schema_catalog, 501 applied_at, 502 build, 503 callbacks, 504 &mut validate_authority, 505 ) 506 .await; 507 self.validate_writer_authority()?; 508 result 509 } 510 511 fn validate_writer_authority(&self) -> Result<(), ServiceSqliteError> { 512 let resources = self.resources.lock().map_err(|_| { 513 connection_error( 514 ServiceSqliteErrorKind::Authority, 515 ConnectionFailureKind::AuthorityMismatch, 516 ) 517 })?; 518 resources 519 .authority 520 .as_ref() 521 .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Authority))? 522 .validate_for(&self.paths)?; 523 self.binding.validate(&self.paths) 524 } 525 526 pub(crate) async fn close(self) -> Option<WriterAuthority> { 527 self.pool.close().await; 528 let mut resources = match self.resources.lock() { 529 Ok(resources) => resources, 530 Err(poisoned) => poisoned.into_inner(), 531 }; 532 resources.inspection_guard.take(); 533 resources.authority.take() 534 } 535 536 /// The outer error means authority release was not proven and close must retry. 537 /// The inner result is terminal and may be cached by the host. 538 pub(crate) async fn close_explicit( 539 &self, 540 failpoints: &crate::failpoint::DurabilityFailpoints, 541 ) -> Result<Result<(), ServiceSqliteError>, ServiceSqliteError> { 542 let mut authority_error = self.validate().err(); 543 let before_drain = failpoints 544 .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeDrain) 545 .map_err(|source| { 546 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) 547 }); 548 authority_error = authority_error.or_else(|| self.validate().err()); 549 if authority_error.is_none() { 550 before_drain?; 551 } 552 self.pool.close().await; 553 let after_drain = failpoints 554 .hit(crate::failpoint::DurabilityFailpoint::CloseAfterDrain) 555 .map_err(|source| { 556 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) 557 }); 558 authority_error = authority_error.or_else(|| self.validate().err()); 559 let terminal = match authority_error { 560 Some(authority) => Err(authority), 561 None => { 562 after_drain?; 563 if self.mode.requires_writer_authority() { 564 self.drive_writable_close(failpoints).await 565 } else { 566 self.validate() 567 } 568 } 569 }; 570 let release_error = self.release_resources(failpoints)?; 571 Ok(match (terminal, release_error) { 572 (Err(error), _) if error.kind() == ServiceSqliteErrorKind::Authority => Err(error), 573 (_, Some(release_error)) => Err(release_error), 574 (terminal, None) => terminal, 575 }) 576 } 577 578 async fn drive_writable_close( 579 &self, 580 failpoints: &crate::failpoint::DurabilityFailpoints, 581 ) -> Result<(), ServiceSqliteError> { 582 let mut driver = self.close_driver.lock().await; 583 loop { 584 match &mut *driver { 585 PrivateCloseDriver::Pending => { 586 if let Err(error) = self.validate_writer_authority() { 587 #[cfg(test)] 588 self.close_phase.store(6, Ordering::Release); 589 *driver = PrivateCloseDriver::Complete(Some(error.kind())); 590 return Err(error); 591 } 592 let options = sqlite_connect_options(&self.paths, self.mode, self.policy); 593 let connect: BoxFuture<'static, Result<SqliteConnection, ServiceSqliteError>> = 594 Box::pin(async move { 595 SqliteConnection::connect_with(&options) 596 .await 597 .map_err(|source| { 598 connection_source(ServiceSqliteErrorKind::Open, source) 599 }) 600 }); 601 #[cfg(test)] 602 self.close_phase.store(1, Ordering::Release); 603 *driver = PrivateCloseDriver::Connecting(connect); 604 } 605 PrivateCloseDriver::Connecting(connect) => { 606 let connected = connect.as_mut().await; 607 match connected { 608 Ok(connection) => { 609 #[cfg(test)] 610 self.close_phase.store(2, Ordering::Release); 611 *driver = PrivateCloseDriver::Connected(connection); 612 } 613 Err(error) => { 614 let authority_error = self.validate_writer_authority().err(); 615 let error = authority_error.unwrap_or(error); 616 #[cfg(test)] 617 self.close_phase.store(6, Ordering::Release); 618 *driver = PrivateCloseDriver::Complete(Some(error.kind())); 619 return Err(error); 620 } 621 } 622 } 623 PrivateCloseDriver::Connected(connection) => { 624 let mut authority_error = self.validate_writer_authority().err(); 625 let mut checkpoint_error = None; 626 if authority_error.is_none() { 627 #[cfg(test)] 628 self.close_phase.store(3, Ordering::Release); 629 checkpoint_error = 630 verify_connection_policy(connection, self.mode, self.policy) 631 .await 632 .map_err(|source| { 633 connection_source(ServiceSqliteErrorKind::Pragma, source) 634 }) 635 .err(); 636 authority_error = 637 authority_error.or_else(|| self.validate_writer_authority().err()); 638 } 639 if checkpoint_error.is_none() && authority_error.is_none() { 640 #[cfg(test)] 641 self.close_phase 642 .store(TEST_CLOSE_PHASE_CHECKPOINT, Ordering::Release); 643 checkpoint_error = failpoints 644 .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeCheckpoint) 645 .map_err(|source| { 646 ServiceSqliteError::with_source( 647 ServiceSqliteErrorKind::Pragma, 648 source, 649 ) 650 }) 651 .err(); 652 authority_error = 653 authority_error.or_else(|| self.validate_writer_authority().err()); 654 } 655 if checkpoint_error.is_none() && authority_error.is_none() { 656 checkpoint_error = 657 sqlx::query_as::<_, (i64, i64, i64)>("PRAGMA wal_checkpoint(TRUNCATE)") 658 .fetch_one(&mut *connection) 659 .await 660 .map_err(|source| { 661 connection_source(ServiceSqliteErrorKind::Pragma, source) 662 }) 663 .and_then(|(busy, _log_frames, _checkpointed_frames)| { 664 if busy == 0 { 665 Ok(()) 666 } else { 667 Err(connection_error( 668 ServiceSqliteErrorKind::Pragma, 669 ConnectionFailureKind::CheckpointBusy, 670 )) 671 } 672 }) 673 .err(); 674 if checkpoint_error.is_none() { 675 checkpoint_error = failpoints 676 .hit(crate::failpoint::DurabilityFailpoint::CloseAfterCheckpoint) 677 .map_err(|source| { 678 ServiceSqliteError::with_source( 679 ServiceSqliteErrorKind::Pragma, 680 source, 681 ) 682 }) 683 .err(); 684 } 685 authority_error = 686 authority_error.or_else(|| self.validate_writer_authority().err()); 687 } 688 689 let connection_close_error = failpoints 690 .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeConnectionClose) 691 .map_err(|source| { 692 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) 693 }) 694 .err(); 695 696 let connected = core::mem::replace(&mut *driver, PrivateCloseDriver::Pending); 697 let PrivateCloseDriver::Connected(connection) = connected else { 698 unreachable!("close driver retains its connected phase") 699 }; 700 let close: BoxFuture<'static, Result<(), ServiceSqliteError>> = 701 Box::pin(async move { 702 connection.close().await.map_err(|source| { 703 connection_source(ServiceSqliteErrorKind::Open, source) 704 }) 705 }); 706 #[cfg(test)] 707 self.close_phase.store(5, Ordering::Release); 708 *driver = PrivateCloseDriver::Closing { 709 future: close, 710 authority_error, 711 checkpoint_error, 712 connection_close_error, 713 }; 714 } 715 PrivateCloseDriver::Closing { 716 future, 717 authority_error, 718 checkpoint_error, 719 connection_close_error, 720 } => { 721 let close_error = future.as_mut().await.err(); 722 let injected_after_close = failpoints 723 .hit(crate::failpoint::DurabilityFailpoint::CloseAfterConnectionClose) 724 .map_err(|source| { 725 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) 726 }) 727 .err(); 728 let authority_error = authority_error 729 .take() 730 .or_else(|| self.validate_writer_authority().err()); 731 let error = authority_error 732 .or_else(|| checkpoint_error.take()) 733 .or_else(|| connection_close_error.take()) 734 .or(close_error) 735 .or(injected_after_close); 736 #[cfg(test)] 737 self.close_phase.store(6, Ordering::Release); 738 *driver = 739 PrivateCloseDriver::Complete(error.as_ref().map(ServiceSqliteError::kind)); 740 return error.map_or(Ok(()), Err); 741 } 742 PrivateCloseDriver::Complete(kind) => { 743 return kind.map_or(Ok(()), |kind| Err(ServiceSqliteError::new(kind))); 744 } 745 } 746 } 747 } 748 749 fn release_resources( 750 &self, 751 failpoints: &crate::failpoint::DurabilityFailpoints, 752 ) -> Result<Option<ServiceSqliteError>, ServiceSqliteError> { 753 failpoints 754 .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeAuthorityRelease) 755 .map_err(|source| { 756 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Authority, source) 757 })?; 758 let mut resources = self.resources.lock().map_err(|_| { 759 connection_error( 760 ServiceSqliteErrorKind::Authority, 761 ConnectionFailureKind::AuthorityMismatch, 762 ) 763 })?; 764 match self.mode { 765 OpenMode::Initialize | OpenMode::ReadWriteExisting => { 766 let authority = resources.authority.as_mut().ok_or_else(|| { 767 connection_error( 768 ServiceSqliteErrorKind::Authority, 769 ConnectionFailureKind::AuthorityMismatch, 770 ) 771 })?; 772 authority.release()?; 773 resources.authority.take(); 774 } 775 OpenMode::ReadOnlyInspection => { 776 let inspection = resources.inspection_guard.as_mut().ok_or_else(|| { 777 inspection_error(ConnectionFailureKind::InspectionUnavailable) 778 })?; 779 inspection.release()?; 780 resources.inspection_guard.take(); 781 } 782 } 783 Ok(failpoints 784 .hit(crate::failpoint::DurabilityFailpoint::CloseAfterAuthorityRelease) 785 .map_err(|source| { 786 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Authority, source) 787 }) 788 .err()) 789 } 790 } 791 792 #[cfg(any(target_os = "linux", target_os = "macos"))] 793 const fn connection_failure_kind( 794 authority: bool, 795 metadata: bool, 796 migration: bool, 797 integrity: bool, 798 pragma: bool, 799 ) -> ServiceSqliteErrorKind { 800 if authority { 801 ServiceSqliteErrorKind::Authority 802 } else if metadata { 803 ServiceSqliteErrorKind::Metadata 804 } else if migration { 805 ServiceSqliteErrorKind::Migration 806 } else if integrity { 807 ServiceSqliteErrorKind::Integrity 808 } else if pragma { 809 ServiceSqliteErrorKind::Pragma 810 } else { 811 ServiceSqliteErrorKind::Open 812 } 813 } 814 815 #[cfg(any(target_os = "linux", target_os = "macos"))] 816 pub(crate) async fn open_existing_connection_pool( 817 paths: &ServiceSqlitePaths, 818 identity: &ServiceDatabaseIdentity, 819 catalog: &MigrationCatalog, 820 schema_catalog: &SchemaCatalog, 821 mode: OpenMode, 822 policy: ServiceSqliteConnectionOptions, 823 ) -> Result<PrivateConnectionPool, ServiceSqliteError> { 824 open_existing_connection_pool_for( 825 paths, 826 ServiceDatabaseExpectation::Exact(identity), 827 catalog, 828 schema_catalog, 829 mode, 830 policy, 831 ) 832 .await 833 } 834 835 #[cfg(any(target_os = "linux", target_os = "macos"))] 836 pub(crate) async fn open_existing_connection_pool_with_intent( 837 paths: &ServiceSqlitePaths, 838 intent: &ExistingServiceDatabaseIntent, 839 catalog: &MigrationCatalog, 840 schema_catalog: &SchemaCatalog, 841 mode: OpenMode, 842 policy: ServiceSqliteConnectionOptions, 843 ) -> Result<PrivateConnectionPool, ServiceSqliteError> { 844 open_existing_connection_pool_for( 845 paths, 846 ServiceDatabaseExpectation::Existing(intent), 847 catalog, 848 schema_catalog, 849 mode, 850 policy, 851 ) 852 .await 853 } 854 855 #[cfg(any(target_os = "linux", target_os = "macos"))] 856 pub(crate) async fn open_existing_connection_pool_with_intent_and_authority( 857 paths: &ServiceSqlitePaths, 858 intent: &ExistingServiceDatabaseIntent, 859 catalog: &MigrationCatalog, 860 schema_catalog: &SchemaCatalog, 861 policy: ServiceSqliteConnectionOptions, 862 authority: WriterAuthority, 863 ) -> Result<PrivateConnectionPool, ServiceSqliteError> { 864 authority.validate_for(paths)?; 865 open_connection_pool( 866 paths, 867 ServiceDatabaseExpectation::Existing(intent), 868 catalog, 869 schema_catalog, 870 OpenMode::ReadWriteExisting, 871 policy, 872 Some(authority), 873 None, 874 ) 875 .await 876 } 877 878 #[cfg(any(target_os = "linux", target_os = "macos"))] 879 async fn open_existing_connection_pool_for( 880 paths: &ServiceSqlitePaths, 881 expectation: ServiceDatabaseExpectation<'_>, 882 catalog: &MigrationCatalog, 883 schema_catalog: &SchemaCatalog, 884 mode: OpenMode, 885 policy: ServiceSqliteConnectionOptions, 886 ) -> Result<PrivateConnectionPool, ServiceSqliteError> { 887 if mode == OpenMode::Initialize { 888 return Err(connection_error( 889 ServiceSqliteErrorKind::Open, 890 ConnectionFailureKind::UnsupportedMode, 891 )); 892 } 893 let (authority, inspection_guard) = match mode { 894 OpenMode::ReadWriteExisting => (WriterAuthority::acquire(paths, mode)?, None), 895 OpenMode::ReadOnlyInspection => (None, Some(ReadOnlyInspectionGuard::acquire(paths)?)), 896 OpenMode::Initialize => unreachable!("initialize mode returned above"), 897 }; 898 open_connection_pool( 899 paths, 900 expectation, 901 catalog, 902 schema_catalog, 903 mode, 904 policy, 905 authority, 906 inspection_guard, 907 ) 908 .await 909 } 910 911 #[cfg(any(target_os = "linux", target_os = "macos"))] 912 pub(crate) async fn open_initialized_connection_pool( 913 paths: &ServiceSqlitePaths, 914 identity: &ServiceDatabaseIdentity, 915 catalog: &MigrationCatalog, 916 schema_catalog: &SchemaCatalog, 917 policy: ServiceSqliteConnectionOptions, 918 authority: WriterAuthority, 919 ) -> Result<PrivateConnectionPool, ServiceSqliteError> { 920 authority.validate_for(paths)?; 921 open_connection_pool( 922 paths, 923 ServiceDatabaseExpectation::Exact(identity), 924 catalog, 925 schema_catalog, 926 OpenMode::Initialize, 927 policy, 928 Some(authority), 929 None, 930 ) 931 .await 932 } 933 934 #[cfg(any(target_os = "linux", target_os = "macos"))] 935 impl PoolConnectionValidation { 936 async fn validate( 937 &self, 938 connection: &mut SqliteConnection, 939 ) -> Result<(), PoolConnectionValidationFailure> { 940 self.validate_authority()?; 941 let policy_result = connection_policy_matches(connection, self.mode, self.policy).await; 942 self.validate_authority()?; 943 if !policy_result.map_err(PoolConnectionValidationFailure::Pragma)? { 944 return Err(PoolConnectionValidationFailure::PolicyMismatch); 945 } 946 let metadata_result = 947 crate::metadata::verify_database_metadata(connection, &self.identity).await; 948 self.validate_authority()?; 949 metadata_result.map_err(|_| PoolConnectionValidationFailure::Metadata)?; 950 let migration_result = crate::migration::verify_migration_history( 951 connection, 952 &self.catalog, 953 &self.schema_catalog, 954 self.mode == OpenMode::ReadOnlyInspection, 955 ) 956 .await; 957 self.validate_authority()?; 958 migration_result.map_err(|error| { 959 if error.kind() == ServiceSqliteErrorKind::Integrity { 960 PoolConnectionValidationFailure::Integrity 961 } else { 962 PoolConnectionValidationFailure::Migration 963 } 964 })?; 965 Ok(()) 966 } 967 968 fn validate_authority(&self) -> Result<(), PoolConnectionValidationFailure> { 969 self.binding 970 .validate(&self.paths) 971 .map_err(|_| PoolConnectionValidationFailure::Authority) 972 } 973 } 974 975 #[cfg(any(target_os = "linux", target_os = "macos"))] 976 #[allow(clippy::too_many_arguments)] 977 async fn open_connection_pool( 978 paths: &ServiceSqlitePaths, 979 expectation: ServiceDatabaseExpectation<'_>, 980 catalog: &MigrationCatalog, 981 schema_catalog: &SchemaCatalog, 982 mode: OpenMode, 983 policy: ServiceSqliteConnectionOptions, 984 authority: Option<WriterAuthority>, 985 inspection_guard: Option<ReadOnlyInspectionGuard>, 986 ) -> Result<PrivateConnectionPool, ServiceSqliteError> { 987 if !expectation.matches_paths(paths) { 988 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata)); 989 } 990 if expectation.supported_state_schema_version().get() != catalog.current_version() { 991 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Migration)); 992 } 993 if !schema_catalog.matches_migrations(catalog) { 994 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity)); 995 } 996 let recovery_guard = match (mode, authority.as_ref(), inspection_guard.as_ref()) { 997 (OpenMode::Initialize, Some(authority), None) => { 998 authority.validate_for(paths)?; 999 let result = crate::restore::refuse_unresolved_recovery(authority.directory()); 1000 authority.validate_for(paths)?; 1001 result 1002 } 1003 (OpenMode::ReadWriteExisting, Some(authority), None) => match expectation { 1004 ServiceDatabaseExpectation::Exact(identity) => { 1005 crate::restore::recover_for_open(paths, identity, authority) 1006 } 1007 ServiceDatabaseExpectation::Existing(intent) => { 1008 crate::restore::recover_for_open_with_intent(paths, intent, authority) 1009 } 1010 }, 1011 (OpenMode::ReadOnlyInspection, None, Some(inspection_guard)) => { 1012 inspection_guard.validate_for(paths)?; 1013 let result = crate::restore::refuse_unresolved_recovery(&inspection_guard.directory); 1014 inspection_guard.validate_for(paths)?; 1015 result 1016 } 1017 _ => Err(connection_error( 1018 ServiceSqliteErrorKind::Authority, 1019 ConnectionFailureKind::AuthorityMismatch, 1020 )), 1021 }; 1022 recovery_guard?; 1023 let binding = match (mode, authority.as_ref(), inspection_guard.as_ref()) { 1024 (OpenMode::Initialize | OpenMode::ReadWriteExisting, Some(authority), None) => { 1025 authority.validate_for(paths)?; 1026 DirectoryBinding::capture(authority.directory(), paths)? 1027 } 1028 (OpenMode::ReadOnlyInspection, None, Some(inspection_guard)) => { 1029 DirectoryBinding::capture(&inspection_guard.directory, paths)? 1030 } 1031 _ => { 1032 return Err(connection_error( 1033 ServiceSqliteErrorKind::Authority, 1034 ConnectionFailureKind::AuthorityMismatch, 1035 )); 1036 } 1037 }; 1038 1039 let connect_options = sqlite_connect_options(paths, mode, policy); 1040 binding.validate(paths)?; 1041 let preflight_result = SqliteConnection::connect_with(&connect_options).await; 1042 binding.validate(paths)?; 1043 let mut preflight = preflight_result 1044 .map_err(|source| connection_source(ServiceSqliteErrorKind::Open, source))?; 1045 let preflight_policy = verify_connection_policy(&mut preflight, mode, policy).await; 1046 binding.validate(paths)?; 1047 preflight_policy.map_err(|source| connection_source(ServiceSqliteErrorKind::Pragma, source))?; 1048 let preflight_metadata = expectation.verify_metadata(&mut preflight).await; 1049 binding.validate(paths)?; 1050 let preflight_metadata = preflight_metadata?; 1051 let identity = expectation.exact_identity(&preflight_metadata); 1052 let preflight_history = crate::migration::verify_migration_history( 1053 &mut preflight, 1054 catalog, 1055 schema_catalog, 1056 mode == OpenMode::ReadOnlyInspection, 1057 ) 1058 .await; 1059 binding.validate(paths)?; 1060 preflight_history?; 1061 let preflight_close = preflight.close().await; 1062 binding.validate(paths)?; 1063 preflight_close.map_err(|source| connection_source(ServiceSqliteErrorKind::Open, source))?; 1064 1065 let retained_binding = binding.clone(); 1066 let pool_binding = binding; 1067 let retained_catalog = catalog.clone(); 1068 let retained_schema_catalog = schema_catalog.clone(); 1069 let retained_identity = identity.clone(); 1070 let authority_failure = Arc::new(AtomicBool::new(false)); 1071 let metadata_failure = Arc::new(AtomicBool::new(false)); 1072 let migration_failure = Arc::new(AtomicBool::new(false)); 1073 let integrity_failure = Arc::new(AtomicBool::new(false)); 1074 let pragma_failure = Arc::new(AtomicBool::new(false)); 1075 let validation = PoolConnectionValidation { 1076 binding: retained_binding.clone(), 1077 paths: paths.clone(), 1078 identity: identity.clone(), 1079 catalog: catalog.clone(), 1080 schema_catalog: schema_catalog.clone(), 1081 mode, 1082 policy, 1083 }; 1084 let after_validation = validation.clone(); 1085 let before_validation = validation; 1086 let flags = PoolConnectionFailureFlags { 1087 authority: Arc::clone(&authority_failure), 1088 metadata: Arc::clone(&metadata_failure), 1089 migration: Arc::clone(&migration_failure), 1090 integrity: Arc::clone(&integrity_failure), 1091 pragma: Arc::clone(&pragma_failure), 1092 }; 1093 let after_flags = flags.clone(); 1094 let before_flags = flags.clone(); 1095 let pool_result = SqlitePoolOptions::new() 1096 .min_connections(1) 1097 .max_connections(policy.max_connections()) 1098 .acquire_timeout(policy.busy_timeout()) 1099 .idle_timeout(None) 1100 .max_lifetime(None) 1101 .test_before_acquire(true) 1102 .after_connect(move |connection, _metadata| { 1103 let validation = after_validation.clone(); 1104 let flags = after_flags.clone(); 1105 Box::pin(async move { 1106 validation.validate(connection).await.map_err(|failure| { 1107 flags.record(&failure); 1108 failure.into_sqlx() 1109 }) 1110 }) 1111 }) 1112 .before_acquire(move |connection, _metadata| { 1113 let validation = before_validation.clone(); 1114 let flags = before_flags.clone(); 1115 Box::pin(async move { 1116 match validation.validate(connection).await { 1117 Ok(()) => Ok(true), 1118 Err(PoolConnectionValidationFailure::PolicyMismatch) => { 1119 flags.record(&PoolConnectionValidationFailure::PolicyMismatch); 1120 Ok(false) 1121 } 1122 Err(failure) => { 1123 flags.record(&failure); 1124 Err(failure.into_sqlx()) 1125 } 1126 } 1127 }) 1128 }) 1129 .connect_with(connect_options) 1130 .await; 1131 pool_binding.validate(paths)?; 1132 let pool = pool_result.map_err(|source| connection_source(flags.kind(), source))?; 1133 1134 Ok(PrivateConnectionPool { 1135 pool, 1136 binding: retained_binding, 1137 paths: paths.clone(), 1138 identity: retained_identity, 1139 catalog: retained_catalog, 1140 schema_catalog: retained_schema_catalog, 1141 mode, 1142 policy, 1143 resources: Arc::new(Mutex::new(PrivateConnectionResources { 1144 authority, 1145 inspection_guard, 1146 })), 1147 close_driver: tokio::sync::Mutex::new(PrivateCloseDriver::Pending), 1148 #[cfg(test)] 1149 close_phase: std::sync::atomic::AtomicU8::new(0), 1150 authority_failure, 1151 metadata_failure, 1152 migration_failure, 1153 integrity_failure, 1154 pragma_failure, 1155 }) 1156 } 1157 1158 #[cfg(any(target_os = "linux", target_os = "macos"))] 1159 #[allow( 1160 dead_code, 1161 reason = "Step 056 keeps SQLx options private until the Step 061 host boundary" 1162 )] 1163 fn sqlite_connect_options( 1164 paths: &ServiceSqlitePaths, 1165 mode: OpenMode, 1166 policy: ServiceSqliteConnectionOptions, 1167 ) -> SqliteConnectOptions { 1168 let mut options = SqliteConnectOptions::new() 1169 .filename(paths.state_database()) 1170 .read_only(mode == OpenMode::ReadOnlyInspection) 1171 .create_if_missing(false) 1172 .foreign_keys(true) 1173 .busy_timeout(policy.busy_timeout()) 1174 .synchronous(SqliteSynchronous::Full) 1175 .pragma("trusted_schema", "OFF") 1176 .pragma( 1177 "query_only", 1178 if mode == OpenMode::ReadOnlyInspection { 1179 "ON" 1180 } else { 1181 "OFF" 1182 }, 1183 ) 1184 .statement_cache_capacity(STATEMENT_CACHE_CAPACITY) 1185 .command_buffer_size(COMMAND_BUFFER_CAPACITY) 1186 .row_buffer_size(ROW_BUFFER_CAPACITY) 1187 .disable_statement_logging(); 1188 if mode != OpenMode::ReadOnlyInspection { 1189 options = options.journal_mode(SqliteJournalMode::Wal); 1190 } else { 1191 options = options.immutable(true); 1192 } 1193 options 1194 } 1195 1196 #[cfg(any(target_os = "linux", target_os = "macos"))] 1197 #[allow( 1198 dead_code, 1199 reason = "Step 056 keeps pragma verification private until the Step 061 host boundary" 1200 )] 1201 async fn verify_connection_policy( 1202 connection: &mut SqliteConnection, 1203 mode: OpenMode, 1204 policy: ServiceSqliteConnectionOptions, 1205 ) -> Result<(), sqlx::Error> { 1206 if connection_policy_matches(connection, mode, policy).await? { 1207 Ok(()) 1208 } else { 1209 Err(sqlx::Error::Protocol( 1210 "SQLite connection policy mismatch".to_owned(), 1211 )) 1212 } 1213 } 1214 1215 #[cfg(any(target_os = "linux", target_os = "macos"))] 1216 #[allow( 1217 dead_code, 1218 reason = "Step 056 keeps pragma verification private until the Step 061 host boundary" 1219 )] 1220 async fn connection_policy_matches( 1221 connection: &mut SqliteConnection, 1222 mode: OpenMode, 1223 policy: ServiceSqliteConnectionOptions, 1224 ) -> Result<bool, sqlx::Error> { 1225 let journal_mode = sqlx::query_scalar::<_, String>("PRAGMA journal_mode") 1226 .fetch_one(&mut *connection) 1227 .await?; 1228 let synchronous = sqlx::query_scalar::<_, i64>("PRAGMA synchronous") 1229 .fetch_one(&mut *connection) 1230 .await?; 1231 let foreign_keys = sqlx::query_scalar::<_, i64>("PRAGMA foreign_keys") 1232 .fetch_one(&mut *connection) 1233 .await?; 1234 let trusted_schema = sqlx::query_scalar::<_, i64>("PRAGMA trusted_schema") 1235 .fetch_one(&mut *connection) 1236 .await?; 1237 let busy_timeout = sqlx::query_scalar::<_, i64>("PRAGMA busy_timeout") 1238 .fetch_one(&mut *connection) 1239 .await?; 1240 let query_only = sqlx::query_scalar::<_, i64>("PRAGMA query_only") 1241 .fetch_one(&mut *connection) 1242 .await?; 1243 Ok(connection_policy_values_match( 1244 ConnectionPolicyValues { 1245 journal_mode: &journal_mode, 1246 synchronous, 1247 foreign_keys, 1248 trusted_schema, 1249 busy_timeout, 1250 query_only, 1251 }, 1252 mode, 1253 policy, 1254 )) 1255 } 1256 1257 #[cfg(any(target_os = "linux", target_os = "macos"))] 1258 #[derive(Clone, Copy)] 1259 struct ConnectionPolicyValues<'a> { 1260 journal_mode: &'a str, 1261 synchronous: i64, 1262 foreign_keys: i64, 1263 trusted_schema: i64, 1264 busy_timeout: i64, 1265 query_only: i64, 1266 } 1267 1268 #[cfg(any(target_os = "linux", target_os = "macos"))] 1269 fn connection_policy_values_match( 1270 values: ConnectionPolicyValues<'_>, 1271 mode: OpenMode, 1272 policy: ServiceSqliteConnectionOptions, 1273 ) -> bool { 1274 // SQLite reports `delete` for immutable handles; the inspection guard 1275 // independently verifies WAL read/write header bytes before this opens. 1276 let journal_mode_matches = if mode == OpenMode::ReadOnlyInspection { 1277 values.journal_mode.eq_ignore_ascii_case("delete") 1278 } else { 1279 values.journal_mode.eq_ignore_ascii_case("wal") 1280 }; 1281 crate::all_constraints([ 1282 journal_mode_matches, 1283 values.synchronous == 2, 1284 values.foreign_keys == 1, 1285 values.trusted_schema == 0, 1286 values.busy_timeout == policy.busy_timeout_milliseconds(), 1287 values.query_only == i64::from(mode == OpenMode::ReadOnlyInspection), 1288 ]) 1289 } 1290 1291 #[cfg(any(target_os = "linux", target_os = "macos"))] 1292 #[allow( 1293 dead_code, 1294 reason = "Step 056 keeps connection failures private until the Step 061 host boundary" 1295 )] 1296 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 1297 enum ConnectionFailureKind { 1298 UnsupportedMode, 1299 AuthorityMismatch, 1300 InspectionUnavailable, 1301 InspectionContended, 1302 CheckpointBusy, 1303 } 1304 1305 #[cfg(any(target_os = "linux", target_os = "macos"))] 1306 impl fmt::Display for ConnectionFailureKind { 1307 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1308 formatter.write_str(match self { 1309 Self::UnsupportedMode => "SQLite initialize mode requires reserved state", 1310 Self::AuthorityMismatch => "SQLite writer authority is missing or mismatched", 1311 Self::InspectionUnavailable => "SQLite inspection authority is unavailable", 1312 Self::InspectionContended => "SQLite inspection requires an offline writer", 1313 Self::CheckpointBusy => "SQLite close checkpoint could not drain active readers", 1314 }) 1315 } 1316 } 1317 1318 #[cfg(any(target_os = "linux", target_os = "macos"))] 1319 impl Error for ConnectionFailureKind {} 1320 1321 #[cfg(any(target_os = "linux", target_os = "macos"))] 1322 #[allow( 1323 dead_code, 1324 reason = "Step 056 keeps connection failures private until the Step 061 host boundary" 1325 )] 1326 fn connection_error( 1327 kind: ServiceSqliteErrorKind, 1328 cause: ConnectionFailureKind, 1329 ) -> ServiceSqliteError { 1330 ServiceSqliteError::with_source(kind, cause) 1331 } 1332 1333 #[cfg(any(target_os = "linux", target_os = "macos"))] 1334 fn require_connection_condition( 1335 condition: bool, 1336 kind: ServiceSqliteErrorKind, 1337 cause: ConnectionFailureKind, 1338 ) -> Result<(), ServiceSqliteError> { 1339 condition 1340 .then_some(()) 1341 .ok_or_else(|| connection_error(kind, cause)) 1342 } 1343 1344 #[cfg(any(target_os = "linux", target_os = "macos"))] 1345 #[allow( 1346 dead_code, 1347 reason = "Step 056 keeps dependency causes private until the Step 061 host boundary" 1348 )] 1349 fn connection_source(kind: ServiceSqliteErrorKind, cause: sqlx::Error) -> ServiceSqliteError { 1350 ServiceSqliteError::with_source(kind, cause) 1351 } 1352 1353 #[cfg(any(target_os = "linux", target_os = "macos"))] 1354 struct ReadOnlyInspectionGuard { 1355 lock: Option<File>, 1356 lock_device: u64, 1357 lock_inode: u64, 1358 directory: File, 1359 directory_device: u64, 1360 directory_inode: u64, 1361 _database: File, 1362 } 1363 1364 #[cfg(any(target_os = "linux", target_os = "macos"))] 1365 impl ReadOnlyInspectionGuard { 1366 fn acquire(paths: &ServiceSqlitePaths) -> Result<Self, ServiceSqliteError> { 1367 use rustix::{ 1368 fs::{AtFlags, FileType, Mode, OFlags, fstat, open, openat, statat}, 1369 process::geteuid, 1370 }; 1371 1372 let directory = open( 1373 paths 1374 .state_lock() 1375 .parent() 1376 .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))?, 1377 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1378 Mode::empty(), 1379 ) 1380 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1381 let directory_status = fstat(&directory) 1382 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1383 require_connection_condition( 1384 crate::native_metadata::secure_directory( 1385 FileType::from_raw_mode(directory_status.st_mode).is_dir(), 1386 directory_status.st_uid, 1387 geteuid().as_raw(), 1388 crate::native_metadata::mode(directory_status.st_mode), 1389 ), 1390 ServiceSqliteErrorKind::Authority, 1391 ConnectionFailureKind::InspectionUnavailable, 1392 )?; 1393 let directory = File::from(directory); 1394 let lock = openat( 1395 &directory, 1396 radroots_runtime_paths::SERVICE_STATE_LOCK_FILE_NAME, 1397 OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK, 1398 Mode::empty(), 1399 ) 1400 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1401 let lock_status = fstat(&lock) 1402 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1403 require_connection_condition( 1404 crate::native_metadata::exact_regular_file( 1405 FileType::from_raw_mode(lock_status.st_mode).is_file(), 1406 crate::native_metadata::link_count(lock_status.st_nlink), 1407 lock_status.st_uid, 1408 geteuid().as_raw(), 1409 crate::native_metadata::mode(lock_status.st_mode), 1410 ), 1411 ServiceSqliteErrorKind::Authority, 1412 ConnectionFailureKind::InspectionUnavailable, 1413 )?; 1414 let lock = File::from(lock); 1415 FileExt::try_lock_shared(&lock).map_err(|error| { 1416 if error.kind() == std::io::ErrorKind::WouldBlock { 1417 inspection_error(ConnectionFailureKind::InspectionContended) 1418 } else { 1419 inspection_error(ConnectionFailureKind::InspectionUnavailable) 1420 } 1421 })?; 1422 crate::restore::refuse_unresolved_recovery(&directory)?; 1423 for sidecar in [WAL_FILE_NAME, SHARED_MEMORY_FILE_NAME] { 1424 match statat(&directory, sidecar, AtFlags::SYMLINK_NOFOLLOW) { 1425 Err(error) if error == rustix::io::Errno::NOENT => {} 1426 Ok(_) | Err(_) => { 1427 return Err(inspection_error( 1428 ConnectionFailureKind::InspectionUnavailable, 1429 )); 1430 } 1431 } 1432 } 1433 let database = openat( 1434 &directory, 1435 radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME, 1436 OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1437 Mode::empty(), 1438 ) 1439 .map_err(|_| { 1440 connection_error( 1441 ServiceSqliteErrorKind::Open, 1442 ConnectionFailureKind::InspectionUnavailable, 1443 ) 1444 })?; 1445 let database_status = fstat(&database).map_err(|_| { 1446 connection_error( 1447 ServiceSqliteErrorKind::Open, 1448 ConnectionFailureKind::InspectionUnavailable, 1449 ) 1450 })?; 1451 require_connection_condition( 1452 crate::native_metadata::exact_regular_file( 1453 FileType::from_raw_mode(database_status.st_mode).is_file(), 1454 crate::native_metadata::link_count(database_status.st_nlink), 1455 database_status.st_uid, 1456 geteuid().as_raw(), 1457 crate::native_metadata::mode(database_status.st_mode), 1458 ), 1459 ServiceSqliteErrorKind::Open, 1460 ConnectionFailureKind::InspectionUnavailable, 1461 )?; 1462 let database = File::from(database); 1463 let mut sqlite_header = [0_u8; 20]; 1464 std::os::unix::fs::FileExt::read_exact_at(&database, &mut sqlite_header, 0).map_err( 1465 |_| { 1466 connection_error( 1467 ServiceSqliteErrorKind::Open, 1468 ConnectionFailureKind::InspectionUnavailable, 1469 ) 1470 }, 1471 )?; 1472 require_connection_condition( 1473 crate::native_metadata::sqlite_wal_header(&sqlite_header), 1474 ServiceSqliteErrorKind::Pragma, 1475 ConnectionFailureKind::InspectionUnavailable, 1476 )?; 1477 Ok(Self { 1478 lock: Some(lock), 1479 lock_device: crate::native_metadata::device(lock_status.st_dev) 1480 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?, 1481 lock_inode: lock_status.st_ino, 1482 directory, 1483 directory_device: crate::native_metadata::device(directory_status.st_dev) 1484 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?, 1485 directory_inode: directory_status.st_ino, 1486 _database: database, 1487 }) 1488 } 1489 1490 fn validate_for(&self, paths: &ServiceSqlitePaths) -> Result<(), ServiceSqliteError> { 1491 use rustix::{ 1492 fs::{AtFlags, FileType, Mode, OFlags, fstat, open, openat, statat}, 1493 process::geteuid, 1494 }; 1495 1496 let directory_path = paths 1497 .state_lock() 1498 .parent() 1499 .filter(|parent| Some(*parent) == paths.state_database().parent()) 1500 .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1501 let directory = open( 1502 directory_path, 1503 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1504 Mode::empty(), 1505 ) 1506 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1507 let directory_status = fstat(&directory) 1508 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1509 let held_directory_status = fstat(&self.directory) 1510 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1511 let directory_device = crate::native_metadata::device(directory_status.st_dev) 1512 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1513 let held_directory_device = crate::native_metadata::device(held_directory_status.st_dev) 1514 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1515 require_connection_condition( 1516 crate::all_constraints([ 1517 crate::native_metadata::secure_directory( 1518 FileType::from_raw_mode(directory_status.st_mode).is_dir(), 1519 directory_status.st_uid, 1520 geteuid().as_raw(), 1521 crate::native_metadata::mode(directory_status.st_mode), 1522 ), 1523 crate::native_metadata::secure_directory( 1524 FileType::from_raw_mode(held_directory_status.st_mode).is_dir(), 1525 held_directory_status.st_uid, 1526 geteuid().as_raw(), 1527 crate::native_metadata::mode(held_directory_status.st_mode), 1528 ), 1529 crate::native_metadata::identity_pair_matches( 1530 held_directory_device, 1531 held_directory_status.st_ino, 1532 directory_device, 1533 directory_status.st_ino, 1534 self.directory_device, 1535 self.directory_inode, 1536 ), 1537 ]), 1538 ServiceSqliteErrorKind::Authority, 1539 ConnectionFailureKind::InspectionUnavailable, 1540 )?; 1541 1542 let lock = openat( 1543 &directory, 1544 radroots_runtime_paths::SERVICE_STATE_LOCK_FILE_NAME, 1545 OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK, 1546 Mode::empty(), 1547 ) 1548 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1549 let lock_status = fstat(&lock) 1550 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1551 let held_lock = self 1552 .lock 1553 .as_ref() 1554 .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1555 let held_lock_status = fstat(held_lock) 1556 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1557 let lock_device = crate::native_metadata::device(lock_status.st_dev) 1558 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1559 let held_lock_device = crate::native_metadata::device(held_lock_status.st_dev) 1560 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1561 require_connection_condition( 1562 crate::all_constraints([ 1563 crate::native_metadata::exact_regular_file( 1564 FileType::from_raw_mode(lock_status.st_mode).is_file(), 1565 crate::native_metadata::link_count(lock_status.st_nlink), 1566 lock_status.st_uid, 1567 geteuid().as_raw(), 1568 crate::native_metadata::mode(lock_status.st_mode), 1569 ), 1570 crate::native_metadata::exact_regular_file( 1571 FileType::from_raw_mode(held_lock_status.st_mode).is_file(), 1572 crate::native_metadata::link_count(held_lock_status.st_nlink), 1573 held_lock_status.st_uid, 1574 geteuid().as_raw(), 1575 crate::native_metadata::mode(held_lock_status.st_mode), 1576 ), 1577 crate::native_metadata::identity_pair_matches( 1578 held_lock_device, 1579 held_lock_status.st_ino, 1580 lock_device, 1581 lock_status.st_ino, 1582 self.lock_device, 1583 self.lock_inode, 1584 ), 1585 ]), 1586 ServiceSqliteErrorKind::Authority, 1587 ConnectionFailureKind::InspectionUnavailable, 1588 )?; 1589 for sidecar in [WAL_FILE_NAME, SHARED_MEMORY_FILE_NAME] { 1590 match statat(&directory, sidecar, AtFlags::SYMLINK_NOFOLLOW) { 1591 Err(error) if error == rustix::io::Errno::NOENT => {} 1592 Ok(_) | Err(_) => { 1593 return Err(inspection_error( 1594 ConnectionFailureKind::InspectionUnavailable, 1595 )); 1596 } 1597 } 1598 } 1599 Ok(()) 1600 } 1601 1602 fn release(&mut self) -> Result<(), ServiceSqliteError> { 1603 let Some(lock) = self.lock.as_ref() else { 1604 return Ok(()); 1605 }; 1606 FileExt::unlock(lock) 1607 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; 1608 self.lock.take(); 1609 Ok(()) 1610 } 1611 } 1612 1613 #[cfg(any(target_os = "linux", target_os = "macos"))] 1614 impl Drop for ReadOnlyInspectionGuard { 1615 fn drop(&mut self) { 1616 let _ = self.release(); 1617 } 1618 } 1619 1620 #[cfg(any(target_os = "linux", target_os = "macos"))] 1621 fn inspection_error(cause: ConnectionFailureKind) -> ServiceSqliteError { 1622 connection_error(ServiceSqliteErrorKind::Authority, cause) 1623 } 1624 1625 #[cfg(any(target_os = "linux", target_os = "macos"))] 1626 #[derive(Clone)] 1627 struct DirectoryBinding { 1628 database_path: PathBuf, 1629 directory: Arc<File>, 1630 directory_device: u64, 1631 directory_inode: u64, 1632 database: Arc<File>, 1633 database_device: u64, 1634 database_inode: u64, 1635 } 1636 1637 #[cfg(any(target_os = "linux", target_os = "macos"))] 1638 impl DirectoryBinding { 1639 fn capture(directory: &File, paths: &ServiceSqlitePaths) -> Result<Self, ServiceSqliteError> { 1640 use rustix::{ 1641 fs::{FileType, Mode, OFlags, fstat, openat}, 1642 process::geteuid, 1643 }; 1644 1645 let directory_status = fstat(directory).map_err(|_| { 1646 connection_error( 1647 ServiceSqliteErrorKind::Authority, 1648 ConnectionFailureKind::AuthorityMismatch, 1649 ) 1650 })?; 1651 let database = openat( 1652 directory, 1653 radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME, 1654 OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1655 Mode::empty(), 1656 ) 1657 .map_err(|error| { 1658 connection_error( 1659 if error == rustix::io::Errno::NOENT { 1660 ServiceSqliteErrorKind::Open 1661 } else { 1662 ServiceSqliteErrorKind::Authority 1663 }, 1664 ConnectionFailureKind::AuthorityMismatch, 1665 ) 1666 })?; 1667 let database_status = fstat(&database).map_err(|_| { 1668 connection_error( 1669 ServiceSqliteErrorKind::Authority, 1670 ConnectionFailureKind::AuthorityMismatch, 1671 ) 1672 })?; 1673 require_connection_condition( 1674 crate::native_metadata::exact_regular_file( 1675 FileType::from_raw_mode(database_status.st_mode).is_file(), 1676 crate::native_metadata::link_count(database_status.st_nlink), 1677 database_status.st_uid, 1678 geteuid().as_raw(), 1679 crate::native_metadata::mode(database_status.st_mode), 1680 ), 1681 ServiceSqliteErrorKind::Authority, 1682 ConnectionFailureKind::AuthorityMismatch, 1683 )?; 1684 Ok(Self { 1685 database_path: paths.state_database().to_path_buf(), 1686 directory: Arc::new(directory.try_clone().map_err(|_| { 1687 connection_error( 1688 ServiceSqliteErrorKind::Authority, 1689 ConnectionFailureKind::AuthorityMismatch, 1690 ) 1691 })?), 1692 directory_device: crate::native_metadata::device(directory_status.st_dev).map_err( 1693 |_| { 1694 connection_error( 1695 ServiceSqliteErrorKind::Authority, 1696 ConnectionFailureKind::AuthorityMismatch, 1697 ) 1698 }, 1699 )?, 1700 directory_inode: directory_status.st_ino, 1701 database: Arc::new(File::from(database)), 1702 database_device: crate::native_metadata::device(database_status.st_dev).map_err( 1703 |_| { 1704 connection_error( 1705 ServiceSqliteErrorKind::Authority, 1706 ConnectionFailureKind::AuthorityMismatch, 1707 ) 1708 }, 1709 )?, 1710 database_inode: database_status.st_ino, 1711 }) 1712 } 1713 1714 fn validate(&self, paths: &ServiceSqlitePaths) -> Result<(), ServiceSqliteError> { 1715 use rustix::{ 1716 fs::{FileType, Mode, OFlags, fstat, open, openat}, 1717 process::geteuid, 1718 }; 1719 1720 require_connection_condition( 1721 self.database_path == paths.state_database(), 1722 ServiceSqliteErrorKind::Authority, 1723 ConnectionFailureKind::AuthorityMismatch, 1724 )?; 1725 let directory = open( 1726 paths.state_database().parent().ok_or_else(|| { 1727 connection_error( 1728 ServiceSqliteErrorKind::Authority, 1729 ConnectionFailureKind::AuthorityMismatch, 1730 ) 1731 })?, 1732 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1733 Mode::empty(), 1734 ) 1735 .map_err(|_| { 1736 connection_error( 1737 ServiceSqliteErrorKind::Authority, 1738 ConnectionFailureKind::AuthorityMismatch, 1739 ) 1740 })?; 1741 let held_directory_status = fstat(&*self.directory).map_err(|_| { 1742 connection_error( 1743 ServiceSqliteErrorKind::Authority, 1744 ConnectionFailureKind::AuthorityMismatch, 1745 ) 1746 })?; 1747 let directory_status = fstat(&directory).map_err(|_| { 1748 connection_error( 1749 ServiceSqliteErrorKind::Authority, 1750 ConnectionFailureKind::AuthorityMismatch, 1751 ) 1752 })?; 1753 let directory_device = 1754 crate::native_metadata::device(directory_status.st_dev).map_err(|_| { 1755 connection_error( 1756 ServiceSqliteErrorKind::Authority, 1757 ConnectionFailureKind::AuthorityMismatch, 1758 ) 1759 })?; 1760 let held_directory_device = crate::native_metadata::device(held_directory_status.st_dev) 1761 .map_err(|_| { 1762 connection_error( 1763 ServiceSqliteErrorKind::Authority, 1764 ConnectionFailureKind::AuthorityMismatch, 1765 ) 1766 })?; 1767 require_connection_condition( 1768 crate::all_constraints([ 1769 crate::native_metadata::secure_directory( 1770 FileType::from_raw_mode(directory_status.st_mode).is_dir(), 1771 directory_status.st_uid, 1772 geteuid().as_raw(), 1773 crate::native_metadata::mode(directory_status.st_mode), 1774 ), 1775 crate::native_metadata::secure_directory( 1776 FileType::from_raw_mode(held_directory_status.st_mode).is_dir(), 1777 held_directory_status.st_uid, 1778 geteuid().as_raw(), 1779 crate::native_metadata::mode(held_directory_status.st_mode), 1780 ), 1781 crate::native_metadata::identity_pair_matches( 1782 held_directory_device, 1783 held_directory_status.st_ino, 1784 directory_device, 1785 directory_status.st_ino, 1786 self.directory_device, 1787 self.directory_inode, 1788 ), 1789 ]), 1790 ServiceSqliteErrorKind::Authority, 1791 ConnectionFailureKind::AuthorityMismatch, 1792 )?; 1793 1794 let database = openat( 1795 &directory, 1796 radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME, 1797 OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1798 Mode::empty(), 1799 ) 1800 .map_err(|_| { 1801 connection_error( 1802 ServiceSqliteErrorKind::Authority, 1803 ConnectionFailureKind::AuthorityMismatch, 1804 ) 1805 })?; 1806 let held_database_status = fstat(&*self.database).map_err(|_| { 1807 connection_error( 1808 ServiceSqliteErrorKind::Authority, 1809 ConnectionFailureKind::AuthorityMismatch, 1810 ) 1811 })?; 1812 let database_status = fstat(&database).map_err(|_| { 1813 connection_error( 1814 ServiceSqliteErrorKind::Authority, 1815 ConnectionFailureKind::AuthorityMismatch, 1816 ) 1817 })?; 1818 let database_device = 1819 crate::native_metadata::device(database_status.st_dev).map_err(|_| { 1820 connection_error( 1821 ServiceSqliteErrorKind::Authority, 1822 ConnectionFailureKind::AuthorityMismatch, 1823 ) 1824 })?; 1825 let held_database_device = crate::native_metadata::device(held_database_status.st_dev) 1826 .map_err(|_| { 1827 connection_error( 1828 ServiceSqliteErrorKind::Authority, 1829 ConnectionFailureKind::AuthorityMismatch, 1830 ) 1831 })?; 1832 require_connection_condition( 1833 crate::all_constraints([ 1834 crate::native_metadata::exact_regular_file( 1835 FileType::from_raw_mode(database_status.st_mode).is_file(), 1836 crate::native_metadata::link_count(database_status.st_nlink), 1837 database_status.st_uid, 1838 geteuid().as_raw(), 1839 crate::native_metadata::mode(database_status.st_mode), 1840 ), 1841 crate::native_metadata::exact_regular_file( 1842 FileType::from_raw_mode(held_database_status.st_mode).is_file(), 1843 crate::native_metadata::link_count(held_database_status.st_nlink), 1844 held_database_status.st_uid, 1845 geteuid().as_raw(), 1846 crate::native_metadata::mode(held_database_status.st_mode), 1847 ), 1848 crate::native_metadata::identity_pair_matches( 1849 held_database_device, 1850 held_database_status.st_ino, 1851 database_device, 1852 database_status.st_ino, 1853 self.database_device, 1854 self.database_inode, 1855 ), 1856 ]), 1857 ServiceSqliteErrorKind::Authority, 1858 ConnectionFailureKind::AuthorityMismatch, 1859 )?; 1860 Ok(()) 1861 } 1862 } 1863 1864 #[cfg(test)] 1865 mod tests { 1866 #[cfg(any(target_os = "linux", target_os = "macos"))] 1867 use std::path::PathBuf; 1868 1869 #[cfg(any(target_os = "linux", target_os = "macos"))] 1870 use std::{ 1871 collections::BTreeMap, 1872 convert::Infallible, 1873 fs, 1874 num::NonZeroU32, 1875 os::unix::fs::{MetadataExt, PermissionsExt, symlink}, 1876 sync::{ 1877 Arc, 1878 atomic::{AtomicUsize, Ordering as AtomicOrdering}, 1879 }, 1880 time::{Duration, SystemTime}, 1881 }; 1882 1883 use radroots_runtime_paths::{ 1884 RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, RadrootsPlatform, 1885 RuntimeContextBootstrap, RuntimeContextSource, 1886 }; 1887 #[cfg(any(target_os = "linux", target_os = "macos"))] 1888 use radroots_storage::event::SourceGeneration; 1889 #[cfg(any(target_os = "linux", target_os = "macos"))] 1890 use sha2::{Digest, Sha256}; 1891 1892 #[cfg(any(target_os = "linux", target_os = "macos"))] 1893 use crate::{ServiceDatabaseMetadata, ServiceSqliteApplicationId}; 1894 #[cfg(any(target_os = "linux", target_os = "macos"))] 1895 use tokio::sync::Notify; 1896 1897 use super::*; 1898 1899 #[cfg(any(target_os = "linux", target_os = "macos"))] 1900 #[test] 1901 fn connection_failure_inventory_is_complete_and_source_free() { 1902 for (kind, message) in [ 1903 ( 1904 ConnectionFailureKind::UnsupportedMode, 1905 "SQLite initialize mode requires reserved state", 1906 ), 1907 ( 1908 ConnectionFailureKind::AuthorityMismatch, 1909 "SQLite writer authority is missing or mismatched", 1910 ), 1911 ( 1912 ConnectionFailureKind::InspectionUnavailable, 1913 "SQLite inspection authority is unavailable", 1914 ), 1915 ( 1916 ConnectionFailureKind::InspectionContended, 1917 "SQLite inspection requires an offline writer", 1918 ), 1919 ( 1920 ConnectionFailureKind::CheckpointBusy, 1921 "SQLite close checkpoint could not drain active readers", 1922 ), 1923 ] { 1924 assert_eq!(kind.to_string(), message); 1925 assert!(kind.source().is_none()); 1926 assert!(format!("{kind:?}").contains(&format!("{kind:?}"))); 1927 let error = connection_error(ServiceSqliteErrorKind::Open, kind); 1928 assert_eq!(error.kind(), ServiceSqliteErrorKind::Open); 1929 assert!(error.source().is_some()); 1930 assert!(require_connection_condition(true, ServiceSqliteErrorKind::Open, kind).is_ok()); 1931 let rejected = 1932 require_connection_condition(false, ServiceSqliteErrorKind::Authority, kind) 1933 .expect_err("false condition"); 1934 assert_eq!(rejected.kind(), ServiceSqliteErrorKind::Authority); 1935 assert!(rejected.source().is_some()); 1936 } 1937 } 1938 1939 #[cfg(any(target_os = "linux", target_os = "macos"))] 1940 #[test] 1941 fn connection_policy_value_matrix_rejects_each_independent_drift() { 1942 let policy = ServiceSqliteConnectionOptions::reviewed(); 1943 let writable = ConnectionPolicyValues { 1944 journal_mode: "wal", 1945 synchronous: 2, 1946 foreign_keys: 1, 1947 trusted_schema: 0, 1948 busy_timeout: policy.busy_timeout_milliseconds(), 1949 query_only: 0, 1950 }; 1951 assert!(connection_policy_values_match( 1952 writable, 1953 OpenMode::Initialize, 1954 policy, 1955 )); 1956 for values in [ 1957 ConnectionPolicyValues { 1958 journal_mode: "delete", 1959 ..writable 1960 }, 1961 ConnectionPolicyValues { 1962 synchronous: 1, 1963 ..writable 1964 }, 1965 ConnectionPolicyValues { 1966 foreign_keys: 0, 1967 ..writable 1968 }, 1969 ConnectionPolicyValues { 1970 trusted_schema: 1, 1971 ..writable 1972 }, 1973 ConnectionPolicyValues { 1974 busy_timeout: 1, 1975 ..writable 1976 }, 1977 ConnectionPolicyValues { 1978 query_only: 1, 1979 ..writable 1980 }, 1981 ] { 1982 assert!(!connection_policy_values_match( 1983 values, 1984 OpenMode::Initialize, 1985 policy, 1986 )); 1987 } 1988 1989 assert!(connection_policy_values_match( 1990 ConnectionPolicyValues { 1991 journal_mode: "DELETE", 1992 query_only: 1, 1993 ..writable 1994 }, 1995 OpenMode::ReadOnlyInspection, 1996 policy, 1997 )); 1998 assert!(!connection_policy_values_match( 1999 ConnectionPolicyValues { 2000 journal_mode: "wal", 2001 query_only: 1, 2002 ..writable 2003 }, 2004 OpenMode::ReadOnlyInspection, 2005 policy, 2006 )); 2007 } 2008 2009 #[cfg(any(target_os = "linux", target_os = "macos"))] 2010 #[derive(Debug, PartialEq, Eq)] 2011 struct FileSnapshot { 2012 bytes: Vec<u8>, 2013 length: u64, 2014 modified: SystemTime, 2015 mode: u32, 2016 } 2017 2018 #[cfg(any(target_os = "linux", target_os = "macos"))] 2019 fn directory_snapshot(directory: &Path) -> BTreeMap<String, FileSnapshot> { 2020 fs::read_dir(directory) 2021 .expect("read state directory") 2022 .map(|entry| { 2023 let entry = entry.expect("state entry"); 2024 let name = entry.file_name().into_string().expect("UTF-8 state entry"); 2025 let metadata = entry.metadata().expect("state metadata"); 2026 let snapshot = FileSnapshot { 2027 bytes: fs::read(entry.path()).expect("state bytes"), 2028 length: metadata.len(), 2029 modified: metadata.modified().expect("modified time"), 2030 mode: metadata.permissions().mode() & 0o777, 2031 }; 2032 (name, snapshot) 2033 }) 2034 .collect() 2035 } 2036 2037 #[cfg(any(target_os = "linux", target_os = "macos"))] 2038 fn database_metadata(paths: &ServiceSqlitePaths) -> ServiceDatabaseMetadata { 2039 ServiceDatabaseMetadata::new( 2040 paths, 2041 SourceGeneration::new([7; 32]).expect("source generation"), 2042 NonZeroU32::new(1).expect("schema version"), 2043 1_700_000_000_000, 2044 ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), 2045 ) 2046 .expect("database metadata") 2047 } 2048 2049 #[cfg(any(target_os = "linux", target_os = "macos"))] 2050 fn restore_expectation(path: &Path) -> crate::restore::RestoreArtifactExpectation { 2051 let metadata = fs::metadata(path).expect("restore artifact metadata"); 2052 crate::restore::RestoreArtifactExpectation::new( 2053 metadata.dev(), 2054 metadata.ino(), 2055 metadata.len(), 2056 Sha256::digest(fs::read(path).expect("restore artifact bytes")).into(), 2057 ) 2058 .expect("restore artifact expectation") 2059 } 2060 2061 #[cfg(any(target_os = "linux", target_os = "macos"))] 2062 fn base_catalog() -> MigrationCatalog { 2063 MigrationCatalog::new([]).expect("empty v1 catalog") 2064 } 2065 2066 #[cfg(any(target_os = "linux", target_os = "macos"))] 2067 fn schema_catalog( 2068 migrations: &MigrationCatalog, 2069 versions: Vec<Vec<crate::SchemaObject>>, 2070 ) -> crate::SchemaCatalog { 2071 let versions = versions 2072 .into_iter() 2073 .enumerate() 2074 .map(|(index, objects)| { 2075 let version = u32::try_from(index + 1).expect("schema version"); 2076 let digest = 2077 crate::SchemaVersionCatalog::computed_digest(version, objects.iter().cloned()) 2078 .expect("schema digest"); 2079 crate::SchemaVersionCatalog::new(version, objects, digest).expect("schema version") 2080 }) 2081 .collect::<Vec<_>>(); 2082 crate::SchemaCatalog::new(migrations, versions).expect("schema catalog") 2083 } 2084 2085 #[cfg(any(target_os = "linux", target_os = "macos"))] 2086 fn base_schema_catalog() -> crate::SchemaCatalog { 2087 schema_catalog(&base_catalog(), vec![Vec::new()]) 2088 } 2089 2090 #[cfg(any(target_os = "linux", target_os = "macos"))] 2091 fn single_table_schema_catalog(name: &'static str, sql: &'static str) -> crate::SchemaCatalog { 2092 let table = crate::SchemaObject::new( 2093 crate::SchemaObjectKind::Table, 2094 name, 2095 name, 2096 sql, 2097 crate::SchemaObject::computed_digest(crate::SchemaObjectKind::Table, name, name, sql) 2098 .expect("schema table digest"), 2099 ) 2100 .expect("schema table"); 2101 schema_catalog(&base_catalog(), vec![vec![table]]) 2102 } 2103 2104 #[cfg(any(target_os = "linux", target_os = "macos"))] 2105 fn migration_catalog() -> MigrationCatalog { 2106 const CREATE: &str = "CREATE TABLE migration_probe (value INTEGER NOT NULL);"; 2107 const CALLBACK_DEFINITION: &[u8] = b"callback:migration_probe:v1"; 2108 MigrationCatalog::new([ 2109 crate::MigrationDescriptor::sql( 2110 2, 2111 "create_migration_probe", 2112 CREATE, 2113 crate::MigrationChecksum::for_sql(CREATE), 2114 ) 2115 .expect("SQL migration"), 2116 crate::MigrationDescriptor::callback( 2117 3, 2118 "populate_migration_probe", 2119 CALLBACK_DEFINITION, 2120 crate::MigrationChecksum::for_callback(CALLBACK_DEFINITION), 2121 ) 2122 .expect("callback migration"), 2123 ]) 2124 .expect("migration catalog") 2125 } 2126 2127 #[cfg(any(target_os = "linux", target_os = "macos"))] 2128 fn migration_schema_catalog() -> crate::SchemaCatalog { 2129 const SQL: &str = "CREATE TABLE migration_probe (value INTEGER NOT NULL)"; 2130 let table = crate::SchemaObject::new( 2131 crate::SchemaObjectKind::Table, 2132 "migration_probe", 2133 "migration_probe", 2134 SQL, 2135 crate::SchemaObject::computed_digest( 2136 crate::SchemaObjectKind::Table, 2137 "migration_probe", 2138 "migration_probe", 2139 SQL, 2140 ) 2141 .expect("migration schema digest"), 2142 ) 2143 .expect("migration schema table"); 2144 schema_catalog( 2145 &migration_catalog(), 2146 vec![Vec::new(), vec![table.clone()], vec![table]], 2147 ) 2148 } 2149 2150 #[cfg(any(target_os = "linux", target_os = "macos"))] 2151 fn migration_build() -> MigrationBuildIdentity { 2152 MigrationBuildIdentity::new( 2153 "0.1.0-alpha", 2154 "0123456789abcdef0123456789abcdef01234567", 2155 "89abcdef0123456789abcdef0123456789abcdef", 2156 "1.97.1", 2157 "x86_64-unknown-linux-gnu", 2158 "service-host", 2159 1, 2160 2, 2161 3, 2162 4, 2163 5, 2164 ) 2165 .expect("migration build") 2166 } 2167 2168 #[cfg(any(target_os = "linux", target_os = "macos"))] 2169 fn migration_callback_binding() -> crate::migration::MigrationCallbackBinding { 2170 let catalog = migration_catalog(); 2171 let descriptor = &catalog.descriptors()[1]; 2172 crate::migration::MigrationCallbackBinding::new( 2173 descriptor.target_version(), 2174 descriptor.name(), 2175 descriptor.checksum(), 2176 migration_callback, 2177 ) 2178 } 2179 2180 #[cfg(any(target_os = "linux", target_os = "macos"))] 2181 fn migration_callback<'a>( 2182 executor: &'a mut crate::migration::MigrationTransactionExecutor<'_>, 2183 ) -> crate::migration::MigrationCallbackFuture<'a> { 2184 Box::pin(async move { 2185 executor 2186 .execute("INSERT INTO migration_probe (value) VALUES (41)") 2187 .await 2188 }) 2189 } 2190 2191 #[cfg(any(target_os = "linux", target_os = "macos"))] 2192 fn counted_migration_callback<'a>( 2193 executor: &'a mut crate::migration::MigrationTransactionExecutor<'_>, 2194 ) -> crate::migration::MigrationCallbackFuture<'a> { 2195 Box::pin(async move { 2196 CONCURRENT_CALLBACK_COUNT.fetch_add(1, AtomicOrdering::SeqCst); 2197 executor 2198 .execute("INSERT INTO migration_probe (value) VALUES (41)") 2199 .await 2200 }) 2201 } 2202 2203 #[cfg(any(target_os = "linux", target_os = "macos"))] 2204 static CONCURRENT_CALLBACK_COUNT: AtomicUsize = AtomicUsize::new(0); 2205 2206 #[cfg(any(target_os = "linux", target_os = "macos"))] 2207 fn yielding_migration_callback<'a>( 2208 executor: &'a mut crate::migration::MigrationTransactionExecutor<'_>, 2209 ) -> crate::migration::MigrationCallbackFuture<'a> { 2210 Box::pin(async move { 2211 AUTHORITY_CALLBACK_COUNT.fetch_add(1, AtomicOrdering::SeqCst); 2212 tokio::task::yield_now().await; 2213 executor 2214 .execute("INSERT INTO migration_probe (value) VALUES (41)") 2215 .await 2216 }) 2217 } 2218 2219 #[cfg(any(target_os = "linux", target_os = "macos"))] 2220 static AUTHORITY_CALLBACK_COUNT: AtomicUsize = AtomicUsize::new(0); 2221 2222 #[cfg(any(target_os = "linux", target_os = "macos"))] 2223 async fn initialized_authority( 2224 root: &Path, 2225 instance: &str, 2226 ) -> (ServiceSqlitePaths, ServiceDatabaseIdentity, WriterAuthority) { 2227 let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context( 2228 RadrootsPathProfile::RepoLocal, 2229 Some(root.to_path_buf()), 2230 "myc", 2231 instance, 2232 )) 2233 .expect("SQLite paths"); 2234 fs::create_dir_all(paths.state_database().parent().expect("state directory")) 2235 .expect("create state directory"); 2236 let metadata = database_metadata(&paths); 2237 let schema_catalog = base_schema_catalog(); 2238 let authority = crate::initialize_database( 2239 &paths, 2240 OpenMode::Initialize, 2241 &metadata, 2242 &schema_catalog, 2243 |_| Box::pin(async move { Ok::<_, Infallible>(()) }), 2244 ) 2245 .await 2246 .expect("initialize database"); 2247 let identity = metadata.identity(); 2248 (paths, identity, authority) 2249 } 2250 2251 #[cfg(any(target_os = "linux", target_os = "macos"))] 2252 async fn initialized_pool( 2253 root: &Path, 2254 policy: ServiceSqliteConnectionOptions, 2255 ) -> (ServiceSqlitePaths, PrivateConnectionPool) { 2256 let (paths, identity, authority) = initialized_authority(root, "primary").await; 2257 let pool = open_initialized_connection_pool( 2258 &paths, 2259 &identity, 2260 &base_catalog(), 2261 &base_schema_catalog(), 2262 policy, 2263 authority, 2264 ) 2265 .await 2266 .expect("open initialized pool"); 2267 (paths, pool) 2268 } 2269 2270 #[cfg(any(target_os = "linux", target_os = "macos"))] 2271 async fn initialized_migration_pool( 2272 root: &Path, 2273 policy: ServiceSqliteConnectionOptions, 2274 catalog: &MigrationCatalog, 2275 ) -> ( 2276 ServiceSqlitePaths, 2277 ServiceDatabaseIdentity, 2278 PrivateConnectionPool, 2279 ) { 2280 let (paths, base_identity, authority) = initialized_authority(root, "migrations").await; 2281 let identity = ServiceDatabaseIdentity::new( 2282 &paths, 2283 base_identity.source_generation(), 2284 NonZeroU32::new(catalog.current_version()).expect("catalog version"), 2285 base_identity.application_id(), 2286 ); 2287 let schema_catalog = migration_schema_catalog(); 2288 let pool = open_initialized_connection_pool( 2289 &paths, 2290 &identity, 2291 catalog, 2292 &schema_catalog, 2293 policy, 2294 authority, 2295 ) 2296 .await 2297 .expect("open migration pool"); 2298 (paths, identity, pool) 2299 } 2300 2301 fn runtime_context( 2302 profile: RadrootsPathProfile, 2303 repo_local_root: Option<PathBuf>, 2304 service: &str, 2305 instance: &str, 2306 ) -> RuntimeContext { 2307 let profile_source = if matches!(profile, RadrootsPathProfile::RepoLocal) { 2308 RuntimeContextSource::BootstrapCli 2309 } else { 2310 RuntimeContextSource::SafeDefault 2311 }; 2312 RuntimeContext::resolve( 2313 &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), 2314 RuntimeContextBootstrap::new( 2315 profile, 2316 repo_local_root, 2317 profile_source, 2318 RuntimeContextSource::BootstrapCli, 2319 ) 2320 .expect("valid bootstrap"), 2321 ServiceId::new(service).expect("valid service"), 2322 InstanceId::new(instance).expect("valid instance"), 2323 ) 2324 .expect("valid runtime context") 2325 } 2326 2327 #[test] 2328 fn paths_bind_exact_service_host_and_repo_local_artifacts() { 2329 let myc = ServiceSqlitePaths::from_runtime_context(&runtime_context( 2330 RadrootsPathProfile::ServiceHost, 2331 None, 2332 "myc", 2333 "primary", 2334 )) 2335 .expect("Myc paths"); 2336 assert_eq!(myc.service().as_str(), "myc"); 2337 assert_eq!(myc.instance().as_str(), "primary"); 2338 assert_eq!( 2339 myc.state_database(), 2340 Path::new("/var/lib/radroots/services/myc/primary/state.sqlite") 2341 ); 2342 assert_eq!( 2343 myc.state_lock(), 2344 Path::new("/var/lib/radroots/services/myc/primary/state.lock") 2345 ); 2346 2347 let rhi = ServiceSqlitePaths::from_runtime_context(&runtime_context( 2348 RadrootsPathProfile::RepoLocal, 2349 Some(PathBuf::from("/repo/.local/radroots")), 2350 "rhi", 2351 "north-01", 2352 )) 2353 .expect("RHI paths"); 2354 assert_eq!(rhi.service().as_str(), "rhi"); 2355 assert_eq!(rhi.instance().as_str(), "north-01"); 2356 assert_eq!( 2357 rhi.state_database(), 2358 Path::new("/repo/.local/radroots/data/services/rhi/north-01/state.sqlite") 2359 ); 2360 assert_eq!( 2361 rhi.state_lock(), 2362 Path::new("/repo/.local/radroots/data/services/rhi/north-01/state.lock") 2363 ); 2364 2365 let second = ServiceSqlitePaths::from_runtime_context(&runtime_context( 2366 RadrootsPathProfile::RepoLocal, 2367 Some(PathBuf::from("/repo/.local/radroots")), 2368 "rhi", 2369 "south-02", 2370 )) 2371 .expect("second RHI paths"); 2372 assert_ne!(rhi, second); 2373 assert_ne!(rhi.state_database(), second.state_database()); 2374 assert_ne!(rhi.state_lock(), second.state_lock()); 2375 } 2376 2377 #[test] 2378 fn path_shape_failures_are_typed_path_free_and_debug_is_redacted() { 2379 assert_eq!( 2380 validate_state_directory(Path::new("relative/state")), 2381 Err(ServiceSqlitePathError::RelativeStateDirectory) 2382 ); 2383 assert_eq!( 2384 validate_state_directory(Path::new("/")), 2385 Err(ServiceSqlitePathError::MissingStateDirectoryParent) 2386 ); 2387 2388 let error = ServiceSqlitePathError::RelativeStateDirectory; 2389 assert_eq!(error.to_string(), "SQLite state directory must be absolute"); 2390 assert_eq!(format!("{error:?}"), "RelativeStateDirectory"); 2391 2392 let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context( 2393 RadrootsPathProfile::RepoLocal, 2394 Some(PathBuf::from("/sensitive/project-root")), 2395 "myc", 2396 "private-instance", 2397 )) 2398 .expect("redacted paths"); 2399 let debug = format!("{paths:?}"); 2400 assert!(debug.contains("service: ServiceId(\"myc\")")); 2401 assert!(debug.contains("instance: InstanceId(\"private-instance\")")); 2402 assert!(debug.contains("state_database: \"[redacted]\"")); 2403 assert!(debug.contains("state_lock: \"[redacted]\"")); 2404 assert!(!debug.contains("sensitive")); 2405 assert!(!debug.contains("project-root")); 2406 assert!(!debug.contains("state.sqlite")); 2407 assert!(!debug.contains("state.lock")); 2408 } 2409 2410 #[test] 2411 fn open_mode_wire_inventory_and_semantics_are_exact() { 2412 let inventory = [ 2413 (OpenMode::Initialize, "initialize", true, false, true), 2414 ( 2415 OpenMode::ReadWriteExisting, 2416 "read_write_existing", 2417 false, 2418 true, 2419 true, 2420 ), 2421 ( 2422 OpenMode::ReadOnlyInspection, 2423 "read_only_inspection", 2424 false, 2425 true, 2426 false, 2427 ), 2428 ]; 2429 for (mode, wire, may_create, requires_existing, requires_writer) in inventory { 2430 assert_eq!( 2431 serde_json::to_string(&mode).unwrap(), 2432 format!(r#""{wire}""#) 2433 ); 2434 assert_eq!(mode.may_create(), may_create); 2435 assert_eq!(mode.requires_existing(), requires_existing); 2436 assert_eq!(mode.requires_writer_authority(), requires_writer); 2437 } 2438 } 2439 2440 #[cfg(any(target_os = "linux", target_os = "macos"))] 2441 #[tokio::test(flavor = "current_thread")] 2442 async fn every_connection_uses_the_exact_reviewed_pragma_policy() { 2443 let directory = tempfile::tempdir().expect("temporary directory"); 2444 let policy = ServiceSqliteConnectionOptions::reviewed(); 2445 let (_paths, pool) = initialized_pool(directory.path(), policy).await; 2446 2447 let mut connections = Vec::with_capacity(8); 2448 for _ in 0..8 { 2449 connections.push(pool.acquire().await.expect("pooled connection")); 2450 } 2451 assert_eq!(pool.pool.size(), 8); 2452 for connection in &mut connections { 2453 assert!( 2454 connection_policy_matches(connection, OpenMode::Initialize, policy) 2455 .await 2456 .unwrap() 2457 ); 2458 } 2459 drop(connections); 2460 assert!(pool.close().await.expect("writer authority").is_held()); 2461 } 2462 2463 #[cfg(any(target_os = "linux", target_os = "macos"))] 2464 #[tokio::test(flavor = "current_thread")] 2465 async fn reused_connection_drift_is_rejected_before_checkout() { 2466 let directory = tempfile::tempdir().expect("temporary directory"); 2467 let policy = ServiceSqliteConnectionOptions::reviewed(); 2468 let (_paths, pool) = initialized_pool(directory.path(), policy).await; 2469 2470 let mut connection = pool.acquire().await.expect("pooled connection"); 2471 sqlx::query("PRAGMA foreign_keys = OFF") 2472 .execute(&mut *connection) 2473 .await 2474 .expect("drift pragma"); 2475 assert_eq!( 2476 sqlx::query_scalar::<_, i64>("PRAGMA foreign_keys") 2477 .fetch_one(&mut *connection) 2478 .await 2479 .unwrap(), 2480 0 2481 ); 2482 drop(connection); 2483 2484 let mut replacement = pool.acquire().await.expect("replacement connection"); 2485 assert!( 2486 connection_policy_matches(&mut replacement, OpenMode::Initialize, policy) 2487 .await 2488 .unwrap() 2489 ); 2490 drop(replacement); 2491 let _authority = pool.close().await; 2492 } 2493 2494 #[cfg(any(target_os = "linux", target_os = "macos"))] 2495 #[tokio::test(flavor = "current_thread")] 2496 async fn metadata_mismatch_fails_open_and_checkout_before_use() { 2497 let directory = tempfile::tempdir().expect("temporary directory"); 2498 let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 1).unwrap(); 2499 let (paths, pool) = initialized_pool(directory.path(), policy).await; 2500 2501 let mut connection = pool.acquire().await.expect("pooled connection"); 2502 sqlx::query("PRAGMA application_id = 1380209490") 2503 .execute(&mut *connection) 2504 .await 2505 .expect("drift application ID"); 2506 drop(connection); 2507 let error = pool 2508 .acquire() 2509 .await 2510 .expect_err("metadata drift must prevent checkout"); 2511 assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata); 2512 let authority = pool.close().await.expect("writer authority retained"); 2513 drop(authority); 2514 2515 let wrong = ServiceDatabaseIdentity::new( 2516 &paths, 2517 SourceGeneration::new([8; 32]).expect("wrong generation"), 2518 NonZeroU32::new(1).expect("schema version"), 2519 ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), 2520 ); 2521 let result = open_existing_connection_pool( 2522 &paths, 2523 &wrong, 2524 &base_catalog(), 2525 &base_schema_catalog(), 2526 OpenMode::ReadWriteExisting, 2527 policy, 2528 ) 2529 .await; 2530 let Err(error) = result else { 2531 panic!("wrong generation must fail open"); 2532 }; 2533 assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata); 2534 } 2535 2536 #[cfg(any(target_os = "linux", target_os = "macos"))] 2537 #[tokio::test(flavor = "current_thread")] 2538 async fn writable_open_recovers_while_read_only_and_initialize_paths_do_not_mutate() { 2539 let directory = tempfile::tempdir().expect("temporary directory"); 2540 let policy = ServiceSqliteConnectionOptions::reviewed(); 2541 let (paths, identity, mut authority) = 2542 initialized_authority(directory.path(), "restore-recovery").await; 2543 authority 2544 .release() 2545 .expect("release initialization authority"); 2546 2547 let staged = paths 2548 .state_database() 2549 .with_file_name(crate::restore::STAGED_FILE_NAME); 2550 fs::copy(paths.state_database(), &staged).expect("copy exact restore stage"); 2551 fs::set_permissions(&staged, fs::Permissions::from_mode(0o600)) 2552 .expect("restrict staged database"); 2553 let authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 2554 .expect("writer authority") 2555 .expect("writable mode retains authority"); 2556 let marker = crate::restore::RestoreRecoveryMarker::prepared( 2557 &database_metadata(&paths), 2558 crate::BackupManifestSha256::from_bytes([23; 32]), 2559 restore_expectation(paths.state_database()), 2560 restore_expectation(&staged), 2561 ) 2562 .expect("prepared restore marker"); 2563 crate::restore::RestoreMarkerBinding::create(&paths, &authority, &marker) 2564 .expect("persist prepared marker"); 2565 drop(authority); 2566 2567 let state_directory = paths.state_database().parent().expect("state directory"); 2568 let before = directory_snapshot(state_directory); 2569 let read_only = open_existing_connection_pool( 2570 &paths, 2571 &identity, 2572 &base_catalog(), 2573 &base_schema_catalog(), 2574 OpenMode::ReadOnlyInspection, 2575 policy, 2576 ) 2577 .await; 2578 let Err(error) = read_only else { 2579 panic!("read-only open must not recover"); 2580 }; 2581 assert_eq!(error.kind(), ServiceSqliteErrorKind::Recovery); 2582 assert_eq!(directory_snapshot(state_directory), before); 2583 2584 let initialize = crate::initialize_database( 2585 &paths, 2586 OpenMode::Initialize, 2587 &database_metadata(&paths), 2588 &base_schema_catalog(), 2589 |_| Box::pin(async { Ok::<_, Infallible>(()) }), 2590 ) 2591 .await 2592 .expect_err("initialize must not recover existing evidence"); 2593 assert_eq!(initialize.kind(), ServiceSqliteErrorKind::Recovery); 2594 assert_eq!(directory_snapshot(state_directory), before); 2595 2596 let initialized_authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 2597 .expect("writer authority") 2598 .expect("writable mode retains authority"); 2599 let initialized = open_initialized_connection_pool( 2600 &paths, 2601 &identity, 2602 &base_catalog(), 2603 &base_schema_catalog(), 2604 policy, 2605 initialized_authority, 2606 ) 2607 .await; 2608 let Err(error) = initialized else { 2609 panic!("initialized open must not recover"); 2610 }; 2611 assert_eq!(error.kind(), ServiceSqliteErrorKind::Recovery); 2612 assert_eq!(directory_snapshot(state_directory), before); 2613 2614 let writable = open_existing_connection_pool( 2615 &paths, 2616 &identity, 2617 &base_catalog(), 2618 &base_schema_catalog(), 2619 OpenMode::ReadWriteExisting, 2620 policy, 2621 ) 2622 .await 2623 .expect("writable open recovers before SQLite"); 2624 assert!(!staged.exists()); 2625 assert!( 2626 !paths 2627 .state_database() 2628 .with_file_name(crate::restore::MARKER_FILE_NAME) 2629 .exists() 2630 ); 2631 drop(writable.close().await); 2632 } 2633 2634 #[cfg(any(target_os = "linux", target_os = "macos"))] 2635 #[test] 2636 fn connection_failure_precedence_is_exact() { 2637 assert_eq!( 2638 connection_failure_kind(false, false, false, false, false), 2639 ServiceSqliteErrorKind::Open 2640 ); 2641 assert_eq!( 2642 connection_failure_kind(false, false, false, false, true), 2643 ServiceSqliteErrorKind::Pragma 2644 ); 2645 assert_eq!( 2646 connection_failure_kind(false, false, false, true, true), 2647 ServiceSqliteErrorKind::Integrity 2648 ); 2649 assert_eq!( 2650 connection_failure_kind(false, false, true, true, true), 2651 ServiceSqliteErrorKind::Migration 2652 ); 2653 assert_eq!( 2654 connection_failure_kind(false, true, true, true, true), 2655 ServiceSqliteErrorKind::Metadata 2656 ); 2657 assert_eq!( 2658 connection_failure_kind(true, true, true, true, true), 2659 ServiceSqliteErrorKind::Authority 2660 ); 2661 2662 for (failure, expected) in [ 2663 ( 2664 PoolConnectionValidationFailure::Authority, 2665 ServiceSqliteErrorKind::Authority, 2666 ), 2667 ( 2668 PoolConnectionValidationFailure::Metadata, 2669 ServiceSqliteErrorKind::Metadata, 2670 ), 2671 ( 2672 PoolConnectionValidationFailure::Migration, 2673 ServiceSqliteErrorKind::Migration, 2674 ), 2675 ( 2676 PoolConnectionValidationFailure::Integrity, 2677 ServiceSqliteErrorKind::Integrity, 2678 ), 2679 ( 2680 PoolConnectionValidationFailure::PolicyMismatch, 2681 ServiceSqliteErrorKind::Pragma, 2682 ), 2683 ( 2684 PoolConnectionValidationFailure::Pragma(sqlx::Error::Protocol( 2685 "test pragma query failure".to_owned(), 2686 )), 2687 ServiceSqliteErrorKind::Pragma, 2688 ), 2689 ] { 2690 let flags = PoolConnectionFailureFlags { 2691 authority: Arc::new(AtomicBool::new(false)), 2692 metadata: Arc::new(AtomicBool::new(false)), 2693 migration: Arc::new(AtomicBool::new(false)), 2694 integrity: Arc::new(AtomicBool::new(false)), 2695 pragma: Arc::new(AtomicBool::new(false)), 2696 }; 2697 flags.record(&failure); 2698 assert_eq!(flags.kind(), expected); 2699 let _ = failure.into_sqlx(); 2700 } 2701 } 2702 2703 #[cfg(any(target_os = "linux", target_os = "macos"))] 2704 #[tokio::test(flavor = "current_thread")] 2705 async fn schema_drift_is_rejected_on_checkout_and_fresh_open() { 2706 let directory = tempfile::tempdir().expect("temporary directory"); 2707 let policy = ServiceSqliteConnectionOptions::reviewed(); 2708 let (paths, pool) = initialized_pool(directory.path(), policy).await; 2709 let identity = database_metadata(&paths).identity(); 2710 let mut connection = pool.acquire().await.expect("connection"); 2711 sqlx::query("CREATE TABLE unlisted (value INTEGER)") 2712 .execute(&mut *connection) 2713 .await 2714 .expect("create unlisted object"); 2715 drop(connection); 2716 2717 let error = pool 2718 .acquire() 2719 .await 2720 .expect_err("checkout must reject schema drift"); 2721 assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); 2722 let authority = pool.close().await.expect("writer authority"); 2723 drop(authority); 2724 2725 let result = open_existing_connection_pool( 2726 &paths, 2727 &identity, 2728 &base_catalog(), 2729 &base_schema_catalog(), 2730 OpenMode::ReadWriteExisting, 2731 policy, 2732 ) 2733 .await; 2734 let Err(error) = result else { 2735 panic!("fresh open must reject schema drift"); 2736 }; 2737 assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); 2738 } 2739 2740 #[cfg(any(target_os = "linux", target_os = "macos"))] 2741 #[tokio::test(flavor = "current_thread")] 2742 async fn migration_entry_preserves_post_open_ledger_drift_classification() { 2743 let directory = tempfile::tempdir().expect("temporary directory"); 2744 let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 1).unwrap(); 2745 let (_paths, pool) = initialized_pool(directory.path(), policy).await; 2746 let build = migration_build(); 2747 let mut connection = pool.acquire().await.expect("pooled connection"); 2748 sqlx::query( 2749 "INSERT INTO schema_migrations ( 2750 version, name, checksum, applied_at_unix_s, 2751 service_version, service_commit, lib_revision, rust_version, target, 2752 feature_profile, config_contract_version, state_contract_version, 2753 admin_contract_version, status_contract_version, provider_contract_version 2754 ) VALUES (2, 'unexpected_row', zeroblob(32), 0, ?, ?, ?, ?, ?, ?, 1, 2, 3, 4, 5)", 2755 ) 2756 .bind(build.service_version()) 2757 .bind(build.service_commit()) 2758 .bind(build.lib_revision()) 2759 .bind(build.rust_version()) 2760 .bind(build.target()) 2761 .bind(build.feature_profile()) 2762 .execute(&mut *connection) 2763 .await 2764 .expect("inject post-open ledger drift"); 2765 drop(connection); 2766 2767 let error = pool 2768 .apply_migrations( 2769 MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(), 2770 &build, 2771 &[], 2772 ) 2773 .await 2774 .expect_err("migration entry must reject drift before execution"); 2775 assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration); 2776 let authority = pool.close().await.expect("writer authority retained"); 2777 drop(authority); 2778 } 2779 2780 #[cfg(any(target_os = "linux", target_os = "macos"))] 2781 #[tokio::test(flavor = "current_thread")] 2782 async fn pending_history_is_not_exposed_and_read_only_requires_current_state() { 2783 let directory = tempfile::tempdir().expect("temporary directory"); 2784 let catalog = migration_catalog(); 2785 let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 2).unwrap(); 2786 let (paths, identity, pending) = 2787 initialized_migration_pool(directory.path(), policy, &catalog).await; 2788 2789 let error = pending 2790 .acquire() 2791 .await 2792 .expect_err("pending schema must not escape the private pool"); 2793 assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration); 2794 let authority = pending.close().await.expect("writer authority"); 2795 drop(authority); 2796 2797 let read_only = open_existing_connection_pool( 2798 &paths, 2799 &identity, 2800 &catalog, 2801 &migration_schema_catalog(), 2802 OpenMode::ReadOnlyInspection, 2803 policy, 2804 ) 2805 .await; 2806 let Err(error) = read_only else { 2807 panic!("read-only pending state must fail closed"); 2808 }; 2809 assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration); 2810 2811 let writable = open_existing_connection_pool( 2812 &paths, 2813 &identity, 2814 &catalog, 2815 &migration_schema_catalog(), 2816 OpenMode::ReadWriteExisting, 2817 policy, 2818 ) 2819 .await 2820 .expect("reopen writable prefix"); 2821 let outcome = writable 2822 .apply_migrations( 2823 MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(), 2824 &migration_build(), 2825 &[migration_callback_binding()], 2826 ) 2827 .await 2828 .expect("apply pending migrations"); 2829 assert_eq!(outcome.initial_version(), 1); 2830 assert_eq!(outcome.final_version(), 3); 2831 assert_eq!(outcome.applied_count(), 2); 2832 let connection = writable.acquire().await.expect("current checkout"); 2833 drop(connection); 2834 let authority = writable.close().await.expect("writer authority"); 2835 drop(authority); 2836 2837 let current = open_existing_connection_pool( 2838 &paths, 2839 &identity, 2840 &catalog, 2841 &migration_schema_catalog(), 2842 OpenMode::ReadOnlyInspection, 2843 policy, 2844 ) 2845 .await 2846 .expect("current read-only state"); 2847 assert!(current.close().await.is_none()); 2848 } 2849 2850 #[cfg(any(target_os = "linux", target_os = "macos"))] 2851 #[tokio::test(flavor = "current_thread")] 2852 async fn concurrent_migration_attempts_serialize_and_execute_callback_once() { 2853 CONCURRENT_CALLBACK_COUNT.store(0, AtomicOrdering::SeqCst); 2854 let directory = tempfile::tempdir().expect("temporary directory"); 2855 let catalog = migration_catalog(); 2856 let policy = ServiceSqliteConnectionOptions::new(Duration::from_secs(2), 2).unwrap(); 2857 let (_paths, _identity, pool) = 2858 initialized_migration_pool(directory.path(), policy, &catalog).await; 2859 let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(); 2860 let build = migration_build(); 2861 let descriptor = &catalog.descriptors()[1]; 2862 let callbacks = [crate::migration::MigrationCallbackBinding::new( 2863 descriptor.target_version(), 2864 descriptor.name(), 2865 descriptor.checksum(), 2866 counted_migration_callback, 2867 )]; 2868 2869 let (first, second) = tokio::join!( 2870 pool.apply_migrations(applied_at, &build, &callbacks), 2871 pool.apply_migrations(applied_at, &build, &callbacks), 2872 ); 2873 let first = first.expect("first migration attempt"); 2874 let second = second.expect("second migration attempt"); 2875 assert_eq!( 2876 first.applied_count() + second.applied_count(), 2877 2, 2878 "exact catalog is committed only once" 2879 ); 2880 assert_eq!(CONCURRENT_CALLBACK_COUNT.load(AtomicOrdering::SeqCst), 1); 2881 let mut connection = pool.acquire().await.expect("current connection"); 2882 assert_eq!( 2883 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM migration_probe") 2884 .fetch_one(&mut *connection) 2885 .await 2886 .unwrap(), 2887 1 2888 ); 2889 drop(connection); 2890 let _authority = pool.close().await; 2891 } 2892 2893 #[cfg(any(target_os = "linux", target_os = "macos"))] 2894 #[tokio::test(flavor = "current_thread")] 2895 async fn authority_replacement_wins_over_in_flight_migration_failure() { 2896 AUTHORITY_CALLBACK_COUNT.store(0, AtomicOrdering::SeqCst); 2897 let directory = tempfile::tempdir().expect("temporary directory"); 2898 let catalog = migration_catalog(); 2899 let policy = ServiceSqliteConnectionOptions::new(Duration::from_secs(2), 1).unwrap(); 2900 let (paths, _identity, pool) = 2901 initialized_migration_pool(directory.path(), policy, &catalog).await; 2902 let descriptor = &catalog.descriptors()[1]; 2903 let callbacks = [crate::migration::MigrationCallbackBinding::new( 2904 descriptor.target_version(), 2905 descriptor.name(), 2906 descriptor.checksum(), 2907 yielding_migration_callback, 2908 )]; 2909 let state_directory = paths.state_database().parent().unwrap().to_path_buf(); 2910 let displaced = directory.path().join("displaced-migration-state"); 2911 let build = migration_build(); 2912 2913 let application = pool.apply_migrations( 2914 MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(), 2915 &build, 2916 &callbacks, 2917 ); 2918 let replace = async { 2919 while AUTHORITY_CALLBACK_COUNT.load(AtomicOrdering::SeqCst) == 0 { 2920 tokio::task::yield_now().await; 2921 } 2922 fs::rename(&state_directory, &displaced).expect("displace state directory"); 2923 fs::create_dir_all(&state_directory).expect("replace state directory"); 2924 fs::copy(displaced.join("state.sqlite"), paths.state_database()) 2925 .expect("copy replacement database"); 2926 fs::set_permissions(paths.state_database(), fs::Permissions::from_mode(0o600)) 2927 .expect("secure replacement database"); 2928 }; 2929 let (result, ()) = tokio::join!(application, replace); 2930 assert_eq!( 2931 result.expect_err("authority replacement must fail").kind(), 2932 ServiceSqliteErrorKind::Authority 2933 ); 2934 let authority = pool.close().await.expect("retained writer authority"); 2935 assert!(authority.is_held()); 2936 } 2937 2938 #[cfg(any(target_os = "linux", target_os = "macos"))] 2939 #[tokio::test(flavor = "current_thread")] 2940 async fn existing_modes_never_create_missing_state_and_enforce_authority() { 2941 let directory = tempfile::tempdir().expect("temporary directory"); 2942 let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context( 2943 RadrootsPathProfile::RepoLocal, 2944 Some(directory.path().to_path_buf()), 2945 "rhi", 2946 "default", 2947 )) 2948 .expect("SQLite paths"); 2949 fs::create_dir_all(paths.state_database().parent().expect("state directory")) 2950 .expect("create state directory"); 2951 let metadata = database_metadata(&paths).identity(); 2952 2953 for mode in [OpenMode::ReadWriteExisting, OpenMode::ReadOnlyInspection] { 2954 let result = open_existing_connection_pool( 2955 &paths, 2956 &metadata, 2957 &base_catalog(), 2958 &base_schema_catalog(), 2959 mode, 2960 ServiceSqliteConnectionOptions::reviewed(), 2961 ) 2962 .await; 2963 let Err(error) = result else { 2964 panic!("missing database must fail"); 2965 }; 2966 assert_eq!(error.kind(), ServiceSqliteErrorKind::Open); 2967 assert!(!paths.state_database().exists()); 2968 } 2969 let result = open_existing_connection_pool( 2970 &paths, 2971 &metadata, 2972 &base_catalog(), 2973 &base_schema_catalog(), 2974 OpenMode::Initialize, 2975 ServiceSqliteConnectionOptions::reviewed(), 2976 ) 2977 .await; 2978 let Err(error) = result else { 2979 panic!("initialize needs reserved state"); 2980 }; 2981 assert_eq!(error.kind(), ServiceSqliteErrorKind::Open); 2982 } 2983 2984 #[cfg(any(target_os = "linux", target_os = "macos"))] 2985 #[tokio::test(flavor = "current_thread")] 2986 async fn initialized_pool_rejects_mismatched_paths_and_rebound_directory() { 2987 let directory = tempfile::tempdir().expect("temporary directory"); 2988 let (_paths, metadata, authority) = 2989 initialized_authority(directory.path(), "primary").await; 2990 let other = ServiceSqlitePaths::from_runtime_context(&runtime_context( 2991 RadrootsPathProfile::RepoLocal, 2992 Some(directory.path().to_path_buf()), 2993 "myc", 2994 "secondary", 2995 )) 2996 .expect("other SQLite paths"); 2997 fs::create_dir_all( 2998 other 2999 .state_database() 3000 .parent() 3001 .expect("other state directory"), 3002 ) 3003 .expect("create other state directory"); 3004 let result = open_initialized_connection_pool( 3005 &other, 3006 &metadata, 3007 &base_catalog(), 3008 &base_schema_catalog(), 3009 ServiceSqliteConnectionOptions::reviewed(), 3010 authority, 3011 ) 3012 .await; 3013 let Err(error) = result else { 3014 panic!("mismatched paths must fail"); 3015 }; 3016 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 3017 3018 let (paths, metadata, authority) = initialized_authority(directory.path(), "rebound").await; 3019 let state_directory = paths.state_database().parent().expect("state directory"); 3020 let displaced = directory.path().join("displaced-state"); 3021 fs::rename(state_directory, &displaced).expect("displace state directory"); 3022 fs::create_dir_all(state_directory).expect("replace state directory"); 3023 fs::copy(displaced.join("state.sqlite"), paths.state_database()) 3024 .expect("copy replacement database"); 3025 let result = open_initialized_connection_pool( 3026 &paths, 3027 &metadata, 3028 &base_catalog(), 3029 &base_schema_catalog(), 3030 ServiceSqliteConnectionOptions::reviewed(), 3031 authority, 3032 ) 3033 .await; 3034 let Err(error) = result else { 3035 panic!("rebound state directory must fail"); 3036 }; 3037 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 3038 } 3039 3040 #[cfg(any(target_os = "linux", target_os = "macos"))] 3041 #[tokio::test(flavor = "current_thread")] 3042 async fn pool_preflight_rejects_identity_version_and_schema_catalog_drift() { 3043 let directory = tempfile::tempdir().expect("temporary directory"); 3044 3045 let (paths, identity, authority) = 3046 initialized_authority(directory.path(), "identity-drift").await; 3047 let other_paths = ServiceSqlitePaths::from_runtime_context(&runtime_context( 3048 RadrootsPathProfile::RepoLocal, 3049 Some(directory.path().to_path_buf()), 3050 "myc", 3051 "other-identity", 3052 )) 3053 .expect("other paths"); 3054 let wrong_identity = ServiceDatabaseIdentity::new( 3055 &other_paths, 3056 identity.source_generation(), 3057 identity.supported_state_schema_version(), 3058 identity.application_id(), 3059 ); 3060 let Err(error) = open_connection_pool( 3061 &paths, 3062 ServiceDatabaseExpectation::Exact(&wrong_identity), 3063 &base_catalog(), 3064 &base_schema_catalog(), 3065 OpenMode::Initialize, 3066 ServiceSqliteConnectionOptions::reviewed(), 3067 Some(authority), 3068 None, 3069 ) 3070 .await 3071 else { 3072 panic!("identity drift must fail"); 3073 }; 3074 assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata); 3075 3076 let (paths, identity, authority) = 3077 initialized_authority(directory.path(), "version-drift").await; 3078 let newer_identity = ServiceDatabaseIdentity::new( 3079 &paths, 3080 identity.source_generation(), 3081 NonZeroU32::new(2).expect("newer schema"), 3082 identity.application_id(), 3083 ); 3084 let Err(error) = open_connection_pool( 3085 &paths, 3086 ServiceDatabaseExpectation::Exact(&newer_identity), 3087 &base_catalog(), 3088 &base_schema_catalog(), 3089 OpenMode::Initialize, 3090 ServiceSqliteConnectionOptions::reviewed(), 3091 Some(authority), 3092 None, 3093 ) 3094 .await 3095 else { 3096 panic!("version drift must fail"); 3097 }; 3098 assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration); 3099 3100 let (paths, identity, authority) = 3101 initialized_authority(directory.path(), "schema-drift-preflight").await; 3102 let Err(error) = open_connection_pool( 3103 &paths, 3104 ServiceDatabaseExpectation::Exact(&identity), 3105 &base_catalog(), 3106 &migration_schema_catalog(), 3107 OpenMode::Initialize, 3108 ServiceSqliteConnectionOptions::reviewed(), 3109 Some(authority), 3110 None, 3111 ) 3112 .await 3113 else { 3114 panic!("schema catalog drift must fail"); 3115 }; 3116 assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); 3117 } 3118 3119 #[cfg(any(target_os = "linux", target_os = "macos"))] 3120 #[tokio::test(flavor = "current_thread")] 3121 async fn pool_checkout_rejects_state_directory_replacement_before_growth() { 3122 let directory = tempfile::tempdir().expect("temporary directory"); 3123 let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 2).unwrap(); 3124 let (paths, pool) = initialized_pool(directory.path(), policy).await; 3125 assert_eq!(pool.pool.size(), 1); 3126 3127 let state_directory = paths.state_database().parent().expect("state directory"); 3128 let displaced = directory.path().join("displaced-live-state"); 3129 fs::rename(state_directory, &displaced).expect("displace live state directory"); 3130 fs::create_dir_all(state_directory).expect("replace live state directory"); 3131 fs::copy(displaced.join("state.sqlite"), paths.state_database()) 3132 .expect("copy replacement database"); 3133 3134 let error = pool 3135 .acquire() 3136 .await 3137 .expect_err("rebound directory must prevent checkout and pool growth"); 3138 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 3139 assert_eq!(pool.pool.size(), 1); 3140 let authority = pool.close().await.expect("writer authority retained"); 3141 assert!(authority.is_held()); 3142 } 3143 3144 #[cfg(any(target_os = "linux", target_os = "macos"))] 3145 #[tokio::test(flavor = "current_thread")] 3146 async fn writable_existing_rejects_database_symlink_and_hardlink() { 3147 let directory = tempfile::tempdir().expect("temporary directory"); 3148 let policy = ServiceSqliteConnectionOptions::reviewed(); 3149 3150 let (symlink_paths, symlink_metadata, symlink_authority) = 3151 initialized_authority(directory.path(), "symlink-database").await; 3152 drop(symlink_authority); 3153 let symlink_backing = symlink_paths 3154 .state_database() 3155 .parent() 3156 .expect("state directory") 3157 .join("backing.sqlite"); 3158 fs::rename(symlink_paths.state_database(), &symlink_backing) 3159 .expect("displace symlink database"); 3160 symlink(&symlink_backing, symlink_paths.state_database()).expect("database symlink"); 3161 let symlink_result = open_existing_connection_pool( 3162 &symlink_paths, 3163 &symlink_metadata, 3164 &base_catalog(), 3165 &base_schema_catalog(), 3166 OpenMode::ReadWriteExisting, 3167 policy, 3168 ) 3169 .await; 3170 let Err(symlink_error) = symlink_result else { 3171 panic!("database symlink must fail"); 3172 }; 3173 assert_eq!(symlink_error.kind(), ServiceSqliteErrorKind::Authority); 3174 3175 let (hardlink_paths, hardlink_metadata, hardlink_authority) = 3176 initialized_authority(directory.path(), "hardlink-database").await; 3177 drop(hardlink_authority); 3178 let hardlink_alias = hardlink_paths 3179 .state_database() 3180 .parent() 3181 .expect("state directory") 3182 .join("alias.sqlite"); 3183 fs::hard_link(hardlink_paths.state_database(), hardlink_alias).expect("database hard link"); 3184 let hardlink_result = open_existing_connection_pool( 3185 &hardlink_paths, 3186 &hardlink_metadata, 3187 &base_catalog(), 3188 &base_schema_catalog(), 3189 OpenMode::ReadWriteExisting, 3190 policy, 3191 ) 3192 .await; 3193 let Err(hardlink_error) = hardlink_result else { 3194 panic!("database hard link must fail"); 3195 }; 3196 assert_eq!(hardlink_error.kind(), ServiceSqliteErrorKind::Authority); 3197 } 3198 3199 #[cfg(any(target_os = "linux", target_os = "macos"))] 3200 #[tokio::test(flavor = "current_thread")] 3201 async fn lazy_pool_growth_rejects_same_directory_database_replacement() { 3202 let directory = tempfile::tempdir().expect("temporary directory"); 3203 let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 2).unwrap(); 3204 let (paths, pool) = initialized_pool(directory.path(), policy).await; 3205 let held = pool.acquire().await.expect("hold initial connection"); 3206 3207 let displaced = paths 3208 .state_database() 3209 .parent() 3210 .expect("state directory") 3211 .join("displaced.sqlite"); 3212 fs::rename(paths.state_database(), &displaced).expect("displace live database"); 3213 fs::copy(&displaced, paths.state_database()).expect("copy replacement database"); 3214 fs::set_permissions(paths.state_database(), fs::Permissions::from_mode(0o600)) 3215 .expect("secure replacement database"); 3216 3217 let error = pool 3218 .acquire() 3219 .await 3220 .expect_err("database replacement must prevent lazy pool growth"); 3221 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 3222 drop(held); 3223 let authority = pool.close().await.expect("writer authority retained"); 3224 assert!(authority.is_held()); 3225 } 3226 3227 #[cfg(any(target_os = "linux", target_os = "macos"))] 3228 #[tokio::test(flavor = "current_thread")] 3229 async fn read_only_inspection_is_offline_query_only_and_side_effect_free() { 3230 let directory = tempfile::tempdir().expect("temporary directory"); 3231 let policy = ServiceSqliteConnectionOptions::reviewed(); 3232 let (paths, writable) = initialized_pool(directory.path(), policy).await; 3233 let metadata = database_metadata(&paths).identity(); 3234 let mut connection = writable.acquire().await.expect("writable connection"); 3235 sqlx::query("CREATE TABLE inspection_fixture (value INTEGER NOT NULL)") 3236 .execute(&mut *connection) 3237 .await 3238 .expect("create fixture"); 3239 sqlx::query("INSERT INTO inspection_fixture (value) VALUES (41)") 3240 .execute(&mut *connection) 3241 .await 3242 .expect("insert fixture"); 3243 drop(connection); 3244 let inspection_schema_catalog = single_table_schema_catalog( 3245 "inspection_fixture", 3246 "CREATE TABLE inspection_fixture (value INTEGER NOT NULL)", 3247 ); 3248 3249 let contended = open_existing_connection_pool( 3250 &paths, 3251 &metadata, 3252 &base_catalog(), 3253 &inspection_schema_catalog, 3254 OpenMode::ReadOnlyInspection, 3255 policy, 3256 ) 3257 .await; 3258 let Err(error) = contended else { 3259 panic!("inspection must reject an active writer"); 3260 }; 3261 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 3262 3263 let authority = writable.close().await.expect("writer authority"); 3264 drop(authority); 3265 let state_directory = paths.state_database().parent().expect("state directory"); 3266 let stale_wal = state_directory.join(WAL_FILE_NAME); 3267 fs::write(&stale_wal, b"stale-wal-evidence").expect("write stale WAL evidence"); 3268 let stale = open_existing_connection_pool( 3269 &paths, 3270 &metadata, 3271 &base_catalog(), 3272 &inspection_schema_catalog, 3273 OpenMode::ReadOnlyInspection, 3274 policy, 3275 ) 3276 .await; 3277 let Err(error) = stale else { 3278 panic!("inspection must reject stale WAL state"); 3279 }; 3280 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 3281 assert_eq!(fs::read(&stale_wal).unwrap(), b"stale-wal-evidence"); 3282 fs::remove_file(stale_wal).expect("remove test WAL evidence"); 3283 let before = directory_snapshot(state_directory); 3284 3285 let read_only = open_existing_connection_pool( 3286 &paths, 3287 &metadata, 3288 &base_catalog(), 3289 &inspection_schema_catalog, 3290 OpenMode::ReadOnlyInspection, 3291 policy, 3292 ) 3293 .await 3294 .expect("offline read-only inspection"); 3295 let mut connection = read_only.acquire().await.expect("inspection connection"); 3296 assert_eq!( 3297 sqlx::query_scalar::<_, i64>("SELECT value FROM inspection_fixture") 3298 .fetch_one(&mut *connection) 3299 .await 3300 .expect("read fixture"), 3301 41 3302 ); 3303 assert!( 3304 sqlx::query("INSERT INTO inspection_fixture (value) VALUES (42)") 3305 .execute(&mut *connection) 3306 .await 3307 .is_err() 3308 ); 3309 assert!( 3310 connection_policy_matches(&mut connection, OpenMode::ReadOnlyInspection, policy) 3311 .await 3312 .unwrap() 3313 ); 3314 drop(connection); 3315 assert!(read_only.close().await.is_none()); 3316 3317 let after = directory_snapshot(state_directory); 3318 assert_eq!( 3319 after.keys().collect::<Vec<_>>(), 3320 before.keys().collect::<Vec<_>>() 3321 ); 3322 for (name, before_file) in &before { 3323 let after_file = after.get(name).expect("same state entry"); 3324 assert!( 3325 after_file.bytes == before_file.bytes, 3326 "{name} bytes changed" 3327 ); 3328 assert_eq!( 3329 after_file.length, before_file.length, 3330 "{name} length changed" 3331 ); 3332 assert_eq!( 3333 after_file.modified, before_file.modified, 3334 "{name} mtime changed" 3335 ); 3336 assert_eq!(after_file.mode, before_file.mode, "{name} mode changed"); 3337 } 3338 assert!(!after.contains_key("state.sqlite-wal")); 3339 assert!(!after.contains_key("state.sqlite-shm")); 3340 } 3341 3342 #[cfg(any(target_os = "linux", target_os = "macos"))] 3343 async fn offline_inspection_fixture() -> (tempfile::TempDir, ServiceSqlitePaths) { 3344 let directory = tempfile::tempdir().expect("temporary directory"); 3345 let (paths, writable) = 3346 initialized_pool(directory.path(), ServiceSqliteConnectionOptions::reviewed()).await; 3347 let authority = writable.close().await.expect("writer authority"); 3348 drop(authority); 3349 (directory, paths) 3350 } 3351 3352 #[cfg(any(target_os = "linux", target_os = "macos"))] 3353 #[tokio::test(flavor = "current_thread")] 3354 async fn read_only_inspection_guard_rejects_each_filesystem_admission_drift() { 3355 let (_directory, paths) = offline_inspection_fixture().await; 3356 let state_directory = paths.state_database().parent().expect("state directory"); 3357 fs::set_permissions(state_directory, fs::Permissions::from_mode(0o775)) 3358 .expect("make state directory writable by group"); 3359 assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err()); 3360 3361 let (_directory, paths) = offline_inspection_fixture().await; 3362 fs::set_permissions(paths.state_lock(), fs::Permissions::from_mode(0o644)) 3363 .expect("weaken state lock mode"); 3364 assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err()); 3365 3366 let (_directory, paths) = offline_inspection_fixture().await; 3367 let lock_alias = paths 3368 .state_lock() 3369 .parent() 3370 .expect("state directory") 3371 .join("state-lock-alias"); 3372 fs::hard_link(paths.state_lock(), lock_alias).expect("hard-link state lock"); 3373 assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err()); 3374 3375 let (_directory, paths) = offline_inspection_fixture().await; 3376 let authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 3377 .expect("writer authority"); 3378 let Err(error) = ReadOnlyInspectionGuard::acquire(&paths) else { 3379 panic!("held writer lock must prevent inspection"); 3380 }; 3381 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 3382 drop(authority); 3383 3384 let (_directory, paths) = offline_inspection_fixture().await; 3385 fs::set_permissions(paths.state_database(), fs::Permissions::from_mode(0o644)) 3386 .expect("weaken database mode"); 3387 assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err()); 3388 3389 let (_directory, paths) = offline_inspection_fixture().await; 3390 let database_alias = paths 3391 .state_database() 3392 .parent() 3393 .expect("state directory") 3394 .join("state-database-alias"); 3395 fs::hard_link(paths.state_database(), database_alias).expect("hard-link database"); 3396 assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err()); 3397 3398 let (_directory, paths) = offline_inspection_fixture().await; 3399 let database = fs::OpenOptions::new() 3400 .write(true) 3401 .open(paths.state_database()) 3402 .expect("open database header"); 3403 std::os::unix::fs::FileExt::write_all_at(&database, &[1], 18) 3404 .expect("corrupt SQLite write version"); 3405 assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err()); 3406 } 3407 3408 #[cfg(any(target_os = "linux", target_os = "macos"))] 3409 #[tokio::test(flavor = "current_thread")] 3410 async fn read_only_inspection_guard_revalidates_live_directory_and_sidecars() { 3411 let (_directory, paths) = offline_inspection_fixture().await; 3412 let state_directory = paths.state_database().parent().expect("state directory"); 3413 let original_mode = fs::metadata(state_directory) 3414 .expect("state directory metadata") 3415 .permissions() 3416 .mode(); 3417 let guard = ReadOnlyInspectionGuard::acquire(&paths).expect("inspection guard"); 3418 fs::set_permissions(state_directory, fs::Permissions::from_mode(0o775)) 3419 .expect("make live directory unsafe"); 3420 assert!(guard.validate_for(&paths).is_err()); 3421 fs::set_permissions( 3422 state_directory, 3423 fs::Permissions::from_mode(original_mode & 0o777), 3424 ) 3425 .expect("restore directory mode"); 3426 fs::write(state_directory.join(WAL_FILE_NAME), b"stale") 3427 .expect("create stale WAL evidence"); 3428 assert!(guard.validate_for(&paths).is_err()); 3429 } 3430 3431 #[cfg(any(target_os = "linux", target_os = "macos"))] 3432 #[tokio::test(flavor = "current_thread")] 3433 async fn pool_saturation_recovers_and_explicit_close_finishes() { 3434 let directory = tempfile::tempdir().expect("temporary directory"); 3435 let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 1).unwrap(); 3436 let (paths, pool) = initialized_pool(directory.path(), policy).await; 3437 let metadata = database_metadata(&paths).identity(); 3438 let observer = pool.pool.clone(); 3439 let held = pool.acquire().await.expect("only connection"); 3440 let saturated = pool.pool.try_acquire(); 3441 assert!(saturated.is_none()); 3442 drop(held); 3443 let recovered = pool.acquire().await.expect("recovered connection"); 3444 drop(recovered); 3445 3446 let authority = pool.close().await.expect("writer authority retained"); 3447 assert!(observer.is_closed()); 3448 assert!(authority.is_held()); 3449 drop(authority); 3450 3451 let read_only = open_existing_connection_pool( 3452 &paths, 3453 &metadata, 3454 &base_catalog(), 3455 &base_schema_catalog(), 3456 OpenMode::ReadOnlyInspection, 3457 ServiceSqliteConnectionOptions::reviewed(), 3458 ) 3459 .await 3460 .expect("read-only inspection"); 3461 assert_eq!(read_only.mode(), OpenMode::ReadOnlyInspection); 3462 assert!(read_only.close().await.is_none()); 3463 } 3464 3465 #[cfg(any(target_os = "linux", target_os = "macos"))] 3466 #[tokio::test] 3467 async fn cancellation_during_explicit_connection_close_retains_the_close_driver() { 3468 let directory = tempfile::tempdir().expect("temporary directory"); 3469 let policy = ServiceSqliteConnectionOptions::reviewed(); 3470 let (paths, pool) = initialized_pool(directory.path(), policy).await; 3471 let connection = SqliteConnection::connect_with(&sqlite_connect_options( 3472 &paths, 3473 OpenMode::Initialize, 3474 policy, 3475 )) 3476 .await 3477 .expect("open private checkpoint connection"); 3478 let entered = Arc::new(Notify::new()); 3479 let release = Arc::new(Notify::new()); 3480 let close: BoxFuture<'static, Result<(), ServiceSqliteError>> = Box::pin({ 3481 let entered = Arc::clone(&entered); 3482 let release = Arc::clone(&release); 3483 async move { 3484 entered.notify_one(); 3485 release.notified().await; 3486 connection 3487 .close() 3488 .await 3489 .map_err(|source| connection_source(ServiceSqliteErrorKind::Open, source)) 3490 } 3491 }); 3492 { 3493 let mut driver = pool.close_driver.lock().await; 3494 *driver = PrivateCloseDriver::Closing { 3495 future: close, 3496 authority_error: None, 3497 checkpoint_error: None, 3498 connection_close_error: None, 3499 }; 3500 } 3501 let pool = Arc::new(pool); 3502 let close_task = tokio::spawn({ 3503 let pool = Arc::clone(&pool); 3504 async move { 3505 pool.close_explicit(&crate::failpoint::DurabilityFailpoints::default()) 3506 .await 3507 } 3508 }); 3509 entered.notified().await; 3510 close_task.abort(); 3511 assert!( 3512 close_task 3513 .await 3514 .expect_err("explicit close task is cancelled") 3515 .is_cancelled() 3516 ); 3517 let retained = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting); 3518 assert!(matches!( 3519 retained, 3520 Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority 3521 )); 3522 3523 release.notify_one(); 3524 pool.close_explicit(&crate::failpoint::DurabilityFailpoints::default()) 3525 .await 3526 .expect("authority release is proven") 3527 .expect("retained connection closes explicitly"); 3528 let mut reacquired = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 3529 .expect("authority reacquisition") 3530 .expect("writer mode yields authority"); 3531 reacquired.release().expect("release reacquired authority"); 3532 } 3533 }