operations_v1.rs (8961B)
1 //! Myc-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::{MycConfigDocumentV1, MycStatusReader}; 18 19 /// Exact Myc TCP operations contract version. 20 pub const MYC_OPERATIONS_CONTRACT_VERSION: u32 = 1; 21 22 /// Exact liveness route exposed by the optional TCP listener. 23 pub const MYC_LIVEZ_PATH: &str = "/livez"; 24 25 /// Exact readiness route exposed by the optional TCP listener. 26 pub const MYC_READYZ_PATH: &str = "/readyz"; 27 28 /// Exact bounded metrics route exposed by the optional TCP listener. 29 pub const MYC_METRICS_PATH: &str = "/metrics"; 30 31 /// Cloneable cooperative cancellation owned by the Myc runtime supervisor. 32 #[derive(Clone, Default)] 33 pub struct MycOperationsCancellationToken { 34 inner: HostCancellationToken, 35 } 36 37 impl MycOperationsCancellationToken { 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 MycOperationsCancellationToken { 55 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 56 formatter 57 .debug_struct("MycOperationsCancellationToken") 58 .field("cancelled", &self.is_cancelled()) 59 .finish() 60 } 61 } 62 63 /// Unbound exact-route Myc TCP operations server. 64 pub struct MycOperationsServer { 65 inner: HostOperationsServer, 66 } 67 68 impl MycOperationsServer { 69 /// Projects the already-validated Myc configuration and passive status cache. 70 pub fn new( 71 config: &MycConfigDocumentV1, 72 status: &MycStatusReader, 73 ) -> Result<Self, MycOperationsError> { 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<MycBoundOperationsServer, MycOperationsError> { 82 self.inner 83 .bind() 84 .await 85 .map(|inner| MycBoundOperationsServer { inner }) 86 .map_err(map_server_error) 87 } 88 } 89 90 impl fmt::Debug for MycOperationsServer { 91 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 92 formatter.write_str("MycOperationsServer([sealed])") 93 } 94 } 95 96 /// Successfully bound exact-route Myc TCP operations server. 97 pub struct MycBoundOperationsServer { 98 inner: HostBoundOperationsServer, 99 } 100 101 impl MycBoundOperationsServer { 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: MycOperationsCancellationToken, 111 ) -> Result<(), MycOperationsError> { 112 self.inner 113 .serve(cancellation.inner) 114 .await 115 .map_err(map_server_error) 116 } 117 } 118 119 impl fmt::Debug for MycBoundOperationsServer { 120 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 121 formatter.write_str("MycBoundOperationsServer([sealed])") 122 } 123 } 124 125 /// Stable source-free Myc operations failure classification. 126 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 127 pub enum MycOperationsErrorKind { 128 Disabled, 129 InvalidConfiguration, 130 Bind, 131 LocalAddress, 132 Accept, 133 ConnectionTaskPanicked, 134 } 135 136 impl MycOperationsErrorKind { 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 Myc operations failure. 151 #[derive(Clone, Copy, PartialEq, Eq)] 152 pub struct MycOperationsError { 153 kind: MycOperationsErrorKind, 154 } 155 156 impl MycOperationsError { 157 const fn new(kind: MycOperationsErrorKind) -> Self { 158 Self { kind } 159 } 160 161 #[must_use] 162 pub const fn kind(self) -> MycOperationsErrorKind { 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 MycOperationsError { 173 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 174 formatter 175 .debug_struct("MycOperationsError") 176 .field("kind", &self.kind) 177 .finish() 178 } 179 } 180 181 impl fmt::Display for MycOperationsError { 182 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 183 formatter.write_str("Myc TCP operations failed") 184 } 185 } 186 187 impl Error for MycOperationsError {} 188 189 fn listener_config( 190 config: &MycConfigDocumentV1, 191 ) -> Result<HostOperationsListenerConfig, MycOperationsError> { 192 let operations = config 193 .normalized() 194 .pointer("/operations") 195 .ok_or_else(invalid_configuration)?; 196 if !boolean(operations, "/enabled")? { 197 return Err(MycOperationsError::new(MycOperationsErrorKind::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 "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, MycOperationsError> { 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, MycOperationsError> { 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, MycOperationsError> { 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, MycOperationsError> { 246 u32::try_from(unsigned(value, pointer)?).map_err(|_| invalid_configuration()) 247 } 248 249 const fn invalid_configuration() -> MycOperationsError { 250 MycOperationsError::new(MycOperationsErrorKind::InvalidConfiguration) 251 } 252 253 const fn map_server_error(error: HostOperationsServerError) -> MycOperationsError { 254 let kind = match error { 255 HostOperationsServerError::Disabled => MycOperationsErrorKind::Disabled, 256 HostOperationsServerError::HeaderLimitBelowParserFloor => { 257 MycOperationsErrorKind::InvalidConfiguration 258 } 259 HostOperationsServerError::Bind { .. } => MycOperationsErrorKind::Bind, 260 HostOperationsServerError::LocalAddress { .. } => MycOperationsErrorKind::LocalAddress, 261 HostOperationsServerError::Accept { .. } => MycOperationsErrorKind::Accept, 262 HostOperationsServerError::ConnectionTaskPanicked => { 263 MycOperationsErrorKind::ConnectionTaskPanicked 264 } 265 }; 266 MycOperationsError::new(kind) 267 }