control_plane_wave_090_a.rs (10411B)
1 //! Native test-only integration gate for RCLD-RSHR-090 wave 090-a. 2 3 use std::{ 4 collections::BTreeMap, 5 net::{Ipv4Addr, SocketAddrV4, TcpListener}, 6 }; 7 8 use tokio::{ 9 io::{AsyncReadExt, AsyncWriteExt}, 10 net::TcpStream, 11 }; 12 13 use crate::{ 14 InstanceId, MycBootstrapProfileV1, MycCliOfflineOperationV1, MycCliPrimaryAuthorityV1, 15 MycConfigProfile, MycConnectionCountsV1, MycDoctorCheckDefinition, MycDoctorCheckId, 16 MycDoctorFuture, MycDoctorObservation, MycDoctorProbe, MycIdentityHealthV1, 17 MycIntegrityStateV1, MycLogRecord, MycOperationsCancellationToken, MycOperationsServer, 18 MycOutboxStatusV1, MycPersistenceHealthV1, MycPersistenceStatusV1, MycProcessResult, 19 MycProviderStatusV1, MycRelayTransportStatusV1, MycServicePhase, MycStatusBuildInfoV1, 20 MycStatusBuildMode, MycStatusCommonV1, MycStatusConfigurationIdentityV1, 21 MycStatusConfigurationSource, MycStatusObservationV1, MycStatusReasonCodes, 22 MycTransportHealthV1, RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, 23 myc_status_cache, parse_myc_cli_v1_from, parse_myc_config_v1, plan_myc_cli_v1, 24 resolve_myc_runtime_context, run_myc_doctor, 25 }; 26 27 const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); 28 const CORPUS: &str = 29 include_str!("../contracts/services_hardening/control_plane_wave_090_a.v1.json"); 30 const SERVICE_REVISION: &str = "0123456789abcdef0123456789abcdef01234567"; 31 const LIB_REVISION: &str = "89abcdef0123456789abcdef0123456789abcdef"; 32 33 struct FixedProbe { 34 observations: BTreeMap<MycDoctorCheckId, MycDoctorObservation>, 35 } 36 37 impl FixedProbe { 38 fn all(observation: MycDoctorObservation) -> Self { 39 Self { 40 observations: crate::myc_doctor_check_definitions() 41 .iter() 42 .map(|definition| (definition.id(), observation)) 43 .collect(), 44 } 45 } 46 47 fn with(mut self, id: MycDoctorCheckId, observation: MycDoctorObservation) -> Self { 48 self.observations.insert(id, observation); 49 self 50 } 51 } 52 53 impl MycDoctorProbe for FixedProbe { 54 fn probe(&self, definition: MycDoctorCheckDefinition) -> MycDoctorFuture<'_> { 55 let observation = self.observations[&definition.id()]; 56 Box::pin(async move { observation }) 57 } 58 } 59 60 fn enabled_config(port: u16) -> String { 61 CONFIG.replacen( 62 "[operations]\nenabled = false", 63 &format!( 64 "[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" 65 ), 66 1, 67 ) 68 } 69 70 fn available_port() -> u16 { 71 TcpListener::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)) 72 .expect("ephemeral listener") 73 .local_addr() 74 .expect("listener address") 75 .port() 76 } 77 78 fn status_observation(phase: MycServicePhase, ready: bool) -> MycStatusObservationV1 { 79 let empty = MycStatusReasonCodes::empty; 80 let build = MycStatusBuildInfoV1::new( 81 MycStatusBuildMode::Release, 82 Some(env!("CARGO_PKG_VERSION")), 83 Some(SERVICE_REVISION), 84 Some(LIB_REVISION), 85 Some("1.97.1"), 86 Some("x86_64-unknown-linux-gnu"), 87 Some("service-host"), 88 ) 89 .expect("build info"); 90 let configuration = MycStatusConfigurationIdentityV1::new( 91 "a".repeat(64), 92 MycStatusConfigurationSource::ExplicitConfig, 93 ) 94 .expect("configuration identity"); 95 let persistence = MycPersistenceStatusV1::new( 96 MycPersistenceHealthV1::Ready, 97 9, 98 1, 99 MycIntegrityStateV1::Verified, 100 empty(), 101 ) 102 .expect("persistence status"); 103 let identity = || MycIdentityHealthV1::new(true, true, empty()).expect("identity health"); 104 let provider = MycProviderStatusV1::new(identity(), identity(), identity(), empty()) 105 .expect("provider status"); 106 let transport = MycRelayTransportStatusV1::new(MycTransportHealthV1::Ready, true, 2, empty()) 107 .expect("transport status"); 108 MycStatusObservationV1::new( 109 MycStatusCommonV1::new(phase, ready, empty(), 1, build, configuration, persistence) 110 .expect("common status"), 111 provider, 112 transport, 113 MycConnectionCountsV1::default(), 114 MycOutboxStatusV1::default(), 115 ) 116 } 117 118 async fn raw_request(address: std::net::SocketAddr, request: &[u8]) -> String { 119 let mut stream = TcpStream::connect(address).await.expect("connect"); 120 stream.write_all(request).await.expect("request write"); 121 let mut response = Vec::new(); 122 stream 123 .read_to_end(&mut response) 124 .await 125 .expect("response read"); 126 String::from_utf8(response).expect("response UTF-8") 127 } 128 129 #[test] 130 fn machine_corpus_freezes_the_accumulated_wave_and_deferrals() { 131 let corpus: serde_json::Value = serde_json::from_str(CORPUS).expect("wave corpus"); 132 assert_eq!(corpus["schema"], "radroots.myc.control-plane-wave-090-a.v1"); 133 assert_eq!(corpus["contract_version"], 1); 134 assert_eq!( 135 corpus["steps"], 136 serde_json::json!([152, 153, 154, 155, 156, 157]) 137 ); 138 assert_eq!(corpus["authority"]["cli_parse_count"], 1); 139 assert_eq!( 140 corpus["authority"]["tcp_routes"], 141 serde_json::json!(["/livez", "/readyz", "/metrics"]) 142 ); 143 assert_eq!(corpus["gate"]["wave"], "090-a"); 144 assert_eq!(corpus["gate"]["complete_after_step"], 157); 145 assert_eq!(corpus["gate"]["rcld_promotion_owner"], 162); 146 assert!( 147 corpus["nonclaims"] 148 .as_array() 149 .expect("nonclaims") 150 .iter() 151 .any(|value| value == "runtime_task_supervision") 152 ); 153 } 154 155 #[tokio::test] 156 async fn one_parse_runtime_doctor_and_diagnostics_share_the_exact_exit_contract() { 157 let root = tempfile::tempdir().expect("temporary root"); 158 let invocation = parse_myc_cli_v1_from([ 159 "myc", 160 "--profile", 161 "repo-local", 162 "--instance", 163 "primary", 164 "--repo-local-root", 165 root.path().to_str().expect("UTF-8 root"), 166 "doctor", 167 ]) 168 .expect("doctor CLI"); 169 assert_eq!(invocation.profile(), MycBootstrapProfileV1::RepoLocal); 170 let plan = plan_myc_cli_v1(&invocation); 171 assert_eq!(plan.primary_authority(), MycCliPrimaryAuthorityV1::Offline); 172 assert_eq!( 173 plan.offline_operation(), 174 Some(MycCliOfflineOperationV1::Doctor) 175 ); 176 let runtime = resolve_myc_runtime_context( 177 &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), 178 &invocation, 179 ) 180 .expect("runtime context"); 181 182 let pass = run_myc_doctor(&runtime, &FixedProbe::all(MycDoctorObservation::Pass)) 183 .await 184 .expect("passing doctor"); 185 assert_eq!(pass.exit_code(), MycProcessResult::Success.exit_code_u8()); 186 187 let fail = run_myc_doctor( 188 &runtime, 189 &FixedProbe::all(MycDoctorObservation::Pass) 190 .with(MycDoctorCheckId::WriterLock, MycDoctorObservation::Fail), 191 ) 192 .await 193 .expect("failing doctor report"); 194 assert_eq!( 195 fail.exit_code(), 196 MycProcessResult::DoctorRequiredCheckFailed.exit_code_u8() 197 ); 198 assert_eq!( 199 MycLogRecord::process_result(MycProcessResult::DoctorRequiredCheckFailed).to_string(), 200 r#"{"schema":"radroots.myc.log.v1","contract_version":1,"service":"myc","level":"error","event":"process_result","code":"doctor_required_check_failed","exit_code":6}"# 201 ); 202 203 let secret = "secret-private-key-path"; 204 let rejected = parse_myc_cli_v1_from(["myc", &format!("--credential={secret}")]) 205 .expect_err("ungoverned argument"); 206 assert!(!format!("{rejected:?} {rejected}").contains(secret)); 207 } 208 209 #[tokio::test] 210 async fn live_status_and_operations_share_only_the_latest_passive_projection() { 211 let live = parse_myc_cli_v1_from([ 212 "myc", 213 "--profile", 214 "service-host", 215 "--instance", 216 "primary", 217 "status", 218 ]) 219 .expect("live status CLI"); 220 let plan = plan_myc_cli_v1(&live); 221 assert_eq!( 222 plan.primary_authority(), 223 MycCliPrimaryAuthorityV1::LiveUnixAdmin 224 ); 225 assert_eq!( 226 plan.offline_operation(), 227 Some(MycCliOfflineOperationV1::StateReadOnly) 228 ); 229 230 let config = parse_myc_config_v1( 231 enabled_config(available_port()).as_bytes(), 232 MycConfigProfile::Production, 233 ) 234 .expect("operations configuration"); 235 let (mut publisher, reader) = myc_status_cache( 236 InstanceId::new("primary").expect("instance"), 237 status_observation(MycServicePhase::Unready, false), 238 ) 239 .expect("status cache"); 240 let initial_detail = reader.snapshot().detailed_status_json().to_vec(); 241 let bound = MycOperationsServer::new(&config, &reader) 242 .expect("operations server") 243 .bind() 244 .await 245 .expect("operations bind"); 246 let address = bound.local_address(); 247 let cancellation = MycOperationsCancellationToken::new(); 248 let task = tokio::spawn(bound.serve(cancellation.clone())); 249 250 let unready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 251 let detailed = raw_request( 252 address, 253 b"GET /v1/status HTTP/1.1\r\nhost: localhost\r\n\r\n", 254 ) 255 .await; 256 assert!(unready.starts_with("HTTP/1.1 503 Service Unavailable\r\n")); 257 assert!(detailed.starts_with("HTTP/1.1 404 Not Found\r\n")); 258 259 publisher 260 .publish(status_observation(MycServicePhase::Ready, true)) 261 .expect("ready publication"); 262 let ready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 263 let metrics = raw_request(address, b"GET /metrics HTTP/1.1\r\nhost: localhost\r\n\r\n").await; 264 let query = raw_request( 265 address, 266 b"GET /readyz?probe=1 HTTP/1.1\r\nhost: localhost\r\n\r\n", 267 ) 268 .await; 269 assert!(ready.starts_with("HTTP/1.1 200 OK\r\n")); 270 assert!(ready.ends_with("ready\n")); 271 assert!(metrics.contains("radroots_myc_service_phase{phase=\"ready\"} 1\n")); 272 assert!(metrics.contains("radroots_myc_service_ready 1\n")); 273 assert!(query.starts_with("HTTP/1.1 404 Not Found\r\n")); 274 assert_ne!(reader.snapshot().detailed_status_json(), initial_detail); 275 276 cancellation.cancel(); 277 assert_eq!(task.await.expect("server join"), Ok(())); 278 }