tangle


git clone https://radroots.dev/git/tangle.git
Log | Files | Refs | README | LICENSE

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 }