system_doctor.rs (13055B)
1 //! Production active-doctor probes composed from existing sealed authorities. 2 3 use radroots_service_host::{SystemWallClock, WallClock}; 4 use radroots_service_sqlite::{ 5 IntegrityCheckOutcome, IntegrityCheckedAtUnixMs, MinimumFreeBytes, 6 PlatformStateFilesystemCapacitySource, inspect_state_filesystem_capacity, 7 }; 8 9 use crate::admin_v1::admin_transport_limits; 10 use crate::transport_nostr_adapter::{build_rhi_nostr_adapters, probe_required_sources}; 11 use crate::{ 12 RhiConfigDocumentV1, RhiDoctorCheckDefinition, RhiDoctorCheckId, RhiDoctorFuture, 13 RhiDoctorObservation, RhiDoctorProbe, RhiIdentityEnvelopeBinding, RhiRuntimeContext, 14 RhiStateHost, open_rhi_encrypted_identity, open_rhi_state_inspection_from_config, 15 resolve_rhi_wrapping_credential, 16 }; 17 18 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 19 enum StateProbeError { 20 Query, 21 } 22 23 pub(crate) struct RhiSystemDoctorProbe<'a> { 24 runtime: &'a RhiRuntimeContext, 25 configuration: &'a RhiConfigDocumentV1, 26 } 27 28 impl<'a> RhiSystemDoctorProbe<'a> { 29 pub(crate) const fn new( 30 runtime: &'a RhiRuntimeContext, 31 configuration: &'a RhiConfigDocumentV1, 32 ) -> Self { 33 Self { 34 runtime, 35 configuration, 36 } 37 } 38 39 async fn run(&self, definition: RhiDoctorCheckDefinition) -> bool { 40 match definition.id() { 41 RhiDoctorCheckId::PathsPermissions => self.probe_paths(), 42 RhiDoctorCheckId::WriterLock 43 | RhiDoctorCheckId::SqliteSchema 44 | RhiDoctorCheckId::SqliteIntegrity 45 | RhiDoctorCheckId::CursorCheckpoint 46 | RhiDoctorCheckId::ReconciliationLeases 47 | RhiDoctorCheckId::ReconciliationBacklog 48 | RhiDoctorCheckId::PublicationInvariants => self.probe_state(definition.id()).await, 49 RhiDoctorCheckId::SqliteFreeSpace => self.probe_free_space(), 50 RhiDoctorCheckId::IdentityBinding => self.probe_identity_binding().await, 51 RhiDoctorCheckId::AdminBindPolicy => self.probe_admin_policy(), 52 RhiDoctorCheckId::OperationsBindPolicy => self.probe_operations_policy(), 53 RhiDoctorCheckId::NetworkPolicy => build_rhi_nostr_adapters(self.configuration).is_ok(), 54 RhiDoctorCheckId::RequiredSources => { 55 let Some(deadline) = absolute_deadline(definition.deadline_ms()) else { 56 return false; 57 }; 58 probe_required_sources(self.configuration, deadline) 59 .await 60 .is_ok() 61 } 62 RhiDoctorCheckId::ClockSkew => false, 63 } 64 } 65 66 fn probe_paths(&self) -> bool { 67 let Ok(paths) = crate::state_host::state_paths(self.runtime) else { 68 return false; 69 }; 70 let Ok(minimum) = MinimumFreeBytes::new(1) else { 71 return false; 72 }; 73 inspect_state_filesystem_capacity(&paths, minimum, &PlatformStateFilesystemCapacitySource) 74 .is_ok() 75 } 76 77 async fn probe_state(&self, check: RhiDoctorCheckId) -> bool { 78 let Ok(state) = 79 open_rhi_state_inspection_from_config(self.runtime, self.configuration).await 80 else { 81 return false; 82 }; 83 let outcome = match check { 84 RhiDoctorCheckId::WriterLock | RhiDoctorCheckId::SqliteSchema => true, 85 RhiDoctorCheckId::SqliteIntegrity => match integrity_time() { 86 Some(checked_at) => state 87 .inspect_integrity(checked_at) 88 .await 89 .is_ok_and(|report| { 90 report.sqlite() == IntegrityCheckOutcome::Verified 91 && report.foreign_keys() == IntegrityCheckOutcome::Verified 92 }), 93 None => false, 94 }, 95 RhiDoctorCheckId::CursorCheckpoint => { 96 run_scalar_probe(&state, CURSOR_CHECKPOINT_INVARIANTS_SQL, None).await 97 } 98 RhiDoctorCheckId::ReconciliationLeases => { 99 run_scalar_probe(&state, RECONCILIATION_LEASE_INVARIANTS_SQL, None).await 100 } 101 RhiDoctorCheckId::ReconciliationBacklog => { 102 let capacity = 103 configuration_u64(self.configuration, "/reconciliation/queue_capacity") 104 .and_then(|value| i64::try_from(value).ok()); 105 match capacity { 106 Some(capacity) => { 107 run_scalar_probe( 108 &state, 109 RECONCILIATION_BACKLOG_INVARIANTS_SQL, 110 Some(capacity), 111 ) 112 .await 113 } 114 None => false, 115 } 116 } 117 RhiDoctorCheckId::PublicationInvariants => { 118 run_scalar_probe(&state, PUBLICATION_INVARIANTS_SQL, None).await 119 } 120 _ => false, 121 }; 122 let closed = state.close().await.is_ok(); 123 outcome && closed 124 } 125 126 fn probe_free_space(&self) -> bool { 127 let Some(minimum) = configuration_u64(self.configuration, "/database/minimum_free_bytes") 128 .and_then(|value| MinimumFreeBytes::new(value).ok()) 129 else { 130 return false; 131 }; 132 let Ok(paths) = crate::state_host::state_paths(self.runtime) else { 133 return false; 134 }; 135 inspect_state_filesystem_capacity(&paths, minimum, &PlatformStateFilesystemCapacitySource) 136 .is_ok_and(|capacity| capacity.allows_authoritative_admission()) 137 } 138 139 async fn probe_identity_binding(&self) -> bool { 140 let Ok(state) = 141 open_rhi_state_inspection_from_config(self.runtime, self.configuration).await 142 else { 143 return false; 144 }; 145 let result = 146 RhiIdentityEnvelopeBinding::from_configuration(self.configuration, state.metadata()) 147 .ok() 148 .and_then(|binding| { 149 let credential = 150 resolve_rhi_wrapping_credential(self.runtime, &binding).ok()?; 151 open_rhi_encrypted_identity(&binding, &credential).ok() 152 }); 153 let closed = state.close().await.is_ok(); 154 result.is_some() && closed 155 } 156 157 fn probe_admin_policy(&self) -> bool { 158 let path = self.runtime.artifacts().admin_socket(); 159 path.is_absolute() 160 && path.to_str().is_some_and(|value| value.len() <= 4_096) 161 && admin_transport_limits(self.configuration).is_ok() 162 } 163 164 fn probe_operations_policy(&self) -> bool { 165 let Some(enabled) = self 166 .configuration 167 .normalized() 168 .pointer("/operations/enabled") 169 .and_then(serde_json::Value::as_bool) 170 else { 171 return false; 172 }; 173 !enabled 174 || (self 175 .configuration 176 .normalized() 177 .pointer("/operations/listen") 178 .and_then(serde_json::Value::as_str) 179 .is_some() 180 && self 181 .configuration 182 .normalized() 183 .pointer("/operations/bind_policy") 184 .and_then(serde_json::Value::as_str) 185 .is_some()) 186 } 187 } 188 189 impl RhiDoctorProbe for RhiSystemDoctorProbe<'_> { 190 fn probe(&self, definition: RhiDoctorCheckDefinition) -> RhiDoctorFuture<'_> { 191 Box::pin(async move { 192 if definition.id() == RhiDoctorCheckId::ClockSkew { 193 RhiDoctorObservation::Skipped 194 } else if self.run(definition).await { 195 RhiDoctorObservation::Pass 196 } else { 197 RhiDoctorObservation::Fail 198 } 199 }) 200 } 201 } 202 203 async fn run_scalar_probe(state: &RhiStateHost, sql: &'static str, bound: Option<i64>) -> bool { 204 state 205 .sqlite_host() 206 .transaction(move |transaction| { 207 Box::pin(async move { 208 let mut query = sqlx::query_scalar::<_, i64>(sql); 209 if let Some(bound) = bound { 210 query = query.bind(bound); 211 } 212 query 213 .fetch_one(&mut *transaction) 214 .await 215 .map(|invalid| invalid == 0) 216 .map_err(|_| StateProbeError::Query) 217 }) 218 }) 219 .await 220 .is_ok_and(|valid| valid) 221 } 222 223 fn configuration_u64(configuration: &RhiConfigDocumentV1, pointer: &str) -> Option<u64> { 224 configuration 225 .normalized() 226 .pointer(pointer) 227 .and_then(serde_json::Value::as_u64) 228 } 229 230 fn wall_time_millis() -> Option<u64> { 231 SystemWallClock 232 .now_utc() 233 .ok() 234 .and_then(|time| time.get().checked_mul(1_000)) 235 .filter(|value| i64::try_from(*value).is_ok()) 236 } 237 238 fn absolute_deadline(duration_ms: u64) -> Option<u64> { 239 wall_time_millis()?.checked_add(duration_ms) 240 } 241 242 fn integrity_time() -> Option<IntegrityCheckedAtUnixMs> { 243 IntegrityCheckedAtUnixMs::new(wall_time_millis()?) 244 } 245 246 const CURSOR_CHECKPOINT_INVARIANTS_SQL: &str = r#"SELECT COUNT(*) 247 FROM relay_checkpoints 248 WHERE typeof(source_id) != 'text' 249 OR length(CAST(source_id AS BLOB)) NOT BETWEEN 1 AND 64 250 OR source_id GLOB '*[^a-z0-9_-]*' 251 OR substr(source_id, 1, 1) NOT GLOB '[a-z]' 252 OR selector_id != 'trade_mutation_lineage_v1' 253 OR typeof(evidence_policy_sha256) != 'blob' 254 OR length(evidence_policy_sha256) != 32 255 OR typeof(trade_id) != 'blob' OR length(trade_id) != 16 256 OR typeof(cursor_event_id) != 'blob' OR length(cursor_event_id) != 32 257 OR cursor_created_at_unix_s NOT BETWEEN 0 AND 9223372036854775807 258 OR revision NOT BETWEEN 1 AND 9223372036854775807 259 OR completed_at_unix_s NOT BETWEEN 1 AND 9223372036854775807"#; 260 261 const RECONCILIATION_LEASE_INVARIANTS_SQL: &str = r#"SELECT COUNT(*) 262 FROM reconciliation_jobs 263 WHERE (state = 'leased' AND ( 264 typeof(lease_owner) != 'blob' OR length(lease_owner) != 16 265 OR lease_expires_unix_ms NOT BETWEEN 1 AND 9223372036854775807 266 OR next_attempt_unix_ms IS NOT NULL)) 267 OR (state != 'leased' AND (lease_owner IS NOT NULL OR lease_expires_unix_ms IS NOT NULL)) 268 OR attempt_count NOT BETWEEN 0 AND max_attempts 269 OR failure_count NOT BETWEEN 0 AND attempt_count 270 OR revision NOT BETWEEN 1 AND 9223372036854775807"#; 271 272 const RECONCILIATION_BACKLOG_INVARIANTS_SQL: &str = r#"SELECT CASE 273 WHEN (SELECT COUNT(*) FROM reconciliation_jobs WHERE state IN ('ready', 'leased')) > ? 274 THEN 1 275 WHEN EXISTS ( 276 SELECT 1 FROM reconciliation_jobs 277 WHERE state = 'ready' AND ( 278 next_attempt_unix_ms IS NULL OR attempt_count >= max_attempts)) 279 THEN 1 280 WHEN EXISTS ( 281 SELECT 1 FROM reconciliation_jobs 282 WHERE state IN ('exhausted', 'superseded', 'completed') 283 AND next_attempt_unix_ms IS NOT NULL) 284 THEN 1 285 ELSE 0 END"#; 286 287 const PUBLICATION_INVARIANTS_SQL: &str = r#"SELECT COUNT(*) 288 FROM publication_outbox AS outbox 289 LEFT JOIN signed_attestation_events AS event ON event.event_id = outbox.event_id 290 WHERE event.event_id IS NULL 291 OR outbox.event_sha256 != event.event_sha256 292 OR outbox.target_count != ( 293 SELECT COUNT(*) FROM publication_targets AS target 294 WHERE target.outbox_id = outbox.outbox_id) 295 OR outbox.required_target_count != ( 296 SELECT COUNT(*) FROM publication_targets AS target 297 WHERE target.outbox_id = outbox.outbox_id AND target.required = 1) 298 OR EXISTS ( 299 SELECT 1 FROM publication_targets AS target 300 WHERE target.outbox_id = outbox.outbox_id 301 AND (target.target_ordinal < 0 OR target.target_ordinal >= outbox.target_count)) 302 OR (outbox.state = 'complete' AND EXISTS ( 303 SELECT 1 FROM publication_targets AS target 304 WHERE target.outbox_id = outbox.outbox_id 305 AND target.required = 1 AND target.state != 'accepted'))"#; 306 307 #[cfg(test)] 308 mod tests { 309 use super::*; 310 311 #[test] 312 fn deadline_math_is_checked_and_clock_skew_remains_unclaimed() { 313 assert!(absolute_deadline(15_000).is_some()); 314 assert!(integrity_time().is_some()); 315 } 316 317 #[test] 318 fn state_queries_are_bounded_scalar_projections() { 319 for query in [ 320 CURSOR_CHECKPOINT_INVARIANTS_SQL, 321 RECONCILIATION_LEASE_INVARIANTS_SQL, 322 RECONCILIATION_BACKLOG_INVARIANTS_SQL, 323 PUBLICATION_INVARIANTS_SQL, 324 ] { 325 assert!(query.starts_with("SELECT")); 326 assert!(!query.contains("SELECT *")); 327 assert!(!query.contains("ORDER BY")); 328 } 329 } 330 331 #[test] 332 fn production_probe_source_retains_no_raw_error_projection() { 333 let source = include_str!("system_doctor.rs") 334 .split("#[cfg(test)]") 335 .next() 336 .expect("production source"); 337 for forbidden in ["format!(\"{error", "to_string()", "source()"] { 338 assert!(!source.contains(forbidden), "found `{forbidden}`"); 339 } 340 } 341 }