rhi

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

services_hardening_operations.rs (9572B)


      1 #![forbid(unsafe_code)]
      2 
      3 use std::error::Error;
      4 use std::net::{Ipv4Addr, SocketAddrV4, TcpListener};
      5 
      6 use rhi::{
      7     InstanceId, RHI_LIVEZ_PATH, RHI_METRICS_PATH, RHI_OPERATIONS_CONTRACT_VERSION, RHI_READYZ_PATH,
      8     RhiConfigProfile, RhiEvidenceTransportStatusV1, RhiIdentityHealthV1, RhiIntegrityStateV1,
      9     RhiOperationsCancellationToken, RhiOperationsErrorKind, RhiOperationsServer,
     10     RhiPersistenceHealthV1, RhiPersistenceStatusV1, RhiPresenceStatusV1, RhiProviderStatusV1,
     11     RhiPublicationStatusV1, RhiReconciliationStatusV1, RhiServicePhase, RhiStatusBuildInfoV1,
     12     RhiStatusBuildMode, RhiStatusCommonV1, RhiStatusConfigurationIdentityV1,
     13     RhiStatusConfigurationSource, RhiStatusObservationV1, RhiStatusReasonCodes,
     14     RhiTransportHealthV1, parse_rhi_config_v1, rhi_status_cache,
     15 };
     16 use tokio::io::{AsyncReadExt, AsyncWriteExt};
     17 use tokio::net::TcpStream;
     18 
     19 const CONTRACT: &str = include_str!("../contracts/services_hardening/tcp_operations.v1.json");
     20 const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml");
     21 const SERVICE_REVISION: &str = "0123456789abcdef0123456789abcdef01234567";
     22 const LIB_REVISION: &str = "89abcdef0123456789abcdef0123456789abcdef";
     23 
     24 fn build_info() -> RhiStatusBuildInfoV1 {
     25     RhiStatusBuildInfoV1::new(
     26         RhiStatusBuildMode::Release,
     27         Some("0.1.0"),
     28         Some(SERVICE_REVISION),
     29         Some(LIB_REVISION),
     30         Some("1.97.1"),
     31         Some("x86_64-unknown-linux-gnu"),
     32         Some("service-host"),
     33     )
     34     .expect("build info")
     35 }
     36 
     37 fn identity(configured: bool, available: bool) -> RhiIdentityHealthV1 {
     38     RhiIdentityHealthV1::new(configured, available, RhiStatusReasonCodes::empty())
     39         .expect("identity health")
     40 }
     41 
     42 fn observation(phase: RhiServicePhase, ready: bool) -> RhiStatusObservationV1 {
     43     let configuration = RhiStatusConfigurationIdentityV1::new(
     44         "a".repeat(64),
     45         RhiStatusConfigurationSource::ExplicitConfig,
     46     )
     47     .expect("configuration");
     48     let persistence = RhiPersistenceStatusV1::new(
     49         RhiPersistenceHealthV1::Ready,
     50         10,
     51         42,
     52         RhiIntegrityStateV1::Verified,
     53         RhiStatusReasonCodes::empty(),
     54     )
     55     .expect("persistence");
     56     let provider = RhiProviderStatusV1::new(identity(true, true), RhiStatusReasonCodes::empty())
     57         .expect("provider");
     58     let transport = RhiEvidenceTransportStatusV1::new(
     59         RhiTransportHealthV1::Ready,
     60         true,
     61         true,
     62         2,
     63         2,
     64         RhiStatusReasonCodes::empty(),
     65     )
     66     .expect("transport");
     67     RhiStatusObservationV1::new(
     68         RhiStatusCommonV1::new(
     69             phase,
     70             ready,
     71             RhiStatusReasonCodes::empty(),
     72             1,
     73             build_info(),
     74             configuration,
     75             persistence,
     76         )
     77         .expect("common status"),
     78         provider,
     79         transport,
     80         RhiReconciliationStatusV1::default(),
     81         RhiPublicationStatusV1::new(0, 0, None),
     82         RhiPresenceStatusV1::default(),
     83     )
     84 }
     85 
     86 fn available_port() -> u16 {
     87     let listener =
     88         TcpListener::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)).expect("ephemeral listener");
     89     listener.local_addr().expect("local address").port()
     90 }
     91 
     92 fn enabled_config(port: u16) -> String {
     93     CONFIG.replacen(
     94         "[operations]\nenabled = false",
     95         &format!(
     96             "[operations]\nenabled = true\nlisten = \"127.0.0.1:{port}\"\nbind_policy = \"loopback_only\"\n\n[operations.limits]\nheader_count = 16\nheader_bytes = 8192\nresponse_body_utf8_bytes = 4096\nconcurrent_connections = 4\nrequest_deadline_ms = 500\nidle_timeout_ms = 500"
     97         ),
     98         1,
     99     )
    100 }
    101 
    102 async fn raw_request(address: std::net::SocketAddr, request: &[u8]) -> Vec<u8> {
    103     let mut stream = TcpStream::connect(address).await.expect("connect");
    104     stream.write_all(request).await.expect("request write");
    105     let mut response = Vec::new();
    106     stream
    107         .read_to_end(&mut response)
    108         .await
    109         .expect("response read");
    110     response
    111 }
    112 
    113 fn response_text(response: &[u8]) -> &str {
    114     std::str::from_utf8(response).expect("response UTF-8")
    115 }
    116 
    117 #[test]
    118 fn machine_contract_freezes_exact_routes_metrics_and_non_authority() {
    119     let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract");
    120     assert_eq!(contract["schema"], "radroots.rhi.tcp-operations.v1");
    121     assert_eq!(
    122         contract["contract_version"],
    123         RHI_OPERATIONS_CONTRACT_VERSION
    124     );
    125     assert_eq!(contract["step"], 212);
    126     assert_eq!(
    127         contract["routes"]
    128             .as_array()
    129             .expect("routes")
    130             .iter()
    131             .map(|route| route["path"].as_str().expect("route path"))
    132             .collect::<Vec<_>>(),
    133         [RHI_LIVEZ_PATH, RHI_READYZ_PATH, RHI_METRICS_PATH]
    134     );
    135     assert_eq!(contract["route_registration_extension"], false);
    136     assert_eq!(contract["metrics"]["high_cardinality_labels"], false);
    137     assert_eq!(contract["metrics"]["arbitrary_labels"], false);
    138 }
    139 
    140 #[tokio::test]
    141 async fn exact_tcp_routes_use_only_latest_cached_lifecycle_and_metrics() {
    142     let config = parse_rhi_config_v1(
    143         enabled_config(available_port()).as_bytes(),
    144         RhiConfigProfile::Production,
    145     )
    146     .expect("enabled config");
    147     let (mut publisher, reader) = rhi_status_cache(
    148         InstanceId::new("primary").expect("instance"),
    149         observation(RhiServicePhase::Ready, true),
    150     )
    151     .expect("status cache");
    152     let detail_pointer = reader.snapshot().detailed_status_json().as_ptr();
    153     let bound = RhiOperationsServer::new(&config, &reader)
    154         .expect("operations server")
    155         .bind()
    156         .await
    157         .expect("bind");
    158     let address = bound.local_address();
    159     let cancellation = RhiOperationsCancellationToken::new();
    160     let task = tokio::spawn(bound.serve(cancellation.clone()));
    161 
    162     let live = raw_request(address, b"GET /livez HTTP/1.1\r\nhost: localhost\r\n\r\n").await;
    163     let ready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await;
    164     let metrics = raw_request(address, b"GET /metrics HTTP/1.1\r\nhost: localhost\r\n\r\n").await;
    165     assert!(response_text(&live).starts_with("HTTP/1.1 200 OK\r\n"));
    166     assert!(response_text(&live).ends_with("live\n"));
    167     assert!(response_text(&ready).ends_with("ready\n"));
    168     assert!(response_text(&metrics).contains("# TYPE radroots_rhi_service_phase gauge\n"));
    169     assert!(response_text(&metrics).contains("radroots_rhi_service_phase{phase=\"ready\"} 1\n"));
    170     assert!(response_text(&metrics).contains("radroots_rhi_service_ready 1\n"));
    171 
    172     for request in [
    173         &b"GET /status HTTP/1.1\r\nhost: localhost\r\n\r\n"[..],
    174         &b"GET /readyz?probe=1 HTTP/1.1\r\nhost: localhost\r\n\r\n"[..],
    175         &b"POST /metrics HTTP/1.1\r\nhost: localhost\r\ncontent-length: 0\r\n\r\n"[..],
    176         &b"GET /v1/status HTTP/1.1\r\nhost: localhost\r\n\r\n"[..],
    177     ] {
    178         let rejected = raw_request(address, request).await;
    179         assert!(response_text(&rejected).starts_with("HTTP/1.1 404 Not Found\r\n"));
    180     }
    181 
    182     publisher
    183         .publish(observation(RhiServicePhase::Unready, false))
    184         .expect("unready publication");
    185     let unready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await;
    186     let metrics = raw_request(address, b"GET /metrics HTTP/1.1\r\nhost: localhost\r\n\r\n").await;
    187     assert!(response_text(&unready).starts_with("HTTP/1.1 503 Service Unavailable\r\n"));
    188     assert!(response_text(&unready).ends_with("unready\n"));
    189     assert!(response_text(&metrics).contains("radroots_rhi_service_phase{phase=\"unready\"} 1\n"));
    190     assert!(response_text(&metrics).contains("radroots_rhi_service_ready 0\n"));
    191     assert_ne!(
    192         reader.snapshot().detailed_status_json().as_ptr(),
    193         detail_pointer
    194     );
    195 
    196     cancellation.cancel();
    197     assert_eq!(task.await.expect("serve task"), Ok(()));
    198 }
    199 
    200 #[tokio::test]
    201 async fn disabled_invalid_and_bind_failures_are_typed_source_free_and_redacted() {
    202     let disabled = parse_rhi_config_v1(CONFIG.as_bytes(), RhiConfigProfile::Production)
    203         .expect("disabled config");
    204     let (_, reader) = rhi_status_cache(
    205         InstanceId::new("primary").expect("instance"),
    206         observation(RhiServicePhase::Ready, true),
    207     )
    208     .expect("status cache");
    209     let disabled_error =
    210         RhiOperationsServer::new(&disabled, &reader).expect_err("disabled operations");
    211     assert_eq!(disabled_error.kind(), RhiOperationsErrorKind::Disabled);
    212     assert_eq!(disabled_error.code(), "operations_disabled");
    213     assert!(Error::source(&disabled_error).is_none());
    214 
    215     let below_floor =
    216         enabled_config(available_port()).replace("header_bytes = 8192", "header_bytes = 8191");
    217     assert!(parse_rhi_config_v1(below_floor.as_bytes(), RhiConfigProfile::Production).is_err());
    218 
    219     let occupied =
    220         TcpListener::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)).expect("occupied listener");
    221     let port = occupied.local_addr().expect("occupied address").port();
    222     let config = parse_rhi_config_v1(
    223         enabled_config(port).as_bytes(),
    224         RhiConfigProfile::Production,
    225     )
    226     .expect("enabled config");
    227     let server = RhiOperationsServer::new(&config, &reader).expect("server");
    228     assert!(!format!("{server:?}").contains(&port.to_string()));
    229     let bind_error = server.bind().await.expect_err("occupied bind");
    230     assert_eq!(bind_error.kind(), RhiOperationsErrorKind::Bind);
    231     assert!(Error::source(&bind_error).is_none());
    232     assert!(!format!("{bind_error:?} {bind_error}").contains(&port.to_string()));
    233 }