lib

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

writer_authority.rs (7629B)


      1 #![cfg(any(target_os = "linux", target_os = "macos"))]
      2 
      3 use std::{
      4     env,
      5     error::Error,
      6     fs,
      7     io::{BufRead, BufReader, Read, Write},
      8     os::unix::{
      9         fs::{MetadataExt, PermissionsExt},
     10         process::ExitStatusExt,
     11     },
     12     path::PathBuf,
     13     process::{Child, Command, ExitStatus, Stdio},
     14     sync::mpsc,
     15     thread,
     16     time::{Duration, Instant},
     17 };
     18 
     19 use radroots_runtime_paths::{
     20     InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver,
     21     RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, ServiceId,
     22 };
     23 use radroots_service_sqlite::{
     24     OpenMode, ServiceSqliteErrorKind, ServiceSqlitePaths, WriterAuthority,
     25 };
     26 
     27 const CHILD_ROOT: &str = "RADROOTS_SERVICE_SQLITE_WRITER_CHILD_ROOT";
     28 const PROCESS_READY: &str = "RSHR_STEP073_WRITER_READY";
     29 const MAX_CHILD_INPUT_BYTES: u64 = 4_096;
     30 
     31 fn paths(root: PathBuf) -> ServiceSqlitePaths {
     32     let context = RuntimeContext::resolve(
     33         &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()),
     34         RuntimeContextBootstrap::new(
     35             RadrootsPathProfile::RepoLocal,
     36             Some(root),
     37             RuntimeContextSource::BootstrapCli,
     38             RuntimeContextSource::BootstrapCli,
     39         )
     40         .expect("bootstrap"),
     41         ServiceId::new("myc").expect("service"),
     42         InstanceId::new("process-contention").expect("instance"),
     43     )
     44     .expect("runtime context");
     45     ServiceSqlitePaths::from_runtime_context(&context).expect("SQLite paths")
     46 }
     47 
     48 #[test]
     49 fn writer_authority_is_exclusive_across_processes() {
     50     let root = tempfile::tempdir().expect("root");
     51     let paths = paths(root.path().to_path_buf());
     52     fs::create_dir_all(paths.state_lock().parent().expect("state directory"))
     53         .expect("create state directory");
     54     let _authority = WriterAuthority::acquire(&paths, OpenMode::Initialize)
     55         .expect("parent authority")
     56         .expect("writer capability");
     57 
     58     let status = Command::new(env::current_exe().expect("test executable"))
     59         .arg("--exact")
     60         .arg("writer_authority_child_probe")
     61         .arg("--nocapture")
     62         .env(CHILD_ROOT, root.path())
     63         .status()
     64         .expect("child probe");
     65     assert!(status.success(), "child must observe writer contention");
     66 }
     67 
     68 #[test]
     69 fn writer_authority_child_probe() {
     70     let Some(root) = env::var_os(CHILD_ROOT) else {
     71         return;
     72     };
     73     let error = WriterAuthority::acquire(&paths(PathBuf::from(root)), OpenMode::Initialize)
     74         .expect_err("parent writer must remain authoritative");
     75     assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
     76     assert_eq!(
     77         error.source().map(ToString::to_string).as_deref(),
     78         Some("another SQLite writer is active")
     79     );
     80 }
     81 
     82 #[test]
     83 #[ignore = "private process child invoked by the parent integration test"]
     84 fn writer_authority_holder_child_probe() {
     85     let root = child_root_from_stdin();
     86     rustix::process::umask(rustix::fs::Mode::empty());
     87     let _authority = WriterAuthority::acquire(&paths(root), OpenMode::Initialize)
     88         .expect("child authority")
     89         .expect("writer capability");
     90     println!("\n{PROCESS_READY}");
     91     std::io::stdout().flush().expect("flush readiness");
     92     loop {
     93         thread::park();
     94     }
     95 }
     96 
     97 #[test]
     98 fn sigkill_releases_process_writer_lock_and_preserves_lock_permissions() {
     99     let root = tempfile::tempdir().expect("root");
    100     let paths = paths(root.path().to_path_buf());
    101     let state_directory = paths.state_lock().parent().expect("state directory");
    102     fs::create_dir_all(state_directory).expect("create state directory");
    103     fs::set_permissions(state_directory, fs::Permissions::from_mode(0o700))
    104         .expect("restrict state directory");
    105     let mut child = Command::new(env::current_exe().expect("test executable"))
    106         .arg("--ignored")
    107         .arg("--exact")
    108         .arg("writer_authority_holder_child_probe")
    109         .arg("--nocapture")
    110         .stdin(Stdio::piped())
    111         .stdout(Stdio::piped())
    112         .spawn()
    113         .expect("spawn authority holder");
    114     writeln!(
    115         child.stdin.take().expect("child stdin"),
    116         "{}",
    117         root.path().display()
    118     )
    119     .expect("send child root");
    120     let ready = readiness_receiver(child.stdout.take().expect("child stdout"));
    121     let mut child = KillOnDrop::new(child);
    122     wait_for_ready(&mut child, &ready);
    123     let contended = WriterAuthority::acquire(&paths, OpenMode::Initialize)
    124         .expect_err("live child retains writer authority");
    125     assert_eq!(contended.kind(), ServiceSqliteErrorKind::Authority);
    126     assert_lock_permissions(&paths);
    127 
    128     rustix::process::kill_process(
    129         rustix::process::Pid::from_child(child.child()),
    130         rustix::process::Signal::KILL,
    131     )
    132     .expect("SIGKILL authority holder");
    133     let status = child.wait().expect("wait for killed authority holder");
    134     assert_eq!(status.signal(), Some(9));
    135 
    136     let mut recovered = WriterAuthority::acquire(&paths, OpenMode::Initialize)
    137         .expect("authority after process death")
    138         .expect("writer capability");
    139     assert_lock_permissions(&paths);
    140     recovered.release().expect("release recovered authority");
    141 }
    142 
    143 fn child_root_from_stdin() -> PathBuf {
    144     let mut input = String::new();
    145     std::io::stdin()
    146         .take(MAX_CHILD_INPUT_BYTES + 1)
    147         .read_to_string(&mut input)
    148         .expect("read child root");
    149     assert!(
    150         input.len() <= MAX_CHILD_INPUT_BYTES as usize,
    151         "bounded child input"
    152     );
    153     let root = input
    154         .strip_suffix('\n')
    155         .expect("newline-terminated child root");
    156     assert!(
    157         !root.is_empty() && !root.contains('\n'),
    158         "single child root"
    159     );
    160     PathBuf::from(root)
    161 }
    162 
    163 fn readiness_receiver(stdout: impl Read + Send + 'static) -> mpsc::Receiver<String> {
    164     let (sender, receiver) = mpsc::channel();
    165     thread::spawn(move || {
    166         for line in BufReader::new(stdout).lines() {
    167             if sender.send(line.unwrap_or_default()).is_err() {
    168                 break;
    169             }
    170         }
    171     });
    172     receiver
    173 }
    174 
    175 fn wait_for_ready(child: &mut KillOnDrop, ready: &mpsc::Receiver<String>) {
    176     let deadline = Instant::now() + Duration::from_secs(20);
    177     loop {
    178         match ready.recv_timeout(deadline.saturating_duration_since(Instant::now())) {
    179             Ok(line) if line == PROCESS_READY => return,
    180             Ok(_) => {}
    181             Err(error) => {
    182                 let status = child.child_mut().try_wait().expect("query child status");
    183                 panic!("writer child readiness failed ({error}); status: {status:?}");
    184             }
    185         }
    186     }
    187 }
    188 
    189 fn assert_lock_permissions(paths: &ServiceSqlitePaths) {
    190     let metadata = fs::metadata(paths.state_lock()).expect("lock metadata");
    191     assert!(metadata.file_type().is_file());
    192     assert_eq!(metadata.mode() & 0o777, 0o600);
    193     assert_eq!(metadata.nlink(), 1);
    194     assert_eq!(metadata.uid(), rustix::process::geteuid().as_raw());
    195 }
    196 
    197 struct KillOnDrop {
    198     child: Child,
    199     finished: bool,
    200 }
    201 
    202 impl KillOnDrop {
    203     fn new(child: Child) -> Self {
    204         Self {
    205             child,
    206             finished: false,
    207         }
    208     }
    209 
    210     fn child(&self) -> &Child {
    211         &self.child
    212     }
    213 
    214     fn child_mut(&mut self) -> &mut Child {
    215         &mut self.child
    216     }
    217 
    218     fn wait(&mut self) -> std::io::Result<ExitStatus> {
    219         let status = self.child.wait()?;
    220         self.finished = true;
    221         Ok(status)
    222     }
    223 }
    224 
    225 impl Drop for KillOnDrop {
    226     fn drop(&mut self) {
    227         if !self.finished {
    228             let _ = self.child.kill();
    229             let _ = self.child.wait();
    230         }
    231     }
    232 }