services_hardening_operations.rs (9572B)
1 #![forbid(unsafe_code)] 2 3 use std::error::Error; 4 use std::net::{Ipv4Addr, SocketAddrV4, TcpListener}; 5 6 use rhi::{ 7 InstanceId, RHI_LIVEZ_PATH, RHI_METRICS_PATH, RHI_OPERATIONS_CONTRACT_VERSION, RHI_READYZ_PATH, 8 RhiConfigProfile, RhiEvidenceTransportStatusV1, RhiIdentityHealthV1, RhiIntegrityStateV1, 9 RhiOperationsCancellationToken, RhiOperationsErrorKind, RhiOperationsServer, 10 RhiPersistenceHealthV1, RhiPersistenceStatusV1, RhiPresenceStatusV1, RhiProviderStatusV1, 11 RhiPublicationStatusV1, RhiReconciliationStatusV1, RhiServicePhase, RhiStatusBuildInfoV1, 12 RhiStatusBuildMode, RhiStatusCommonV1, RhiStatusConfigurationIdentityV1, 13 RhiStatusConfigurationSource, RhiStatusObservationV1, RhiStatusReasonCodes, 14 RhiTransportHealthV1, parse_rhi_config_v1, rhi_status_cache, 15 }; 16 use tokio::io::{AsyncReadExt, AsyncWriteExt}; 17 use tokio::net::TcpStream; 18 19 const CONTRACT: &str = include_str!("../contracts/services_hardening/tcp_operations.v1.json"); 20 const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); 21 const SERVICE_REVISION: &str = "0123456789abcdef0123456789abcdef01234567"; 22 const LIB_REVISION: &str = "89abcdef0123456789abcdef0123456789abcdef"; 23 24 fn build_info() -> RhiStatusBuildInfoV1 { 25 RhiStatusBuildInfoV1::new( 26 RhiStatusBuildMode::Release, 27 Some("0.1.0"), 28 Some(SERVICE_REVISION), 29 Some(LIB_REVISION), 30 Some("1.97.1"), 31 Some("x86_64-unknown-linux-gnu"), 32 Some("service-host"), 33 ) 34 .expect("build info") 35 } 36 37 fn identity(configured: bool, available: bool) -> RhiIdentityHealthV1 { 38 RhiIdentityHealthV1::new(configured, available, RhiStatusReasonCodes::empty()) 39 .expect("identity health") 40 } 41 42 fn observation(phase: RhiServicePhase, ready: bool) -> RhiStatusObservationV1 { 43 let configuration = RhiStatusConfigurationIdentityV1::new( 44 "a".repeat(64), 45 RhiStatusConfigurationSource::ExplicitConfig, 46 ) 47 .expect("configuration"); 48 let persistence = RhiPersistenceStatusV1::new( 49 RhiPersistenceHealthV1::Ready, 50 10, 51 42, 52 RhiIntegrityStateV1::Verified, 53 RhiStatusReasonCodes::empty(), 54 ) 55 .expect("persistence"); 56 let provider = RhiProviderStatusV1::new(identity(true, true), RhiStatusReasonCodes::empty()) 57 .expect("provider"); 58 let transport = RhiEvidenceTransportStatusV1::new( 59 RhiTransportHealthV1::Ready, 60 true, 61 true, 62 2, 63 2, 64 RhiStatusReasonCodes::empty(), 65 ) 66 .expect("transport"); 67 RhiStatusObservationV1::new( 68 RhiStatusCommonV1::new( 69 phase, 70 ready, 71 RhiStatusReasonCodes::empty(), 72 1, 73 build_info(), 74 configuration, 75 persistence, 76 ) 77 .expect("common status"), 78 provider, 79 transport, 80 RhiReconciliationStatusV1::default(), 81 RhiPublicationStatusV1::new(0, 0, None), 82 RhiPresenceStatusV1::default(), 83 ) 84 } 85 86 fn available_port() -> u16 { 87 let listener = 88 TcpListener::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)).expect("ephemeral listener"); 89 listener.local_addr().expect("local address").port() 90 } 91 92 fn enabled_config(port: u16) -> String { 93 CONFIG.replacen( 94 "[operations]\nenabled = false", 95 &format!( 96 "[operations]\nenabled = true\nlisten = \"127.0.0.1:{port}\"\nbind_policy = \"loopback_only\"\n\n[operations.limits]\nheader_count = 16\nheader_bytes = 8192\nresponse_body_utf8_bytes = 4096\nconcurrent_connections = 4\nrequest_deadline_ms = 500\nidle_timeout_ms = 500" 97 ), 98 1, 99 ) 100 } 101 102 async fn raw_request(address: std::net::SocketAddr, request: &[u8]) -> Vec<u8> { 103 let mut stream = TcpStream::connect(address).await.expect("connect"); 104 stream.write_all(request).await.expect("request write"); 105 let mut response = Vec::new(); 106 stream 107 .read_to_end(&mut response) 108 .await 109 .expect("response read"); 110 response 111 } 112 113 fn response_text(response: &[u8]) -> &str { 114 std::str::from_utf8(response).expect("response UTF-8") 115 } 116 117 #[test] 118 fn machine_contract_freezes_exact_routes_metrics_and_non_authority() { 119 let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract"); 120 assert_eq!(contract["schema"], "radroots.rhi.tcp-operations.v1"); 121 assert_eq!( 122 contract["contract_version"], 123 RHI_OPERATIONS_CONTRACT_VERSION 124 ); 125 assert_eq!(contract["step"], 212); 126 assert_eq!( 127 contract["routes"] 128 .as_array() 129 .expect("routes") 130 .iter() 131 .map(|route| route["path"].as_str().expect("route path")) 132 .collect::<Vec<_>>(), 133 [RHI_LIVEZ_PATH, RHI_READYZ_PATH, RHI_METRICS_PATH] 134 ); 135 assert_eq!(contract["route_registration_extension"], false); 136 assert_eq!(contract["metrics"]["high_cardinality_labels"], false); 137 assert_eq!(contract["metrics"]["arbitrary_labels"], false); 138 } 139 140 #[tokio::test] 141 async fn exact_tcp_routes_use_only_latest_cached_lifecycle_and_metrics() { 142 let config = parse_rhi_config_v1( 143 enabled_config(available_port()).as_bytes(), 144 RhiConfigProfile::Production, 145 ) 146 .expect("enabled config"); 147 let (mut publisher, reader) = rhi_status_cache( 148 InstanceId::new("primary").expect("instance"), 149 observation(RhiServicePhase::Ready, true), 150 ) 151 .expect("status cache"); 152 let detail_pointer = reader.snapshot().detailed_status_json().as_ptr(); 153 let bound = RhiOperationsServer::new(&config, &reader) 154 .expect("operations server") 155 .bind() 156 .await 157 .expect("bind"); 158 let address = bound.local_address(); 159 let cancellation = RhiOperationsCancellationToken::new(); 160 let task = tokio::spawn(bound.serve(cancellation.clone())); 161 162 let live = raw_request(address, b"GET /livez HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 163 let ready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 164 let metrics = raw_request(address, b"GET /metrics HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 165 assert!(response_text(&live).starts_with("HTTP/1.1 200 OK\r\n")); 166 assert!(response_text(&live).ends_with("live\n")); 167 assert!(response_text(&ready).ends_with("ready\n")); 168 assert!(response_text(&metrics).contains("# TYPE radroots_rhi_service_phase gauge\n")); 169 assert!(response_text(&metrics).contains("radroots_rhi_service_phase{phase=\"ready\"} 1\n")); 170 assert!(response_text(&metrics).contains("radroots_rhi_service_ready 1\n")); 171 172 for request in [ 173 &b"GET /status HTTP/1.1\r\nhost: localhost\r\n\r\n"[..], 174 &b"GET /readyz?probe=1 HTTP/1.1\r\nhost: localhost\r\n\r\n"[..], 175 &b"POST /metrics HTTP/1.1\r\nhost: localhost\r\ncontent-length: 0\r\n\r\n"[..], 176 &b"GET /v1/status HTTP/1.1\r\nhost: localhost\r\n\r\n"[..], 177 ] { 178 let rejected = raw_request(address, request).await; 179 assert!(response_text(&rejected).starts_with("HTTP/1.1 404 Not Found\r\n")); 180 } 181 182 publisher 183 .publish(observation(RhiServicePhase::Unready, false)) 184 .expect("unready publication"); 185 let unready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 186 let metrics = raw_request(address, b"GET /metrics HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 187 assert!(response_text(&unready).starts_with("HTTP/1.1 503 Service Unavailable\r\n")); 188 assert!(response_text(&unready).ends_with("unready\n")); 189 assert!(response_text(&metrics).contains("radroots_rhi_service_phase{phase=\"unready\"} 1\n")); 190 assert!(response_text(&metrics).contains("radroots_rhi_service_ready 0\n")); 191 assert_ne!( 192 reader.snapshot().detailed_status_json().as_ptr(), 193 detail_pointer 194 ); 195 196 cancellation.cancel(); 197 assert_eq!(task.await.expect("serve task"), Ok(())); 198 } 199 200 #[tokio::test] 201 async fn disabled_invalid_and_bind_failures_are_typed_source_free_and_redacted() { 202 let disabled = parse_rhi_config_v1(CONFIG.as_bytes(), RhiConfigProfile::Production) 203 .expect("disabled config"); 204 let (_, reader) = rhi_status_cache( 205 InstanceId::new("primary").expect("instance"), 206 observation(RhiServicePhase::Ready, true), 207 ) 208 .expect("status cache"); 209 let disabled_error = 210 RhiOperationsServer::new(&disabled, &reader).expect_err("disabled operations"); 211 assert_eq!(disabled_error.kind(), RhiOperationsErrorKind::Disabled); 212 assert_eq!(disabled_error.code(), "operations_disabled"); 213 assert!(Error::source(&disabled_error).is_none()); 214 215 let below_floor = 216 enabled_config(available_port()).replace("header_bytes = 8192", "header_bytes = 8191"); 217 assert!(parse_rhi_config_v1(below_floor.as_bytes(), RhiConfigProfile::Production).is_err()); 218 219 let occupied = 220 TcpListener::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)).expect("occupied listener"); 221 let port = occupied.local_addr().expect("occupied address").port(); 222 let config = parse_rhi_config_v1( 223 enabled_config(port).as_bytes(), 224 RhiConfigProfile::Production, 225 ) 226 .expect("enabled config"); 227 let server = RhiOperationsServer::new(&config, &reader).expect("server"); 228 assert!(!format!("{server:?}").contains(&port.to_string())); 229 let bind_error = server.bind().await.expect_err("occupied bind"); 230 assert_eq!(bind_error.kind(), RhiOperationsErrorKind::Bind); 231 assert!(Error::source(&bind_error).is_none()); 232 assert!(!format!("{bind_error:?} {bind_error}").contains(&port.to_string())); 233 }