rhi

Coordinated trade for connected markets
git clone https://radroots.dev/git/rhi.git
Log | Files | Refs | README | LICENSE

operations_v1.rs (8970B)


      1 //! Rhi-owned adapter for the exact passive TCP operations surface.
      2 
      3 use core::{fmt, time::Duration};
      4 use std::{error::Error, net::SocketAddr};
      5 
      6 use radroots_service_host::{
      7     BoundOperationsServer as HostBoundOperationsServer, CancellationToken as HostCancellationToken,
      8     OperationsBindPolicy as HostOperationsBindPolicy,
      9     OperationsListenAddress as HostOperationsListenAddress,
     10     OperationsListenerConfig as HostOperationsListenerConfig,
     11     OperationsServer as HostOperationsServer, OperationsServerError as HostOperationsServerError,
     12     OperationsTransportLimitValues as HostOperationsTransportLimitValues,
     13     OperationsTransportLimits as HostOperationsTransportLimits,
     14 };
     15 use serde_json::Value;
     16 
     17 use crate::{RhiConfigDocumentV1, RhiStatusReader};
     18 
     19 /// Exact Rhi TCP operations contract version.
     20 pub const RHI_OPERATIONS_CONTRACT_VERSION: u32 = 1;
     21 
     22 /// Exact liveness route exposed by the optional TCP listener.
     23 pub const RHI_LIVEZ_PATH: &str = "/livez";
     24 
     25 /// Exact readiness route exposed by the optional TCP listener.
     26 pub const RHI_READYZ_PATH: &str = "/readyz";
     27 
     28 /// Exact bounded metrics route exposed by the optional TCP listener.
     29 pub const RHI_METRICS_PATH: &str = "/metrics";
     30 
     31 /// Cloneable cooperative cancellation owned by the Rhi runtime supervisor.
     32 #[derive(Clone, Default)]
     33 pub struct RhiOperationsCancellationToken {
     34     inner: HostCancellationToken,
     35 }
     36 
     37 impl RhiOperationsCancellationToken {
     38     #[must_use]
     39     pub fn new() -> Self {
     40         Self::default()
     41     }
     42 
     43     /// Requests cancellation. Repeated requests have no additional effect.
     44     pub fn cancel(&self) {
     45         self.inner.cancel();
     46     }
     47 
     48     #[must_use]
     49     pub fn is_cancelled(&self) -> bool {
     50         self.inner.is_cancelled()
     51     }
     52 }
     53 
     54 impl fmt::Debug for RhiOperationsCancellationToken {
     55     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     56         formatter
     57             .debug_struct("RhiOperationsCancellationToken")
     58             .field("cancelled", &self.is_cancelled())
     59             .finish()
     60     }
     61 }
     62 
     63 /// Unbound exact-route Rhi TCP operations server.
     64 pub struct RhiOperationsServer {
     65     inner: HostOperationsServer,
     66 }
     67 
     68 impl RhiOperationsServer {
     69     /// Projects the already-validated Rhi configuration and passive status cache.
     70     pub fn new(
     71         config: &RhiConfigDocumentV1,
     72         status: &RhiStatusReader,
     73     ) -> Result<Self, RhiOperationsError> {
     74         let listener = listener_config(config)?;
     75         HostOperationsServer::new(listener, status.operations_cache())
     76             .map(|inner| Self { inner })
     77             .map_err(map_server_error)
     78     }
     79 
     80     /// Binds the exact configured address without starting admission.
     81     pub async fn bind(self) -> Result<RhiBoundOperationsServer, RhiOperationsError> {
     82         self.inner
     83             .bind()
     84             .await
     85             .map(|inner| RhiBoundOperationsServer { inner })
     86             .map_err(map_server_error)
     87     }
     88 }
     89 
     90 impl fmt::Debug for RhiOperationsServer {
     91     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     92         formatter.write_str("RhiOperationsServer([sealed])")
     93     }
     94 }
     95 
     96 /// Successfully bound exact-route Rhi TCP operations server.
     97 pub struct RhiBoundOperationsServer {
     98     inner: HostBoundOperationsServer,
     99 }
    100 
    101 impl RhiBoundOperationsServer {
    102     #[must_use]
    103     pub fn local_address(&self) -> SocketAddr {
    104         self.inner.local_address()
    105     }
    106 
    107     /// Serves until explicit supervisor cancellation, then drains owned work.
    108     pub async fn serve(
    109         self,
    110         cancellation: RhiOperationsCancellationToken,
    111     ) -> Result<(), RhiOperationsError> {
    112         self.inner
    113             .serve(cancellation.inner)
    114             .await
    115             .map_err(map_server_error)
    116     }
    117 }
    118 
    119 impl fmt::Debug for RhiBoundOperationsServer {
    120     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    121         formatter.write_str("RhiBoundOperationsServer([sealed])")
    122     }
    123 }
    124 
    125 /// Stable source-free Rhi operations failure classification.
    126 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    127 pub enum RhiOperationsErrorKind {
    128     Disabled,
    129     InvalidConfiguration,
    130     Bind,
    131     LocalAddress,
    132     Accept,
    133     ConnectionTaskPanicked,
    134 }
    135 
    136 impl RhiOperationsErrorKind {
    137     #[must_use]
    138     pub const fn code(self) -> &'static str {
    139         match self {
    140             Self::Disabled => "operations_disabled",
    141             Self::InvalidConfiguration => "operations_configuration_invalid",
    142             Self::Bind => "operations_bind_failed",
    143             Self::LocalAddress => "operations_local_address_failed",
    144             Self::Accept => "operations_accept_failed",
    145             Self::ConnectionTaskPanicked => "operations_connection_task_panicked",
    146         }
    147     }
    148 }
    149 
    150 /// One redacted source-free Rhi operations failure.
    151 #[derive(Clone, Copy, PartialEq, Eq)]
    152 pub struct RhiOperationsError {
    153     kind: RhiOperationsErrorKind,
    154 }
    155 
    156 impl RhiOperationsError {
    157     const fn new(kind: RhiOperationsErrorKind) -> Self {
    158         Self { kind }
    159     }
    160 
    161     #[must_use]
    162     pub const fn kind(self) -> RhiOperationsErrorKind {
    163         self.kind
    164     }
    165 
    166     #[must_use]
    167     pub const fn code(self) -> &'static str {
    168         self.kind.code()
    169     }
    170 }
    171 
    172 impl fmt::Debug for RhiOperationsError {
    173     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    174         formatter
    175             .debug_struct("RhiOperationsError")
    176             .field("kind", &self.kind)
    177             .finish()
    178     }
    179 }
    180 
    181 impl fmt::Display for RhiOperationsError {
    182     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    183         formatter.write_str("Rhi TCP operations failed")
    184     }
    185 }
    186 
    187 impl Error for RhiOperationsError {}
    188 
    189 fn listener_config(
    190     config: &RhiConfigDocumentV1,
    191 ) -> Result<HostOperationsListenerConfig, RhiOperationsError> {
    192     let operations = config
    193         .normalized()
    194         .pointer("/operations")
    195         .ok_or_else(invalid_configuration)?;
    196     if !boolean(operations, "/enabled")? {
    197         return Err(RhiOperationsError::new(RhiOperationsErrorKind::Disabled));
    198     }
    199     let listen = string(operations, "/listen")?
    200         .parse::<SocketAddr>()
    201         .map_err(|_| invalid_configuration())?;
    202     let listen = HostOperationsListenAddress::new(listen).map_err(|_| invalid_configuration())?;
    203     let bind_policy = match string(operations, "/bind_policy")? {
    204         "loopback_only" => HostOperationsBindPolicy::LoopbackOnly,
    205         "explicit_public" => HostOperationsBindPolicy::Public,
    206         _ => return Err(invalid_configuration()),
    207     };
    208     let values = HostOperationsTransportLimitValues {
    209         header_count: unsigned_u32(operations, "/limits/header_count")?,
    210         header_bytes: unsigned_u32(operations, "/limits/header_bytes")?,
    211         response_body_utf8_bytes: unsigned_u32(operations, "/limits/response_body_utf8_bytes")?,
    212         concurrent_connections: unsigned_u32(operations, "/limits/concurrent_connections")?,
    213         request_deadline: Duration::from_millis(unsigned(
    214             operations,
    215             "/limits/request_deadline_ms",
    216         )?),
    217         idle_timeout: Duration::from_millis(unsigned(operations, "/limits/idle_timeout_ms")?),
    218     };
    219     let limits = HostOperationsTransportLimits::new(values).map_err(|_| invalid_configuration())?;
    220     HostOperationsListenerConfig::enabled(listen, bind_policy, limits)
    221         .map_err(|_| invalid_configuration())
    222 }
    223 
    224 fn boolean(value: &Value, pointer: &str) -> Result<bool, RhiOperationsError> {
    225     value
    226         .pointer(pointer)
    227         .and_then(Value::as_bool)
    228         .ok_or_else(invalid_configuration)
    229 }
    230 
    231 fn string<'a>(value: &'a Value, pointer: &str) -> Result<&'a str, RhiOperationsError> {
    232     value
    233         .pointer(pointer)
    234         .and_then(Value::as_str)
    235         .ok_or_else(invalid_configuration)
    236 }
    237 
    238 fn unsigned(value: &Value, pointer: &str) -> Result<u64, RhiOperationsError> {
    239     value
    240         .pointer(pointer)
    241         .and_then(Value::as_u64)
    242         .ok_or_else(invalid_configuration)
    243 }
    244 
    245 fn unsigned_u32(value: &Value, pointer: &str) -> Result<u32, RhiOperationsError> {
    246     u32::try_from(unsigned(value, pointer)?).map_err(|_| invalid_configuration())
    247 }
    248 
    249 const fn invalid_configuration() -> RhiOperationsError {
    250     RhiOperationsError::new(RhiOperationsErrorKind::InvalidConfiguration)
    251 }
    252 
    253 const fn map_server_error(error: HostOperationsServerError) -> RhiOperationsError {
    254     let kind = match error {
    255         HostOperationsServerError::Disabled => RhiOperationsErrorKind::Disabled,
    256         HostOperationsServerError::HeaderLimitBelowParserFloor => {
    257             RhiOperationsErrorKind::InvalidConfiguration
    258         }
    259         HostOperationsServerError::Bind { .. } => RhiOperationsErrorKind::Bind,
    260         HostOperationsServerError::LocalAddress { .. } => RhiOperationsErrorKind::LocalAddress,
    261         HostOperationsServerError::Accept { .. } => RhiOperationsErrorKind::Accept,
    262         HostOperationsServerError::ConnectionTaskPanicked => {
    263             RhiOperationsErrorKind::ConnectionTaskPanicked
    264         }
    265     };
    266     RhiOperationsError::new(kind)
    267 }