lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

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 }