lib

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

commit d7a3cb5a6c3e6f2f7c88458fd514dd354906b8ea
parent 6f70340ad71fa51c135059f408ca6926ea0821ae
Author: triesap <tyson@radroots.org>
Date:   Wed, 12 Aug 2026 02:57:06 +0000

service-sqlite: test process recovery

- add deterministic SIGKILL barriers for writer authority and restore durability edges
- verify orphan refusal plus exact rollback, roll-forward, and terminal recovery
- bind owner-only recovery artifact modes under a permissive child umask
- qualify the process harness on Linux and preserve unsupported-target boundaries

Diffstat:
MAGENTS.md | 11++++++++++-
Mcrates/service_sqlite/README.md | 17++++++++++++++++-
Mcrates/service_sqlite/src/failpoint.rs | 171++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mcrates/service_sqlite/src/restore/mod.rs | 3+++
Acrates/service_sqlite/src/restore/process_tests.rs | 526+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/service_sqlite/tests/package_boundary.rs | 83++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/service_sqlite/tests/writer_authority.rs | 170++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
7 files changed, 971 insertions(+), 10 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -284,7 +284,16 @@ Before editing code: ordinary controllers have zero behavior. Never export failpoint types, use process-global failpoint state, select a point from environment or service configuration, or add a Cargo feature that alters production behavior. - Process-level crash and signal qualification remains a separate layer. + Process-level crash and signal qualification must remain in private Cargo + test binaries. Pass only one bounded temporary root over stdin, require a + fixed stdout token from an occurrence-aware failpoint barrier before + `SIGKILL`, and retain a parent kill-on-drop watchdog. Cover writer-lock death + and the exact pre-marker, prepared, marker-scratch, installed-replacement, + and terminal-marker restore topologies under a permissive child umask. + Require Linux execution for OS-level qualification; macOS is developer + evidence only. Do not ship a helper binary, add production signal/process + behavior, poll filesystem state for crash timing, or claim abrupt power-loss + durability from process-death tests. - Runtime-management flows consume a sealed `RuntimeContext` for every service instance. They must not reconstruct service paths from raw identifiers, ambient selectors, or manager-owned roots, and registries must not persist diff --git a/crates/service_sqlite/README.md b/crates/service_sqlite/README.md @@ -260,7 +260,22 @@ exported from the crate root. There is no process-global failpoint state, environment or configuration selector, Cargo feature, hidden task, timer, panic, or process-exit behavior. These deterministic in-process edges qualify error ordering, rollback, cleanup, recovery evidence, and one-shot retry -semantics; process crashes and signals remain a separate qualification layer. +semantics. + +Process-crash qualification remains test-only and reuses Cargo's private +library and integration-test binaries; the crate ships no helper executable or +signal handler. A parent sends one bounded temporary root over stdin, waits for +a fixed stdout readiness token from an occurrence-aware failpoint barrier, and +then issues `SIGKILL`. The suite proves cross-process writer contention and +lock release plus five restore boundaries: orphan-stage refusal before a +durable marker, prepared rollback, interrupted marker-scratch promotion, +installed-replacement recovery, and terminal-marker cleanup. A permissive +child umask cannot broaden the fixed `0700` state directory or `0600` lock, +database, stage, backup, marker, and marker-scratch artifacts. Linux execution +is required for OS-level qualification; macOS execution is developer evidence, +and unsupported targets compile the production library without the native +process harness. These tests exercise process death at named durable edges and +do not claim abrupt power-loss or storage-device durability behavior. The crate owns mechanics only. Service-specific tables, SQL, repositories, backup content policy, identity material, process lifecycle, and readiness diff --git a/crates/service_sqlite/src/failpoint.rs b/crates/service_sqlite/src/failpoint.rs @@ -3,7 +3,13 @@ use core::fmt; #[cfg(test)] -use std::sync::{Arc, Mutex}; +use std::{ + io::{self, Write}, + sync::{Arc, Condvar, Mutex}, +}; + +#[cfg(test)] +const PROCESS_BARRIER_READY: &[u8] = b"\nRSHR_STEP073_READY\n"; /// Closed inventory of durability edges exercised by the crash-boundary harness. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -126,18 +132,63 @@ struct TestState { fired: bool, reached: Vec<DurabilityFailpoint>, observations: Vec<(DurabilityFailpoint, u8)>, + process_barrier: Option<TestProcessBarrier>, +} + +#[cfg(test)] +#[derive(Clone)] +struct TestProcessBarrier { + point: DurabilityFailpoint, + occurrence: u8, + seen: u8, + gate: Arc<TestProcessBarrierGate>, +} + +#[cfg(test)] +#[derive(Default)] +struct TestProcessBarrierGate { + state: Mutex<TestProcessBarrierGateState>, + changed: Condvar, +} + +#[cfg(test)] +#[derive(Default)] +struct TestProcessBarrierGateState { + ready: bool, + released: bool, } impl DurabilityFailpoints { pub(crate) fn hit(&self, point: DurabilityFailpoint) -> Result<(), DurabilityFailpointError> { #[cfg(test)] { - let mut state = self.state.lock().map_err(|_| DurabilityFailpointError)?; - if state.reached.len() < DurabilityFailpoint::ALL.len() { - state.reached.push(point); + let (injected, process_gate) = { + let mut state = self.state.lock().map_err(|_| DurabilityFailpointError)?; + if state.reached.len() < DurabilityFailpoint::ALL.len() { + state.reached.push(point); + } + let injected = state.armed == Some(point) && !state.fired; + let process_gate = if !state.fired { + state.process_barrier.as_mut().and_then(|barrier| { + if barrier.point != point { + return None; + } + barrier.seen = barrier.seen.saturating_add(1); + (barrier.seen == barrier.occurrence).then(|| Arc::clone(&barrier.gate)) + }) + } else { + None + }; + if injected || process_gate.is_some() { + state.fired = true; + } + (injected, process_gate) + }; + if let Some(process_gate) = process_gate { + process_gate.notify_and_wait()?; + return Err(DurabilityFailpointError); } - if state.armed == Some(point) && !state.fired { - state.fired = true; + if injected { return Err(DurabilityFailpointError); } } @@ -154,11 +205,60 @@ impl DurabilityFailpoints { fired: false, reached: Vec::new(), observations: Vec::new(), + process_barrier: None, + })), + } + } + + #[cfg(test)] + pub(crate) fn process_barrier(point: DurabilityFailpoint, occurrence: u8) -> Self { + assert!( + matches!(occurrence, 1 | 2), + "process occurrence must be 1 or 2" + ); + Self { + state: Arc::new(Mutex::new(TestState { + armed: None, + fired: false, + reached: Vec::new(), + observations: Vec::new(), + process_barrier: Some(TestProcessBarrier { + point, + occurrence, + seen: 0, + gate: Arc::new(TestProcessBarrierGate::default()), + }), })), } } #[cfg(test)] + fn wait_for_process_barrier(&self) { + let gate = self + .state + .lock() + .expect("durability failpoint state") + .process_barrier + .as_ref() + .map(|barrier| Arc::clone(&barrier.gate)) + .expect("process barrier"); + gate.wait_until_ready(); + } + + #[cfg(test)] + fn release_process_barrier(&self) { + let gate = self + .state + .lock() + .expect("durability failpoint state") + .process_barrier + .as_ref() + .map(|barrier| Arc::clone(&barrier.gate)) + .expect("process barrier"); + gate.release(); + } + + #[cfg(test)] pub(crate) fn arm(&self, point: DurabilityFailpoint) { let mut state = self.state.lock().expect("durability failpoint state"); state.armed = Some(point); @@ -202,6 +302,44 @@ impl DurabilityFailpoints { } } +#[cfg(test)] +impl TestProcessBarrierGate { + fn notify_and_wait(&self) -> Result<(), DurabilityFailpointError> { + { + let mut state = self.state.lock().map_err(|_| DurabilityFailpointError)?; + state.ready = true; + self.changed.notify_all(); + } + let mut stdout = io::stdout().lock(); + stdout + .write_all(PROCESS_BARRIER_READY) + .and_then(|()| stdout.flush()) + .map_err(|_| DurabilityFailpointError)?; + let state = self.state.lock().map_err(|_| DurabilityFailpointError)?; + drop( + self.changed + .wait_while(state, |state| !state.released) + .map_err(|_| DurabilityFailpointError)?, + ); + Ok(()) + } + + fn wait_until_ready(&self) { + let state = self.state.lock().expect("process barrier gate"); + drop( + self.changed + .wait_while(state, |state| !state.ready) + .expect("process barrier ready"), + ); + } + + fn release(&self) { + let mut state = self.state.lock().expect("process barrier gate"); + state.released = true; + self.changed.notify_all(); + } +} + impl fmt::Debug for DurabilityFailpoints { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter.write_str("DurabilityFailpoints([redacted])") @@ -223,6 +361,7 @@ impl std::error::Error for DurabilityFailpointError {} #[cfg(test)] mod tests { use super::*; + use std::thread; #[test] fn every_closed_point_fires_once_on_its_owned_plan() { @@ -252,4 +391,24 @@ mod tests { Ok(()) ); } + + #[test] + fn process_barrier_waits_for_the_selected_occurrence_and_releases_once() { + let point = DurabilityFailpoint::MarkerAdvanceAfterDirectorySync; + let plan = DurabilityFailpoints::process_barrier(point, 2); + let worker_plan = plan.clone(); + let worker = thread::spawn(move || { + assert_eq!(worker_plan.hit(point), Ok(())); + worker_plan.hit(point) + }); + + plan.wait_for_process_barrier(); + assert!(plan.fired()); + plan.release_process_barrier(); + assert_eq!( + worker.join().expect("barrier worker"), + Err(DurabilityFailpointError) + ); + assert_eq!(plan.reached(), [point, point]); + } } diff --git a/crates/service_sqlite/src/restore/mod.rs b/crates/service_sqlite/src/restore/mod.rs @@ -5,6 +5,9 @@ mod marker; mod recover; mod stage; +#[cfg(all(test, any(target_os = "linux", target_os = "macos")))] +mod process_tests; + pub use finalize::finalize_staged_restore; pub use stage::{StagedServiceRestore, stage_verified_restore}; diff --git a/crates/service_sqlite/src/restore/process_tests.rs b/crates/service_sqlite/src/restore/process_tests.rs @@ -0,0 +1,526 @@ +use core::num::{NonZeroU32, NonZeroU64}; +use std::{ + env, fs, + io::{BufRead, BufReader, Read, Write}, + os::unix::{ + fs::{MetadataExt, PermissionsExt}, + process::ExitStatusExt, + }, + path::{Path, PathBuf}, + process::{Child, Command, ExitStatus, Stdio}, + sync::mpsc, + thread, + time::{Duration, Instant}, +}; + +use radroots_runtime_paths::{ + InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, + RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, ServiceId, +}; +use radroots_storage::event::SourceGeneration; +use sha2::{Digest, Sha256}; +use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions}; + +use super::{ + BACKUP_FILE_NAME, MARKER_FILE_NAME, MARKER_NEXT_FILE_NAME, STAGED_FILE_NAME, + finalize::test_finalize_with_failpoint, +}; +use crate::failpoint::{DurabilityFailpoint, DurabilityFailpoints}; +use crate::{ + BackupCreatedAtUnixMs, BackupMemberSha256, MigrationAppliedAtUnixSeconds, + MigrationBuildIdentity, MigrationCatalog, OpenMode, SchemaCatalog, SchemaVersionCatalog, + ServiceBackupManifest, ServiceDatabaseIdentity, ServiceDatabaseMetadata, + ServiceSqliteApplicationId, ServiceSqliteConnectionOptions, ServiceSqliteErrorKind, + ServiceSqliteHost, ServiceSqlitePaths, initialize_database, stage_verified_restore, + verify_backup_bundle, +}; + +const PROCESS_READY: &str = "RSHR_STEP073_READY"; +const MANIFEST_FILE_NAME: &str = "process-restore-manifest.v1.json"; +const BUNDLE_DIRECTORY_NAME: &str = "process-restore-bundle"; +const MAX_CHILD_INPUT_BYTES: u64 = 4_096; +const REPLACEMENT_USER_VERSION: i64 = 73; +const PROCESS_TIMEOUT: Duration = Duration::from_secs(30); + +struct Fixture { + root: tempfile::TempDir, + paths: ServiceSqlitePaths, + identity: ServiceDatabaseIdentity, + migrations: MigrationCatalog, + schema: SchemaCatalog, +} + +impl Fixture { + async fn new() -> Self { + let root = tempfile::tempdir().expect("process fixture root"); + let paths = paths(root.path()); + let state_directory = paths.state_database().parent().expect("state directory"); + fs::create_dir_all(state_directory).expect("create state directory"); + fs::set_permissions(state_directory, fs::Permissions::from_mode(0o700)) + .expect("restrict state directory"); + let metadata = metadata(&paths); + let (migrations, schema) = catalogs(); + let mut authority = initialize_database( + &paths, + OpenMode::Initialize, + &metadata, + &schema, + |path| async move { + let options = SqliteConnectOptions::new() + .filename(path) + .create_if_missing(false) + .disable_statement_logging(); + let connection = sqlx::SqliteConnection::connect_with(&options).await?; + connection.close().await + }, + ) + .await + .expect("initialize process database"); + authority + .release() + .expect("release initialization authority"); + { + let connection = rusqlite::Connection::open(paths.state_database()) + .expect("open live database for WAL posture"); + connection + .pragma_update(None, "journal_mode", "WAL") + .expect("set WAL posture"); + connection + .pragma_update(None, "wal_checkpoint", "TRUNCATE") + .expect("checkpoint WAL posture"); + } + + let bundle = root.path().join(BUNDLE_DIRECTORY_NAME); + fs::create_dir(&bundle).expect("create process bundle"); + fs::set_permissions(&bundle, fs::Permissions::from_mode(0o700)) + .expect("restrict process bundle"); + let member = bundle.join(crate::BACKUP_STATE_MEMBER_NAME); + fs::copy(paths.state_database(), &member).expect("copy process member"); + fs::set_permissions(&member, fs::Permissions::from_mode(0o600)) + .expect("restrict process member"); + { + let connection = rusqlite::Connection::open(&member).expect("open process member"); + connection + .pragma_update(None, "user_version", REPLACEMENT_USER_VERSION) + .expect("set replacement probe"); + connection + .pragma_update(None, "wal_checkpoint", "TRUNCATE") + .expect("checkpoint replacement probe"); + } + let bytes = fs::read(&member).expect("read process member"); + let manifest = ServiceBackupManifest::from_capture( + &metadata, + BackupCreatedAtUnixMs::new(1_700_000_073_000).expect("capture time"), + u64::try_from(bytes.len()).expect("member length"), + BackupMemberSha256::from_bytes(Sha256::digest(&bytes).into()), + ) + .expect("process manifest"); + fs::write( + root.path().join(MANIFEST_FILE_NAME), + manifest.canonical_bytes(), + ) + .expect("write process manifest"); + + let identity = metadata.identity(); + Self { + root, + paths, + identity, + migrations, + schema, + } + } + + fn root(&self) -> &Path { + self.root.path() + } +} + +#[tokio::test(flavor = "current_thread")] +#[ignore = "private process child invoked by the parent crash test"] +async fn child_before_prepared_marker() { + run_child(DurabilityFailpoint::MarkerBeforeCreate, 1).await; +} + +#[tokio::test(flavor = "current_thread")] +#[ignore = "private process child invoked by the parent crash test"] +async fn child_after_prepared_marker() { + run_child(DurabilityFailpoint::MarkerAfterDirectorySync, 1).await; +} + +#[tokio::test(flavor = "current_thread")] +#[ignore = "private process child invoked by the parent crash test"] +async fn child_after_live_retained_scratch_sync() { + run_child(DurabilityFailpoint::MarkerAdvanceAfterWriteAndFileSync, 1).await; +} + +#[tokio::test(flavor = "current_thread")] +#[ignore = "private process child invoked by the parent crash test"] +async fn child_after_replacement_install_sync() { + run_child(DurabilityFailpoint::RestoreAfterInstallStageSync, 1).await; +} + +#[tokio::test(flavor = "current_thread")] +#[ignore = "private process child invoked by the parent crash test"] +async fn child_after_terminal_marker_sync() { + run_child(DurabilityFailpoint::MarkerAdvanceAfterDirectorySync, 2).await; +} + +async fn run_child(point: DurabilityFailpoint, occurrence: u8) { + let root = child_root_from_stdin(); + rustix::process::umask(rustix::fs::Mode::empty()); + let paths = paths(&root); + let metadata = metadata(&paths); + let identity = metadata.identity(); + let (migrations, schema) = catalogs(); + let manifest = ServiceBackupManifest::from_canonical_bytes( + &fs::read(root.join(MANIFEST_FILE_NAME)).expect("read child manifest"), + ) + .expect("parse child manifest"); + let verified = verify_backup_bundle( + manifest.canonical_bytes(), + manifest.digest(), + &root.join(BUNDLE_DIRECTORY_NAME), + &identity, + NonZeroU64::new(16 * 1_024 * 1_024).expect("member limit"), + ) + .expect("verify child bundle"); + let staged = stage_verified_restore(&paths, &identity, &migrations, &schema, verified) + .await + .expect("stage child restore"); + let failpoints = DurabilityFailpoints::process_barrier(point, occurrence); + let result = test_finalize_with_failpoint(staged, failpoints).await; + panic!("process durability barrier returned before SIGKILL: {result:?}"); +} + +#[tokio::test(flavor = "current_thread")] +async fn sigkill_restore_boundaries_recover_exact_topologies_and_preserve_permissions() { + for scenario in [ + Scenario::unresolved("child_before_prepared_marker"), + Scenario::recovered("child_after_prepared_marker", 0), + Scenario::recovered( + "child_after_live_retained_scratch_sync", + REPLACEMENT_USER_VERSION, + ), + Scenario::recovered( + "child_after_replacement_install_sync", + REPLACEMENT_USER_VERSION, + ), + Scenario::recovered("child_after_terminal_marker_sync", REPLACEMENT_USER_VERSION), + ] { + run_parent_scenario(scenario).await; + } +} + +struct Scenario { + child_name: &'static str, + expected_user_version: Option<i64>, +} + +impl Scenario { + const fn unresolved(child_name: &'static str) -> Self { + Self { + child_name, + expected_user_version: None, + } + } + + const fn recovered(child_name: &'static str, expected_user_version: i64) -> Self { + Self { + child_name, + expected_user_version: Some(expected_user_version), + } + } +} + +async fn run_parent_scenario(scenario: Scenario) { + let fixture = Fixture::new().await; + let original_live = file_snapshot(fixture.paths.state_database()); + let mut child = spawn_restore_child(scenario.child_name, fixture.root()); + wait_for_ready(&mut child); + assert_recovery_permissions(&fixture.paths); + rustix::process::kill_process( + rustix::process::Pid::from_child(child.child()), + rustix::process::Signal::KILL, + ) + .expect("SIGKILL restore child"); + let status = child.wait().expect("wait for restore child"); + assert_eq!(status.signal(), Some(9), "scenario {}", scenario.child_name); + assert_recovery_permissions(&fixture.paths); + + if let Some(expected_user_version) = scenario.expected_user_version { + let (host, outcome) = ServiceSqliteHost::open_read_write_existing( + &fixture.paths, + &fixture.identity, + &fixture.migrations, + &fixture.schema, + ServiceSqliteConnectionOptions::reviewed(), + MigrationAppliedAtUnixSeconds::new(1_700_000_073).expect("migration time"), + &build_identity(), + &[], + ) + .await + .expect("reopen and reconcile interrupted restore"); + assert_eq!(outcome.applied_count(), 0); + host.close().await.expect("close recovered host"); + assert_no_recovery_evidence(&fixture.paths); + assert_eq!(database_user_version(&fixture.paths), expected_user_version); + assert_live_permissions(&fixture.paths); + } else { + let staged = recovery_path(&fixture.paths, STAGED_FILE_NAME); + let staged_before = file_snapshot(&staged); + let error = ServiceSqliteHost::open_read_write_existing( + &fixture.paths, + &fixture.identity, + &fixture.migrations, + &fixture.schema, + ServiceSqliteConnectionOptions::reviewed(), + MigrationAppliedAtUnixSeconds::new(1_700_000_073).expect("migration time"), + &build_identity(), + &[], + ) + .await + .expect_err("orphan stage without a marker must fail closed"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Recovery); + assert_eq!(file_snapshot(fixture.paths.state_database()), original_live); + assert_eq!(file_snapshot(&staged), staged_before); + assert!(!recovery_path(&fixture.paths, MARKER_FILE_NAME).exists()); + } +} + +fn spawn_restore_child(child_name: &str, root: &Path) -> KillOnDrop { + let exact_name = format!("restore::process_tests::{child_name}"); + let mut child = Command::new(env::current_exe().expect("test executable")) + .arg("--ignored") + .arg("--exact") + .arg(exact_name) + .arg("--nocapture") + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .spawn() + .expect("spawn restore child"); + writeln!( + child.stdin.take().expect("child stdin"), + "{}", + root.display() + ) + .expect("send child root"); + let ready = readiness_receiver(child.stdout.take().expect("child stdout")); + KillOnDrop::new(child, ready) +} + +fn wait_for_ready(child: &mut KillOnDrop) { + let deadline = Instant::now() + PROCESS_TIMEOUT; + loop { + let remaining = deadline.saturating_duration_since(Instant::now()); + match child.ready.recv_timeout(remaining) { + Ok(line) if line == PROCESS_READY => return, + Ok(_) => {} + Err(error) => { + let status = child.child.try_wait().expect("query child status"); + panic!("restore child readiness failed ({error}); status: {status:?}"); + } + } + } +} + +fn readiness_receiver(stdout: impl Read + Send + 'static) -> mpsc::Receiver<String> { + let (sender, receiver) = mpsc::channel(); + thread::spawn(move || { + for line in BufReader::new(stdout).lines() { + if sender.send(line.unwrap_or_default()).is_err() { + break; + } + } + }); + receiver +} + +fn child_root_from_stdin() -> PathBuf { + let mut input = String::new(); + std::io::stdin() + .take(MAX_CHILD_INPUT_BYTES + 1) + .read_to_string(&mut input) + .expect("read child root"); + assert!( + input.len() <= MAX_CHILD_INPUT_BYTES as usize, + "bounded child input" + ); + let root = input + .strip_suffix('\n') + .expect("newline-terminated child root"); + assert!( + !root.is_empty() && !root.contains('\n'), + "single child root" + ); + PathBuf::from(root) +} + +fn paths(root: &Path) -> ServiceSqlitePaths { + let context = RuntimeContext::resolve( + &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), + RuntimeContextBootstrap::new( + RadrootsPathProfile::RepoLocal, + Some(root.to_path_buf()), + RuntimeContextSource::BootstrapCli, + RuntimeContextSource::BootstrapCli, + ) + .expect("bootstrap"), + ServiceId::new("myc").expect("service"), + InstanceId::new("process-recovery").expect("instance"), + ) + .expect("runtime context"); + ServiceSqlitePaths::from_runtime_context(&context).expect("SQLite paths") +} + +fn metadata(paths: &ServiceSqlitePaths) -> ServiceDatabaseMetadata { + ServiceDatabaseMetadata::new( + paths, + SourceGeneration::new([7; 32]).expect("generation"), + NonZeroU32::new(1).expect("schema"), + 1_700_000_000_000, + ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), + ) + .expect("metadata") +} + +fn catalogs() -> (MigrationCatalog, SchemaCatalog) { + let migrations = MigrationCatalog::new([]).expect("migration catalog"); + let digest = SchemaVersionCatalog::computed_digest(1, []).expect("schema digest"); + let version = SchemaVersionCatalog::new(1, [], digest).expect("schema version"); + let schema = SchemaCatalog::new(&migrations, [version]).expect("schema catalog"); + (migrations, schema) +} + +fn build_identity() -> MigrationBuildIdentity { + MigrationBuildIdentity::new( + "1.0.0", + "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + "rustc-1.97.0", + "process-test", + "test", + 1, + 1, + 1, + 1, + 1, + ) + .expect("build identity") +} + +fn recovery_path(paths: &ServiceSqlitePaths, name: &str) -> PathBuf { + paths + .state_database() + .parent() + .expect("state directory") + .join(name) +} + +fn assert_recovery_permissions(paths: &ServiceSqlitePaths) { + let state_directory = paths.state_database().parent().expect("state directory"); + let directory = fs::metadata(state_directory).expect("state directory metadata"); + assert!(directory.is_dir()); + assert_eq!(directory.mode() & 0o777, 0o700); + assert_eq!(directory.uid(), rustix::process::geteuid().as_raw()); + for path in [ + paths.state_database().to_path_buf(), + paths.state_lock().to_path_buf(), + recovery_path(paths, STAGED_FILE_NAME), + recovery_path(paths, BACKUP_FILE_NAME), + recovery_path(paths, MARKER_FILE_NAME), + recovery_path(paths, MARKER_NEXT_FILE_NAME), + ] { + if path.exists() { + assert_file_permissions(&path); + } + } +} + +fn assert_live_permissions(paths: &ServiceSqlitePaths) { + assert_file_permissions(paths.state_database()); + assert_file_permissions(paths.state_lock()); +} + +fn assert_file_permissions(path: &Path) { + let metadata = fs::symlink_metadata(path).expect("artifact metadata"); + assert!(metadata.file_type().is_file()); + assert_eq!(metadata.mode() & 0o777, 0o600); + assert_eq!(metadata.nlink(), 1); + assert_eq!(metadata.uid(), rustix::process::geteuid().as_raw()); +} + +fn assert_no_recovery_evidence(paths: &ServiceSqlitePaths) { + for name in [ + STAGED_FILE_NAME, + BACKUP_FILE_NAME, + MARKER_FILE_NAME, + MARKER_NEXT_FILE_NAME, + ] { + assert!(!recovery_path(paths, name).exists(), "retained {name}"); + } +} + +fn database_user_version(paths: &ServiceSqlitePaths) -> i64 { + let connection = rusqlite::Connection::open_with_flags( + paths.state_database(), + rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX, + ) + .expect("open recovered database"); + connection + .pragma_query_value(None, "user_version", |row| row.get(0)) + .expect("read recovery probe") +} + +#[derive(Debug, PartialEq, Eq)] +struct FileSnapshot { + device: u64, + inode: u64, + mode: u32, + bytes: Vec<u8>, +} + +fn file_snapshot(path: &Path) -> FileSnapshot { + let metadata = fs::metadata(path).expect("snapshot metadata"); + FileSnapshot { + device: metadata.dev(), + inode: metadata.ino(), + mode: metadata.mode(), + bytes: fs::read(path).expect("snapshot bytes"), + } +} + +struct KillOnDrop { + child: Child, + ready: mpsc::Receiver<String>, + finished: bool, +} + +impl KillOnDrop { + fn new(child: Child, ready: mpsc::Receiver<String>) -> Self { + Self { + child, + ready, + finished: false, + } + } + + fn child(&self) -> &Child { + &self.child + } + + fn wait(&mut self) -> std::io::Result<ExitStatus> { + let status = self.child.wait()?; + self.finished = true; + Ok(status) + } +} + +impl Drop for KillOnDrop { + fn drop(&mut self) { + if !self.finished { + let _ = self.child.kill(); + let _ = self.child.wait(); + } + } +} diff --git a/crates/service_sqlite/tests/package_boundary.rs b/crates/service_sqlite/tests/package_boundary.rs @@ -22,10 +22,12 @@ const RESTORE_MARKER_SOURCE: &str = include_str!("../src/restore/marker.rs"); const RESTORE_FINALIZE_SOURCE: &str = include_str!("../src/restore/finalize.rs"); const RESTORE_RECOVER_SOURCE: &str = include_str!("../src/restore/recover.rs"); const RESTORE_ROOT_SOURCE: &str = include_str!("../src/restore/mod.rs"); +const RESTORE_PROCESS_TEST_SOURCE: &str = include_str!("../src/restore/process_tests.rs"); const RESTORE_STAGE_SOURCE: &str = include_str!("../src/restore/stage.rs"); const STATUS_SOURCE: &str = include_str!("../src/status/mod.rs"); const DISK_SOURCE: &str = include_str!("../src/status/disk.rs"); const TRANSACTION_CONTROL_SOURCE: &str = include_str!("../src/transaction_control.rs"); +const WRITER_PROCESS_TEST_SOURCE: &str = include_str!("writer_authority.rs"); #[test] fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { @@ -248,7 +250,14 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "no failpoint type or selector is exported from the crate root", "no process-global failpoint state, environment or configuration selector, Cargo feature", "hidden task, timer, panic, or process-exit behavior", - "process crashes and signals remain a separate qualification layer", + "Process-crash qualification remains test-only", + "one bounded temporary root over stdin", + "fixed stdout readiness token from an occurrence-aware failpoint barrier", + "orphan-stage refusal before a durable marker", + "interrupted marker-scratch promotion", + "permissive child umask cannot broaden", + "Linux execution is required for OS-level qualification", + "do not claim abrupt power-loss or storage-device durability behavior", ] { assert!( readme_words.contains(required), @@ -482,6 +491,78 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "Step 072 integration source is missing `{required}`" ); } + for required in [ + "pub(crate) fn process_barrier", + "occurrence: u8", + "RSHR_STEP073_READY", + "Condvar", + "wait_while", + ] { + assert!( + FAILPOINT_SOURCE.contains(required), + "Step 073 process barrier is missing `{required}`" + ); + } + for required in [ + "mod process_tests;", + "#[cfg(all(test, any(target_os = \"linux\", target_os = \"macos\")))]", + ] { + assert!( + RESTORE_ROOT_SOURCE.contains(required), + "Step 073 restore test boundary is missing `{required}`" + ); + } + for required in [ + "child_before_prepared_marker", + "child_after_prepared_marker", + "child_after_live_retained_scratch_sync", + "child_after_replacement_install_sync", + "child_after_terminal_marker_sync", + "MarkerBeforeCreate", + "MarkerAdvanceAfterWriteAndFileSync", + "MarkerAdvanceAfterDirectorySync, 2", + "--ignored", + "--exact", + "Stdio::piped()", + "Signal::KILL", + "status.signal(), Some(9)", + "assert_recovery_permissions", + "ServiceSqliteErrorKind::Recovery", + ] { + assert!( + RESTORE_PROCESS_TEST_SOURCE.contains(required), + "Step 073 restore process harness is missing `{required}`" + ); + } + for required in [ + "writer_authority_holder_child_probe", + "RSHR_STEP073_WRITER_READY", + "--ignored", + "Stdio::piped()", + "Signal::KILL", + "status.signal(), Some(9)", + "assert_lock_permissions", + ] { + assert!( + WRITER_PROCESS_TEST_SOURCE.contains(required), + "Step 073 writer process harness is missing `{required}`" + ); + } + for forbidden in [ + "std::env::var", + "env::var_os", + ".env(", + "ready.is_file", + "thread::sleep", + "remove_dir_all", + "process::exit", + ] { + assert!( + !RESTORE_PROCESS_TEST_SOURCE.contains(forbidden), + "Step 073 restore process harness contains forbidden `{forbidden}`" + ); + } + assert!(!ROOT.contains("process_tests")); for required in [ "radroots.service-backup", diff --git a/crates/service_sqlite/tests/writer_authority.rs b/crates/service_sqlite/tests/writer_authority.rs @@ -1,6 +1,20 @@ #![cfg(any(target_os = "linux", target_os = "macos"))] -use std::{env, error::Error, fs, path::PathBuf, process::Command}; +use std::{ + env, + error::Error, + fs, + io::{BufRead, BufReader, Read, Write}, + os::unix::{ + fs::{MetadataExt, PermissionsExt}, + process::ExitStatusExt, + }, + path::PathBuf, + process::{Child, Command, ExitStatus, Stdio}, + sync::mpsc, + thread, + time::{Duration, Instant}, +}; use radroots_runtime_paths::{ InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, @@ -11,6 +25,8 @@ use radroots_service_sqlite::{ }; const CHILD_ROOT: &str = "RADROOTS_SERVICE_SQLITE_WRITER_CHILD_ROOT"; +const PROCESS_READY: &str = "RSHR_STEP073_WRITER_READY"; +const MAX_CHILD_INPUT_BYTES: u64 = 4_096; fn paths(root: PathBuf) -> ServiceSqlitePaths { let context = RuntimeContext::resolve( @@ -62,3 +78,155 @@ fn writer_authority_child_probe() { Some("another SQLite writer is active") ); } + +#[test] +#[ignore = "private process child invoked by the parent integration test"] +fn writer_authority_holder_child_probe() { + let root = child_root_from_stdin(); + rustix::process::umask(rustix::fs::Mode::empty()); + let _authority = WriterAuthority::acquire(&paths(root), OpenMode::Initialize) + .expect("child authority") + .expect("writer capability"); + println!("\n{PROCESS_READY}"); + std::io::stdout().flush().expect("flush readiness"); + loop { + thread::park(); + } +} + +#[test] +fn sigkill_releases_process_writer_lock_and_preserves_lock_permissions() { + let root = tempfile::tempdir().expect("root"); + let paths = paths(root.path().to_path_buf()); + let state_directory = paths.state_lock().parent().expect("state directory"); + fs::create_dir_all(state_directory).expect("create state directory"); + fs::set_permissions(state_directory, fs::Permissions::from_mode(0o700)) + .expect("restrict state directory"); + let mut child = Command::new(env::current_exe().expect("test executable")) + .arg("--ignored") + .arg("--exact") + .arg("writer_authority_holder_child_probe") + .arg("--nocapture") + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .spawn() + .expect("spawn authority holder"); + writeln!( + child.stdin.take().expect("child stdin"), + "{}", + root.path().display() + ) + .expect("send child root"); + let ready = readiness_receiver(child.stdout.take().expect("child stdout")); + let mut child = KillOnDrop::new(child); + wait_for_ready(&mut child, &ready); + let contended = WriterAuthority::acquire(&paths, OpenMode::Initialize) + .expect_err("live child retains writer authority"); + assert_eq!(contended.kind(), ServiceSqliteErrorKind::Authority); + assert_lock_permissions(&paths); + + rustix::process::kill_process( + rustix::process::Pid::from_child(child.child()), + rustix::process::Signal::KILL, + ) + .expect("SIGKILL authority holder"); + let status = child.wait().expect("wait for killed authority holder"); + assert_eq!(status.signal(), Some(9)); + + let mut recovered = WriterAuthority::acquire(&paths, OpenMode::Initialize) + .expect("authority after process death") + .expect("writer capability"); + assert_lock_permissions(&paths); + recovered.release().expect("release recovered authority"); +} + +fn child_root_from_stdin() -> PathBuf { + let mut input = String::new(); + std::io::stdin() + .take(MAX_CHILD_INPUT_BYTES + 1) + .read_to_string(&mut input) + .expect("read child root"); + assert!( + input.len() <= MAX_CHILD_INPUT_BYTES as usize, + "bounded child input" + ); + let root = input + .strip_suffix('\n') + .expect("newline-terminated child root"); + assert!( + !root.is_empty() && !root.contains('\n'), + "single child root" + ); + PathBuf::from(root) +} + +fn readiness_receiver(stdout: impl Read + Send + 'static) -> mpsc::Receiver<String> { + let (sender, receiver) = mpsc::channel(); + thread::spawn(move || { + for line in BufReader::new(stdout).lines() { + if sender.send(line.unwrap_or_default()).is_err() { + break; + } + } + }); + receiver +} + +fn wait_for_ready(child: &mut KillOnDrop, ready: &mpsc::Receiver<String>) { + let deadline = Instant::now() + Duration::from_secs(20); + loop { + match ready.recv_timeout(deadline.saturating_duration_since(Instant::now())) { + Ok(line) if line == PROCESS_READY => return, + Ok(_) => {} + Err(error) => { + let status = child.child_mut().try_wait().expect("query child status"); + panic!("writer child readiness failed ({error}); status: {status:?}"); + } + } + } +} + +fn assert_lock_permissions(paths: &ServiceSqlitePaths) { + let metadata = fs::metadata(paths.state_lock()).expect("lock metadata"); + assert!(metadata.file_type().is_file()); + assert_eq!(metadata.mode() & 0o777, 0o600); + assert_eq!(metadata.nlink(), 1); + assert_eq!(metadata.uid(), rustix::process::geteuid().as_raw()); +} + +struct KillOnDrop { + child: Child, + finished: bool, +} + +impl KillOnDrop { + fn new(child: Child) -> Self { + Self { + child, + finished: false, + } + } + + fn child(&self) -> &Child { + &self.child + } + + fn child_mut(&mut self) -> &mut Child { + &mut self.child + } + + fn wait(&mut self) -> std::io::Result<ExitStatus> { + let status = self.child.wait()?; + self.finished = true; + Ok(status) + } +} + +impl Drop for KillOnDrop { + fn drop(&mut self) { + if !self.finished { + let _ = self.child.kill(); + let _ = self.child.wait(); + } + } +}