capture.rs (71292B)
1 //! Cancellation-aware online capture into a caller-owned staging directory. 2 3 use core::fmt; 4 use std::{ 5 error::Error, 6 ffi::{OsStr, OsString}, 7 fs::File, 8 io::{self, Read, Seek, SeekFrom}, 9 os::unix::ffi::OsStrExt, 10 path::{Component, Path, PathBuf}, 11 sync::{ 12 Arc, 13 atomic::{AtomicBool, Ordering}, 14 }, 15 thread, 16 }; 17 18 #[cfg(test)] 19 use core::sync::atomic::AtomicU8; 20 21 use rustix::{ 22 fs::{AtFlags, FileType, Mode, OFlags, fchmod, fstat, mkdirat, open, openat, statat, unlinkat}, 23 process::geteuid, 24 }; 25 use sha2::{Digest, Sha256}; 26 use sqlx::{ 27 ConnectOptions, Connection as _, Row, Sqlite, SqliteConnection, pool::PoolConnection, 28 sqlite::SqliteConnectOptions, 29 }; 30 31 use crate::{ 32 BackupCreatedAtUnixMs, BackupMemberSha256, OpenMode, ServiceBackupManifest, 33 ServiceDatabaseMetadata, ServiceSqliteError, ServiceSqliteErrorKind, 34 open::{BackupSourceValidator, PrivateConnectionPool}, 35 sqlite_native_backup::{NativeBackup, NativeBackupStep}, 36 }; 37 38 const MAX_STAGING_PATH_BYTES: usize = 4_096; 39 const BACKUP_PAGES_PER_STEP: i32 = 64; 40 const HASH_BUFFER_BYTES: usize = 16 * 1_024; 41 const STATE_FILE_NAME: &str = radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME; 42 const KNOWN_SIDECARS: [&str; 3] = [ 43 "state.sqlite-wal", 44 "state.sqlite-shm", 45 "state.sqlite-journal", 46 ]; 47 48 pub(crate) const TEST_CAPTURE_PHASE_BEFORE_CREATE: u8 = 1; 49 pub(crate) const TEST_CAPTURE_PHASE_STAGING_CREATED: u8 = 2; 50 pub(crate) const TEST_CAPTURE_PHASE_BACKUP_STEPPED: u8 = 3; 51 pub(crate) const TEST_CAPTURE_PHASE_POST_COPY: u8 = 4; 52 pub(crate) const TEST_CAPTURE_PHASE_PRE_FINAL_SYNC: u8 = 5; 53 pub(crate) const TEST_CAPTURE_PHASE_METADATA_AWAITED: u8 = 10; 54 pub(crate) const TEST_CAPTURE_PHASE_JOIN_AWAITED: u8 = 11; 55 #[cfg(test)] 56 static TEST_CAPTURE_PHASE: AtomicU8 = AtomicU8::new(0); 57 #[cfg(test)] 58 static TEST_CAPTURE_BLOCK_PHASE: AtomicU8 = AtomicU8::new(0); 59 #[cfg(test)] 60 static TEST_CAPTURE_INJECT_METADATA_FAILURE: AtomicBool = AtomicBool::new(false); 61 #[cfg(test)] 62 static TEST_CAPTURE_PANIC_WORKER: AtomicBool = AtomicBool::new(false); 63 64 #[cfg(test)] 65 pub(crate) fn test_capture_phase() -> u8 { 66 TEST_CAPTURE_PHASE.load(Ordering::Acquire) 67 } 68 69 #[cfg(test)] 70 pub(crate) fn test_capture_block_phase(phase: u8) { 71 TEST_CAPTURE_BLOCK_PHASE.store(phase, Ordering::Release); 72 } 73 74 #[cfg(test)] 75 pub(crate) fn test_capture_reset() { 76 TEST_CAPTURE_BLOCK_PHASE.store(0, Ordering::Release); 77 TEST_CAPTURE_PHASE.store(0, Ordering::Release); 78 TEST_CAPTURE_INJECT_METADATA_FAILURE.store(false, Ordering::Release); 79 TEST_CAPTURE_PANIC_WORKER.store(false, Ordering::Release); 80 } 81 82 #[cfg(test)] 83 pub(crate) fn test_capture_inject_metadata_failure(enabled: bool) { 84 TEST_CAPTURE_INJECT_METADATA_FAILURE.store(enabled, Ordering::Release); 85 } 86 87 #[cfg(test)] 88 pub(crate) fn test_capture_panic_worker(enabled: bool) { 89 TEST_CAPTURE_PANIC_WORKER.store(enabled, Ordering::Release); 90 } 91 92 async fn test_async_phase(phase: u8) { 93 #[cfg(test)] 94 { 95 TEST_CAPTURE_PHASE.store(phase, Ordering::Release); 96 while TEST_CAPTURE_BLOCK_PHASE.load(Ordering::Acquire) == phase { 97 tokio::task::yield_now().await; 98 } 99 } 100 #[cfg(not(test))] 101 let _ = phase; 102 } 103 104 trait CaptureOperations: Send + Sync { 105 fn sync_state(&self, state: &File) -> io::Result<()>; 106 fn sync_staging(&self, staging: &File) -> io::Result<()>; 107 fn sync_parent(&self, parent: &File) -> io::Result<()>; 108 } 109 110 struct SystemCaptureOperations; 111 112 impl CaptureOperations for SystemCaptureOperations { 113 #[cfg_attr(coverage_nightly, coverage(off))] 114 fn sync_state(&self, state: &File) -> io::Result<()> { 115 state.sync_all() 116 } 117 118 #[cfg_attr(coverage_nightly, coverage(off))] 119 fn sync_staging(&self, staging: &File) -> io::Result<()> { 120 staging.sync_all() 121 } 122 123 #[cfg_attr(coverage_nightly, coverage(off))] 124 fn sync_parent(&self, parent: &File) -> io::Result<()> { 125 parent.sync_all() 126 } 127 } 128 129 #[cfg(test)] 130 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 131 pub(crate) enum TestCaptureSyncFailure { 132 State, 133 Staging, 134 FinalParent, 135 } 136 137 #[cfg(test)] 138 struct FailingCaptureOperations { 139 failure: TestCaptureSyncFailure, 140 parent_syncs: core::sync::atomic::AtomicU8, 141 } 142 143 #[cfg(test)] 144 impl CaptureOperations for FailingCaptureOperations { 145 fn sync_state(&self, state: &File) -> io::Result<()> { 146 if self.failure == TestCaptureSyncFailure::State { 147 Err(crate::failpoint::storage_full_error()) 148 } else { 149 state.sync_all() 150 } 151 } 152 153 fn sync_staging(&self, staging: &File) -> io::Result<()> { 154 if self.failure == TestCaptureSyncFailure::Staging { 155 Err(crate::failpoint::storage_full_error()) 156 } else { 157 staging.sync_all() 158 } 159 } 160 161 fn sync_parent(&self, parent: &File) -> io::Result<()> { 162 let occurrence = self.parent_syncs.fetch_add(1, Ordering::AcqRel); 163 if self.failure == TestCaptureSyncFailure::FinalParent && occurrence == 1 { 164 Err(crate::failpoint::storage_full_error()) 165 } else { 166 parent.sync_all() 167 } 168 } 169 } 170 171 pub(crate) async fn capture_online_backup( 172 pool: &PrivateConnectionPool, 173 closing: &AtomicBool, 174 active: &Arc<AtomicBool>, 175 staging_directory: &Path, 176 created_at_unix_ms: BackupCreatedAtUnixMs, 177 failpoints: &crate::failpoint::DurabilityFailpoints, 178 ) -> Result<ServiceBackupManifest, ServiceSqliteError> { 179 capture_online_backup_with_operations( 180 pool, 181 closing, 182 active, 183 staging_directory, 184 created_at_unix_ms, 185 Arc::new(SystemCaptureOperations), 186 failpoints, 187 ) 188 .await 189 } 190 191 #[cfg(test)] 192 pub(crate) async fn test_capture_online_backup_with_sync_failure( 193 pool: &PrivateConnectionPool, 194 closing: &AtomicBool, 195 active: &Arc<AtomicBool>, 196 staging_directory: &Path, 197 created_at_unix_ms: BackupCreatedAtUnixMs, 198 failure: TestCaptureSyncFailure, 199 ) -> Result<ServiceBackupManifest, ServiceSqliteError> { 200 capture_online_backup_with_operations( 201 pool, 202 closing, 203 active, 204 staging_directory, 205 created_at_unix_ms, 206 Arc::new(FailingCaptureOperations { 207 failure, 208 parent_syncs: core::sync::atomic::AtomicU8::new(0), 209 }), 210 &crate::failpoint::DurabilityFailpoints::default(), 211 ) 212 .await 213 } 214 215 async fn capture_online_backup_with_operations( 216 pool: &PrivateConnectionPool, 217 closing: &AtomicBool, 218 active: &Arc<AtomicBool>, 219 staging_directory: &Path, 220 created_at_unix_ms: BackupCreatedAtUnixMs, 221 operations: Arc<dyn CaptureOperations>, 222 failpoints: &crate::failpoint::DurabilityFailpoints, 223 ) -> Result<ServiceBackupManifest, ServiceSqliteError> { 224 crate::require_condition( 225 matches!( 226 pool.mode(), 227 OpenMode::Initialize | OpenMode::ReadWriteExisting 228 ) && !closing.load(Ordering::Acquire), 229 ServiceSqliteErrorKind::Open, 230 )?; 231 let staging = StagingPath::new(staging_directory)?; 232 let permit = CapturePermit::acquire(Arc::clone(active))?; 233 pool.validate()?; 234 let mut admission = pool.acquire().await?; 235 pool.validate()?; 236 crate::require_condition( 237 !closing.load(Ordering::Acquire), 238 ServiceSqliteErrorKind::Open, 239 )?; 240 let metadata = crate::metadata::verify_database_metadata(&mut admission, pool.identity()).await; 241 #[cfg(test)] 242 let metadata = if TEST_CAPTURE_INJECT_METADATA_FAILURE.load(Ordering::Acquire) { 243 Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata)) 244 } else { 245 metadata 246 }; 247 test_async_phase(TEST_CAPTURE_PHASE_METADATA_AWAITED).await; 248 pool.validate()?; 249 let metadata = metadata?; 250 let validator = pool.backup_source_validator(); 251 validator.validate()?; 252 253 let cancellation = Arc::new(AtomicBool::new(false)); 254 let cancellation_guard = CaptureCancellation::new(Arc::clone(&cancellation)); 255 let worker = CaptureWorker { 256 admission: Some(admission), 257 _permit: permit, 258 validator, 259 metadata, 260 staging, 261 created_at_unix_ms, 262 cancellation, 263 operations, 264 failpoints: failpoints.clone(), 265 runtime: tokio::runtime::Handle::current(), 266 }; 267 let joined = tokio::task::spawn_blocking(move || worker.run()).await; 268 test_async_phase(TEST_CAPTURE_PHASE_JOIN_AWAITED).await; 269 cancellation_guard.complete(); 270 pool.validate()?; 271 let result = joined.map_err(|source| backup_source(BackupFailureKind::Join, source))?; 272 result.map(PendingCapture::commit) 273 } 274 275 struct CapturePermit { 276 active: Arc<AtomicBool>, 277 } 278 279 impl CapturePermit { 280 fn acquire(active: Arc<AtomicBool>) -> Result<Self, ServiceSqliteError> { 281 active 282 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire) 283 .map_err(|_| backup_error(BackupFailureKind::AlreadyActive))?; 284 Ok(Self { active }) 285 } 286 } 287 288 impl Drop for CapturePermit { 289 fn drop(&mut self) { 290 self.active.store(false, Ordering::Release); 291 } 292 } 293 294 struct CaptureCancellation { 295 cancelled: Arc<AtomicBool>, 296 complete: AtomicBool, 297 } 298 299 impl CaptureCancellation { 300 fn new(cancelled: Arc<AtomicBool>) -> Self { 301 Self { 302 cancelled, 303 complete: AtomicBool::new(false), 304 } 305 } 306 307 fn complete(&self) { 308 self.complete.store(true, Ordering::Release); 309 } 310 } 311 312 impl Drop for CaptureCancellation { 313 fn drop(&mut self) { 314 if !self.complete.load(Ordering::Acquire) { 315 self.cancelled.store(true, Ordering::Release); 316 } 317 } 318 } 319 320 struct CaptureWorker { 321 admission: Option<PoolConnection<Sqlite>>, 322 _permit: CapturePermit, 323 validator: BackupSourceValidator, 324 metadata: ServiceDatabaseMetadata, 325 staging: StagingPath, 326 created_at_unix_ms: BackupCreatedAtUnixMs, 327 cancellation: Arc<AtomicBool>, 328 operations: Arc<dyn CaptureOperations>, 329 failpoints: crate::failpoint::DurabilityFailpoints, 330 runtime: tokio::runtime::Handle, 331 } 332 333 impl CaptureWorker { 334 fn run(mut self) -> Result<PendingCapture, ServiceSqliteError> { 335 self.check_cancelled()?; 336 self.validator.validate()?; 337 self.test_phase(TEST_CAPTURE_PHASE_BEFORE_CREATE); 338 self.check_cancelled()?; 339 self.hit_checked( 340 None, 341 crate::failpoint::DurabilityFailpoint::BackupBeforeCreate, 342 BackupFailureKind::CreateStaging, 343 )?; 344 let mut staging = StagingGuard::create(&self.staging, self.operations.as_ref())?; 345 self.hit_checked( 346 Some(&staging), 347 crate::failpoint::DurabilityFailpoint::BackupAfterCreate, 348 BackupFailureKind::CreateStaging, 349 )?; 350 self.test_phase(TEST_CAPTURE_PHASE_STAGING_CREATED); 351 #[cfg(test)] 352 if TEST_CAPTURE_PANIC_WORKER.load(Ordering::Acquire) { 353 panic!("injected backup worker failure"); 354 } 355 self.check_cancelled()?; 356 self.validator.validate()?; 357 staging.validate()?; 358 359 let source = self 360 .admission 361 .as_mut() 362 .ok_or_else(|| backup_error(BackupFailureKind::Capture))?; 363 self.runtime.block_on(verify_database_inventory(source))?; 364 self.runtime 365 .block_on(verify_database_metadata(source, &self.metadata))?; 366 self.validator.validate()?; 367 staging.validate()?; 368 369 let mut destination = self.open_sqlx_destination(&staging)?; 370 staging.record_sidecars(); 371 372 self.hit_checked( 373 Some(&staging), 374 crate::failpoint::DurabilityFailpoint::BackupBeforeCopy, 375 BackupFailureKind::Capture, 376 )?; 377 let capture_result = self.copy_with_locked_sqlx_handles(&mut destination, &mut staging); 378 let close_result = self 379 .runtime 380 .block_on(destination.close()) 381 .map_err(|source| backup_source(BackupFailureKind::Capture, source)); 382 staging.record_sidecars(); 383 self.validator.validate()?; 384 staging.validate()?; 385 capture_result?; 386 close_result?; 387 self.hit_checked( 388 Some(&staging), 389 crate::failpoint::DurabilityFailpoint::BackupAfterCopy, 390 BackupFailureKind::Capture, 391 )?; 392 staging.record_sidecars(); 393 394 self.test_phase(TEST_CAPTURE_PHASE_POST_COPY); 395 self.check_cancelled()?; 396 self.validator.validate()?; 397 staging.validate()?; 398 let mut destination = self.open_inspection_destination(&staging)?; 399 self.runtime 400 .block_on(verify_database_inventory(&mut destination))?; 401 self.runtime 402 .block_on(verify_database_metadata(&mut destination, &self.metadata))?; 403 self.runtime.block_on(verify_integrity(&mut destination))?; 404 staging.record_sidecars(); 405 self.check_cancelled()?; 406 self.validator.validate()?; 407 staging.validate()?; 408 409 self.runtime 410 .block_on(destination.close()) 411 .map_err(|source| backup_source(BackupFailureKind::Capture, source))?; 412 staging.record_sidecars(); 413 self.validator.validate()?; 414 staging.validate()?; 415 staging.validate_inventory()?; 416 self.check_cancelled()?; 417 418 self.hit_checked( 419 Some(&staging), 420 crate::failpoint::DurabilityFailpoint::BackupBeforeFileSync, 421 BackupFailureKind::SyncState, 422 )?; 423 staging.sync_state(self.operations.as_ref())?; 424 self.hit_checked( 425 Some(&staging), 426 crate::failpoint::DurabilityFailpoint::BackupAfterFileSync, 427 BackupFailureKind::SyncState, 428 )?; 429 let (byte_length, digest) = staging.hash_state(&self.cancellation)?; 430 self.test_phase(TEST_CAPTURE_PHASE_PRE_FINAL_SYNC); 431 self.check_cancelled()?; 432 self.hit_checked( 433 Some(&staging), 434 crate::failpoint::DurabilityFailpoint::BackupBeforeDirectorySync, 435 BackupFailureKind::SyncStaging, 436 )?; 437 staging.sync_directories(self.operations.as_ref())?; 438 self.hit_checked( 439 Some(&staging), 440 crate::failpoint::DurabilityFailpoint::BackupAfterDirectorySync, 441 BackupFailureKind::SyncParent, 442 )?; 443 self.validator.validate()?; 444 staging.validate()?; 445 staging.validate_inventory()?; 446 self.check_cancelled()?; 447 448 let manifest = ServiceBackupManifest::from_capture( 449 &self.metadata, 450 self.created_at_unix_ms, 451 byte_length, 452 BackupMemberSha256::from_bytes(digest), 453 ) 454 .map_err(|source| backup_source(BackupFailureKind::Manifest, source))?; 455 Ok(PendingCapture { 456 staging, 457 _admission: self 458 .admission 459 .take() 460 .ok_or_else(|| backup_error(BackupFailureKind::Capture))?, 461 _permit: self._permit, 462 manifest, 463 }) 464 } 465 466 fn copy_with_locked_sqlx_handles( 467 &mut self, 468 destination: &mut SqliteConnection, 469 staging: &mut StagingGuard, 470 ) -> Result<(), ServiceSqliteError> { 471 let mut admission = self 472 .admission 473 .take() 474 .ok_or_else(|| backup_error(BackupFailureKind::Capture))?; 475 let result = (|| { 476 let mut source_handle = self 477 .runtime 478 .block_on(admission.lock_handle()) 479 .map_err(|source| backup_source(BackupFailureKind::Capture, source))?; 480 let mut destination_handle = self 481 .runtime 482 .block_on(destination.lock_handle()) 483 .map_err(|source| backup_source(BackupFailureKind::Capture, source))?; 484 let mut backup = NativeBackup::start(&mut destination_handle, &mut source_handle) 485 .map_err(|source| backup_source(BackupFailureKind::Capture, source))?; 486 loop { 487 self.check_cancelled()?; 488 self.validator.validate()?; 489 staging.validate()?; 490 let step = backup 491 .step(BACKUP_PAGES_PER_STEP) 492 .map_err(|source| backup_source(BackupFailureKind::Capture, source)); 493 staging.record_sidecars(); 494 self.test_phase(TEST_CAPTURE_PHASE_BACKUP_STEPPED); 495 let step = step?; 496 self.validator.validate()?; 497 staging.validate()?; 498 match step { 499 NativeBackupStep::Done => { 500 return backup 501 .finish() 502 .map_err(|source| backup_source(BackupFailureKind::Capture, source)); 503 } 504 NativeBackupStep::More => {} 505 NativeBackupStep::Busy | NativeBackupStep::Locked => thread::yield_now(), 506 } 507 } 508 })(); 509 self.admission = Some(admission); 510 result 511 } 512 513 fn open_sqlx_destination( 514 &self, 515 staging: &StagingGuard, 516 ) -> Result<SqliteConnection, ServiceSqliteError> { 517 staging.validate()?; 518 let options = SqliteConnectOptions::new() 519 .filename(staging.state_path()) 520 .create_if_missing(false) 521 .foreign_keys(false) 522 .disable_statement_logging(); 523 let result = self 524 .runtime 525 .block_on(SqliteConnection::connect_with(&options)); 526 staging.validate()?; 527 result.map_err(|source| backup_source(BackupFailureKind::Capture, source)) 528 } 529 530 fn open_inspection_destination( 531 &self, 532 staging: &StagingGuard, 533 ) -> Result<SqliteConnection, ServiceSqliteError> { 534 staging.validate()?; 535 let options = SqliteConnectOptions::new() 536 .filename(staging.state_path()) 537 .create_if_missing(false) 538 .foreign_keys(false) 539 .disable_statement_logging(); 540 let result = self 541 .runtime 542 .block_on(SqliteConnection::connect_with(&options)); 543 staging.validate()?; 544 result.map_err(|source| backup_source(BackupFailureKind::Capture, source)) 545 } 546 547 fn check_cancelled(&self) -> Result<(), ServiceSqliteError> { 548 if self.cancellation.load(Ordering::Acquire) { 549 Err(backup_error(BackupFailureKind::Cancelled)) 550 } else { 551 Ok(()) 552 } 553 } 554 555 fn hit_checked( 556 &self, 557 staging: Option<&StagingGuard>, 558 point: crate::failpoint::DurabilityFailpoint, 559 failure: BackupFailureKind, 560 ) -> Result<(), ServiceSqliteError> { 561 #[cfg(test)] 562 self.failpoints 563 .observe(point, TEST_CAPTURE_PHASE.load(Ordering::Acquire)); 564 let injected = self 565 .failpoints 566 .hit(point) 567 .map_err(|source| backup_source(failure, source)); 568 self.validator.validate()?; 569 if let Some(staging) = staging { 570 staging.validate()?; 571 } 572 injected 573 } 574 575 fn test_phase(&self, phase: u8) { 576 #[cfg(test)] 577 { 578 TEST_CAPTURE_PHASE.store(phase, Ordering::Release); 579 while TEST_CAPTURE_BLOCK_PHASE.load(Ordering::Acquire) == phase 580 && !self.cancellation.load(Ordering::Acquire) 581 { 582 thread::yield_now(); 583 } 584 } 585 #[cfg(not(test))] 586 let _ = phase; 587 } 588 } 589 590 struct PendingCapture { 591 staging: StagingGuard, 592 _admission: PoolConnection<Sqlite>, 593 _permit: CapturePermit, 594 manifest: ServiceBackupManifest, 595 } 596 597 impl PendingCapture { 598 fn commit(mut self) -> ServiceBackupManifest { 599 self.staging.commit(); 600 self.manifest 601 } 602 } 603 604 #[derive(Clone)] 605 struct StagingPath { 606 full: PathBuf, 607 parent: PathBuf, 608 name: OsString, 609 } 610 611 impl StagingPath { 612 fn new(path: &Path) -> Result<Self, ServiceSqliteError> { 613 if path.as_os_str().as_bytes().len() > MAX_STAGING_PATH_BYTES || !path.is_absolute() { 614 return Err(backup_error(BackupFailureKind::InvalidStagingPath)); 615 } 616 if path 617 .components() 618 .any(|component| matches!(component, Component::CurDir | Component::ParentDir)) 619 { 620 return Err(backup_error(BackupFailureKind::InvalidStagingPath)); 621 } 622 let name = match path.components().next_back() { 623 Some(Component::Normal(name)) if !name.as_bytes().is_empty() => name.to_os_string(), 624 _ => return Err(backup_error(BackupFailureKind::InvalidStagingPath)), 625 }; 626 let parent = path 627 .parent() 628 .filter(|parent| parent.is_absolute()) 629 .ok_or_else(|| backup_error(BackupFailureKind::InvalidStagingPath))? 630 .to_path_buf(); 631 Ok(Self { 632 full: path.to_path_buf(), 633 parent, 634 name, 635 }) 636 } 637 } 638 639 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 640 struct FileIdentity { 641 device: u64, 642 inode: u64, 643 } 644 645 struct StagingGuard { 646 path: StagingPath, 647 parent: File, 648 parent_identity: FileIdentity, 649 directory: File, 650 directory_identity: FileIdentity, 651 state: File, 652 state_identity: FileIdentity, 653 sidecar_identities: [Option<FileIdentity>; KNOWN_SIDECARS.len()], 654 committed: bool, 655 } 656 657 impl StagingGuard { 658 fn create( 659 path: &StagingPath, 660 operations: &dyn CaptureOperations, 661 ) -> Result<Self, ServiceSqliteError> { 662 let parent = File::from( 663 open( 664 &path.parent, 665 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 666 Mode::empty(), 667 ) 668 .map_err(|source| backup_source(BackupFailureKind::InvalidStagingParent, source))?, 669 ); 670 let parent_identity = validate_directory_descriptor(&parent, false)?; 671 mkdirat(&parent, &path.name, Mode::RUSR | Mode::WUSR | Mode::XUSR).map_err(|source| { 672 backup_source( 673 if source == rustix::io::Errno::EXIST { 674 BackupFailureKind::StagingCollision 675 } else { 676 BackupFailureKind::CreateStaging 677 }, 678 source, 679 ) 680 })?; 681 let created_directory_identity = match created_directory_identity(&parent, &path.name) { 682 Ok(identity) => identity, 683 Err(error) => { 684 // Without a proven identity, the current entry must be preserved. 685 let _ = parent.sync_all(); 686 return Err(error); 687 } 688 }; 689 let directory = match openat( 690 &parent, 691 &path.name, 692 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 693 Mode::empty(), 694 ) { 695 Ok(directory) => File::from(directory), 696 Err(source) => { 697 cleanup_partial_staging( 698 &parent, 699 &path.name, 700 Some(created_directory_identity), 701 None, 702 None, 703 ); 704 return Err(backup_source(BackupFailureKind::CreateStaging, source)); 705 } 706 }; 707 let opened_directory_identity = match descriptor_identity(&directory) { 708 Ok(identity) => identity, 709 Err(error) => { 710 cleanup_partial_staging( 711 &parent, 712 &path.name, 713 Some(created_directory_identity), 714 Some(&directory), 715 None, 716 ); 717 return Err(error); 718 } 719 }; 720 if let Err(error) = require_backup_condition( 721 opened_directory_identity == created_directory_identity, 722 BackupFailureKind::StagingReplaced, 723 ) { 724 cleanup_partial_staging( 725 &parent, 726 &path.name, 727 Some(created_directory_identity), 728 Some(&directory), 729 None, 730 ); 731 return Err(error); 732 } 733 if let Err(source) = fchmod(&directory, Mode::RUSR | Mode::WUSR | Mode::XUSR) { 734 cleanup_partial_staging( 735 &parent, 736 &path.name, 737 Some(created_directory_identity), 738 Some(&directory), 739 None, 740 ); 741 return Err(backup_source(BackupFailureKind::CreateStaging, source)); 742 } 743 let directory_identity = match validate_directory_descriptor(&directory, true) { 744 Ok(identity) => identity, 745 Err(error) => { 746 cleanup_partial_staging( 747 &parent, 748 &path.name, 749 Some(created_directory_identity), 750 Some(&directory), 751 None, 752 ); 753 return Err(error); 754 } 755 }; 756 let state = match openat( 757 &directory, 758 STATE_FILE_NAME, 759 OFlags::RDWR | OFlags::CREATE | OFlags::EXCL | OFlags::NOFOLLOW | OFlags::CLOEXEC, 760 Mode::RUSR | Mode::WUSR, 761 ) { 762 Ok(state) => File::from(state), 763 Err(source) => { 764 cleanup_partial_staging( 765 &parent, 766 &path.name, 767 Some(created_directory_identity), 768 Some(&directory), 769 None, 770 ); 771 return Err(backup_source(BackupFailureKind::CreateState, source)); 772 } 773 }; 774 let created_state_identity = match descriptor_identity(&state) { 775 Ok(identity) => identity, 776 Err(error) => { 777 cleanup_partial_staging( 778 &parent, 779 &path.name, 780 Some(created_directory_identity), 781 Some(&directory), 782 None, 783 ); 784 return Err(error); 785 } 786 }; 787 if let Err(source) = fchmod(&state, Mode::RUSR | Mode::WUSR) { 788 cleanup_partial_staging( 789 &parent, 790 &path.name, 791 Some(created_directory_identity), 792 Some(&directory), 793 Some(created_state_identity), 794 ); 795 return Err(backup_source(BackupFailureKind::CreateState, source)); 796 } 797 let state_identity = match validate_file_descriptor(&state) { 798 Ok(identity) => identity, 799 Err(error) => { 800 cleanup_partial_staging( 801 &parent, 802 &path.name, 803 Some(created_directory_identity), 804 Some(&directory), 805 Some(created_state_identity), 806 ); 807 return Err(error); 808 } 809 }; 810 let staging = Self { 811 path: path.clone(), 812 parent, 813 parent_identity, 814 directory, 815 directory_identity, 816 state, 817 state_identity, 818 sidecar_identities: [None; KNOWN_SIDECARS.len()], 819 committed: false, 820 }; 821 staging.validate()?; 822 operations 823 .sync_parent(&staging.parent) 824 .map_err(|source| backup_source(BackupFailureKind::SyncParent, source))?; 825 Ok(staging) 826 } 827 828 fn state_path(&self) -> PathBuf { 829 self.path.full.join(STATE_FILE_NAME) 830 } 831 832 fn validate(&self) -> Result<(), ServiceSqliteError> { 833 validate_reopened_directory(&self.path.parent, &self.parent, self.parent_identity, false)?; 834 validate_directory_entry( 835 &self.parent, 836 &self.path.name, 837 &self.directory, 838 self.directory_identity, 839 )?; 840 validate_file_entry( 841 &self.directory, 842 OsStr::new(STATE_FILE_NAME), 843 &self.state, 844 self.state_identity, 845 ) 846 } 847 848 fn validate_inventory(&self) -> Result<(), ServiceSqliteError> { 849 self.validate()?; 850 let mut entries = std::fs::read_dir(&self.path.full) 851 .map_err(|source| backup_source(BackupFailureKind::InvalidStagingInventory, source))?; 852 let first = entries 853 .next() 854 .transpose() 855 .map_err(|source| backup_source(BackupFailureKind::InvalidStagingInventory, source))? 856 .ok_or_else(|| backup_error(BackupFailureKind::InvalidStagingInventory))?; 857 require_backup_condition( 858 crate::all_constraints([ 859 first.file_name() == OsStr::new(STATE_FILE_NAME), 860 entries.next().is_none(), 861 ]), 862 BackupFailureKind::InvalidStagingInventory, 863 )?; 864 Ok(()) 865 } 866 867 fn sync_state(&self, operations: &dyn CaptureOperations) -> Result<(), ServiceSqliteError> { 868 self.validate()?; 869 operations 870 .sync_state(&self.state) 871 .map_err(|source| backup_source(BackupFailureKind::SyncState, source))?; 872 self.validate() 873 } 874 875 fn hash_state(&self, cancellation: &AtomicBool) -> Result<(u64, [u8; 32]), ServiceSqliteError> { 876 self.validate()?; 877 let mut state = self 878 .state 879 .try_clone() 880 .map_err(|source| backup_source(BackupFailureKind::HashState, source))?; 881 state 882 .seek(SeekFrom::Start(0)) 883 .map_err(|source| backup_source(BackupFailureKind::HashState, source))?; 884 let mut hasher = Sha256::new(); 885 let mut length = 0_u64; 886 let mut buffer = [0_u8; HASH_BUFFER_BYTES]; 887 loop { 888 require_backup_condition( 889 !cancellation.load(Ordering::Acquire), 890 BackupFailureKind::Cancelled, 891 )?; 892 let count = state 893 .read(&mut buffer) 894 .map_err(|source| backup_source(BackupFailureKind::HashState, source))?; 895 if count == 0 { 896 break; 897 } 898 length = length 899 .checked_add( 900 u64::try_from(count).map_err(|_| backup_error(BackupFailureKind::HashState))?, 901 ) 902 .ok_or_else(|| backup_error(BackupFailureKind::HashState))?; 903 require_backup_condition(length <= i64::MAX as u64, BackupFailureKind::HashState)?; 904 hasher.update(&buffer[..count]); 905 } 906 require_backup_condition(length != 0, BackupFailureKind::HashState)?; 907 self.validate()?; 908 Ok((length, hasher.finalize().into())) 909 } 910 911 fn sync_directories( 912 &self, 913 operations: &dyn CaptureOperations, 914 ) -> Result<(), ServiceSqliteError> { 915 self.validate()?; 916 operations 917 .sync_staging(&self.directory) 918 .map_err(|source| backup_source(BackupFailureKind::SyncStaging, source))?; 919 self.validate()?; 920 operations 921 .sync_parent(&self.parent) 922 .map_err(|source| backup_source(BackupFailureKind::SyncParent, source))?; 923 self.validate() 924 } 925 926 fn record_sidecars(&mut self) { 927 for (index, name) in KNOWN_SIDECARS.iter().enumerate() { 928 if self.sidecar_identities[index].is_none() { 929 self.sidecar_identities[index] = safe_sidecar_identity(&self.directory, name); 930 } 931 } 932 } 933 934 fn commit(&mut self) { 935 self.committed = true; 936 } 937 938 fn cleanup(&mut self) { 939 if self.committed { 940 return; 941 } 942 if validate_directory_entry( 943 &self.parent, 944 &self.path.name, 945 &self.directory, 946 self.directory_identity, 947 ) 948 .is_err() 949 { 950 return; 951 } 952 if current_entry_identity(&self.directory, OsStr::new(STATE_FILE_NAME)) 953 == Some(self.state_identity) 954 { 955 let _ = unlinkat(&self.directory, STATE_FILE_NAME, AtFlags::empty()); 956 } 957 for (sidecar, identity) in KNOWN_SIDECARS.iter().zip(self.sidecar_identities) { 958 if identity.is_some() && safe_sidecar_identity(&self.directory, sidecar) == identity { 959 let _ = unlinkat(&self.directory, *sidecar, AtFlags::empty()); 960 } 961 } 962 if current_entry_identity(&self.parent, &self.path.name) == Some(self.directory_identity) { 963 let _ = unlinkat(&self.parent, &self.path.name, AtFlags::REMOVEDIR); 964 } 965 let _ = self.parent.sync_all(); 966 self.committed = true; 967 } 968 } 969 970 fn cleanup_partial_staging( 971 parent: &File, 972 name: &OsStr, 973 directory_identity: Option<FileIdentity>, 974 directory: Option<&File>, 975 state_identity: Option<FileIdentity>, 976 ) { 977 if let (Some(directory), Some(state_identity)) = (directory, state_identity) 978 && current_entry_identity(directory, OsStr::new(STATE_FILE_NAME)) == Some(state_identity) 979 { 980 let _ = unlinkat(directory, STATE_FILE_NAME, AtFlags::empty()); 981 } 982 if directory_identity.is_some() && current_entry_identity(parent, name) == directory_identity { 983 let _ = unlinkat(parent, name, AtFlags::REMOVEDIR); 984 } 985 let _ = parent.sync_all(); 986 } 987 988 impl Drop for StagingGuard { 989 fn drop(&mut self) { 990 self.cleanup(); 991 } 992 } 993 994 fn created_directory_identity( 995 parent: &File, 996 name: &OsStr, 997 ) -> Result<FileIdentity, ServiceSqliteError> { 998 let status = statat(parent, name, AtFlags::SYMLINK_NOFOLLOW) 999 .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?; 1000 if !crate::native_metadata::secure_directory( 1001 FileType::from_raw_mode(status.st_mode).is_dir(), 1002 status.st_uid, 1003 geteuid().as_raw(), 1004 crate::native_metadata::mode(status.st_mode), 1005 ) { 1006 return Err(backup_error(BackupFailureKind::StagingReplaced)); 1007 } 1008 identity(&status) 1009 } 1010 1011 fn descriptor_identity(descriptor: &File) -> Result<FileIdentity, ServiceSqliteError> { 1012 let status = fstat(descriptor) 1013 .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?; 1014 identity(&status) 1015 } 1016 1017 fn validate_directory_descriptor( 1018 directory: &File, 1019 exact_owner_mode: bool, 1020 ) -> Result<FileIdentity, ServiceSqliteError> { 1021 let status = fstat(directory) 1022 .map_err(|source| backup_source(BackupFailureKind::InvalidStagingParent, source))?; 1023 let mode = crate::native_metadata::mode(status.st_mode) & 0o777; 1024 let valid = if exact_owner_mode { 1025 crate::native_metadata::exact_directory( 1026 FileType::from_raw_mode(status.st_mode).is_dir(), 1027 status.st_uid, 1028 geteuid().as_raw(), 1029 mode, 1030 ) 1031 } else { 1032 crate::native_metadata::secure_directory( 1033 FileType::from_raw_mode(status.st_mode).is_dir(), 1034 status.st_uid, 1035 geteuid().as_raw(), 1036 mode, 1037 ) 1038 }; 1039 if !valid { 1040 return Err(backup_error(BackupFailureKind::InvalidStagingParent)); 1041 } 1042 identity(&status) 1043 } 1044 1045 fn validate_file_descriptor(file: &File) -> Result<FileIdentity, ServiceSqliteError> { 1046 let status = fstat(file) 1047 .map_err(|source| backup_source(BackupFailureKind::InvalidStagingInventory, source))?; 1048 if !crate::native_metadata::exact_regular_file( 1049 FileType::from_raw_mode(status.st_mode).is_file(), 1050 crate::native_metadata::link_count(status.st_nlink), 1051 status.st_uid, 1052 geteuid().as_raw(), 1053 crate::native_metadata::mode(status.st_mode), 1054 ) { 1055 return Err(backup_error(BackupFailureKind::InvalidStagingInventory)); 1056 } 1057 identity(&status) 1058 } 1059 1060 fn validate_reopened_directory( 1061 path: &Path, 1062 held: &File, 1063 expected: FileIdentity, 1064 exact_owner_mode: bool, 1065 ) -> Result<(), ServiceSqliteError> { 1066 let current = File::from( 1067 open( 1068 path, 1069 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1070 Mode::empty(), 1071 ) 1072 .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?, 1073 ); 1074 require_backup_condition( 1075 validate_directory_descriptor(¤t, exact_owner_mode)? == expected, 1076 BackupFailureKind::StagingReplaced, 1077 )?; 1078 require_backup_condition( 1079 validate_directory_descriptor(held, exact_owner_mode)? == expected, 1080 BackupFailureKind::StagingReplaced, 1081 )?; 1082 Ok(()) 1083 } 1084 1085 fn validate_directory_entry( 1086 parent: &File, 1087 name: &OsStr, 1088 held: &File, 1089 expected: FileIdentity, 1090 ) -> Result<(), ServiceSqliteError> { 1091 let current = File::from( 1092 openat( 1093 parent, 1094 name, 1095 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1096 Mode::empty(), 1097 ) 1098 .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?, 1099 ); 1100 require_backup_condition( 1101 validate_directory_descriptor(¤t, true)? == expected, 1102 BackupFailureKind::StagingReplaced, 1103 )?; 1104 require_backup_condition( 1105 validate_directory_descriptor(held, true)? == expected, 1106 BackupFailureKind::StagingReplaced, 1107 )?; 1108 Ok(()) 1109 } 1110 1111 fn validate_file_entry( 1112 directory: &File, 1113 name: &OsStr, 1114 held: &File, 1115 expected: FileIdentity, 1116 ) -> Result<(), ServiceSqliteError> { 1117 let current = File::from( 1118 openat( 1119 directory, 1120 name, 1121 OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1122 Mode::empty(), 1123 ) 1124 .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?, 1125 ); 1126 require_backup_condition( 1127 validate_file_descriptor(¤t)? == expected, 1128 BackupFailureKind::StagingReplaced, 1129 )?; 1130 require_backup_condition( 1131 validate_file_descriptor(held)? == expected, 1132 BackupFailureKind::StagingReplaced, 1133 )?; 1134 Ok(()) 1135 } 1136 1137 fn current_entry_identity(directory: &File, name: &OsStr) -> Option<FileIdentity> { 1138 let status = statat(directory, name, AtFlags::SYMLINK_NOFOLLOW).ok()?; 1139 identity(&status).ok() 1140 } 1141 1142 fn safe_sidecar_identity(directory: &File, name: &str) -> Option<FileIdentity> { 1143 let status = statat(directory, name, AtFlags::SYMLINK_NOFOLLOW).ok()?; 1144 if !crate::native_metadata::regular_owner_single_link( 1145 FileType::from_raw_mode(status.st_mode).is_file(), 1146 crate::native_metadata::link_count(status.st_nlink), 1147 status.st_uid, 1148 geteuid().as_raw(), 1149 ) { 1150 return None; 1151 } 1152 identity(&status).ok() 1153 } 1154 1155 fn identity(status: &rustix::fs::Stat) -> Result<FileIdentity, ServiceSqliteError> { 1156 Ok(FileIdentity { 1157 device: crate::native_metadata::device(status.st_dev) 1158 .map_err(|_| backup_error(BackupFailureKind::StagingReplaced))?, 1159 inode: status.st_ino, 1160 }) 1161 } 1162 1163 async fn verify_database_inventory( 1164 connection: &mut SqliteConnection, 1165 ) -> Result<(), ServiceSqliteError> { 1166 let rows = sqlx::query( 1167 "SELECT 1168 seq, 1169 typeof(name) = 'text' AS name_type_ok, 1170 length(CAST(name AS BLOB)) AS name_length, 1171 substr(CAST(name AS BLOB), 1, 5) AS name_prefix 1172 FROM pragma_database_list 1173 LIMIT 2", 1174 ) 1175 .fetch_all(connection) 1176 .await 1177 .map_err(|source| backup_source(BackupFailureKind::Capture, source))?; 1178 let first = rows 1179 .first() 1180 .ok_or_else(|| backup_error(BackupFailureKind::Capture))?; 1181 let sequence = first 1182 .try_get::<i64, _>(0) 1183 .map_err(|source| backup_source(BackupFailureKind::Capture, source))?; 1184 let name = crate::persisted_value::bounded_utf8( 1185 first, 1186 "name_type_ok", 1187 "name_length", 1188 "name_prefix", 1189 1, 1190 4, 1191 ) 1192 .ok_or_else(|| backup_error(BackupFailureKind::Capture))?; 1193 require_backup_condition( 1194 database_inventory_matches(sequence, name, rows.len() > 1), 1195 BackupFailureKind::Capture, 1196 )?; 1197 Ok(()) 1198 } 1199 1200 fn database_inventory_matches(sequence: i64, name: &str, has_extra: bool) -> bool { 1201 crate::all_constraints([sequence == 0, name == "main", !has_extra]) 1202 } 1203 1204 async fn verify_database_metadata( 1205 connection: &mut SqliteConnection, 1206 expected: &ServiceDatabaseMetadata, 1207 ) -> Result<(), ServiceSqliteError> { 1208 let application_id = sqlx::query_scalar::<_, i64>("PRAGMA application_id") 1209 .fetch_one(&mut *connection) 1210 .await 1211 .map_err(metadata_source)?; 1212 let row_count = sqlx::query_scalar::<_, i64>( 1213 "SELECT COUNT(*) FROM (SELECT 1 FROM radroots_service_metadata LIMIT 2)", 1214 ) 1215 .fetch_one(&mut *connection) 1216 .await 1217 .map_err(metadata_source)?; 1218 crate::require_condition( 1219 crate::all_constraints([ 1220 row_count == 1, 1221 application_id == i64::from(expected.application_id().get()), 1222 ]), 1223 ServiceSqliteErrorKind::Metadata, 1224 )?; 1225 let row = sqlx::query( 1226 "SELECT 1227 typeof(service_id) = 'text' AS service_id_type_ok, 1228 length(CAST(service_id AS BLOB)) AS service_id_length, 1229 substr(CAST(service_id AS BLOB), 1, 129) AS service_id_prefix, 1230 typeof(instance_id) = 'text' AS instance_id_type_ok, 1231 length(CAST(instance_id AS BLOB)) AS instance_id_length, 1232 substr(CAST(instance_id AS BLOB), 1, 129) AS instance_id_prefix, 1233 typeof(source_generation) = 'blob' AS source_generation_type_ok, 1234 length(source_generation) AS source_generation_length, 1235 substr(source_generation, 1, 33) AS source_generation_prefix, 1236 CASE WHEN typeof(state_schema_version) = 'integer' 1237 THEN state_schema_version END, 1238 CASE WHEN typeof(created_at_unix_ms) = 'integer' 1239 THEN created_at_unix_ms END 1240 FROM radroots_service_metadata 1241 WHERE singleton = 1 1242 LIMIT 1", 1243 ) 1244 .fetch_optional(&mut *connection) 1245 .await 1246 .map_err(metadata_source)?; 1247 let Some(row) = row else { 1248 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata)); 1249 }; 1250 let service = crate::persisted_value::bounded_utf8( 1251 &row, 1252 "service_id_type_ok", 1253 "service_id_length", 1254 "service_id_prefix", 1255 1, 1256 crate::persisted_value::MAX_IDENTIFIER_UTF8_BYTES, 1257 ) 1258 .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata))?; 1259 let instance = crate::persisted_value::bounded_utf8( 1260 &row, 1261 "instance_id_type_ok", 1262 "instance_id_length", 1263 "instance_id_prefix", 1264 1, 1265 crate::persisted_value::MAX_IDENTIFIER_UTF8_BYTES, 1266 ) 1267 .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata))?; 1268 let generation = crate::persisted_value::bounded_bytes( 1269 &row, 1270 "source_generation_type_ok", 1271 "source_generation_length", 1272 "source_generation_prefix", 1273 32, 1274 32, 1275 ) 1276 .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata))?; 1277 let schema = row.try_get::<Option<i64>, _>(9).map_err(metadata_source)?; 1278 let created_at = row.try_get::<Option<i64>, _>(10).map_err(metadata_source)?; 1279 let (Some(schema), Some(created_at)) = (schema, created_at) else { 1280 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata)); 1281 }; 1282 crate::require_condition( 1283 crate::all_constraints([ 1284 service == expected.service().as_str(), 1285 instance == expected.instance().as_str(), 1286 generation == expected.source_generation().as_bytes(), 1287 schema == i64::from(expected.state_schema_version().get()), 1288 created_at == i64::try_from(expected.created_at_unix_ms()).unwrap_or(-1), 1289 ]), 1290 ServiceSqliteErrorKind::Metadata, 1291 )?; 1292 Ok(()) 1293 } 1294 1295 async fn verify_integrity(connection: &mut SqliteConnection) -> Result<(), ServiceSqliteError> { 1296 let rows = sqlx::query(crate::persisted_value::INTEGRITY_CHECK_SQL) 1297 .fetch_all(&mut *connection) 1298 .await 1299 .map_err(integrity_source)?; 1300 let row = rows 1301 .first() 1302 .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?; 1303 let value = crate::persisted_value::bounded_integrity_bytes(row); 1304 crate::require_condition( 1305 integrity_projection_is_ok(value) && rows.len() == 1, 1306 ServiceSqliteErrorKind::Integrity, 1307 )?; 1308 let violation = sqlx::query_scalar::<_, i64>("SELECT 1 FROM pragma_foreign_key_check LIMIT 1") 1309 .fetch_optional(connection) 1310 .await 1311 .map_err(integrity_source)?; 1312 crate::require_condition(violation.is_none(), ServiceSqliteErrorKind::Integrity)?; 1313 Ok(()) 1314 } 1315 1316 fn integrity_projection_is_ok(value: Option<&[u8]>) -> bool { 1317 value == Some(b"ok") 1318 } 1319 1320 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 1321 enum BackupFailureKind { 1322 InvalidStagingPath, 1323 InvalidStagingParent, 1324 StagingCollision, 1325 CreateStaging, 1326 CreateState, 1327 StagingReplaced, 1328 InvalidStagingInventory, 1329 AlreadyActive, 1330 Capture, 1331 Cancelled, 1332 HashState, 1333 SyncState, 1334 SyncStaging, 1335 SyncParent, 1336 Manifest, 1337 Join, 1338 } 1339 1340 struct BackupFailure { 1341 kind: BackupFailureKind, 1342 source: Option<Box<dyn Error + Send + Sync + 'static>>, 1343 } 1344 1345 impl fmt::Debug for BackupFailure { 1346 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1347 formatter 1348 .debug_struct("BackupFailure") 1349 .field("kind", &self.kind) 1350 .field("source", &self.source.as_ref().map(|_| "[redacted]")) 1351 .finish() 1352 } 1353 } 1354 1355 impl fmt::Display for BackupFailure { 1356 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1357 formatter.write_str(match self.kind { 1358 BackupFailureKind::InvalidStagingPath => "backup staging path is invalid", 1359 BackupFailureKind::InvalidStagingParent => "backup staging parent is invalid", 1360 BackupFailureKind::StagingCollision => "backup staging destination already exists", 1361 BackupFailureKind::CreateStaging => "backup staging directory could not be created", 1362 BackupFailureKind::CreateState => "backup state member could not be created", 1363 BackupFailureKind::StagingReplaced => "backup staging identity changed", 1364 BackupFailureKind::InvalidStagingInventory => "backup staging inventory is invalid", 1365 BackupFailureKind::AlreadyActive => "another backup capture is active", 1366 BackupFailureKind::Capture => "online backup capture failed", 1367 BackupFailureKind::Cancelled => "online backup capture was cancelled", 1368 BackupFailureKind::HashState => "backup state member could not be hashed", 1369 BackupFailureKind::SyncState => "backup state member could not be synchronized", 1370 BackupFailureKind::SyncStaging => "backup staging directory could not be synchronized", 1371 BackupFailureKind::SyncParent => "backup staging parent could not be synchronized", 1372 BackupFailureKind::Manifest => "backup manifest could not be constructed", 1373 BackupFailureKind::Join => "backup worker could not be joined", 1374 }) 1375 } 1376 } 1377 1378 impl Error for BackupFailure { 1379 fn source(&self) -> Option<&(dyn Error + 'static)> { 1380 self.source 1381 .as_deref() 1382 .map(|source| source as &(dyn Error + 'static)) 1383 } 1384 } 1385 1386 fn backup_error(kind: BackupFailureKind) -> ServiceSqliteError { 1387 ServiceSqliteError::with_source( 1388 ServiceSqliteErrorKind::Backup, 1389 BackupFailure { kind, source: None }, 1390 ) 1391 } 1392 1393 fn require_backup_condition( 1394 condition: bool, 1395 kind: BackupFailureKind, 1396 ) -> Result<(), ServiceSqliteError> { 1397 if condition { 1398 Ok(()) 1399 } else { 1400 Err(backup_error(kind)) 1401 } 1402 } 1403 1404 fn backup_source( 1405 kind: BackupFailureKind, 1406 source: impl Error + Send + Sync + 'static, 1407 ) -> ServiceSqliteError { 1408 ServiceSqliteError::with_source( 1409 ServiceSqliteErrorKind::Backup, 1410 BackupFailure { 1411 kind, 1412 source: Some(Box::new(source)), 1413 }, 1414 ) 1415 } 1416 1417 fn integrity_source(source: sqlx::Error) -> ServiceSqliteError { 1418 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Integrity, source) 1419 } 1420 1421 fn metadata_source(source: sqlx::Error) -> ServiceSqliteError { 1422 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Metadata, source) 1423 } 1424 1425 #[cfg(test)] 1426 mod tests { 1427 use core::num::NonZeroU32; 1428 use std::os::unix::fs::{PermissionsExt, symlink}; 1429 1430 use radroots_runtime_paths::{InstanceId, ServiceId}; 1431 use radroots_storage::event::SourceGeneration; 1432 use tempfile::tempdir; 1433 1434 use super::*; 1435 1436 #[test] 1437 fn backup_failure_inventory_is_complete_and_source_aware() { 1438 let cases = [ 1439 ( 1440 BackupFailureKind::InvalidStagingPath, 1441 "backup staging path is invalid", 1442 ), 1443 ( 1444 BackupFailureKind::InvalidStagingParent, 1445 "backup staging parent is invalid", 1446 ), 1447 ( 1448 BackupFailureKind::StagingCollision, 1449 "backup staging destination already exists", 1450 ), 1451 ( 1452 BackupFailureKind::CreateStaging, 1453 "backup staging directory could not be created", 1454 ), 1455 ( 1456 BackupFailureKind::CreateState, 1457 "backup state member could not be created", 1458 ), 1459 ( 1460 BackupFailureKind::StagingReplaced, 1461 "backup staging identity changed", 1462 ), 1463 ( 1464 BackupFailureKind::InvalidStagingInventory, 1465 "backup staging inventory is invalid", 1466 ), 1467 ( 1468 BackupFailureKind::AlreadyActive, 1469 "another backup capture is active", 1470 ), 1471 (BackupFailureKind::Capture, "online backup capture failed"), 1472 ( 1473 BackupFailureKind::Cancelled, 1474 "online backup capture was cancelled", 1475 ), 1476 ( 1477 BackupFailureKind::HashState, 1478 "backup state member could not be hashed", 1479 ), 1480 ( 1481 BackupFailureKind::SyncState, 1482 "backup state member could not be synchronized", 1483 ), 1484 ( 1485 BackupFailureKind::SyncStaging, 1486 "backup staging directory could not be synchronized", 1487 ), 1488 ( 1489 BackupFailureKind::SyncParent, 1490 "backup staging parent could not be synchronized", 1491 ), 1492 ( 1493 BackupFailureKind::Manifest, 1494 "backup manifest could not be constructed", 1495 ), 1496 (BackupFailureKind::Join, "backup worker could not be joined"), 1497 ]; 1498 for (kind, message) in cases { 1499 let plain = BackupFailure { kind, source: None }; 1500 assert_eq!(plain.to_string(), message); 1501 assert!(plain.source().is_none()); 1502 assert!(format!("{plain:?}").contains("source: None")); 1503 1504 let sourced = BackupFailure { 1505 kind, 1506 source: Some(Box::new(std::io::Error::other("private-cause"))), 1507 }; 1508 assert_eq!(sourced.to_string(), message); 1509 assert_eq!( 1510 sourced.source().expect("source").to_string(), 1511 "private-cause" 1512 ); 1513 let debug = format!("{sourced:?}"); 1514 assert!(debug.contains("[redacted]")); 1515 assert!(!debug.contains("private-cause")); 1516 assert!(require_backup_condition(true, kind).is_ok()); 1517 assert_eq!( 1518 require_backup_condition(false, kind) 1519 .expect_err("false condition") 1520 .kind(), 1521 ServiceSqliteErrorKind::Backup 1522 ); 1523 } 1524 } 1525 1526 #[test] 1527 fn staging_path_rejects_relative_parent_and_oversize_inputs() { 1528 assert!(StagingPath::new(Path::new("relative/stage")).is_err()); 1529 assert!(StagingPath::new(Path::new("/tmp/../stage")).is_err()); 1530 assert!(StagingPath::new(Path::new("/")).is_err()); 1531 let large = format!("/tmp/{}", "x".repeat(MAX_STAGING_PATH_BYTES)); 1532 assert!(StagingPath::new(Path::new(&large)).is_err()); 1533 } 1534 1535 #[test] 1536 fn staging_guard_creates_exact_modes_and_cleans_uncommitted_state() { 1537 let root = tempdir().expect("root"); 1538 std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700)) 1539 .expect("parent mode"); 1540 let path = StagingPath::new(&root.path().join("backup-stage")).expect("path"); 1541 { 1542 let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging"); 1543 staging.validate().expect("validate"); 1544 assert_eq!( 1545 std::fs::metadata(&path.full) 1546 .expect("directory metadata") 1547 .permissions() 1548 .mode() 1549 & 0o777, 1550 0o700 1551 ); 1552 assert_eq!( 1553 std::fs::metadata(staging.state_path()) 1554 .expect("state metadata") 1555 .permissions() 1556 .mode() 1557 & 0o777, 1558 0o600 1559 ); 1560 } 1561 assert!(!path.full.exists()); 1562 } 1563 1564 #[test] 1565 fn staging_guard_success_inventory_hash_sync_commit_and_empty_failure_are_exercised() { 1566 let root = tempdir().expect("root"); 1567 std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700)) 1568 .expect("parent mode"); 1569 1570 let empty_path = StagingPath::new(&root.path().join("empty-stage")).expect("empty path"); 1571 let empty = StagingGuard::create(&empty_path, &SystemCaptureOperations).expect("empty"); 1572 assert!(empty.hash_state(&AtomicBool::new(false)).is_err()); 1573 drop(empty); 1574 1575 let path = StagingPath::new(&root.path().join("complete-stage")).expect("path"); 1576 let mut staging = 1577 StagingGuard::create(&path, &SystemCaptureOperations).expect("complete stage"); 1578 std::fs::write(staging.state_path(), b"captured-state").expect("state bytes"); 1579 staging.validate_inventory().expect("singleton inventory"); 1580 let (length, digest) = staging 1581 .hash_state(&AtomicBool::new(false)) 1582 .expect("hash state"); 1583 assert_eq!(length, 14); 1584 assert_eq!(digest, Sha256::digest(b"captured-state").as_slice()); 1585 staging 1586 .sync_state(&SystemCaptureOperations) 1587 .expect("sync state"); 1588 staging 1589 .sync_directories(&SystemCaptureOperations) 1590 .expect("sync directories"); 1591 staging.commit(); 1592 drop(staging); 1593 assert!(path.full.exists()); 1594 std::fs::remove_file(path.full.join(STATE_FILE_NAME)).expect("remove state"); 1595 std::fs::remove_dir(path.full).expect("remove stage"); 1596 } 1597 1598 #[test] 1599 fn staging_inventory_and_sidecar_cleanup_bind_exact_entries() { 1600 let root = tempdir().expect("root"); 1601 std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700)) 1602 .expect("parent mode"); 1603 let path = StagingPath::new(&root.path().join("backup-stage")).expect("path"); 1604 let mut staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging"); 1605 std::fs::write(staging.state_path(), b"state").expect("state bytes"); 1606 let wal = path.full.join(KNOWN_SIDECARS[0]); 1607 std::fs::write(&wal, b"sidecar").expect("sidecar"); 1608 std::fs::set_permissions(&wal, std::fs::Permissions::from_mode(0o600)) 1609 .expect("sidecar mode"); 1610 assert!(staging.validate_inventory().is_err()); 1611 staging.record_sidecars(); 1612 drop(staging); 1613 assert!(!path.full.exists()); 1614 1615 let path = StagingPath::new(&root.path().join("unsafe-sidecar-stage")).expect("path"); 1616 let mut staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging"); 1617 std::fs::write(staging.state_path(), b"state").expect("state bytes"); 1618 let wal = path.full.join(KNOWN_SIDECARS[0]); 1619 std::fs::write(&wal, b"unsafe-sidecar").expect("sidecar"); 1620 std::fs::set_permissions(&wal, std::fs::Permissions::from_mode(0o666)) 1621 .expect("unsafe sidecar mode"); 1622 let outside_wal = root.path().join("unsafe-sidecar-alias"); 1623 std::fs::hard_link(&wal, &outside_wal).expect("make sidecar unsafe by link count"); 1624 staging.record_sidecars(); 1625 drop(staging); 1626 assert_eq!( 1627 std::fs::read(wal).expect("unsafe sidecar preserved"), 1628 b"unsafe-sidecar" 1629 ); 1630 assert_eq!( 1631 std::fs::read(outside_wal).expect("outside sidecar link preserved"), 1632 b"unsafe-sidecar" 1633 ); 1634 } 1635 1636 #[test] 1637 fn hash_cancellation_partial_cleanup_and_directory_replacement_are_exact() { 1638 let root = tempdir().expect("root"); 1639 std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700)) 1640 .expect("parent mode"); 1641 1642 let cancelled_path = StagingPath::new(&root.path().join("cancelled-stage")).expect("path"); 1643 let cancelled = 1644 StagingGuard::create(&cancelled_path, &SystemCaptureOperations).expect("staging"); 1645 std::fs::write(cancelled.state_path(), b"state").expect("state bytes"); 1646 assert!(cancelled.hash_state(&AtomicBool::new(true)).is_err()); 1647 drop(cancelled); 1648 1649 let exact_path = StagingPath::new(&root.path().join("exact-stage")).expect("path"); 1650 let exact = StagingGuard::create(&exact_path, &SystemCaptureOperations).expect("staging"); 1651 cleanup_partial_staging( 1652 &exact.parent, 1653 &exact.path.name, 1654 Some(exact.directory_identity), 1655 Some(&exact.directory), 1656 Some(exact.state_identity), 1657 ); 1658 assert!(!exact_path.full.exists()); 1659 drop(exact); 1660 1661 let replaced_path = StagingPath::new(&root.path().join("replaced-stage")).expect("path"); 1662 let replaced = 1663 StagingGuard::create(&replaced_path, &SystemCaptureOperations).expect("staging"); 1664 let retired = root.path().join("retired-stage"); 1665 std::fs::rename(&replaced_path.full, &retired).expect("retire governed directory"); 1666 std::fs::create_dir(&replaced_path.full).expect("replacement directory"); 1667 std::fs::set_permissions(&replaced_path.full, std::fs::Permissions::from_mode(0o700)) 1668 .expect("replacement mode"); 1669 std::fs::write(replaced_path.full.join("foreign"), b"foreign").expect("foreign entry"); 1670 drop(replaced); 1671 assert_eq!( 1672 std::fs::read(replaced_path.full.join("foreign")).expect("replacement survives"), 1673 b"foreign" 1674 ); 1675 } 1676 1677 #[test] 1678 fn staging_guard_rejects_collision_and_preserves_replacement() { 1679 let root = tempdir().expect("root"); 1680 std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700)) 1681 .expect("parent mode"); 1682 let full = root.path().join("backup-stage"); 1683 std::fs::create_dir(&full).expect("collision"); 1684 assert!( 1685 StagingGuard::create( 1686 &StagingPath::new(&full).expect("path"), 1687 &SystemCaptureOperations, 1688 ) 1689 .is_err() 1690 ); 1691 std::fs::remove_dir(&full).expect("remove collision"); 1692 1693 let path = StagingPath::new(&full).expect("path"); 1694 let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging"); 1695 let original = staging.state_path(); 1696 let replacement = full.join("replacement"); 1697 std::fs::write(&replacement, b"foreign").expect("replacement"); 1698 std::fs::set_permissions(&replacement, std::fs::Permissions::from_mode(0o600)) 1699 .expect("replacement mode"); 1700 std::fs::rename(&replacement, &original).expect("replace state"); 1701 drop(staging); 1702 assert_eq!(std::fs::read(&original).expect("preserved"), b"foreign"); 1703 assert!(full.exists()); 1704 } 1705 1706 #[test] 1707 fn staging_guard_rejects_insecure_or_symlinked_parent_without_mutation() { 1708 let root = tempdir().expect("root"); 1709 std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700)) 1710 .expect("root mode"); 1711 let insecure = root.path().join("insecure"); 1712 std::fs::create_dir(&insecure).expect("insecure parent"); 1713 std::fs::set_permissions(&insecure, std::fs::Permissions::from_mode(0o770)) 1714 .expect("insecure mode"); 1715 let insecure_stage = insecure.join("stage"); 1716 let insecure_error = match StagingGuard::create( 1717 &StagingPath::new(&insecure_stage).expect("insecure staging path"), 1718 &SystemCaptureOperations, 1719 ) { 1720 Ok(_) => panic!("group-writable parent must be rejected"), 1721 Err(error) => error, 1722 }; 1723 assert_eq!(insecure_error.kind(), ServiceSqliteErrorKind::Backup); 1724 assert!(!insecure_stage.exists()); 1725 1726 let parent = root.path().join("real-parent"); 1727 std::fs::create_dir(&parent).expect("real parent"); 1728 std::fs::set_permissions(&parent, std::fs::Permissions::from_mode(0o700)) 1729 .expect("real parent mode"); 1730 let alias = root.path().join("parent-alias"); 1731 symlink(&parent, &alias).expect("parent symlink"); 1732 let symlink_stage = alias.join("stage"); 1733 assert!( 1734 StagingGuard::create( 1735 &StagingPath::new(&symlink_stage).expect("symlink staging path"), 1736 &SystemCaptureOperations, 1737 ) 1738 .is_err() 1739 ); 1740 assert!(!parent.join("stage").exists()); 1741 } 1742 1743 #[test] 1744 fn staging_guard_rejects_hardlinked_state_and_preserves_other_link() { 1745 let root = tempdir().expect("root"); 1746 std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700)) 1747 .expect("parent mode"); 1748 let path = StagingPath::new(&root.path().join("backup-stage")).expect("path"); 1749 let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging"); 1750 std::fs::write(staging.state_path(), b"captured").expect("state bytes"); 1751 let outside = root.path().join("outside-link"); 1752 std::fs::hard_link(staging.state_path(), &outside).expect("hard link"); 1753 assert!(staging.validate().is_err()); 1754 drop(staging); 1755 assert_eq!( 1756 std::fs::read(&outside).expect("other link survives"), 1757 b"captured" 1758 ); 1759 assert!(!path.full.exists()); 1760 } 1761 1762 #[test] 1763 fn partial_cleanup_preserves_replaced_directory_and_state_entries() { 1764 let root = tempdir().expect("root"); 1765 std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700)) 1766 .expect("root mode"); 1767 let parent = File::from( 1768 open( 1769 root.path(), 1770 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1771 Mode::empty(), 1772 ) 1773 .expect("open parent"), 1774 ); 1775 mkdirat(&parent, "stage", Mode::RUSR | Mode::WUSR | Mode::XUSR) 1776 .expect("create original stage"); 1777 let original_identity = 1778 created_directory_identity(&parent, OsStr::new("stage")).expect("original identity"); 1779 let original = File::from( 1780 openat( 1781 &parent, 1782 "stage", 1783 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, 1784 Mode::empty(), 1785 ) 1786 .expect("open original stage"), 1787 ); 1788 std::fs::rename(root.path().join("stage"), root.path().join("retired")) 1789 .expect("retire original stage"); 1790 std::fs::create_dir(root.path().join("stage")).expect("replacement stage"); 1791 std::fs::set_permissions( 1792 root.path().join("stage"), 1793 std::fs::Permissions::from_mode(0o700), 1794 ) 1795 .expect("replacement mode"); 1796 std::fs::write(root.path().join("stage/foreign"), b"foreign").expect("replacement member"); 1797 cleanup_partial_staging( 1798 &parent, 1799 OsStr::new("stage"), 1800 Some(original_identity), 1801 Some(&original), 1802 None, 1803 ); 1804 assert_eq!( 1805 std::fs::read(root.path().join("stage/foreign")).expect("foreign survives"), 1806 b"foreign" 1807 ); 1808 1809 let path = StagingPath::new(&root.path().join("state-stage")).expect("path"); 1810 let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging"); 1811 let retired_state = path.full.join("retired-state"); 1812 std::fs::rename(staging.state_path(), &retired_state).expect("retire state"); 1813 std::fs::write(staging.state_path(), b"foreign-state").expect("foreign state"); 1814 std::fs::set_permissions(staging.state_path(), std::fs::Permissions::from_mode(0o600)) 1815 .expect("foreign state mode"); 1816 cleanup_partial_staging( 1817 &staging.parent, 1818 &staging.path.name, 1819 Some(staging.directory_identity), 1820 Some(&staging.directory), 1821 Some(staging.state_identity), 1822 ); 1823 assert_eq!( 1824 std::fs::read(staging.state_path()).expect("replacement state survives"), 1825 b"foreign-state" 1826 ); 1827 drop(staging); 1828 assert!(path.full.exists()); 1829 } 1830 1831 #[test] 1832 fn integrity_projection_bounds_corrupt_text_before_semantic_acceptance() { 1833 assert!(integrity_projection_is_ok(Some(b"ok"))); 1834 assert!(!integrity_projection_is_ok(Some(b""))); 1835 assert!(!integrity_projection_is_ok(Some(b"not-ok"))); 1836 assert!(!integrity_projection_is_ok(None)); 1837 } 1838 1839 #[test] 1840 fn capture_database_inventory_rejects_each_independent_projection_drift() { 1841 assert!(database_inventory_matches(0, "main", false)); 1842 assert!(!database_inventory_matches(1, "main", false)); 1843 assert!(!database_inventory_matches(0, "temp", false)); 1844 assert!(!database_inventory_matches(0, "main", true)); 1845 } 1846 1847 async fn metadata_fixture() -> (SqliteConnection, ServiceDatabaseMetadata) { 1848 let metadata = ServiceDatabaseMetadata::from_verified_backup( 1849 ServiceId::new("myc").expect("service"), 1850 InstanceId::new("primary").expect("instance"), 1851 SourceGeneration::new([7; 32]).expect("generation"), 1852 NonZeroU32::new(1).expect("schema"), 1853 1_234, 1854 crate::ServiceSqliteApplicationId::new(0x5244_5254).expect("application ID"), 1855 ) 1856 .expect("metadata"); 1857 let mut connection = SqliteConnection::connect("sqlite::memory:") 1858 .await 1859 .expect("database"); 1860 let schema = format!( 1861 "PRAGMA application_id = {}; 1862 CREATE TABLE radroots_service_metadata ( 1863 singleton INTEGER, 1864 service_id TEXT, 1865 instance_id TEXT, 1866 source_generation BLOB, 1867 state_schema_version INTEGER, 1868 created_at_unix_ms INTEGER 1869 );", 1870 metadata.application_id().get() 1871 ); 1872 sqlx::raw_sql(sqlx::AssertSqlSafe(schema.as_str())) 1873 .execute(&mut connection) 1874 .await 1875 .expect("metadata schema"); 1876 sqlx::query("INSERT INTO radroots_service_metadata VALUES (1, ?, ?, ?, ?, ?)") 1877 .bind(metadata.service().as_str()) 1878 .bind(metadata.instance().as_str()) 1879 .bind(metadata.source_generation().as_bytes().as_slice()) 1880 .bind(i64::from(metadata.state_schema_version().get())) 1881 .bind(i64::try_from(metadata.created_at_unix_ms()).expect("time")) 1882 .execute(&mut connection) 1883 .await 1884 .expect("metadata row"); 1885 (connection, metadata) 1886 } 1887 1888 #[tokio::test(flavor = "current_thread")] 1889 async fn capture_database_inventory_metadata_and_integrity_accept_exact_state() { 1890 let (mut connection, metadata) = metadata_fixture().await; 1891 verify_database_inventory(&mut connection) 1892 .await 1893 .expect("main-only inventory"); 1894 verify_database_metadata(&mut connection, &metadata) 1895 .await 1896 .expect("exact metadata"); 1897 verify_integrity(&mut connection) 1898 .await 1899 .expect("healthy database"); 1900 1901 sqlx::query("ATTACH DATABASE ':memory:' AS extra") 1902 .execute(&mut connection) 1903 .await 1904 .expect("attach extra"); 1905 assert!(verify_database_inventory(&mut connection).await.is_err()); 1906 } 1907 1908 #[tokio::test(flavor = "current_thread")] 1909 async fn capture_metadata_rejects_every_independent_identity_drift() { 1910 for statement in [ 1911 "PRAGMA application_id = 1", 1912 "INSERT INTO radroots_service_metadata SELECT 2, service_id, instance_id, source_generation, state_schema_version, created_at_unix_ms FROM radroots_service_metadata", 1913 "UPDATE radroots_service_metadata SET service_id = 'rhi'", 1914 "UPDATE radroots_service_metadata SET instance_id = 'secondary'", 1915 "UPDATE radroots_service_metadata SET source_generation = zeroblob(32)", 1916 "UPDATE radroots_service_metadata SET state_schema_version = 2", 1917 "UPDATE radroots_service_metadata SET created_at_unix_ms = 1235", 1918 "UPDATE radroots_service_metadata SET service_id = NULL", 1919 ] { 1920 let (mut connection, metadata) = metadata_fixture().await; 1921 sqlx::raw_sql(statement) 1922 .execute(&mut connection) 1923 .await 1924 .expect("apply drift"); 1925 assert!( 1926 verify_database_metadata(&mut connection, &metadata) 1927 .await 1928 .is_err(), 1929 "drift must fail: {statement}" 1930 ); 1931 } 1932 } 1933 1934 #[tokio::test(flavor = "current_thread")] 1935 async fn capture_integrity_rejects_foreign_key_violations() { 1936 let mut connection = SqliteConnection::connect("sqlite::memory:") 1937 .await 1938 .expect("database"); 1939 sqlx::raw_sql( 1940 "PRAGMA foreign_keys = OFF; 1941 CREATE TABLE parent(id INTEGER PRIMARY KEY); 1942 CREATE TABLE child(parent_id INTEGER REFERENCES parent(id)); 1943 INSERT INTO child(parent_id) VALUES (41);", 1944 ) 1945 .execute(&mut connection) 1946 .await 1947 .expect("foreign-key violation fixture"); 1948 let error = verify_integrity(&mut connection) 1949 .await 1950 .expect_err("foreign-key drift must fail"); 1951 assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); 1952 } 1953 1954 #[test] 1955 fn active_capture_permit_is_exclusive_and_recoverable() { 1956 let active = Arc::new(AtomicBool::new(false)); 1957 let first = CapturePermit::acquire(Arc::clone(&active)).expect("first permit"); 1958 assert!(CapturePermit::acquire(Arc::clone(&active)).is_err()); 1959 drop(first); 1960 assert!(CapturePermit::acquire(active).is_ok()); 1961 } 1962 }