process_tests.rs (19128B)
1 use core::num::{NonZeroU32, NonZeroU64}; 2 use std::{ 3 env, fs, 4 io::{BufRead, BufReader, Read, Write}, 5 os::unix::{ 6 fs::{MetadataExt, PermissionsExt}, 7 process::ExitStatusExt, 8 }, 9 path::{Path, PathBuf}, 10 process::{Child, Command, ExitStatus, Stdio}, 11 sync::mpsc, 12 thread, 13 time::{Duration, Instant}, 14 }; 15 16 use radroots_runtime_paths::{ 17 InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, 18 RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, ServiceId, 19 }; 20 use radroots_storage::event::SourceGeneration; 21 use sha2::{Digest, Sha256}; 22 use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions}; 23 24 use super::{ 25 BACKUP_FILE_NAME, MARKER_FILE_NAME, MARKER_NEXT_FILE_NAME, STAGED_FILE_NAME, 26 finalize::test_finalize_with_failpoint, 27 }; 28 use crate::failpoint::{DurabilityFailpoint, DurabilityFailpoints}; 29 use crate::{ 30 BackupCreatedAtUnixMs, BackupMemberSha256, MigrationAppliedAtUnixSeconds, 31 MigrationBuildIdentity, MigrationCatalog, OpenMode, SchemaCatalog, SchemaVersionCatalog, 32 ServiceBackupManifest, ServiceDatabaseIdentity, ServiceDatabaseMetadata, 33 ServiceSqliteApplicationId, ServiceSqliteConnectionOptions, ServiceSqliteErrorKind, 34 ServiceSqliteHost, ServiceSqlitePaths, initialize_database, stage_verified_restore, 35 verify_backup_bundle, 36 }; 37 38 const PROCESS_READY: &str = "RSHR_STEP073_READY"; 39 const MANIFEST_FILE_NAME: &str = "process-restore-manifest.v1.json"; 40 const BUNDLE_DIRECTORY_NAME: &str = "process-restore-bundle"; 41 const MAX_CHILD_INPUT_BYTES: u64 = 4_096; 42 const REPLACEMENT_USER_VERSION: i64 = 73; 43 const PROCESS_TIMEOUT: Duration = Duration::from_secs(30); 44 45 struct Fixture { 46 root: tempfile::TempDir, 47 paths: ServiceSqlitePaths, 48 identity: ServiceDatabaseIdentity, 49 migrations: MigrationCatalog, 50 schema: SchemaCatalog, 51 } 52 53 impl Fixture { 54 async fn new() -> Self { 55 let root = tempfile::tempdir().expect("process fixture root"); 56 let paths = paths(root.path()); 57 let state_directory = paths.state_database().parent().expect("state directory"); 58 fs::create_dir_all(state_directory).expect("create state directory"); 59 fs::set_permissions(state_directory, fs::Permissions::from_mode(0o700)) 60 .expect("restrict state directory"); 61 let metadata = metadata(&paths); 62 let (migrations, schema) = catalogs(); 63 let mut authority = 64 initialize_database(&paths, OpenMode::Initialize, &metadata, &schema, |_| { 65 Box::pin(async move { Ok::<(), sqlx::Error>(()) }) 66 }) 67 .await 68 .expect("initialize process database"); 69 authority 70 .release() 71 .expect("release initialization authority"); 72 { 73 let mut connection = sqlx::SqliteConnection::connect_with( 74 &SqliteConnectOptions::new() 75 .filename(paths.state_database()) 76 .create_if_missing(false) 77 .disable_statement_logging(), 78 ) 79 .await 80 .expect("open live database for WAL posture"); 81 sqlx::query("PRAGMA journal_mode = WAL") 82 .execute(&mut connection) 83 .await 84 .expect("set WAL posture"); 85 sqlx::query("PRAGMA wal_checkpoint(TRUNCATE)") 86 .execute(&mut connection) 87 .await 88 .expect("checkpoint WAL posture"); 89 connection.close().await.expect("close live database"); 90 } 91 92 let bundle = root.path().join(BUNDLE_DIRECTORY_NAME); 93 fs::create_dir(&bundle).expect("create process bundle"); 94 fs::set_permissions(&bundle, fs::Permissions::from_mode(0o700)) 95 .expect("restrict process bundle"); 96 let member = bundle.join(crate::BACKUP_STATE_MEMBER_NAME); 97 fs::copy(paths.state_database(), &member).expect("copy process member"); 98 fs::set_permissions(&member, fs::Permissions::from_mode(0o600)) 99 .expect("restrict process member"); 100 { 101 let mut connection = sqlx::SqliteConnection::connect_with( 102 &SqliteConnectOptions::new() 103 .filename(&member) 104 .create_if_missing(false) 105 .disable_statement_logging(), 106 ) 107 .await 108 .expect("open process member"); 109 let user_version = format!("PRAGMA user_version = {REPLACEMENT_USER_VERSION}"); 110 sqlx::query(sqlx::AssertSqlSafe(user_version.as_str())) 111 .execute(&mut connection) 112 .await 113 .expect("set replacement probe"); 114 sqlx::query("PRAGMA wal_checkpoint(TRUNCATE)") 115 .execute(&mut connection) 116 .await 117 .expect("checkpoint replacement probe"); 118 connection.close().await.expect("close process member"); 119 } 120 let bytes = fs::read(&member).expect("read process member"); 121 let manifest = ServiceBackupManifest::from_capture( 122 &metadata, 123 BackupCreatedAtUnixMs::new(1_700_000_073_000).expect("capture time"), 124 u64::try_from(bytes.len()).expect("member length"), 125 BackupMemberSha256::from_bytes(Sha256::digest(&bytes).into()), 126 ) 127 .expect("process manifest"); 128 fs::write( 129 root.path().join(MANIFEST_FILE_NAME), 130 manifest.canonical_bytes(), 131 ) 132 .expect("write process manifest"); 133 134 let identity = metadata.identity(); 135 Self { 136 root, 137 paths, 138 identity, 139 migrations, 140 schema, 141 } 142 } 143 144 fn root(&self) -> &Path { 145 self.root.path() 146 } 147 } 148 149 #[tokio::test(flavor = "current_thread")] 150 #[ignore = "private process child invoked by the parent crash test"] 151 async fn child_before_prepared_marker() { 152 run_child(DurabilityFailpoint::MarkerBeforeCreate, 1).await; 153 } 154 155 #[tokio::test(flavor = "current_thread")] 156 #[ignore = "private process child invoked by the parent crash test"] 157 async fn child_after_prepared_marker() { 158 run_child(DurabilityFailpoint::MarkerAfterDirectorySync, 1).await; 159 } 160 161 #[tokio::test(flavor = "current_thread")] 162 #[ignore = "private process child invoked by the parent crash test"] 163 async fn child_after_live_retained_scratch_sync() { 164 run_child(DurabilityFailpoint::MarkerAdvanceAfterWriteAndFileSync, 1).await; 165 } 166 167 #[tokio::test(flavor = "current_thread")] 168 #[ignore = "private process child invoked by the parent crash test"] 169 async fn child_after_replacement_install_sync() { 170 run_child(DurabilityFailpoint::RestoreAfterInstallStageSync, 1).await; 171 } 172 173 #[tokio::test(flavor = "current_thread")] 174 #[ignore = "private process child invoked by the parent crash test"] 175 async fn child_after_terminal_marker_sync() { 176 run_child(DurabilityFailpoint::MarkerAdvanceAfterDirectorySync, 2).await; 177 } 178 179 async fn run_child(point: DurabilityFailpoint, occurrence: u8) { 180 let root = child_root_from_stdin(); 181 rustix::process::umask(rustix::fs::Mode::empty()); 182 let paths = paths(&root); 183 let metadata = metadata(&paths); 184 let identity = metadata.identity(); 185 let (migrations, schema) = catalogs(); 186 let manifest = ServiceBackupManifest::from_canonical_bytes( 187 &fs::read(root.join(MANIFEST_FILE_NAME)).expect("read child manifest"), 188 ) 189 .expect("parse child manifest"); 190 let verified = verify_backup_bundle( 191 manifest.canonical_bytes(), 192 manifest.digest(), 193 &root.join(BUNDLE_DIRECTORY_NAME), 194 &identity, 195 NonZeroU64::new(16 * 1_024 * 1_024).expect("member limit"), 196 ) 197 .expect("verify child bundle"); 198 let staged = stage_verified_restore(&paths, &identity, &migrations, &schema, verified) 199 .await 200 .expect("stage child restore"); 201 let failpoints = DurabilityFailpoints::process_barrier(point, occurrence); 202 let result = test_finalize_with_failpoint(staged, failpoints).await; 203 panic!("process durability barrier returned before SIGKILL: {result:?}"); 204 } 205 206 #[tokio::test(flavor = "current_thread")] 207 async fn sigkill_restore_boundaries_recover_exact_topologies_and_preserve_permissions() { 208 for scenario in [ 209 Scenario::unresolved("child_before_prepared_marker"), 210 Scenario::recovered("child_after_prepared_marker", 0), 211 Scenario::recovered( 212 "child_after_live_retained_scratch_sync", 213 REPLACEMENT_USER_VERSION, 214 ), 215 Scenario::recovered( 216 "child_after_replacement_install_sync", 217 REPLACEMENT_USER_VERSION, 218 ), 219 Scenario::recovered("child_after_terminal_marker_sync", REPLACEMENT_USER_VERSION), 220 ] { 221 run_parent_scenario(scenario).await; 222 } 223 } 224 225 struct Scenario { 226 child_name: &'static str, 227 expected_user_version: Option<i64>, 228 } 229 230 impl Scenario { 231 const fn unresolved(child_name: &'static str) -> Self { 232 Self { 233 child_name, 234 expected_user_version: None, 235 } 236 } 237 238 const fn recovered(child_name: &'static str, expected_user_version: i64) -> Self { 239 Self { 240 child_name, 241 expected_user_version: Some(expected_user_version), 242 } 243 } 244 } 245 246 async fn run_parent_scenario(scenario: Scenario) { 247 let fixture = Fixture::new().await; 248 let original_live = file_snapshot(fixture.paths.state_database()); 249 let mut child = spawn_restore_child(scenario.child_name, fixture.root()); 250 wait_for_ready(&mut child); 251 assert_recovery_permissions(&fixture.paths); 252 rustix::process::kill_process( 253 rustix::process::Pid::from_child(child.child()), 254 rustix::process::Signal::KILL, 255 ) 256 .expect("SIGKILL restore child"); 257 let status = child.wait().expect("wait for restore child"); 258 assert_eq!(status.signal(), Some(9), "scenario {}", scenario.child_name); 259 assert_recovery_permissions(&fixture.paths); 260 261 if let Some(expected_user_version) = scenario.expected_user_version { 262 let (host, outcome) = ServiceSqliteHost::open_read_write_existing( 263 &fixture.paths, 264 &fixture.identity, 265 &fixture.migrations, 266 &fixture.schema, 267 ServiceSqliteConnectionOptions::reviewed(), 268 MigrationAppliedAtUnixSeconds::new(1_700_000_073).expect("migration time"), 269 &build_identity(), 270 &[], 271 ) 272 .await 273 .expect("reopen and reconcile interrupted restore"); 274 assert_eq!(outcome.applied_count(), 0); 275 host.close().await.expect("close recovered host"); 276 assert_no_recovery_evidence(&fixture.paths); 277 assert_eq!( 278 database_user_version(&fixture.paths).await, 279 expected_user_version 280 ); 281 assert_live_permissions(&fixture.paths); 282 } else { 283 let staged = recovery_path(&fixture.paths, STAGED_FILE_NAME); 284 let staged_before = file_snapshot(&staged); 285 let error = ServiceSqliteHost::open_read_write_existing( 286 &fixture.paths, 287 &fixture.identity, 288 &fixture.migrations, 289 &fixture.schema, 290 ServiceSqliteConnectionOptions::reviewed(), 291 MigrationAppliedAtUnixSeconds::new(1_700_000_073).expect("migration time"), 292 &build_identity(), 293 &[], 294 ) 295 .await 296 .expect_err("orphan stage without a marker must fail closed"); 297 assert_eq!(error.kind(), ServiceSqliteErrorKind::Recovery); 298 assert_eq!(file_snapshot(fixture.paths.state_database()), original_live); 299 assert_eq!(file_snapshot(&staged), staged_before); 300 assert!(!recovery_path(&fixture.paths, MARKER_FILE_NAME).exists()); 301 } 302 } 303 304 fn spawn_restore_child(child_name: &str, root: &Path) -> KillOnDrop { 305 let exact_name = format!("restore::process_tests::{child_name}"); 306 let mut child = Command::new(env::current_exe().expect("test executable")) 307 .arg("--ignored") 308 .arg("--exact") 309 .arg(exact_name) 310 .arg("--nocapture") 311 .stdin(Stdio::piped()) 312 .stdout(Stdio::piped()) 313 .spawn() 314 .expect("spawn restore child"); 315 writeln!( 316 child.stdin.take().expect("child stdin"), 317 "{}", 318 root.display() 319 ) 320 .expect("send child root"); 321 let ready = readiness_receiver(child.stdout.take().expect("child stdout")); 322 KillOnDrop::new(child, ready) 323 } 324 325 fn wait_for_ready(child: &mut KillOnDrop) { 326 let deadline = Instant::now() + PROCESS_TIMEOUT; 327 loop { 328 let remaining = deadline.saturating_duration_since(Instant::now()); 329 match child.ready.recv_timeout(remaining) { 330 Ok(line) if line == PROCESS_READY => return, 331 Ok(_) => {} 332 Err(error) => { 333 let status = child.child.try_wait().expect("query child status"); 334 panic!("restore child readiness failed ({error}); status: {status:?}"); 335 } 336 } 337 } 338 } 339 340 fn readiness_receiver(stdout: impl Read + Send + 'static) -> mpsc::Receiver<String> { 341 let (sender, receiver) = mpsc::channel(); 342 thread::spawn(move || { 343 for line in BufReader::new(stdout).lines() { 344 if sender.send(line.unwrap_or_default()).is_err() { 345 break; 346 } 347 } 348 }); 349 receiver 350 } 351 352 fn child_root_from_stdin() -> PathBuf { 353 let mut input = String::new(); 354 std::io::stdin() 355 .take(MAX_CHILD_INPUT_BYTES + 1) 356 .read_to_string(&mut input) 357 .expect("read child root"); 358 assert!( 359 input.len() <= MAX_CHILD_INPUT_BYTES as usize, 360 "bounded child input" 361 ); 362 let root = input 363 .strip_suffix('\n') 364 .expect("newline-terminated child root"); 365 assert!( 366 !root.is_empty() && !root.contains('\n'), 367 "single child root" 368 ); 369 PathBuf::from(root) 370 } 371 372 fn paths(root: &Path) -> ServiceSqlitePaths { 373 let context = RuntimeContext::resolve( 374 &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), 375 RuntimeContextBootstrap::new( 376 RadrootsPathProfile::RepoLocal, 377 Some(root.to_path_buf()), 378 RuntimeContextSource::BootstrapCli, 379 RuntimeContextSource::BootstrapCli, 380 ) 381 .expect("bootstrap"), 382 ServiceId::new("myc").expect("service"), 383 InstanceId::new("process-recovery").expect("instance"), 384 ) 385 .expect("runtime context"); 386 ServiceSqlitePaths::from_runtime_context(&context).expect("SQLite paths") 387 } 388 389 fn metadata(paths: &ServiceSqlitePaths) -> ServiceDatabaseMetadata { 390 ServiceDatabaseMetadata::new( 391 paths, 392 SourceGeneration::new([7; 32]).expect("generation"), 393 NonZeroU32::new(1).expect("schema"), 394 1_700_000_000_000, 395 ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), 396 ) 397 .expect("metadata") 398 } 399 400 fn catalogs() -> (MigrationCatalog, SchemaCatalog) { 401 let migrations = MigrationCatalog::new([]).expect("migration catalog"); 402 let digest = SchemaVersionCatalog::computed_digest(1, []).expect("schema digest"); 403 let version = SchemaVersionCatalog::new(1, [], digest).expect("schema version"); 404 let schema = SchemaCatalog::new(&migrations, [version]).expect("schema catalog"); 405 (migrations, schema) 406 } 407 408 fn build_identity() -> MigrationBuildIdentity { 409 MigrationBuildIdentity::new( 410 "1.0.0", 411 "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", 412 "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", 413 "rustc-1.97.0", 414 "process-test", 415 "test", 416 1, 417 1, 418 1, 419 1, 420 1, 421 ) 422 .expect("build identity") 423 } 424 425 fn recovery_path(paths: &ServiceSqlitePaths, name: &str) -> PathBuf { 426 paths 427 .state_database() 428 .parent() 429 .expect("state directory") 430 .join(name) 431 } 432 433 fn assert_recovery_permissions(paths: &ServiceSqlitePaths) { 434 let state_directory = paths.state_database().parent().expect("state directory"); 435 let directory = fs::metadata(state_directory).expect("state directory metadata"); 436 assert!(directory.is_dir()); 437 assert_eq!(directory.mode() & 0o777, 0o700); 438 assert_eq!(directory.uid(), rustix::process::geteuid().as_raw()); 439 for path in [ 440 paths.state_database().to_path_buf(), 441 paths.state_lock().to_path_buf(), 442 recovery_path(paths, STAGED_FILE_NAME), 443 recovery_path(paths, BACKUP_FILE_NAME), 444 recovery_path(paths, MARKER_FILE_NAME), 445 recovery_path(paths, MARKER_NEXT_FILE_NAME), 446 ] { 447 if path.exists() { 448 assert_file_permissions(&path); 449 } 450 } 451 } 452 453 fn assert_live_permissions(paths: &ServiceSqlitePaths) { 454 assert_file_permissions(paths.state_database()); 455 assert_file_permissions(paths.state_lock()); 456 } 457 458 fn assert_file_permissions(path: &Path) { 459 let metadata = fs::symlink_metadata(path).expect("artifact metadata"); 460 assert!(metadata.file_type().is_file()); 461 assert_eq!(metadata.mode() & 0o777, 0o600); 462 assert_eq!(metadata.nlink(), 1); 463 assert_eq!(metadata.uid(), rustix::process::geteuid().as_raw()); 464 } 465 466 fn assert_no_recovery_evidence(paths: &ServiceSqlitePaths) { 467 for name in [ 468 STAGED_FILE_NAME, 469 BACKUP_FILE_NAME, 470 MARKER_FILE_NAME, 471 MARKER_NEXT_FILE_NAME, 472 ] { 473 assert!(!recovery_path(paths, name).exists(), "retained {name}"); 474 } 475 } 476 477 async fn database_user_version(paths: &ServiceSqlitePaths) -> i64 { 478 let mut connection = sqlx::SqliteConnection::connect_with( 479 &SqliteConnectOptions::new() 480 .filename(paths.state_database()) 481 .read_only(true) 482 .create_if_missing(false) 483 .disable_statement_logging(), 484 ) 485 .await 486 .expect("open recovered database"); 487 let version = sqlx::query_scalar::<_, i64>("PRAGMA user_version") 488 .fetch_one(&mut connection) 489 .await 490 .expect("read recovery probe"); 491 connection.close().await.expect("close recovered database"); 492 version 493 } 494 495 #[derive(Debug, PartialEq, Eq)] 496 struct FileSnapshot { 497 device: u64, 498 inode: u64, 499 mode: u32, 500 bytes: Vec<u8>, 501 } 502 503 fn file_snapshot(path: &Path) -> FileSnapshot { 504 let metadata = fs::metadata(path).expect("snapshot metadata"); 505 FileSnapshot { 506 device: metadata.dev(), 507 inode: metadata.ino(), 508 mode: metadata.mode(), 509 bytes: fs::read(path).expect("snapshot bytes"), 510 } 511 } 512 513 struct KillOnDrop { 514 child: Child, 515 ready: mpsc::Receiver<String>, 516 finished: bool, 517 } 518 519 impl KillOnDrop { 520 fn new(child: Child, ready: mpsc::Receiver<String>) -> Self { 521 Self { 522 child, 523 ready, 524 finished: false, 525 } 526 } 527 528 fn child(&self) -> &Child { 529 &self.child 530 } 531 532 fn wait(&mut self) -> std::io::Result<ExitStatus> { 533 let status = self.child.wait()?; 534 self.finished = true; 535 Ok(status) 536 } 537 } 538 539 impl Drop for KillOnDrop { 540 fn drop(&mut self) { 541 if !self.finished { 542 let _ = self.child.kill(); 543 let _ = self.child.wait(); 544 } 545 } 546 }