status.rs (8633B)
1 //! Context-bound status-v1 access over the shared Unix admin client. 2 3 use core::fmt; 4 5 use radroots_runtime_paths::{ 6 InstanceId, RuntimeContext, ServiceId, default_service_instance_artifacts, 7 }; 8 use radroots_service_host::{ 9 AdminClient, AdminClientTarget, AdminTransportLimits, SERVICE_STATUS_CONTRACT_VERSION, 10 }; 11 use serde::{Serialize, de::DeserializeOwned}; 12 13 use crate::RadrootsRuntimeManagerError; 14 15 const STATUS_V1_TARGET: &str = "/v1/status"; 16 17 /// Identity projection required from a service-owned status-v1 response. 18 /// 19 /// Myc and RHI retain ownership of their typed status details. This trait lets 20 /// the manager validate only the shared contract and selected runtime identity 21 /// without accepting an untyped JSON payload. 22 pub trait ManagedServiceStatusV1: Serialize + DeserializeOwned { 23 fn contract_version(&self) -> u32; 24 fn service_id(&self) -> &ServiceId; 25 fn instance_id(&self) -> &InstanceId; 26 } 27 28 /// Bounded status-v1 client sealed to one [`RuntimeContext`] admin socket. 29 pub struct ManagedRuntimeStatusClient { 30 expected_service: ServiceId, 31 expected_instance: InstanceId, 32 client: AdminClient, 33 target: AdminClientTarget, 34 } 35 36 impl ManagedRuntimeStatusClient { 37 pub(crate) fn for_context( 38 context: &RuntimeContext, 39 limits: AdminTransportLimits, 40 ) -> Result<Self, RadrootsRuntimeManagerError> { 41 let artifacts = default_service_instance_artifacts(context.paths()); 42 let client = AdminClient::new(artifacts.admin_socket(), limits) 43 .map_err(|_| RadrootsRuntimeManagerError::AdminClient)?; 44 let target = AdminClientTarget::new(STATUS_V1_TARGET) 45 .map_err(|_| RadrootsRuntimeManagerError::AdminClient)?; 46 Ok(Self { 47 expected_service: context.service().clone(), 48 expected_instance: context.instance().clone(), 49 client, 50 target, 51 }) 52 } 53 54 /// Requests and identity-validates the service-owned typed status model. 55 pub async fn get<S>(&self) -> Result<S, RadrootsRuntimeManagerError> 56 where 57 S: ManagedServiceStatusV1, 58 { 59 let status = self 60 .client 61 .get::<S>(&self.target) 62 .await 63 .map_err(|_| RadrootsRuntimeManagerError::AdminRequest)? 64 .into_result(); 65 if status.contract_version() != SERVICE_STATUS_CONTRACT_VERSION 66 || status.service_id() != &self.expected_service 67 || status.instance_id() != &self.expected_instance 68 { 69 return Err(RadrootsRuntimeManagerError::StatusContractMismatch); 70 } 71 Ok(status) 72 } 73 } 74 75 impl fmt::Debug for ManagedRuntimeStatusClient { 76 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 77 formatter 78 .debug_struct("ManagedRuntimeStatusClient") 79 .field("expected_service", &self.expected_service) 80 .field("expected_instance", &self.expected_instance) 81 .field("socket_path", &"[redacted]") 82 .field("target", &STATUS_V1_TARGET) 83 .finish() 84 } 85 } 86 87 #[cfg(test)] 88 mod tests { 89 use std::path::PathBuf; 90 91 use radroots_runtime_paths::{ 92 InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, 93 RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, ServiceId, 94 }; 95 use radroots_service_host::AdminTransportLimits; 96 use serde::{Deserialize, Serialize}; 97 use tempfile::Builder; 98 use tokio::{ 99 io::{AsyncReadExt, AsyncWriteExt}, 100 net::UnixListener, 101 }; 102 103 use super::{ManagedRuntimeStatusClient, ManagedServiceStatusV1}; 104 105 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)] 106 struct TestStatus { 107 contract_version: u32, 108 service: ServiceId, 109 instance: InstanceId, 110 phase: String, 111 } 112 113 impl ManagedServiceStatusV1 for TestStatus { 114 fn contract_version(&self) -> u32 { 115 self.contract_version 116 } 117 118 fn service_id(&self) -> &ServiceId { 119 &self.service 120 } 121 122 fn instance_id(&self) -> &InstanceId { 123 &self.instance 124 } 125 } 126 127 fn context(base: PathBuf) -> RuntimeContext { 128 RuntimeContext::resolve( 129 &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), 130 RuntimeContextBootstrap::new( 131 RadrootsPathProfile::RepoLocal, 132 Some(base), 133 RuntimeContextSource::BootstrapCli, 134 RuntimeContextSource::BootstrapCli, 135 ) 136 .expect("bootstrap"), 137 ServiceId::new("myc").expect("service"), 138 InstanceId::new("primary").expect("instance"), 139 ) 140 .expect("context") 141 } 142 143 async fn serve_status_once(socket: PathBuf, result: serde_json::Value) { 144 if socket.exists() { 145 std::fs::remove_file(&socket).expect("remove prior socket"); 146 } 147 let listener = UnixListener::bind(socket).expect("bind status socket"); 148 let (mut stream, _) = listener.accept().await.expect("accept status request"); 149 let mut request = [0_u8; 4_096]; 150 let read = stream.read(&mut request).await.expect("read request"); 151 let request = std::str::from_utf8(&request[..read]).expect("UTF-8 request"); 152 assert!(request.starts_with("GET /v1/status HTTP/1.1\r\n")); 153 154 let body = serde_json::to_vec(&serde_json::json!({ 155 "contract_version": 1, 156 "ok": true, 157 "correlation_id": "status-test-01", 158 "result": result, 159 })) 160 .expect("response body"); 161 let head = format!( 162 "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", 163 body.len() 164 ); 165 stream 166 .write_all(head.as_bytes()) 167 .await 168 .expect("write response head"); 169 stream.write_all(&body).await.expect("write response body"); 170 stream.shutdown().await.expect("shutdown response"); 171 } 172 173 #[test] 174 fn construction_is_context_bound_and_debug_redacts_the_socket() { 175 let client = ManagedRuntimeStatusClient::for_context( 176 &context(PathBuf::from("/sensitive/project-root")), 177 AdminTransportLimits::DEFAULT, 178 ) 179 .expect("client"); 180 let debug = format!("{client:?}"); 181 assert!(debug.contains("/v1/status")); 182 assert!(!debug.contains("sensitive")); 183 assert!(!debug.contains("admin.sock")); 184 } 185 186 #[tokio::test(flavor = "current_thread")] 187 async fn bounded_unix_status_round_trip_validates_the_complete_common_identity() { 188 let temp = Builder::new() 189 .prefix("rrm") 190 .tempdir_in("/tmp") 191 .expect("short temp root"); 192 let context = context(temp.path().to_path_buf()); 193 let socket = radroots_runtime_paths::default_service_instance_artifacts(context.paths()) 194 .admin_socket() 195 .to_path_buf(); 196 std::fs::create_dir_all(socket.parent().expect("socket parent")) 197 .expect("create socket parent"); 198 let client = 199 ManagedRuntimeStatusClient::for_context(&context, AdminTransportLimits::DEFAULT) 200 .expect("client"); 201 202 let success_server = tokio::spawn(serve_status_once( 203 socket.clone(), 204 serde_json::json!({ 205 "contract_version": 1, 206 "service": "myc", 207 "instance": "primary", 208 "phase": "ready", 209 }), 210 )); 211 tokio::task::yield_now().await; 212 let status = client.get::<TestStatus>().await.expect("status"); 213 assert_eq!(status.phase, "ready"); 214 success_server.await.expect("success server"); 215 216 for (field, replacement) in [ 217 ("contract_version", serde_json::json!(2)), 218 ("service", serde_json::json!("rhi")), 219 ("instance", serde_json::json!("secondary")), 220 ] { 221 let mut result = serde_json::json!({ 222 "contract_version": 1, 223 "service": "myc", 224 "instance": "primary", 225 "phase": "ready", 226 }); 227 result[field] = replacement; 228 let server = tokio::spawn(serve_status_once(socket.clone(), result)); 229 tokio::task::yield_now().await; 230 assert_eq!( 231 client.get::<TestStatus>().await, 232 Err(crate::RadrootsRuntimeManagerError::StatusContractMismatch) 233 ); 234 server.await.expect("mismatch server"); 235 } 236 } 237 }