lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

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 }