runtime_graph.rs (107189B)
1 //! Final native Myc daemon graph, readiness, reconnect, and bounded shutdown. 2 3 use core::{fmt, future::Future, future::pending, pin::Pin, time::Duration}; 4 use std::{ 5 collections::VecDeque, 6 error::Error, 7 path::PathBuf, 8 sync::{ 9 Arc, 10 atomic::{AtomicBool, Ordering}, 11 }, 12 }; 13 14 use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD}; 15 use radroots_service_host::{ 16 EntropySource, GracefulShutdown, HostError, HostErrorKind, MetricValue, MonotonicClock, 17 ProcessSignal, ProcessSignalAdapter, ProcessSignalFuture, ProcessSignalSource, 18 ShutdownDisposition, ShutdownPhase, ShutdownPhaseFuture, ShutdownPhaseHandler, 19 SupervisedTaskExitStatus, SystemEntropy, SystemMonotonicClock, SystemWallClock, 20 TaskClassification, TaskMetadata, TaskName, TaskSupervisor, WallClock, 21 }; 22 use radroots_service_sqlite::{ 23 BackupCreatedAtUnixMs, MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, 24 }; 25 use radroots_transport::{BoxSubscription, SubscriptionNext, source::SubscriptionCheckpoint}; 26 use serde_json::{Value, json}; 27 use sha2::{Digest, Sha256}; 28 use tokio::sync::mpsc; 29 30 use crate::delivery_worker::{MycDeliveryExecutionEvidence, MycDeliveryWorker}; 31 use crate::provider_executor::MycProviderExecutor; 32 use crate::runtime_nip46::{ 33 MycNip46DispatchErrorKind, MycRuntimeNip46AdmissionEvidence, MycRuntimeNip46Coordinator, 34 }; 35 use crate::state_admin::AdminJournalOperationError; 36 use crate::state_connection::{ 37 MycAdminConnectionAction, apply_admin_challenge_authorize, apply_admin_challenge_require, 38 apply_admin_connection_mutation, 39 }; 40 use crate::state_discovery::{apply_admin_discovery_publish, prepare_discovery_signing_bytes}; 41 use crate::state_governance::MycAdminAuditQuery; 42 use crate::transport_nostr_adapter::MycNostrIngressAdapter; 43 use crate::{ 44 MycAdminCancellationToken, MycAdminFuture, MycAdminHandler, MycAdminHandlerError, 45 MycAdminHandlerErrorKind, MycAdminOperationAdmission, MycAdminOperationErrorKind, 46 MycAdminOperationJournalPolicy, MycAdminOperationTimeUnixMs, MycAdminRequestDocument, 47 MycAdminResponseDocument, MycAdminRoute, MycAdminServer, MycAuditCorrelationId, MycAuditKind, 48 MycAuditOutcome, MycAuditPageLimit, MycAuthorizationChallengeId, 49 MycAuthorizationChallengeNonce, MycAuthorizationChallengeRequest, 50 MycAuthorizationChallengeState, MycAuthorizationChallengeUrl, MycBootstrapProfileV1, 51 MycBoundAdminServer, MycBoundOperationsServer, MycConfigDocumentV1, MycConnectionCountsV1, 52 MycConnectionId, MycConnectionPermission, MycConnectionPermissionSet, 53 MycConnectionPolicyGeneration, MycConnectionStatus, MycConnectionTimeUnixMs, 54 MycDeliveryAttemptNonce, MycDeliveryRecoveryEntropy, MycDeliveryTimeUnixMs, 55 MycDiscoveryCommitRequest, MycIdentityHealthV1, MycOperationsCancellationToken, 56 MycOperationsServer, MycOutboxStatusV1, MycPersistenceHealthV1, MycPersistenceStatusV1, 57 MycProcessResult, MycProcessSignal, MycProcessSignalSource, MycProviderCorrelationId, 58 MycProviderDeadlineUnixMs, MycProviderOperation, MycProviderOperationId, 59 MycProviderOperationInput, MycProviderResponseObservedAtUnixMs, MycProviderRole, 60 MycProviderStatusV1, MycRateLimitClass, MycRelayTransportStatusV1, MycRuntimeContext, 61 MycServicePhase, MycSignerOperationId, MycStateHost, MycStatusBuildInfoV1, MycStatusBuildMode, 62 MycStatusCommonV1, MycStatusConfigurationIdentityV1, MycStatusConfigurationSource, 63 MycStatusObservationV1, MycStatusPublisher, MycStatusReader, MycStatusReasonCode, 64 MycStatusReasonCodes, MycTaskCancellation, MycTransportHealthV1, myc_status_cache, 65 open_myc_state_read_write_from_config, 66 }; 67 68 const TASK_ADMIN_SERVER: &str = "admin_server"; 69 const TASK_OPERATIONS_SERVER: &str = "operations_server"; 70 const TASK_RELAY_INGRESS: &str = "relay_ingress"; 71 const TASK_PROVIDER_DISPATCH: &str = "provider_dispatch"; 72 const TASK_DELIVERY_OUTBOX: &str = "delivery_outbox"; 73 74 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 75 enum MycDaemonErrorKind { 76 State, 77 Provider, 78 Relay, 79 Admin, 80 Operations, 81 Runtime, 82 } 83 84 struct MycDaemonError { 85 kind: MycDaemonErrorKind, 86 } 87 88 impl MycDaemonError { 89 const fn new(kind: MycDaemonErrorKind) -> Self { 90 Self { kind } 91 } 92 93 const fn process_result(&self) -> MycProcessResult { 94 match self.kind { 95 MycDaemonErrorKind::State => MycProcessResult::StateOrIdentityUnavailable, 96 MycDaemonErrorKind::Provider 97 | MycDaemonErrorKind::Relay 98 | MycDaemonErrorKind::Admin 99 | MycDaemonErrorKind::Operations => MycProcessResult::ServiceOrDependencyUnavailable, 100 MycDaemonErrorKind::Runtime => MycProcessResult::UnexpectedInternal, 101 } 102 } 103 } 104 105 impl fmt::Debug for MycDaemonError { 106 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 107 formatter 108 .debug_struct("MycDaemonError") 109 .field("kind", &self.kind) 110 .finish() 111 } 112 } 113 114 impl fmt::Display for MycDaemonError { 115 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 116 formatter.write_str("Myc daemon failed") 117 } 118 } 119 120 impl Error for MycDaemonError {} 121 122 struct HostSignalSource<S> { 123 inner: S, 124 } 125 126 impl<S> HostSignalSource<S> { 127 const fn new(inner: S) -> Self { 128 Self { inner } 129 } 130 } 131 132 impl<S> ProcessSignalSource for HostSignalSource<S> 133 where 134 S: MycProcessSignalSource, 135 { 136 fn next_signal(&mut self) -> ProcessSignalFuture<'_> { 137 Box::pin(async move { 138 self.inner.next_signal().await.map(|signal| match signal { 139 MycProcessSignal::Interrupt => ProcessSignal::Interrupt, 140 #[cfg(unix)] 141 MycProcessSignal::Terminate => ProcessSignal::Terminate, 142 }) 143 }) 144 } 145 } 146 147 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 148 enum HealthEvent { 149 Providers(bool), 150 RequiredRelays(bool), 151 StateChanged, 152 } 153 154 struct RuntimeHealth { 155 providers: bool, 156 required_relays: bool, 157 operations: bool, 158 } 159 160 struct RuntimeIngressItem { 161 raw_event: Box<[u8]>, 162 relay_id: crate::MycRateRelayId, 163 observed_at_unix_ms: u64, 164 } 165 166 struct RuntimeIngressSlot { 167 index: usize, 168 required: bool, 169 checkpoints: Vec<SubscriptionCheckpoint>, 170 state: RuntimeIngressSlotState, 171 retry_delay_ms: u64, 172 } 173 174 enum RuntimeIngressSlotState { 175 Active(BoxSubscription), 176 Waiting, 177 Connecting( 178 Pin< 179 Box< 180 dyn Future< 181 Output = Result< 182 BoxSubscription, 183 crate::transport_nostr_adapter::MycRelayAdapterError, 184 >, 185 > + Send, 186 >, 187 >, 188 ), 189 } 190 191 enum RuntimeIngressAction { 192 Next(Result<SubscriptionNext, radroots_transport::Error>), 193 Retry, 194 Connected(Result<BoxSubscription, crate::transport_nostr_adapter::MycRelayAdapterError>), 195 } 196 197 impl RuntimeHealth { 198 const fn ready(operations: bool) -> Self { 199 Self { 200 providers: true, 201 required_relays: true, 202 operations, 203 } 204 } 205 206 fn observe(&mut self, event: HealthEvent) -> bool { 207 match event { 208 HealthEvent::Providers(value) => { 209 let changed = self.providers != value; 210 self.providers = value; 211 changed 212 } 213 HealthEvent::RequiredRelays(value) => { 214 let changed = self.required_relays != value; 215 self.required_relays = value; 216 changed 217 } 218 HealthEvent::StateChanged => true, 219 } 220 } 221 } 222 223 struct RuntimeStatusContext { 224 runtime: MycRuntimeContext, 225 configuration: Arc<MycConfigDocumentV1>, 226 state: Arc<MycStateHost>, 227 clock: SystemMonotonicClock, 228 generation: u64, 229 connection_counts: MycConnectionCountsV1, 230 outbox: MycOutboxStatusV1, 231 } 232 233 impl RuntimeStatusContext { 234 async fn refresh_state(&mut self) -> Result<(), MycDaemonError> { 235 let connection_counts = self 236 .state 237 .repository() 238 .read_runtime_connection_counts() 239 .await 240 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?; 241 let outbox = self 242 .state 243 .repository() 244 .read_runtime_outbox_status() 245 .await 246 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?; 247 self.connection_counts = connection_counts; 248 self.outbox = outbox; 249 Ok(()) 250 } 251 252 fn observation( 253 &self, 254 phase: MycServicePhase, 255 health: &RuntimeHealth, 256 ) -> Result<MycStatusObservationV1, MycDaemonError> { 257 let ready = matches!(phase, MycServicePhase::Ready | MycServicePhase::Degraded) 258 && health.providers 259 && health.required_relays; 260 let mut reasons = Vec::new(); 261 if !health.providers { 262 reasons.push(MycStatusReasonCode::SignerProviderUnavailable); 263 } 264 if !health.required_relays { 265 reasons.push(MycStatusReasonCode::RequiredRelayUnavailable); 266 reasons.push(MycStatusReasonCode::SubscriberNotActive); 267 } 268 if !health.operations { 269 reasons.push(MycStatusReasonCode::OperationsListenerFailed); 270 } 271 if phase == MycServicePhase::Stopping { 272 reasons.push(MycStatusReasonCode::ShutdownInProgress); 273 } 274 let reasons = MycStatusReasonCodes::new(reasons) 275 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 276 let build = runtime_build_info()?; 277 let configuration = MycStatusConfigurationIdentityV1::new( 278 hex::encode(self.state.metadata().configuration_digest().as_bytes()), 279 if self.runtime.profile() == MycBootstrapProfileV1::RepoLocal { 280 MycStatusConfigurationSource::DerivedRepoLocal 281 } else { 282 MycStatusConfigurationSource::ExplicitConfig 283 }, 284 ) 285 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 286 let persistence = MycPersistenceStatusV1::new( 287 MycPersistenceHealthV1::Ready, 288 self.state 289 .metadata() 290 .initial_database_metadata() 291 .state_schema_version() 292 .get(), 293 self.generation, 294 crate::MycIntegrityStateV1::Verified, 295 MycStatusReasonCodes::empty(), 296 ) 297 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 298 let common = MycStatusCommonV1::new( 299 phase, 300 ready, 301 reasons.clone(), 302 u64::try_from( 303 self.clock 304 .now_monotonic() 305 .duration_since_origin() 306 .as_millis(), 307 ) 308 .unwrap_or(u64::MAX), 309 build, 310 configuration, 311 persistence, 312 ) 313 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 314 let identity = |configured: bool| { 315 MycIdentityHealthV1::new(configured, configured && health.providers, reasons.clone()) 316 }; 317 let provider = MycProviderStatusV1::new( 318 identity(true).map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, 319 identity(true).map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, 320 identity( 321 self.configuration 322 .provider_contract() 323 .binding(MycProviderRole::Discovery) 324 .is_some(), 325 ) 326 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, 327 reasons.clone(), 328 ) 329 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 330 let transport = MycRelayTransportStatusV1::new( 331 if health.required_relays { 332 MycTransportHealthV1::Ready 333 } else { 334 MycTransportHealthV1::Unavailable 335 }, 336 health.required_relays, 337 if health.required_relays { 338 u64::try_from(self.configuration.relay_count()).unwrap_or(u64::MAX) 339 } else { 340 0 341 }, 342 reasons, 343 ) 344 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 345 Ok(MycStatusObservationV1::new( 346 common, 347 provider, 348 transport, 349 self.connection_counts, 350 self.outbox, 351 )) 352 } 353 } 354 355 struct RuntimeAdminHandler { 356 state: Arc<MycStateHost>, 357 configuration: Arc<MycConfigDocumentV1>, 358 providers: Arc<MycProviderExecutor>, 359 status: MycStatusReader, 360 accepting_mutations: Arc<AtomicBool>, 361 cursor_key: [u8; 32], 362 health: mpsc::Sender<HealthEvent>, 363 } 364 365 impl RuntimeAdminHandler { 366 async fn handle_inner( 367 &self, 368 request: MycAdminRequestDocument, 369 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 370 if request.route().is_mutation() && !self.accepting_mutations.load(Ordering::Acquire) { 371 return Err(handler_error(MycAdminHandlerErrorKind::Unavailable)); 372 } 373 match request.route() { 374 MycAdminRoute::Status => MycAdminResponseDocument::from_canonical_bytes( 375 request.route(), 376 self.status.snapshot().detailed_status_json(), 377 ) 378 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)), 379 MycAdminRoute::EffectiveConfig => MycAdminResponseDocument::from_canonical_bytes( 380 request.route(), 381 self.configuration.effective().canonical_json().as_bytes(), 382 ) 383 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)), 384 MycAdminRoute::IdentityStatus | MycAdminRoute::IdentityPublic => { 385 self.identity_response(request).await 386 } 387 MycAdminRoute::StateStatus => self.state_status(request).await, 388 MycAdminRoute::StateBackup => self.backup(request).await, 389 MycAdminRoute::MetricsSnapshot => self.metrics_snapshot(request), 390 MycAdminRoute::ConnectionsList => self.connections(request).await, 391 MycAdminRoute::AuditEvents => self.audit_events(request).await, 392 MycAdminRoute::AuditSummary => self.audit_summary(request).await, 393 MycAdminRoute::DiscoveryDesired => self.discovery_desired(request).await, 394 MycAdminRoute::ConnectionApprove 395 | MycAdminRoute::ConnectionReject 396 | MycAdminRoute::ConnectionRevoke => self.connection_mutation(request).await, 397 MycAdminRoute::ChallengeRequire => self.challenge_require(request).await, 398 MycAdminRoute::ChallengeAuthorize => self.challenge_authorize(request).await, 399 MycAdminRoute::DiscoveryRender 400 | MycAdminRoute::DiscoveryRefresh 401 | MycAdminRoute::DiscoveryPublish => self.discovery_mutation(request).await, 402 } 403 } 404 405 async fn identity_response( 406 &self, 407 request: MycAdminRequestDocument, 408 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 409 let model: Value = serde_json::from_slice(request.model_bytes()) 410 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 411 let role_name = model 412 .pointer("/role") 413 .and_then(Value::as_str) 414 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 415 let role = match role_name { 416 "transport" => MycProviderRole::Transport, 417 "user" => MycProviderRole::User, 418 "discovery" => MycProviderRole::Discovery, 419 _ => return Err(handler_error(MycAdminHandlerErrorKind::Internal)), 420 }; 421 let binding = self.configuration.provider_contract().binding(role); 422 let expected = match role { 423 MycProviderRole::Transport => { 424 Some(self.state.metadata().expected_identities().transport()) 425 } 426 MycProviderRole::User => Some(self.state.metadata().expected_identities().user()), 427 MycProviderRole::Discovery => self.state.metadata().expected_identities().discovery(), 428 }; 429 let generation = u64::from( 430 self.state 431 .repository() 432 .current_configuration_generation() 433 .await 434 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 435 ); 436 let value = if request.route() == MycAdminRoute::IdentityPublic { 437 let expected = 438 expected.ok_or_else(|| handler_error(MycAdminHandlerErrorKind::NotFound))?; 439 json!({"generation":generation,"public_key":expected.as_hex(),"role":role_name}) 440 } else { 441 let mut value = json!({ 442 "available": binding.is_some(), 443 "configured": binding.is_some(), 444 "generation": generation, 445 "reason_codes": [], 446 "role": role_name, 447 }); 448 if let Some(binding) = binding { 449 value["provider"] = Value::String(binding.kind().as_str().to_owned()); 450 } 451 if let Some(expected) = expected { 452 value["public_key"] = Value::String(expected.as_hex().to_owned()); 453 } 454 value 455 }; 456 response(request.route(), value) 457 } 458 459 async fn state_status( 460 &self, 461 request: MycAdminRequestDocument, 462 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 463 let generation = self 464 .state 465 .repository() 466 .current_configuration_generation() 467 .await 468 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 469 response( 470 request.route(), 471 json!({ 472 "backup_eligible": true, 473 "generation": u64::from(generation), 474 "integrity": "verified", 475 "reason_codes": [], 476 "schema_version": self.state.metadata().initial_database_metadata().state_schema_version().get(), 477 "writer_lock": "held_by_daemon", 478 }), 479 ) 480 } 481 482 async fn backup( 483 &self, 484 request: MycAdminRequestDocument, 485 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 486 let now_seconds = SystemWallClock 487 .now_utc() 488 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))? 489 .get(); 490 let now_millis = now_seconds 491 .checked_mul(1_000) 492 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 493 let prepared_at = MycAdminOperationTimeUnixMs::new(now_millis) 494 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 495 let model: Value = serde_json::from_slice(request.model_bytes()) 496 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 497 let expected_generation = model 498 .pointer("/expected_generation") 499 .and_then(Value::as_u64) 500 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 501 let actual_generation = u64::from( 502 self.state 503 .repository() 504 .current_configuration_generation() 505 .await 506 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 507 ); 508 if expected_generation != actual_generation { 509 return Err(handler_error(MycAdminHandlerErrorKind::Conflict)); 510 } 511 let prepared = match self 512 .state 513 .repository() 514 .prepare_admin_operation(&request, prepared_at) 515 .await 516 .map_err(map_admin_journal_error)? 517 { 518 MycAdminOperationAdmission::ExactReplay(response) => return Ok(response), 519 MycAdminOperationAdmission::Prepared(prepared) => prepared, 520 }; 521 let target = model 522 .pointer("/target_path") 523 .and_then(Value::as_str) 524 .map(PathBuf::from) 525 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 526 let manifest = self 527 .state 528 .capture_online_backup( 529 &target, 530 BackupCreatedAtUnixMs::new(now_millis) 531 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 532 ) 533 .await 534 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))?; 535 let completed_at_seconds = SystemWallClock 536 .now_utc() 537 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))? 538 .get(); 539 let completed_at = MycAdminOperationTimeUnixMs::new( 540 completed_at_seconds 541 .checked_mul(1_000) 542 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?, 543 ) 544 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 545 let operation_id = request 546 .operation_id() 547 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 548 let response = response( 549 request.route(), 550 json!({ 551 "completed_at_utc": completed_at_seconds, 552 "manifest_digest": hex::encode(manifest.digest().as_bytes()), 553 "operation_id": operation_id, 554 "snapshot_generation": actual_generation, 555 }), 556 )?; 557 self.state 558 .repository() 559 .complete_admin_operation( 560 &prepared, 561 &response, 562 completed_at, 563 MycAdminOperationJournalPolicy::seven_days(), 564 ) 565 .await 566 .map_err(map_admin_journal_error)?; 567 Ok(response) 568 } 569 570 fn metrics_snapshot( 571 &self, 572 request: MycAdminRequestDocument, 573 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 574 let captured_at = SystemWallClock 575 .now_utc() 576 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))? 577 .get(); 578 let snapshot = self.status.operations_cache().snapshot(); 579 let mut metrics = serde_json::Map::new(); 580 for sample in snapshot.metrics().samples() { 581 let mut key = sample.name().as_str().to_owned(); 582 for label in sample.labels() { 583 key.push('_'); 584 key.push_str(label.value()); 585 } 586 if !safe_metric_key(&key) { 587 return Err(handler_error(MycAdminHandlerErrorKind::Internal)); 588 } 589 let value = match sample.value() { 590 MetricValue::Counter(value) => value, 591 MetricValue::Gauge(value) => u64::try_from(value) 592 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 593 }; 594 if metrics.insert(key, Value::from(value)).is_some() { 595 return Err(handler_error(MycAdminHandlerErrorKind::Internal)); 596 } 597 } 598 response( 599 request.route(), 600 json!({"captured_at_utc":captured_at,"metrics":metrics}), 601 ) 602 } 603 604 async fn connections( 605 &self, 606 request: MycAdminRequestDocument, 607 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 608 let model = request_model(&request)?; 609 let limit = model_u64(&model, "/limit") 610 .and_then(|value| u16::try_from(value).ok()) 611 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 612 let status = model 613 .pointer("/state") 614 .and_then(Value::as_str) 615 .map(|value| match value { 616 "pending" => Ok(MycConnectionStatus::Pending), 617 "approved" => Ok(MycConnectionStatus::Active), 618 "rejected" => Ok(MycConnectionStatus::Denied), 619 "revoked" => Ok(MycConnectionStatus::Expired), 620 _ => Err(handler_error(MycAdminHandlerErrorKind::Internal)), 621 }) 622 .transpose()?; 623 let cursor = model 624 .pointer("/cursor") 625 .and_then(Value::as_str) 626 .map(|value| decode_connection_cursor(value, status, &self.cursor_key)) 627 .transpose()?; 628 let snapshot = match cursor.as_ref() { 629 Some(cursor) => cursor.snapshot, 630 None => connection_time_now_for_admin()?, 631 }; 632 let before = cursor.map(|cursor| (cursor.before, cursor.id)); 633 let page = self 634 .state 635 .repository() 636 .read_admin_connection_page(limit, status, snapshot, before) 637 .await 638 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 639 let items = page 640 .items() 641 .iter() 642 .map(|record| { 643 json!({ 644 "client_public_key": record.client_public_key().as_hex(), 645 "connection_id": hex::encode(record.id().as_bytes()), 646 "created_at_utc": record.created_at().get() / 1_000, 647 "generation": record.policy_generation().get(), 648 "permissions": record.admin_permissions(), 649 "state": record.status().admin_state(), 650 "updated_at_utc": record.updated_at().get() / 1_000, 651 }) 652 }) 653 .collect::<Vec<_>>(); 654 let mut value = json!({ 655 "items": items, 656 "snapshot_generation": snapshot.get(), 657 }); 658 if let Some((before, id)) = page.next() { 659 value["next_cursor"] = Value::String(encode_connection_cursor( 660 snapshot, 661 before, 662 id, 663 status, 664 &self.cursor_key, 665 )); 666 } 667 response(request.route(), value) 668 } 669 670 async fn audit_events( 671 &self, 672 request: MycAdminRequestDocument, 673 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 674 let model = request_model(&request)?; 675 let limit = model_u64(&model, "/limit") 676 .and_then(|value| u16::try_from(value).ok()) 677 .and_then(|value| MycAuditPageLimit::new(value).ok()) 678 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 679 let from_ms = optional_seconds_as_millis(&model, "/from_utc", false)?; 680 let to_ms = optional_seconds_as_millis(&model, "/to_utc", true)?; 681 let kind = model 682 .pointer("/kind") 683 .and_then(Value::as_str) 684 .map(|value| { 685 MycAuditKind::parse(value) 686 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) 687 }) 688 .transpose()?; 689 let outcome = model 690 .pointer("/outcome") 691 .and_then(Value::as_str) 692 .map(|value| { 693 MycAuditOutcome::parse(value) 694 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) 695 }) 696 .transpose()?; 697 let query_digest = audit_query_digest(from_ms, to_ms, kind, outcome); 698 let cursor = model 699 .pointer("/cursor") 700 .and_then(Value::as_str) 701 .map(|value| decode_audit_cursor(value, query_digest, &self.cursor_key)) 702 .transpose()?; 703 let page = self 704 .state 705 .repository() 706 .read_admin_audit_page( 707 MycAdminAuditQuery::new( 708 limit, 709 cursor.as_ref().map(|value| value.snapshot), 710 cursor.as_ref().map(|value| value.before), 711 from_ms, 712 to_ms, 713 kind, 714 outcome, 715 ) 716 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 717 ) 718 .await 719 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 720 let items = page 721 .items() 722 .iter() 723 .map(|record| { 724 json!({ 725 "audit_id": format!("audit-{}", record.sequence()), 726 "correlation_id": hex::encode(record.correlation_id().as_bytes()), 727 "kind": record.kind().as_str(), 728 "occurred_at_utc": record.occurred_at().get() / 1_000, 729 "outcome": record.outcome().as_str(), 730 "reason_code": record.reason().as_str(), 731 }) 732 }) 733 .collect::<Vec<_>>(); 734 let mut value = json!({ 735 "items": items, 736 "snapshot_generation": page.snapshot_sequence(), 737 }); 738 if let Some(before) = page.next_before_sequence() { 739 value["next_cursor"] = Value::String(encode_audit_cursor( 740 page.snapshot_sequence(), 741 before, 742 query_digest, 743 &self.cursor_key, 744 )); 745 } 746 response(request.route(), value) 747 } 748 749 async fn audit_summary( 750 &self, 751 request: MycAdminRequestDocument, 752 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 753 let model = request_model(&request)?; 754 let from = model_u64(&model, "/from_utc") 755 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 756 let to = model_u64(&model, "/to_utc") 757 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 758 let from_ms = seconds_as_millis(from, false)?; 759 let to_ms = seconds_as_millis(to, true)?; 760 let counts = self 761 .state 762 .repository() 763 .read_admin_audit_summary(from_ms, to_ms) 764 .await 765 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))? 766 .iter() 767 .map(|(key, value)| (key.clone(), Value::from(*value))) 768 .collect::<serde_json::Map<_, _>>(); 769 response( 770 request.route(), 771 json!({ 772 "counts": counts, 773 "from_utc": from, 774 "to_utc": to, 775 }), 776 ) 777 } 778 779 async fn connection_mutation( 780 &self, 781 request: MycAdminRequestDocument, 782 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 783 let route = request.route(); 784 let model = request_model(&request)?; 785 let operation_id = request_operation_id(&request)?.to_owned(); 786 let connection_id = 787 digest_id_parameter(&request, "connection_id").map(MycConnectionId::from_bytes)?; 788 let generation = model_u64(&model, "/expected_generation") 789 .and_then(|value| MycConnectionPolicyGeneration::new(value).ok()) 790 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 791 let observed_at = connection_time_now_for_admin()?; 792 let correlation = admin_audit_correlation(&request); 793 let action = match route { 794 MycAdminRoute::ConnectionApprove => { 795 let permissions = model 796 .pointer("/permissions") 797 .and_then(Value::as_str) 798 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 799 let permissions = if permissions.is_empty() { 800 Vec::new() 801 } else { 802 permissions 803 .split(',') 804 .map(|value| { 805 MycConnectionPermission::parse(value) 806 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) 807 }) 808 .collect::<Result<Vec<_>, _>>()? 809 }; 810 let permissions = MycConnectionPermissionSet::new(&permissions) 811 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 812 let authorized_until = self 813 .configuration 814 .normalized() 815 .pointer("/policy/challenges/enabled") 816 .and_then(Value::as_bool) 817 .unwrap_or(false) 818 .then(|| { 819 let lifetime = configuration_integer( 820 &self.configuration, 821 "/policy/challenges/authorized_lifetime_ms", 822 )?; 823 let until = observed_at 824 .get() 825 .checked_add(lifetime) 826 .ok_or_else(|| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 827 MycConnectionTimeUnixMs::new(until) 828 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 829 }) 830 .transpose() 831 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 832 MycAdminConnectionAction::Approve { 833 permissions, 834 authorized_until, 835 } 836 } 837 MycAdminRoute::ConnectionReject => MycAdminConnectionAction::Reject, 838 MycAdminRoute::ConnectionRevoke => MycAdminConnectionAction::Revoke, 839 _ => return Err(handler_error(MycAdminHandlerErrorKind::Internal)), 840 }; 841 let completed_at = admin_operation_time(observed_at)?; 842 self.state 843 .repository() 844 .execute_database_admin_operation( 845 &request, 846 completed_at, 847 MycAdminOperationJournalPolicy::seven_days(), 848 move |transaction| { 849 Box::pin(async move { 850 let (before, after) = apply_admin_connection_mutation( 851 transaction, 852 connection_id, 853 generation, 854 observed_at, 855 correlation, 856 action, 857 ) 858 .await?; 859 response( 860 route, 861 json!({ 862 "connection_id": hex::encode(after.id().as_bytes()), 863 "current_state": after.status().admin_state(), 864 "generation": after.policy_generation().get(), 865 "operation_id": operation_id, 866 "previous_state": before.status().admin_state(), 867 }), 868 ) 869 .map_err(|_| AdminJournalOperationError::Binding) 870 }) 871 }, 872 ) 873 .await 874 .map_err(map_admin_journal_error) 875 } 876 877 async fn challenge_require( 878 &self, 879 request: MycAdminRequestDocument, 880 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 881 let route = request.route(); 882 let model = request_model(&request)?; 883 let connection_id = model 884 .pointer("/connection_id") 885 .and_then(Value::as_str) 886 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) 887 .and_then(digest_id) 888 .map(MycConnectionId::from_bytes)?; 889 let operation_id = model 890 .pointer("/request_identity") 891 .and_then(Value::as_str) 892 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) 893 .and_then(digest_id) 894 .map(MycSignerOperationId::from_persisted)?; 895 let generation = model_u64(&model, "/expected_policy_generation") 896 .and_then(|value| MycConnectionPolicyGeneration::new(value).ok()) 897 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 898 let issued_at = connection_time_now_for_admin()?; 899 let lifetime = configuration_integer( 900 &self.configuration, 901 "/policy/challenges/pending_lifetime_ms", 902 ) 903 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 904 let expires_at = MycConnectionTimeUnixMs::new( 905 issued_at 906 .get() 907 .checked_add(lifetime) 908 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?, 909 ) 910 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 911 let url = self 912 .configuration 913 .normalized() 914 .pointer("/policy/challenges/url") 915 .and_then(Value::as_str) 916 .and_then(|value| MycAuthorizationChallengeUrl::new(value).ok()) 917 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Unavailable))?; 918 let mut nonce = [0_u8; 32]; 919 SystemEntropy 920 .fill_bytes(&mut nonce) 921 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))?; 922 let challenge = MycAuthorizationChallengeRequest::new( 923 operation_id, 924 connection_id, 925 generation, 926 url, 927 MycAuthorizationChallengeNonce::from_injected_entropy(nonce), 928 issued_at, 929 expires_at, 930 ) 931 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 932 if !self 933 .state 934 .metadata() 935 .admits_authorization_challenge_request(&challenge) 936 { 937 return Err(handler_error(MycAdminHandlerErrorKind::Conflict)); 938 } 939 let rate = self 940 .state 941 .metadata() 942 .governance_rate_policy(MycRateLimitClass::ChallengeCreation); 943 let completed_at = admin_operation_time(issued_at)?; 944 self.state 945 .repository() 946 .execute_database_admin_operation( 947 &request, 948 completed_at, 949 MycAdminOperationJournalPolicy::seven_days(), 950 move |transaction| { 951 Box::pin(async move { 952 let record = 953 apply_admin_challenge_require(transaction, challenge, rate).await?; 954 response( 955 route, 956 json!({ 957 "challenge_id": hex::encode(record.id().as_bytes()), 958 "challenge_url": record.url().as_str(), 959 "connection_id": hex::encode(record.connection_id().as_bytes()), 960 "expires_at_utc": record.expires_at().get() / 1_000, 961 "issued_at_utc": record.issued_at().get() / 1_000, 962 "request_identity": hex::encode(record.operation_id().as_bytes()), 963 "state": challenge_admin_state(record.state()), 964 }), 965 ) 966 .map_err(|_| AdminJournalOperationError::Binding) 967 }) 968 }, 969 ) 970 .await 971 .map_err(map_admin_journal_error) 972 } 973 974 async fn challenge_authorize( 975 &self, 976 request: MycAdminRequestDocument, 977 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 978 let route = request.route(); 979 let model = request_model(&request)?; 980 let challenge_id = digest_id_parameter(&request, "challenge_id") 981 .map(MycAuthorizationChallengeId::from_bytes)?; 982 let generation = model_u64(&model, "/expected_generation") 983 .and_then(|value| MycConnectionPolicyGeneration::new(value).ok()) 984 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 985 let observed_at = connection_time_now_for_admin()?; 986 let completed_at = admin_operation_time(observed_at)?; 987 let operation_id = request_operation_id(&request)?.to_owned(); 988 let rate = self 989 .state 990 .metadata() 991 .governance_rate_policy(MycRateLimitClass::ChallengeAuthorization); 992 self.state 993 .repository() 994 .execute_database_admin_operation( 995 &request, 996 completed_at, 997 MycAdminOperationJournalPolicy::seven_days(), 998 move |transaction| { 999 Box::pin(async move { 1000 let record = apply_admin_challenge_authorize( 1001 transaction, 1002 challenge_id, 1003 generation, 1004 observed_at, 1005 rate, 1006 ) 1007 .await?; 1008 response( 1009 route, 1010 json!({ 1011 "challenge_id": hex::encode(record.id().as_bytes()), 1012 "generation": record.policy_generation().get(), 1013 "operation_id": operation_id, 1014 "request_identity": hex::encode(record.operation_id().as_bytes()), 1015 "state": challenge_admin_state(record.state()), 1016 }), 1017 ) 1018 .map_err(|_| AdminJournalOperationError::Binding) 1019 }) 1020 }, 1021 ) 1022 .await 1023 .map_err(map_admin_journal_error) 1024 } 1025 1026 async fn discovery_mutation( 1027 &self, 1028 request: MycAdminRequestDocument, 1029 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 1030 let route = request.route(); 1031 let model = request_model(&request)?; 1032 let expected_generation = model_u64(&model, "/expected_generation") 1033 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 1034 let actual_generation = u64::from( 1035 self.state 1036 .repository() 1037 .current_configuration_generation() 1038 .await 1039 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 1040 ); 1041 if expected_generation != actual_generation { 1042 return Err(handler_error(MycAdminHandlerErrorKind::Conflict)); 1043 } 1044 let now_seconds = SystemWallClock 1045 .now_utc() 1046 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))? 1047 .get(); 1048 let now_millis = now_seconds 1049 .checked_mul(1_000) 1050 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 1051 let prepared_at = MycAdminOperationTimeUnixMs::new(now_millis) 1052 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 1053 let prepared = match self 1054 .state 1055 .repository() 1056 .prepare_admin_operation(&request, prepared_at) 1057 .await 1058 .map_err(map_admin_journal_error)? 1059 { 1060 MycAdminOperationAdmission::ExactReplay(response) => return Ok(response), 1061 MycAdminOperationAdmission::Prepared(prepared) => prepared, 1062 }; 1063 let commit = self 1064 .prepare_discovery_commit(&request, now_seconds, now_millis) 1065 .await?; 1066 let completed_at = admin_operation_time(connection_time_now_for_admin()?)?; 1067 let operation_id = request_operation_id(&request)?.to_owned(); 1068 match route { 1069 MycAdminRoute::DiscoveryRender => { 1070 let response = response( 1071 route, 1072 json!({ 1073 "document_digests": discovery_digests(&commit), 1074 "generation": actual_generation, 1075 "operation_id": operation_id, 1076 }), 1077 )?; 1078 self.state 1079 .repository() 1080 .complete_admin_operation( 1081 &prepared, 1082 &response, 1083 completed_at, 1084 MycAdminOperationJournalPolicy::seven_days(), 1085 ) 1086 .await 1087 .map_err(map_admin_journal_error)?; 1088 Ok(response) 1089 } 1090 MycAdminRoute::DiscoveryRefresh => { 1091 let state = self 1092 .state 1093 .repository() 1094 .read_discovery_publication_state() 1095 .await 1096 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 1097 let diff = if state 1098 .as_ref() 1099 .is_some_and(|state| state.desired_generation_id() == commit.generation_id()) 1100 { 1101 "in_sync" 1102 } else { 1103 "drift" 1104 }; 1105 let response = response( 1106 route, 1107 json!({ 1108 "diff_state": diff, 1109 "generation": actual_generation, 1110 "operation_id": operation_id, 1111 "source_completion": "complete", 1112 }), 1113 )?; 1114 self.state 1115 .repository() 1116 .complete_admin_operation( 1117 &prepared, 1118 &response, 1119 completed_at, 1120 MycAdminOperationJournalPolicy::seven_days(), 1121 ) 1122 .await 1123 .map_err(map_admin_journal_error)?; 1124 Ok(response) 1125 } 1126 MycAdminRoute::DiscoveryPublish => { 1127 let delivery_policy = self.state.metadata().delivery_policies().clone(); 1128 let discovery_policy = self 1129 .state 1130 .metadata() 1131 .discovery_policies() 1132 .cloned() 1133 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Unavailable))?; 1134 self.state 1135 .repository() 1136 .complete_prepared_database_admin_operation( 1137 &prepared, 1138 completed_at, 1139 MycAdminOperationJournalPolicy::seven_days(), 1140 move |transaction| { 1141 Box::pin(async move { 1142 let admission = apply_admin_discovery_publish( 1143 transaction, 1144 &commit, 1145 &delivery_policy, 1146 &discovery_policy, 1147 ) 1148 .await?; 1149 let record = admission.record(); 1150 response( 1151 route, 1152 json!({ 1153 "exact_bytes_digest": hex::encode(commit.event_digest().as_bytes()), 1154 "generation": actual_generation, 1155 "operation_id": operation_id, 1156 "target_count": u64::try_from(record.job().targets().len()).map_err(|_| AdminJournalOperationError::Binding)?, 1157 "workflow_id": hex::encode(record.job().id().as_bytes()), 1158 }), 1159 ) 1160 .map_err(|_| AdminJournalOperationError::Binding) 1161 }) 1162 }, 1163 ) 1164 .await 1165 .map_err(map_admin_journal_error) 1166 } 1167 _ => Err(handler_error(MycAdminHandlerErrorKind::Internal)), 1168 } 1169 } 1170 1171 async fn prepare_discovery_commit( 1172 &self, 1173 request: &MycAdminRequestDocument, 1174 now_seconds: u64, 1175 now_millis: u64, 1176 ) -> Result<MycDiscoveryCommitRequest, MycAdminHandlerError> { 1177 let unsigned = prepare_discovery_signing_bytes(self.state.metadata(), now_seconds) 1178 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))?; 1179 let binding = self 1180 .configuration 1181 .provider_contract() 1182 .binding(MycProviderRole::Discovery) 1183 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Unavailable))?; 1184 let timeout = configuration_integer( 1185 &self.configuration, 1186 "/transport/ingress/subscription_deadline_ms", 1187 ) 1188 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 1189 let deadline = MycProviderDeadlineUnixMs::new( 1190 now_millis 1191 .checked_add(timeout) 1192 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?, 1193 ) 1194 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 1195 let operation = MycProviderOperation::new( 1196 binding, 1197 MycProviderOperationId::from_bytes(admin_identity( 1198 b"radroots.myc.admin.discovery.operation.v1\0", 1199 request_operation_id(request)?, 1200 )), 1201 MycProviderCorrelationId::from_bytes(admin_identity( 1202 b"radroots.myc.admin.discovery.correlation.v1\0", 1203 request.correlation_id(), 1204 )), 1205 deadline, 1206 MycProviderOperationInput::sign_event(&unsigned) 1207 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 1208 ) 1209 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 1210 let verified = self 1211 .providers 1212 .execute( 1213 operation, 1214 MycProviderResponseObservedAtUnixMs::new(now_millis) 1215 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 1216 &MycTaskCancellation::uncancelled(), 1217 ) 1218 .await 1219 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))?; 1220 let event = verified 1221 .signed_event_bytes() 1222 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; 1223 MycDiscoveryCommitRequest::new( 1224 self.state.metadata(), 1225 event, 1226 MycDeliveryTimeUnixMs::new(now_millis) 1227 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 1228 ) 1229 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) 1230 } 1231 1232 async fn discovery_desired( 1233 &self, 1234 request: MycAdminRequestDocument, 1235 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 1236 let generation = self 1237 .state 1238 .repository() 1239 .current_configuration_generation() 1240 .await 1241 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 1242 let enabled = self.state.metadata().discovery_policies().is_some(); 1243 let publication = self 1244 .state 1245 .repository() 1246 .read_discovery_publication_state() 1247 .await 1248 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 1249 let document = match publication.as_ref() { 1250 Some(publication) => self 1251 .state 1252 .repository() 1253 .read_discovery_document_for_job(publication.desired_job_id()) 1254 .await 1255 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, 1256 None => None, 1257 }; 1258 let mut digests = serde_json::Map::new(); 1259 if let Some(document) = document { 1260 digests.insert( 1261 "handler_event".to_owned(), 1262 Value::String(hex::encode(document.event_digest().as_bytes())), 1263 ); 1264 digests.insert( 1265 "nip05_projection".to_owned(), 1266 Value::String(hex::encode(document.nip05_projection_digest().as_bytes())), 1267 ); 1268 } 1269 let publication_state = if !enabled { 1270 "disabled" 1271 } else if publication.as_ref().is_some_and(|state| { 1272 state.current_generation_id() == Some(state.desired_generation_id()) 1273 }) { 1274 "delivered" 1275 } else { 1276 "pending" 1277 }; 1278 response( 1279 request.route(), 1280 json!({ 1281 "document_digests": digests, 1282 "enabled": enabled, 1283 "generation": u64::from(generation), 1284 "publication_state": publication_state, 1285 }), 1286 ) 1287 } 1288 } 1289 1290 impl MycAdminHandler for RuntimeAdminHandler { 1291 fn handle<'a>(&'a self, request: MycAdminRequestDocument) -> MycAdminFuture<'a> { 1292 let mutation = request.route().is_mutation(); 1293 Box::pin(async move { 1294 let result = self.handle_inner(request).await; 1295 if mutation && result.is_ok() { 1296 let _ = self.health.try_send(HealthEvent::StateChanged); 1297 } 1298 result 1299 }) 1300 } 1301 } 1302 1303 struct ConnectionCursor { 1304 snapshot: MycConnectionTimeUnixMs, 1305 before: MycConnectionTimeUnixMs, 1306 id: MycConnectionId, 1307 } 1308 1309 struct AuditCursor { 1310 snapshot: u64, 1311 before: u64, 1312 } 1313 1314 fn request_model(request: &MycAdminRequestDocument) -> Result<Value, MycAdminHandlerError> { 1315 serde_json::from_slice(request.model_bytes()) 1316 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) 1317 } 1318 1319 fn request_operation_id(request: &MycAdminRequestDocument) -> Result<&str, MycAdminHandlerError> { 1320 request 1321 .operation_id() 1322 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) 1323 } 1324 1325 fn model_u64(model: &Value, pointer: &str) -> Option<u64> { 1326 model.pointer(pointer).and_then(Value::as_u64) 1327 } 1328 1329 fn digest_id(value: &str) -> Result<[u8; 32], MycAdminHandlerError> { 1330 let mut bytes = [0_u8; 32]; 1331 if value.len() != 64 || hex::decode_to_slice(value, &mut bytes).is_err() { 1332 return Err(handler_error(MycAdminHandlerErrorKind::NotFound)); 1333 } 1334 Ok(bytes) 1335 } 1336 1337 fn digest_id_parameter( 1338 request: &MycAdminRequestDocument, 1339 name: &str, 1340 ) -> Result<[u8; 32], MycAdminHandlerError> { 1341 request 1342 .parameter(name) 1343 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::NotFound)) 1344 .and_then(digest_id) 1345 } 1346 1347 fn connection_time_now_for_admin() -> Result<MycConnectionTimeUnixMs, MycAdminHandlerError> { 1348 MycConnectionTimeUnixMs::new( 1349 SystemWallClock 1350 .now_utc() 1351 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))? 1352 .get() 1353 .checked_mul(1_000) 1354 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?, 1355 ) 1356 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) 1357 } 1358 1359 fn admin_operation_time( 1360 time: MycConnectionTimeUnixMs, 1361 ) -> Result<MycAdminOperationTimeUnixMs, MycAdminHandlerError> { 1362 MycAdminOperationTimeUnixMs::new(time.get()) 1363 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) 1364 } 1365 1366 fn admin_identity(domain: &[u8], value: &str) -> [u8; 32] { 1367 let mut hasher = Sha256::new(); 1368 hasher.update(domain); 1369 hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes()); 1370 hasher.update(value.as_bytes()); 1371 hasher.finalize().into() 1372 } 1373 1374 fn admin_audit_correlation(request: &MycAdminRequestDocument) -> MycAuditCorrelationId { 1375 MycAuditCorrelationId::new(admin_identity( 1376 b"radroots.myc.admin.audit_correlation.v1\0", 1377 request.correlation_id(), 1378 )) 1379 } 1380 1381 const fn challenge_admin_state(state: MycAuthorizationChallengeState) -> &'static str { 1382 match state { 1383 MycAuthorizationChallengeState::Pending => "pending", 1384 MycAuthorizationChallengeState::Authorized => "authorized", 1385 MycAuthorizationChallengeState::Expired => "expired", 1386 } 1387 } 1388 1389 fn discovery_digests(request: &MycDiscoveryCommitRequest) -> Value { 1390 json!({ 1391 "desired": hex::encode(request.desired_digest().as_bytes()), 1392 "handler_event": hex::encode(request.event_digest().as_bytes()), 1393 "nip05_projection": hex::encode(request.projection_digest().as_bytes()), 1394 }) 1395 } 1396 1397 fn seconds_as_millis(seconds: u64, include_full_second: bool) -> Result<u64, MycAdminHandlerError> { 1398 seconds 1399 .checked_mul(1_000) 1400 .and_then(|value| { 1401 if include_full_second { 1402 value.checked_add(999) 1403 } else { 1404 Some(value) 1405 } 1406 }) 1407 .filter(|value| i64::try_from(*value).is_ok()) 1408 .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) 1409 } 1410 1411 fn optional_seconds_as_millis( 1412 model: &Value, 1413 pointer: &str, 1414 include_full_second: bool, 1415 ) -> Result<Option<u64>, MycAdminHandlerError> { 1416 model_u64(model, pointer) 1417 .map(|value| seconds_as_millis(value, include_full_second)) 1418 .transpose() 1419 } 1420 1421 const fn connection_state_code(status: Option<MycConnectionStatus>) -> u8 { 1422 match status { 1423 None => 0, 1424 Some(MycConnectionStatus::Pending) => 1, 1425 Some(MycConnectionStatus::Active) => 2, 1426 Some(MycConnectionStatus::Denied) => 3, 1427 Some(MycConnectionStatus::Expired) => 4, 1428 } 1429 } 1430 1431 fn encode_connection_cursor( 1432 snapshot: MycConnectionTimeUnixMs, 1433 before: MycConnectionTimeUnixMs, 1434 id: MycConnectionId, 1435 status: Option<MycConnectionStatus>, 1436 key: &[u8; 32], 1437 ) -> String { 1438 let mut payload = Vec::with_capacity(83); 1439 payload.extend_from_slice(&[1, 1]); 1440 payload.extend_from_slice(&snapshot.get().to_be_bytes()); 1441 payload.extend_from_slice(&before.get().to_be_bytes()); 1442 payload.extend_from_slice(id.as_bytes()); 1443 payload.push(connection_state_code(status)); 1444 payload.extend_from_slice(&cursor_mac(key, &payload)); 1445 URL_SAFE_NO_PAD.encode(payload) 1446 } 1447 1448 fn decode_connection_cursor( 1449 encoded: &str, 1450 status: Option<MycConnectionStatus>, 1451 key: &[u8; 32], 1452 ) -> Result<ConnectionCursor, MycAdminHandlerError> { 1453 let bytes = URL_SAFE_NO_PAD 1454 .decode(encoded) 1455 .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?; 1456 if bytes.len() != 83 1457 || bytes[0..2] != [1, 1] 1458 || bytes[50] != connection_state_code(status) 1459 || !constant_time_equal(&bytes[51..], &cursor_mac(key, &bytes[..51])) 1460 { 1461 return Err(handler_error(MycAdminHandlerErrorKind::InvalidCursor)); 1462 } 1463 let snapshot = MycConnectionTimeUnixMs::new(u64::from_be_bytes( 1464 bytes[2..10] 1465 .try_into() 1466 .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, 1467 )) 1468 .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?; 1469 let before = MycConnectionTimeUnixMs::new(u64::from_be_bytes( 1470 bytes[10..18] 1471 .try_into() 1472 .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, 1473 )) 1474 .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?; 1475 let id = MycConnectionId::from_bytes( 1476 bytes[18..50] 1477 .try_into() 1478 .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, 1479 ); 1480 Ok(ConnectionCursor { 1481 snapshot, 1482 before, 1483 id, 1484 }) 1485 } 1486 1487 fn audit_query_digest( 1488 from: Option<u64>, 1489 to: Option<u64>, 1490 kind: Option<MycAuditKind>, 1491 outcome: Option<MycAuditOutcome>, 1492 ) -> [u8; 32] { 1493 let mut hasher = Sha256::new(); 1494 hasher.update(b"radroots.myc.admin.audit_query.v1\0"); 1495 hasher.update(from.unwrap_or(u64::MAX).to_be_bytes()); 1496 hasher.update(to.unwrap_or(u64::MAX).to_be_bytes()); 1497 hasher.update( 1498 kind.map(MycAuditKind::as_str) 1499 .unwrap_or_default() 1500 .as_bytes(), 1501 ); 1502 hasher.update([0]); 1503 hasher.update( 1504 outcome 1505 .map(MycAuditOutcome::as_str) 1506 .unwrap_or_default() 1507 .as_bytes(), 1508 ); 1509 hasher.finalize().into() 1510 } 1511 1512 fn encode_audit_cursor( 1513 snapshot: u64, 1514 before: u64, 1515 query_digest: [u8; 32], 1516 key: &[u8; 32], 1517 ) -> String { 1518 let mut payload = Vec::with_capacity(82); 1519 payload.extend_from_slice(&[1, 2]); 1520 payload.extend_from_slice(&snapshot.to_be_bytes()); 1521 payload.extend_from_slice(&before.to_be_bytes()); 1522 payload.extend_from_slice(&query_digest); 1523 payload.extend_from_slice(&cursor_mac(key, &payload)); 1524 URL_SAFE_NO_PAD.encode(payload) 1525 } 1526 1527 fn decode_audit_cursor( 1528 encoded: &str, 1529 query_digest: [u8; 32], 1530 key: &[u8; 32], 1531 ) -> Result<AuditCursor, MycAdminHandlerError> { 1532 let bytes = URL_SAFE_NO_PAD 1533 .decode(encoded) 1534 .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?; 1535 if bytes.len() != 82 1536 || bytes[0..2] != [1, 2] 1537 || bytes[18..50] != query_digest 1538 || !constant_time_equal(&bytes[50..], &cursor_mac(key, &bytes[..50])) 1539 { 1540 return Err(handler_error(MycAdminHandlerErrorKind::InvalidCursor)); 1541 } 1542 Ok(AuditCursor { 1543 snapshot: u64::from_be_bytes( 1544 bytes[2..10] 1545 .try_into() 1546 .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, 1547 ), 1548 before: u64::from_be_bytes( 1549 bytes[10..18] 1550 .try_into() 1551 .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, 1552 ), 1553 }) 1554 } 1555 1556 fn cursor_mac(key: &[u8; 32], payload: &[u8]) -> [u8; 32] { 1557 let mut hasher = Sha256::new(); 1558 hasher.update(b"radroots.myc.admin.cursor.v1\0"); 1559 hasher.update(key); 1560 hasher.update( 1561 u64::try_from(payload.len()) 1562 .unwrap_or(u64::MAX) 1563 .to_be_bytes(), 1564 ); 1565 hasher.update(payload); 1566 hasher.finalize().into() 1567 } 1568 1569 fn constant_time_equal(left: &[u8], right: &[u8]) -> bool { 1570 left.len() == right.len() 1571 && left 1572 .iter() 1573 .zip(right) 1574 .fold(0_u8, |difference, (left, right)| { 1575 difference | (left ^ right) 1576 }) 1577 == 0 1578 } 1579 1580 fn safe_metric_key(value: &str) -> bool { 1581 (1..=64).contains(&value.len()) 1582 && value.as_bytes().first().is_some_and(u8::is_ascii_lowercase) 1583 && value 1584 .bytes() 1585 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_') 1586 } 1587 1588 struct RuntimeShutdownHandler { 1589 accepting_mutations: Arc<AtomicBool>, 1590 state: Arc<MycStateHost>, 1591 status: MycStatusPublisher, 1592 status_context: RuntimeStatusContext, 1593 health: RuntimeHealth, 1594 } 1595 1596 impl RuntimeShutdownHandler { 1597 fn publish(&mut self, phase: MycServicePhase) -> Result<(), HostError> { 1598 let observation = self 1599 .status_context 1600 .observation(phase, &self.health) 1601 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 1602 self.status 1603 .publish(observation) 1604 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error)) 1605 } 1606 } 1607 1608 impl ShutdownPhaseHandler for RuntimeShutdownHandler { 1609 fn enter(&mut self, phase: ShutdownPhase) -> ShutdownPhaseFuture<'_> { 1610 Box::pin(async move { 1611 match phase { 1612 ShutdownPhase::RejectNewMutations => { 1613 self.accepting_mutations.store(false, Ordering::Release); 1614 self.publish(MycServicePhase::Stopping)?; 1615 } 1616 ShutdownPhase::PersistRecoverableWork => { 1617 self.state 1618 .repository() 1619 .verify_delivery_invariants() 1620 .await 1621 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 1622 } 1623 ShutdownPhase::CloseSqlite => { 1624 self.state 1625 .close() 1626 .await 1627 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 1628 } 1629 ShutdownPhase::CancelIngress 1630 | ShutdownPhase::DrainOperations 1631 | ShutdownPhase::CloseNetwork 1632 | ShutdownPhase::CloseSockets => {} 1633 } 1634 Ok(()) 1635 }) 1636 } 1637 } 1638 1639 pub(crate) async fn run_myc_daemon<S>( 1640 runtime: MycRuntimeContext, 1641 configuration: MycConfigDocumentV1, 1642 applied_at: MigrationAppliedAtUnixSeconds, 1643 build: &MigrationBuildIdentity, 1644 signals: S, 1645 ) -> MycProcessResult 1646 where 1647 S: MycProcessSignalSource + 'static, 1648 { 1649 run_myc_daemon_inner(runtime, configuration, applied_at, build, signals) 1650 .await 1651 .unwrap_or_else(|error| error.process_result()) 1652 } 1653 1654 async fn run_myc_daemon_inner<S>( 1655 runtime: MycRuntimeContext, 1656 configuration: MycConfigDocumentV1, 1657 applied_at: MigrationAppliedAtUnixSeconds, 1658 build: &MigrationBuildIdentity, 1659 signals: S, 1660 ) -> Result<MycProcessResult, MycDaemonError> 1661 where 1662 S: MycProcessSignalSource + 'static, 1663 { 1664 let configuration = Arc::new(configuration); 1665 let state = Arc::new( 1666 open_myc_state_read_write_from_config(&runtime, &configuration, applied_at, build) 1667 .await 1668 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?, 1669 ); 1670 let now_millis = wall_time_millis()?; 1671 let mut recovery_entropy = [0_u8; 32]; 1672 SystemEntropy 1673 .fill_bytes(&mut recovery_entropy) 1674 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 1675 state 1676 .repository() 1677 .recover_delivery_state( 1678 MycDeliveryTimeUnixMs::new(now_millis) 1679 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, 1680 MycDeliveryRecoveryEntropy::from_injected_entropy(recovery_entropy), 1681 ) 1682 .await 1683 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?; 1684 state 1685 .repository() 1686 .verify_delivery_invariants() 1687 .await 1688 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?; 1689 1690 let startup_cancellation = MycTaskCancellation::uncancelled(); 1691 let providers = Arc::new( 1692 MycProviderExecutor::open(&runtime, &configuration, &startup_cancellation) 1693 .await 1694 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Provider))?, 1695 ); 1696 let mut provider_seed = [0_u8; 32]; 1697 SystemEntropy 1698 .fill_bytes(&mut provider_seed) 1699 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 1700 providers 1701 .probe_all(now_millis, provider_seed, &startup_cancellation) 1702 .await 1703 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Provider))?; 1704 let ingress_adapter = Arc::new( 1705 MycNostrIngressAdapter::from_configuration(&configuration) 1706 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Relay))?, 1707 ); 1708 let ingress_slots = open_initial_subscriptions(&ingress_adapter, &configuration).await?; 1709 1710 let generation = u64::from( 1711 state 1712 .repository() 1713 .current_configuration_generation() 1714 .await 1715 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?, 1716 ); 1717 let mut status_context = RuntimeStatusContext { 1718 runtime: runtime.clone(), 1719 configuration: Arc::clone(&configuration), 1720 state: Arc::clone(&state), 1721 clock: SystemMonotonicClock::new(), 1722 generation, 1723 connection_counts: MycConnectionCountsV1::default(), 1724 outbox: MycOutboxStatusV1::default(), 1725 }; 1726 status_context.refresh_state().await?; 1727 let operations_enabled = configuration 1728 .normalized() 1729 .pointer("/operations/enabled") 1730 .and_then(Value::as_bool) 1731 .unwrap_or(false); 1732 let health = RuntimeHealth::ready(true); 1733 let (publisher, status_reader) = myc_status_cache( 1734 runtime.context().instance().clone(), 1735 status_context.observation(MycServicePhase::Starting, &health)?, 1736 ) 1737 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 1738 let accepting_mutations = Arc::new(AtomicBool::new(true)); 1739 let mut cursor_key = [0_u8; 32]; 1740 SystemEntropy 1741 .fill_bytes(&mut cursor_key) 1742 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 1743 let (health_tx, mut health_rx) = mpsc::channel(8); 1744 let admin_handler = Arc::new(RuntimeAdminHandler { 1745 state: Arc::clone(&state), 1746 configuration: Arc::clone(&configuration), 1747 providers: Arc::clone(&providers), 1748 status: status_reader.clone(), 1749 accepting_mutations: Arc::clone(&accepting_mutations), 1750 cursor_key, 1751 health: health_tx.clone(), 1752 }); 1753 let admin = MycAdminServer::new(&configuration, admin_handler) 1754 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Admin))? 1755 .bind(&runtime) 1756 .await 1757 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Admin))?; 1758 let operations = if operations_enabled { 1759 Some( 1760 MycOperationsServer::new(&configuration, &status_reader) 1761 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Operations))? 1762 .bind() 1763 .await 1764 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Operations))?, 1765 ) 1766 } else { 1767 None 1768 }; 1769 1770 let mut shutdown_handler = RuntimeShutdownHandler { 1771 accepting_mutations, 1772 state: Arc::clone(&state), 1773 status: publisher, 1774 status_context, 1775 health, 1776 }; 1777 let ingress_capacity = configuration_usize(&configuration, "/resource_limits/queues/ingress")?; 1778 let provider_capacity = 1779 configuration_usize(&configuration, "/resource_limits/queues/provider")?; 1780 let (ingress_tx, ingress_rx) = mpsc::channel(ingress_capacity); 1781 let mut supervisor = TaskSupervisor::new(); 1782 spawn_admin(&mut supervisor, admin)?; 1783 if let Some(operations) = operations { 1784 spawn_operations(&mut supervisor, operations)?; 1785 } 1786 spawn_relay_ingress( 1787 &mut supervisor, 1788 ingress_adapter, 1789 ingress_slots, 1790 Arc::clone(&configuration), 1791 health_tx.clone(), 1792 ingress_tx, 1793 )?; 1794 spawn_provider_dispatch( 1795 &mut supervisor, 1796 Arc::clone(&providers), 1797 Arc::clone(&configuration), 1798 Arc::clone(&state), 1799 health_tx.clone(), 1800 ingress_rx, 1801 provider_capacity, 1802 )?; 1803 spawn_delivery_outbox( 1804 &mut supervisor, 1805 Arc::clone(&state), 1806 Arc::clone(&configuration), 1807 health_tx, 1808 )?; 1809 shutdown_handler 1810 .publish(MycServicePhase::Ready) 1811 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 1812 1813 let mut signals = ProcessSignalAdapter::new(HostSignalSource::new(signals)); 1814 let signal_initiated = loop { 1815 tokio::select! { 1816 action = signals.next_action() => { 1817 match action { 1818 Ok(action) if !action.forces_termination() => break true, 1819 Ok(_) | Err(_) => break false, 1820 } 1821 } 1822 joined = supervisor.join_next() => { 1823 match joined { 1824 Some(Ok(exit)) if exit.status() == SupervisedTaskExitStatus::OptionalFailure 1825 || exit.metadata().name().as_str() == TASK_OPERATIONS_SERVER => { 1826 shutdown_handler.health.operations = false; 1827 shutdown_handler.publish(MycServicePhase::Degraded) 1828 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 1829 } 1830 Some(Ok(_)) | Some(Err(_)) | None => break false, 1831 } 1832 } 1833 event = health_rx.recv() => { 1834 let Some(event) = event else { break false; }; 1835 if shutdown_handler.health.observe(event) { 1836 if event == HealthEvent::StateChanged { 1837 shutdown_handler.status_context.refresh_state().await?; 1838 } 1839 let phase = if shutdown_handler.health.providers 1840 && shutdown_handler.health.required_relays { 1841 if shutdown_handler.health.operations { 1842 MycServicePhase::Ready 1843 } else { 1844 MycServicePhase::Degraded 1845 } 1846 } else { 1847 MycServicePhase::Unready 1848 }; 1849 shutdown_handler.publish(phase) 1850 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 1851 } 1852 } 1853 } 1854 }; 1855 let grace = Duration::from_millis(configuration_integer( 1856 &configuration, 1857 "/service/shutdown_grace_ms", 1858 )?); 1859 let mut shutdown = GracefulShutdown::new(grace) 1860 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 1861 let clock = SystemMonotonicClock::new(); 1862 let summary = if signal_initiated { 1863 shutdown 1864 .run(&clock, &mut supervisor, &mut shutdown_handler, async { 1865 let _ = signals.next_action().await; 1866 }) 1867 .await 1868 } else { 1869 shutdown 1870 .run( 1871 &clock, 1872 &mut supervisor, 1873 &mut shutdown_handler, 1874 pending::<()>(), 1875 ) 1876 .await 1877 } 1878 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 1879 if signal_initiated && summary.disposition() == ShutdownDisposition::Completed { 1880 Ok(MycProcessResult::Success) 1881 } else { 1882 Err(MycDaemonError::new(MycDaemonErrorKind::Runtime)) 1883 } 1884 } 1885 1886 fn spawn_admin( 1887 supervisor: &mut TaskSupervisor, 1888 server: MycBoundAdminServer, 1889 ) -> Result<(), MycDaemonError> { 1890 supervisor 1891 .spawn(task_metadata(TASK_ADMIN_SERVER, TaskClassification::Critical, ShutdownPhase::CloseSockets)?, move |cancellation| async move { 1892 let token = MycAdminCancellationToken::new(); 1893 let serve_token = token.clone(); 1894 let serve = server.serve(serve_token); 1895 tokio::pin!(serve); 1896 tokio::select! { 1897 result = serve.as_mut() => result.map_err(|error| HostError::with_source(HostErrorKind::AdminTransport, error)), 1898 () = cancellation.cancelled() => { 1899 token.cancel(); 1900 serve.await.map_err(|error| HostError::with_source(HostErrorKind::AdminTransport, error)) 1901 } 1902 } 1903 }) 1904 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 1905 } 1906 1907 fn spawn_operations( 1908 supervisor: &mut TaskSupervisor, 1909 server: MycBoundOperationsServer, 1910 ) -> Result<(), MycDaemonError> { 1911 supervisor 1912 .spawn(task_metadata(TASK_OPERATIONS_SERVER, TaskClassification::Optional, ShutdownPhase::CloseSockets)?, move |cancellation| async move { 1913 let token = MycOperationsCancellationToken::new(); 1914 let serve_token = token.clone(); 1915 let serve = server.serve(serve_token); 1916 tokio::pin!(serve); 1917 tokio::select! { 1918 result = serve.as_mut() => result.map_err(|error| HostError::with_source(HostErrorKind::OperationsServe, error)), 1919 () = cancellation.cancelled() => { 1920 token.cancel(); 1921 serve.await.map_err(|error| HostError::with_source(HostErrorKind::OperationsServe, error)) 1922 } 1923 } 1924 }) 1925 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 1926 } 1927 1928 fn spawn_relay_ingress( 1929 supervisor: &mut TaskSupervisor, 1930 adapter: Arc<MycNostrIngressAdapter>, 1931 slots: Vec<RuntimeIngressSlot>, 1932 configuration: Arc<MycConfigDocumentV1>, 1933 health: mpsc::Sender<HealthEvent>, 1934 ingress: mpsc::Sender<RuntimeIngressItem>, 1935 ) -> Result<(), MycDaemonError> { 1936 supervisor 1937 .spawn( 1938 task_metadata( 1939 TASK_RELAY_INGRESS, 1940 TaskClassification::Critical, 1941 ShutdownPhase::CancelIngress, 1942 )?, 1943 move |cancellation| async move { 1944 run_relay_ingress(cancellation, adapter, slots, configuration, health, ingress) 1945 .await 1946 }, 1947 ) 1948 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 1949 } 1950 1951 fn spawn_provider_dispatch( 1952 supervisor: &mut TaskSupervisor, 1953 providers: Arc<MycProviderExecutor>, 1954 configuration: Arc<MycConfigDocumentV1>, 1955 state: Arc<MycStateHost>, 1956 health: mpsc::Sender<HealthEvent>, 1957 ingress: mpsc::Receiver<RuntimeIngressItem>, 1958 provider_capacity: usize, 1959 ) -> Result<(), MycDaemonError> { 1960 supervisor 1961 .spawn( 1962 task_metadata( 1963 TASK_PROVIDER_DISPATCH, 1964 TaskClassification::Critical, 1965 ShutdownPhase::DrainOperations, 1966 )?, 1967 move |cancellation| async move { 1968 run_provider_dispatch( 1969 cancellation, 1970 configuration, 1971 state, 1972 providers, 1973 health, 1974 ingress, 1975 provider_capacity, 1976 ) 1977 .await 1978 }, 1979 ) 1980 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 1981 } 1982 1983 async fn open_initial_subscriptions( 1984 adapter: &Arc<MycNostrIngressAdapter>, 1985 configuration: &MycConfigDocumentV1, 1986 ) -> Result<Vec<RuntimeIngressSlot>, MycDaemonError> { 1987 let initial = 1988 configuration_integer(configuration, "/transport/publish_retry/initial_backoff_ms")?; 1989 let mut slots = Vec::with_capacity(adapter.group_count()); 1990 for index in 0..adapter.group_count() { 1991 let required = adapter 1992 .group_is_required(index) 1993 .ok_or_else(|| MycDaemonError::new(MycDaemonErrorKind::Relay))?; 1994 let duration = 1995 configuration_integer(configuration, "/transport/ingress/subscription_deadline_ms")?; 1996 let result = subscribe_ingress(Arc::clone(adapter), index, Vec::new(), duration).await; 1997 let state = match result { 1998 Ok(subscription) => RuntimeIngressSlotState::Active(subscription), 1999 Err(_) if required => return Err(MycDaemonError::new(MycDaemonErrorKind::Relay)), 2000 Err(_) => RuntimeIngressSlotState::Waiting, 2001 }; 2002 slots.push(RuntimeIngressSlot { 2003 index, 2004 required, 2005 checkpoints: Vec::new(), 2006 state, 2007 retry_delay_ms: initial, 2008 }); 2009 } 2010 if slots.is_empty() 2011 || slots 2012 .iter() 2013 .any(|slot| slot.required && !matches!(slot.state, RuntimeIngressSlotState::Active(_))) 2014 { 2015 return Err(MycDaemonError::new(MycDaemonErrorKind::Relay)); 2016 } 2017 Ok(slots) 2018 } 2019 2020 async fn subscribe_ingress( 2021 adapter: Arc<MycNostrIngressAdapter>, 2022 index: usize, 2023 checkpoints: Vec<SubscriptionCheckpoint>, 2024 subscription_duration_ms: u64, 2025 ) -> Result<BoxSubscription, crate::transport_nostr_adapter::MycRelayAdapterError> { 2026 let deadline = wall_time_millis() 2027 .ok() 2028 .and_then(|now| now.checked_add(subscription_duration_ms)) 2029 .ok_or_else(|| { 2030 crate::transport_nostr_adapter::runtime_relay_adapter_error( 2031 crate::transport_nostr_adapter::MycRelayAdapterErrorKind::Configuration, 2032 ) 2033 })?; 2034 adapter 2035 .subscribe(index, 1_000, deadline, &checkpoints) 2036 .await 2037 } 2038 2039 async fn run_relay_ingress( 2040 cancellation: radroots_service_host::CancellationToken, 2041 adapter: Arc<MycNostrIngressAdapter>, 2042 mut slots: Vec<RuntimeIngressSlot>, 2043 configuration: Arc<MycConfigDocumentV1>, 2044 health: mpsc::Sender<HealthEvent>, 2045 ingress: mpsc::Sender<RuntimeIngressItem>, 2046 ) -> Result<(), HostError> { 2047 if slots.is_empty() || slots.len() > 2 { 2048 return Err(HostError::new(HostErrorKind::TaskFailure)); 2049 } 2050 loop { 2051 if cancellation.is_cancelled() { 2052 cancel_ingress_slots(&mut slots).await; 2053 return Ok(()); 2054 } 2055 let selected = if slots.len() == 1 { 2056 tokio::select! { 2057 () = cancellation.cancelled() => { 2058 cancel_ingress_slots(&mut slots).await; 2059 return Ok(()); 2060 } 2061 action = poll_ingress_slot(&mut slots[0]) => (0, action), 2062 } 2063 } else { 2064 let (left, right) = slots.split_at_mut(1); 2065 tokio::select! { 2066 () = cancellation.cancelled() => { 2067 cancel_ingress_slots(&mut slots).await; 2068 return Ok(()); 2069 } 2070 action = poll_ingress_slot(&mut left[0]) => (0, action), 2071 action = poll_ingress_slot(&mut right[0]) => (1, action), 2072 } 2073 }; 2074 handle_ingress_action( 2075 selected.0, 2076 selected.1, 2077 &adapter, 2078 &configuration, 2079 &mut slots, 2080 &health, 2081 &ingress, 2082 &cancellation, 2083 ) 2084 .await?; 2085 } 2086 } 2087 2088 async fn poll_ingress_slot(slot: &mut RuntimeIngressSlot) -> RuntimeIngressAction { 2089 match &mut slot.state { 2090 RuntimeIngressSlotState::Active(subscription) => { 2091 RuntimeIngressAction::Next(subscription.next().await) 2092 } 2093 RuntimeIngressSlotState::Waiting => { 2094 tokio::time::sleep(Duration::from_millis(slot.retry_delay_ms)).await; 2095 RuntimeIngressAction::Retry 2096 } 2097 RuntimeIngressSlotState::Connecting(connecting) => { 2098 RuntimeIngressAction::Connected(connecting.await) 2099 } 2100 } 2101 } 2102 2103 #[allow(clippy::too_many_arguments)] 2104 async fn handle_ingress_action( 2105 selected: usize, 2106 action: RuntimeIngressAction, 2107 adapter: &Arc<MycNostrIngressAdapter>, 2108 configuration: &MycConfigDocumentV1, 2109 slots: &mut [RuntimeIngressSlot], 2110 health: &mpsc::Sender<HealthEvent>, 2111 ingress: &mpsc::Sender<RuntimeIngressItem>, 2112 cancellation: &radroots_service_host::CancellationToken, 2113 ) -> Result<(), HostError> { 2114 let initial = 2115 configuration_integer(configuration, "/transport/publish_retry/initial_backoff_ms") 2116 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 2117 let maximum = 2118 configuration_integer(configuration, "/transport/publish_retry/maximum_backoff_ms") 2119 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 2120 match action { 2121 RuntimeIngressAction::Next(Ok(SubscriptionNext::Event(event))) => { 2122 let observed = event.observed(); 2123 let relay_id = adapter 2124 .relay_id(selected, observed.provenance().target()) 2125 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 2126 replace_checkpoint(&mut slots[selected].checkpoints, event.checkpoint().clone()); 2127 let item = RuntimeIngressItem { 2128 raw_event: Box::from(observed.event().raw_json().as_bytes()), 2129 relay_id, 2130 observed_at_unix_ms: observed.provenance().observed_at_unix_ms(), 2131 }; 2132 tokio::select! { 2133 result = ingress.send(item) => result.map_err(|_| HostError::new(HostErrorKind::TaskFailure))?, 2134 () = cancellation.cancelled() => return Ok(()), 2135 } 2136 } 2137 RuntimeIngressAction::Next(Ok(SubscriptionNext::End(end))) => { 2138 slots[selected].checkpoints = end.checkpoints().to_vec(); 2139 slots[selected].state = RuntimeIngressSlotState::Waiting; 2140 send_required_relay_health(required_relays_ready(slots), health, cancellation).await?; 2141 } 2142 RuntimeIngressAction::Next(Err(_)) => { 2143 slots[selected].state = RuntimeIngressSlotState::Waiting; 2144 send_required_relay_health(required_relays_ready(slots), health, cancellation).await?; 2145 } 2146 RuntimeIngressAction::Retry => { 2147 let adapter = Arc::clone(adapter); 2148 let index = slots[selected].index; 2149 let checkpoints = slots[selected].checkpoints.clone(); 2150 let duration = 2151 configuration_integer(configuration, "/transport/ingress/subscription_deadline_ms") 2152 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 2153 slots[selected].state = RuntimeIngressSlotState::Connecting(Box::pin(async move { 2154 subscribe_ingress(adapter, index, checkpoints, duration).await 2155 })); 2156 } 2157 RuntimeIngressAction::Connected(Ok(subscription)) => { 2158 slots[selected].state = RuntimeIngressSlotState::Active(subscription); 2159 slots[selected].retry_delay_ms = initial; 2160 send_required_relay_health(required_relays_ready(slots), health, cancellation).await?; 2161 } 2162 RuntimeIngressAction::Connected(Err(_)) => { 2163 slots[selected].state = RuntimeIngressSlotState::Waiting; 2164 slots[selected].retry_delay_ms = slots[selected] 2165 .retry_delay_ms 2166 .saturating_mul(2) 2167 .min(maximum); 2168 send_required_relay_health(required_relays_ready(slots), health, cancellation).await?; 2169 } 2170 } 2171 Ok(()) 2172 } 2173 2174 fn replace_checkpoint( 2175 checkpoints: &mut Vec<SubscriptionCheckpoint>, 2176 checkpoint: SubscriptionCheckpoint, 2177 ) { 2178 if let Some(existing) = checkpoints 2179 .iter_mut() 2180 .find(|existing| existing.target() == checkpoint.target()) 2181 { 2182 *existing = checkpoint; 2183 } else { 2184 checkpoints.push(checkpoint); 2185 } 2186 } 2187 2188 fn required_relays_ready(slots: &[RuntimeIngressSlot]) -> bool { 2189 slots 2190 .iter() 2191 .all(|slot| !slot.required || matches!(slot.state, RuntimeIngressSlotState::Active(_))) 2192 } 2193 2194 async fn send_required_relay_health( 2195 ready: bool, 2196 health: &mpsc::Sender<HealthEvent>, 2197 cancellation: &radroots_service_host::CancellationToken, 2198 ) -> Result<(), HostError> { 2199 tokio::select! { 2200 result = health.send(HealthEvent::RequiredRelays(ready)) => { 2201 result.map_err(|_| HostError::new(HostErrorKind::TaskFailure)) 2202 } 2203 () = cancellation.cancelled() => Ok(()), 2204 } 2205 } 2206 2207 async fn cancel_ingress_slots(slots: &mut [RuntimeIngressSlot]) { 2208 for slot in slots { 2209 if let RuntimeIngressSlotState::Active(subscription) = &mut slot.state { 2210 let _ = subscription.cancel().await; 2211 } 2212 } 2213 } 2214 2215 async fn run_provider_dispatch( 2216 cancellation: radroots_service_host::CancellationToken, 2217 configuration: Arc<MycConfigDocumentV1>, 2218 state: Arc<MycStateHost>, 2219 providers: Arc<MycProviderExecutor>, 2220 health: mpsc::Sender<HealthEvent>, 2221 mut ingress: mpsc::Receiver<RuntimeIngressItem>, 2222 provider_capacity: usize, 2223 ) -> Result<(), HostError> { 2224 if provider_capacity == 0 { 2225 return Err(HostError::new(HostErrorKind::TaskFailure)); 2226 } 2227 let coordinator = MycRuntimeNip46Coordinator::new(configuration.clone(), state, providers) 2228 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 2229 let task_cancellation = MycTaskCancellation::from_host(cancellation.clone()); 2230 let initial = configuration_integer( 2231 &configuration, 2232 "/transport/publish_retry/initial_backoff_ms", 2233 ) 2234 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 2235 let maximum = configuration_integer( 2236 &configuration, 2237 "/transport/publish_retry/maximum_backoff_ms", 2238 ) 2239 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 2240 let mut queue = VecDeque::with_capacity(provider_capacity); 2241 loop { 2242 if queue.is_empty() { 2243 tokio::select! { 2244 item = ingress.recv() => match item { 2245 Some(item) => queue.push_back(item), 2246 None => return Err(HostError::new(HostErrorKind::TaskFailure)), 2247 }, 2248 () = cancellation.cancelled() => return Ok(()), 2249 } 2250 } 2251 while queue.len() < provider_capacity { 2252 match ingress.try_recv() { 2253 Ok(item) => queue.push_back(item), 2254 Err(mpsc::error::TryRecvError::Empty) => break, 2255 Err(mpsc::error::TryRecvError::Disconnected) if queue.is_empty() => { 2256 return Err(HostError::new(HostErrorKind::TaskFailure)); 2257 } 2258 Err(mpsc::error::TryRecvError::Disconnected) => break, 2259 } 2260 } 2261 let item = queue.pop_front().expect("provider queue is nonempty"); 2262 let admission_evidence = runtime_nip46_admission_evidence() 2263 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 2264 let mut retry = initial; 2265 loop { 2266 if cancellation.is_cancelled() { 2267 return Ok(()); 2268 } 2269 match coordinator 2270 .process( 2271 &item.raw_event, 2272 item.relay_id.clone(), 2273 item.observed_at_unix_ms, 2274 admission_evidence, 2275 &task_cancellation, 2276 ) 2277 .await 2278 { 2279 Ok(_) => { 2280 let _ = health.try_send(HealthEvent::StateChanged); 2281 send_provider_health(&health, true, &cancellation).await?; 2282 break; 2283 } 2284 Err(error) if error.kind() == MycNip46DispatchErrorKind::Provider => { 2285 send_provider_health(&health, false, &cancellation).await?; 2286 tokio::select! { 2287 () = cancellation.cancelled() => return Ok(()), 2288 () = tokio::time::sleep(Duration::from_millis(retry)) => {} 2289 } 2290 retry = retry.saturating_mul(2).min(maximum); 2291 } 2292 Err(error) => { 2293 return Err(HostError::with_source(HostErrorKind::TaskFailure, error)); 2294 } 2295 } 2296 } 2297 } 2298 } 2299 2300 fn runtime_nip46_admission_evidence() -> Result<MycRuntimeNip46AdmissionEvidence, MycDaemonError> { 2301 let mut request_nonce = [0_u8; 32]; 2302 SystemEntropy 2303 .fill_bytes(&mut request_nonce) 2304 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; 2305 MycRuntimeNip46AdmissionEvidence::new(request_nonce, wall_time_millis()?) 2306 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 2307 } 2308 2309 async fn send_provider_health( 2310 health: &mpsc::Sender<HealthEvent>, 2311 ready: bool, 2312 cancellation: &radroots_service_host::CancellationToken, 2313 ) -> Result<(), HostError> { 2314 tokio::select! { 2315 result = health.send(HealthEvent::Providers(ready)) => { 2316 result.map_err(|_| HostError::new(HostErrorKind::TaskFailure)) 2317 } 2318 () = cancellation.cancelled() => Ok(()), 2319 } 2320 } 2321 2322 fn spawn_delivery_outbox( 2323 supervisor: &mut TaskSupervisor, 2324 state: Arc<MycStateHost>, 2325 configuration: Arc<MycConfigDocumentV1>, 2326 health: mpsc::Sender<HealthEvent>, 2327 ) -> Result<(), MycDaemonError> { 2328 supervisor 2329 .spawn( 2330 task_metadata( 2331 TASK_DELIVERY_OUTBOX, 2332 TaskClassification::Critical, 2333 ShutdownPhase::PersistRecoverableWork, 2334 )?, 2335 move |cancellation| async move { 2336 let worker = MycDeliveryWorker::from_configuration(&configuration) 2337 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 2338 let task_cancellation = MycTaskCancellation::from_host(cancellation.clone()); 2339 let idle = Duration::from_millis( 2340 configuration_integer( 2341 &configuration, 2342 "/transport/publish_retry/initial_backoff_ms", 2343 ) 2344 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?, 2345 ); 2346 loop { 2347 if cancellation.is_cancelled() { 2348 return Ok(()); 2349 } 2350 let now = wall_time_millis().map_err(|error| { 2351 HostError::with_source(HostErrorKind::TaskFailure, error) 2352 })?; 2353 let now = MycDeliveryTimeUnixMs::new(now).map_err(|error| { 2354 HostError::with_source(HostErrorKind::TaskFailure, error) 2355 })?; 2356 let next = state 2357 .repository() 2358 .next_ready_delivery_target(now) 2359 .await 2360 .map_err(|error| { 2361 HostError::with_source(HostErrorKind::TaskFailure, error) 2362 })?; 2363 let Some((job, relay)) = next else { 2364 tokio::select! { 2365 () = cancellation.cancelled() => return Ok(()), 2366 () = tokio::time::sleep(idle) => continue, 2367 } 2368 }; 2369 let mut nonce = [0_u8; 32]; 2370 SystemEntropy.fill_bytes(&mut nonce).map_err(|error| { 2371 HostError::with_source(HostErrorKind::TaskFailure, error) 2372 })?; 2373 let mut retry_entropy = [0_u8; 8]; 2374 SystemEntropy 2375 .fill_bytes(&mut retry_entropy) 2376 .map_err(|error| { 2377 HostError::with_source(HostErrorKind::TaskFailure, error) 2378 })?; 2379 let evidence = MycDeliveryExecutionEvidence { 2380 claimed_at: now, 2381 submitted_at: now, 2382 observed_at: now, 2383 retry_entropy, 2384 }; 2385 worker 2386 .run_one( 2387 &state.repository(), 2388 job, 2389 &relay, 2390 MycDeliveryAttemptNonce::from_injected_entropy(nonce), 2391 evidence, 2392 &task_cancellation, 2393 ) 2394 .await 2395 .map_err(|error| { 2396 HostError::with_source(HostErrorKind::TaskFailure, error) 2397 })?; 2398 let _ = health.try_send(HealthEvent::StateChanged); 2399 } 2400 }, 2401 ) 2402 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 2403 } 2404 2405 fn task_metadata( 2406 name: &'static str, 2407 classification: TaskClassification, 2408 phase: ShutdownPhase, 2409 ) -> Result<TaskMetadata, MycDaemonError> { 2410 TaskMetadata::new( 2411 TaskName::new(name).map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, 2412 classification, 2413 Some(phase), 2414 ) 2415 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 2416 } 2417 2418 fn runtime_build_info() -> Result<MycStatusBuildInfoV1, MycDaemonError> { 2419 let service_revision = option_env!("RADROOTS_SERVICE_REVISION"); 2420 let lib_revision = option_env!("RADROOTS_LIB_REVISION"); 2421 let rust_version = option_env!("RADROOTS_RUST_VERSION"); 2422 let target = option_env!("RADROOTS_BUILD_TARGET"); 2423 let mode = if service_revision.is_some() 2424 && lib_revision.is_some() 2425 && rust_version.is_some() 2426 && target.is_some() 2427 { 2428 MycStatusBuildMode::Release 2429 } else { 2430 MycStatusBuildMode::Development 2431 }; 2432 MycStatusBuildInfoV1::new( 2433 mode, 2434 Some(env!("CARGO_PKG_VERSION")), 2435 service_revision, 2436 lib_revision, 2437 rust_version, 2438 target, 2439 Some("service-host"), 2440 ) 2441 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 2442 } 2443 2444 fn wall_time_millis() -> Result<u64, MycDaemonError> { 2445 SystemWallClock 2446 .now_utc() 2447 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))? 2448 .get() 2449 .checked_mul(1_000) 2450 .filter(|value| *value <= i64::MAX as u64) 2451 .ok_or_else(|| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 2452 } 2453 2454 fn configuration_integer( 2455 configuration: &MycConfigDocumentV1, 2456 pointer: &str, 2457 ) -> Result<u64, MycDaemonError> { 2458 configuration 2459 .normalized() 2460 .pointer(pointer) 2461 .and_then(Value::as_u64) 2462 .ok_or_else(|| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 2463 } 2464 2465 fn configuration_usize( 2466 configuration: &MycConfigDocumentV1, 2467 pointer: &str, 2468 ) -> Result<usize, MycDaemonError> { 2469 configuration_integer(configuration, pointer)? 2470 .try_into() 2471 .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) 2472 } 2473 2474 fn response( 2475 route: MycAdminRoute, 2476 value: Value, 2477 ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { 2478 let bytes = serde_json::to_vec(&value) 2479 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; 2480 MycAdminResponseDocument::from_canonical_bytes(route, &bytes) 2481 .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) 2482 } 2483 2484 const fn handler_error(kind: MycAdminHandlerErrorKind) -> MycAdminHandlerError { 2485 MycAdminHandlerError::new(kind) 2486 } 2487 2488 const fn map_admin_journal_error(error: crate::MycAdminOperationError) -> MycAdminHandlerError { 2489 match error.kind() { 2490 MycAdminOperationErrorKind::OperationConflict => { 2491 handler_error(MycAdminHandlerErrorKind::OperationIdConflict) 2492 } 2493 MycAdminOperationErrorKind::OperationOutcomeUnknown 2494 | MycAdminOperationErrorKind::CommitOutcomeUnknown => { 2495 handler_error(MycAdminHandlerErrorKind::Conflict) 2496 } 2497 MycAdminOperationErrorKind::ResourceExhausted => { 2498 handler_error(MycAdminHandlerErrorKind::Unavailable) 2499 } 2500 MycAdminOperationErrorKind::InvalidMode 2501 | MycAdminOperationErrorKind::InvalidInput 2502 | MycAdminOperationErrorKind::Binding 2503 | MycAdminOperationErrorKind::Transaction => { 2504 handler_error(MycAdminHandlerErrorKind::Internal) 2505 } 2506 } 2507 } 2508 2509 #[cfg(test)] 2510 mod tests { 2511 use std::{fs, os::unix::fs::PermissionsExt as _}; 2512 2513 use super::*; 2514 2515 #[test] 2516 fn exact_task_inventory_and_shutdown_phases_are_frozen() { 2517 let expected = [ 2518 ( 2519 TASK_ADMIN_SERVER, 2520 TaskClassification::Critical, 2521 ShutdownPhase::CloseSockets, 2522 ), 2523 ( 2524 TASK_OPERATIONS_SERVER, 2525 TaskClassification::Optional, 2526 ShutdownPhase::CloseSockets, 2527 ), 2528 ( 2529 TASK_RELAY_INGRESS, 2530 TaskClassification::Critical, 2531 ShutdownPhase::CancelIngress, 2532 ), 2533 ( 2534 TASK_PROVIDER_DISPATCH, 2535 TaskClassification::Critical, 2536 ShutdownPhase::DrainOperations, 2537 ), 2538 ( 2539 TASK_DELIVERY_OUTBOX, 2540 TaskClassification::Critical, 2541 ShutdownPhase::PersistRecoverableWork, 2542 ), 2543 ]; 2544 for (name, classification, phase) in expected { 2545 let metadata = task_metadata(name, classification, phase).expect("task metadata"); 2546 assert_eq!(metadata.name().as_str(), name); 2547 assert_eq!(metadata.classification(), classification); 2548 assert_eq!(metadata.shutdown_phase(), Some(phase)); 2549 } 2550 } 2551 2552 #[test] 2553 fn runtime_errors_are_safe_and_have_stable_exit_classes() { 2554 for (kind, result) in [ 2555 ( 2556 MycDaemonErrorKind::State, 2557 MycProcessResult::StateOrIdentityUnavailable, 2558 ), 2559 ( 2560 MycDaemonErrorKind::Provider, 2561 MycProcessResult::ServiceOrDependencyUnavailable, 2562 ), 2563 ( 2564 MycDaemonErrorKind::Relay, 2565 MycProcessResult::ServiceOrDependencyUnavailable, 2566 ), 2567 ( 2568 MycDaemonErrorKind::Admin, 2569 MycProcessResult::ServiceOrDependencyUnavailable, 2570 ), 2571 ( 2572 MycDaemonErrorKind::Operations, 2573 MycProcessResult::ServiceOrDependencyUnavailable, 2574 ), 2575 ( 2576 MycDaemonErrorKind::Runtime, 2577 MycProcessResult::UnexpectedInternal, 2578 ), 2579 ] { 2580 let error = MycDaemonError::new(kind); 2581 assert_eq!(error.process_result(), result); 2582 assert_eq!(error.to_string(), "Myc daemon failed"); 2583 assert!(Error::source(&error).is_none()); 2584 } 2585 } 2586 2587 #[test] 2588 fn backoff_inputs_are_bounded_by_the_admitted_configuration() { 2589 let configuration = crate::parse_myc_config_v1( 2590 include_bytes!("../contracts/services_hardening/config.v1.example.toml"), 2591 crate::MycConfigProfile::Production, 2592 ) 2593 .expect("configuration"); 2594 assert_eq!( 2595 configuration_integer( 2596 &configuration, 2597 "/transport/publish_retry/initial_backoff_ms" 2598 ) 2599 .expect("initial"), 2600 250 2601 ); 2602 assert_eq!( 2603 configuration_integer( 2604 &configuration, 2605 "/transport/publish_retry/maximum_backoff_ms" 2606 ) 2607 .expect("maximum"), 2608 30_000 2609 ); 2610 } 2611 2612 #[test] 2613 fn admin_cursors_are_authenticated_query_bound_and_round_trip_exactly() { 2614 let key = [0x42; 32]; 2615 let snapshot = MycConnectionTimeUnixMs::new(10_000).expect("snapshot"); 2616 let before = MycConnectionTimeUnixMs::new(9_000).expect("before"); 2617 let id = MycConnectionId::from_bytes([0x24; 32]); 2618 let encoded = encode_connection_cursor( 2619 snapshot, 2620 before, 2621 id, 2622 Some(MycConnectionStatus::Active), 2623 &key, 2624 ); 2625 let decoded = decode_connection_cursor(&encoded, Some(MycConnectionStatus::Active), &key) 2626 .expect("connection cursor"); 2627 assert_eq!(decoded.snapshot, snapshot); 2628 assert_eq!(decoded.before, before); 2629 assert_eq!(decoded.id, id); 2630 assert!( 2631 decode_connection_cursor(&encoded, Some(MycConnectionStatus::Denied), &key).is_err() 2632 ); 2633 2634 let query = audit_query_digest( 2635 Some(1_000), 2636 Some(9_999), 2637 Some(MycAuditKind::ChallengeAuthorization), 2638 Some(MycAuditOutcome::Succeeded), 2639 ); 2640 let encoded = encode_audit_cursor(17, 11, query, &key); 2641 let decoded = decode_audit_cursor(&encoded, query, &key).expect("audit cursor"); 2642 assert_eq!(decoded.snapshot, 17); 2643 assert_eq!(decoded.before, 11); 2644 let other_query = audit_query_digest(None, None, None, None); 2645 assert!(decode_audit_cursor(&encoded, other_query, &key).is_err()); 2646 2647 let mut tampered = URL_SAFE_NO_PAD.decode(&encoded).expect("cursor bytes"); 2648 tampered[10] ^= 1; 2649 assert!(decode_audit_cursor(&URL_SAFE_NO_PAD.encode(tampered), query, &key).is_err()); 2650 } 2651 2652 #[test] 2653 fn admin_time_and_identity_helpers_are_bounded_and_domain_separated() { 2654 assert_eq!(seconds_as_millis(1, false).expect("start"), 1_000); 2655 assert_eq!(seconds_as_millis(1, true).expect("end"), 1_999); 2656 assert!(seconds_as_millis(u64::MAX, false).is_err()); 2657 let value = "caller-01"; 2658 let operation = admin_identity(b"operation\0", value); 2659 let correlation = admin_identity(b"correlation\0", value); 2660 assert_ne!(operation, correlation); 2661 assert!(!format!("{:?}", MycAuditCorrelationId::new(operation)).contains(value)); 2662 } 2663 2664 #[test] 2665 fn admin_metric_keys_are_closed_and_bounded() { 2666 for value in [ 2667 "radroots_myc_service_ready", 2668 "radroots_myc_service_phase_starting", 2669 "radroots_myc_service_phase_degraded", 2670 ] { 2671 assert!(safe_metric_key(value)); 2672 } 2673 for value in ["", "Uppercase", "hyphen-key", "path/key", "secret:key"] { 2674 assert!(!safe_metric_key(value)); 2675 } 2676 assert!(safe_metric_key(&"x".repeat(64))); 2677 assert!(!safe_metric_key(&"x".repeat(65))); 2678 } 2679 2680 #[tokio::test] 2681 async fn runtime_status_queries_real_connection_and_outbox_state() { 2682 let directory = tempfile::tempdir().expect("temporary root"); 2683 let runtime = crate::nip46_wave_080_a::runtime(directory.path()); 2684 fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); 2685 fs::set_permissions( 2686 runtime.context().paths().state(), 2687 fs::Permissions::from_mode(0o700), 2688 ) 2689 .expect("state mode"); 2690 let metadata = crate::nip46_wave_080_a::metadata(&runtime); 2691 let configuration = crate::nip46_wave_080_a::configuration(); 2692 let (applied_at, build) = crate::nip46_wave_080_a::migration_evidence(); 2693 crate::initialize_myc_state(&runtime, &metadata, applied_at, &build) 2694 .await 2695 .expect("initialize"); 2696 let host = crate::open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 2697 .await 2698 .expect("host"); 2699 let repository = host.repository(); 2700 let empty = repository 2701 .read_runtime_connection_counts() 2702 .await 2703 .expect("empty connection counts"); 2704 assert_eq!(empty, MycConnectionCountsV1::default()); 2705 let pending = crate::nip46_wave_080_a::admit_connect( 2706 &repository, 2707 &configuration, 2708 20, 2709 "runtime-status-connect", 2710 crate::nip46_wave_080_a::OBSERVED_AT_SECONDS, 2711 crate::nip46_wave_080_a::RECEIVED_AT_MS, 2712 21, 2713 crate::nip46_wave_080_a::RECEIVED_AT_MS + 1, 2714 ) 2715 .await; 2716 assert!(pending.record().is_some()); 2717 let counts = repository 2718 .read_runtime_connection_counts() 2719 .await 2720 .expect("connection counts"); 2721 assert_eq!(counts.pending(), 1); 2722 assert_eq!(counts.active(), 0); 2723 assert_eq!(counts.denied(), 0); 2724 assert_eq!(counts.expired(), 0); 2725 assert_eq!( 2726 repository 2727 .read_runtime_outbox_status() 2728 .await 2729 .expect("empty outbox"), 2730 MycOutboxStatusV1::default() 2731 ); 2732 host.close().await.expect("close"); 2733 2734 let inspection = crate::open_myc_state_inspection(&runtime, &metadata) 2735 .await 2736 .expect("inspection"); 2737 let repository = inspection.repository(); 2738 assert_eq!( 2739 repository 2740 .read_runtime_outbox_status() 2741 .await 2742 .expect("inspection outbox"), 2743 MycOutboxStatusV1::default() 2744 ); 2745 assert_eq!( 2746 repository 2747 .read_runtime_connection_counts() 2748 .await 2749 .expect("inspection connection counts") 2750 .pending(), 2751 1 2752 ); 2753 inspection 2754 .inspect_integrity( 2755 radroots_service_sqlite::IntegrityCheckedAtUnixMs::new(1).expect("integrity time"), 2756 ) 2757 .await 2758 .expect("inspection integrity"); 2759 inspection.close().await.expect("inspection close"); 2760 } 2761 }