connection.rs (160648B)
1 //! Narrow service-store host and transaction execution boundary. 2 3 use core::fmt; 4 use std::{ 5 error::Error, 6 future::Future, 7 path::Path, 8 pin::Pin, 9 sync::{ 10 Arc, 11 atomic::{AtomicBool, Ordering}, 12 }, 13 }; 14 15 use futures::{future::BoxFuture, stream::BoxStream}; 16 use sqlx::{ 17 Either, Execute, Executor, SqlStr, Sqlite, SqliteConnection, 18 sqlite::{SqliteQueryResult, SqliteRow, SqliteStatement, SqliteTypeInfo}, 19 }; 20 21 use crate::{ 22 ExistingServiceDatabaseIntent, MigrationApplicationOutcome, MigrationAppliedAtUnixSeconds, 23 MigrationBuildIdentity, MigrationCallbackBinding, MigrationCatalog, OpenMode, SchemaCatalog, 24 ServiceDatabaseIdentity, ServiceDatabaseMetadata, ServiceSqliteConnectionOptions, 25 ServiceSqliteError, ServiceSqliteErrorKind, ServiceSqliteIntegrityReport, ServiceSqlitePaths, 26 WriterAuthority, 27 }; 28 29 #[cfg(any(target_os = "linux", target_os = "macos"))] 30 use sqlx::{Connection, pool::PoolConnection}; 31 32 /// One service-owned SQLite host whose raw pool remains inaccessible. 33 /// 34 /// The host intentionally has no raw-pool accessor: 35 /// 36 /// ```compile_fail 37 /// use radroots_service_sqlite::ServiceSqliteHost; 38 /// 39 /// fn leak_pool(host: &ServiceSqliteHost) { 40 /// let _ = host.pool(); 41 /// } 42 /// ``` 43 pub struct ServiceSqliteHost { 44 mode: OpenMode, 45 #[cfg(any(target_os = "linux", target_os = "macos"))] 46 pool: crate::open::PrivateConnectionPool, 47 #[cfg(any(target_os = "linux", target_os = "macos"))] 48 closing: AtomicBool, 49 #[cfg(any(target_os = "linux", target_os = "macos"))] 50 close_state: tokio::sync::Mutex<ServiceSqliteHostCloseState>, 51 #[cfg(any(target_os = "linux", target_os = "macos"))] 52 backup_active: Arc<AtomicBool>, 53 #[cfg(any(target_os = "linux", target_os = "macos"))] 54 integrity_driver: tokio::sync::Mutex<IntegrityInspectionDriver>, 55 #[cfg(any(target_os = "linux", target_os = "macos"))] 56 failpoints: crate::failpoint::DurabilityFailpoints, 57 } 58 59 /// Database opened under retained authority with its verified metadata. 60 /// 61 /// This result cannot be assembled independently from a host and metadata: 62 /// 63 /// ```compile_fail 64 /// use radroots_service_sqlite::{OpenedServiceDatabase, ServiceSqliteHost}; 65 /// 66 /// fn forge(host: ServiceSqliteHost) { 67 /// let _ = OpenedServiceDatabase { host }; 68 /// } 69 /// ``` 70 pub struct OpenedServiceDatabase { 71 host: ServiceSqliteHost, 72 metadata: ServiceDatabaseMetadata, 73 } 74 75 impl OpenedServiceDatabase { 76 fn new(host: ServiceSqliteHost, metadata: ServiceDatabaseMetadata) -> Self { 77 Self { host, metadata } 78 } 79 80 /// Borrows the authority-retaining host. 81 #[must_use] 82 pub const fn host(&self) -> &ServiceSqliteHost { 83 &self.host 84 } 85 86 /// Borrows the metadata discovered and verified by the retained open. 87 #[must_use] 88 pub const fn database_metadata(&self) -> &ServiceDatabaseMetadata { 89 &self.metadata 90 } 91 92 /// Consumes the binding into the authority-retaining host and actual metadata. 93 #[must_use] 94 pub fn into_parts(self) -> (ServiceSqliteHost, ServiceDatabaseMetadata) { 95 (self.host, self.metadata) 96 } 97 } 98 99 impl fmt::Debug for OpenedServiceDatabase { 100 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 101 formatter 102 .debug_struct("OpenedServiceDatabase") 103 .field("mode", &self.host.mode()) 104 .field("database_metadata", &"[redacted]") 105 .finish() 106 } 107 } 108 109 #[cfg(any(target_os = "linux", target_os = "macos"))] 110 enum ServiceSqliteHostCloseState { 111 Pending, 112 Complete(Option<ServiceSqliteErrorKind>), 113 } 114 115 #[cfg(any(target_os = "linux", target_os = "macos"))] 116 enum IntegrityInspectionDriver { 117 Idle, 118 Connected(QuarantinedConnection), 119 Closing(BoxFuture<'static, Result<(), sqlx::Error>>), 120 } 121 122 #[cfg(any(target_os = "linux", target_os = "macos"))] 123 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 124 enum IntegrityInspectionDriverFailure { 125 Invariant, 126 ConnectionClose, 127 } 128 129 #[cfg(any(target_os = "linux", target_os = "macos"))] 130 fn integrity_driver_close_result( 131 result: Result<(), sqlx::Error>, 132 injected_failure: bool, 133 ) -> Result<(), IntegrityInspectionDriverFailure> { 134 if injected_failure { 135 Err(IntegrityInspectionDriverFailure::ConnectionClose) 136 } else { 137 result.map_err(|_| IntegrityInspectionDriverFailure::ConnectionClose) 138 } 139 } 140 141 #[cfg(any(target_os = "linux", target_os = "macos"))] 142 fn final_connection_policy_matches( 143 initial: &crate::migration::MigrationConnectionPolicy, 144 final_policy: &crate::migration::MigrationConnectionPolicy, 145 ) -> Result<(), ServiceSqliteError> { 146 (final_policy == initial) 147 .then_some(()) 148 .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Pragma)) 149 } 150 151 #[cfg(any(target_os = "linux", target_os = "macos"))] 152 fn unconfirmed_rollback_error( 153 rollback: Option<ServiceSqliteError>, 154 rollback_was_confirmed: bool, 155 ) -> Option<ServiceSqliteError> { 156 rollback.filter(|_| !rollback_was_confirmed) 157 } 158 159 #[cfg(any(target_os = "linux", target_os = "macos"))] 160 fn precondition_rollback_failure( 161 authority: Option<ServiceSqliteError>, 162 rollback: Option<ServiceSqliteError>, 163 rollback_was_confirmed: bool, 164 hook_removal: Option<ServiceSqliteError>, 165 ) -> Option<ServiceSqliteError> { 166 authority 167 .or_else(|| unconfirmed_rollback_error(rollback, rollback_was_confirmed)) 168 .or(hook_removal) 169 } 170 171 #[cfg(any(target_os = "linux", target_os = "macos"))] 172 fn authority_drift_rollback_failure( 173 rollback: Option<ServiceSqliteError>, 174 rollback_was_confirmed: bool, 175 hook_removal: Option<ServiceSqliteError>, 176 ) -> Option<ServiceSqliteError> { 177 unconfirmed_rollback_error(rollback, rollback_was_confirmed).or(hook_removal) 178 } 179 180 #[cfg(any(target_os = "linux", target_os = "macos"))] 181 fn operation_rollback_failure( 182 rollback: Option<ServiceSqliteError>, 183 rollback_was_confirmed: bool, 184 hook_removal: Option<ServiceSqliteError>, 185 authority: Option<ServiceSqliteError>, 186 ) -> Option<ServiceSqliteError> { 187 unconfirmed_rollback_error(rollback, rollback_was_confirmed) 188 .or(hook_removal) 189 .or(authority) 190 } 191 192 #[cfg(any(target_os = "linux", target_os = "macos"))] 193 impl IntegrityInspectionDriver { 194 async fn close_retained(&mut self) -> Result<(), IntegrityInspectionDriverFailure> { 195 loop { 196 match self { 197 Self::Idle => return Ok(()), 198 Self::Connected(_) => { 199 let Self::Connected(connection) = core::mem::replace(self, Self::Idle) else { 200 return Err(IntegrityInspectionDriverFailure::Invariant); 201 }; 202 *self = Self::Closing( 203 connection 204 .into_close_future() 205 .ok_or(IntegrityInspectionDriverFailure::Invariant)?, 206 ); 207 } 208 Self::Closing(close) => { 209 #[cfg(test)] 210 crate::integrity::integrity_test_seam::pause( 211 crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING, 212 ) 213 .await; 214 let result = close.await; 215 *self = Self::Idle; 216 #[cfg(test)] 217 let injected_failure = 218 crate::integrity::integrity_test_seam::take_connection_close_failure(); 219 #[cfg(not(test))] 220 let injected_failure = false; 221 return integrity_driver_close_result(result, injected_failure); 222 } 223 } 224 } 225 } 226 227 fn connection_mut(&mut self) -> Result<&mut SqliteConnection, ServiceSqliteError> { 228 match self { 229 Self::Connected(connection) => Ok(connection), 230 Self::Idle | Self::Closing(_) => { 231 Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity)) 232 } 233 } 234 } 235 236 fn return_to_pool(&mut self) -> Result<(), ServiceSqliteError> { 237 let Self::Connected(mut connection) = core::mem::replace(self, Self::Idle) else { 238 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity)); 239 }; 240 connection.trust(); 241 drop(connection); 242 Ok(()) 243 } 244 245 #[cfg(test)] 246 const fn is_idle(&self) -> bool { 247 matches!(self, Self::Idle) 248 } 249 } 250 251 impl ServiceSqliteHost { 252 /// Opens existing writable state and finishes every pending governed migration. 253 /// 254 /// Before opening SQLite, this path holds exclusive writer authority and 255 /// synchronously reconciles any exact interrupted-restore topology. The 256 /// recovery sequence has no await point: cancellation cannot split one 257 /// filesystem step from its authority check. If the surrounding open is 258 /// cancelled later, a retry re-reads the already durable filesystem state. 259 #[allow(clippy::too_many_arguments)] 260 pub async fn open_read_write_existing( 261 paths: &ServiceSqlitePaths, 262 identity: &ServiceDatabaseIdentity, 263 migrations: &MigrationCatalog, 264 schema: &SchemaCatalog, 265 options: ServiceSqliteConnectionOptions, 266 applied_at: MigrationAppliedAtUnixSeconds, 267 build: &MigrationBuildIdentity, 268 callbacks: &[MigrationCallbackBinding], 269 ) -> Result<(Self, MigrationApplicationOutcome), ServiceSqliteError> { 270 #[cfg(any(target_os = "linux", target_os = "macos"))] 271 { 272 let pool = crate::open::open_existing_connection_pool( 273 paths, 274 identity, 275 migrations, 276 schema, 277 OpenMode::ReadWriteExisting, 278 options, 279 ) 280 .await?; 281 match pool.apply_migrations(applied_at, build, callbacks).await { 282 Ok(outcome) => Ok((Self::from_pool(OpenMode::ReadWriteExisting, pool), outcome)), 283 Err(error) => { 284 drop(pool.close().await); 285 Err(error) 286 } 287 } 288 } 289 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 290 { 291 let _ = ( 292 paths, identity, migrations, schema, options, applied_at, build, callbacks, 293 ); 294 Err(unsupported_host()) 295 } 296 } 297 298 /// Opens existing writable state and discovers its stored generation under authority. 299 /// 300 /// The intent binds service, instance, application ID, and the supported 301 /// schema ceiling before filesystem or SQLite admission. The returned 302 /// metadata is read from the same authority-retaining host after governed 303 /// migrations finish, so callers never guess a source generation or reopen 304 /// the database between discovery and use. 305 #[allow(clippy::too_many_arguments)] 306 pub async fn open_read_write_existing_with_intent( 307 paths: &ServiceSqlitePaths, 308 intent: &ExistingServiceDatabaseIntent, 309 migrations: &MigrationCatalog, 310 schema: &SchemaCatalog, 311 options: ServiceSqliteConnectionOptions, 312 applied_at: MigrationAppliedAtUnixSeconds, 313 build: &MigrationBuildIdentity, 314 callbacks: &[MigrationCallbackBinding], 315 ) -> Result<(OpenedServiceDatabase, MigrationApplicationOutcome), ServiceSqliteError> { 316 #[cfg(any(target_os = "linux", target_os = "macos"))] 317 { 318 let pool = crate::open::open_existing_connection_pool_with_intent( 319 paths, 320 intent, 321 migrations, 322 schema, 323 OpenMode::ReadWriteExisting, 324 options, 325 ) 326 .await?; 327 let outcome = match pool.apply_migrations(applied_at, build, callbacks).await { 328 Ok(outcome) => outcome, 329 Err(error) => { 330 drop(pool.close().await); 331 return Err(error); 332 } 333 }; 334 let metadata = match pool.database_metadata().await { 335 Ok(metadata) => metadata, 336 Err(error) => { 337 drop(pool.close().await); 338 return Err(error); 339 } 340 }; 341 let host = Self::from_pool(OpenMode::ReadWriteExisting, pool); 342 Ok((OpenedServiceDatabase::new(host, metadata), outcome)) 343 } 344 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 345 { 346 let _ = ( 347 paths, intent, migrations, schema, options, applied_at, build, callbacks, 348 ); 349 Err(unsupported_host()) 350 } 351 } 352 353 /// Atomically creates interactive state or opens the exact existing database. 354 /// 355 /// The runner holds writer authority while an exclusive create decides the 356 /// branch. Callers never probe the filesystem or inspect error text. On the 357 /// create branch, the sealed initializer, shared metadata, empty v1 ledger, 358 /// and schema-catalog verification commit together before the host opens and 359 /// applies governed migrations. On exact create collision, the same retained 360 /// authority is transferred to an existing-only open and the initialization 361 /// callback is never invoked. Success binds the retained host to the actual 362 /// verified metadata selected by that atomic decision. 363 #[allow(clippy::too_many_arguments)] 364 pub async fn open_or_initialize<F, E>( 365 paths: &ServiceSqlitePaths, 366 initialization_metadata: &ServiceDatabaseMetadata, 367 migrations: &MigrationCatalog, 368 schema: &SchemaCatalog, 369 options: ServiceSqliteConnectionOptions, 370 applied_at: MigrationAppliedAtUnixSeconds, 371 build: &MigrationBuildIdentity, 372 callbacks: &[MigrationCallbackBinding], 373 initialize_schema: F, 374 ) -> Result<(OpenedServiceDatabase, MigrationApplicationOutcome), ServiceSqliteError> 375 where 376 F: for<'a> FnOnce( 377 &'a mut crate::ServiceSqliteInitializer<'_>, 378 ) -> crate::ServiceSqliteInitializerFuture<'a, E>, 379 E: Error + Send + Sync + 'static, 380 { 381 #[cfg(any(target_os = "linux", target_os = "macos"))] 382 { 383 if !schema.matches_migrations(migrations) { 384 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity)); 385 } 386 let supported_version = core::num::NonZeroU32::new(migrations.current_version()) 387 .expect("migration catalogs always have a nonzero current version"); 388 let initialized = crate::initialize::initialize_or_existing_database( 389 paths, 390 initialization_metadata, 391 schema, 392 initialize_schema, 393 ) 394 .await?; 395 let (mode, pool) = match initialized { 396 crate::initialize::InitializeDatabaseOutcome::Initialized(authority) => { 397 let identity = ServiceDatabaseIdentity::new( 398 paths, 399 initialization_metadata.source_generation(), 400 supported_version, 401 initialization_metadata.application_id(), 402 ); 403 let pool = crate::open::open_initialized_connection_pool( 404 paths, &identity, migrations, schema, options, authority, 405 ) 406 .await?; 407 (OpenMode::Initialize, pool) 408 } 409 crate::initialize::InitializeDatabaseOutcome::Existing(authority) => { 410 let intent = ExistingServiceDatabaseIntent::new( 411 paths, 412 supported_version, 413 initialization_metadata.application_id(), 414 ); 415 let pool = 416 crate::open::open_existing_connection_pool_with_intent_and_authority( 417 paths, &intent, migrations, schema, options, authority, 418 ) 419 .await?; 420 (OpenMode::ReadWriteExisting, pool) 421 } 422 }; 423 let outcome = match pool.apply_migrations(applied_at, build, callbacks).await { 424 Ok(outcome) => outcome, 425 Err(error) => { 426 drop(pool.close().await); 427 return Err(error); 428 } 429 }; 430 let metadata = match pool.database_metadata().await { 431 Ok(metadata) => metadata, 432 Err(error) => { 433 drop(pool.close().await); 434 return Err(error); 435 } 436 }; 437 let host = Self::from_pool(mode, pool); 438 Ok((OpenedServiceDatabase::new(host, metadata), outcome)) 439 } 440 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 441 { 442 drop(( 443 paths, 444 initialization_metadata, 445 migrations, 446 schema, 447 options, 448 applied_at, 449 build, 450 callbacks, 451 initialize_schema, 452 )); 453 Err(unsupported_host()) 454 } 455 } 456 457 /// Opens state created under a retained initialization writer authority. 458 #[allow(clippy::too_many_arguments)] 459 pub async fn open_initialized( 460 paths: &ServiceSqlitePaths, 461 identity: &ServiceDatabaseIdentity, 462 migrations: &MigrationCatalog, 463 schema: &SchemaCatalog, 464 options: ServiceSqliteConnectionOptions, 465 authority: WriterAuthority, 466 applied_at: MigrationAppliedAtUnixSeconds, 467 build: &MigrationBuildIdentity, 468 callbacks: &[MigrationCallbackBinding], 469 ) -> Result<(Self, MigrationApplicationOutcome), ServiceSqliteError> { 470 #[cfg(any(target_os = "linux", target_os = "macos"))] 471 { 472 let pool = crate::open::open_initialized_connection_pool( 473 paths, identity, migrations, schema, options, authority, 474 ) 475 .await?; 476 match pool.apply_migrations(applied_at, build, callbacks).await { 477 Ok(outcome) => Ok((Self::from_pool(OpenMode::Initialize, pool), outcome)), 478 Err(error) => { 479 drop(pool.close().await); 480 Err(error) 481 } 482 } 483 } 484 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 485 { 486 let _ = ( 487 paths, identity, migrations, schema, options, authority, applied_at, build, 488 callbacks, 489 ); 490 Err(unsupported_host()) 491 } 492 } 493 494 /// Opens an immutable, current-schema inspection host without writer authority. 495 pub async fn open_read_only_inspection( 496 paths: &ServiceSqlitePaths, 497 identity: &ServiceDatabaseIdentity, 498 migrations: &MigrationCatalog, 499 schema: &SchemaCatalog, 500 options: ServiceSqliteConnectionOptions, 501 ) -> Result<Self, ServiceSqliteError> { 502 #[cfg(any(target_os = "linux", target_os = "macos"))] 503 { 504 let pool = crate::open::open_existing_connection_pool( 505 paths, 506 identity, 507 migrations, 508 schema, 509 OpenMode::ReadOnlyInspection, 510 options, 511 ) 512 .await?; 513 Ok(Self::from_pool(OpenMode::ReadOnlyInspection, pool)) 514 } 515 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 516 { 517 let _ = (paths, identity, migrations, schema, options); 518 Err(unsupported_host()) 519 } 520 } 521 522 /// Opens existing state for immutable inspection and discovers its metadata. 523 pub async fn open_read_only_inspection_with_intent( 524 paths: &ServiceSqlitePaths, 525 intent: &ExistingServiceDatabaseIntent, 526 migrations: &MigrationCatalog, 527 schema: &SchemaCatalog, 528 options: ServiceSqliteConnectionOptions, 529 ) -> Result<OpenedServiceDatabase, ServiceSqliteError> { 530 #[cfg(any(target_os = "linux", target_os = "macos"))] 531 { 532 let pool = crate::open::open_existing_connection_pool_with_intent( 533 paths, 534 intent, 535 migrations, 536 schema, 537 OpenMode::ReadOnlyInspection, 538 options, 539 ) 540 .await?; 541 let metadata = match pool.database_metadata().await { 542 Ok(metadata) => metadata, 543 Err(error) => { 544 drop(pool.close().await); 545 return Err(error); 546 } 547 }; 548 let host = Self::from_pool(OpenMode::ReadOnlyInspection, pool); 549 Ok(OpenedServiceDatabase::new(host, metadata)) 550 } 551 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 552 { 553 let _ = (paths, intent, migrations, schema, options); 554 Err(unsupported_host()) 555 } 556 } 557 558 /// Returns the fixed mode selected when the host was opened. 559 #[must_use] 560 pub const fn mode(&self) -> OpenMode { 561 self.mode 562 } 563 564 /// Closes all connections and explicitly releases retained instance authority. 565 /// 566 /// Close rejects new transactions as soon as it starts and waits for already 567 /// admitted transactions to finish. Writable hosts then perform the fixed 568 /// governed `TRUNCATE` WAL checkpoint before releasing writer authority; 569 /// read-only inspection performs no checkpoint or filesystem mutation. 570 /// 571 /// Cancelling this future leaves the host permanently non-admitting and retains 572 /// authority until a later call resumes close. A completed result is cached, so 573 /// sequential or concurrent later calls return the same stable outer outcome. 574 pub async fn close(&self) -> Result<(), ServiceSqliteError> { 575 #[cfg(any(target_os = "linux", target_os = "macos"))] 576 { 577 self.closing.store(true, Ordering::Release); 578 let mut state = self.close_state.lock().await; 579 if let ServiceSqliteHostCloseState::Complete(kind) = *state { 580 return kind.map_or(Ok(()), |kind| Err(ServiceSqliteError::new(kind))); 581 } 582 let (integrity_cleanup, integrity_validation) = { 583 let mut driver = self.integrity_driver.lock().await; 584 let cleanup = driver 585 .close_retained() 586 .await 587 .map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Open)); 588 let validation = self.pool.validate(); 589 (cleanup, validation) 590 }; 591 let close = self.pool.close_explicit(&self.failpoints).await; 592 match close { 593 Err(retryable) => Err(retryable), 594 Ok(terminal) => { 595 let terminal = integrity_validation.and(terminal).and(integrity_cleanup); 596 *state = ServiceSqliteHostCloseState::Complete( 597 terminal.as_ref().err().map(ServiceSqliteError::kind), 598 ); 599 terminal 600 } 601 } 602 } 603 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 604 { 605 Err(unsupported_host()) 606 } 607 } 608 609 /// Captures one point-in-time SQLite backup into a new staging directory. 610 /// 611 /// The caller supplies the exact new absolute staging-directory path and an 612 /// injected creation time. Capture is available only on writable hosts and 613 /// admits at most one active capture per host. The canonical manifest and a 614 /// successful result are returned only after the visible staging directory's 615 /// sole `state.sqlite` member has passed metadata, integrity, digest, and 616 /// durability checks. The manifest remains in memory and is not written into 617 /// the staging directory. 618 /// 619 /// Dropping this future requests cancellation. The admitted worker retains 620 /// host authority until it has closed SQLite handles and either completed or 621 /// cleaned the exact staging artifacts, so `close` drains that work before 622 /// releasing writer authority. 623 pub async fn capture_online_backup( 624 &self, 625 staging_directory: &Path, 626 created_at_unix_ms: crate::BackupCreatedAtUnixMs, 627 ) -> Result<crate::ServiceBackupManifest, ServiceSqliteError> { 628 #[cfg(any(target_os = "linux", target_os = "macos"))] 629 { 630 crate::backup::capture_online_backup( 631 &self.pool, 632 &self.closing, 633 &self.backup_active, 634 staging_directory, 635 created_at_unix_ms, 636 &self.failpoints, 637 ) 638 .await 639 } 640 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 641 { 642 let _ = (staging_directory, created_at_unix_ms); 643 Err(unsupported_host()) 644 } 645 } 646 647 /// Runs one explicit bounded integrity inspection over a single read snapshot. 648 /// 649 /// The caller injects the wall-clock completion time and owns any monotonic 650 /// deadline. The host admits at most one inspection at a time. Dropping this 651 /// future before it returns publishes no report, persists no status, and 652 /// leaves the checked-out connection in a host-owned explicit-close driver. 653 /// Retry or host close finishes that close before another check or authority 654 /// release; a retry must inject a new time. 655 /// Completed SQLite and foreign-key failures are returned only as fixed safe 656 /// diagnostic codes. An inability to execute or decode either check is an 657 /// `Integrity` error. 658 pub async fn inspect_integrity( 659 &self, 660 checked_at: crate::IntegrityCheckedAtUnixMs, 661 ) -> Result<ServiceSqliteIntegrityReport, ServiceSqliteError> { 662 #[cfg(any(target_os = "linux", target_os = "macos"))] 663 { 664 self.inspect_integrity_supported(checked_at).await 665 } 666 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 667 { 668 let _ = checked_at; 669 Err(unsupported_host()) 670 } 671 } 672 673 /// Executes one runner-owned transaction without exposing its connection or pool. 674 /// 675 /// Dropping this future before the runner enables its outer commit quarantines 676 /// the connection and leaves no authoritative transaction effect. An operation 677 /// error is returned as `OperationRolledBack` only after rollback is confirmed; 678 /// an unconfirmed rollback is `RollbackFailed`. Once outer commit begins, 679 /// cancelling the future yields no result and must be treated as an unknown 680 /// commit outcome. Callers receiving `CommitOutcomeUnknown`, or cancelling after 681 /// commit begins, must reread authoritative state before any idempotent retry. 682 pub async fn transaction<T, E, F>( 683 &self, 684 operation: F, 685 ) -> Result<T, ServiceSqliteTransactionError<E>> 686 where 687 T: Send + 'static, 688 E: Send + 'static, 689 F: for<'a> FnOnce( 690 &'a mut ServiceSqliteTransaction<'_>, 691 ) -> ServiceSqliteTransactionFuture<'a, T, E> 692 + Send, 693 { 694 #[cfg(any(target_os = "linux", target_os = "macos"))] 695 { 696 self.transaction_supported(operation).await 697 } 698 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 699 { 700 drop(operation); 701 Err(ServiceSqliteTransactionError::not_committed( 702 unsupported_host(), 703 )) 704 } 705 } 706 707 #[cfg(any(target_os = "linux", target_os = "macos"))] 708 async fn transaction_supported<T, E, F>( 709 &self, 710 operation: F, 711 ) -> Result<T, ServiceSqliteTransactionError<E>> 712 where 713 T: Send + 'static, 714 E: Send + 'static, 715 F: for<'a> FnOnce( 716 &'a mut ServiceSqliteTransaction<'_>, 717 ) -> ServiceSqliteTransactionFuture<'a, T, E> 718 + Send, 719 { 720 crate::require_condition( 721 !self.closing.load(Ordering::Acquire), 722 ServiceSqliteErrorKind::Open, 723 ) 724 .map_err(ServiceSqliteTransactionError::not_committed)?; 725 self.pool 726 .validate() 727 .map_err(ServiceSqliteTransactionError::not_committed)?; 728 let connection = self 729 .pool 730 .acquire() 731 .await 732 .map_err(ServiceSqliteTransactionError::not_committed)?; 733 let mut connection = QuarantinedConnection::new(connection); 734 self.pool 735 .validate() 736 .map_err(ServiceSqliteTransactionError::not_committed)?; 737 let initial_policy = crate::migration::read_connection_policy(&mut connection) 738 .await 739 .map_err(ServiceSqliteTransactionError::not_committed)?; 740 self.pool 741 .validate() 742 .map_err(ServiceSqliteTransactionError::not_committed)?; 743 let gate = crate::transaction_control::TransactionControlGate::install(&mut connection) 744 .await 745 .map_err(|source| { 746 ServiceSqliteTransactionError::not_committed(sqlite_source(source)) 747 })?; 748 self.pool 749 .validate() 750 .map_err(ServiceSqliteTransactionError::not_committed)?; 751 let before_begin = self 752 .failpoints 753 .hit(crate::failpoint::DurabilityFailpoint::TransactionBeforeBegin) 754 .map_err(|source| { 755 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) 756 }); 757 self.pool 758 .validate() 759 .map_err(ServiceSqliteTransactionError::not_committed)?; 760 before_begin.map_err(ServiceSqliteTransactionError::not_committed)?; 761 let mut transaction = match match self.pool.mode() { 762 OpenMode::Initialize | OpenMode::ReadWriteExisting => { 763 connection.begin_with("BEGIN IMMEDIATE").await 764 } 765 OpenMode::ReadOnlyInspection => connection.begin().await, 766 } { 767 Ok(transaction) => transaction, 768 Err(source) => { 769 return Err(ServiceSqliteTransactionError::not_committed(sqlite_source( 770 source, 771 ))); 772 } 773 }; 774 let injected_after_begin = self 775 .failpoints 776 .hit(crate::failpoint::DurabilityFailpoint::TransactionAfterBegin) 777 .map_err(|source| { 778 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) 779 }); 780 let after_begin = self.pool.validate().and(injected_after_begin); 781 if let Err(error) = after_begin { 782 let permit = gate.permit_runner_rollback(); 783 let rollback = transaction.rollback().await.map_err(sqlite_source); 784 drop(permit); 785 let rollback_was_confirmed = 786 gate.rejected_commit_rolled_back() && !connection.is_in_transaction(); 787 let remove = gate.remove(&mut connection).await.map_err(sqlite_source); 788 let authority = self.pool.validate(); 789 if let Some(rollback_error) = precondition_rollback_failure( 790 authority.err(), 791 rollback.err(), 792 rollback_was_confirmed, 793 remove.err(), 794 ) { 795 return Err(ServiceSqliteTransactionError::rollback_failed( 796 None, 797 rollback_error, 798 )); 799 } 800 return Err(ServiceSqliteTransactionError::not_committed(error)); 801 } 802 803 let operation_result = { 804 let statement_control_rejected = Arc::new(AtomicBool::new(false)); 805 let mut executor = ServiceSqliteTransaction { 806 connection: &mut transaction, 807 statement_control_rejected: Arc::clone(&statement_control_rejected), 808 }; 809 (operation(&mut executor).await, statement_control_rejected) 810 }; 811 let (operation_result, statement_control_rejected) = operation_result; 812 if let Err(error) = self.pool.validate() { 813 let operation_error = operation_result.err(); 814 let permit = gate.permit_runner_rollback(); 815 let rollback = transaction.rollback().await.map_err(sqlite_source); 816 drop(permit); 817 let rollback_was_confirmed = 818 gate.rejected_commit_rolled_back() && !connection.is_in_transaction(); 819 let remove = gate.remove(&mut connection).await.map_err(sqlite_source); 820 let rollback_error = authority_drift_rollback_failure( 821 rollback.err(), 822 rollback_was_confirmed, 823 remove.err(), 824 ); 825 return Err(match rollback_error { 826 Some(rollback_error) => { 827 ServiceSqliteTransactionError::rollback_failed(operation_error, rollback_error) 828 } 829 None => ServiceSqliteTransactionError::not_committed_with_operation( 830 operation_error, 831 error, 832 ), 833 }); 834 } 835 let value = match operation_result { 836 Ok(value) => value, 837 Err(operation_error) => { 838 let permit = gate.permit_runner_rollback(); 839 let rollback = transaction.rollback().await.map_err(sqlite_source); 840 drop(permit); 841 let rollback_was_confirmed = 842 gate.rejected_commit_rolled_back() && !connection.is_in_transaction(); 843 let remove = gate.remove(&mut connection).await.map_err(sqlite_source); 844 let authority = self.pool.validate(); 845 if let Some(error) = operation_rollback_failure( 846 rollback.err(), 847 rollback_was_confirmed, 848 remove.err(), 849 authority.err(), 850 ) { 851 return Err(ServiceSqliteTransactionError::rollback_failed( 852 Some(operation_error), 853 error, 854 )); 855 } 856 return Err(ServiceSqliteTransactionError::operation_rolled_back( 857 operation_error, 858 )); 859 } 860 }; 861 862 let precommit = self 863 .verify_before_commit( 864 &mut transaction, 865 &gate, 866 &initial_policy, 867 &statement_control_rejected, 868 ) 869 .await; 870 let precommit = match precommit { 871 Ok(()) => { 872 let injected = self 873 .failpoints 874 .hit(crate::failpoint::DurabilityFailpoint::TransactionBeforeCommit) 875 .map_err(|source| { 876 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) 877 }); 878 self.pool.validate().and(injected) 879 } 880 Err(error) => Err(error), 881 }; 882 if let Err(error) = precommit { 883 let permit = gate.permit_runner_rollback(); 884 let rollback = transaction.rollback().await.map_err(sqlite_source); 885 drop(permit); 886 let rollback_was_confirmed = 887 gate.rejected_commit_rolled_back() && !connection.is_in_transaction(); 888 let remove = gate.remove(&mut connection).await.map_err(sqlite_source); 889 let authority = self.pool.validate(); 890 if let Some(rollback_error) = precondition_rollback_failure( 891 authority.err(), 892 rollback.err(), 893 rollback_was_confirmed, 894 remove.err(), 895 ) { 896 return Err(ServiceSqliteTransactionError::rollback_failed( 897 None, 898 rollback_error, 899 )); 900 } 901 return Err(ServiceSqliteTransactionError::not_committed(error)); 902 } 903 904 let permit = gate.permit_outer_commit(); 905 let commit = transaction.commit().await.map_err(sqlite_source); 906 drop(permit); 907 let injected_after_commit = if commit.is_ok() { 908 self.failpoints 909 .hit(crate::failpoint::DurabilityFailpoint::TransactionAfterCommit) 910 .map_err(|source| { 911 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) 912 }) 913 } else { 914 Ok(()) 915 }; 916 let remove = gate.remove(&mut connection).await.map_err(sqlite_source); 917 self.pool 918 .validate() 919 .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 920 commit.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 921 remove.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 922 injected_after_commit.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 923 let final_policy = crate::migration::read_connection_policy(&mut connection) 924 .await 925 .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 926 self.pool 927 .validate() 928 .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 929 final_connection_policy_matches(&initial_policy, &final_policy) 930 .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 931 crate::metadata::verify_database_metadata(&mut connection, self.pool.identity()) 932 .await 933 .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 934 self.pool 935 .validate() 936 .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 937 crate::migration::verify_migration_history( 938 &mut connection, 939 self.pool.catalog(), 940 self.pool.schema_catalog(), 941 true, 942 ) 943 .await 944 .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 945 self.pool 946 .validate() 947 .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; 948 connection.trust(); 949 Ok(value) 950 } 951 952 #[cfg(any(target_os = "linux", target_os = "macos"))] 953 async fn inspect_integrity_supported( 954 &self, 955 checked_at: crate::IntegrityCheckedAtUnixMs, 956 ) -> Result<ServiceSqliteIntegrityReport, ServiceSqliteError> { 957 crate::require_condition( 958 !self.closing.load(Ordering::Acquire), 959 ServiceSqliteErrorKind::Open, 960 )?; 961 let mut driver = self 962 .integrity_driver 963 .try_lock() 964 .map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?; 965 let cleanup = driver.close_retained().await; 966 self.pool.validate()?; 967 cleanup.map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?; 968 crate::require_condition( 969 !self.closing.load(Ordering::Acquire), 970 ServiceSqliteErrorKind::Open, 971 )?; 972 self.pool.validate()?; 973 let connection = self.pool.acquire().await; 974 self.pool.validate()?; 975 *driver = IntegrityInspectionDriver::Connected(QuarantinedConnection::new(connection?)); 976 crate::require_condition( 977 !self.closing.load(Ordering::Acquire), 978 ServiceSqliteErrorKind::Open, 979 )?; 980 let report = crate::integrity::inspect_database_integrity( 981 driver.connection_mut()?, 982 checked_at, 983 || self.pool.validate(), 984 ) 985 .await; 986 let validation = self.pool.validate(); 987 match (validation, report) { 988 (Err(error), _) => { 989 let cleanup = driver.close_retained().await; 990 let validation = self.pool.validate(); 991 if error.kind() == ServiceSqliteErrorKind::Authority { 992 return Err(error); 993 } 994 validation?; 995 cleanup.map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?; 996 Err(error) 997 } 998 (Ok(()), Err(error)) => { 999 let cleanup = driver.close_retained().await; 1000 self.pool.validate()?; 1001 cleanup.map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?; 1002 Err(error) 1003 } 1004 (Ok(()), Ok(report)) => { 1005 driver.return_to_pool()?; 1006 Ok(report) 1007 } 1008 } 1009 } 1010 1011 #[cfg(any(target_os = "linux", target_os = "macos"))] 1012 fn from_pool(mode: OpenMode, pool: crate::open::PrivateConnectionPool) -> Self { 1013 Self { 1014 mode, 1015 pool, 1016 closing: AtomicBool::new(false), 1017 close_state: tokio::sync::Mutex::new(ServiceSqliteHostCloseState::Pending), 1018 backup_active: Arc::new(AtomicBool::new(false)), 1019 integrity_driver: tokio::sync::Mutex::new(IntegrityInspectionDriver::Idle), 1020 failpoints: crate::failpoint::DurabilityFailpoints::default(), 1021 } 1022 } 1023 1024 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 1025 fn arm_durability_failpoint(&self, point: crate::failpoint::DurabilityFailpoint) { 1026 self.failpoints.arm(point); 1027 } 1028 1029 fn lifecycle_state(&self) -> &'static str { 1030 #[cfg(any(target_os = "linux", target_os = "macos"))] 1031 { 1032 if self.closing.load(Ordering::Acquire) { 1033 "closing_or_closed" 1034 } else { 1035 "open" 1036 } 1037 } 1038 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 1039 { 1040 "closing_or_closed" 1041 } 1042 } 1043 1044 #[cfg(any(target_os = "linux", target_os = "macos"))] 1045 async fn verify_before_commit( 1046 &self, 1047 connection: &mut SqliteConnection, 1048 gate: &crate::transaction_control::TransactionControlGate, 1049 initial_policy: &crate::migration::MigrationConnectionPolicy, 1050 statement_control_rejected: &AtomicBool, 1051 ) -> Result<(), ServiceSqliteError> { 1052 self.pool.validate()?; 1053 crate::require_condition( 1054 !gate.control_violation_observed() 1055 && !statement_control_rejected.load(Ordering::Acquire), 1056 ServiceSqliteErrorKind::Open, 1057 )?; 1058 crate::migration::assert_governed_transaction(connection).await?; 1059 self.pool.validate()?; 1060 crate::require_condition( 1061 &crate::migration::read_connection_policy(connection).await? == initial_policy, 1062 ServiceSqliteErrorKind::Pragma, 1063 )?; 1064 self.pool.validate()?; 1065 crate::metadata::verify_database_metadata(connection, self.pool.identity()).await?; 1066 self.pool.validate()?; 1067 crate::migration::verify_migration_history_snapshot( 1068 connection, 1069 self.pool.catalog(), 1070 self.pool.schema_catalog(), 1071 true, 1072 ) 1073 .await?; 1074 self.pool.validate()?; 1075 crate::migration::assert_governed_transaction(connection).await?; 1076 crate::require_condition( 1077 !gate.control_violation_observed(), 1078 ServiceSqliteErrorKind::Open, 1079 )?; 1080 Ok(()) 1081 } 1082 } 1083 1084 impl fmt::Debug for ServiceSqliteHost { 1085 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1086 formatter 1087 .debug_struct("ServiceSqliteHost") 1088 .field("mode", &self.mode) 1089 .field("pool", &"[redacted]") 1090 .field("lifecycle", &self.lifecycle_state()) 1091 .finish() 1092 } 1093 } 1094 1095 /// A sealed transaction executor that never exposes its raw SQLite connection. 1096 /// 1097 /// Service repositories may use ordinary typed SQLx queries through the 1098 /// borrowed executor: 1099 /// 1100 /// ``` 1101 /// use radroots_service_sqlite::ServiceSqliteTransaction; 1102 /// 1103 /// async fn row_count( 1104 /// transaction: &mut ServiceSqliteTransaction<'_>, 1105 /// ) -> Result<i64, sqlx::Error> { 1106 /// sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM service_items") 1107 /// .fetch_one(transaction) 1108 /// .await 1109 /// } 1110 /// ``` 1111 /// 1112 /// Transaction control remains runner-owned: 1113 /// 1114 /// ```compile_fail 1115 /// use radroots_service_sqlite::ServiceSqliteTransaction; 1116 /// 1117 /// async fn bypass(transaction: ServiceSqliteTransaction<'_>) { 1118 /// transaction.commit().await.unwrap(); 1119 /// } 1120 /// ``` 1121 pub struct ServiceSqliteTransaction<'connection> { 1122 connection: &'connection mut SqliteConnection, 1123 statement_control_rejected: Arc<AtomicBool>, 1124 } 1125 1126 struct RestrictedExecute<Q> { 1127 query: Q, 1128 statement_control_rejected: Arc<AtomicBool>, 1129 } 1130 1131 impl<'query, Q> Execute<'query, Sqlite> for RestrictedExecute<Q> 1132 where 1133 Q: Execute<'query, Sqlite>, 1134 { 1135 fn sql(self) -> SqlStr { 1136 restricted_sql(self.query.sql(), &self.statement_control_rejected) 1137 } 1138 1139 fn statement(&self) -> Option<&SqliteStatement> { 1140 None 1141 } 1142 1143 fn take_arguments( 1144 &mut self, 1145 ) -> Result<Option<<Sqlite as sqlx::Database>::Arguments>, sqlx::error::BoxDynError> { 1146 self.query.take_arguments() 1147 } 1148 1149 fn persistent(&self) -> bool { 1150 self.query.persistent() 1151 } 1152 } 1153 1154 fn restricted_sql(sql: SqlStr, statement_control_rejected: &AtomicBool) -> SqlStr { 1155 if crate::statement_policy::contains_forbidden_statement_control(sql.as_str()) { 1156 statement_control_rejected.store(true, Ordering::Release); 1157 SqlStr::from_static("RADROOTS_FORBIDDEN_STATEMENT_CONTROL") 1158 } else { 1159 sql 1160 } 1161 } 1162 1163 impl fmt::Debug for ServiceSqliteTransaction<'_> { 1164 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1165 formatter.write_str("ServiceSqliteTransaction([redacted])") 1166 } 1167 } 1168 1169 impl<'executor, 'connection> Executor<'executor> 1170 for &'executor mut ServiceSqliteTransaction<'connection> 1171 where 1172 'connection: 'executor, 1173 { 1174 type Database = Sqlite; 1175 1176 fn fetch_many<'e, 'q: 'e, Q>( 1177 self, 1178 query: Q, 1179 ) -> BoxStream<'e, Result<Either<SqliteQueryResult, SqliteRow>, sqlx::Error>> 1180 where 1181 'executor: 'e, 1182 Q: 'q + Execute<'q, Self::Database>, 1183 { 1184 (&mut *self.connection).fetch_many(RestrictedExecute { 1185 query, 1186 statement_control_rejected: Arc::clone(&self.statement_control_rejected), 1187 }) 1188 } 1189 1190 fn fetch_optional<'e, 'q: 'e, Q>( 1191 self, 1192 query: Q, 1193 ) -> BoxFuture<'e, Result<Option<SqliteRow>, sqlx::Error>> 1194 where 1195 'executor: 'e, 1196 Q: 'q + Execute<'q, Self::Database>, 1197 { 1198 (&mut *self.connection).fetch_optional(RestrictedExecute { 1199 query, 1200 statement_control_rejected: Arc::clone(&self.statement_control_rejected), 1201 }) 1202 } 1203 1204 fn prepare_with<'e>( 1205 self, 1206 sql: SqlStr, 1207 parameters: &'e [SqliteTypeInfo], 1208 ) -> BoxFuture<'e, Result<SqliteStatement, sqlx::Error>> 1209 where 1210 'executor: 'e, 1211 { 1212 (&mut *self.connection).prepare_with( 1213 restricted_sql(sql, &self.statement_control_rejected), 1214 parameters, 1215 ) 1216 } 1217 } 1218 1219 /// Boxed callback future tied to the borrowed transaction executor. 1220 pub type ServiceSqliteTransactionFuture<'a, T, E> = 1221 Pin<Box<dyn Future<Output = Result<T, E>> + Send + 'a>>; 1222 1223 /// Stable transaction completion phases without expanding SQLite error kinds. 1224 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 1225 pub enum ServiceSqliteTransactionErrorKind { 1226 NotCommitted, 1227 OperationRolledBack, 1228 RollbackFailed, 1229 CommitOutcomeUnknown, 1230 } 1231 1232 /// Transaction failure retaining trusted details behind redacted diagnostics. 1233 pub struct ServiceSqliteTransactionError<E> { 1234 kind: ServiceSqliteTransactionErrorKind, 1235 operation_error: Option<E>, 1236 sqlite_error: Option<ServiceSqliteError>, 1237 } 1238 1239 impl<E> ServiceSqliteTransactionError<E> { 1240 fn not_committed(error: ServiceSqliteError) -> Self { 1241 Self::not_committed_with_operation(None, error) 1242 } 1243 1244 fn not_committed_with_operation( 1245 operation_error: Option<E>, 1246 sqlite_error: ServiceSqliteError, 1247 ) -> Self { 1248 Self { 1249 kind: ServiceSqliteTransactionErrorKind::NotCommitted, 1250 operation_error, 1251 sqlite_error: Some(sqlite_error), 1252 } 1253 } 1254 1255 #[cfg(any(target_os = "linux", target_os = "macos"))] 1256 fn operation_rolled_back(error: E) -> Self { 1257 Self { 1258 kind: ServiceSqliteTransactionErrorKind::OperationRolledBack, 1259 operation_error: Some(error), 1260 sqlite_error: None, 1261 } 1262 } 1263 1264 #[cfg(any(target_os = "linux", target_os = "macos"))] 1265 fn rollback_failed(operation_error: Option<E>, sqlite_error: ServiceSqliteError) -> Self { 1266 Self { 1267 kind: ServiceSqliteTransactionErrorKind::RollbackFailed, 1268 operation_error, 1269 sqlite_error: Some(sqlite_error), 1270 } 1271 } 1272 1273 #[cfg(any(target_os = "linux", target_os = "macos"))] 1274 fn commit_outcome_unknown(error: ServiceSqliteError) -> Self { 1275 Self { 1276 kind: ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown, 1277 operation_error: None, 1278 sqlite_error: Some(error), 1279 } 1280 } 1281 1282 #[must_use] 1283 pub const fn kind(&self) -> ServiceSqliteTransactionErrorKind { 1284 self.kind 1285 } 1286 1287 #[must_use] 1288 pub const fn operation_error(&self) -> Option<&E> { 1289 self.operation_error.as_ref() 1290 } 1291 1292 #[must_use] 1293 pub const fn sqlite_error(&self) -> Option<&ServiceSqliteError> { 1294 self.sqlite_error.as_ref() 1295 } 1296 } 1297 1298 impl<E> fmt::Debug for ServiceSqliteTransactionError<E> { 1299 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1300 formatter 1301 .debug_struct("ServiceSqliteTransactionError") 1302 .field("kind", &self.kind) 1303 .field( 1304 "operation_error", 1305 &self.operation_error.as_ref().map(|_| "[redacted]"), 1306 ) 1307 .field( 1308 "sqlite_error", 1309 &self.sqlite_error.as_ref().map(|_| "[redacted]"), 1310 ) 1311 .finish() 1312 } 1313 } 1314 1315 impl<E> fmt::Display for ServiceSqliteTransactionError<E> { 1316 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1317 formatter.write_str(match self.kind { 1318 ServiceSqliteTransactionErrorKind::NotCommitted => { 1319 "SQLite transaction did not reach commit" 1320 } 1321 ServiceSqliteTransactionErrorKind::OperationRolledBack => { 1322 "SQLite transaction operation was rolled back" 1323 } 1324 ServiceSqliteTransactionErrorKind::RollbackFailed => { 1325 "SQLite transaction rollback could not be confirmed" 1326 } 1327 ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown => { 1328 "SQLite transaction commit outcome is unknown" 1329 } 1330 }) 1331 } 1332 } 1333 1334 impl<E: Send + Sync + 'static> Error for ServiceSqliteTransactionError<E> {} 1335 1336 #[cfg(any(target_os = "linux", target_os = "macos"))] 1337 struct QuarantinedConnection { 1338 connection: Option<PoolConnection<Sqlite>>, 1339 trusted: bool, 1340 } 1341 1342 #[cfg(any(target_os = "linux", target_os = "macos"))] 1343 impl QuarantinedConnection { 1344 fn new(connection: PoolConnection<Sqlite>) -> Self { 1345 Self { 1346 connection: Some(connection), 1347 trusted: false, 1348 } 1349 } 1350 1351 fn trust(&mut self) { 1352 self.trusted = true; 1353 } 1354 1355 fn into_close_future(mut self) -> Option<BoxFuture<'static, Result<(), sqlx::Error>>> { 1356 let connection = self.connection.take()?; 1357 Some(Box::pin(async move { connection.close().await })) 1358 } 1359 } 1360 1361 #[cfg(any(target_os = "linux", target_os = "macos"))] 1362 impl core::ops::Deref for QuarantinedConnection { 1363 type Target = SqliteConnection; 1364 1365 fn deref(&self) -> &Self::Target { 1366 self.connection.as_deref().expect("connection is retained") 1367 } 1368 } 1369 1370 #[cfg(any(target_os = "linux", target_os = "macos"))] 1371 impl core::ops::DerefMut for QuarantinedConnection { 1372 fn deref_mut(&mut self) -> &mut Self::Target { 1373 self.connection 1374 .as_deref_mut() 1375 .expect("connection is retained") 1376 } 1377 } 1378 1379 #[cfg(any(target_os = "linux", target_os = "macos"))] 1380 impl Drop for QuarantinedConnection { 1381 fn drop(&mut self) { 1382 if !self.trusted 1383 && let Some(connection) = self.connection.as_mut() 1384 { 1385 connection.close_on_drop(); 1386 } 1387 } 1388 } 1389 1390 #[cfg(any(target_os = "linux", target_os = "macos"))] 1391 fn sqlite_source(source: sqlx::Error) -> ServiceSqliteError { 1392 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) 1393 } 1394 1395 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 1396 fn unsupported_host() -> ServiceSqliteError { 1397 ServiceSqliteError::new(ServiceSqliteErrorKind::Open) 1398 } 1399 1400 #[cfg(test)] 1401 mod tests { 1402 #[cfg(any(target_os = "linux", target_os = "macos"))] 1403 use std::{ 1404 collections::BTreeMap, 1405 convert::Infallible, 1406 fs, 1407 num::NonZeroU32, 1408 os::unix::fs::PermissionsExt, 1409 path::Path, 1410 sync::{ 1411 Arc, 1412 atomic::{AtomicBool, Ordering}, 1413 }, 1414 time::{Duration, SystemTime}, 1415 }; 1416 1417 #[cfg(any(target_os = "linux", target_os = "macos"))] 1418 use radroots_runtime_paths::{ 1419 InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, 1420 RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, ServiceId, 1421 }; 1422 #[cfg(any(target_os = "linux", target_os = "macos"))] 1423 use radroots_storage::event::SourceGeneration; 1424 #[cfg(any(target_os = "linux", target_os = "macos"))] 1425 use sqlx::{Connection, sqlite::SqliteConnectOptions}; 1426 #[cfg(any(target_os = "linux", target_os = "macos"))] 1427 use tokio::sync::Notify; 1428 1429 #[cfg(any(target_os = "linux", target_os = "macos"))] 1430 static CAPTURE_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); 1431 1432 use super::*; 1433 1434 #[cfg(any(target_os = "linux", target_os = "macos"))] 1435 const HOST_TABLE_SQL: &str = "CREATE TABLE host_probe (value INTEGER NOT NULL)"; 1436 1437 #[cfg(any(target_os = "linux", target_os = "macos"))] 1438 fn runtime_context(root: &std::path::Path) -> RuntimeContext { 1439 RuntimeContext::resolve( 1440 &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), 1441 RuntimeContextBootstrap::new( 1442 RadrootsPathProfile::RepoLocal, 1443 Some(root.to_path_buf()), 1444 RuntimeContextSource::BootstrapCli, 1445 RuntimeContextSource::BootstrapCli, 1446 ) 1447 .expect("valid bootstrap"), 1448 ServiceId::new("myc").expect("service ID"), 1449 InstanceId::new("host-boundary").expect("instance ID"), 1450 ) 1451 .expect("runtime context") 1452 } 1453 1454 #[cfg(any(target_os = "linux", target_os = "macos"))] 1455 #[test] 1456 fn rollback_failure_selection_preserves_each_exact_precedence() { 1457 let error = || ServiceSqliteError::new(ServiceSqliteErrorKind::Open); 1458 for authority in [false, true] { 1459 for rollback in [false, true] { 1460 for confirmed in [false, true] { 1461 for removal in [false, true] { 1462 let expected_precondition = if authority { 1463 1 1464 } else if rollback && !confirmed { 1465 2 1466 } else if removal { 1467 3 1468 } else { 1469 0 1470 }; 1471 let precondition = precondition_rollback_failure( 1472 authority.then(error), 1473 rollback.then(error), 1474 confirmed, 1475 removal.then(error), 1476 ); 1477 assert_eq!( 1478 usize::from(precondition.is_some()), 1479 usize::from(expected_precondition != 0) 1480 ); 1481 1482 let expected_drift = (rollback && !confirmed) || removal; 1483 assert_eq!( 1484 authority_drift_rollback_failure( 1485 rollback.then(error), 1486 confirmed, 1487 removal.then(error), 1488 ) 1489 .is_some(), 1490 expected_drift 1491 ); 1492 1493 let expected_operation = (rollback && !confirmed) || removal || authority; 1494 assert_eq!( 1495 operation_rollback_failure( 1496 rollback.then(error), 1497 confirmed, 1498 removal.then(error), 1499 authority.then(error), 1500 ) 1501 .is_some(), 1502 expected_operation 1503 ); 1504 } 1505 } 1506 } 1507 } 1508 assert!(unconfirmed_rollback_error(Some(error()), false).is_some()); 1509 assert!(unconfirmed_rollback_error(Some(error()), true).is_none()); 1510 assert!(unconfirmed_rollback_error(None, false).is_none()); 1511 1512 assert!(integrity_driver_close_result(Ok(()), false).is_ok()); 1513 assert!(matches!( 1514 integrity_driver_close_result(Ok(()), true), 1515 Err(IntegrityInspectionDriverFailure::ConnectionClose) 1516 )); 1517 assert!(matches!( 1518 integrity_driver_close_result(Err(sqlx::Error::Protocol("close".to_owned())), false), 1519 Err(IntegrityInspectionDriverFailure::ConnectionClose) 1520 )); 1521 } 1522 1523 #[cfg(any(target_os = "linux", target_os = "macos"))] 1524 #[tokio::test(flavor = "current_thread")] 1525 async fn final_connection_policy_classifier_preserves_pragma_kind() { 1526 let mut first = 1527 SqliteConnection::connect_with(&SqliteConnectOptions::new().filename(":memory:")) 1528 .await 1529 .expect("first connection"); 1530 let mut second = 1531 SqliteConnection::connect_with(&SqliteConnectOptions::new().filename(":memory:")) 1532 .await 1533 .expect("second connection"); 1534 let initial = crate::migration::read_connection_policy(&mut first) 1535 .await 1536 .expect("initial policy"); 1537 let same = crate::migration::read_connection_policy(&mut second) 1538 .await 1539 .expect("same policy"); 1540 assert!(final_connection_policy_matches(&initial, &same).is_ok()); 1541 sqlx::query("PRAGMA query_only = ON") 1542 .execute(&mut second) 1543 .await 1544 .expect("change policy"); 1545 let changed = crate::migration::read_connection_policy(&mut second) 1546 .await 1547 .expect("changed policy"); 1548 let error = final_connection_policy_matches(&initial, &changed).expect_err("policy drift"); 1549 assert_eq!(error.kind(), ServiceSqliteErrorKind::Pragma); 1550 } 1551 1552 #[cfg(any(target_os = "linux", target_os = "macos"))] 1553 fn migration_catalog() -> MigrationCatalog { 1554 MigrationCatalog::new([]).expect("empty v1 migration catalog") 1555 } 1556 1557 #[cfg(any(target_os = "linux", target_os = "macos"))] 1558 fn schema_catalog(migrations: &MigrationCatalog) -> SchemaCatalog { 1559 let table = crate::SchemaObject::new( 1560 crate::SchemaObjectKind::Table, 1561 "host_probe", 1562 "host_probe", 1563 HOST_TABLE_SQL, 1564 crate::SchemaObject::computed_digest( 1565 crate::SchemaObjectKind::Table, 1566 "host_probe", 1567 "host_probe", 1568 HOST_TABLE_SQL, 1569 ) 1570 .expect("table digest"), 1571 ) 1572 .expect("table descriptor"); 1573 let version_digest = crate::SchemaVersionCatalog::computed_digest(1, [table.clone()]) 1574 .expect("version digest"); 1575 let version = 1576 crate::SchemaVersionCatalog::new(1, [table], version_digest).expect("schema version"); 1577 SchemaCatalog::new(migrations, [version]).expect("schema catalog") 1578 } 1579 1580 #[cfg(any(target_os = "linux", target_os = "macos"))] 1581 fn build_identity() -> MigrationBuildIdentity { 1582 MigrationBuildIdentity::new( 1583 "0.1.0-alpha", 1584 "0123456789abcdef0123456789abcdef01234567", 1585 "89abcdef0123456789abcdef0123456789abcdef", 1586 "1.97.1", 1587 "x86_64-unknown-linux-gnu", 1588 "service-host", 1589 1, 1590 2, 1591 3, 1592 4, 1593 5, 1594 ) 1595 .expect("build identity") 1596 } 1597 1598 #[cfg(any(target_os = "linux", target_os = "macos"))] 1599 async fn initialized_host() -> ( 1600 tempfile::TempDir, 1601 ServiceSqlitePaths, 1602 ServiceDatabaseIdentity, 1603 MigrationCatalog, 1604 SchemaCatalog, 1605 ServiceSqliteHost, 1606 ) { 1607 initialized_host_with_options(ServiceSqliteConnectionOptions::reviewed()).await 1608 } 1609 1610 #[cfg(any(target_os = "linux", target_os = "macos"))] 1611 async fn initialized_host_with_options( 1612 options: ServiceSqliteConnectionOptions, 1613 ) -> ( 1614 tempfile::TempDir, 1615 ServiceSqlitePaths, 1616 ServiceDatabaseIdentity, 1617 MigrationCatalog, 1618 SchemaCatalog, 1619 ServiceSqliteHost, 1620 ) { 1621 let root = tempfile::tempdir().expect("temporary root"); 1622 let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(root.path())) 1623 .expect("SQLite paths"); 1624 fs::create_dir_all(paths.state_database().parent().expect("state directory")) 1625 .expect("create state directory"); 1626 let metadata = crate::ServiceDatabaseMetadata::new( 1627 &paths, 1628 SourceGeneration::new([9; 32]).expect("source generation"), 1629 NonZeroU32::new(1).expect("schema version"), 1630 1_700_000_000_000, 1631 crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), 1632 ) 1633 .expect("database metadata"); 1634 let migrations = migration_catalog(); 1635 let schema = schema_catalog(&migrations); 1636 let authority = crate::initialize_database( 1637 &paths, 1638 OpenMode::Initialize, 1639 &metadata, 1640 &schema, 1641 |initializer| { 1642 Box::pin(async move { 1643 sqlx::query(HOST_TABLE_SQL) 1644 .execute(initializer) 1645 .await 1646 .expect("create host table"); 1647 Ok::<_, Infallible>(()) 1648 }) 1649 }, 1650 ) 1651 .await 1652 .expect("initialize database"); 1653 let identity = metadata.identity(); 1654 let (host, outcome) = ServiceSqliteHost::open_initialized( 1655 &paths, 1656 &identity, 1657 &migrations, 1658 &schema, 1659 options, 1660 authority, 1661 MigrationAppliedAtUnixSeconds::new(1_700_000_000).expect("migration time"), 1662 &build_identity(), 1663 &[], 1664 ) 1665 .await 1666 .expect("open initialized host"); 1667 assert_eq!(outcome.initial_version(), 1); 1668 assert_eq!(outcome.final_version(), 1); 1669 assert_eq!(outcome.applied_count(), 0); 1670 (root, paths, identity, migrations, schema, host) 1671 } 1672 1673 #[cfg(any(target_os = "linux", target_os = "macos"))] 1674 async fn row_count(host: &ServiceSqliteHost) -> i64 { 1675 host.transaction(|transaction| { 1676 Box::pin(async move { 1677 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") 1678 .fetch_one(&mut *transaction) 1679 .await 1680 }) 1681 }) 1682 .await 1683 .expect("count rows") 1684 } 1685 1686 #[cfg(any(target_os = "linux", target_os = "macos"))] 1687 fn integrity_checked_at(value: u64) -> crate::IntegrityCheckedAtUnixMs { 1688 crate::IntegrityCheckedAtUnixMs::new(value).expect("integrity inspection time") 1689 } 1690 1691 #[cfg(any(target_os = "linux", target_os = "macos"))] 1692 #[tokio::test] 1693 async fn existing_intent_returns_actual_metadata_without_releasing_authority() { 1694 let (_root, paths, identity, migrations, schema, initialized) = initialized_host().await; 1695 initialized.close().await.expect("close initialized host"); 1696 let intent = ExistingServiceDatabaseIntent::new( 1697 &paths, 1698 identity.supported_state_schema_version(), 1699 identity.application_id(), 1700 ); 1701 1702 let wrong_application = ExistingServiceDatabaseIntent::new( 1703 &paths, 1704 identity.supported_state_schema_version(), 1705 crate::ServiceSqliteApplicationId::new(7).expect("other application"), 1706 ); 1707 let error = ServiceSqliteHost::open_read_write_existing_with_intent( 1708 &paths, 1709 &wrong_application, 1710 &migrations, 1711 &schema, 1712 ServiceSqliteConnectionOptions::reviewed(), 1713 MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"), 1714 &build_identity(), 1715 &[], 1716 ) 1717 .await 1718 .expect_err("application mismatch"); 1719 assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata); 1720 1721 let (opened, outcome) = ServiceSqliteHost::open_read_write_existing_with_intent( 1722 &paths, 1723 &intent, 1724 &migrations, 1725 &schema, 1726 ServiceSqliteConnectionOptions::reviewed(), 1727 MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"), 1728 &build_identity(), 1729 &[], 1730 ) 1731 .await 1732 .expect("open writable from intent"); 1733 assert_eq!(outcome.applied_count(), 0); 1734 assert_eq!( 1735 opened.database_metadata().source_generation(), 1736 identity.source_generation() 1737 ); 1738 assert_eq!(opened.host().mode(), OpenMode::ReadWriteExisting); 1739 let debug = format!("{opened:?}"); 1740 assert!(debug.contains("OpenedServiceDatabase")); 1741 assert!(!debug.contains("09090909")); 1742 assert!(WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting).is_err()); 1743 let (writable, actual) = opened.into_parts(); 1744 assert_eq!(actual.source_generation(), identity.source_generation()); 1745 assert_eq!(row_count(&writable).await, 0); 1746 writable.close().await.expect("close writable host"); 1747 1748 let inspected = ServiceSqliteHost::open_read_only_inspection_with_intent( 1749 &paths, 1750 &intent, 1751 &migrations, 1752 &schema, 1753 ServiceSqliteConnectionOptions::reviewed(), 1754 ) 1755 .await 1756 .expect("open inspection from intent"); 1757 assert_eq!(inspected.host().mode(), OpenMode::ReadOnlyInspection); 1758 assert_eq!( 1759 inspected.database_metadata().source_generation(), 1760 identity.source_generation() 1761 ); 1762 assert!(WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting).is_err()); 1763 let (inspection, actual) = inspected.into_parts(); 1764 assert_eq!(actual.source_generation(), identity.source_generation()); 1765 inspection.close().await.expect("close inspection host"); 1766 } 1767 1768 #[cfg(any(target_os = "linux", target_os = "macos"))] 1769 #[tokio::test(flavor = "current_thread")] 1770 async fn open_or_initialize_uses_exclusive_create_without_filesystem_probing() { 1771 let root = tempfile::tempdir().expect("temporary root"); 1772 let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(root.path())) 1773 .expect("SQLite paths"); 1774 fs::create_dir_all(paths.state_database().parent().expect("state directory")) 1775 .expect("create state directory"); 1776 let metadata = ServiceDatabaseMetadata::new( 1777 &paths, 1778 SourceGeneration::new([31; 32]).expect("source generation"), 1779 NonZeroU32::new(1).expect("schema version"), 1780 1_700_000_000_000, 1781 crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), 1782 ) 1783 .expect("database metadata"); 1784 let migrations = migration_catalog(); 1785 let schema = schema_catalog(&migrations); 1786 let initialized_callback = Arc::new(AtomicBool::new(false)); 1787 let called = Arc::clone(&initialized_callback); 1788 let (initialized, outcome) = ServiceSqliteHost::open_or_initialize( 1789 &paths, 1790 &metadata, 1791 &migrations, 1792 &schema, 1793 ServiceSqliteConnectionOptions::reviewed(), 1794 MigrationAppliedAtUnixSeconds::new(1_700_000_000).expect("migration time"), 1795 &build_identity(), 1796 &[], 1797 move |initializer| { 1798 called.store(true, Ordering::Release); 1799 Box::pin(async move { 1800 sqlx::query(HOST_TABLE_SQL).execute(initializer).await?; 1801 Ok::<(), sqlx::Error>(()) 1802 }) 1803 }, 1804 ) 1805 .await 1806 .expect("initialize missing state"); 1807 assert!(initialized_callback.load(Ordering::Acquire)); 1808 assert_eq!(initialized.host().mode(), OpenMode::Initialize); 1809 assert_eq!(initialized.database_metadata(), &metadata); 1810 assert_eq!(outcome.applied_count(), 0); 1811 let (initialized, initialized_metadata) = initialized.into_parts(); 1812 assert_eq!(initialized_metadata, metadata); 1813 initialized.close().await.expect("close initialized host"); 1814 1815 let existing_callback = Arc::new(AtomicBool::new(false)); 1816 let called = Arc::clone(&existing_callback); 1817 let alternative_metadata = ServiceDatabaseMetadata::new( 1818 &paths, 1819 SourceGeneration::new([32; 32]).expect("alternative source generation"), 1820 NonZeroU32::new(1).expect("schema version"), 1821 1_700_000_000_001, 1822 crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), 1823 ) 1824 .expect("alternative database metadata"); 1825 let (existing, outcome) = ServiceSqliteHost::open_or_initialize( 1826 &paths, 1827 &alternative_metadata, 1828 &migrations, 1829 &schema, 1830 ServiceSqliteConnectionOptions::reviewed(), 1831 MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"), 1832 &build_identity(), 1833 &[], 1834 move |_| { 1835 called.store(true, Ordering::Release); 1836 Box::pin(async { Ok::<(), Infallible>(()) }) 1837 }, 1838 ) 1839 .await 1840 .expect("open exact existing state"); 1841 assert!(!existing_callback.load(Ordering::Acquire)); 1842 assert_eq!(existing.host().mode(), OpenMode::ReadWriteExisting); 1843 assert_eq!(existing.database_metadata(), &metadata); 1844 assert_eq!(outcome.applied_count(), 0); 1845 let (existing, existing_metadata) = existing.into_parts(); 1846 assert_eq!(existing_metadata, metadata); 1847 assert_eq!(row_count(&existing).await, 0); 1848 existing.close().await.expect("close existing host"); 1849 } 1850 1851 #[cfg(any(target_os = "linux", target_os = "macos"))] 1852 #[tokio::test(flavor = "current_thread")] 1853 async fn open_or_initialize_rejects_catalog_mismatch_before_reservation() { 1854 const MIGRATION_SQL: &str = "CREATE TABLE later_probe (value INTEGER NOT NULL)"; 1855 1856 let root = tempfile::tempdir().expect("temporary root"); 1857 let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(root.path())) 1858 .expect("SQLite paths"); 1859 fs::create_dir_all(paths.state_database().parent().expect("state directory")) 1860 .expect("create state directory"); 1861 let metadata = ServiceDatabaseMetadata::new( 1862 &paths, 1863 SourceGeneration::new([32; 32]).expect("source generation"), 1864 NonZeroU32::new(1).expect("schema version"), 1865 1_700_000_000_000, 1866 crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), 1867 ) 1868 .expect("database metadata"); 1869 let base_migrations = migration_catalog(); 1870 let schema = schema_catalog(&base_migrations); 1871 let different_migrations = MigrationCatalog::new([crate::MigrationDescriptor::sql( 1872 2, 1873 "create_later_probe", 1874 MIGRATION_SQL, 1875 crate::MigrationChecksum::for_sql(MIGRATION_SQL), 1876 ) 1877 .expect("different migration")]) 1878 .expect("different migration catalog"); 1879 let callback_called = Arc::new(AtomicBool::new(false)); 1880 let called = Arc::clone(&callback_called); 1881 1882 let error = ServiceSqliteHost::open_or_initialize( 1883 &paths, 1884 &metadata, 1885 &different_migrations, 1886 &schema, 1887 ServiceSqliteConnectionOptions::reviewed(), 1888 MigrationAppliedAtUnixSeconds::new(1_700_000_002).expect("migration time"), 1889 &build_identity(), 1890 &[], 1891 move |_| { 1892 called.store(true, Ordering::Release); 1893 Box::pin(async { Ok::<(), Infallible>(()) }) 1894 }, 1895 ) 1896 .await 1897 .expect_err("catalog mismatch must fail before reservation"); 1898 1899 assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); 1900 assert!(!callback_called.load(Ordering::Acquire)); 1901 assert!(!paths.state_database().exists()); 1902 } 1903 1904 #[cfg(any(target_os = "linux", target_os = "macos"))] 1905 #[tokio::test] 1906 async fn integrity_inspection_is_explicit_safe_and_available_in_every_host_mode() { 1907 let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await; 1908 crate::integrity::integrity_test_seam::release(); 1909 let (_root, paths, identity, migrations, schema, initialized) = initialized_host().await; 1910 1911 let initialized_report = initialized 1912 .inspect_integrity(integrity_checked_at(1_700_000_000_500)) 1913 .await 1914 .expect("inspect initialized host"); 1915 assert_eq!(initialized.mode(), OpenMode::Initialize); 1916 assert_eq!( 1917 initialized_report.sqlite(), 1918 crate::IntegrityCheckOutcome::Verified 1919 ); 1920 assert_eq!( 1921 initialized_report.foreign_keys(), 1922 crate::IntegrityCheckOutcome::Verified 1923 ); 1924 assert!(initialized_report.diagnostics().is_empty()); 1925 assert_eq!( 1926 initialized_report.storage_integrity(), 1927 crate::StorageIntegrity::Verified 1928 ); 1929 initialized.close().await.expect("close initialized host"); 1930 1931 let (writable, outcome) = ServiceSqliteHost::open_read_write_existing( 1932 &paths, 1933 &identity, 1934 &migrations, 1935 &schema, 1936 ServiceSqliteConnectionOptions::reviewed(), 1937 MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"), 1938 &build_identity(), 1939 &[], 1940 ) 1941 .await 1942 .expect("open writable host"); 1943 assert_eq!(outcome.applied_count(), 0); 1944 let writable_report = writable 1945 .inspect_integrity(integrity_checked_at(1_700_000_000_501)) 1946 .await 1947 .expect("inspect writable host"); 1948 assert_eq!(writable.mode(), OpenMode::ReadWriteExisting); 1949 assert_eq!( 1950 writable_report.storage_integrity(), 1951 crate::StorageIntegrity::Verified 1952 ); 1953 writable.close().await.expect("close writable host"); 1954 1955 let read_only = ServiceSqliteHost::open_read_only_inspection( 1956 &paths, 1957 &identity, 1958 &migrations, 1959 &schema, 1960 ServiceSqliteConnectionOptions::reviewed(), 1961 ) 1962 .await 1963 .expect("open read-only host"); 1964 let read_only_report = read_only 1965 .inspect_integrity(integrity_checked_at(1_700_000_000_502)) 1966 .await 1967 .expect("inspect read-only host"); 1968 assert_eq!(read_only.mode(), OpenMode::ReadOnlyInspection); 1969 assert_eq!( 1970 read_only_report.storage_integrity(), 1971 crate::StorageIntegrity::Verified 1972 ); 1973 read_only.close().await.expect("close read-only host"); 1974 } 1975 1976 #[cfg(any(target_os = "linux", target_os = "macos"))] 1977 #[tokio::test] 1978 async fn integrity_inspection_is_single_admission_cancel_safe_and_close_drained() { 1979 let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await; 1980 crate::integrity::integrity_test_seam::release(); 1981 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 1982 let host = Arc::new(host); 1983 1984 crate::integrity::integrity_test_seam::block( 1985 crate::integrity::integrity_test_seam::PHASE_BEFORE_SQLITE, 1986 ); 1987 let first = tokio::spawn({ 1988 let host = Arc::clone(&host); 1989 async move { 1990 host.inspect_integrity(integrity_checked_at(1_700_000_000_510)) 1991 .await 1992 } 1993 }); 1994 while crate::integrity::integrity_test_seam::reached() 1995 != crate::integrity::integrity_test_seam::PHASE_BEFORE_SQLITE 1996 { 1997 tokio::task::yield_now().await; 1998 } 1999 let concurrent = host 2000 .inspect_integrity(integrity_checked_at(1_700_000_000_511)) 2001 .await 2002 .expect_err("second integrity inspection is rejected"); 2003 assert_eq!(concurrent.kind(), ServiceSqliteErrorKind::Integrity); 2004 first.abort(); 2005 assert!( 2006 first 2007 .await 2008 .expect_err("inspection is cancelled") 2009 .is_cancelled() 2010 ); 2011 crate::integrity::integrity_test_seam::release(); 2012 assert!( 2013 !host.integrity_driver.lock().await.is_idle(), 2014 "cancelled SQLx connection remains host-owned until explicit cleanup" 2015 ); 2016 2017 let recovered = host 2018 .inspect_integrity(integrity_checked_at(1_700_000_000_512)) 2019 .await 2020 .expect("inspection recovers after cancellation"); 2021 assert_eq!( 2022 recovered.storage_integrity(), 2023 crate::StorageIntegrity::Verified 2024 ); 2025 2026 crate::integrity::integrity_test_seam::block( 2027 crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS, 2028 ); 2029 let timed_out = tokio::time::timeout( 2030 Duration::from_millis(20), 2031 host.inspect_integrity(integrity_checked_at(1_700_000_000_513)), 2032 ) 2033 .await; 2034 assert!(timed_out.is_err(), "caller deadline cancels inspection"); 2035 crate::integrity::integrity_test_seam::release(); 2036 assert!(!host.integrity_driver.lock().await.is_idle()); 2037 2038 crate::integrity::integrity_test_seam::block( 2039 crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK, 2040 ); 2041 let pre_rollback = tokio::spawn({ 2042 let host = Arc::clone(&host); 2043 async move { 2044 host.inspect_integrity(integrity_checked_at(1_700_000_000_514)) 2045 .await 2046 } 2047 }); 2048 while crate::integrity::integrity_test_seam::reached() 2049 != crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK 2050 { 2051 tokio::task::yield_now().await; 2052 } 2053 pre_rollback.abort(); 2054 assert!( 2055 pre_rollback 2056 .await 2057 .expect_err("pre-rollback inspection is cancelled") 2058 .is_cancelled() 2059 ); 2060 crate::integrity::integrity_test_seam::release(); 2061 assert!(!host.integrity_driver.lock().await.is_idle()); 2062 host.inspect_integrity(integrity_checked_at(1_700_000_000_515)) 2063 .await 2064 .expect("inspection recovers after pre-rollback cancellation"); 2065 2066 crate::integrity::integrity_test_seam::block( 2067 crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK, 2068 ); 2069 let admitted = tokio::spawn({ 2070 let host = Arc::clone(&host); 2071 async move { 2072 host.inspect_integrity(integrity_checked_at(1_700_000_000_516)) 2073 .await 2074 } 2075 }); 2076 while crate::integrity::integrity_test_seam::reached() 2077 != crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK 2078 { 2079 tokio::task::yield_now().await; 2080 } 2081 let close = tokio::spawn({ 2082 let host = Arc::clone(&host); 2083 async move { host.close().await } 2084 }); 2085 while !host.closing.load(Ordering::Acquire) { 2086 tokio::task::yield_now().await; 2087 } 2088 let rejected = host 2089 .inspect_integrity(integrity_checked_at(1_700_000_000_517)) 2090 .await 2091 .expect_err("closing host rejects inspection"); 2092 assert_eq!(rejected.kind(), ServiceSqliteErrorKind::Open); 2093 assert!(!close.is_finished()); 2094 crate::integrity::integrity_test_seam::release(); 2095 admitted 2096 .await 2097 .expect("admitted inspection task joins") 2098 .expect("admitted inspection completes"); 2099 close 2100 .await 2101 .expect("close task joins") 2102 .expect("close drains admitted inspection"); 2103 let closed = host 2104 .inspect_integrity(integrity_checked_at(1_700_000_000_518)) 2105 .await 2106 .expect_err("closed host rejects inspection"); 2107 assert_eq!(closed.kind(), ServiceSqliteErrorKind::Open); 2108 } 2109 2110 #[cfg(any(target_os = "linux", target_os = "macos"))] 2111 #[tokio::test] 2112 async fn integrity_inspection_preserves_authority_precedence_after_await() { 2113 let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await; 2114 crate::integrity::integrity_test_seam::release(); 2115 let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 2116 let host = Arc::new(host); 2117 crate::integrity::integrity_test_seam::block( 2118 crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS, 2119 ); 2120 let inspection = tokio::spawn({ 2121 let host = Arc::clone(&host); 2122 async move { 2123 host.inspect_integrity(integrity_checked_at(1_700_000_000_520)) 2124 .await 2125 } 2126 }); 2127 while crate::integrity::integrity_test_seam::reached() 2128 != crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS 2129 { 2130 tokio::task::yield_now().await; 2131 } 2132 let retired_lock = paths 2133 .state_lock() 2134 .parent() 2135 .expect("state directory") 2136 .join("retired-integrity-state.lock"); 2137 fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock"); 2138 crate::integrity::integrity_test_seam::release(); 2139 let error = inspection 2140 .await 2141 .expect("inspection task joins") 2142 .expect_err("authority drift rejects inspection"); 2143 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 2144 fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock"); 2145 host.close().await.expect("close host after restored lock"); 2146 } 2147 2148 #[cfg(any(target_os = "linux", target_os = "macos"))] 2149 #[tokio::test] 2150 async fn observed_authority_drift_precedes_transient_restore_and_close_failure() { 2151 let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await; 2152 crate::integrity::integrity_test_seam::release(); 2153 crate::integrity::integrity_test_seam::inject_connection_close_failure(false); 2154 let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 2155 let host = Arc::new(host); 2156 2157 crate::integrity::integrity_test_seam::block( 2158 crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS, 2159 ); 2160 let inspection = tokio::spawn({ 2161 let host = Arc::clone(&host); 2162 async move { 2163 host.inspect_integrity(integrity_checked_at(1_700_000_000_525)) 2164 .await 2165 } 2166 }); 2167 while crate::integrity::integrity_test_seam::reached() 2168 != crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS 2169 { 2170 tokio::task::yield_now().await; 2171 } 2172 2173 let retired_lock = paths 2174 .state_lock() 2175 .parent() 2176 .expect("state directory") 2177 .join("transient-integrity-state.lock"); 2178 fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock"); 2179 crate::integrity::integrity_test_seam::block( 2180 crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING, 2181 ); 2182 while crate::integrity::integrity_test_seam::reached() 2183 != crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING 2184 { 2185 tokio::task::yield_now().await; 2186 } 2187 2188 fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock"); 2189 crate::integrity::integrity_test_seam::inject_connection_close_failure(true); 2190 crate::integrity::integrity_test_seam::release(); 2191 let error = inspection 2192 .await 2193 .expect("inspection task joins") 2194 .expect_err("observed authority drift remains terminal"); 2195 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 2196 assert!(host.integrity_driver.lock().await.is_idle()); 2197 host.close().await.expect("close host after restored lock"); 2198 crate::integrity::integrity_test_seam::inject_connection_close_failure(false); 2199 } 2200 2201 #[cfg(any(target_os = "linux", target_os = "macos"))] 2202 #[tokio::test] 2203 async fn cancelled_real_sqlite_work_is_explicitly_closed_before_retry_or_host_close() { 2204 let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await; 2205 crate::integrity::integrity_test_seam::release(); 2206 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 2207 let host = Arc::new(host); 2208 2209 crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true); 2210 let cancelled = tokio::spawn({ 2211 let host = Arc::clone(&host); 2212 async move { 2213 host.inspect_integrity(integrity_checked_at(1_700_000_000_530)) 2214 .await 2215 } 2216 }); 2217 while crate::integrity::integrity_test_seam::reached() 2218 != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING 2219 { 2220 tokio::task::yield_now().await; 2221 } 2222 assert!(!cancelled.is_finished(), "SQLite probe remains in flight"); 2223 cancelled.abort(); 2224 assert!( 2225 cancelled 2226 .await 2227 .expect_err("real SQLite inspection is cancelled") 2228 .is_cancelled() 2229 ); 2230 crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false); 2231 assert!(!host.integrity_driver.lock().await.is_idle()); 2232 2233 let cancelled_cleanup = tokio::spawn({ 2234 let host = Arc::clone(&host); 2235 async move { 2236 host.inspect_integrity(integrity_checked_at(1_700_000_000_531)) 2237 .await 2238 } 2239 }); 2240 while crate::integrity::integrity_test_seam::reached() 2241 != crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING 2242 { 2243 tokio::task::yield_now().await; 2244 } 2245 cancelled_cleanup.abort(); 2246 assert!( 2247 cancelled_cleanup 2248 .await 2249 .expect_err("retained close retry is cancelled") 2250 .is_cancelled() 2251 ); 2252 assert!(!host.integrity_driver.lock().await.is_idle()); 2253 2254 let recovered = host 2255 .inspect_integrity(integrity_checked_at(1_700_000_000_532)) 2256 .await 2257 .expect("retry explicitly closes prior SQLite worker before inspecting"); 2258 assert_eq!( 2259 recovered.storage_integrity(), 2260 crate::StorageIntegrity::Verified 2261 ); 2262 assert!(host.integrity_driver.lock().await.is_idle()); 2263 2264 crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true); 2265 let cancelled = tokio::spawn({ 2266 let host = Arc::clone(&host); 2267 async move { 2268 host.inspect_integrity(integrity_checked_at(1_700_000_000_533)) 2269 .await 2270 } 2271 }); 2272 while crate::integrity::integrity_test_seam::reached() 2273 != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING 2274 { 2275 tokio::task::yield_now().await; 2276 } 2277 cancelled.abort(); 2278 assert!( 2279 cancelled 2280 .await 2281 .expect_err("second real SQLite inspection is cancelled") 2282 .is_cancelled() 2283 ); 2284 crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false); 2285 assert!(!host.integrity_driver.lock().await.is_idle()); 2286 host.close() 2287 .await 2288 .expect("host close explicitly terminates retained SQLite worker"); 2289 assert!(host.integrity_driver.lock().await.is_idle()); 2290 } 2291 2292 #[cfg(any(target_os = "linux", target_os = "macos"))] 2293 #[tokio::test] 2294 async fn retained_integrity_close_releases_authority_and_caches_concurrent_lock_drift() { 2295 let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await; 2296 crate::integrity::integrity_test_seam::release(); 2297 let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 2298 let host = Arc::new(host); 2299 2300 crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true); 2301 let inspection = tokio::spawn({ 2302 let host = Arc::clone(&host); 2303 async move { 2304 host.inspect_integrity(integrity_checked_at(1_700_000_000_540)) 2305 .await 2306 } 2307 }); 2308 while crate::integrity::integrity_test_seam::reached() 2309 != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING 2310 { 2311 tokio::task::yield_now().await; 2312 } 2313 inspection.abort(); 2314 assert!( 2315 inspection 2316 .await 2317 .expect_err("real integrity work is cancelled") 2318 .is_cancelled() 2319 ); 2320 crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false); 2321 2322 let close = tokio::spawn({ 2323 let host = Arc::clone(&host); 2324 async move { host.close().await } 2325 }); 2326 while crate::integrity::integrity_test_seam::reached() 2327 != crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING 2328 { 2329 tokio::task::yield_now().await; 2330 } 2331 let retired_lock = paths 2332 .state_lock() 2333 .parent() 2334 .expect("state directory") 2335 .join("retired-integrity-close-state.lock"); 2336 fs::rename(paths.state_lock(), &retired_lock).expect("retire held writer lock"); 2337 fs::write(paths.state_lock(), b"").expect("create replacement writer lock"); 2338 fs::set_permissions(paths.state_lock(), fs::Permissions::from_mode(0o600)) 2339 .expect("replacement writer lock mode"); 2340 2341 let error = close 2342 .await 2343 .expect("close task joins") 2344 .expect_err("authority drift is terminal"); 2345 assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); 2346 let repeated = host.close().await.expect_err("terminal result is cached"); 2347 assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Authority); 2348 assert!(host.integrity_driver.lock().await.is_idle()); 2349 2350 fs::remove_file(paths.state_lock()).expect("remove replacement writer lock"); 2351 fs::rename(&retired_lock, paths.state_lock()).expect("restore original writer lock"); 2352 let mut authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 2353 .expect("authority reacquisition") 2354 .expect("writer authority"); 2355 authority.release().expect("release reacquired authority"); 2356 } 2357 2358 #[cfg(any(target_os = "linux", target_os = "macos"))] 2359 #[tokio::test] 2360 async fn retained_integrity_close_failure_is_cached_as_open_and_releases_authority() { 2361 let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await; 2362 crate::integrity::integrity_test_seam::release(); 2363 crate::integrity::integrity_test_seam::inject_connection_close_failure(false); 2364 let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 2365 let host = Arc::new(host); 2366 2367 crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true); 2368 let inspection = tokio::spawn({ 2369 let host = Arc::clone(&host); 2370 async move { 2371 host.inspect_integrity(integrity_checked_at(1_700_000_000_550)) 2372 .await 2373 } 2374 }); 2375 while crate::integrity::integrity_test_seam::reached() 2376 != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING 2377 { 2378 tokio::task::yield_now().await; 2379 } 2380 inspection.abort(); 2381 assert!( 2382 inspection 2383 .await 2384 .expect_err("real integrity work is cancelled") 2385 .is_cancelled() 2386 ); 2387 crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false); 2388 assert!(!host.integrity_driver.lock().await.is_idle()); 2389 2390 crate::integrity::integrity_test_seam::inject_connection_close_failure(true); 2391 let error = host 2392 .close() 2393 .await 2394 .expect_err("retained connection close failure is terminal"); 2395 assert_eq!(error.kind(), ServiceSqliteErrorKind::Open); 2396 let repeated = host.close().await.expect_err("terminal result is cached"); 2397 assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Open); 2398 assert!(host.integrity_driver.lock().await.is_idle()); 2399 2400 let mut authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 2401 .expect("authority reacquisition") 2402 .expect("writer authority"); 2403 authority.release().expect("release reacquired authority"); 2404 crate::integrity::integrity_test_seam::inject_connection_close_failure(false); 2405 } 2406 2407 #[cfg(any(target_os = "linux", target_os = "macos"))] 2408 #[tokio::test] 2409 async fn online_backup_captures_exact_member_manifest_and_preserves_source() { 2410 use sha2::Digest; 2411 2412 let _serial = CAPTURE_TEST_LOCK.lock().await; 2413 crate::backup::test_capture_reset(); 2414 let (root, paths, identity, migrations, schema, host) = initialized_host().await; 2415 host.transaction(|transaction| { 2416 Box::pin(async move { 2417 sqlx::query("INSERT INTO host_probe (value) VALUES (41), (42)") 2418 .execute(&mut *transaction) 2419 .await 2420 .map(|_| ()) 2421 }) 2422 }) 2423 .await 2424 .expect("seed live WAL state"); 2425 2426 let output = root.path().join("backup-output"); 2427 fs::create_dir(&output).expect("backup parent"); 2428 fs::set_permissions(&output, fs::Permissions::from_mode(0o700)) 2429 .expect("backup parent mode"); 2430 let collision = output.join("collision"); 2431 fs::create_dir(&collision).expect("preexisting collision"); 2432 fs::write(collision.join("foreign"), b"preserve").expect("foreign collision member"); 2433 let collision_error = host 2434 .capture_online_backup( 2435 &collision, 2436 crate::BackupCreatedAtUnixMs::new(1_700_000_000_100).expect("backup creation time"), 2437 ) 2438 .await 2439 .expect_err("existing destination is rejected"); 2440 assert_eq!(collision_error.kind(), ServiceSqliteErrorKind::Backup); 2441 assert_eq!( 2442 fs::read(collision.join("foreign")).expect("preserved collision"), 2443 b"preserve" 2444 ); 2445 2446 let source_database_before = fs::read(paths.state_database()).expect("source bytes"); 2447 let source_inventory_before = 2448 fs::read_dir(paths.state_database().parent().expect("source directory")) 2449 .expect("source inventory") 2450 .map(|entry| entry.expect("source entry").file_name()) 2451 .collect::<std::collections::BTreeSet<_>>(); 2452 let stage = output.join("successful"); 2453 let created_at = 2454 crate::BackupCreatedAtUnixMs::new(1_700_000_000_101).expect("backup creation time"); 2455 let manifest = host 2456 .capture_online_backup(&stage, created_at) 2457 .await 2458 .expect("online backup"); 2459 2460 let verified = crate::verify_backup_bundle( 2461 manifest.canonical_bytes(), 2462 manifest.digest(), 2463 &stage, 2464 &identity, 2465 std::num::NonZeroU64::new(manifest.members()[0].byte_length()) 2466 .expect("positive captured member length"), 2467 ) 2468 .expect("independently verify captured bundle"); 2469 assert_eq!(verified.manifest(), &manifest); 2470 assert_eq!(verified.database_metadata().service(), identity.service()); 2471 assert_eq!(verified.database_metadata().instance(), identity.instance()); 2472 2473 assert_eq!(manifest.service(), paths.service()); 2474 assert_eq!(manifest.instance(), paths.instance()); 2475 assert_eq!(manifest.source_generation(), identity.source_generation()); 2476 assert_eq!( 2477 manifest.state_schema_version(), 2478 identity.supported_state_schema_version() 2479 ); 2480 assert_eq!(manifest.created_at_unix_ms(), created_at); 2481 assert_eq!(manifest.members().len(), 1); 2482 assert!(!manifest.protected_material_included()); 2483 assert_eq!(manifest.integrity().sqlite(), "ok"); 2484 assert_eq!(manifest.integrity().foreign_keys(), "ok"); 2485 2486 let state = stage.join("state.sqlite"); 2487 let bytes = fs::read(&state).expect("captured state bytes"); 2488 let digest: [u8; 32] = sha2::Sha256::digest(&bytes).into(); 2489 assert_eq!(manifest.members()[0].byte_length(), bytes.len() as u64); 2490 assert_eq!(manifest.members()[0].sha256().as_bytes(), &digest); 2491 assert_eq!( 2492 fs::metadata(&stage) 2493 .expect("stage metadata") 2494 .permissions() 2495 .mode() 2496 & 0o777, 2497 0o700 2498 ); 2499 assert_eq!( 2500 fs::metadata(&state) 2501 .expect("state metadata") 2502 .permissions() 2503 .mode() 2504 & 0o777, 2505 0o600 2506 ); 2507 assert_eq!( 2508 fs::read_dir(&stage) 2509 .expect("stage inventory") 2510 .map(|entry| entry.expect("stage entry").file_name()) 2511 .collect::<Vec<_>>(), 2512 vec![std::ffi::OsString::from("state.sqlite")] 2513 ); 2514 2515 let mut backup = SqliteConnection::connect_with( 2516 &SqliteConnectOptions::new().filename(&state).read_only(true), 2517 ) 2518 .await 2519 .expect("open captured database"); 2520 assert_eq!( 2521 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") 2522 .fetch_one(&mut backup) 2523 .await 2524 .expect("captured row count"), 2525 2 2526 ); 2527 backup.close().await.expect("close captured database"); 2528 2529 assert_eq!( 2530 fs::read(paths.state_database()).expect("source bytes after capture"), 2531 source_database_before 2532 ); 2533 assert_eq!( 2534 fs::read_dir(paths.state_database().parent().expect("source directory")) 2535 .expect("source inventory after") 2536 .map(|entry| entry.expect("source entry").file_name()) 2537 .collect::<std::collections::BTreeSet<_>>(), 2538 source_inventory_before 2539 ); 2540 2541 host.close().await.expect("close writable host"); 2542 let inspection = ServiceSqliteHost::open_read_only_inspection( 2543 &paths, 2544 &identity, 2545 &migrations, 2546 &schema, 2547 ServiceSqliteConnectionOptions::reviewed(), 2548 ) 2549 .await 2550 .expect("open inspection"); 2551 let forbidden = output.join("read-only-forbidden"); 2552 let error = inspection 2553 .capture_online_backup(&forbidden, created_at) 2554 .await 2555 .expect_err("read-only capture is unavailable"); 2556 assert_eq!(error.kind(), ServiceSqliteErrorKind::Open); 2557 assert!(!forbidden.exists()); 2558 inspection.close().await.expect("close inspection"); 2559 } 2560 2561 #[cfg(any(target_os = "linux", target_os = "macos"))] 2562 #[tokio::test] 2563 async fn cancelled_capture_cleans_up_close_drains_and_next_capture_recovers() { 2564 let _serial = CAPTURE_TEST_LOCK.lock().await; 2565 crate::backup::test_capture_reset(); 2566 let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 2567 let host = Arc::new(host); 2568 let output = root.path().join("cancel-output"); 2569 fs::create_dir(&output).expect("backup parent"); 2570 fs::set_permissions(&output, fs::Permissions::from_mode(0o700)) 2571 .expect("backup parent mode"); 2572 let cancelled_stage = output.join("cancelled"); 2573 crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED); 2574 let capture = tokio::spawn({ 2575 let host = Arc::clone(&host); 2576 let stage = cancelled_stage.clone(); 2577 async move { 2578 host.capture_online_backup( 2579 &stage, 2580 crate::BackupCreatedAtUnixMs::new(1_700_000_000_200) 2581 .expect("backup creation time"), 2582 ) 2583 .await 2584 } 2585 }); 2586 while crate::backup::test_capture_phase() 2587 != crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED 2588 { 2589 tokio::task::yield_now().await; 2590 } 2591 2592 let concurrent_stage = output.join("concurrent"); 2593 let concurrent = host 2594 .capture_online_backup( 2595 &concurrent_stage, 2596 crate::BackupCreatedAtUnixMs::new(1_700_000_000_201).expect("backup creation time"), 2597 ) 2598 .await 2599 .expect_err("second capture is rejected"); 2600 assert_eq!(concurrent.kind(), ServiceSqliteErrorKind::Backup); 2601 assert!(!concurrent_stage.exists()); 2602 2603 capture.abort(); 2604 assert!( 2605 capture 2606 .await 2607 .expect_err("capture task is cancelled") 2608 .is_cancelled() 2609 ); 2610 while host.backup_active.load(Ordering::Acquire) { 2611 tokio::task::yield_now().await; 2612 } 2613 crate::backup::test_capture_reset(); 2614 assert!(!cancelled_stage.exists()); 2615 2616 let recovered_stage = output.join("recovered"); 2617 host.capture_online_backup( 2618 &recovered_stage, 2619 crate::BackupCreatedAtUnixMs::new(1_700_000_000_202).expect("backup creation time"), 2620 ) 2621 .await 2622 .expect("capture recovers after cancellation"); 2623 assert!(recovered_stage.join("state.sqlite").exists()); 2624 2625 crate::backup::test_capture_reset(); 2626 crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED); 2627 let close_drained_stage = output.join("close-drained"); 2628 let capture = tokio::spawn({ 2629 let host = Arc::clone(&host); 2630 let stage = close_drained_stage.clone(); 2631 async move { 2632 host.capture_online_backup( 2633 &stage, 2634 crate::BackupCreatedAtUnixMs::new(1_700_000_000_203) 2635 .expect("backup creation time"), 2636 ) 2637 .await 2638 } 2639 }); 2640 while crate::backup::test_capture_phase() 2641 != crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED 2642 { 2643 tokio::task::yield_now().await; 2644 } 2645 capture.abort(); 2646 assert!( 2647 capture 2648 .await 2649 .expect_err("capture task is cancelled") 2650 .is_cancelled() 2651 ); 2652 host.close() 2653 .await 2654 .expect("close drains every admitted capture"); 2655 crate::backup::test_capture_reset(); 2656 assert!(!close_drained_stage.exists()); 2657 assert!(!host.backup_active.load(Ordering::Acquire)); 2658 } 2659 2660 #[cfg(any(target_os = "linux", target_os = "macos"))] 2661 #[tokio::test] 2662 async fn capture_cancellation_is_cleanup_safe_at_every_governed_phase() { 2663 let _serial = CAPTURE_TEST_LOCK.lock().await; 2664 2665 for (index, phase) in [ 2666 crate::backup::TEST_CAPTURE_PHASE_BEFORE_CREATE, 2667 crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED, 2668 crate::backup::TEST_CAPTURE_PHASE_POST_COPY, 2669 crate::backup::TEST_CAPTURE_PHASE_PRE_FINAL_SYNC, 2670 ] 2671 .into_iter() 2672 .enumerate() 2673 { 2674 crate::backup::test_capture_reset(); 2675 let (root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 2676 let host = Arc::new(host); 2677 if phase == crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED { 2678 host.transaction(|transaction| { 2679 Box::pin(async move { 2680 sqlx::raw_sql( 2681 "WITH RECURSIVE sequence(value) AS ( 2682 VALUES(1) 2683 UNION ALL 2684 SELECT value + 1 FROM sequence WHERE value < 100000 2685 ) 2686 INSERT INTO host_probe(value) SELECT 0 FROM sequence", 2687 ) 2688 .execute(&mut *transaction) 2689 .await 2690 .map(|_| ()) 2691 }) 2692 }) 2693 .await 2694 .expect("seed multi-batch cancellation fixture"); 2695 } 2696 2697 let output = root.path().join(format!("phase-cancel-output-{index}")); 2698 fs::create_dir(&output).expect("backup parent"); 2699 fs::set_permissions(&output, fs::Permissions::from_mode(0o700)) 2700 .expect("backup parent mode"); 2701 let stage = output.join("cancelled"); 2702 crate::backup::test_capture_block_phase(phase); 2703 let capture = tokio::spawn({ 2704 let host = Arc::clone(&host); 2705 let stage = stage.clone(); 2706 async move { 2707 host.capture_online_backup( 2708 &stage, 2709 crate::BackupCreatedAtUnixMs::new(1_700_000_000_220 + index as u64) 2710 .expect("backup creation time"), 2711 ) 2712 .await 2713 } 2714 }); 2715 while crate::backup::test_capture_phase() != phase { 2716 tokio::task::yield_now().await; 2717 } 2718 2719 capture.abort(); 2720 assert!( 2721 capture 2722 .await 2723 .expect_err("capture task is cancelled") 2724 .is_cancelled() 2725 ); 2726 host.close() 2727 .await 2728 .expect("close drains phase-cancelled capture cleanup"); 2729 crate::backup::test_capture_reset(); 2730 2731 assert!(!stage.exists(), "phase {phase} must leave no staging tree"); 2732 assert!(!host.backup_active.load(Ordering::Acquire)); 2733 let mut authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 2734 .expect("writer authority can be reacquired") 2735 .expect("writable mode returns authority"); 2736 authority.release().expect("release reacquired authority"); 2737 } 2738 } 2739 2740 #[cfg(any(target_os = "linux", target_os = "macos"))] 2741 #[tokio::test] 2742 async fn online_backup_remains_consistent_with_a_concurrent_wal_writer() { 2743 let _serial = CAPTURE_TEST_LOCK.lock().await; 2744 crate::backup::test_capture_reset(); 2745 let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 2746 let host = Arc::new(host); 2747 host.transaction(|transaction| { 2748 Box::pin(async move { 2749 sqlx::raw_sql( 2750 "WITH RECURSIVE sequence(value) AS ( 2751 VALUES(1) UNION ALL SELECT value + 1 FROM sequence WHERE value < 100000 2752 ) 2753 INSERT INTO host_probe(value) SELECT 0 FROM sequence", 2754 ) 2755 .execute(&mut *transaction) 2756 .await 2757 .map(|_| ()) 2758 }) 2759 }) 2760 .await 2761 .expect("seed a multi-batch database"); 2762 2763 let output = root.path().join("concurrent-output"); 2764 fs::create_dir(&output).expect("backup parent"); 2765 fs::set_permissions(&output, fs::Permissions::from_mode(0o700)) 2766 .expect("backup parent mode"); 2767 let stage = output.join("consistent"); 2768 crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED); 2769 let capture = tokio::spawn({ 2770 let host = Arc::clone(&host); 2771 let stage = stage.clone(); 2772 async move { 2773 host.capture_online_backup( 2774 &stage, 2775 crate::BackupCreatedAtUnixMs::new(1_700_000_000_300) 2776 .expect("backup creation time"), 2777 ) 2778 .await 2779 } 2780 }); 2781 while crate::backup::test_capture_phase() 2782 != crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED 2783 { 2784 tokio::task::yield_now().await; 2785 } 2786 host.transaction(|transaction| { 2787 Box::pin(async move { 2788 sqlx::query("UPDATE host_probe SET value = 1 WHERE rowid IN (1, 2)") 2789 .execute(&mut *transaction) 2790 .await 2791 .map(|_| ()) 2792 }) 2793 }) 2794 .await 2795 .expect("commit concurrent WAL transaction"); 2796 crate::backup::test_capture_block_phase(0); 2797 capture 2798 .await 2799 .expect("capture joins") 2800 .expect("capture remains consistent"); 2801 crate::backup::test_capture_reset(); 2802 2803 let mut backup = SqliteConnection::connect_with( 2804 &SqliteConnectOptions::new() 2805 .filename(stage.join("state.sqlite")) 2806 .read_only(true), 2807 ) 2808 .await 2809 .expect("open backup"); 2810 let updated = 2811 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe WHERE value = 1") 2812 .fetch_one(&mut backup) 2813 .await 2814 .expect("count transaction projection"); 2815 assert!( 2816 updated == 0 || updated == 2, 2817 "backup must not tear a transaction" 2818 ); 2819 assert_eq!( 2820 sqlx::query_scalar::<_, String>("PRAGMA integrity_check(1)") 2821 .fetch_one(&mut backup) 2822 .await 2823 .expect("backup integrity"), 2824 "ok" 2825 ); 2826 backup.close().await.expect("close backup"); 2827 host.close().await.expect("close writer host"); 2828 } 2829 2830 #[cfg(any(target_os = "linux", target_os = "macos"))] 2831 #[tokio::test] 2832 async fn backup_storage_full_sync_failures_cleanup_and_leave_host_recoverable() { 2833 let _serial = CAPTURE_TEST_LOCK.lock().await; 2834 crate::backup::test_capture_reset(); 2835 let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 2836 let output = root.path().join("sync-failure-output"); 2837 fs::create_dir(&output).expect("backup parent"); 2838 fs::set_permissions(&output, fs::Permissions::from_mode(0o700)) 2839 .expect("backup parent mode"); 2840 2841 for (index, failure) in [ 2842 crate::backup::TestCaptureSyncFailure::State, 2843 crate::backup::TestCaptureSyncFailure::Staging, 2844 crate::backup::TestCaptureSyncFailure::FinalParent, 2845 ] 2846 .into_iter() 2847 .enumerate() 2848 { 2849 let stage = output.join(format!("failure-{index}")); 2850 let error = crate::backup::test_capture_online_backup_with_sync_failure( 2851 &host.pool, 2852 &host.closing, 2853 &host.backup_active, 2854 &stage, 2855 crate::BackupCreatedAtUnixMs::new(1_700_000_000_400 + index as u64) 2856 .expect("backup creation time"), 2857 failure, 2858 ) 2859 .await 2860 .expect_err("injected synchronization failure must reject capture"); 2861 assert_eq!(error.kind(), ServiceSqliteErrorKind::Backup); 2862 let storage = error 2863 .source() 2864 .and_then(Error::source) 2865 .and_then(|source| source.downcast_ref::<std::io::Error>()) 2866 .expect("storage-full cause"); 2867 assert_eq!(storage.kind(), std::io::ErrorKind::StorageFull); 2868 assert!(!stage.exists()); 2869 assert!(!host.backup_active.load(Ordering::Acquire)); 2870 assert_eq!(row_count(&host).await, 0); 2871 } 2872 2873 let recovered = output.join("recovered"); 2874 host.capture_online_backup( 2875 &recovered, 2876 crate::BackupCreatedAtUnixMs::new(1_700_000_000_410).expect("backup creation time"), 2877 ) 2878 .await 2879 .expect("host remains usable after sync failures"); 2880 assert!(recovered.join("state.sqlite").exists()); 2881 host.close().await.expect("close host"); 2882 } 2883 2884 #[cfg(any(target_os = "linux", target_os = "macos"))] 2885 #[tokio::test] 2886 async fn backup_await_boundaries_preserve_authority_precedence() { 2887 let _serial = CAPTURE_TEST_LOCK.lock().await; 2888 crate::backup::test_capture_reset(); 2889 let (root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 2890 let host = Arc::new(host); 2891 let output = root.path().join("precedence-output"); 2892 fs::create_dir(&output).expect("backup parent"); 2893 fs::set_permissions(&output, fs::Permissions::from_mode(0o700)) 2894 .expect("backup parent mode"); 2895 2896 crate::backup::test_capture_inject_metadata_failure(true); 2897 crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_METADATA_AWAITED); 2898 let metadata_stage = output.join("metadata"); 2899 let capture = tokio::spawn({ 2900 let host = Arc::clone(&host); 2901 let stage = metadata_stage.clone(); 2902 async move { 2903 host.capture_online_backup( 2904 &stage, 2905 crate::BackupCreatedAtUnixMs::new(1_700_000_000_420) 2906 .expect("backup creation time"), 2907 ) 2908 .await 2909 } 2910 }); 2911 while crate::backup::test_capture_phase() 2912 != crate::backup::TEST_CAPTURE_PHASE_METADATA_AWAITED 2913 { 2914 tokio::task::yield_now().await; 2915 } 2916 let retired_lock = paths 2917 .state_lock() 2918 .parent() 2919 .expect("state directory") 2920 .join("retired-backup-precedence.lock"); 2921 fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock"); 2922 crate::backup::test_capture_block_phase(0); 2923 let metadata_error = capture 2924 .await 2925 .expect("metadata capture joins") 2926 .expect_err("authority overrides metadata failure"); 2927 assert_eq!(metadata_error.kind(), ServiceSqliteErrorKind::Authority); 2928 fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock"); 2929 crate::backup::test_capture_reset(); 2930 assert!(!metadata_stage.exists()); 2931 assert_eq!(row_count(&host).await, 0); 2932 2933 crate::backup::test_capture_panic_worker(true); 2934 crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_JOIN_AWAITED); 2935 let join_stage = output.join("join"); 2936 let capture = tokio::spawn({ 2937 let host = Arc::clone(&host); 2938 let stage = join_stage.clone(); 2939 async move { 2940 host.capture_online_backup( 2941 &stage, 2942 crate::BackupCreatedAtUnixMs::new(1_700_000_000_421) 2943 .expect("backup creation time"), 2944 ) 2945 .await 2946 } 2947 }); 2948 while crate::backup::test_capture_phase() != crate::backup::TEST_CAPTURE_PHASE_JOIN_AWAITED 2949 { 2950 tokio::task::yield_now().await; 2951 } 2952 fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock again"); 2953 crate::backup::test_capture_block_phase(0); 2954 let join_error = capture 2955 .await 2956 .expect("join-failure capture joins") 2957 .expect_err("authority overrides worker join failure"); 2958 assert_eq!(join_error.kind(), ServiceSqliteErrorKind::Authority); 2959 fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock again"); 2960 crate::backup::test_capture_reset(); 2961 assert!(!join_stage.exists()); 2962 assert_eq!(row_count(&host).await, 0); 2963 host.close().await.expect("close host"); 2964 } 2965 2966 #[cfg(any(target_os = "linux", target_os = "macos"))] 2967 #[derive(Debug, PartialEq, Eq)] 2968 struct StateFileSnapshot { 2969 bytes: Vec<u8>, 2970 length: u64, 2971 modified: SystemTime, 2972 mode: u32, 2973 } 2974 2975 #[cfg(any(target_os = "linux", target_os = "macos"))] 2976 fn state_directory_snapshot(directory: &Path) -> BTreeMap<String, StateFileSnapshot> { 2977 fs::read_dir(directory) 2978 .expect("read state directory") 2979 .map(|entry| { 2980 let entry = entry.expect("state entry"); 2981 let name = entry.file_name().into_string().expect("UTF-8 state entry"); 2982 let metadata = entry.metadata().expect("state entry metadata"); 2983 ( 2984 name, 2985 StateFileSnapshot { 2986 bytes: fs::read(entry.path()).expect("state entry bytes"), 2987 length: metadata.len(), 2988 modified: metadata.modified().expect("state entry mtime"), 2989 mode: metadata.permissions().mode() & 0o777, 2990 }, 2991 ) 2992 }) 2993 .collect() 2994 } 2995 2996 #[cfg(any(target_os = "linux", target_os = "macos"))] 2997 #[tokio::test] 2998 async fn close_drains_admitted_work_rejects_new_work_and_is_idempotent() { 2999 let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 3000 let host = Arc::new(host); 3001 let entered = Arc::new(Notify::new()); 3002 let release = Arc::new(Notify::new()); 3003 let transaction = tokio::spawn({ 3004 let host = Arc::clone(&host); 3005 let entered = Arc::clone(&entered); 3006 let release = Arc::clone(&release); 3007 async move { 3008 host.transaction::<i64, Infallible, _>(|transaction| { 3009 Box::pin(async move { 3010 let count = sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") 3011 .fetch_one(&mut *transaction) 3012 .await 3013 .expect("read in admitted transaction"); 3014 entered.notify_one(); 3015 release.notified().await; 3016 Ok(count) 3017 }) 3018 }) 3019 .await 3020 } 3021 }); 3022 entered.notified().await; 3023 3024 let first_close = tokio::spawn({ 3025 let host = Arc::clone(&host); 3026 async move { host.close().await } 3027 }); 3028 while !host.closing.load(Ordering::Acquire) { 3029 tokio::task::yield_now().await; 3030 } 3031 let second_close = tokio::spawn({ 3032 let host = Arc::clone(&host); 3033 async move { host.close().await } 3034 }); 3035 let rejected = host 3036 .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) 3037 .await 3038 .expect_err("close admission must reject new work"); 3039 assert_eq!( 3040 rejected.kind(), 3041 ServiceSqliteTransactionErrorKind::NotCommitted 3042 ); 3043 assert_eq!( 3044 rejected.sqlite_error().map(ServiceSqliteError::kind), 3045 Some(ServiceSqliteErrorKind::Open) 3046 ); 3047 let contended = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting); 3048 assert!(matches!( 3049 contended, 3050 Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority 3051 )); 3052 3053 release.notify_one(); 3054 assert_eq!( 3055 transaction 3056 .await 3057 .expect("transaction task joins") 3058 .expect("admitted transaction"), 3059 0 3060 ); 3061 first_close 3062 .await 3063 .expect("first close task joins") 3064 .expect("first close succeeds"); 3065 second_close 3066 .await 3067 .expect("second close task joins") 3068 .expect("concurrent close is idempotent"); 3069 host.close().await.expect("sequential close is idempotent"); 3070 3071 let mut next = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 3072 .expect("authority can be reacquired") 3073 .expect("writer mode yields authority"); 3074 next.release().expect("release reacquired authority"); 3075 } 3076 3077 #[cfg(any(target_os = "linux", target_os = "macos"))] 3078 #[tokio::test] 3079 async fn cancelled_close_retains_authority_and_retry_finishes() { 3080 let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 3081 let host = Arc::new(host); 3082 let entered = Arc::new(Notify::new()); 3083 let release = Arc::new(Notify::new()); 3084 let transaction = tokio::spawn({ 3085 let host = Arc::clone(&host); 3086 let entered = Arc::clone(&entered); 3087 let release = Arc::clone(&release); 3088 async move { 3089 host.transaction::<(), Infallible, _>(|transaction| { 3090 Box::pin(async move { 3091 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") 3092 .fetch_one(&mut *transaction) 3093 .await 3094 .expect("admitted read"); 3095 entered.notify_one(); 3096 release.notified().await; 3097 Ok(()) 3098 }) 3099 }) 3100 .await 3101 } 3102 }); 3103 entered.notified().await; 3104 let close_task = tokio::spawn({ 3105 let host = Arc::clone(&host); 3106 async move { host.close().await } 3107 }); 3108 while !host.closing.load(Ordering::Acquire) { 3109 tokio::task::yield_now().await; 3110 } 3111 close_task.abort(); 3112 assert!( 3113 close_task 3114 .await 3115 .expect_err("close task is cancelled") 3116 .is_cancelled() 3117 ); 3118 let retained = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting); 3119 assert!(matches!( 3120 retained, 3121 Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority 3122 )); 3123 let rejected = host 3124 .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) 3125 .await 3126 .expect_err("cancelled close remains non-admitting"); 3127 assert_eq!( 3128 rejected.kind(), 3129 ServiceSqliteTransactionErrorKind::NotCommitted 3130 ); 3131 3132 release.notify_one(); 3133 transaction 3134 .await 3135 .expect("transaction task joins") 3136 .expect("admitted transaction finishes"); 3137 host.close().await.expect("close retry succeeds"); 3138 assert!( 3139 WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 3140 .expect("authority reacquisition after retry") 3141 .is_some() 3142 ); 3143 } 3144 3145 #[cfg(any(target_os = "linux", target_os = "macos"))] 3146 #[tokio::test] 3147 async fn writable_close_checkpoints_and_read_only_close_is_side_effect_free() { 3148 let (_root, paths, identity, migrations, schema, host) = initialized_host().await; 3149 host.transaction(|transaction| { 3150 Box::pin(async move { 3151 sqlx::query("INSERT INTO host_probe (value) VALUES (41)") 3152 .execute(&mut *transaction) 3153 .await 3154 .map(|_| ()) 3155 }) 3156 }) 3157 .await 3158 .expect("write WAL frame"); 3159 host.close().await.expect("writable close checkpoints"); 3160 let state_directory = paths.state_database().parent().expect("state directory"); 3161 assert!(!state_directory.join("state.sqlite-wal").exists()); 3162 assert!(!state_directory.join("state.sqlite-shm").exists()); 3163 let before = state_directory_snapshot(state_directory); 3164 3165 let inspection = ServiceSqliteHost::open_read_only_inspection( 3166 &paths, 3167 &identity, 3168 &migrations, 3169 &schema, 3170 ServiceSqliteConnectionOptions::reviewed(), 3171 ) 3172 .await 3173 .expect("open read-only inspection"); 3174 assert_eq!(row_count(&inspection).await, 1); 3175 inspection 3176 .close() 3177 .await 3178 .expect("close read-only inspection"); 3179 assert_eq!(state_directory_snapshot(state_directory), before); 3180 3181 let (reopened, outcome) = ServiceSqliteHost::open_read_write_existing( 3182 &paths, 3183 &identity, 3184 &migrations, 3185 &schema, 3186 ServiceSqliteConnectionOptions::reviewed(), 3187 MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"), 3188 &build_identity(), 3189 &[], 3190 ) 3191 .await 3192 .expect("reopen writer after read-only close"); 3193 assert_eq!(outcome.applied_count(), 0); 3194 assert_eq!(row_count(&reopened).await, 1); 3195 reopened.close().await.expect("close reopened writer"); 3196 } 3197 3198 #[cfg(any(target_os = "linux", target_os = "macos"))] 3199 #[tokio::test] 3200 async fn read_only_close_revalidates_after_drain_before_releasing_stale_authority() { 3201 let (_root, paths, identity, migrations, schema, writer) = initialized_host().await; 3202 writer.close().await.expect("close writer host"); 3203 let inspection = Arc::new( 3204 ServiceSqliteHost::open_read_only_inspection( 3205 &paths, 3206 &identity, 3207 &migrations, 3208 &schema, 3209 ServiceSqliteConnectionOptions::reviewed(), 3210 ) 3211 .await 3212 .expect("open read-only inspection"), 3213 ); 3214 let entered = Arc::new(Notify::new()); 3215 let release = Arc::new(Notify::new()); 3216 let transaction = tokio::spawn({ 3217 let inspection = Arc::clone(&inspection); 3218 let entered = Arc::clone(&entered); 3219 let release = Arc::clone(&release); 3220 async move { 3221 inspection 3222 .transaction::<(), Infallible, _>(|transaction| { 3223 Box::pin(async move { 3224 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") 3225 .fetch_one(&mut *transaction) 3226 .await 3227 .expect("read through retained inspection"); 3228 entered.notify_one(); 3229 release.notified().await; 3230 Ok(()) 3231 }) 3232 }) 3233 .await 3234 } 3235 }); 3236 entered.notified().await; 3237 let close_task = tokio::spawn({ 3238 let inspection = Arc::clone(&inspection); 3239 async move { inspection.close().await } 3240 }); 3241 while !inspection.closing.load(Ordering::Acquire) { 3242 tokio::task::yield_now().await; 3243 } 3244 3245 let retired_lock = paths 3246 .state_lock() 3247 .parent() 3248 .expect("state directory") 3249 .join("retired-inspection-close.lock"); 3250 fs::rename(paths.state_lock(), retired_lock).expect("replace inspection lock"); 3251 let mut replacement = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 3252 .expect("replacement acquisition") 3253 .expect("replacement authority"); 3254 release.notify_one(); 3255 let transaction_error = transaction 3256 .await 3257 .expect("inspection transaction task joins") 3258 .expect_err("binding drift revokes admitted inspection"); 3259 assert_eq!( 3260 transaction_error.kind(), 3261 ServiceSqliteTransactionErrorKind::NotCommitted 3262 ); 3263 let close_error = close_task 3264 .await 3265 .expect("inspection close task joins") 3266 .expect_err("close must report stale inspection authority"); 3267 assert_eq!(close_error.kind(), ServiceSqliteErrorKind::Authority); 3268 assert_eq!( 3269 inspection 3270 .close() 3271 .await 3272 .expect_err("terminal authority result is cached") 3273 .kind(), 3274 ServiceSqliteErrorKind::Authority 3275 ); 3276 replacement 3277 .release() 3278 .expect("release replacement authority"); 3279 } 3280 3281 #[cfg(any(target_os = "linux", target_os = "macos"))] 3282 #[tokio::test] 3283 async fn cancelled_checkpoint_resumes_to_terminal_error_and_releases_authority() { 3284 let options = ServiceSqliteConnectionOptions::new(Duration::from_secs(1), 1) 3285 .expect("short reviewed limits"); 3286 let (_root, paths, _identity, _migrations, _schema, host) = 3287 initialized_host_with_options(options).await; 3288 let host = Arc::new(host); 3289 host.transaction(|transaction| { 3290 Box::pin(async move { 3291 sqlx::query("INSERT INTO host_probe (value) VALUES (1)") 3292 .execute(&mut *transaction) 3293 .await 3294 .map(|_| ()) 3295 }) 3296 }) 3297 .await 3298 .expect("seed reader snapshot"); 3299 3300 let mut reader = SqliteConnection::connect_with( 3301 &SqliteConnectOptions::new() 3302 .filename(paths.state_database()) 3303 .read_only(true) 3304 .create_if_missing(false), 3305 ) 3306 .await 3307 .expect("open external reader"); 3308 let mut reader_transaction = reader.begin().await.expect("begin external read"); 3309 assert_eq!( 3310 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") 3311 .fetch_one(&mut *reader_transaction) 3312 .await 3313 .expect("establish reader snapshot"), 3314 1 3315 ); 3316 host.transaction(|transaction| { 3317 Box::pin(async move { 3318 sqlx::query("INSERT INTO host_probe (value) VALUES (2)") 3319 .execute(&mut *transaction) 3320 .await 3321 .map(|_| ()) 3322 }) 3323 }) 3324 .await 3325 .expect("append frame after reader snapshot"); 3326 3327 let close_task = tokio::spawn({ 3328 let host = Arc::clone(&host); 3329 async move { host.close().await } 3330 }); 3331 while host.pool.close_phase() != crate::open::TEST_CLOSE_PHASE_CHECKPOINT { 3332 tokio::task::yield_now().await; 3333 } 3334 close_task.abort(); 3335 assert!( 3336 close_task 3337 .await 3338 .expect_err("checkpoint close task is cancelled") 3339 .is_cancelled() 3340 ); 3341 let retained = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting); 3342 assert!(matches!( 3343 retained, 3344 Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority 3345 )); 3346 3347 let first = host 3348 .close() 3349 .await 3350 .expect_err("active reader prevents TRUNCATE"); 3351 assert_eq!(first.kind(), ServiceSqliteErrorKind::Pragma); 3352 assert!(!first.to_string().contains("state.sqlite")); 3353 let repeated = host.close().await.expect_err("terminal result is cached"); 3354 assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Pragma); 3355 let mut replacement = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 3356 .expect("close releases authority despite checkpoint failure") 3357 .expect("writer mode yields authority"); 3358 replacement 3359 .release() 3360 .expect("release replacement authority"); 3361 3362 reader_transaction 3363 .rollback() 3364 .await 3365 .expect("release external snapshot"); 3366 reader.close().await.expect("close external reader"); 3367 } 3368 3369 #[cfg(any(target_os = "linux", target_os = "macos"))] 3370 #[tokio::test] 3371 async fn close_authority_drift_is_cached_and_releases_the_stale_lock() { 3372 let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 3373 let retired_lock = paths 3374 .state_lock() 3375 .parent() 3376 .expect("state directory") 3377 .join("retired-close-state.lock"); 3378 fs::rename(paths.state_lock(), retired_lock).expect("replace canonical lock"); 3379 let mut replacement = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 3380 .expect("replacement acquisition") 3381 .expect("replacement authority"); 3382 3383 let first = host 3384 .close() 3385 .await 3386 .expect_err("close detects authority drift"); 3387 assert_eq!(first.kind(), ServiceSqliteErrorKind::Authority); 3388 let repeated = host.close().await.expect_err("authority result is cached"); 3389 assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Authority); 3390 assert!(!format!("{host:?}").contains(paths.state_database().to_string_lossy().as_ref())); 3391 replacement 3392 .release() 3393 .expect("release replacement authority"); 3394 } 3395 3396 #[cfg(any(target_os = "linux", target_os = "macos"))] 3397 #[tokio::test] 3398 async fn typed_execution_commits_and_operation_failure_rolls_back() { 3399 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3400 assert_eq!(host.mode(), OpenMode::Initialize); 3401 3402 let inserted = host 3403 .transaction(|transaction| { 3404 Box::pin(async move { 3405 sqlx::query("INSERT INTO host_probe (value) VALUES (?)") 3406 .bind(41_i64) 3407 .execute(&mut *transaction) 3408 .await?; 3409 sqlx::query_scalar::<_, i64>("SELECT value FROM host_probe") 3410 .fetch_one(&mut *transaction) 3411 .await 3412 }) 3413 }) 3414 .await 3415 .expect("commit typed operation"); 3416 assert_eq!(inserted, 41); 3417 3418 let error = host 3419 .transaction(|transaction| { 3420 Box::pin(async move { 3421 sqlx::query("INSERT INTO host_probe (value) VALUES (99)") 3422 .execute(&mut *transaction) 3423 .await 3424 .map_err(|_| "query-failure-secret")?; 3425 Err::<(), _>("operation-secret") 3426 }) 3427 }) 3428 .await 3429 .expect_err("operation rejection must roll back"); 3430 assert_eq!( 3431 error.kind(), 3432 ServiceSqliteTransactionErrorKind::OperationRolledBack 3433 ); 3434 assert_eq!(error.operation_error(), Some(&"operation-secret")); 3435 assert!(error.sqlite_error().is_none()); 3436 assert!(!format!("{error:?}").contains("operation-secret")); 3437 assert!(!error.to_string().contains("operation-secret")); 3438 assert_eq!(row_count(&host).await, 1); 3439 } 3440 3441 #[cfg(any(target_os = "linux", target_os = "macos"))] 3442 #[tokio::test] 3443 async fn transaction_durability_edges_preserve_exact_commit_semantics() { 3444 use crate::failpoint::DurabilityFailpoint; 3445 3446 for (point, expected_reached) in [ 3447 ( 3448 DurabilityFailpoint::TransactionBeforeBegin, 3449 &[DurabilityFailpoint::TransactionBeforeBegin][..], 3450 ), 3451 ( 3452 DurabilityFailpoint::TransactionAfterBegin, 3453 &[ 3454 DurabilityFailpoint::TransactionBeforeBegin, 3455 DurabilityFailpoint::TransactionAfterBegin, 3456 ][..], 3457 ), 3458 ( 3459 DurabilityFailpoint::TransactionBeforeCommit, 3460 &[ 3461 DurabilityFailpoint::TransactionBeforeBegin, 3462 DurabilityFailpoint::TransactionAfterBegin, 3463 DurabilityFailpoint::TransactionBeforeCommit, 3464 ][..], 3465 ), 3466 ( 3467 DurabilityFailpoint::TransactionAfterCommit, 3468 &[ 3469 DurabilityFailpoint::TransactionBeforeBegin, 3470 DurabilityFailpoint::TransactionAfterBegin, 3471 DurabilityFailpoint::TransactionBeforeCommit, 3472 DurabilityFailpoint::TransactionAfterCommit, 3473 ][..], 3474 ), 3475 ] { 3476 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3477 host.arm_durability_failpoint(point); 3478 let error = host 3479 .transaction(|transaction| { 3480 Box::pin(async move { 3481 sqlx::query("INSERT INTO host_probe (value) VALUES (72)") 3482 .execute(&mut *transaction) 3483 .await 3484 .map(|_| ()) 3485 }) 3486 }) 3487 .await 3488 .expect_err("injected transaction edge"); 3489 let reached = host.failpoints.reached(); 3490 if point == DurabilityFailpoint::TransactionAfterCommit { 3491 assert_eq!( 3492 error.kind(), 3493 ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown 3494 ); 3495 assert_eq!(row_count(&host).await, 1); 3496 } else { 3497 assert_eq!( 3498 error.kind(), 3499 ServiceSqliteTransactionErrorKind::NotCommitted 3500 ); 3501 assert_eq!(row_count(&host).await, 0); 3502 } 3503 assert!(host.failpoints.fired()); 3504 assert_eq!(reached, expected_reached); 3505 host.close() 3506 .await 3507 .expect("close host after transaction edge"); 3508 } 3509 3510 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3511 host.arm_durability_failpoint(DurabilityFailpoint::TransactionBeforeCommit); 3512 let error = host 3513 .transaction(|transaction| { 3514 Box::pin(async move { 3515 let _ = sqlx::raw_sql( 3516 "PRAGMA trusted_schema=ON; INSERT INTO host_probe (value) VALUES (73)", 3517 ) 3518 .execute(&mut *transaction) 3519 .await; 3520 Ok::<_, Infallible>(()) 3521 }) 3522 }) 3523 .await 3524 .expect_err("precommit policy rejection must precede the commit-edge hook"); 3525 assert_eq!( 3526 error.kind(), 3527 ServiceSqliteTransactionErrorKind::NotCommitted 3528 ); 3529 assert!(!host.failpoints.fired()); 3530 assert_eq!( 3531 host.failpoints.reached(), 3532 [ 3533 DurabilityFailpoint::TransactionBeforeBegin, 3534 DurabilityFailpoint::TransactionAfterBegin, 3535 ] 3536 ); 3537 host.failpoints.disarm(); 3538 assert_eq!(row_count(&host).await, 0); 3539 host.close().await.expect("close policy-drift host"); 3540 } 3541 3542 #[cfg(any(target_os = "linux", target_os = "macos"))] 3543 #[tokio::test] 3544 async fn backup_durability_edges_fail_once_clean_exact_stage_and_recover() { 3545 use crate::failpoint::DurabilityFailpoint; 3546 3547 let _serial = CAPTURE_TEST_LOCK.lock().await; 3548 for (index, (point, expected_phase)) in [ 3549 (DurabilityFailpoint::BackupBeforeCreate, 1), 3550 (DurabilityFailpoint::BackupAfterCreate, 1), 3551 (DurabilityFailpoint::BackupBeforeCopy, 2), 3552 (DurabilityFailpoint::BackupAfterCopy, 3), 3553 (DurabilityFailpoint::BackupBeforeFileSync, 4), 3554 (DurabilityFailpoint::BackupAfterFileSync, 4), 3555 (DurabilityFailpoint::BackupBeforeDirectorySync, 5), 3556 (DurabilityFailpoint::BackupAfterDirectorySync, 5), 3557 ] 3558 .into_iter() 3559 .enumerate() 3560 { 3561 let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3562 let stage = root.path().join(format!("failpoint-backup-{index}")); 3563 host.arm_durability_failpoint(point); 3564 let error = host 3565 .capture_online_backup( 3566 &stage, 3567 crate::BackupCreatedAtUnixMs::new(1_700_000_072_000).expect("capture time"), 3568 ) 3569 .await 3570 .expect_err("injected backup edge"); 3571 assert_eq!(error.kind(), ServiceSqliteErrorKind::Backup); 3572 assert!(host.failpoints.fired()); 3573 assert_eq!(host.failpoints.reached().last(), Some(&point)); 3574 assert_eq!( 3575 host.failpoints.observation(point), 3576 Some(expected_phase), 3577 "backup failpoint must fire in its named lifecycle phase" 3578 ); 3579 assert!(!stage.exists(), "owned failed stage is cleaned"); 3580 let recovery = root.path().join(format!("recovered-backup-{index}")); 3581 host.capture_online_backup( 3582 &recovery, 3583 crate::BackupCreatedAtUnixMs::new(1_700_000_072_001).expect("capture time"), 3584 ) 3585 .await 3586 .expect("one-shot failpoint permits retry"); 3587 assert!(recovery.join(crate::BACKUP_STATE_MEMBER_NAME).is_file()); 3588 host.close().await.expect("close backup host"); 3589 } 3590 } 3591 3592 #[cfg(any(target_os = "linux", target_os = "macos"))] 3593 #[tokio::test] 3594 async fn close_durability_edges_are_once_only_retryable_or_terminal() { 3595 use crate::failpoint::DurabilityFailpoint; 3596 3597 let all_close_edges = [ 3598 DurabilityFailpoint::CloseBeforeDrain, 3599 DurabilityFailpoint::CloseAfterDrain, 3600 DurabilityFailpoint::CloseBeforeCheckpoint, 3601 DurabilityFailpoint::CloseAfterCheckpoint, 3602 DurabilityFailpoint::CloseBeforeConnectionClose, 3603 DurabilityFailpoint::CloseAfterConnectionClose, 3604 DurabilityFailpoint::CloseBeforeAuthorityRelease, 3605 DurabilityFailpoint::CloseAfterAuthorityRelease, 3606 ]; 3607 for (point, expected, retryable, expected_phase, expected_reached) in [ 3608 ( 3609 DurabilityFailpoint::CloseBeforeDrain, 3610 ServiceSqliteErrorKind::Open, 3611 true, 3612 0, 3613 &all_close_edges[..1], 3614 ), 3615 ( 3616 DurabilityFailpoint::CloseAfterDrain, 3617 ServiceSqliteErrorKind::Open, 3618 true, 3619 0, 3620 &all_close_edges[..2], 3621 ), 3622 ( 3623 DurabilityFailpoint::CloseBeforeCheckpoint, 3624 ServiceSqliteErrorKind::Pragma, 3625 false, 3626 6, 3627 &[ 3628 DurabilityFailpoint::CloseBeforeDrain, 3629 DurabilityFailpoint::CloseAfterDrain, 3630 DurabilityFailpoint::CloseBeforeCheckpoint, 3631 DurabilityFailpoint::CloseBeforeConnectionClose, 3632 DurabilityFailpoint::CloseAfterConnectionClose, 3633 DurabilityFailpoint::CloseBeforeAuthorityRelease, 3634 DurabilityFailpoint::CloseAfterAuthorityRelease, 3635 ][..], 3636 ), 3637 ( 3638 DurabilityFailpoint::CloseAfterCheckpoint, 3639 ServiceSqliteErrorKind::Pragma, 3640 false, 3641 6, 3642 &all_close_edges[..], 3643 ), 3644 ( 3645 DurabilityFailpoint::CloseBeforeConnectionClose, 3646 ServiceSqliteErrorKind::Open, 3647 false, 3648 6, 3649 &all_close_edges[..], 3650 ), 3651 ( 3652 DurabilityFailpoint::CloseAfterConnectionClose, 3653 ServiceSqliteErrorKind::Open, 3654 false, 3655 6, 3656 &all_close_edges[..], 3657 ), 3658 ( 3659 DurabilityFailpoint::CloseBeforeAuthorityRelease, 3660 ServiceSqliteErrorKind::Authority, 3661 true, 3662 6, 3663 &all_close_edges[..7], 3664 ), 3665 ( 3666 DurabilityFailpoint::CloseAfterAuthorityRelease, 3667 ServiceSqliteErrorKind::Authority, 3668 false, 3669 6, 3670 &all_close_edges[..], 3671 ), 3672 ] { 3673 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3674 host.arm_durability_failpoint(point); 3675 let error = host.close().await.expect_err("injected close edge"); 3676 assert_eq!(error.kind(), expected); 3677 assert!(host.failpoints.fired()); 3678 assert_eq!(host.failpoints.reached(), expected_reached); 3679 assert_eq!( 3680 host.pool.close_phase(), 3681 expected_phase, 3682 "close failpoint must fire in its named driver phase" 3683 ); 3684 if retryable { 3685 host.close().await.expect("one-shot close edge resumes"); 3686 } else { 3687 assert_eq!( 3688 host.close() 3689 .await 3690 .expect_err("terminal close result is cached") 3691 .kind(), 3692 expected 3693 ); 3694 } 3695 } 3696 } 3697 3698 #[cfg(any(target_os = "linux", target_os = "macos"))] 3699 #[tokio::test] 3700 async fn cancellation_quarantines_connection_and_pool_recovers() { 3701 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3702 let host = Arc::new(host); 3703 let entered = Arc::new(AtomicBool::new(false)); 3704 let task = tokio::spawn({ 3705 let host = Arc::clone(&host); 3706 let entered = Arc::clone(&entered); 3707 async move { 3708 host.transaction::<(), Infallible, _>(|transaction| { 3709 Box::pin(async move { 3710 sqlx::query("INSERT INTO host_probe (value) VALUES (77)") 3711 .execute(&mut *transaction) 3712 .await 3713 .expect("tentative insert"); 3714 entered.store(true, Ordering::Release); 3715 std::future::pending::<()>().await; 3716 Ok(()) 3717 }) 3718 }) 3719 .await 3720 } 3721 }); 3722 while !entered.load(Ordering::Acquire) { 3723 tokio::task::yield_now().await; 3724 } 3725 task.abort(); 3726 assert!( 3727 task.await 3728 .expect_err("task must be cancelled") 3729 .is_cancelled() 3730 ); 3731 assert_eq!(row_count(&host).await, 0); 3732 } 3733 3734 #[cfg(any(target_os = "linux", target_os = "macos"))] 3735 #[tokio::test] 3736 async fn complete_statement_control_inventory_is_sticky_and_fails_closed() { 3737 for statement in [ 3738 "INSERT INTO host_probe (value) VALUES (0); /* policy */ PrAgMa\ntrusted_schema=ON", 3739 "INSERT INTO host_probe (value) VALUES (1); ATTACH DATABASE ':memory:' AS extra", 3740 "INSERT INTO host_probe (value) VALUES (2); DETACH DATABASE extra", 3741 "INSERT INTO host_probe (value) VALUES (3); BEGIN DEFERRED", 3742 "INSERT INTO host_probe (value) VALUES (4); COMMIT", 3743 "INSERT INTO host_probe (value) VALUES (5); END", 3744 "INSERT INTO host_probe (value) VALUES (6); ROLLBACK", 3745 "INSERT INTO host_probe (value) VALUES (7); SAVEPOINT escaped", 3746 "INSERT INTO host_probe (value) VALUES (8); RELEASE SAVEPOINT escaped", 3747 ] { 3748 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3749 let error = host 3750 .transaction(|transaction| { 3751 Box::pin(async move { 3752 let _ = sqlx::raw_sql(statement).execute(&mut *transaction).await; 3753 Ok::<_, Infallible>(()) 3754 }) 3755 }) 3756 .await 3757 .expect_err("escape attempt must not commit"); 3758 assert_eq!( 3759 error.kind(), 3760 ServiceSqliteTransactionErrorKind::NotCommitted 3761 ); 3762 assert!(error.sqlite_error().is_some()); 3763 assert_eq!(row_count(&host).await, 0); 3764 } 3765 3766 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3767 let error = host 3768 .transaction(|transaction| { 3769 Box::pin(async move { 3770 let _ = sqlx::raw_sql("INSERT INTO host_probe (value) VALUES (4); COMMIT") 3771 .execute(&mut *transaction) 3772 .await; 3773 let _ = 3774 sqlx::raw_sql("BEGIN DEFERRED; INSERT INTO host_probe (value) VALUES (5)") 3775 .execute(&mut *transaction) 3776 .await; 3777 Ok::<_, Infallible>(()) 3778 }) 3779 }) 3780 .await 3781 .expect_err("replacement transaction after denied COMMIT must not escape"); 3782 assert_eq!( 3783 error.kind(), 3784 ServiceSqliteTransactionErrorKind::NotCommitted 3785 ); 3786 assert_eq!(row_count(&host).await, 0); 3787 } 3788 3789 #[cfg(any(target_os = "linux", target_os = "macos"))] 3790 #[tokio::test] 3791 async fn prepared_query_policy_rejection_is_sticky_and_rolls_back_prior_work() { 3792 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3793 let error = host 3794 .transaction(|transaction| { 3795 Box::pin(async move { 3796 sqlx::query("INSERT INTO host_probe (value) VALUES (11)") 3797 .execute(&mut *transaction) 3798 .await 3799 .expect("ordinary service statement"); 3800 let _ = (&mut *transaction) 3801 .prepare(SqlStr::from_static( 3802 "SELECT 1; /* ignored */ PRAGMA trusted_schema=ON", 3803 )) 3804 .await; 3805 Ok::<_, Infallible>(()) 3806 }) 3807 }) 3808 .await 3809 .expect_err("ignored prepared-query rejection must block commit"); 3810 assert_eq!( 3811 error.kind(), 3812 ServiceSqliteTransactionErrorKind::NotCommitted 3813 ); 3814 assert_eq!(row_count(&host).await, 0); 3815 } 3816 3817 #[cfg(any(target_os = "linux", target_os = "macos"))] 3818 #[tokio::test] 3819 async fn control_words_in_values_and_case_expressions_remain_available() { 3820 let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3821 let value = host 3822 .transaction(|transaction| { 3823 Box::pin(async move { 3824 let value = sqlx::query_scalar::<_, String>( 3825 "SELECT CASE WHEN 1 = 1 THEN 'commit' ELSE 'end' END", 3826 ) 3827 .fetch_one(&mut *transaction) 3828 .await?; 3829 sqlx::query("INSERT INTO host_probe (value) VALUES (12)") 3830 .execute(&mut *transaction) 3831 .await?; 3832 Ok::<_, sqlx::Error>(value) 3833 }) 3834 }) 3835 .await 3836 .expect("ordinary expression commits"); 3837 assert_eq!(value, "commit"); 3838 assert_eq!(row_count(&host).await, 1); 3839 } 3840 3841 #[cfg(any(target_os = "linux", target_os = "macos"))] 3842 #[tokio::test] 3843 async fn attach_detach_is_rejected_before_it_can_create_external_state() { 3844 let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; 3845 let external_database = root.path().join("forbidden-attachment.sqlite"); 3846 let statement = format!( 3847 "ATTACH DATABASE '{}' AS extra; DETACH DATABASE extra", 3848 external_database.display() 3849 ); 3850 let error = host 3851 .transaction(|transaction| { 3852 Box::pin(async move { 3853 sqlx::raw_sql(sqlx::AssertSqlSafe(statement)) 3854 .execute(&mut *transaction) 3855 .await 3856 .map(|_| ()) 3857 }) 3858 }) 3859 .await 3860 .expect_err("ATTACH and DETACH must be rejected before SQLite compilation"); 3861 assert_eq!( 3862 error.kind(), 3863 ServiceSqliteTransactionErrorKind::OperationRolledBack 3864 ); 3865 assert!(!external_database.exists()); 3866 assert_eq!(row_count(&host).await, 0); 3867 } 3868 3869 #[cfg(any(target_os = "linux", target_os = "macos"))] 3870 #[tokio::test] 3871 async fn writer_lock_replacement_and_insecure_directory_revoke_live_host() { 3872 let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 3873 let retired_lock = paths 3874 .state_lock() 3875 .parent() 3876 .expect("state directory") 3877 .join("retired-state.lock"); 3878 fs::rename(paths.state_lock(), &retired_lock).expect("replace canonical lock name"); 3879 let replacement_authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) 3880 .expect("replacement authority acquisition") 3881 .expect("new writer authority"); 3882 let replaced = host 3883 .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) 3884 .await 3885 .expect_err("old writer authority must reject the replacement lock"); 3886 assert_eq!( 3887 replaced.kind(), 3888 ServiceSqliteTransactionErrorKind::NotCommitted 3889 ); 3890 assert_eq!( 3891 replaced.sqlite_error().map(ServiceSqliteError::kind), 3892 Some(ServiceSqliteErrorKind::Authority) 3893 ); 3894 drop(replacement_authority); 3895 drop(host); 3896 3897 let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; 3898 let directory = paths.state_database().parent().expect("state directory"); 3899 let original_mode = fs::metadata(directory) 3900 .expect("state directory metadata") 3901 .permissions() 3902 .mode() 3903 & 0o777; 3904 fs::set_permissions(directory, fs::Permissions::from_mode(0o770)) 3905 .expect("make directory insecure"); 3906 let insecure = host 3907 .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) 3908 .await 3909 .expect_err("insecure live directory must revoke authority"); 3910 assert_eq!( 3911 insecure.kind(), 3912 ServiceSqliteTransactionErrorKind::NotCommitted 3913 ); 3914 assert_eq!( 3915 insecure.sqlite_error().map(ServiceSqliteError::kind), 3916 Some(ServiceSqliteErrorKind::Authority) 3917 ); 3918 fs::set_permissions(directory, fs::Permissions::from_mode(original_mode)) 3919 .expect("restore directory mode"); 3920 } 3921 3922 #[cfg(any(target_os = "linux", target_os = "macos"))] 3923 #[tokio::test] 3924 async fn inspection_lock_replacement_and_new_writer_revoke_live_inspection() { 3925 let (_root, paths, identity, migrations, schema, host) = initialized_host().await; 3926 host.close().await.expect("close writer host"); 3927 drop(host); 3928 let inspection = ServiceSqliteHost::open_read_only_inspection( 3929 &paths, 3930 &identity, 3931 &migrations, 3932 &schema, 3933 ServiceSqliteConnectionOptions::reviewed(), 3934 ) 3935 .await 3936 .expect("open inspection host"); 3937 let retired_lock = paths 3938 .state_lock() 3939 .parent() 3940 .expect("state directory") 3941 .join("inspection-state.lock"); 3942 fs::rename(paths.state_lock(), retired_lock).expect("replace inspection lock name"); 3943 let (writer, _outcome) = ServiceSqliteHost::open_read_write_existing( 3944 &paths, 3945 &identity, 3946 &migrations, 3947 &schema, 3948 ServiceSqliteConnectionOptions::reviewed(), 3949 MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"), 3950 &build_identity(), 3951 &[], 3952 ) 3953 .await 3954 .expect("open replacement writer"); 3955 writer 3956 .transaction(|transaction| { 3957 Box::pin(async move { 3958 sqlx::query("INSERT INTO host_probe (value) VALUES (88)") 3959 .execute(&mut *transaction) 3960 .await 3961 .map(|_| ()) 3962 }) 3963 }) 3964 .await 3965 .expect("replacement writer commits"); 3966 let stale = inspection 3967 .transaction(|transaction| { 3968 Box::pin(async move { 3969 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") 3970 .fetch_one(&mut *transaction) 3971 .await 3972 }) 3973 }) 3974 .await 3975 .expect_err("stale inspection authority must refuse work"); 3976 assert_eq!( 3977 stale.kind(), 3978 ServiceSqliteTransactionErrorKind::NotCommitted 3979 ); 3980 assert_eq!( 3981 stale.sqlite_error().map(ServiceSqliteError::kind), 3982 Some(ServiceSqliteErrorKind::Authority) 3983 ); 3984 } 3985 3986 #[cfg(any(target_os = "linux", target_os = "macos"))] 3987 #[tokio::test] 3988 async fn read_only_writes_roll_back_and_internal_pool_close_refuses_new_work() { 3989 let (root, paths, identity, migrations, schema, host) = initialized_host().await; 3990 host.close().await.expect("close writer host"); 3991 let closed = host 3992 .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) 3993 .await 3994 .expect_err("closed pool must refuse work"); 3995 assert_eq!( 3996 closed.kind(), 3997 ServiceSqliteTransactionErrorKind::NotCommitted 3998 ); 3999 drop(host); 4000 4001 let inspection = ServiceSqliteHost::open_read_only_inspection( 4002 &paths, 4003 &identity, 4004 &migrations, 4005 &schema, 4006 ServiceSqliteConnectionOptions::reviewed(), 4007 ) 4008 .await 4009 .expect("open read-only host"); 4010 let error = inspection 4011 .transaction(|transaction| { 4012 Box::pin(async move { 4013 sqlx::query("INSERT INTO host_probe (value) VALUES (5)") 4014 .execute(&mut *transaction) 4015 .await 4016 .map(|_| ()) 4017 }) 4018 }) 4019 .await 4020 .expect_err("read-only write must fail"); 4021 assert_eq!( 4022 error.kind(), 4023 ServiceSqliteTransactionErrorKind::OperationRolledBack 4024 ); 4025 assert_eq!(row_count(&inspection).await, 0); 4026 drop(inspection); 4027 drop(root); 4028 } 4029 4030 #[test] 4031 fn host_and_transaction_errors_are_redacted_and_source_free() { 4032 let error = ServiceSqliteTransactionError::not_committed_with_operation( 4033 Some("operation-secret"), 4034 ServiceSqliteError::new(ServiceSqliteErrorKind::Open), 4035 ); 4036 let debug = format!("{error:?}"); 4037 assert!(!debug.contains("operation-secret")); 4038 assert!(!debug.contains("state.sqlite")); 4039 assert!(error.source().is_none()); 4040 assert_eq!(error.operation_error(), Some(&"operation-secret")); 4041 assert_eq!( 4042 error.sqlite_error().map(ServiceSqliteError::kind), 4043 Some(ServiceSqliteErrorKind::Open) 4044 ); 4045 } 4046 }