ops.rs (18742B)
1 #![forbid(unsafe_code)] 2 3 use axum::{Json, Router, extract::State, response::IntoResponse, routing::get}; 4 use http::{StatusCode, header}; 5 use serde::{Deserialize, Serialize}; 6 use std::sync::{Arc, RwLock}; 7 8 use crate::runtime::TangleRuntimeMetrics; 9 10 const PROMETHEUS_CONTENT_TYPE: &str = "text/plain; version=0.0.4; charset=utf-8"; 11 12 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 13 pub enum BaseRelayReadinessCheckStatus { 14 Ready, 15 NotReady, 16 } 17 18 impl BaseRelayReadinessCheckStatus { 19 pub fn as_str(self) -> &'static str { 20 match self { 21 Self::Ready => "ready", 22 Self::NotReady => "not_ready", 23 } 24 } 25 26 pub fn is_ready(self) -> bool { 27 self == Self::Ready 28 } 29 } 30 31 #[derive(Debug, Clone, PartialEq, Eq)] 32 pub struct BaseRelayReadinessState { 33 config: BaseRelayReadinessCheckStatus, 34 server_bind: BaseRelayReadinessCheckStatus, 35 relay_identity: BaseRelayReadinessCheckStatus, 36 pocket_storage: BaseRelayReadinessCheckStatus, 37 group_projection: BaseRelayReadinessCheckStatus, 38 group_outbox_replay: BaseRelayReadinessCheckStatus, 39 event_bus: BaseRelayReadinessCheckStatus, 40 } 41 42 impl BaseRelayReadinessState { 43 pub fn new( 44 config: BaseRelayReadinessCheckStatus, 45 server_bind: BaseRelayReadinessCheckStatus, 46 relay_identity: BaseRelayReadinessCheckStatus, 47 pocket_storage: BaseRelayReadinessCheckStatus, 48 group_projection: BaseRelayReadinessCheckStatus, 49 group_outbox_replay: BaseRelayReadinessCheckStatus, 50 event_bus: BaseRelayReadinessCheckStatus, 51 ) -> Self { 52 Self { 53 config, 54 server_bind, 55 relay_identity, 56 pocket_storage, 57 group_projection, 58 group_outbox_replay, 59 event_bus, 60 } 61 } 62 63 pub fn ready() -> Self { 64 Self::new( 65 BaseRelayReadinessCheckStatus::Ready, 66 BaseRelayReadinessCheckStatus::Ready, 67 BaseRelayReadinessCheckStatus::Ready, 68 BaseRelayReadinessCheckStatus::Ready, 69 BaseRelayReadinessCheckStatus::Ready, 70 BaseRelayReadinessCheckStatus::Ready, 71 BaseRelayReadinessCheckStatus::Ready, 72 ) 73 } 74 75 pub fn runtime_ready_before_bind() -> Self { 76 Self::new( 77 BaseRelayReadinessCheckStatus::Ready, 78 BaseRelayReadinessCheckStatus::NotReady, 79 BaseRelayReadinessCheckStatus::Ready, 80 BaseRelayReadinessCheckStatus::Ready, 81 BaseRelayReadinessCheckStatus::Ready, 82 BaseRelayReadinessCheckStatus::Ready, 83 BaseRelayReadinessCheckStatus::Ready, 84 ) 85 } 86 87 pub fn with_server_bind(mut self, server_bind: BaseRelayReadinessCheckStatus) -> Self { 88 self.server_bind = server_bind; 89 self 90 } 91 92 pub fn is_ready(&self) -> bool { 93 [ 94 self.config, 95 self.server_bind, 96 self.relay_identity, 97 self.pocket_storage, 98 self.group_projection, 99 self.group_outbox_replay, 100 self.event_bus, 101 ] 102 .into_iter() 103 .all(BaseRelayReadinessCheckStatus::is_ready) 104 } 105 106 pub fn response(&self) -> BaseRelayReadinessDocument { 107 BaseRelayReadinessDocument { 108 status: if self.is_ready() { 109 "ready".to_owned() 110 } else { 111 "not_ready".to_owned() 112 }, 113 checks: BaseRelayReadinessChecksDocument { 114 config: self.config.as_str().to_owned(), 115 server_bind: self.server_bind.as_str().to_owned(), 116 relay_identity: self.relay_identity.as_str().to_owned(), 117 pocket_storage: self.pocket_storage.as_str().to_owned(), 118 group_projection: self.group_projection.as_str().to_owned(), 119 group_outbox_replay: self.group_outbox_replay.as_str().to_owned(), 120 event_bus: self.event_bus.as_str().to_owned(), 121 }, 122 } 123 } 124 } 125 126 #[derive(Debug, Clone)] 127 pub struct BaseRelayReadinessHandle { 128 inner: Arc<RwLock<BaseRelayReadinessState>>, 129 } 130 131 impl BaseRelayReadinessHandle { 132 pub fn new(state: BaseRelayReadinessState) -> Self { 133 Self { 134 inner: Arc::new(RwLock::new(state)), 135 } 136 } 137 138 pub fn snapshot(&self) -> BaseRelayReadinessState { 139 match self.inner.read() { 140 Ok(state) => state.clone(), 141 Err(poisoned) => poisoned.into_inner().clone(), 142 } 143 } 144 145 pub fn set_server_bind(&self, status: BaseRelayReadinessCheckStatus) { 146 let mut state = match self.inner.write() { 147 Ok(state) => state, 148 Err(poisoned) => poisoned.into_inner(), 149 }; 150 state.server_bind = status; 151 } 152 } 153 154 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] 155 pub struct BaseRelayHealthDocument { 156 pub status: String, 157 } 158 159 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] 160 pub struct BaseRelayReadinessDocument { 161 pub status: String, 162 pub checks: BaseRelayReadinessChecksDocument, 163 } 164 165 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] 166 pub struct BaseRelayReadinessChecksDocument { 167 pub config: String, 168 pub server_bind: String, 169 pub relay_identity: String, 170 pub pocket_storage: String, 171 pub group_projection: String, 172 pub group_outbox_replay: String, 173 pub event_bus: String, 174 } 175 176 #[derive(Debug, Clone)] 177 struct BaseRelayOpsState { 178 readiness: BaseRelayReadinessHandle, 179 metrics: TangleRuntimeMetrics, 180 } 181 182 pub fn base_relay_ops_router( 183 readiness: BaseRelayReadinessHandle, 184 metrics: TangleRuntimeMetrics, 185 ) -> Router { 186 Router::new() 187 .route("/healthz", get(base_relay_healthz)) 188 .route("/readyz", get(base_relay_readyz)) 189 .route("/metricsz", get(base_relay_metricsz)) 190 .with_state(BaseRelayOpsState { readiness, metrics }) 191 } 192 193 async fn base_relay_healthz() -> Json<BaseRelayHealthDocument> { 194 Json(BaseRelayHealthDocument { 195 status: "ok".to_owned(), 196 }) 197 } 198 199 async fn base_relay_readyz( 200 State(state): State<BaseRelayOpsState>, 201 ) -> (StatusCode, Json<BaseRelayReadinessDocument>) { 202 let readiness = state.readiness.snapshot(); 203 let status = if readiness.is_ready() { 204 StatusCode::OK 205 } else { 206 StatusCode::SERVICE_UNAVAILABLE 207 }; 208 (status, Json(readiness.response())) 209 } 210 211 async fn base_relay_metricsz(State(state): State<BaseRelayOpsState>) -> impl IntoResponse { 212 let readiness = state.readiness.snapshot(); 213 ( 214 [(header::CONTENT_TYPE, PROMETHEUS_CONTENT_TYPE)], 215 state 216 .metrics 217 .snapshot_with_readiness(readiness.is_ready()) 218 .prometheus_text(), 219 ) 220 } 221 222 #[cfg(test)] 223 mod tests { 224 use super::{ 225 BaseRelayReadinessCheckStatus, BaseRelayReadinessHandle, BaseRelayReadinessState, 226 base_relay_ops_router, 227 }; 228 use crate::relay::core::BaseRelayQueryMetrics; 229 use crate::runtime::TangleRuntimeMetrics; 230 use axum::body::to_bytes; 231 use http::{Request, StatusCode, header}; 232 use std::collections::{BTreeMap, BTreeSet}; 233 use tower::ServiceExt; 234 235 #[tokio::test] 236 async fn base_relay_ops_router_reports_health_and_readiness() { 237 let metrics = TangleRuntimeMetrics::new(); 238 let readiness = BaseRelayReadinessHandle::new(BaseRelayReadinessState::ready()); 239 let health = base_relay_ops_router(readiness.clone(), metrics.clone()) 240 .oneshot( 241 Request::builder() 242 .uri("/healthz") 243 .body(axum::body::Body::empty()) 244 .expect("request"), 245 ) 246 .await 247 .expect("health"); 248 249 assert_eq!(health.status(), StatusCode::OK); 250 let health_body = to_bytes(health.into_body(), usize::MAX) 251 .await 252 .expect("body"); 253 let health_value = serde_json::from_slice::<serde_json::Value>(&health_body).expect("json"); 254 assert_eq!(health_value["status"], "ok"); 255 256 metrics.record_session_opened(); 257 metrics.record_client_message(crate::runtime::TangleClientMessageMetricKind::Req); 258 metrics.record_subscription_opened(); 259 metrics.record_auth_success(); 260 metrics.record_auth_failure(); 261 metrics.record_event_admission(); 262 metrics.record_event_rejection(); 263 metrics.record_group_read_denial(); 264 metrics.record_group_write_denial(); 265 metrics.record_event_bus_receivers(2); 266 metrics.record_event_bus_publish(2); 267 metrics.record_event_bus_lagged(3); 268 metrics.record_outbound_queue_full_close(); 269 metrics.record_outbox_pending_events(5); 270 metrics.record_outbox_replayed_event(); 271 metrics.record_disk_used_bytes(89); 272 metrics.record_event_admission_latency(11); 273 metrics.record_query_latency(17); 274 metrics.record_query_metrics(BaseRelayQueryMetrics::new(5, 3, 2)); 275 metrics.record_count_refusal(); 276 metrics.record_broad_query_rejection(); 277 let ready = base_relay_ops_router(readiness.clone(), metrics.clone()) 278 .oneshot( 279 Request::builder() 280 .uri("/readyz") 281 .body(axum::body::Body::empty()) 282 .expect("request"), 283 ) 284 .await 285 .expect("ready"); 286 287 assert_eq!(ready.status(), StatusCode::OK); 288 let ready_body = to_bytes(ready.into_body(), usize::MAX).await.expect("body"); 289 let ready_value = serde_json::from_slice::<serde_json::Value>(&ready_body).expect("json"); 290 assert_eq!(ready_value["status"], "ready"); 291 assert_eq!(ready_value["checks"]["server_bind"], "ready"); 292 assert_eq!(ready_value["checks"]["group_outbox_replay"], "ready"); 293 assert_eq!(ready_value["checks"]["event_bus"], "ready"); 294 let metrics_response = base_relay_ops_router(readiness, metrics) 295 .oneshot( 296 Request::builder() 297 .uri("/metricsz") 298 .body(axum::body::Body::empty()) 299 .expect("request"), 300 ) 301 .await 302 .expect("metrics"); 303 304 assert_eq!(metrics_response.status(), StatusCode::OK); 305 assert_eq!( 306 metrics_response.headers()[header::CONTENT_TYPE], 307 super::PROMETHEUS_CONTENT_TYPE 308 ); 309 let metrics_body = to_bytes(metrics_response.into_body(), usize::MAX) 310 .await 311 .expect("body"); 312 let metrics_text = String::from_utf8(metrics_body.to_vec()).expect("utf8"); 313 let mut metric_types = BTreeMap::new(); 314 let mut metric_values = BTreeMap::new(); 315 for line in metrics_text.lines() { 316 if let Some(type_line) = line.strip_prefix("# TYPE ") { 317 let (name, metric_type) = type_line.split_once(' ').expect("metric type"); 318 assert!(metric_types.insert(name, metric_type).is_none()); 319 } else if !line.is_empty() { 320 let (name, value) = line.split_once(' ').expect("metric sample"); 321 assert!(metric_values.insert(name, value).is_none()); 322 } 323 } 324 let keys = metric_values.keys().copied().collect::<BTreeSet<_>>(); 325 assert_eq!( 326 keys, 327 [ 328 "tangle_auth_failure_total", 329 "tangle_auth_messages_total", 330 "tangle_auth_success_total", 331 "tangle_client_messages_total", 332 "tangle_close_messages_total", 333 "tangle_count_messages_total", 334 "tangle_count_refusals_total", 335 "tangle_disk_used_bytes", 336 "tangle_event_admission_latency_count", 337 "tangle_event_admission_latency_total_micros", 338 "tangle_event_admitted_total", 339 "tangle_event_bus_lagged_offsets_total", 340 "tangle_event_bus_lagged_receivers_total", 341 "tangle_event_bus_published_offsets_total", 342 "tangle_event_bus_receivers_current", 343 "tangle_event_messages_total", 344 "tangle_event_rejected_total", 345 "tangle_group_read_denied_total", 346 "tangle_group_write_denied_total", 347 "tangle_outbound_queue_full_closes_total", 348 "tangle_outbox_pending_events", 349 "tangle_outbox_replayed_events_total", 350 "tangle_broad_query_rejections_total", 351 "tangle_query_candidates_scanned_total", 352 "tangle_query_latency_count", 353 "tangle_query_latency_total_micros", 354 "tangle_query_redacted_events_total", 355 "tangle_query_returned_events_total", 356 "tangle_rate_limit_rejections_total", 357 "tangle_readiness_ready", 358 "tangle_req_messages_total", 359 "tangle_runtime_uptime_seconds", 360 "tangle_stored_event_offsets_total", 361 "tangle_subscriptions_current", 362 "tangle_subscriptions_closed_total", 363 "tangle_subscriptions_opened_total", 364 "tangle_ws_connections_current", 365 "tangle_ws_connections_total", 366 ] 367 .into_iter() 368 .into_iter() 369 .collect::<BTreeSet<_>>() 370 ); 371 assert_eq!(metric_types.keys().copied().collect::<BTreeSet<_>>(), keys); 372 assert_eq!(metric_types["tangle_readiness_ready"], "gauge"); 373 assert_eq!(metric_types["tangle_ws_connections_total"], "counter"); 374 assert_eq!(metric_values["tangle_readiness_ready"], "1"); 375 assert_eq!(metric_values["tangle_ws_connections_current"], "1"); 376 assert_eq!(metric_values["tangle_ws_connections_total"], "1"); 377 assert_eq!(metric_values["tangle_client_messages_total"], "1"); 378 assert_eq!(metric_values["tangle_req_messages_total"], "1"); 379 assert_eq!(metric_values["tangle_subscriptions_current"], "1"); 380 assert_eq!(metric_values["tangle_subscriptions_opened_total"], "1"); 381 assert_eq!(metric_values["tangle_auth_success_total"], "1"); 382 assert_eq!(metric_values["tangle_auth_failure_total"], "1"); 383 assert_eq!(metric_values["tangle_event_admitted_total"], "1"); 384 assert_eq!(metric_values["tangle_event_rejected_total"], "1"); 385 assert_eq!(metric_values["tangle_group_read_denied_total"], "1"); 386 assert_eq!(metric_values["tangle_group_write_denied_total"], "1"); 387 assert_eq!(metric_values["tangle_event_bus_receivers_current"], "2"); 388 assert_eq!( 389 metric_values["tangle_event_bus_published_offsets_total"], 390 "1" 391 ); 392 assert_eq!( 393 metric_values["tangle_event_bus_lagged_receivers_total"], 394 "1" 395 ); 396 assert_eq!(metric_values["tangle_event_bus_lagged_offsets_total"], "3"); 397 assert_eq!( 398 metric_values["tangle_outbound_queue_full_closes_total"], 399 "1" 400 ); 401 assert_eq!(metric_values["tangle_outbox_pending_events"], "5"); 402 assert_eq!(metric_values["tangle_outbox_replayed_events_total"], "1"); 403 assert_eq!(metric_values["tangle_disk_used_bytes"], "89"); 404 assert_eq!( 405 metric_values["tangle_event_admission_latency_total_micros"], 406 "11" 407 ); 408 assert_eq!(metric_values["tangle_event_admission_latency_count"], "1"); 409 assert_eq!(metric_values["tangle_query_latency_total_micros"], "17"); 410 assert_eq!(metric_values["tangle_query_latency_count"], "1"); 411 assert_eq!(metric_values["tangle_query_candidates_scanned_total"], "5"); 412 assert_eq!(metric_values["tangle_query_returned_events_total"], "3"); 413 assert_eq!(metric_values["tangle_query_redacted_events_total"], "2"); 414 assert_eq!(metric_values["tangle_count_refusals_total"], "1"); 415 assert_eq!(metric_values["tangle_broad_query_rejections_total"], "1"); 416 assert!(!metrics_text.contains("relay_secret")); 417 assert!(!metrics_text.contains("invite")); 418 419 let not_ready = BaseRelayReadinessState::new( 420 BaseRelayReadinessCheckStatus::Ready, 421 BaseRelayReadinessCheckStatus::Ready, 422 BaseRelayReadinessCheckStatus::Ready, 423 BaseRelayReadinessCheckStatus::Ready, 424 BaseRelayReadinessCheckStatus::NotReady, 425 BaseRelayReadinessCheckStatus::Ready, 426 BaseRelayReadinessCheckStatus::Ready, 427 ); 428 let rejected = base_relay_ops_router( 429 BaseRelayReadinessHandle::new(not_ready), 430 TangleRuntimeMetrics::new(), 431 ) 432 .oneshot( 433 Request::builder() 434 .uri("/readyz") 435 .body(axum::body::Body::empty()) 436 .expect("request"), 437 ) 438 .await 439 .expect("not ready"); 440 441 assert_eq!(rejected.status(), StatusCode::SERVICE_UNAVAILABLE); 442 let rejected_body = to_bytes(rejected.into_body(), usize::MAX) 443 .await 444 .expect("body"); 445 let rejected_value = 446 serde_json::from_slice::<serde_json::Value>(&rejected_body).expect("json"); 447 assert_eq!(rejected_value["status"], "not_ready"); 448 assert_eq!(rejected_value["checks"]["server_bind"], "ready"); 449 assert_eq!(rejected_value["checks"]["group_projection"], "not_ready"); 450 } 451 452 #[tokio::test] 453 async fn base_relay_ops_router_reports_live_readiness_state() { 454 let readiness = 455 BaseRelayReadinessHandle::new(BaseRelayReadinessState::runtime_ready_before_bind()); 456 let router = base_relay_ops_router(readiness.clone(), TangleRuntimeMetrics::new()); 457 let not_ready = router 458 .clone() 459 .oneshot( 460 Request::builder() 461 .uri("/readyz") 462 .body(axum::body::Body::empty()) 463 .expect("request"), 464 ) 465 .await 466 .expect("not ready"); 467 468 assert_eq!(not_ready.status(), StatusCode::SERVICE_UNAVAILABLE); 469 readiness.set_server_bind(BaseRelayReadinessCheckStatus::Ready); 470 let ready = router 471 .oneshot( 472 Request::builder() 473 .uri("/readyz") 474 .body(axum::body::Body::empty()) 475 .expect("request"), 476 ) 477 .await 478 .expect("ready"); 479 480 assert_eq!(ready.status(), StatusCode::OK); 481 let body = to_bytes(ready.into_body(), usize::MAX).await.expect("body"); 482 let value = serde_json::from_slice::<serde_json::Value>(&body).expect("json"); 483 assert_eq!(value["status"], "ready"); 484 assert_eq!(value["checks"]["event_bus"], "ready"); 485 } 486 }