rhi

Coordinated trade for connected markets
git clone https://radroots.dev/git/rhi.git
Log | Files | Refs | README | LICENSE

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 }