admin_v1.rs (85615B)
1 //! Exact RHI v1 Unix-admin route and model boundary. 2 3 #[cfg(any(target_os = "linux", target_os = "macos"))] 4 use core::time::Duration; 5 use core::{fmt, future::Future, pin::Pin}; 6 use std::{collections::BTreeSet, error::Error, sync::OnceLock}; 7 8 #[cfg(any(target_os = "linux", target_os = "macos"))] 9 use std::sync::Arc; 10 11 use radroots_service_host::{AdminCorrelationId, AdminOperationId, CancellationToken}; 12 #[cfg(any(target_os = "linux", target_os = "macos"))] 13 use radroots_service_host::{ 14 AdminError, AdminErrorCode, AdminErrorMessage, AdminHttpMethod, AdminMutationRequest, 15 AdminRequest, AdminRouteFailure, AdminRouteFailureStatus, AdminRouteOutcome, 16 AdminRouter as SharedAdminRouter, AdminServer as SharedAdminServer, 17 AdminServerError as SharedAdminServerError, AdminTransportLimitValues, AdminTransportLimits, 18 UnixAdminSocketBinding, UnixAdminSocketWriterAuthority, 19 }; 20 use serde::de::{self, DeserializeSeed, MapAccess, SeqAccess, Visitor}; 21 use serde_json::{Map, Value}; 22 23 // Original model admission is never permitted to exceed the complete response 24 // body cap enforced again by the shared host after envelope encoding. 25 const RHI_ADMIN_RESPONSE_BODY_MAX_UTF8_BYTES: usize = 1_048_576; 26 27 const OPERATOR_CONTRACT: &str = 28 include_str!("../contracts/services_hardening/operator_contract.v1.json"); 29 30 /// Closed RHI v1 admin method vocabulary. 31 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 32 pub enum RhiAdminMethod { 33 Get, 34 Post, 35 } 36 37 /// Closed RHI v1 Unix-admin route inventory. 38 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 39 pub enum RhiAdminRoute { 40 Status, 41 EffectiveConfig, 42 IdentityStatus, 43 IdentityPublic, 44 StateStatus, 45 StateBackup, 46 MetricsSnapshot, 47 ReconciliationStatus, 48 ReconciliationJobs, 49 ReconciliationRefresh, 50 Sources, 51 TradeProjection, 52 TradeReportCurrent, 53 TradeReports, 54 PublicationBacklog, 55 PublicationTargets, 56 PublicationRetry, 57 PresenceDesired, 58 PresenceRender, 59 PresenceRefresh, 60 } 61 62 impl RhiAdminRoute { 63 /// Complete final machine-governed route inventory. 64 pub const ALL: [Self; 20] = [ 65 Self::Status, 66 Self::EffectiveConfig, 67 Self::IdentityStatus, 68 Self::IdentityPublic, 69 Self::StateStatus, 70 Self::StateBackup, 71 Self::MetricsSnapshot, 72 Self::ReconciliationStatus, 73 Self::ReconciliationJobs, 74 Self::ReconciliationRefresh, 75 Self::Sources, 76 Self::TradeProjection, 77 Self::TradeReportCurrent, 78 Self::TradeReports, 79 Self::PublicationBacklog, 80 Self::PublicationTargets, 81 Self::PublicationRetry, 82 Self::PresenceDesired, 83 Self::PresenceRender, 84 Self::PresenceRefresh, 85 ]; 86 87 /// Common routes admitted by the Step 206 partial control surface. 88 pub const COMMON: [Self; 7] = [ 89 Self::Status, 90 Self::EffectiveConfig, 91 Self::IdentityStatus, 92 Self::IdentityPublic, 93 Self::StateStatus, 94 Self::StateBackup, 95 Self::MetricsSnapshot, 96 ]; 97 98 /// Domain routes admitted by the Step 207 control surface. 99 pub const DOMAIN: [Self; 13] = [ 100 Self::ReconciliationStatus, 101 Self::ReconciliationJobs, 102 Self::ReconciliationRefresh, 103 Self::Sources, 104 Self::TradeProjection, 105 Self::TradeReportCurrent, 106 Self::TradeReports, 107 Self::PublicationBacklog, 108 Self::PublicationTargets, 109 Self::PublicationRetry, 110 Self::PresenceDesired, 111 Self::PresenceRender, 112 Self::PresenceRefresh, 113 ]; 114 115 /// Routes admitted through Step 208, in final machine-contract order. 116 pub const ACTIVE: [Self; 20] = Self::ALL; 117 118 #[must_use] 119 pub const fn method(self) -> RhiAdminMethod { 120 match self { 121 Self::Status 122 | Self::EffectiveConfig 123 | Self::IdentityStatus 124 | Self::IdentityPublic 125 | Self::StateStatus 126 | Self::MetricsSnapshot 127 | Self::ReconciliationStatus 128 | Self::ReconciliationJobs 129 | Self::Sources 130 | Self::TradeProjection 131 | Self::TradeReportCurrent 132 | Self::TradeReports 133 | Self::PublicationBacklog 134 | Self::PublicationTargets 135 | Self::PresenceDesired => RhiAdminMethod::Get, 136 Self::StateBackup 137 | Self::ReconciliationRefresh 138 | Self::PublicationRetry 139 | Self::PresenceRender 140 | Self::PresenceRefresh => RhiAdminMethod::Post, 141 } 142 } 143 144 #[must_use] 145 pub const fn path(self) -> &'static str { 146 match self { 147 Self::Status => "/v1/status", 148 Self::EffectiveConfig => "/v1/config/effective", 149 Self::IdentityStatus => "/v1/identity/status", 150 Self::IdentityPublic => "/v1/identity/public", 151 Self::StateStatus => "/v1/state/status", 152 Self::StateBackup => "/v1/state/backup", 153 Self::MetricsSnapshot => "/v1/metrics/snapshot", 154 Self::ReconciliationStatus => "/v1/reconciliation/status", 155 Self::ReconciliationJobs => "/v1/reconciliation/jobs", 156 Self::ReconciliationRefresh => "/v1/reconciliation/refresh", 157 Self::Sources => "/v1/sources", 158 Self::TradeProjection => "/v1/trades/{trade_id}/projection", 159 Self::TradeReportCurrent => "/v1/trades/{trade_id}/reports/current", 160 Self::TradeReports => "/v1/trades/{trade_id}/reports", 161 Self::PublicationBacklog => "/v1/publication/backlog", 162 Self::PublicationTargets => "/v1/publication/targets", 163 Self::PublicationRetry => "/v1/publication/retry", 164 Self::PresenceDesired => "/v1/presence/desired", 165 Self::PresenceRender => "/v1/presence/render", 166 Self::PresenceRefresh => "/v1/presence/refresh", 167 } 168 } 169 170 #[must_use] 171 pub const fn operation_id(self) -> &'static str { 172 match self { 173 Self::Status => "radroots.rhi.status.get.v1", 174 Self::EffectiveConfig => "radroots.rhi.config.effective.get.v1", 175 Self::IdentityStatus => "radroots.rhi.identity.status.get.v1", 176 Self::IdentityPublic => "radroots.rhi.identity.public.get.v1", 177 Self::StateStatus => "radroots.rhi.state.status.get.v1", 178 Self::StateBackup => "radroots.rhi.state.backup.create.v1", 179 Self::MetricsSnapshot => "radroots.rhi.metrics.snapshot.get.v1", 180 Self::ReconciliationStatus => "radroots.rhi.reconciliation.status.get.v1", 181 Self::ReconciliationJobs => "radroots.rhi.reconciliation.jobs.list.v1", 182 Self::ReconciliationRefresh => "radroots.rhi.reconciliation.refresh.v1", 183 Self::Sources => "radroots.rhi.sources.list.v1", 184 Self::TradeProjection => "radroots.rhi.trade.projection.get.v1", 185 Self::TradeReportCurrent => "radroots.rhi.trade.report.current.get.v1", 186 Self::TradeReports => "radroots.rhi.trade.reports.list.v1", 187 Self::PublicationBacklog => "radroots.rhi.publication.backlog.list.v1", 188 Self::PublicationTargets => "radroots.rhi.publication.targets.list.v1", 189 Self::PublicationRetry => "radroots.rhi.publication.retry.v1", 190 Self::PresenceDesired => "radroots.rhi.presence.desired.get.v1", 191 Self::PresenceRender => "radroots.rhi.presence.render.v1", 192 Self::PresenceRefresh => "radroots.rhi.presence.refresh.v1", 193 } 194 } 195 196 #[must_use] 197 pub const fn request_model(self) -> &'static str { 198 match self { 199 Self::Status 200 | Self::EffectiveConfig 201 | Self::StateStatus 202 | Self::MetricsSnapshot 203 | Self::ReconciliationStatus 204 | Self::TradeProjection 205 | Self::TradeReportCurrent 206 | Self::PresenceDesired => "empty", 207 Self::IdentityStatus => "identity_status_query_v1", 208 Self::IdentityPublic => "identity_public_query_v1", 209 Self::StateBackup => "state_backup_request_v1", 210 Self::ReconciliationJobs => "reconciliation_jobs_query_v1", 211 Self::ReconciliationRefresh => "reconciliation_refresh_request_v1", 212 Self::Sources => "sources_query_v1", 213 Self::TradeReports => "reports_query_v1", 214 Self::PublicationBacklog => "publication_backlog_query_v1", 215 Self::PublicationTargets => "publication_targets_query_v1", 216 Self::PublicationRetry => "publication_retry_request_v1", 217 Self::PresenceRender => "presence_render_request_v1", 218 Self::PresenceRefresh => "presence_refresh_request_v1", 219 } 220 } 221 222 #[must_use] 223 pub const fn response_model(self) -> &'static str { 224 match self { 225 Self::Status => "service_status_v1", 226 Self::EffectiveConfig => "effective_config_v1", 227 Self::IdentityStatus => "identity_status_v1", 228 Self::IdentityPublic => "identity_public_v1", 229 Self::StateStatus => "state_status_v1", 230 Self::StateBackup => "state_backup_receipt_v1", 231 Self::MetricsSnapshot => "metrics_snapshot_v1", 232 Self::ReconciliationStatus => "reconciliation_status_v1", 233 Self::ReconciliationJobs => "reconciliation_jobs_page_v1", 234 Self::ReconciliationRefresh => "reconciliation_refresh_receipt_v1", 235 Self::Sources => "sources_page_v1", 236 Self::TradeProjection => "trade_projection_v1", 237 Self::TradeReportCurrent => "report_detail_v1", 238 Self::TradeReports => "reports_page_v1", 239 Self::PublicationBacklog => "publication_backlog_page_v1", 240 Self::PublicationTargets => "publication_targets_page_v1", 241 Self::PublicationRetry => "publication_retry_receipt_v1", 242 Self::PresenceDesired => "presence_desired_v1", 243 Self::PresenceRender => "presence_render_receipt_v1", 244 Self::PresenceRefresh => "presence_refresh_receipt_v1", 245 } 246 } 247 248 #[must_use] 249 pub const fn is_mutation(self) -> bool { 250 matches!(self.method(), RhiAdminMethod::Post) 251 } 252 253 #[cfg(any(target_os = "linux", target_os = "macos"))] 254 const fn host_method(self) -> AdminHttpMethod { 255 match self.method() { 256 RhiAdminMethod::Get => AdminHttpMethod::Get, 257 RhiAdminMethod::Post => AdminHttpMethod::Post, 258 } 259 } 260 261 #[cfg(any(target_os = "linux", target_os = "macos"))] 262 const fn parameter_binding(self) -> Option<(&'static str, &'static str)> { 263 match self { 264 Self::TradeProjection | Self::TradeReportCurrent | Self::TradeReports => { 265 Some(("trade_id", "trade_id")) 266 } 267 _ => None, 268 } 269 } 270 } 271 272 /// One route-bound, already validated RHI admin request. 273 /// 274 /// Construction is sealed to an admitted request from the shared server: 275 /// 276 /// ```compile_fail 277 /// use rhi::{RhiAdminRequestDocument, RhiAdminRoute}; 278 /// 279 /// let _ = RhiAdminRequestDocument { 280 /// route: RhiAdminRoute::Status, 281 /// operation_id: None, 282 /// correlation_id: todo!(), 283 /// parameter: None, 284 /// model_bytes: Box::new([]), 285 /// }; 286 /// ``` 287 pub struct RhiAdminRequestDocument { 288 route: RhiAdminRoute, 289 operation_id: Option<AdminOperationId>, 290 correlation_id: AdminCorrelationId, 291 parameter: Option<(&'static str, Box<str>)>, 292 model_bytes: Box<[u8]>, 293 } 294 295 impl RhiAdminRequestDocument { 296 #[must_use] 297 pub const fn route(&self) -> RhiAdminRoute { 298 self.route 299 } 300 301 /// Returns the caller's durable idempotency identity for a mutation. 302 #[must_use] 303 pub fn operation_id(&self) -> Option<&str> { 304 self.operation_id.as_ref().map(AdminOperationId::as_str) 305 } 306 307 #[must_use] 308 pub fn correlation_id(&self) -> &str { 309 self.correlation_id.as_str() 310 } 311 312 /// Returns a validated percent-decoded path parameter when this route has one. 313 #[must_use] 314 pub fn parameter(&self, name: &str) -> Option<&str> { 315 self.parameter 316 .as_ref() 317 .filter(|(parameter_name, _)| *parameter_name == name) 318 .map(|(_, value)| value.as_ref()) 319 } 320 321 /// Returns compact canonical JSON for the route's exact request model. 322 #[must_use] 323 pub fn model_bytes(&self) -> &[u8] { 324 &self.model_bytes 325 } 326 } 327 328 impl fmt::Debug for RhiAdminRequestDocument { 329 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 330 formatter 331 .debug_struct("RhiAdminRequestDocument") 332 .field("route", &self.route) 333 .field("mutation", &self.operation_id.is_some()) 334 .field("has_parameter", &self.parameter.is_some()) 335 .field("model", &"[redacted]") 336 .finish() 337 } 338 } 339 340 /// One exact validated response model for a fixed route. 341 pub struct RhiAdminResponseDocument { 342 route: RhiAdminRoute, 343 canonical_bytes: Box<[u8]>, 344 #[cfg(any(target_os = "linux", target_os = "macos"))] 345 value: Value, 346 } 347 348 impl RhiAdminResponseDocument { 349 /// Admits only compact canonical JSON matching the route's response model. 350 pub fn from_canonical_bytes( 351 route: RhiAdminRoute, 352 bytes: &[u8], 353 ) -> Result<Self, RhiAdminDocumentError> { 354 if bytes.is_empty() { 355 return Err(RhiAdminDocumentError::new( 356 RhiAdminDocumentErrorKind::Malformed, 357 )); 358 } 359 if bytes.len() > RHI_ADMIN_RESPONSE_BODY_MAX_UTF8_BYTES { 360 return Err(RhiAdminDocumentError::new( 361 RhiAdminDocumentErrorKind::TooLarge, 362 )); 363 } 364 let value = strict_json(bytes)?; 365 validate_model(route.response_model(), &value)?; 366 let canonical = serde_json::to_vec(&value) 367 .map_err(|_| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::Malformed))?; 368 if canonical.as_slice() != bytes { 369 return Err(RhiAdminDocumentError::new( 370 RhiAdminDocumentErrorKind::NonCanonical, 371 )); 372 } 373 Ok(Self { 374 route, 375 canonical_bytes: canonical.into_boxed_slice(), 376 #[cfg(any(target_os = "linux", target_os = "macos"))] 377 value, 378 }) 379 } 380 381 #[must_use] 382 pub const fn route(&self) -> RhiAdminRoute { 383 self.route 384 } 385 386 #[must_use] 387 pub fn canonical_bytes(&self) -> &[u8] { 388 &self.canonical_bytes 389 } 390 } 391 392 impl fmt::Debug for RhiAdminResponseDocument { 393 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 394 formatter 395 .debug_struct("RhiAdminResponseDocument") 396 .field("route", &self.route) 397 .field("model", &"[redacted]") 398 .finish() 399 } 400 } 401 402 /// Stable classification for a rejected RHI admin document. 403 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 404 pub enum RhiAdminDocumentErrorKind { 405 TooLarge, 406 Malformed, 407 DuplicateField, 408 NullForbidden, 409 NonCanonical, 410 InvalidModel, 411 } 412 413 /// Source-free and content-free RHI admin document error. 414 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 415 pub struct RhiAdminDocumentError { 416 kind: RhiAdminDocumentErrorKind, 417 } 418 419 impl RhiAdminDocumentError { 420 const fn new(kind: RhiAdminDocumentErrorKind) -> Self { 421 Self { kind } 422 } 423 424 #[must_use] 425 pub const fn kind(self) -> RhiAdminDocumentErrorKind { 426 self.kind 427 } 428 } 429 430 impl fmt::Display for RhiAdminDocumentError { 431 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 432 formatter.write_str("RHI admin document is invalid") 433 } 434 } 435 436 impl Error for RhiAdminDocumentError {} 437 438 /// Stable route-handler failure mapped to a bounded safe admin response. 439 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 440 pub enum RhiAdminHandlerErrorKind { 441 InvalidCursor, 442 OperationIdConflict, 443 NotFound, 444 Conflict, 445 Unavailable, 446 Internal, 447 } 448 449 /// Source-free route-handler error. 450 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 451 pub struct RhiAdminHandlerError { 452 kind: RhiAdminHandlerErrorKind, 453 } 454 455 impl RhiAdminHandlerError { 456 #[must_use] 457 pub const fn new(kind: RhiAdminHandlerErrorKind) -> Self { 458 Self { kind } 459 } 460 461 #[must_use] 462 pub const fn kind(self) -> RhiAdminHandlerErrorKind { 463 self.kind 464 } 465 } 466 467 impl fmt::Display for RhiAdminHandlerError { 468 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 469 formatter.write_str("RHI admin operation failed") 470 } 471 } 472 473 impl Error for RhiAdminHandlerError {} 474 475 /// Boxed route future used by the sealed RHI admin adapter. 476 pub type RhiAdminFuture<'a> = Pin< 477 Box<dyn Future<Output = Result<RhiAdminResponseDocument, RhiAdminHandlerError>> + Send + 'a>, 478 >; 479 480 /// Domain port behind the exact RHI admin transport. 481 /// 482 /// Implementations own authoritative local commit, durable operation-ID 483 /// replay/conflict handling, and any provider or outbox orchestration. Returning 484 /// success means that the operation's contract-defined local effect is already 485 /// committed; relay submission or delivery is not implied. Pagination cursors 486 /// must be authenticated and bound to the same route, filters, and snapshot; 487 /// mismatched or invalid cursors return [`RhiAdminHandlerErrorKind::InvalidCursor`]. 488 /// Reusing an operation ID with identical canonical request bytes returns the 489 /// original committed response; reuse with different bytes returns 490 /// [`RhiAdminHandlerErrorKind::OperationIdConflict`]. 491 pub trait RhiAdminHandler: Send + Sync + 'static { 492 fn handle<'a>(&'a self, request: RhiAdminRequestDocument) -> RhiAdminFuture<'a>; 493 } 494 495 /// Source-free router-construction failure. 496 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 497 pub struct RhiAdminRouterError; 498 499 impl fmt::Display for RhiAdminRouterError { 500 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 501 formatter.write_str("RHI admin router could not be constructed") 502 } 503 } 504 505 impl Error for RhiAdminRouterError {} 506 507 /// Opaque final RHI v1 router capability through Step 209. 508 /// 509 /// The underlying shared-host router remains an implementation detail. The 510 /// later runtime-composition checkpoint consumes this capability without 511 /// exposing raw listener or transport authority. 512 /// 513 /// ```compile_fail 514 /// use rhi::RhiAdminRouter; 515 /// 516 /// let _ = RhiAdminRouter { inner: todo!() }; 517 /// ``` 518 #[cfg(any(target_os = "linux", target_os = "macos"))] 519 pub struct RhiAdminRouter { 520 inner: SharedAdminRouter, 521 } 522 523 #[cfg(any(target_os = "linux", target_os = "macos"))] 524 impl fmt::Debug for RhiAdminRouter { 525 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 526 let Self { inner } = self; 527 let _ = inner; 528 formatter.write_str("RhiAdminRouter") 529 } 530 } 531 532 #[cfg(any(target_os = "linux", target_os = "macos"))] 533 impl RhiAdminRouter { 534 fn into_inner(self) -> SharedAdminRouter { 535 self.inner 536 } 537 } 538 539 /// Cloneable cooperative cancellation for the RHI Unix-admin server. 540 #[derive(Clone, Default)] 541 pub struct RhiAdminCancellationToken { 542 inner: CancellationToken, 543 } 544 545 impl RhiAdminCancellationToken { 546 #[must_use] 547 pub fn new() -> Self { 548 Self::default() 549 } 550 551 /// Requests cancellation. Repeated requests have no additional effect. 552 pub fn cancel(&self) { 553 self.inner.cancel(); 554 } 555 556 #[must_use] 557 pub fn is_cancelled(&self) -> bool { 558 self.inner.is_cancelled() 559 } 560 } 561 562 impl fmt::Debug for RhiAdminCancellationToken { 563 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 564 formatter 565 .debug_struct("RhiAdminCancellationToken") 566 .field("cancelled", &self.is_cancelled()) 567 .finish() 568 } 569 } 570 571 /// Stable source-free RHI Unix-admin server failure classification. 572 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 573 pub enum RhiAdminServerErrorKind { 574 InvalidConfiguration, 575 Router, 576 ServerConfiguration, 577 WriterAuthority, 578 Bind, 579 Listener, 580 Accept, 581 ConnectionTaskPanicked, 582 } 583 584 impl RhiAdminServerErrorKind { 585 #[must_use] 586 pub const fn code(self) -> &'static str { 587 match self { 588 Self::InvalidConfiguration => "admin_configuration_invalid", 589 Self::Router => "admin_router_invalid", 590 Self::ServerConfiguration => "admin_server_configuration_invalid", 591 Self::WriterAuthority => "admin_writer_authority_unavailable", 592 Self::Bind => "admin_bind_failed", 593 Self::Listener => "admin_listener_failed", 594 Self::Accept => "admin_accept_failed", 595 Self::ConnectionTaskPanicked => "admin_connection_task_panicked", 596 } 597 } 598 } 599 600 /// One redacted source-free RHI Unix-admin server failure. 601 #[derive(Clone, Copy, PartialEq, Eq)] 602 pub struct RhiAdminServerError { 603 kind: RhiAdminServerErrorKind, 604 } 605 606 impl RhiAdminServerError { 607 #[cfg(any(target_os = "linux", target_os = "macos"))] 608 const fn new(kind: RhiAdminServerErrorKind) -> Self { 609 Self { kind } 610 } 611 612 #[must_use] 613 pub const fn kind(self) -> RhiAdminServerErrorKind { 614 self.kind 615 } 616 617 #[must_use] 618 pub const fn code(self) -> &'static str { 619 self.kind.code() 620 } 621 } 622 623 impl fmt::Debug for RhiAdminServerError { 624 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 625 formatter 626 .debug_struct("RhiAdminServerError") 627 .field("kind", &self.kind) 628 .finish() 629 } 630 } 631 632 impl fmt::Display for RhiAdminServerError { 633 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 634 formatter.write_str("RHI Unix-admin server failed") 635 } 636 } 637 638 impl Error for RhiAdminServerError {} 639 640 /// Unbound final RHI Unix-admin server through Step 209. 641 /// 642 /// Construction projects only the already-admitted Rhi configuration, seals 643 /// the exact route inventory around the supplied domain handler, and uses the 644 /// shared host's system entropy. The raw shared router and server never cross 645 /// this boundary. 646 #[cfg(any(target_os = "linux", target_os = "macos"))] 647 pub struct RhiAdminServer { 648 inner: SharedAdminServer, 649 } 650 651 #[cfg(any(target_os = "linux", target_os = "macos"))] 652 impl RhiAdminServer { 653 pub fn new<H>( 654 configuration: &crate::RhiConfigDocumentV1, 655 handler: Arc<H>, 656 ) -> Result<Self, RhiAdminServerError> 657 where 658 H: RhiAdminHandler, 659 { 660 let limits = admin_transport_limits(configuration)?; 661 let router = build_rhi_admin_router(handler) 662 .map_err(|_| RhiAdminServerError::new(RhiAdminServerErrorKind::Router))?; 663 let inner = SharedAdminServer::with_system_entropy(router.into_inner(), limits) 664 .map_err(|_| RhiAdminServerError::new(RhiAdminServerErrorKind::ServerConfiguration))?; 665 Ok(Self { inner }) 666 } 667 668 /// Acquires the canonical runtime-directory authority and binds `admin.sock`. 669 /// 670 /// Binding does not spawn a task or begin request admission. Unit 15 owns 671 /// the final supervised server task and its shutdown phase. 672 pub async fn bind( 673 self, 674 runtime: &crate::RhiRuntimeContext, 675 ) -> Result<RhiBoundAdminServer, RhiAdminServerError> { 676 let authority = UnixAdminSocketWriterAuthority::acquire(runtime.context().paths().run()) 677 .map_err(|_| RhiAdminServerError::new(RhiAdminServerErrorKind::WriterAuthority))?; 678 let binding = UnixAdminSocketBinding::bind(authority, runtime.artifacts().admin_socket()) 679 .await 680 .map_err(|_| RhiAdminServerError::new(RhiAdminServerErrorKind::Bind))?; 681 Ok(RhiBoundAdminServer { 682 inner: self.inner, 683 binding, 684 }) 685 } 686 } 687 688 #[cfg(any(target_os = "linux", target_os = "macos"))] 689 impl fmt::Debug for RhiAdminServer { 690 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 691 formatter.write_str("RhiAdminServer([sealed])") 692 } 693 } 694 695 /// Bound final RHI Unix-admin server through Step 209. 696 #[cfg(any(target_os = "linux", target_os = "macos"))] 697 pub struct RhiBoundAdminServer { 698 inner: SharedAdminServer, 699 binding: UnixAdminSocketBinding, 700 } 701 702 #[cfg(any(target_os = "linux", target_os = "macos"))] 703 impl RhiBoundAdminServer { 704 /// Serves until supervisor cancellation and then drains bounded connection work. 705 pub async fn serve( 706 self, 707 cancellation: RhiAdminCancellationToken, 708 ) -> Result<(), RhiAdminServerError> { 709 self.inner 710 .serve(self.binding, cancellation.inner) 711 .await 712 .map_err(map_admin_server_error) 713 } 714 } 715 716 #[cfg(any(target_os = "linux", target_os = "macos"))] 717 impl fmt::Debug for RhiBoundAdminServer { 718 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 719 formatter.write_str("RhiBoundAdminServer([sealed])") 720 } 721 } 722 723 /// Registers the final seven common and thirteen domain routes through Step 209. 724 /// 725 /// Live identity rekey and replace are absent by final offline-only policy. 726 #[cfg(any(target_os = "linux", target_os = "macos"))] 727 pub fn build_rhi_admin_router<H>(handler: Arc<H>) -> Result<RhiAdminRouter, RhiAdminRouterError> 728 where 729 H: RhiAdminHandler, 730 { 731 if !operator_route_inventory_is_exact() { 732 return Err(RhiAdminRouterError); 733 } 734 let mut router = SharedAdminRouter::new(); 735 for route in RhiAdminRoute::ACTIVE { 736 let handler = Arc::clone(&handler); 737 router 738 .route(route.host_method(), route.path(), move |request| { 739 let handler = Arc::clone(&handler); 740 async move { dispatch_route(route, handler, request).await } 741 }) 742 .map_err(|_| RhiAdminRouterError)?; 743 } 744 Ok(RhiAdminRouter { inner: router }) 745 } 746 747 #[cfg(any(target_os = "linux", target_os = "macos"))] 748 pub(crate) fn admin_transport_limits( 749 configuration: &crate::RhiConfigDocumentV1, 750 ) -> Result<AdminTransportLimits, RhiAdminServerError> { 751 let admin = configuration 752 .normalized() 753 .pointer("/resource_limits/admin") 754 .ok_or_else(invalid_admin_configuration)?; 755 let values = AdminTransportLimitValues { 756 header_count: admin_u32(admin, "/header_count")?, 757 header_bytes: admin_u32(admin, "/header_bytes")?, 758 request_body_utf8_bytes: admin_u32(admin, "/request_body_utf8_bytes")?, 759 response_body_utf8_bytes: admin_u32(admin, "/response_body_utf8_bytes")?, 760 concurrent_connections: admin_u32(admin, "/concurrent_connections")?, 761 request_deadline: Duration::from_millis(admin_u64(admin, "/request_deadline_ms")?), 762 idle_timeout: Duration::from_millis(admin_u64(admin, "/idle_timeout_ms")?), 763 query_items: admin_u32(admin, "/query_items")?, 764 }; 765 AdminTransportLimits::new(values).map_err(|_| invalid_admin_configuration()) 766 } 767 768 #[cfg(any(target_os = "linux", target_os = "macos"))] 769 pub(crate) fn admit_admin_response_value( 770 route: RhiAdminRoute, 771 value: &Value, 772 ) -> Result<Box<[u8]>, RhiAdminDocumentError> { 773 let bytes = serde_json::to_vec(value) 774 .map_err(|_| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::Malformed))?; 775 RhiAdminResponseDocument::from_canonical_bytes(route, &bytes) 776 .map(|document| document.canonical_bytes) 777 } 778 779 #[cfg(any(target_os = "linux", target_os = "macos"))] 780 fn admin_u64(value: &Value, pointer: &str) -> Result<u64, RhiAdminServerError> { 781 value 782 .pointer(pointer) 783 .and_then(Value::as_u64) 784 .ok_or_else(invalid_admin_configuration) 785 } 786 787 #[cfg(any(target_os = "linux", target_os = "macos"))] 788 fn admin_u32(value: &Value, pointer: &str) -> Result<u32, RhiAdminServerError> { 789 u32::try_from(admin_u64(value, pointer)?).map_err(|_| invalid_admin_configuration()) 790 } 791 792 #[cfg(any(target_os = "linux", target_os = "macos"))] 793 const fn invalid_admin_configuration() -> RhiAdminServerError { 794 RhiAdminServerError::new(RhiAdminServerErrorKind::InvalidConfiguration) 795 } 796 797 #[cfg(any(target_os = "linux", target_os = "macos"))] 798 const fn map_admin_server_error(error: SharedAdminServerError) -> RhiAdminServerError { 799 let kind = match error { 800 SharedAdminServerError::ListenerClone { .. } 801 | SharedAdminServerError::ListenerRegistration { .. } => RhiAdminServerErrorKind::Listener, 802 SharedAdminServerError::Accept { .. } => RhiAdminServerErrorKind::Accept, 803 SharedAdminServerError::ConnectionTaskPanicked => { 804 RhiAdminServerErrorKind::ConnectionTaskPanicked 805 } 806 }; 807 RhiAdminServerError::new(kind) 808 } 809 810 #[cfg(any(target_os = "linux", target_os = "macos"))] 811 async fn dispatch_route<H>( 812 route: RhiAdminRoute, 813 handler: Arc<H>, 814 request: AdminRequest, 815 ) -> AdminRouteOutcome 816 where 817 H: RhiAdminHandler, 818 { 819 let document = match request_document(route, &request) { 820 Ok(document) => document, 821 Err(_) => return failure(RhiAdminHandlerErrorKind::Conflict, true), 822 }; 823 match handler.handle(document).await { 824 Ok(response) if response.route == route => match request.success(&response.value) { 825 Ok(outcome) => outcome, 826 Err(_) => failure(RhiAdminHandlerErrorKind::Internal, false), 827 }, 828 Ok(_) => failure(RhiAdminHandlerErrorKind::Internal, false), 829 Err(error) => failure(error.kind, false), 830 } 831 } 832 833 #[cfg(any(target_os = "linux", target_os = "macos"))] 834 fn failure(kind: RhiAdminHandlerErrorKind, invalid_request: bool) -> AdminRouteOutcome { 835 let (status, code, message) = if invalid_request { 836 ( 837 AdminRouteFailureStatus::BadRequest, 838 "invalid_request", 839 "admin request does not match the route model", 840 ) 841 } else { 842 match kind { 843 RhiAdminHandlerErrorKind::InvalidCursor => ( 844 AdminRouteFailureStatus::BadRequest, 845 "invalid_cursor", 846 "admin pagination cursor is invalid", 847 ), 848 RhiAdminHandlerErrorKind::OperationIdConflict => ( 849 AdminRouteFailureStatus::Conflict, 850 "operation_id_conflict", 851 "admin operation identity conflicts with retained state", 852 ), 853 RhiAdminHandlerErrorKind::NotFound => ( 854 AdminRouteFailureStatus::NotFound, 855 "not_found", 856 "admin resource was not found", 857 ), 858 RhiAdminHandlerErrorKind::Conflict => ( 859 AdminRouteFailureStatus::Conflict, 860 "operation_conflict", 861 "admin operation conflicts with current state", 862 ), 863 RhiAdminHandlerErrorKind::Unavailable => ( 864 AdminRouteFailureStatus::Unavailable, 865 "service_unavailable", 866 "admin operation is temporarily unavailable", 867 ), 868 RhiAdminHandlerErrorKind::Internal => ( 869 AdminRouteFailureStatus::Internal, 870 "internal_error", 871 "admin operation failed internally", 872 ), 873 } 874 }; 875 let code = AdminErrorCode::new(code).expect("fixed admin error code"); 876 let message = AdminErrorMessage::new(message).expect("fixed admin error message"); 877 AdminRouteOutcome::failure(AdminRouteFailure::new( 878 status, 879 AdminError::new(code, message), 880 )) 881 } 882 883 #[cfg(any(target_os = "linux", target_os = "macos"))] 884 fn request_document( 885 route: RhiAdminRoute, 886 request: &AdminRequest, 887 ) -> Result<RhiAdminRequestDocument, RhiAdminDocumentError> { 888 let (operation_id, value) = match route.method() { 889 RhiAdminMethod::Get => (None, query_model(route, request.query())?), 890 RhiAdminMethod::Post => { 891 let envelope = request 892 .decode_json::<AdminMutationRequest<Value>>() 893 .map_err(|_| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::Malformed))?; 894 ( 895 Some(envelope.operation_id().clone()), 896 envelope.into_request(), 897 ) 898 } 899 }; 900 validate_model(route.request_model(), &value)?; 901 let model_bytes = serde_json::to_vec(&value) 902 .map_err(|_| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::Malformed))?; 903 let parameter = route 904 .parameter_binding() 905 .map(|(name, type_name)| { 906 let value = request 907 .parameter(name) 908 .ok_or_else(|| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel))? 909 .to_owned() 910 .into_boxed_str(); 911 validate_type(type_name, &Value::String(value.to_string()), 0)?; 912 Ok((name, value)) 913 }) 914 .transpose()?; 915 Ok(RhiAdminRequestDocument { 916 route, 917 operation_id, 918 correlation_id: request.correlation_id().clone(), 919 parameter, 920 model_bytes: model_bytes.into_boxed_slice(), 921 }) 922 } 923 924 #[cfg(any(target_os = "linux", target_os = "macos"))] 925 fn query_model(route: RhiAdminRoute, query: Option<&str>) -> Result<Value, RhiAdminDocumentError> { 926 let Some(query) = query else { 927 return Ok(Value::Object(Map::new())); 928 }; 929 if query.is_empty() { 930 return Err(RhiAdminDocumentError::new( 931 RhiAdminDocumentErrorKind::InvalidModel, 932 )); 933 } 934 let fields = model_fields(route.request_model())?; 935 let mut output = Map::new(); 936 for item in query.split('&') { 937 let (raw_key, raw_value) = item 938 .split_once('=') 939 .ok_or_else(|| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel))?; 940 if raw_key.is_empty() 941 || !valid_percent_encoding(raw_key) 942 || !valid_percent_encoding(raw_value) 943 { 944 return Err(RhiAdminDocumentError::new( 945 RhiAdminDocumentErrorKind::InvalidModel, 946 )); 947 } 948 let key = percent_decode(raw_key)?; 949 let value = percent_decode(raw_value)?; 950 if output.contains_key(&key) { 951 return Err(RhiAdminDocumentError::new( 952 RhiAdminDocumentErrorKind::DuplicateField, 953 )); 954 } 955 let descriptor = fields 956 .get(&key) 957 .ok_or_else(|| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel))?; 958 let type_name = descriptor 959 .get("type") 960 .and_then(Value::as_str) 961 .ok_or_else(|| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel))?; 962 output.insert(key, query_scalar(type_name, value)?); 963 } 964 Ok(Value::Object(output)) 965 } 966 967 #[cfg(any(target_os = "linux", target_os = "macos"))] 968 fn query_scalar(type_name: &str, value: String) -> Result<Value, RhiAdminDocumentError> { 969 let descriptor = type_descriptor(type_name)?; 970 match descriptor.get("kind").and_then(Value::as_str) { 971 Some("integer") => { 972 let canonical = value == "0" 973 || (value.as_bytes().first().is_some_and(u8::is_ascii_digit) 974 && !value.starts_with('0') 975 && value.bytes().all(|byte| byte.is_ascii_digit())); 976 if !canonical { 977 return Err(RhiAdminDocumentError::new( 978 RhiAdminDocumentErrorKind::InvalidModel, 979 )); 980 } 981 value 982 .parse::<u64>() 983 .map(Value::from) 984 .map_err(|_| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel)) 985 } 986 Some("boolean") => match value.as_str() { 987 "true" => Ok(Value::Bool(true)), 988 "false" => Ok(Value::Bool(false)), 989 _ => Err(RhiAdminDocumentError::new( 990 RhiAdminDocumentErrorKind::InvalidModel, 991 )), 992 }, 993 _ => Ok(Value::String(value)), 994 } 995 } 996 997 #[cfg(any(target_os = "linux", target_os = "macos"))] 998 fn percent_decode(value: &str) -> Result<String, RhiAdminDocumentError> { 999 let bytes = value.as_bytes(); 1000 let mut decoded = Vec::with_capacity(bytes.len()); 1001 let mut index = 0; 1002 while index < bytes.len() { 1003 match bytes[index] { 1004 b'%' => { 1005 let high = hex_nibble(bytes[index + 1]).ok_or_else(invalid_model_error)?; 1006 let low = hex_nibble(bytes[index + 2]).ok_or_else(invalid_model_error)?; 1007 decoded.push((high << 4) | low); 1008 index += 3; 1009 } 1010 b'+' => { 1011 decoded.push(b' '); 1012 index += 1; 1013 } 1014 byte => { 1015 decoded.push(byte); 1016 index += 1; 1017 } 1018 } 1019 } 1020 String::from_utf8(decoded).map_err(|_| invalid_model_error()) 1021 } 1022 1023 #[cfg(any(target_os = "linux", target_os = "macos"))] 1024 const fn hex_nibble(byte: u8) -> Option<u8> { 1025 match byte { 1026 b'0'..=b'9' => Some(byte - b'0'), 1027 b'a'..=b'f' => Some(byte - b'a' + 10), 1028 b'A'..=b'F' => Some(byte - b'A' + 10), 1029 _ => None, 1030 } 1031 } 1032 1033 #[cfg(any(target_os = "linux", target_os = "macos"))] 1034 fn valid_percent_encoding(value: &str) -> bool { 1035 let bytes = value.as_bytes(); 1036 let mut index = 0; 1037 while index < bytes.len() { 1038 if bytes[index] == b'%' { 1039 if index + 2 >= bytes.len() 1040 || !bytes[index + 1].is_ascii_hexdigit() 1041 || !bytes[index + 2].is_ascii_hexdigit() 1042 { 1043 return false; 1044 } 1045 index += 3; 1046 } else { 1047 index += 1; 1048 } 1049 } 1050 true 1051 } 1052 1053 fn operator_contract() -> Result<&'static Value, RhiAdminDocumentError> { 1054 static CONTRACT: OnceLock<Option<Value>> = OnceLock::new(); 1055 CONTRACT 1056 .get_or_init(|| serde_json::from_str(OPERATOR_CONTRACT).ok()) 1057 .as_ref() 1058 .ok_or_else(invalid_model_error) 1059 } 1060 1061 #[cfg(any(target_os = "linux", target_os = "macos"))] 1062 fn operator_route_inventory_is_exact() -> bool { 1063 let Ok(contract) = operator_contract() else { 1064 return false; 1065 }; 1066 let Some(admin) = contract.get("admin") else { 1067 return false; 1068 }; 1069 let Some(routes) = admin.get("routes").and_then(Value::as_array) else { 1070 return false; 1071 }; 1072 let Some(models) = admin.get("models").and_then(Value::as_object) else { 1073 return false; 1074 }; 1075 routes.len() == RhiAdminRoute::ALL.len() 1076 && models.len() == 33 1077 && admin 1078 .pointer("/model_wire_contract/response_body_max_utf8_bytes") 1079 .and_then(Value::as_u64) 1080 == Some(RHI_ADMIN_RESPONSE_BODY_MAX_UTF8_BYTES as u64) 1081 && routes.iter().zip(RhiAdminRoute::ALL).all(|(wire, route)| { 1082 wire.get("method").and_then(Value::as_str) 1083 == Some(match route.method() { 1084 RhiAdminMethod::Get => "GET", 1085 RhiAdminMethod::Post => "POST", 1086 }) 1087 && wire.get("path").and_then(Value::as_str) == Some(route.path()) 1088 && wire.get("operation_id").and_then(Value::as_str) == Some(route.operation_id()) 1089 && wire.get("request_model").and_then(Value::as_str) == Some(route.request_model()) 1090 && wire.get("response_model").and_then(Value::as_str) 1091 == Some(route.response_model()) 1092 && wire.get("mutation").and_then(Value::as_bool) == Some(route.is_mutation()) 1093 }) 1094 } 1095 1096 fn model_fields(model_name: &str) -> Result<&'static Map<String, Value>, RhiAdminDocumentError> { 1097 operator_contract()? 1098 .pointer(&format!("/admin/models/{model_name}/fields")) 1099 .and_then(Value::as_object) 1100 .ok_or_else(|| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel)) 1101 } 1102 1103 fn type_descriptor(type_name: &str) -> Result<&'static Value, RhiAdminDocumentError> { 1104 operator_contract()? 1105 .pointer(&format!("/admin/types/{type_name}")) 1106 .ok_or_else(|| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel)) 1107 } 1108 1109 fn validate_model(model_name: &str, value: &Value) -> Result<(), RhiAdminDocumentError> { 1110 let object = value 1111 .as_object() 1112 .ok_or_else(|| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel))?; 1113 let fields = model_fields(model_name)?; 1114 for key in object.keys() { 1115 if !fields.contains_key(key) { 1116 return Err(RhiAdminDocumentError::new( 1117 RhiAdminDocumentErrorKind::InvalidModel, 1118 )); 1119 } 1120 } 1121 for (name, field) in fields { 1122 let required = field.get("presence").and_then(Value::as_str) == Some("required"); 1123 let type_name = field 1124 .get("type") 1125 .and_then(Value::as_str) 1126 .ok_or_else(|| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel))?; 1127 match object.get(name) { 1128 Some(value) => validate_type(type_name, value, 0)?, 1129 None if required => { 1130 return Err(RhiAdminDocumentError::new( 1131 RhiAdminDocumentErrorKind::InvalidModel, 1132 )); 1133 } 1134 None => {} 1135 } 1136 } 1137 Ok(()) 1138 } 1139 1140 fn validate_type( 1141 type_name: &str, 1142 value: &Value, 1143 depth: usize, 1144 ) -> Result<(), RhiAdminDocumentError> { 1145 if depth > 24 || value.is_null() { 1146 return Err(RhiAdminDocumentError::new(if value.is_null() { 1147 RhiAdminDocumentErrorKind::NullForbidden 1148 } else { 1149 RhiAdminDocumentErrorKind::InvalidModel 1150 })); 1151 } 1152 let descriptor = type_descriptor(type_name)?; 1153 match descriptor.get("kind").and_then(Value::as_str) { 1154 Some("literal") => { 1155 if descriptor.get("value") != Some(value) { 1156 return invalid_model(); 1157 } 1158 } 1159 Some("boolean") => { 1160 if !value.is_boolean() { 1161 return invalid_model(); 1162 } 1163 } 1164 Some("integer") => validate_integer(descriptor, value)?, 1165 Some("string") => validate_string(descriptor, value)?, 1166 Some("enum") => { 1167 if !descriptor 1168 .get("values") 1169 .and_then(Value::as_array) 1170 .is_some_and(|values| values.contains(value)) 1171 { 1172 return invalid_model(); 1173 } 1174 } 1175 Some("string_union") => validate_string_union(descriptor, value)?, 1176 Some("array") => validate_array(descriptor, value, depth + 1)?, 1177 Some("canonical_delimited_set") => { 1178 validate_delimited_set(descriptor, value, depth + 1)?; 1179 } 1180 Some("map") => validate_map(descriptor, value, depth + 1)?, 1181 Some("closed_object") => validate_closed_object(descriptor, value, depth + 1)?, 1182 Some("tagged_union") => validate_tagged_union(descriptor, value, depth + 1)?, 1183 Some("alias") => validate_type( 1184 descriptor 1185 .get("target") 1186 .and_then(Value::as_str) 1187 .ok_or_else(invalid_model_error)?, 1188 value, 1189 depth + 1, 1190 )?, 1191 Some("optional") => validate_type( 1192 descriptor 1193 .get("value") 1194 .and_then(Value::as_str) 1195 .ok_or_else(invalid_model_error)?, 1196 value, 1197 depth + 1, 1198 )?, 1199 Some("canonical_json_object") => { 1200 if !value.is_object() 1201 || serde_json::to_vec(value).ok().is_none_or(|bytes| { 1202 bytes.len() 1203 > descriptor 1204 .get("maximum_utf8_bytes") 1205 .and_then(Value::as_u64) 1206 .and_then(|value| usize::try_from(value).ok()) 1207 .unwrap_or(0) 1208 }) 1209 { 1210 return invalid_model(); 1211 } 1212 } 1213 _ => return invalid_model(), 1214 } 1215 Ok(()) 1216 } 1217 1218 fn validate_integer(descriptor: &Value, value: &Value) -> Result<(), RhiAdminDocumentError> { 1219 let number = value.as_u64().ok_or_else(invalid_model_error)?; 1220 let minimum = descriptor 1221 .get("minimum") 1222 .and_then(Value::as_u64) 1223 .unwrap_or(0); 1224 let maximum = descriptor 1225 .get("maximum") 1226 .and_then(Value::as_u64) 1227 .unwrap_or(u64::MAX); 1228 if !(minimum..=maximum).contains(&number) { 1229 return invalid_model(); 1230 } 1231 Ok(()) 1232 } 1233 1234 fn validate_string(descriptor: &Value, value: &Value) -> Result<(), RhiAdminDocumentError> { 1235 let string = value.as_str().ok_or_else(invalid_model_error)?; 1236 let exact = descriptor 1237 .get("utf8_bytes") 1238 .and_then(Value::as_u64) 1239 .and_then(|value| usize::try_from(value).ok()); 1240 let minimum = descriptor 1241 .get("minimum_utf8_bytes") 1242 .and_then(Value::as_u64) 1243 .and_then(|value| usize::try_from(value).ok()) 1244 .unwrap_or(0); 1245 let maximum = descriptor 1246 .get("maximum_utf8_bytes") 1247 .and_then(Value::as_u64) 1248 .and_then(|value| usize::try_from(value).ok()) 1249 .unwrap_or(usize::MAX); 1250 if exact.is_some_and(|exact| string.len() != exact) 1251 || !(minimum..=maximum).contains(&string.len()) 1252 || string.chars().any(char::is_control) 1253 { 1254 return invalid_model(); 1255 } 1256 if let Some(pattern) = descriptor.get("pattern").and_then(Value::as_str) { 1257 let valid = match pattern { 1258 "^[A-Za-z0-9][A-Za-z0-9._:-]*$" => bounded_id(string), 1259 "^[a-z][a-z0-9_]*$" => safe_code(string), 1260 "^[a-z][a-z0-9_-]*$" => stable_id(string), 1261 "^[0-9a-f]{32}$" => lower_hex(string, 32), 1262 "^[0-9a-f]{64}$" => lower_hex(string, 64), 1263 "^[0-9a-f]{40}$" => lower_hex(string, 40), 1264 "^/" => string.starts_with('/'), 1265 _ => false, 1266 }; 1267 if !valid { 1268 return invalid_model(); 1269 } 1270 } 1271 if descriptor.get("encoding").and_then(Value::as_str) == Some("canonical_base64url_no_padding") 1272 && !canonical_base64url_no_padding(string) 1273 { 1274 return invalid_model(); 1275 } 1276 if let Some(schemes) = descriptor.get("allowed_schemes").and_then(Value::as_array) { 1277 let parsed = url::Url::parse(string).map_err(|_| invalid_model_error())?; 1278 if !schemes 1279 .iter() 1280 .any(|scheme| scheme.as_str() == Some(parsed.scheme())) 1281 { 1282 return invalid_model(); 1283 } 1284 } 1285 Ok(()) 1286 } 1287 1288 fn validate_string_union(descriptor: &Value, value: &Value) -> Result<(), RhiAdminDocumentError> { 1289 let string = value.as_str().ok_or_else(invalid_model_error)?; 1290 if descriptor 1291 .get("simple_values") 1292 .and_then(Value::as_array) 1293 .is_some_and(|values| values.iter().any(|value| value.as_str() == Some(string))) 1294 { 1295 return Ok(()); 1296 } 1297 let kind = string 1298 .strip_prefix("sign_event:kind:") 1299 .ok_or_else(invalid_model_error)?; 1300 if kind.is_empty() 1301 || (kind.len() > 1 && kind.starts_with('0')) 1302 || !kind.bytes().all(|byte| byte.is_ascii_digit()) 1303 || kind.parse::<u32>().is_err() 1304 || string.len() 1305 > descriptor 1306 .get("maximum_utf8_bytes") 1307 .and_then(Value::as_u64) 1308 .and_then(|value| usize::try_from(value).ok()) 1309 .unwrap_or(0) 1310 { 1311 return invalid_model(); 1312 } 1313 Ok(()) 1314 } 1315 1316 fn validate_array( 1317 descriptor: &Value, 1318 value: &Value, 1319 depth: usize, 1320 ) -> Result<(), RhiAdminDocumentError> { 1321 let values = value.as_array().ok_or_else(invalid_model_error)?; 1322 let maximum = descriptor 1323 .get("maximum_items") 1324 .and_then(Value::as_u64) 1325 .and_then(|value| usize::try_from(value).ok()) 1326 .unwrap_or(0); 1327 if values.len() > maximum { 1328 return invalid_model(); 1329 } 1330 let item_type = descriptor 1331 .get("items") 1332 .and_then(Value::as_str) 1333 .ok_or_else(invalid_model_error)?; 1334 for item in values { 1335 validate_type(item_type, item, depth)?; 1336 } 1337 if descriptor.get("unique").and_then(Value::as_bool) == Some(true) { 1338 let mut unique = BTreeSet::new(); 1339 for item in values { 1340 let encoded = serde_json::to_vec(item).map_err(|_| invalid_model_error())?; 1341 if !unique.insert(encoded) { 1342 return invalid_model(); 1343 } 1344 } 1345 } 1346 if descriptor.get("canonical_sort").is_some() { 1347 let rendered = values 1348 .iter() 1349 .map(|value| value.as_str().ok_or_else(invalid_model_error)) 1350 .collect::<Result<Vec<_>, _>>()?; 1351 if !rendered.windows(2).all(|pair| pair[0] < pair[1]) { 1352 return invalid_model(); 1353 } 1354 } 1355 Ok(()) 1356 } 1357 1358 fn validate_delimited_set( 1359 descriptor: &Value, 1360 value: &Value, 1361 depth: usize, 1362 ) -> Result<(), RhiAdminDocumentError> { 1363 let rendered = value.as_str().ok_or_else(invalid_model_error)?; 1364 let maximum_bytes = descriptor 1365 .get("maximum_utf8_bytes") 1366 .and_then(Value::as_u64) 1367 .and_then(|value| usize::try_from(value).ok()) 1368 .unwrap_or(0); 1369 if rendered.len() > maximum_bytes { 1370 return invalid_model(); 1371 } 1372 if rendered.is_empty() { 1373 return if descriptor.get("empty_allowed").and_then(Value::as_bool) == Some(true) { 1374 Ok(()) 1375 } else { 1376 invalid_model() 1377 }; 1378 } 1379 let delimiter = descriptor 1380 .get("delimiter") 1381 .and_then(Value::as_str) 1382 .filter(|delimiter| delimiter.len() == 1) 1383 .ok_or_else(invalid_model_error)?; 1384 let items = rendered.split(delimiter).collect::<Vec<_>>(); 1385 let maximum_items = descriptor 1386 .get("maximum_items") 1387 .and_then(Value::as_u64) 1388 .and_then(|value| usize::try_from(value).ok()) 1389 .unwrap_or(0); 1390 if items.is_empty() 1391 || items.len() > maximum_items 1392 || items.iter().any(|item| item.is_empty()) 1393 || !items.windows(2).all(|pair| pair[0] < pair[1]) 1394 { 1395 return invalid_model(); 1396 } 1397 let item_type = descriptor 1398 .get("items") 1399 .and_then(Value::as_str) 1400 .ok_or_else(invalid_model_error)?; 1401 for item in items { 1402 validate_type(item_type, &Value::String(item.to_owned()), depth)?; 1403 } 1404 Ok(()) 1405 } 1406 1407 fn validate_map( 1408 descriptor: &Value, 1409 value: &Value, 1410 depth: usize, 1411 ) -> Result<(), RhiAdminDocumentError> { 1412 let values = value.as_object().ok_or_else(invalid_model_error)?; 1413 let maximum_entries = descriptor 1414 .get("maximum_entries") 1415 .and_then(Value::as_u64) 1416 .and_then(|value| usize::try_from(value).ok()) 1417 .unwrap_or(0); 1418 if values.len() > maximum_entries { 1419 return invalid_model(); 1420 } 1421 let key_type = descriptor 1422 .get("key") 1423 .and_then(Value::as_str) 1424 .ok_or_else(invalid_model_error)?; 1425 let value_type = descriptor 1426 .get("value") 1427 .and_then(Value::as_str) 1428 .ok_or_else(invalid_model_error)?; 1429 for (key, value) in values { 1430 validate_type(key_type, &Value::String(key.clone()), depth)?; 1431 validate_type(value_type, value, depth)?; 1432 } 1433 Ok(()) 1434 } 1435 1436 fn validate_closed_object( 1437 descriptor: &Value, 1438 value: &Value, 1439 depth: usize, 1440 ) -> Result<(), RhiAdminDocumentError> { 1441 let values = value.as_object().ok_or_else(invalid_model_error)?; 1442 let fields = descriptor 1443 .get("fields") 1444 .and_then(Value::as_object) 1445 .ok_or_else(invalid_model_error)?; 1446 if values.keys().any(|key| !fields.contains_key(key)) { 1447 return invalid_model(); 1448 } 1449 for (field, type_name) in fields { 1450 let type_name = type_name.as_str().ok_or_else(invalid_model_error)?; 1451 match values.get(field) { 1452 Some(value) => validate_type(type_name, value, depth)?, 1453 None if type_descriptor(type_name)? 1454 .get("kind") 1455 .and_then(Value::as_str) 1456 == Some("optional") => {} 1457 None => return invalid_model(), 1458 } 1459 } 1460 Ok(()) 1461 } 1462 1463 fn validate_tagged_union( 1464 descriptor: &Value, 1465 value: &Value, 1466 depth: usize, 1467 ) -> Result<(), RhiAdminDocumentError> { 1468 let object = value.as_object().ok_or_else(invalid_model_error)?; 1469 let discriminator = descriptor 1470 .get("discriminator") 1471 .and_then(Value::as_str) 1472 .ok_or_else(invalid_model_error)?; 1473 let selected = object.get(discriminator).ok_or_else(invalid_model_error)?; 1474 let variants = descriptor 1475 .get("variants") 1476 .and_then(Value::as_array) 1477 .ok_or_else(invalid_model_error)?; 1478 let mut matched = None; 1479 for variant in variants { 1480 let variant = variant.as_str().ok_or_else(invalid_model_error)?; 1481 let fields = type_descriptor(variant)? 1482 .get("fields") 1483 .and_then(Value::as_object) 1484 .ok_or_else(invalid_model_error)?; 1485 let discriminator_type = fields 1486 .get(discriminator) 1487 .and_then(Value::as_str) 1488 .ok_or_else(invalid_model_error)?; 1489 if type_descriptor(discriminator_type)?.get("value") == Some(selected) 1490 && matched.replace(variant).is_some() 1491 { 1492 return invalid_model(); 1493 } 1494 } 1495 validate_type(matched.ok_or_else(invalid_model_error)?, value, depth) 1496 } 1497 1498 fn bounded_id(value: &str) -> bool { 1499 (1..=128).contains(&value.len()) 1500 && value 1501 .as_bytes() 1502 .first() 1503 .is_some_and(u8::is_ascii_alphanumeric) 1504 && value 1505 .bytes() 1506 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b':' | b'-')) 1507 } 1508 1509 fn safe_code(value: &str) -> bool { 1510 (1..=64).contains(&value.len()) 1511 && value.as_bytes().first().is_some_and(u8::is_ascii_lowercase) 1512 && value 1513 .bytes() 1514 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_') 1515 } 1516 1517 fn stable_id(value: &str) -> bool { 1518 (1..=64).contains(&value.len()) 1519 && value.as_bytes().first().is_some_and(u8::is_ascii_lowercase) 1520 && value.bytes().all(|byte| { 1521 byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-') 1522 }) 1523 } 1524 1525 fn lower_hex(value: &str, exact: usize) -> bool { 1526 value.len() == exact 1527 && value 1528 .bytes() 1529 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) 1530 } 1531 1532 fn canonical_base64url_no_padding(value: &str) -> bool { 1533 if value.is_empty() || value.len() % 4 == 1 { 1534 return false; 1535 } 1536 let sextet = |byte: u8| match byte { 1537 b'A'..=b'Z' => Some(byte - b'A'), 1538 b'a'..=b'z' => Some(byte - b'a' + 26), 1539 b'0'..=b'9' => Some(byte - b'0' + 52), 1540 b'-' => Some(62), 1541 b'_' => Some(63), 1542 _ => None, 1543 }; 1544 let bytes = value.as_bytes(); 1545 if bytes.iter().copied().any(|byte| sextet(byte).is_none()) { 1546 return false; 1547 } 1548 let Some(last) = sextet(bytes[bytes.len() - 1]) else { 1549 return false; 1550 }; 1551 match value.len() % 4 { 1552 0 => true, 1553 2 => last & 0x0f == 0, 1554 3 => last & 0x03 == 0, 1555 _ => false, 1556 } 1557 } 1558 1559 fn invalid_model<T>() -> Result<T, RhiAdminDocumentError> { 1560 Err(invalid_model_error()) 1561 } 1562 1563 const fn invalid_model_error() -> RhiAdminDocumentError { 1564 RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel) 1565 } 1566 1567 const DUPLICATE_MARKER: &str = "rhi_admin_duplicate_field"; 1568 const NULL_MARKER: &str = "rhi_admin_null_forbidden"; 1569 1570 fn strict_json(bytes: &[u8]) -> Result<Value, RhiAdminDocumentError> { 1571 let mut deserializer = serde_json::Deserializer::from_slice(bytes); 1572 let value = StrictValueSeed 1573 .deserialize(&mut deserializer) 1574 .map_err(|error| { 1575 let message = error.to_string(); 1576 if message.contains(DUPLICATE_MARKER) { 1577 RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::DuplicateField) 1578 } else if message.contains(NULL_MARKER) { 1579 RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::NullForbidden) 1580 } else { 1581 RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::Malformed) 1582 } 1583 })?; 1584 deserializer 1585 .end() 1586 .map_err(|_| RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::Malformed))?; 1587 Ok(value) 1588 } 1589 1590 struct StrictValueSeed; 1591 1592 impl<'de> DeserializeSeed<'de> for StrictValueSeed { 1593 type Value = Value; 1594 1595 fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error> 1596 where 1597 D: serde::Deserializer<'de>, 1598 { 1599 deserializer.deserialize_any(StrictValueVisitor) 1600 } 1601 } 1602 1603 struct StrictValueVisitor; 1604 1605 impl<'de> Visitor<'de> for StrictValueVisitor { 1606 type Value = Value; 1607 1608 fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1609 formatter.write_str("a non-null JSON value") 1610 } 1611 1612 fn visit_bool<E>(self, value: bool) -> Result<Self::Value, E> { 1613 Ok(Value::Bool(value)) 1614 } 1615 1616 fn visit_i64<E>(self, value: i64) -> Result<Self::Value, E> { 1617 Ok(Value::from(value)) 1618 } 1619 1620 fn visit_u64<E>(self, value: u64) -> Result<Self::Value, E> { 1621 Ok(Value::from(value)) 1622 } 1623 1624 fn visit_f64<E>(self, value: f64) -> Result<Self::Value, E> 1625 where 1626 E: de::Error, 1627 { 1628 serde_json::Number::from_f64(value) 1629 .map(Value::Number) 1630 .ok_or_else(|| E::custom("invalid JSON number")) 1631 } 1632 1633 fn visit_str<E>(self, value: &str) -> Result<Self::Value, E> { 1634 Ok(Value::String(value.to_owned())) 1635 } 1636 1637 fn visit_string<E>(self, value: String) -> Result<Self::Value, E> { 1638 Ok(Value::String(value)) 1639 } 1640 1641 fn visit_none<E>(self) -> Result<Self::Value, E> 1642 where 1643 E: de::Error, 1644 { 1645 Err(E::custom(NULL_MARKER)) 1646 } 1647 1648 fn visit_unit<E>(self) -> Result<Self::Value, E> 1649 where 1650 E: de::Error, 1651 { 1652 Err(E::custom(NULL_MARKER)) 1653 } 1654 1655 fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error> 1656 where 1657 A: SeqAccess<'de>, 1658 { 1659 let mut values = Vec::new(); 1660 while let Some(value) = sequence.next_element_seed(StrictValueSeed)? { 1661 values.push(value); 1662 } 1663 Ok(Value::Array(values)) 1664 } 1665 1666 fn visit_map<A>(self, mut object: A) -> Result<Self::Value, A::Error> 1667 where 1668 A: MapAccess<'de>, 1669 { 1670 let mut values = Map::new(); 1671 while let Some(key) = object.next_key::<String>()? { 1672 if values.contains_key(&key) { 1673 return Err(de::Error::custom(DUPLICATE_MARKER)); 1674 } 1675 let value = object.next_value_seed(StrictValueSeed)?; 1676 values.insert(key, value); 1677 } 1678 Ok(Value::Object(values)) 1679 } 1680 } 1681 1682 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 1683 mod tests { 1684 use super::*; 1685 1686 fn sample_model(model_name: &str) -> Value { 1687 let fields = model_fields(model_name).expect("governed model"); 1688 let mut object = Map::new(); 1689 for (name, field) in fields { 1690 if field.get("presence").and_then(Value::as_str) == Some("required") { 1691 let type_name = field 1692 .get("type") 1693 .and_then(Value::as_str) 1694 .expect("field type"); 1695 object.insert(name.clone(), sample_type(type_name).expect("required type")); 1696 } 1697 } 1698 Value::Object(object) 1699 } 1700 1701 fn sample_type(type_name: &str) -> Option<Value> { 1702 let descriptor = type_descriptor(type_name).expect("governed type"); 1703 match descriptor 1704 .get("kind") 1705 .and_then(Value::as_str) 1706 .expect("type kind") 1707 { 1708 "literal" => Some(descriptor.get("value").expect("literal value").clone()), 1709 "boolean" => Some(Value::Bool(false)), 1710 "integer" => Some(Value::from( 1711 descriptor 1712 .get("minimum") 1713 .and_then(Value::as_u64) 1714 .unwrap_or(0), 1715 )), 1716 "string" => { 1717 let value = if descriptor.get("encoding").is_some() { 1718 "AA".to_owned() 1719 } else if descriptor.get("allowed_schemes").is_some() { 1720 "https://example.test/".to_owned() 1721 } else if descriptor.get("pattern").and_then(Value::as_str) == Some("^/") { 1722 "/var/lib/radroots/rhi".to_owned() 1723 } else if let Some(length) = descriptor.get("utf8_bytes").and_then(Value::as_u64) { 1724 "0".repeat(usize::try_from(length).expect("bounded sample length")) 1725 } else { 1726 "x".to_owned() 1727 }; 1728 Some(Value::String(value)) 1729 } 1730 "enum" => descriptor 1731 .get("values") 1732 .and_then(Value::as_array) 1733 .and_then(|values| values.first()) 1734 .cloned(), 1735 "string_union" => descriptor 1736 .get("simple_values") 1737 .and_then(Value::as_array) 1738 .and_then(|values| values.first()) 1739 .cloned(), 1740 "array" => Some(Value::Array(Vec::new())), 1741 "canonical_delimited_set" => Some(Value::String(String::new())), 1742 "map" | "canonical_json_object" => Some(Value::Object(Map::new())), 1743 "closed_object" => { 1744 let mut object = Map::new(); 1745 for (field, field_type) in descriptor 1746 .get("fields") 1747 .and_then(Value::as_object) 1748 .expect("closed fields") 1749 { 1750 if let Some(value) = sample_type(field_type.as_str().expect("closed type")) { 1751 object.insert(field.clone(), value); 1752 } 1753 } 1754 Some(Value::Object(object)) 1755 } 1756 "tagged_union" => descriptor 1757 .get("variants") 1758 .and_then(Value::as_array) 1759 .and_then(|variants| variants.first()) 1760 .and_then(Value::as_str) 1761 .and_then(sample_type), 1762 "alias" => descriptor 1763 .get("target") 1764 .and_then(Value::as_str) 1765 .and_then(sample_type), 1766 "optional" => None, 1767 kind => panic!("unsupported governed kind {kind}"), 1768 } 1769 } 1770 1771 fn response_document(route: RhiAdminRoute) -> RhiAdminResponseDocument { 1772 let value = sample_model(route.response_model()); 1773 validate_model(route.response_model(), &value).expect("generated response model"); 1774 let bytes = serde_json::to_vec(&value).expect("canonical response bytes"); 1775 RhiAdminResponseDocument::from_canonical_bytes(route, &bytes) 1776 .expect("admitted response document") 1777 } 1778 1779 #[test] 1780 fn complete_route_and_model_inventory_matches_the_machine_contract() { 1781 assert!(operator_route_inventory_is_exact()); 1782 assert_eq!(RhiAdminRoute::ALL.len(), 20); 1783 assert_eq!(RhiAdminRoute::COMMON.len(), 7); 1784 assert_eq!(RhiAdminRoute::DOMAIN.len(), 13); 1785 assert_eq!(RhiAdminRoute::ACTIVE.len(), 20); 1786 assert_eq!( 1787 RhiAdminRoute::ACTIVE, 1788 RhiAdminRoute::COMMON 1789 .into_iter() 1790 .chain(RhiAdminRoute::DOMAIN) 1791 .collect::<Vec<_>>() 1792 .as_slice() 1793 ); 1794 let referenced = RhiAdminRoute::ALL 1795 .into_iter() 1796 .flat_map(|route| [route.request_model(), route.response_model()]) 1797 .collect::<BTreeSet<_>>(); 1798 let governed = operator_contract().expect("operator contract")["admin"]["models"] 1799 .as_object() 1800 .expect("model inventory") 1801 .keys() 1802 .map(String::as_str) 1803 .collect::<BTreeSet<_>>(); 1804 assert_eq!(referenced, governed); 1805 assert_eq!(governed.len(), 33); 1806 assert!( 1807 RhiAdminRoute::COMMON 1808 .into_iter() 1809 .all(|route| route.parameter_binding().is_none()) 1810 ); 1811 for route in [ 1812 RhiAdminRoute::TradeProjection, 1813 RhiAdminRoute::TradeReportCurrent, 1814 RhiAdminRoute::TradeReports, 1815 ] { 1816 assert_eq!(route.parameter_binding(), Some(("trade_id", "trade_id"))); 1817 } 1818 for model in governed { 1819 let value = sample_model(model); 1820 validate_model(model, &value).expect("minimum exact model"); 1821 } 1822 } 1823 1824 #[test] 1825 fn strict_response_admission_rejects_duplicates_null_and_noncanonical_bytes() { 1826 let route = RhiAdminRoute::IdentityPublic; 1827 let valid = response_document(route); 1828 assert_eq!(valid.route(), route); 1829 assert!(!valid.canonical_bytes().is_empty()); 1830 1831 let duplicate = br#"{"generation":0,"generation":1,"public_key":"0000000000000000000000000000000000000000000000000000000000000000","role":"service"}"#; 1832 assert_eq!( 1833 RhiAdminResponseDocument::from_canonical_bytes(route, duplicate) 1834 .expect_err("duplicate key") 1835 .kind(), 1836 RhiAdminDocumentErrorKind::DuplicateField 1837 ); 1838 let nested_null = br#"{"generation":0,"public_key":null,"role":"service"}"#; 1839 assert_eq!( 1840 RhiAdminResponseDocument::from_canonical_bytes(route, nested_null) 1841 .expect_err("nested null") 1842 .kind(), 1843 RhiAdminDocumentErrorKind::NullForbidden 1844 ); 1845 let noncanonical = format!(" {}", String::from_utf8_lossy(valid.canonical_bytes())); 1846 assert_eq!( 1847 RhiAdminResponseDocument::from_canonical_bytes(route, noncanonical.as_bytes()) 1848 .expect_err("whitespace") 1849 .kind(), 1850 RhiAdminDocumentErrorKind::NonCanonical 1851 ); 1852 let invalid = br#"{"generation":0,"public_key":"0000000000000000000000000000000000000000000000000000000000000000","role":"administrator"}"#; 1853 assert_eq!( 1854 RhiAdminResponseDocument::from_canonical_bytes(route, invalid) 1855 .expect_err("invalid model") 1856 .kind(), 1857 RhiAdminDocumentErrorKind::InvalidModel 1858 ); 1859 assert_eq!( 1860 RhiAdminResponseDocument::from_canonical_bytes(route, b"") 1861 .expect_err("empty response") 1862 .kind(), 1863 RhiAdminDocumentErrorKind::Malformed 1864 ); 1865 let oversized = vec![b' '; RHI_ADMIN_RESPONSE_BODY_MAX_UTF8_BYTES + 1]; 1866 assert_eq!( 1867 RhiAdminResponseDocument::from_canonical_bytes(route, &oversized) 1868 .expect_err("oversized response") 1869 .kind(), 1870 RhiAdminDocumentErrorKind::TooLarge 1871 ); 1872 assert_eq!( 1873 format!("{valid:?}"), 1874 "RhiAdminResponseDocument { route: IdentityPublic, model: \"[redacted]\" }" 1875 ); 1876 } 1877 1878 #[test] 1879 fn query_and_nested_type_admission_is_exact_and_bounded() { 1880 let identity = query_model(RhiAdminRoute::IdentityStatus, Some("role=service")) 1881 .expect("common identity query"); 1882 validate_model("identity_status_query_v1", &identity).expect("identity query model"); 1883 for invalid in [ 1884 "", 1885 "role=service&role=service", 1886 "role=transport", 1887 "unknown=x", 1888 ] { 1889 let result = query_model(RhiAdminRoute::IdentityStatus, Some(invalid)) 1890 .and_then(|value| validate_model("identity_status_query_v1", &value)); 1891 assert!( 1892 result.is_err(), 1893 "common query `{invalid}` unexpectedly passed" 1894 ); 1895 } 1896 1897 let jobs = query_model( 1898 RhiAdminRoute::ReconciliationJobs, 1899 Some("cursor=AA&limit=200&state=pending"), 1900 ) 1901 .expect("canonical query"); 1902 validate_model("reconciliation_jobs_query_v1", &jobs).expect("query model"); 1903 for invalid in [ 1904 "limit=01", 1905 "limit=201", 1906 "limit=1&limit=2", 1907 "unknown=x", 1908 "cursor=%", 1909 "cursor=A&limit=1", 1910 "cursor=AB&limit=1", 1911 "cursor=A%3D%3D&limit=1", 1912 "state=unknown", 1913 ] { 1914 let result = query_model(RhiAdminRoute::ReconciliationJobs, Some(invalid)) 1915 .and_then(|value| validate_model("reconciliation_jobs_query_v1", &value)); 1916 assert!(result.is_err(), "query `{invalid}` unexpectedly passed"); 1917 } 1918 validate_type("stable_id", &Value::String("source-1".to_owned()), 0) 1919 .expect("stable identifier"); 1920 assert!(validate_type("stable_id", &Value::String("Source".to_owned()), 0).is_err()); 1921 validate_type("trade_id", &Value::String("0".repeat(32)), 0).expect("trade identifier"); 1922 assert!(validate_type("trade_id", &Value::String("0".repeat(31)), 0).is_err()); 1923 validate_type( 1924 "absolute_path", 1925 &Value::String("/var/lib/radroots/rhi".to_owned()), 1926 0, 1927 ) 1928 .expect("Unix absolute path"); 1929 assert!( 1930 validate_type( 1931 "absolute_path", 1932 &Value::String("C:\\rhi\\state".to_owned()), 1933 0, 1934 ) 1935 .is_err() 1936 ); 1937 assert!(validate_type("bounded_id", &Value::String("../escape".to_owned()), 0).is_err()); 1938 } 1939 1940 #[test] 1941 fn public_diagnostics_are_source_free_and_content_free() { 1942 let document = RhiAdminDocumentError::new(RhiAdminDocumentErrorKind::InvalidModel); 1943 let handler = RhiAdminHandlerError::new(RhiAdminHandlerErrorKind::Internal); 1944 let server = RhiAdminServerError::new(RhiAdminServerErrorKind::Bind); 1945 for rendered in [ 1946 format!("{document}"), 1947 format!("{document:?}"), 1948 format!("{handler}"), 1949 format!("{handler:?}"), 1950 format!("{server}"), 1951 format!("{server:?}"), 1952 ] { 1953 assert!(!rendered.contains("/tmp/protected")); 1954 assert!(!rendered.contains("credential")); 1955 } 1956 assert!(Error::source(&document).is_none()); 1957 assert!(Error::source(&handler).is_none()); 1958 assert!(Error::source(&server).is_none()); 1959 assert_eq!(server.code(), "admin_bind_failed"); 1960 } 1961 1962 #[cfg(any(target_os = "linux", target_os = "macos"))] 1963 mod native { 1964 use std::collections::BTreeMap; 1965 use std::fs; 1966 use std::path::Path; 1967 use std::sync::Mutex; 1968 1969 use radroots_service_host::{AdminClient, AdminClientTarget, AdminTransportLimits}; 1970 1971 use super::*; 1972 use tokio::io::{AsyncReadExt, AsyncWriteExt}; 1973 1974 const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); 1975 1976 type FixtureCall = (RhiAdminRoute, Option<String>, Option<String>, Box<[u8]>); 1977 type FixtureOperation = (RhiAdminRoute, Box<[u8]>); 1978 1979 struct FixtureHandler { 1980 calls: Mutex<Vec<FixtureCall>>, 1981 operations: Mutex<BTreeMap<String, FixtureOperation>>, 1982 } 1983 1984 impl FixtureHandler { 1985 fn new() -> Self { 1986 Self { 1987 calls: Mutex::new(Vec::new()), 1988 operations: Mutex::new(BTreeMap::new()), 1989 } 1990 } 1991 } 1992 1993 impl RhiAdminHandler for FixtureHandler { 1994 fn handle<'a>(&'a self, request: RhiAdminRequestDocument) -> RhiAdminFuture<'a> { 1995 Box::pin(async move { 1996 if request 1997 .model_bytes() 1998 .windows(15) 1999 .any(|bytes| bytes == b"\"cursor\":\"AAAA\"") 2000 { 2001 return Err(RhiAdminHandlerError::new( 2002 RhiAdminHandlerErrorKind::InvalidCursor, 2003 )); 2004 } 2005 if let Some(operation_id) = request.operation_id() { 2006 let mut operations = self.operations.lock().expect("operations"); 2007 if let Some((route, request_bytes)) = operations.get(operation_id) { 2008 if *route != request.route() 2009 || request_bytes.as_ref() != request.model_bytes() 2010 { 2011 return Err(RhiAdminHandlerError::new( 2012 RhiAdminHandlerErrorKind::OperationIdConflict, 2013 )); 2014 } 2015 return Ok(response_document(request.route())); 2016 } 2017 operations.insert( 2018 operation_id.to_owned(), 2019 (request.route(), request.model_bytes().into()), 2020 ); 2021 } 2022 self.calls.lock().expect("calls").push(( 2023 request.route(), 2024 request.operation_id().map(str::to_owned), 2025 request.parameter("trade_id").map(str::to_owned), 2026 request.model_bytes().into(), 2027 )); 2028 Ok(response_document(request.route())) 2029 }) 2030 } 2031 } 2032 2033 fn runtime_context() -> ( 2034 tempfile::TempDir, 2035 crate::RhiRuntimeContext, 2036 crate::RhiConfigDocumentV1, 2037 ) { 2038 let root = tempfile::Builder::new() 2039 .prefix("rhi-admin-") 2040 .tempdir_in("/tmp") 2041 .expect("short runtime root"); 2042 let root_path = root.path().to_str().expect("UTF-8 test root"); 2043 let invocation = crate::parse_rhi_cli_v1_from([ 2044 "rhi", 2045 "--profile", 2046 "repo-local", 2047 "--instance", 2048 "primary", 2049 "--repo-local-root", 2050 root_path, 2051 "run", 2052 ]) 2053 .expect("test CLI"); 2054 let resolver = crate::RadrootsPathResolver::new( 2055 crate::RadrootsPlatform::Linux, 2056 crate::RadrootsHostEnvironment::default(), 2057 ); 2058 let runtime = crate::resolve_rhi_runtime_context(&resolver, &invocation) 2059 .expect("test runtime context"); 2060 fs::create_dir_all(runtime.context().paths().run()).expect("runtime directory"); 2061 let configuration = 2062 crate::parse_rhi_config_v1(CONFIG.as_bytes(), crate::RhiConfigProfile::RepoLocal) 2063 .expect("test configuration"); 2064 (root, runtime, configuration) 2065 } 2066 2067 fn target_for(route: RhiAdminRoute, request: &Value) -> AdminClientTarget { 2068 let path = route 2069 .path() 2070 .replace("{trade_id}", "00000000000000000000000000000000"); 2071 if !matches!(route.method(), RhiAdminMethod::Get) 2072 || request.as_object().is_some_and(Map::is_empty) 2073 { 2074 return AdminClientTarget::new(path).expect("route target"); 2075 } 2076 let mut serializer = url::form_urlencoded::Serializer::new(String::new()); 2077 for (name, value) in request.as_object().expect("query object") { 2078 let rendered = value 2079 .as_str() 2080 .map(str::to_owned) 2081 .unwrap_or_else(|| value.to_string()); 2082 serializer.append_pair(name, &rendered); 2083 } 2084 AdminClientTarget::new(format!("{path}?{}", serializer.finish())).expect("query target") 2085 } 2086 2087 async fn raw_post(socket: &Path, path: &str, body: &str) -> String { 2088 let mut stream = tokio::net::UnixStream::connect(socket) 2089 .await 2090 .expect("raw connection"); 2091 let request = format!( 2092 "POST {path} HTTP/1.1\r\nHost: localhost\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", 2093 body.len() 2094 ); 2095 match stream.write_all(request.as_bytes()).await { 2096 Ok(()) => {} 2097 Err(error) 2098 if matches!( 2099 error.kind(), 2100 std::io::ErrorKind::BrokenPipe | std::io::ErrorKind::ConnectionReset 2101 ) => {} 2102 Err(error) => panic!("raw request: {error}"), 2103 } 2104 let mut response = Vec::new(); 2105 match stream.read_to_end(&mut response).await { 2106 Ok(_) => {} 2107 Err(error) if error.kind() == std::io::ErrorKind::ConnectionReset => {} 2108 Err(error) => panic!("raw response: {error}"), 2109 } 2110 String::from_utf8(response).expect("HTTP response") 2111 } 2112 2113 #[tokio::test] 2114 async fn twenty_active_routes_round_trip_over_the_hardened_unix_boundary() { 2115 let (_root, runtime, configuration) = runtime_context(); 2116 let socket = runtime.artifacts().admin_socket().to_path_buf(); 2117 let handler = Arc::new(FixtureHandler::new()); 2118 let server = RhiAdminServer::new(&configuration, Arc::clone(&handler)) 2119 .expect("production admin server") 2120 .bind(&runtime) 2121 .await 2122 .expect("canonical admin binding"); 2123 let cancellation = RhiAdminCancellationToken::new(); 2124 let server_cancellation = cancellation.clone(); 2125 let task = tokio::spawn(async move { 2126 server 2127 .serve(server_cancellation) 2128 .await 2129 .expect("serve RHI admin"); 2130 }); 2131 let client = 2132 AdminClient::new(&socket, AdminTransportLimits::DEFAULT).expect("admin client"); 2133 2134 use std::os::unix::fs::PermissionsExt as _; 2135 assert_eq!( 2136 fs::metadata(&socket) 2137 .expect("admin socket metadata") 2138 .permissions() 2139 .mode() 2140 & 0o777, 2141 0o600 2142 ); 2143 2144 for (index, route) in RhiAdminRoute::ACTIVE.into_iter().enumerate() { 2145 let request = sample_model(route.request_model()); 2146 let target = target_for(route, &request); 2147 let response = match route.method() { 2148 RhiAdminMethod::Get => client.get::<Value>(&target).await.expect("GET route"), 2149 RhiAdminMethod::Post => { 2150 let operation_id = AdminOperationId::new(format!("operation-{index}")) 2151 .expect("operation ID"); 2152 client 2153 .mutate::<_, Value>(&target, operation_id, None, &request) 2154 .await 2155 .expect("POST route") 2156 } 2157 }; 2158 validate_model(route.response_model(), response.result()) 2159 .expect("route response model"); 2160 } 2161 2162 let replay_target = 2163 AdminClientTarget::new("/v1/state/backup").expect("mutation target"); 2164 let replay_operation = AdminOperationId::new("operation-5").expect("operation ID"); 2165 let replay_request = sample_model("state_backup_request_v1"); 2166 client 2167 .mutate::<_, Value>( 2168 &replay_target, 2169 replay_operation.clone(), 2170 None, 2171 &replay_request, 2172 ) 2173 .await 2174 .expect("exact operation replay"); 2175 let mut conflicting_request = replay_request; 2176 conflicting_request 2177 .as_object_mut() 2178 .expect("backup request") 2179 .insert( 2180 "target_path".to_owned(), 2181 Value::String("/tmp/different-backup".to_owned()), 2182 ); 2183 let conflict_error = client 2184 .mutate::<_, Value>(&replay_target, replay_operation, None, &conflicting_request) 2185 .await 2186 .expect_err("conflicting operation reuse"); 2187 assert_eq!( 2188 conflict_error 2189 .failure() 2190 .expect("failure envelope") 2191 .error() 2192 .code() 2193 .as_str(), 2194 "operation_id_conflict" 2195 ); 2196 2197 for body in [ 2198 r#"{"contract_version":1,"operation_id":"duplicate","request":{"confirmation":"confirm","expected_generation":0,"expected_generation":1,"target_path":"/tmp/backup"}}"#, 2199 r#"{"contract_version":1,"operation_id":"null","request":{"confirmation":"confirm","expected_generation":null,"target_path":"/tmp/backup"}}"#, 2200 ] { 2201 let response = raw_post(&socket, "/v1/state/backup", body).await; 2202 assert!(response.starts_with("HTTP/1.1 400 "), "{response}"); 2203 } 2204 2205 for (index, route) in RhiAdminRoute::ACTIVE 2206 .into_iter() 2207 .filter(|route| route.is_mutation()) 2208 .enumerate() 2209 { 2210 let request = sample_model(route.request_model()); 2211 for version in [0, 2] { 2212 let body = serde_json::json!({ 2213 "contract_version": version, 2214 "operation_id": format!("invalid-version-{index}-{version}"), 2215 "request": request, 2216 }) 2217 .to_string(); 2218 let response = raw_post(&socket, route.path(), &body).await; 2219 assert!(response.starts_with("HTTP/1.1 400 "), "{response}"); 2220 } 2221 let missing_version = serde_json::json!({ 2222 "operation_id": format!("missing-version-{index}"), 2223 "request": request, 2224 }) 2225 .to_string(); 2226 let response = raw_post(&socket, route.path(), &missing_version).await; 2227 assert!(response.starts_with("HTTP/1.1 400 "), "{response}"); 2228 } 2229 2230 let request_limit = 2231 usize::try_from(AdminTransportLimits::DEFAULT.request_body_utf8_bytes()) 2232 .expect("request limit"); 2233 let oversized = " ".repeat(request_limit + 1); 2234 let response = raw_post(&socket, "/v1/state/backup", &oversized).await; 2235 assert!(response.starts_with("HTTP/1.1 413 "), "{response}"); 2236 2237 let domain_target = AdminClientTarget::new("/v1/reconciliation/jobs?limit=201") 2238 .expect("bounded domain route"); 2239 assert!(client.get::<Value>(&domain_target).await.is_err()); 2240 2241 let invalid_cursor_target = 2242 AdminClientTarget::new("/v1/reconciliation/jobs?cursor=AAAA&limit=1") 2243 .expect("authenticated cursor route"); 2244 let invalid_cursor = client 2245 .get::<Value>(&invalid_cursor_target) 2246 .await 2247 .expect_err("domain handler rejects unbound cursor"); 2248 assert_eq!( 2249 invalid_cursor 2250 .failure() 2251 .expect("failure envelope") 2252 .error() 2253 .code() 2254 .as_str(), 2255 "invalid_cursor" 2256 ); 2257 2258 let invalid_trade = AdminClientTarget::new("/v1/trades/not-hex/projection") 2259 .expect("trade route target"); 2260 assert!(client.get::<Value>(&invalid_trade).await.is_err()); 2261 2262 assert!(AdminClientTarget::new("/v2/status").is_err()); 2263 2264 for (index, path) in ["/v1/identity/rekey", "/v1/identity/replace"] 2265 .into_iter() 2266 .enumerate() 2267 { 2268 let removed_target = 2269 AdminClientTarget::new(path).expect("removed live identity route"); 2270 assert!( 2271 client 2272 .mutate::<_, Value>( 2273 &removed_target, 2274 AdminOperationId::new(format!("removed-identity-{index}")) 2275 .expect("operation ID"), 2276 None, 2277 serde_json::json!({}), 2278 ) 2279 .await 2280 .is_err() 2281 ); 2282 } 2283 2284 { 2285 let calls = handler.calls.lock().expect("calls"); 2286 assert_eq!(calls.len(), 20); 2287 for (index, (route, operation_id, parameter, request)) in calls.iter().enumerate() { 2288 assert_eq!(*route, RhiAdminRoute::ACTIVE[index]); 2289 assert_eq!(operation_id.is_some(), route.is_mutation()); 2290 assert_eq!( 2291 parameter.is_some(), 2292 matches!( 2293 route, 2294 RhiAdminRoute::TradeProjection 2295 | RhiAdminRoute::TradeReportCurrent 2296 | RhiAdminRoute::TradeReports 2297 ) 2298 ); 2299 validate_model( 2300 route.request_model(), 2301 &serde_json::from_slice(request).expect("retained request model"), 2302 ) 2303 .expect("retained exact model"); 2304 } 2305 } 2306 cancellation.cancel(); 2307 task.await.expect("server task"); 2308 assert!(!socket.exists()); 2309 } 2310 2311 #[test] 2312 fn production_server_projects_exact_validated_admin_limits() { 2313 let (_root, _runtime, configuration) = runtime_context(); 2314 assert_eq!( 2315 admin_transport_limits(&configuration).expect("admin limits"), 2316 AdminTransportLimits::DEFAULT 2317 ); 2318 } 2319 } 2320 }