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 }