services_hardening_operations.rs (9591B)
1 #![forbid(unsafe_code)] 2 3 use std::error::Error; 4 use std::net::{Ipv4Addr, SocketAddrV4, TcpListener}; 5 6 use myc::{ 7 InstanceId, MYC_LIVEZ_PATH, MYC_METRICS_PATH, MYC_OPERATIONS_CONTRACT_VERSION, MYC_READYZ_PATH, 8 MYC_STATE_SCHEMA_VERSION, MycConfigProfile, MycConnectionCountsV1, MycIdentityHealthV1, 9 MycIntegrityStateV1, MycOperationsCancellationToken, MycOperationsErrorKind, 10 MycOperationsServer, MycOutboxStatusV1, MycPersistenceHealthV1, MycPersistenceStatusV1, 11 MycProviderStatusV1, MycRelayTransportStatusV1, MycServicePhase, MycStatusBuildInfoV1, 12 MycStatusBuildMode, MycStatusCommonV1, MycStatusConfigurationIdentityV1, 13 MycStatusConfigurationSource, MycStatusObservationV1, MycStatusReasonCodes, 14 MycTransportHealthV1, myc_status_cache, parse_myc_config_v1, 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() -> MycStatusBuildInfoV1 { 25 MycStatusBuildInfoV1::new( 26 MycStatusBuildMode::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) -> MycIdentityHealthV1 { 38 MycIdentityHealthV1::new(configured, available, MycStatusReasonCodes::empty()) 39 .expect("identity health") 40 } 41 42 fn observation(phase: MycServicePhase, ready: bool) -> MycStatusObservationV1 { 43 let configuration = MycStatusConfigurationIdentityV1::new( 44 "a".repeat(64), 45 MycStatusConfigurationSource::ExplicitConfig, 46 ) 47 .expect("configuration"); 48 let persistence = MycPersistenceStatusV1::new( 49 MycPersistenceHealthV1::Ready, 50 MYC_STATE_SCHEMA_VERSION, 51 42, 52 MycIntegrityStateV1::Verified, 53 MycStatusReasonCodes::empty(), 54 ) 55 .expect("persistence"); 56 let provider = MycProviderStatusV1::new( 57 identity(true, true), 58 identity(true, true), 59 identity(false, false), 60 MycStatusReasonCodes::empty(), 61 ) 62 .expect("provider"); 63 let transport = MycRelayTransportStatusV1::new( 64 MycTransportHealthV1::Ready, 65 true, 66 2, 67 MycStatusReasonCodes::empty(), 68 ) 69 .expect("transport"); 70 MycStatusObservationV1::new( 71 MycStatusCommonV1::new( 72 phase, 73 ready, 74 MycStatusReasonCodes::empty(), 75 1, 76 build_info(), 77 configuration, 78 persistence, 79 ) 80 .expect("common status"), 81 provider, 82 transport, 83 MycConnectionCountsV1::default(), 84 MycOutboxStatusV1::new(0, 0, None), 85 ) 86 } 87 88 fn available_port() -> u16 { 89 let listener = 90 TcpListener::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)).expect("ephemeral listener"); 91 listener.local_addr().expect("local address").port() 92 } 93 94 fn enabled_config(port: u16) -> String { 95 CONFIG.replacen( 96 "[operations]\nenabled = false", 97 &format!( 98 "[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" 99 ), 100 1, 101 ) 102 } 103 104 async fn raw_request(address: std::net::SocketAddr, request: &[u8]) -> Vec<u8> { 105 let mut stream = TcpStream::connect(address).await.expect("connect"); 106 stream.write_all(request).await.expect("request write"); 107 let mut response = Vec::new(); 108 stream 109 .read_to_end(&mut response) 110 .await 111 .expect("response read"); 112 response 113 } 114 115 fn response_text(response: &[u8]) -> &str { 116 std::str::from_utf8(response).expect("response UTF-8") 117 } 118 119 #[test] 120 fn machine_contract_freezes_exact_routes_metrics_and_non_authority() { 121 let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract"); 122 assert_eq!(contract["schema"], "radroots.myc.tcp-operations.v1"); 123 assert_eq!( 124 contract["contract_version"], 125 MYC_OPERATIONS_CONTRACT_VERSION 126 ); 127 assert_eq!(contract["step"], 155); 128 assert_eq!( 129 contract["routes"] 130 .as_array() 131 .expect("routes") 132 .iter() 133 .map(|route| route["path"].as_str().expect("route path")) 134 .collect::<Vec<_>>(), 135 [MYC_LIVEZ_PATH, MYC_READYZ_PATH, MYC_METRICS_PATH] 136 ); 137 assert_eq!(contract["route_registration_extension"], false); 138 assert_eq!(contract["metrics"]["high_cardinality_labels"], false); 139 assert_eq!(contract["metrics"]["arbitrary_labels"], false); 140 } 141 142 #[tokio::test] 143 async fn exact_tcp_routes_use_only_latest_cached_lifecycle_and_metrics() { 144 let config = parse_myc_config_v1( 145 enabled_config(available_port()).as_bytes(), 146 MycConfigProfile::Production, 147 ) 148 .expect("enabled config"); 149 let (mut publisher, reader) = myc_status_cache( 150 InstanceId::new("primary").expect("instance"), 151 observation(MycServicePhase::Ready, true), 152 ) 153 .expect("status cache"); 154 let detail_pointer = reader.snapshot().detailed_status_json().as_ptr(); 155 let bound = MycOperationsServer::new(&config, &reader) 156 .expect("operations server") 157 .bind() 158 .await 159 .expect("bind"); 160 let address = bound.local_address(); 161 let cancellation = MycOperationsCancellationToken::new(); 162 let task = tokio::spawn(bound.serve(cancellation.clone())); 163 164 let live = raw_request(address, b"GET /livez HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 165 let ready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 166 let metrics = raw_request(address, b"GET /metrics HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 167 assert!(response_text(&live).starts_with("HTTP/1.1 200 OK\r\n")); 168 assert!(response_text(&live).ends_with("live\n")); 169 assert!(response_text(&ready).ends_with("ready\n")); 170 assert!(response_text(&metrics).contains("# TYPE radroots_myc_service_phase gauge\n")); 171 assert!(response_text(&metrics).contains("radroots_myc_service_phase{phase=\"ready\"} 1\n")); 172 assert!(response_text(&metrics).contains("radroots_myc_service_ready 1\n")); 173 174 for request in [ 175 &b"GET /status HTTP/1.1\r\nhost: localhost\r\n\r\n"[..], 176 &b"GET /readyz?probe=1 HTTP/1.1\r\nhost: localhost\r\n\r\n"[..], 177 &b"POST /metrics HTTP/1.1\r\nhost: localhost\r\ncontent-length: 0\r\n\r\n"[..], 178 &b"GET /v1/status HTTP/1.1\r\nhost: localhost\r\n\r\n"[..], 179 ] { 180 let rejected = raw_request(address, request).await; 181 assert!(response_text(&rejected).starts_with("HTTP/1.1 404 Not Found\r\n")); 182 } 183 184 publisher 185 .publish(observation(MycServicePhase::Unready, false)) 186 .expect("unready publication"); 187 let unready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 188 let metrics = raw_request(address, b"GET /metrics HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 189 assert!(response_text(&unready).starts_with("HTTP/1.1 503 Service Unavailable\r\n")); 190 assert!(response_text(&unready).ends_with("unready\n")); 191 assert!(response_text(&metrics).contains("radroots_myc_service_phase{phase=\"unready\"} 1\n")); 192 assert!(response_text(&metrics).contains("radroots_myc_service_ready 0\n")); 193 assert_ne!( 194 reader.snapshot().detailed_status_json().as_ptr(), 195 detail_pointer 196 ); 197 198 cancellation.cancel(); 199 assert_eq!(task.await.expect("serve task"), Ok(())); 200 } 201 202 #[tokio::test] 203 async fn disabled_invalid_and_bind_failures_are_typed_source_free_and_redacted() { 204 let disabled = parse_myc_config_v1(CONFIG.as_bytes(), MycConfigProfile::Production) 205 .expect("disabled config"); 206 let (_, reader) = myc_status_cache( 207 InstanceId::new("primary").expect("instance"), 208 observation(MycServicePhase::Ready, true), 209 ) 210 .expect("status cache"); 211 let disabled_error = 212 MycOperationsServer::new(&disabled, &reader).expect_err("disabled operations"); 213 assert_eq!(disabled_error.kind(), MycOperationsErrorKind::Disabled); 214 assert_eq!(disabled_error.code(), "operations_disabled"); 215 assert!(Error::source(&disabled_error).is_none()); 216 217 let below_floor = 218 enabled_config(available_port()).replace("header_bytes = 8192", "header_bytes = 8191"); 219 assert!(parse_myc_config_v1(below_floor.as_bytes(), MycConfigProfile::Production).is_err()); 220 221 let occupied = 222 TcpListener::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)).expect("occupied listener"); 223 let port = occupied.local_addr().expect("occupied address").port(); 224 let config = parse_myc_config_v1( 225 enabled_config(port).as_bytes(), 226 MycConfigProfile::Production, 227 ) 228 .expect("enabled config"); 229 let server = MycOperationsServer::new(&config, &reader).expect("server"); 230 assert!(!format!("{server:?}").contains(&port.to_string())); 231 let bind_error = server.bind().await.expect_err("occupied bind"); 232 assert_eq!(bind_error.kind(), MycOperationsErrorKind::Bind); 233 assert!(Error::source(&bind_error).is_none()); 234 assert!(!format!("{bind_error:?} {bind_error}").contains(&port.to_string())); 235 }