runtime_graph.rs (68755B)
1 //! Binary-invoked RHI daemon graph with fixed task ownership and bounded shutdown. 2 3 use core::{fmt, future::pending, time::Duration}; 4 use std::{ 5 collections::{BTreeMap, BTreeSet}, 6 error::Error, 7 sync::{ 8 Arc, 9 atomic::{AtomicBool, Ordering}, 10 }, 11 }; 12 13 use radroots_event::id::TradeId; 14 use radroots_service_host::{ 15 GracefulShutdown, HostError, HostErrorKind, ProcessSignal, ProcessSignalAdapter, 16 ProcessSignalFuture, ProcessSignalSource, ShutdownDisposition, ShutdownPhase, 17 ShutdownPhaseFuture, ShutdownPhaseHandler, SupervisedTaskExitStatus, TaskClassification, 18 TaskMetadata, TaskName, 19 }; 20 use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity}; 21 use radroots_transport::{ 22 FetchRequest, SubscriptionNext, SubscriptionRequest, Target, TargetSet, 23 outcome::FetchTargetState, 24 source::{ 25 FETCH_PAGE_MAX_EVENTS, FetchBounds, FetchCursor, FetchSelector, NextPage, 26 SubscriptionBounds, 27 }, 28 }; 29 use serde_json::Value; 30 use sha2::{Digest as _, Sha256}; 31 use sqlx::Row; 32 33 use crate::transport_nostr_adapter::{RhiNostrExactSink, build_rhi_nostr_adapters}; 34 use crate::{ 35 RhiAdminCancellationToken, RhiAdminServer, RhiAdmittedTradeMutationEvent, RhiBoundAdminServer, 36 RhiBoundOperationsServer, RhiConfigDocumentV1, RhiEvidenceAttestationSupersession, 37 RhiEvidenceTransportStatusV1, RhiIdentityHealthV1, RhiIntegrityStateV1, 38 RhiJitterBoundMilliseconds, RhiOperationsCancellationToken, RhiOperationsServer, 39 RhiPersistenceHealthV1, RhiPersistenceStatusV1, RhiPresenceDesiredAuthority, 40 RhiPresenceLeaseOwner, RhiPresenceStatusV1, RhiProcessResult, RhiProcessSignal, 41 RhiProcessSignalSource, RhiProviderStatusV1, RhiPublicationAuthority, RhiPublicationLeaseOwner, 42 RhiPublicationStatusV1, RhiReconciliationAttemptPlan, RhiReconciliationJobPolicy, 43 RhiReconciliationLease, RhiReconciliationLeaseOwner, RhiReconciliationRetryDelayMilliseconds, 44 RhiReconciliationScopePrerequisites, RhiReconciliationSourceReplay, 45 RhiReconciliationSourceReplayPlan, RhiReconciliationSourceRequest, RhiReconciliationStatusV1, 46 RhiReconciliationUnixMilliseconds, RhiRuntimeFoundation, RhiServicePhase, RhiStateHost, 47 RhiStatusBuildInfoV1, RhiStatusBuildMode, RhiStatusCommonV1, RhiStatusConfigurationIdentityV1, 48 RhiStatusConfigurationSource, RhiStatusObservationV1, RhiStatusPublisher, RhiStatusReasonCode, 49 RhiStatusReasonCodes, RhiStatusUnixSeconds, RhiTimeEntropyAdapters, 50 RhiTradeMutationAdmissionLimits, RhiTradeMutationAuthoredTimePolicy, 51 RhiTradeMutationObservedAtUnixSeconds, RhiTradeSourceAttempt, RhiTradeSourceCompletion, 52 RhiTransportHealthV1, admit_rhi_trade_mutation_event, build_rhi_signed_evidence_attestation, 53 build_rhi_signed_presence_documents, evaluate_rhi_reconciliation_claim, 54 open_rhi_runtime_foundation, reduce_rhi_reconciliation_manifest, rhi_status_cache, 55 }; 56 use crate::{reconciliation_replay, source_ingest, state_config}; 57 58 const TASK_ADMIN_SERVER: &str = "admin_server"; 59 const TASK_OPERATIONS_SERVER: &str = "operations_server"; 60 const TASK_SOURCE_SUBSCRIPTION: &str = "source_subscription"; 61 const TASK_RECONCILIATION_WORKER: &str = "reconciliation_worker"; 62 const TASK_PUBLICATION_WORKER: &str = "publication_worker"; 63 const TASK_PRESENCE_WORKER: &str = "presence_worker"; 64 const TRADE_EVENT_KINDS: [u32; 5] = [3470, 3471, 3472, 3473, 3474]; 65 66 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 67 enum RhiDaemonErrorKind { 68 State, 69 Identity, 70 Transport, 71 Admin, 72 Operations, 73 Runtime, 74 } 75 76 struct RhiDaemonError { 77 kind: RhiDaemonErrorKind, 78 } 79 80 impl RhiDaemonError { 81 const fn new(kind: RhiDaemonErrorKind) -> Self { 82 Self { kind } 83 } 84 85 const fn process_result(&self) -> RhiProcessResult { 86 match self.kind { 87 RhiDaemonErrorKind::State | RhiDaemonErrorKind::Identity => { 88 RhiProcessResult::StateOrIdentityUnavailable 89 } 90 RhiDaemonErrorKind::Transport 91 | RhiDaemonErrorKind::Admin 92 | RhiDaemonErrorKind::Operations => RhiProcessResult::ServiceOrDependencyUnavailable, 93 RhiDaemonErrorKind::Runtime => RhiProcessResult::UnexpectedInternal, 94 } 95 } 96 } 97 98 impl fmt::Debug for RhiDaemonError { 99 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 100 formatter 101 .debug_struct("RhiDaemonError") 102 .field("kind", &self.kind) 103 .finish() 104 } 105 } 106 107 impl fmt::Display for RhiDaemonError { 108 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 109 formatter.write_str("RHI daemon failed") 110 } 111 } 112 113 impl Error for RhiDaemonError {} 114 115 struct HostSignalSource<S> { 116 inner: S, 117 } 118 119 impl<S> HostSignalSource<S> { 120 const fn new(inner: S) -> Self { 121 Self { inner } 122 } 123 } 124 125 impl<S> ProcessSignalSource for HostSignalSource<S> 126 where 127 S: RhiProcessSignalSource, 128 { 129 fn next_signal(&mut self) -> ProcessSignalFuture<'_> { 130 Box::pin(async move { 131 self.inner.next_signal().await.map(|signal| match signal { 132 RhiProcessSignal::Interrupt => ProcessSignal::Interrupt, 133 #[cfg(unix)] 134 RhiProcessSignal::Terminate => ProcessSignal::Terminate, 135 }) 136 }) 137 } 138 } 139 140 struct RuntimeShutdownHandler { 141 accepting_mutations: Arc<AtomicBool>, 142 status: RhiStatusPublisher, 143 status_context: RuntimeStatusContext, 144 transport_ready: bool, 145 operations_ready: bool, 146 } 147 148 impl RuntimeShutdownHandler { 149 fn publish(&mut self, phase: RhiServicePhase) -> Result<(), HostError> { 150 self.status 151 .publish(self.status_context.observation( 152 phase, 153 self.transport_ready, 154 self.operations_ready, 155 )?) 156 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error)) 157 } 158 159 fn running_phase(&self) -> RhiServicePhase { 160 if self.transport_ready && self.operations_ready { 161 RhiServicePhase::Ready 162 } else { 163 RhiServicePhase::Degraded 164 } 165 } 166 } 167 168 impl ShutdownPhaseHandler for RuntimeShutdownHandler { 169 fn enter(&mut self, phase: ShutdownPhase) -> ShutdownPhaseFuture<'_> { 170 Box::pin(async move { 171 if phase == ShutdownPhase::RejectNewMutations { 172 self.accepting_mutations.store(false, Ordering::Release); 173 self.publish(RhiServicePhase::Stopping)?; 174 } 175 Ok(()) 176 }) 177 } 178 } 179 180 struct RuntimeStatusContext { 181 configuration_digest: String, 182 configuration_source: RhiStatusConfigurationSource, 183 schema_version: u32, 184 generation: u64, 185 source_count: u64, 186 started_at: radroots_service_host::MonotonicTime, 187 time_entropy: RhiTimeEntropyAdapters, 188 reconciliation: RhiReconciliationStatusV1, 189 publication: RhiPublicationStatusV1, 190 presence: RhiPresenceStatusV1, 191 } 192 193 impl RuntimeStatusContext { 194 async fn refresh_state(&mut self, state: &RhiStateHost) -> Result<(), HostError> { 195 let summary = state 196 .sqlite_host() 197 .transaction(|transaction| { 198 Box::pin(async move { 199 let row = sqlx::query(RUNTIME_STATUS_SQL) 200 .fetch_one(&mut *transaction) 201 .await 202 .map_err(|_| ())?; 203 Ok::<_, ()>(( 204 row.try_get::<i64, _>("reconciliation_pending") 205 .map_err(|_| ())?, 206 row.try_get::<i64, _>("reconciliation_leased") 207 .map_err(|_| ())?, 208 row.try_get::<i64, _>("reconciliation_exhausted") 209 .map_err(|_| ())?, 210 row.try_get::<Option<i64>, _>("reconciliation_oldest") 211 .map_err(|_| ())?, 212 row.try_get::<i64, _>("publication_pending") 213 .map_err(|_| ())?, 214 row.try_get::<i64, _>("publication_unknown") 215 .map_err(|_| ())?, 216 row.try_get::<Option<i64>, _>("publication_oldest") 217 .map_err(|_| ())?, 218 row.try_get::<i64, _>("presence_pending").map_err(|_| ())?, 219 row.try_get::<i64, _>("presence_unknown").map_err(|_| ())?, 220 )) 221 }) 222 }) 223 .await 224 .map_err(|_| HostError::new(HostErrorKind::Lifecycle))?; 225 self.reconciliation = RhiReconciliationStatusV1::new( 226 status_count(summary.0)?, 227 status_count(summary.1)?, 228 status_count(summary.2)?, 229 status_time(summary.3)?, 230 ); 231 self.publication = RhiPublicationStatusV1::new( 232 status_count(summary.4)?, 233 status_count(summary.5)?, 234 status_time(summary.6)?, 235 ); 236 self.presence = 237 RhiPresenceStatusV1::new(status_count(summary.7)?, status_count(summary.8)?); 238 Ok(()) 239 } 240 241 fn observation( 242 &self, 243 phase: RhiServicePhase, 244 transport_ready: bool, 245 operations_ready: bool, 246 ) -> Result<RhiStatusObservationV1, HostError> { 247 let ready = 248 matches!(phase, RhiServicePhase::Ready | RhiServicePhase::Degraded) && transport_ready; 249 let mut transport_reasons = Vec::new(); 250 if !transport_ready { 251 transport_reasons.push(RhiStatusReasonCode::SourceUnavailable); 252 transport_reasons.push(RhiStatusReasonCode::SubscriptionInactive); 253 } 254 let transport_reasons = RhiStatusReasonCodes::new(transport_reasons) 255 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 256 let mut lifecycle_reasons = transport_reasons.as_slice().to_vec(); 257 if !operations_ready { 258 lifecycle_reasons.push(RhiStatusReasonCode::OperationsListenerFailed); 259 } 260 if phase == RhiServicePhase::Stopping { 261 lifecycle_reasons.push(RhiStatusReasonCode::ShutdownInProgress); 262 } 263 let lifecycle_reasons = RhiStatusReasonCodes::new(lifecycle_reasons) 264 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 265 let transport_health = if transport_ready { 266 RhiTransportHealthV1::Ready 267 } else { 268 RhiTransportHealthV1::Unavailable 269 }; 270 let uptime = self 271 .time_entropy 272 .now_monotonic() 273 .duration_since_origin() 274 .saturating_sub(self.started_at.duration_since_origin()) 275 .as_millis(); 276 let uptime = u64::try_from(uptime).map_err(|_| HostError::new(HostErrorKind::Lifecycle))?; 277 let build = runtime_build_info()?; 278 let configuration = RhiStatusConfigurationIdentityV1::new( 279 &self.configuration_digest, 280 self.configuration_source, 281 ) 282 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 283 let persistence = RhiPersistenceStatusV1::new( 284 RhiPersistenceHealthV1::Ready, 285 self.schema_version, 286 self.generation, 287 RhiIntegrityStateV1::Verified, 288 RhiStatusReasonCodes::empty(), 289 ) 290 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 291 let identity = RhiIdentityHealthV1::new(true, true, RhiStatusReasonCodes::empty()) 292 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 293 let provider = RhiProviderStatusV1::new(identity, RhiStatusReasonCodes::empty()) 294 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 295 let transport = RhiEvidenceTransportStatusV1::new( 296 transport_health, 297 transport_ready, 298 transport_ready, 299 self.source_count, 300 if transport_ready { 301 self.source_count 302 } else { 303 0 304 }, 305 transport_reasons, 306 ) 307 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 308 let common = RhiStatusCommonV1::new( 309 phase, 310 ready, 311 lifecycle_reasons, 312 uptime, 313 build, 314 configuration, 315 persistence, 316 ) 317 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; 318 Ok(RhiStatusObservationV1::new( 319 common, 320 provider, 321 transport, 322 self.reconciliation, 323 self.publication, 324 self.presence, 325 )) 326 } 327 } 328 329 const RUNTIME_STATUS_SQL: &str = r#"SELECT 330 (SELECT COUNT(*) FROM reconciliation_jobs WHERE state = 'ready') 331 AS reconciliation_pending, 332 (SELECT COUNT(*) FROM reconciliation_jobs WHERE state = 'leased') 333 AS reconciliation_leased, 334 (SELECT COUNT(*) FROM reconciliation_jobs WHERE state = 'exhausted') 335 AS reconciliation_exhausted, 336 (SELECT MIN(created_at_unix_ms / 1000) FROM reconciliation_jobs WHERE state = 'ready') 337 AS reconciliation_oldest, 338 (SELECT COUNT(*) FROM publication_outbox WHERE state IN ('pending', 'leased')) 339 AS publication_pending, 340 (SELECT COUNT(*) FROM publication_outbox AS outbox 341 WHERE outbox.state = 'blocked' OR EXISTS ( 342 SELECT 1 FROM publication_targets AS target 343 WHERE target.outbox_id = outbox.outbox_id AND target.state = 'unknown')) 344 AS publication_unknown, 345 (SELECT MIN(created_at_unix_ms / 1000) FROM publication_outbox 346 WHERE state IN ('pending', 'leased')) AS publication_oldest, 347 (SELECT COUNT(*) FROM presence_outbox WHERE state IN ('pending', 'leased')) 348 AS presence_pending, 349 (SELECT COUNT(*) FROM presence_outbox AS outbox 350 WHERE outbox.state = 'blocked' OR EXISTS ( 351 SELECT 1 FROM presence_targets AS target 352 WHERE target.outbox_id = outbox.outbox_id AND target.state = 'unknown')) 353 AS presence_unknown"#; 354 355 fn status_count(value: i64) -> Result<u64, HostError> { 356 u64::try_from(value).map_err(|_| HostError::new(HostErrorKind::Lifecycle)) 357 } 358 359 fn status_time(value: Option<i64>) -> Result<Option<RhiStatusUnixSeconds>, HostError> { 360 value 361 .map(|value| { 362 u64::try_from(value) 363 .map_err(|_| HostError::new(HostErrorKind::Lifecycle)) 364 .and_then(|value| { 365 RhiStatusUnixSeconds::new(value) 366 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error)) 367 }) 368 }) 369 .transpose() 370 } 371 372 struct SourceBinding { 373 source_id: Box<str>, 374 target: Target, 375 } 376 377 struct InitialSubscription { 378 request: SubscriptionRequest, 379 subscription: radroots_transport::BoxSubscription, 380 sources: Box<[SourceBinding]>, 381 } 382 383 /// Executes the real RHI daemon and converts every failure to the frozen process result. 384 pub(crate) async fn run_rhi_daemon<S>( 385 runtime: crate::RhiRuntimeContext, 386 configuration: RhiConfigDocumentV1, 387 applied_at: MigrationAppliedAtUnixSeconds, 388 build: &MigrationBuildIdentity, 389 signals: S, 390 ) -> RhiProcessResult 391 where 392 S: RhiProcessSignalSource + 'static, 393 { 394 run_rhi_daemon_inner(runtime, configuration, applied_at, build, signals) 395 .await 396 .unwrap_or_else(|error| error.process_result()) 397 } 398 399 async fn run_rhi_daemon_inner<S>( 400 runtime: crate::RhiRuntimeContext, 401 configuration: RhiConfigDocumentV1, 402 applied_at: MigrationAppliedAtUnixSeconds, 403 build: &MigrationBuildIdentity, 404 signals: S, 405 ) -> Result<RhiProcessResult, RhiDaemonError> 406 where 407 S: RhiProcessSignalSource + 'static, 408 { 409 let (transport, exact_sink) = build_rhi_nostr_adapters(&configuration) 410 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Transport))?; 411 let adapters = crate::RhiRuntimeAdapters::new( 412 crate::RhiTimeEntropyAdapters::system(), 413 transport, 414 crate::RhiIdentityCredentialAdapters::canonical(), 415 ); 416 let mut foundation = 417 open_rhi_runtime_foundation(runtime, configuration, adapters, applied_at, build) 418 .await 419 .map_err(|error| match error.kind() { 420 crate::RhiRuntimeFoundationErrorKind::IdentityAccess 421 | crate::RhiRuntimeFoundationErrorKind::IdentityBinding => { 422 RhiDaemonError::new(RhiDaemonErrorKind::Identity) 423 } 424 _ => RhiDaemonError::new(RhiDaemonErrorKind::State), 425 })?; 426 427 let configuration = foundation.configuration_arc(); 428 let state = foundation.state(); 429 let publication = Arc::new( 430 RhiPublicationAuthority::from_config(&configuration) 431 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?, 432 ); 433 let presence = Arc::new( 434 RhiPresenceDesiredAuthority::from_config(&configuration) 435 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?, 436 ); 437 initialize_presence(&foundation, &presence) 438 .await 439 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::State))?; 440 441 let initial_subscription = open_initial_subscription(&foundation, &configuration) 442 .await 443 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Transport))?; 444 let generation = u64::from( 445 state_config::current_generation(&state) 446 .await 447 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::State))?, 448 ); 449 let source_count = u64::try_from(initial_subscription.sources.len()) 450 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 451 let time_entropy = foundation.adapters().time_entropy().clone(); 452 let started_at = time_entropy.now_monotonic(); 453 let mut status_context = RuntimeStatusContext { 454 configuration_digest: lower_hex(state.metadata().configuration_digest().as_bytes()), 455 configuration_source: if foundation.runtime_context().profile() 456 == crate::RhiBootstrapProfileV1::RepoLocal 457 { 458 RhiStatusConfigurationSource::DerivedRepoLocal 459 } else { 460 RhiStatusConfigurationSource::ExplicitConfig 461 }, 462 schema_version: crate::RHI_STATE_SCHEMA_VERSION, 463 generation, 464 source_count, 465 started_at, 466 time_entropy: time_entropy.clone(), 467 reconciliation: RhiReconciliationStatusV1::default(), 468 publication: RhiPublicationStatusV1::default(), 469 presence: RhiPresenceStatusV1::default(), 470 }; 471 status_context 472 .refresh_state(&state) 473 .await 474 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::State))?; 475 let (status, status_reader) = rhi_status_cache( 476 foundation.runtime_context().context().instance().clone(), 477 status_context 478 .observation(RhiServicePhase::Starting, true, true) 479 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?, 480 ) 481 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 482 let accepting_mutations = Arc::new(AtomicBool::new(true)); 483 let mut cursor_key = [0_u8; 32]; 484 time_entropy 485 .entropy() 486 .fill_bytes(&mut cursor_key) 487 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 488 let admin_handler = Arc::new(crate::runtime_admin::RuntimeAdminHandler { 489 state: Arc::clone(&state), 490 configuration: Arc::clone(&configuration), 491 identity: foundation.identity_arc(), 492 publication: Arc::clone(&publication), 493 presence: Arc::clone(&presence), 494 status: status_reader.clone(), 495 accepting_mutations: Arc::clone(&accepting_mutations), 496 cursor_key, 497 time_entropy: time_entropy.clone(), 498 }); 499 let admin = RhiAdminServer::new(&configuration, admin_handler) 500 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Admin))? 501 .bind(foundation.runtime_context()) 502 .await 503 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Admin))?; 504 let operations = if operations_enabled(&configuration)? { 505 Some( 506 RhiOperationsServer::new(&configuration, &status_reader) 507 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Operations))? 508 .bind() 509 .await 510 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Operations))?, 511 ) 512 } else { 513 None 514 }; 515 516 spawn_admin(foundation.supervisor_mut(), admin)?; 517 if let Some(operations) = operations { 518 spawn_operations(foundation.supervisor_mut(), operations)?; 519 } 520 let source_transport = foundation.adapters().transport().clone(); 521 let (transport_health_sender, mut transport_health) = tokio::sync::watch::channel(true); 522 spawn_source_subscription( 523 foundation.supervisor_mut(), 524 Arc::clone(&state), 525 Arc::clone(&configuration), 526 source_transport, 527 time_entropy.clone(), 528 initial_subscription, 529 transport_health_sender, 530 )?; 531 let reconciliation_transport = foundation.adapters().transport().clone(); 532 let reconciliation_time = foundation.adapters().time_entropy().clone(); 533 let reconciliation_identity = foundation.identity_arc(); 534 spawn_reconciliation_worker( 535 foundation.supervisor_mut(), 536 Arc::clone(&state), 537 Arc::clone(&configuration), 538 reconciliation_transport, 539 reconciliation_time, 540 reconciliation_identity, 541 Arc::clone(&publication), 542 )?; 543 spawn_publication_worker( 544 foundation.supervisor_mut(), 545 Arc::clone(&state), 546 Arc::clone(&configuration), 547 Arc::clone(&publication), 548 Arc::clone(&exact_sink), 549 time_entropy.clone(), 550 )?; 551 spawn_presence_worker( 552 foundation.supervisor_mut(), 553 Arc::clone(&state), 554 Arc::clone(&configuration), 555 Arc::clone(&exact_sink), 556 time_entropy, 557 )?; 558 559 let mut shutdown_handler = RuntimeShutdownHandler { 560 accepting_mutations, 561 status, 562 status_context, 563 transport_ready: true, 564 operations_ready: true, 565 }; 566 shutdown_handler 567 .publish(RhiServicePhase::Ready) 568 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 569 570 let mut signals = ProcessSignalAdapter::new(HostSignalSource::new(signals)); 571 let refresh_period = worker_idle(&configuration) 572 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 573 let mut status_refresh = tokio::time::interval(refresh_period); 574 status_refresh.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); 575 status_refresh.tick().await; 576 let signal_initiated = loop { 577 tokio::select! { 578 action = signals.next_action() => { 579 match action { 580 Ok(action) if !action.forces_termination() => break true, 581 Ok(_) | Err(_) => break false, 582 } 583 } 584 changed = transport_health.changed() => { 585 if changed.is_err() { 586 break false; 587 } 588 let ready = *transport_health.borrow_and_update(); 589 shutdown_handler.transport_ready = ready; 590 shutdown_handler.publish(shutdown_handler.running_phase()) 591 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 592 } 593 joined = foundation.supervisor_mut().join_next() => { 594 match joined { 595 Some(Ok(exit)) 596 if exit.status() == SupervisedTaskExitStatus::OptionalFailure 597 || exit.metadata().name().as_str() == TASK_OPERATIONS_SERVER => { 598 shutdown_handler.operations_ready = false; 599 shutdown_handler.publish(shutdown_handler.running_phase()) 600 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 601 } 602 Some(Ok(_)) | Some(Err(_)) | None => break false, 603 } 604 } 605 _ = status_refresh.tick() => { 606 shutdown_handler.status_context.refresh_state(&state).await 607 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::State))?; 608 shutdown_handler.publish(shutdown_handler.running_phase()) 609 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 610 } 611 } 612 }; 613 let grace = Duration::from_millis(configuration_integer( 614 &configuration, 615 "/service/shutdown_grace_ms", 616 )?); 617 let mut shutdown = GracefulShutdown::new(grace) 618 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 619 let clock = radroots_service_host::SystemMonotonicClock::new(); 620 let summary = if signal_initiated { 621 shutdown 622 .run( 623 &clock, 624 foundation.supervisor_mut(), 625 &mut shutdown_handler, 626 async { 627 let _ = signals.next_action().await; 628 }, 629 ) 630 .await 631 } else { 632 shutdown 633 .run( 634 &clock, 635 foundation.supervisor_mut(), 636 &mut shutdown_handler, 637 pending::<()>(), 638 ) 639 .await 640 } 641 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?; 642 drop(shutdown_handler); 643 drop(state); 644 foundation 645 .shutdown() 646 .await 647 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::State))?; 648 if signal_initiated && summary.disposition() == ShutdownDisposition::Completed { 649 Ok(RhiProcessResult::Success) 650 } else { 651 Err(RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 652 } 653 } 654 655 async fn initialize_presence( 656 foundation: &RhiRuntimeFoundation, 657 authority: &RhiPresenceDesiredAuthority, 658 ) -> Result<(), HostError> { 659 let state = foundation.state(); 660 let desired = state 661 .repositories() 662 .desired_presence() 663 .commit(authority) 664 .await 665 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 666 if desired.changed() { 667 let documents = build_rhi_signed_presence_documents( 668 desired, 669 authority, 670 foundation.identity(), 671 foundation 672 .adapters() 673 .time_entropy() 674 .now_utc() 675 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?, 676 foundation.adapters().time_entropy().entropy(), 677 ) 678 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 679 let now = crate::RhiPresenceUnixMilliseconds::new( 680 foundation 681 .adapters() 682 .time_entropy() 683 .now_utc_milliseconds() 684 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?, 685 ) 686 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 687 state 688 .repositories() 689 .presence_outbox() 690 .commit_signed_presence(&documents, now) 691 .await 692 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 693 } 694 Ok(()) 695 } 696 697 async fn open_initial_subscription( 698 foundation: &RhiRuntimeFoundation, 699 configuration: &RhiConfigDocumentV1, 700 ) -> Result<InitialSubscription, HostError> { 701 let sources = configured_sources(configuration)?; 702 let (request, subscription) = subscribe_sources( 703 configuration, 704 foundation.adapters().transport(), 705 foundation.adapters().time_entropy(), 706 &sources, 707 ) 708 .await?; 709 Ok(InitialSubscription { 710 request, 711 subscription, 712 sources: sources.into_boxed_slice(), 713 }) 714 } 715 716 async fn subscribe_sources( 717 configuration: &RhiConfigDocumentV1, 718 transport: &crate::RhiTransportAdapters, 719 time_entropy: &RhiTimeEntropyAdapters, 720 sources: &[SourceBinding], 721 ) -> Result<(SubscriptionRequest, radroots_transport::BoxSubscription), HostError> { 722 let targets = TargetSet::new(sources.iter().map(|source| source.target.clone()).collect()) 723 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 724 let deadline = time_entropy 725 .now_utc_milliseconds() 726 .ok() 727 .and_then(|now| now.checked_add(maximum_source_deadline(configuration).ok()?)) 728 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))?; 729 let selector = FetchSelector::all() 730 .with_kinds(TRADE_EVENT_KINDS.to_vec()) 731 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 732 let request = SubscriptionRequest::new( 733 "rhi-source-subscription", 734 targets, 735 SubscriptionBounds::new(1_000, deadline) 736 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?, 737 ) 738 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))? 739 .with_selector(selector); 740 let subscription = transport 741 .evidence_subscriber() 742 .subscribe(request.clone()) 743 .await 744 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 745 Ok((request, subscription)) 746 } 747 748 fn configured_sources( 749 configuration: &RhiConfigDocumentV1, 750 ) -> Result<Vec<SourceBinding>, HostError> { 751 let relays = configuration 752 .normalized() 753 .pointer("/relays") 754 .and_then(Value::as_array) 755 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))?; 756 let relay_urls = relays 757 .iter() 758 .map(|relay| { 759 let id = relay.pointer("/id").and_then(Value::as_str)?; 760 let url = relay.pointer("/url").and_then(Value::as_str)?; 761 Some((id, url)) 762 }) 763 .collect::<Option<BTreeMap<_, _>>>() 764 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))?; 765 configuration 766 .normalized() 767 .pointer("/evidence/sources") 768 .and_then(Value::as_array) 769 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))? 770 .iter() 771 .map(|source| { 772 let source_id = source 773 .pointer("/source_id") 774 .and_then(Value::as_str) 775 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))?; 776 let relay_id = source 777 .pointer("/relay_id") 778 .and_then(Value::as_str) 779 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))?; 780 let url = relay_urls 781 .get(relay_id) 782 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))?; 783 Ok(SourceBinding { 784 source_id: source_id.into(), 785 target: Target::nostr_relay(url) 786 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?, 787 }) 788 }) 789 .collect() 790 } 791 792 fn spawn_admin( 793 supervisor: &mut radroots_service_host::TaskSupervisor, 794 server: RhiBoundAdminServer, 795 ) -> Result<(), RhiDaemonError> { 796 supervisor 797 .spawn( 798 task_metadata(TASK_ADMIN_SERVER, TaskClassification::Critical, ShutdownPhase::CloseSockets)?, 799 move |cancellation| async move { 800 let token = RhiAdminCancellationToken::new(); 801 let serve_token = token.clone(); 802 let serve = server.serve(serve_token); 803 tokio::pin!(serve); 804 tokio::select! { 805 result = serve.as_mut() => result.map_err(|error| HostError::with_source(HostErrorKind::AdminTransport, error)), 806 () = cancellation.cancelled() => { 807 token.cancel(); 808 serve.await.map_err(|error| HostError::with_source(HostErrorKind::AdminTransport, error)) 809 } 810 } 811 }, 812 ) 813 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 814 } 815 816 fn spawn_operations( 817 supervisor: &mut radroots_service_host::TaskSupervisor, 818 server: RhiBoundOperationsServer, 819 ) -> Result<(), RhiDaemonError> { 820 supervisor 821 .spawn( 822 task_metadata(TASK_OPERATIONS_SERVER, TaskClassification::Optional, ShutdownPhase::CloseSockets)?, 823 move |cancellation| async move { 824 let token = RhiOperationsCancellationToken::new(); 825 let serve_token = token.clone(); 826 let serve = server.serve(serve_token); 827 tokio::pin!(serve); 828 tokio::select! { 829 result = serve.as_mut() => result.map_err(|error| HostError::with_source(HostErrorKind::OperationsServe, error)), 830 () = cancellation.cancelled() => { 831 token.cancel(); 832 serve.await.map_err(|error| HostError::with_source(HostErrorKind::OperationsServe, error)) 833 } 834 } 835 }, 836 ) 837 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 838 } 839 840 fn spawn_source_subscription( 841 supervisor: &mut radroots_service_host::TaskSupervisor, 842 state: Arc<RhiStateHost>, 843 configuration: Arc<RhiConfigDocumentV1>, 844 transport: crate::RhiTransportAdapters, 845 time_entropy: RhiTimeEntropyAdapters, 846 initial: InitialSubscription, 847 health: tokio::sync::watch::Sender<bool>, 848 ) -> Result<(), RhiDaemonError> { 849 supervisor 850 .spawn( 851 task_metadata( 852 TASK_SOURCE_SUBSCRIPTION, 853 TaskClassification::Critical, 854 ShutdownPhase::CancelIngress, 855 )?, 856 move |cancellation| async move { 857 run_source_subscription( 858 cancellation, 859 state, 860 configuration, 861 transport, 862 time_entropy, 863 initial, 864 health, 865 ) 866 .await 867 }, 868 ) 869 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 870 } 871 872 async fn run_source_subscription( 873 cancellation: radroots_service_host::CancellationToken, 874 state: Arc<RhiStateHost>, 875 configuration: Arc<RhiConfigDocumentV1>, 876 transport: crate::RhiTransportAdapters, 877 time_entropy: RhiTimeEntropyAdapters, 878 mut active: InitialSubscription, 879 health: tokio::sync::watch::Sender<bool>, 880 ) -> Result<(), HostError> { 881 let policy = RhiReconciliationJobPolicy::from_configuration(&configuration) 882 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 883 let retry_delay = worker_idle(&configuration)?; 884 loop { 885 let next = tokio::select! { 886 () = cancellation.cancelled() => { 887 active.subscription.cancel().await 888 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 889 return Ok(()); 890 } 891 next = active.subscription.next() => next, 892 }; 893 let reconnect = match next { 894 Ok(SubscriptionNext::Event(event)) => { 895 event 896 .validate_for_request(&active.request) 897 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 898 let Some(source) = active.sources.iter().find(|source| { 899 source.target.fingerprint() == event.observed().provenance().target() 900 }) else { 901 return Err(HostError::new(HostErrorKind::TaskFailure)); 902 }; 903 let now = time_entropy 904 .now_utc() 905 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))? 906 .get(); 907 let observed = RhiTradeMutationObservedAtUnixSeconds::new(now) 908 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 909 let limits = RhiTradeMutationAdmissionLimits::from_config(&configuration) 910 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 911 let admitted = match admit_rhi_trade_mutation_event( 912 limits, 913 event.observed().event().raw_json().as_bytes(), 914 observed, 915 RhiTradeMutationAuthoredTimePolicy::new(0).map_err(|error| { 916 HostError::with_source(HostErrorKind::TaskFailure, error) 917 })?, 918 ) { 919 Ok(admitted) => admitted, 920 Err(_) => continue, 921 }; 922 let trade_id = admitted.mutation().trade_id; 923 let attempt = RhiTradeSourceAttempt::new( 924 subscription_attempt_id(event.checkpoint().cursor().as_str()), 925 radroots_service_host::UnixTimeSeconds::new(now), 926 observed, 927 RhiTradeMutationAuthoredTimePolicy::new(0).map_err(|error| { 928 HostError::with_source(HostErrorKind::TaskFailure, error) 929 })?, 930 ) 931 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 932 let outcome = source_ingest::ingest_rhi_subscribed_trade_event( 933 &state.repositories(), 934 &configuration, 935 &source.source_id, 936 admitted, 937 attempt, 938 ) 939 .await 940 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 941 if outcome.dirty_generation_advanced() { 942 state 943 .repositories() 944 .reconciliation_jobs() 945 .schedule_trade( 946 trade_id, 947 policy, 948 RhiReconciliationUnixMilliseconds::new( 949 now.checked_mul(1_000) 950 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))?, 951 ) 952 .map_err(|error| { 953 HostError::with_source(HostErrorKind::TaskFailure, error) 954 })?, 955 ) 956 .await 957 .map_err(|error| { 958 HostError::with_source(HostErrorKind::TaskFailure, error) 959 })?; 960 } 961 false 962 } 963 Ok(SubscriptionNext::End(end)) => { 964 end.validate_for_request(&active.request) 965 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 966 true 967 } 968 Err(_) => true, 969 }; 970 if !reconnect { 971 continue; 972 } 973 health.send_replace(false); 974 loop { 975 tokio::select! { 976 () = cancellation.cancelled() => return Ok(()), 977 () = tokio::time::sleep(retry_delay) => {} 978 } 979 match subscribe_sources(&configuration, &transport, &time_entropy, &active.sources) 980 .await 981 { 982 Ok((request, subscription)) => { 983 active.request = request; 984 active.subscription = subscription; 985 health.send_replace(true); 986 break; 987 } 988 Err(_) => continue, 989 } 990 } 991 } 992 } 993 994 fn spawn_reconciliation_worker( 995 supervisor: &mut radroots_service_host::TaskSupervisor, 996 state: Arc<RhiStateHost>, 997 configuration: Arc<RhiConfigDocumentV1>, 998 transport: crate::RhiTransportAdapters, 999 time_entropy: RhiTimeEntropyAdapters, 1000 identity: Arc<crate::RhiDecryptedIdentity>, 1001 publication: Arc<RhiPublicationAuthority>, 1002 ) -> Result<(), RhiDaemonError> { 1003 supervisor 1004 .spawn( 1005 task_metadata( 1006 TASK_RECONCILIATION_WORKER, 1007 TaskClassification::Critical, 1008 ShutdownPhase::DrainOperations, 1009 )?, 1010 move |cancellation| async move { 1011 run_reconciliation_worker( 1012 cancellation, 1013 state, 1014 configuration, 1015 transport, 1016 time_entropy, 1017 identity, 1018 publication, 1019 ) 1020 .await 1021 }, 1022 ) 1023 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 1024 } 1025 1026 async fn run_reconciliation_worker( 1027 cancellation: radroots_service_host::CancellationToken, 1028 state: Arc<RhiStateHost>, 1029 configuration: Arc<RhiConfigDocumentV1>, 1030 transport: crate::RhiTransportAdapters, 1031 time_entropy: RhiTimeEntropyAdapters, 1032 identity: Arc<crate::RhiDecryptedIdentity>, 1033 publication: Arc<RhiPublicationAuthority>, 1034 ) -> Result<(), HostError> { 1035 let owner = RhiReconciliationLeaseOwner::from_bytes(nonzero_entropy_16(&time_entropy)?) 1036 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1037 let idle = Duration::from_millis( 1038 configuration_integer(&configuration, "/reconciliation/initial_backoff_ms") 1039 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?, 1040 ); 1041 loop { 1042 if cancellation.is_cancelled() { 1043 return Ok(()); 1044 } 1045 let now = reconciliation_now_with(&time_entropy)?; 1046 let lease = state 1047 .repositories() 1048 .reconciliation_jobs() 1049 .claim_next(owner, now) 1050 .await 1051 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1052 let Some(lease) = lease else { 1053 tokio::select! { 1054 () = cancellation.cancelled() => return Ok(()), 1055 () = tokio::time::sleep(idle) => continue, 1056 } 1057 }; 1058 let completed = execute_reconciliation_attempt( 1059 &cancellation, 1060 &state, 1061 &configuration, 1062 &transport, 1063 &time_entropy, 1064 &identity, 1065 &publication, 1066 lease, 1067 now, 1068 ) 1069 .await; 1070 if completed.is_err() && !cancellation.is_cancelled() { 1071 let maximum = RhiJitterBoundMilliseconds::new(lease.retry_delay_upper_bound()) 1072 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1073 let delay = time_entropy 1074 .sample_full_jitter(maximum) 1075 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1076 state 1077 .repositories() 1078 .reconciliation_jobs() 1079 .record_failure( 1080 lease, 1081 reconciliation_now_with(&time_entropy)?, 1082 RhiReconciliationRetryDelayMilliseconds::new(delay.get()).map_err(|error| { 1083 HostError::with_source(HostErrorKind::TaskFailure, error) 1084 })?, 1085 ) 1086 .await 1087 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1088 } else if cancellation.is_cancelled() { 1089 return Ok(()); 1090 } 1091 } 1092 } 1093 1094 #[allow(clippy::too_many_arguments)] 1095 async fn execute_reconciliation_attempt( 1096 cancellation: &radroots_service_host::CancellationToken, 1097 state: &RhiStateHost, 1098 configuration: &RhiConfigDocumentV1, 1099 transport: &crate::RhiTransportAdapters, 1100 time_entropy: &RhiTimeEntropyAdapters, 1101 identity: &crate::RhiDecryptedIdentity, 1102 publication: &RhiPublicationAuthority, 1103 lease: RhiReconciliationLease, 1104 started_at: RhiReconciliationUnixMilliseconds, 1105 ) -> Result<(), HostError> { 1106 let plan = RhiReconciliationAttemptPlan::from_claim(lease, configuration, started_at) 1107 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1108 let mut replays = Vec::with_capacity(plan.requests().len()); 1109 for request in plan.requests() { 1110 if cancellation.is_cancelled() { 1111 return Err(HostError::new(HostErrorKind::TaskFailure)); 1112 } 1113 let prior = reconciliation_replay::read_committed_reconciliation_cursor( 1114 &state.repositories(), 1115 request, 1116 plan.evidence_policy_digest(), 1117 ) 1118 .await 1119 .map_err(|()| HostError::new(HostErrorKind::TaskFailure))?; 1120 let replay = 1121 RhiReconciliationSourceReplayPlan::from_request(&plan, request, configuration, prior) 1122 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1123 replays.push( 1124 fetch_reconciliation_source( 1125 cancellation, 1126 transport, 1127 time_entropy, 1128 configuration, 1129 request, 1130 replay, 1131 ) 1132 .await?, 1133 ); 1134 } 1135 let committed = state 1136 .repositories() 1137 .reconciliation_attempts() 1138 .commit_source_replays(lease, plan, replays) 1139 .await 1140 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1141 let observed_at = time_entropy 1142 .now_utc() 1143 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1144 let manifest = committed 1145 .into_evidence_manifest(observed_at, RhiReconciliationScopePrerequisites::Satisfied) 1146 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1147 let projection = reduce_rhi_reconciliation_manifest(manifest) 1148 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1149 let claim = *projection 1150 .root_mutation_id() 1151 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))?; 1152 let evaluation = evaluate_rhi_reconciliation_claim(projection, claim); 1153 let now = reconciliation_now_with(time_entropy)?; 1154 let fence = state 1155 .repositories() 1156 .reconciliation_attempts() 1157 .prepare_finalization(lease, evaluation, now) 1158 .await 1159 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1160 let supersession = 1161 current_verified_supersession(state, fence.evaluation().projection().trade_id()).await?; 1162 let attestation = build_rhi_signed_evidence_attestation( 1163 fence, 1164 identity, 1165 observed_at, 1166 time_entropy.entropy(), 1167 supersession, 1168 ) 1169 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1170 state 1171 .repositories() 1172 .reconciliation_attempts() 1173 .commit_finalization(&attestation, publication, now) 1174 .await 1175 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1176 Ok(()) 1177 } 1178 1179 async fn current_verified_supersession( 1180 state: &RhiStateHost, 1181 trade_id: &TradeId, 1182 ) -> Result<Option<RhiEvidenceAttestationSupersession>, HostError> { 1183 let trade_id = *trade_id; 1184 state 1185 .sqlite_host() 1186 .transaction(move |transaction| { 1187 Box::pin(async move { 1188 let rows = sqlx::query( 1189 r#"SELECT 1190 CASE WHEN typeof(report.statement_sha256) = 'blob' 1191 AND length(report.statement_sha256) = 32 1192 THEN report.statement_sha256 ELSE NULL END AS statement_sha256, 1193 CASE WHEN typeof(event.event_id) = 'blob' AND length(event.event_id) = 32 1194 THEN event.event_id ELSE NULL END AS event_id, 1195 CASE WHEN typeof(report.canonical_report) = 'blob' 1196 AND length(report.canonical_report) BETWEEN 1 AND 16384 1197 THEN report.canonical_report ELSE NULL END AS canonical_report, 1198 CASE WHEN typeof(event.canonical_event_json) = 'blob' 1199 AND length(event.canonical_event_json) BETWEEN 1 AND 32768 1200 THEN event.canonical_event_json ELSE NULL END AS canonical_event_json 1201 FROM attestation_reports AS report 1202 JOIN signed_attestation_events AS event 1203 ON event.statement_sha256 = report.statement_sha256 1204 WHERE report.trade_id = ? AND NOT EXISTS ( 1205 SELECT 1 FROM attestation_reports AS successor 1206 WHERE successor.supersedes_statement_sha256 = report.statement_sha256 1207 ) 1208 ORDER BY report.observed_at_unix_s DESC, report.statement_sha256 DESC 1209 LIMIT 2"#, 1210 ) 1211 .bind(trade_id.as_bytes().as_slice()) 1212 .fetch_all(&mut *transaction) 1213 .await 1214 .map_err(|_| ())?; 1215 if rows.is_empty() { 1216 return Ok(None); 1217 } 1218 if rows.len() != 1 { 1219 return Err(()); 1220 } 1221 let row = &rows[0]; 1222 let statement = bounded_blob::<32>(row, "statement_sha256")?; 1223 let event_id = bounded_blob::<32>(row, "event_id")?; 1224 let canonical_report = row 1225 .try_get::<Option<Vec<u8>>, _>("canonical_report") 1226 .map_err(|_| ())? 1227 .ok_or(())?; 1228 let canonical_event = row 1229 .try_get::<Option<Vec<u8>>, _>("canonical_event_json") 1230 .map_err(|_| ())? 1231 .ok_or(())?; 1232 RhiEvidenceAttestationSupersession::from_persisted( 1233 &trade_id, 1234 statement, 1235 event_id, 1236 &canonical_report, 1237 &canonical_event, 1238 ) 1239 .map(Some) 1240 .map_err(|_| ()) 1241 }) 1242 }) 1243 .await 1244 .map_err(|_| HostError::new(HostErrorKind::TaskFailure)) 1245 } 1246 1247 fn bounded_blob<const N: usize>( 1248 row: &sqlx::sqlite::SqliteRow, 1249 column: &str, 1250 ) -> Result<[u8; N], ()> { 1251 row.try_get::<Option<Vec<u8>>, _>(column) 1252 .map_err(|_| ())? 1253 .ok_or(())? 1254 .try_into() 1255 .map_err(|_| ()) 1256 } 1257 1258 async fn fetch_reconciliation_source( 1259 cancellation: &radroots_service_host::CancellationToken, 1260 transport: &crate::RhiTransportAdapters, 1261 time_entropy: &RhiTimeEntropyAdapters, 1262 configuration: &RhiConfigDocumentV1, 1263 source_request: &RhiReconciliationSourceRequest, 1264 replay: RhiReconciliationSourceReplayPlan, 1265 ) -> Result<RhiReconciliationSourceReplay, HostError> { 1266 let source = configured_sources(configuration)? 1267 .into_iter() 1268 .find(|source| source.source_id.as_ref() == source_request.source_id()) 1269 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure))?; 1270 let target_fingerprint = source.target.fingerprint().clone(); 1271 let targets = TargetSet::new(vec![source.target]) 1272 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 1273 let selector = FetchSelector::all() 1274 .with_kinds(TRADE_EVENT_KINDS.to_vec()) 1275 .and_then(|selector| selector.with_exact_tag_value('d', source_request.trade_id().to_hex())) 1276 .and_then(|selector| selector.with_since_unix_seconds(replay.since_unix_seconds())) 1277 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 1278 let started_at = reconciliation_now_with(time_entropy)?; 1279 let request_id = format!( 1280 "rhi-reconcile-{}", 1281 lower_hex(source_request.id().as_bytes()) 1282 ); 1283 let maximum_events = usize::try_from(source_request.maximum_events()) 1284 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 1285 let mut adapter_cursor = None::<FetchCursor>; 1286 let mut seen_cursors = BTreeSet::new(); 1287 let mut events = Vec::<RhiAdmittedTradeMutationEvent>::new(); 1288 let mut original_bytes = 0_u64; 1289 let mut completion = 'pages: loop { 1290 if cancellation.is_cancelled() { 1291 break RhiTradeSourceCompletion::IncompleteTimeout; 1292 } 1293 let remaining = maximum_events.saturating_sub(events.len()); 1294 let limit = usize::min( 1295 usize::from(FETCH_PAGE_MAX_EVENTS), 1296 remaining.saturating_add(1), 1297 ); 1298 let bounds = FetchBounds::new( 1299 u16::try_from(limit).map_err(|_| HostError::new(HostErrorKind::TaskFailure))?, 1300 source_request.deadline().get(), 1301 ) 1302 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 1303 let mut request = FetchRequest::new(request_id.clone(), targets.clone(), bounds) 1304 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))? 1305 .with_selector(selector.clone()); 1306 if let Some(cursor) = adapter_cursor.take() { 1307 request = request.with_cursor(cursor); 1308 } 1309 let fetched = tokio::select! { 1310 () = cancellation.cancelled() => { 1311 break 'pages RhiTradeSourceCompletion::IncompleteTimeout; 1312 } 1313 fetched = transport.evidence_source().fetch(request.clone()) => fetched, 1314 }; 1315 let page = match fetched { 1316 Ok(page) => page, 1317 Err(radroots_transport::Error::UnsupportedOperation) => { 1318 events.clear(); 1319 break RhiTradeSourceCompletion::Unsupported; 1320 } 1321 Err(_) => break RhiTradeSourceCompletion::IncompleteUnavailable, 1322 }; 1323 if page.validate_for_request(&request).is_err() { 1324 break RhiTradeSourceCompletion::IncompleteUnknown; 1325 } 1326 let Some(outcome) = page 1327 .target_outcomes() 1328 .iter() 1329 .find(|outcome| outcome.target() == &target_fingerprint) 1330 .filter(|_| page.target_outcomes().len() == 1) 1331 else { 1332 break RhiTradeSourceCompletion::IncompleteUnknown; 1333 }; 1334 match outcome.state() { 1335 FetchTargetState::Complete | FetchTargetState::Partial => {} 1336 FetchTargetState::Unavailable | FetchTargetState::FailedRetryable => { 1337 break RhiTradeSourceCompletion::IncompleteUnavailable; 1338 } 1339 FetchTargetState::Cancelled => { 1340 break RhiTradeSourceCompletion::IncompleteTimeout; 1341 } 1342 FetchTargetState::FailedTerminal => { 1343 break RhiTradeSourceCompletion::IncompleteUnknown; 1344 } 1345 } 1346 for observed in page.events() { 1347 let bytes = u64::try_from(observed.event().raw_json().len()) 1348 .map_err(|_| HostError::new(HostErrorKind::TaskFailure))?; 1349 original_bytes = original_bytes.saturating_add(bytes); 1350 if events.len() >= maximum_events || original_bytes > source_request.maximum_bytes() { 1351 break 'pages RhiTradeSourceCompletion::IncompleteResourceLimit; 1352 } 1353 let observed_at = RhiTradeMutationObservedAtUnixSeconds::new( 1354 observed.provenance().observed_at_unix_ms() / 1_000, 1355 ) 1356 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1357 if let Ok(admitted) = admit_rhi_trade_mutation_event( 1358 RhiTradeMutationAdmissionLimits::from_config(configuration) 1359 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?, 1360 observed.event().raw_json().as_bytes(), 1361 observed_at, 1362 RhiTradeMutationAuthoredTimePolicy::new(0) 1363 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?, 1364 ) && admitted.mutation().trade_id == source_request.trade_id() 1365 { 1366 events.push(admitted); 1367 } 1368 } 1369 if outcome.state() == FetchTargetState::Partial { 1370 break RhiTradeSourceCompletion::IncompleteUnknown; 1371 } 1372 match page.next_page() { 1373 NextPage::Complete => break RhiTradeSourceCompletion::Complete, 1374 NextPage::Cancelled { .. } => break RhiTradeSourceCompletion::IncompleteTimeout, 1375 NextPage::Cursor(cursor) => { 1376 if page.events().is_empty() || !seen_cursors.insert(cursor.as_str().to_owned()) { 1377 break RhiTradeSourceCompletion::IncompleteUnknown; 1378 } 1379 adapter_cursor = Some(cursor.clone()); 1380 } 1381 } 1382 }; 1383 let mut finished_at = reconciliation_now_with(time_entropy)?; 1384 if finished_at >= source_request.deadline() { 1385 completion = RhiTradeSourceCompletion::IncompleteTimeout; 1386 finished_at = source_request.deadline(); 1387 } 1388 replay 1389 .finish(source_request, completion, started_at, finished_at, events) 1390 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error)) 1391 } 1392 1393 fn spawn_publication_worker( 1394 supervisor: &mut radroots_service_host::TaskSupervisor, 1395 state: Arc<RhiStateHost>, 1396 configuration: Arc<RhiConfigDocumentV1>, 1397 authority: Arc<RhiPublicationAuthority>, 1398 sink: Arc<RhiNostrExactSink>, 1399 time_entropy: RhiTimeEntropyAdapters, 1400 ) -> Result<(), RhiDaemonError> { 1401 supervisor 1402 .spawn( 1403 task_metadata( 1404 TASK_PUBLICATION_WORKER, 1405 TaskClassification::Critical, 1406 ShutdownPhase::PersistRecoverableWork, 1407 )?, 1408 move |cancellation| async move { 1409 let owner = 1410 RhiPublicationLeaseOwner::from_bytes(nonzero_entropy_16(&time_entropy)?) 1411 .map_err(|error| { 1412 HostError::with_source(HostErrorKind::TaskFailure, error) 1413 })?; 1414 let idle = worker_idle(&configuration)?; 1415 loop { 1416 if cancellation.is_cancelled() { 1417 return Ok(()); 1418 } 1419 let result = state 1420 .repositories() 1421 .publication_outbox() 1422 .execute_next_publication(owner, &time_entropy, sink.as_ref(), &authority) 1423 .await 1424 .map_err(|error| { 1425 HostError::with_source(HostErrorKind::TaskFailure, error) 1426 })?; 1427 if result.is_none() { 1428 tokio::select! { 1429 () = cancellation.cancelled() => return Ok(()), 1430 () = tokio::time::sleep(idle) => {} 1431 } 1432 } 1433 } 1434 }, 1435 ) 1436 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 1437 } 1438 1439 fn spawn_presence_worker( 1440 supervisor: &mut radroots_service_host::TaskSupervisor, 1441 state: Arc<RhiStateHost>, 1442 configuration: Arc<RhiConfigDocumentV1>, 1443 sink: Arc<RhiNostrExactSink>, 1444 time_entropy: RhiTimeEntropyAdapters, 1445 ) -> Result<(), RhiDaemonError> { 1446 supervisor 1447 .spawn( 1448 task_metadata( 1449 TASK_PRESENCE_WORKER, 1450 TaskClassification::Critical, 1451 ShutdownPhase::PersistRecoverableWork, 1452 )?, 1453 move |cancellation| async move { 1454 let owner = RhiPresenceLeaseOwner::from_bytes(nonzero_entropy_16(&time_entropy)?) 1455 .map_err(|error| { 1456 HostError::with_source(HostErrorKind::TaskFailure, error) 1457 })?; 1458 let idle = worker_idle(&configuration)?; 1459 loop { 1460 if cancellation.is_cancelled() { 1461 return Ok(()); 1462 } 1463 let result = state 1464 .repositories() 1465 .presence_outbox() 1466 .execute_next_presence(owner, &time_entropy, sink.as_ref()) 1467 .await 1468 .map_err(|error| { 1469 HostError::with_source(HostErrorKind::TaskFailure, error) 1470 })?; 1471 if result.is_none() { 1472 tokio::select! { 1473 () = cancellation.cancelled() => return Ok(()), 1474 () = tokio::time::sleep(idle) => {} 1475 } 1476 } 1477 } 1478 }, 1479 ) 1480 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 1481 } 1482 1483 fn task_metadata( 1484 name: &'static str, 1485 classification: TaskClassification, 1486 phase: ShutdownPhase, 1487 ) -> Result<TaskMetadata, RhiDaemonError> { 1488 TaskMetadata::new( 1489 TaskName::new(name).map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime))?, 1490 classification, 1491 Some(phase), 1492 ) 1493 .map_err(|_| RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 1494 } 1495 1496 fn runtime_build_info() -> Result<RhiStatusBuildInfoV1, HostError> { 1497 let service_revision = option_env!("RADROOTS_SERVICE_REVISION"); 1498 let lib_revision = option_env!("RADROOTS_LIB_REVISION"); 1499 let rust_version = option_env!("RADROOTS_RUST_VERSION"); 1500 let target = option_env!("RADROOTS_BUILD_TARGET"); 1501 let mode = if service_revision.is_some() 1502 && lib_revision.is_some() 1503 && rust_version.is_some() 1504 && target.is_some() 1505 { 1506 RhiStatusBuildMode::Release 1507 } else { 1508 RhiStatusBuildMode::Development 1509 }; 1510 RhiStatusBuildInfoV1::new( 1511 mode, 1512 Some(env!("CARGO_PKG_VERSION")), 1513 service_revision, 1514 lib_revision, 1515 rust_version, 1516 target, 1517 Some("service-host"), 1518 ) 1519 .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error)) 1520 } 1521 1522 fn operations_enabled(configuration: &RhiConfigDocumentV1) -> Result<bool, RhiDaemonError> { 1523 configuration 1524 .normalized() 1525 .pointer("/operations/enabled") 1526 .and_then(Value::as_bool) 1527 .ok_or_else(|| RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 1528 } 1529 1530 fn maximum_source_deadline(configuration: &RhiConfigDocumentV1) -> Result<u64, HostError> { 1531 configuration 1532 .normalized() 1533 .pointer("/evidence/sources") 1534 .and_then(Value::as_array) 1535 .and_then(|sources| { 1536 sources 1537 .iter() 1538 .filter_map(|source| source.pointer("/deadline_ms").and_then(Value::as_u64)) 1539 .max() 1540 }) 1541 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure)) 1542 } 1543 1544 fn configuration_integer( 1545 configuration: &RhiConfigDocumentV1, 1546 pointer: &str, 1547 ) -> Result<u64, RhiDaemonError> { 1548 configuration 1549 .normalized() 1550 .pointer(pointer) 1551 .and_then(Value::as_u64) 1552 .ok_or_else(|| RhiDaemonError::new(RhiDaemonErrorKind::Runtime)) 1553 } 1554 1555 fn worker_idle(configuration: &RhiConfigDocumentV1) -> Result<Duration, HostError> { 1556 configuration 1557 .normalized() 1558 .pointer("/reconciliation/initial_backoff_ms") 1559 .and_then(Value::as_u64) 1560 .map(Duration::from_millis) 1561 .ok_or_else(|| HostError::new(HostErrorKind::TaskFailure)) 1562 } 1563 1564 fn reconciliation_now_with( 1565 time_entropy: &RhiTimeEntropyAdapters, 1566 ) -> Result<RhiReconciliationUnixMilliseconds, HostError> { 1567 let milliseconds = time_entropy 1568 .now_utc_milliseconds() 1569 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1570 RhiReconciliationUnixMilliseconds::new(milliseconds) 1571 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error)) 1572 } 1573 1574 fn subscription_attempt_id(cursor: &str) -> String { 1575 let digest = Sha256::digest(cursor.as_bytes()); 1576 format!("subscription-{}", lower_hex(&digest)) 1577 } 1578 1579 fn nonzero_entropy_16(time_entropy: &RhiTimeEntropyAdapters) -> Result<[u8; 16], HostError> { 1580 for _ in 0..4 { 1581 let mut bytes = [0_u8; 16]; 1582 time_entropy 1583 .entropy() 1584 .fill_bytes(&mut bytes) 1585 .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; 1586 if bytes.iter().any(|byte| *byte != 0) { 1587 return Ok(bytes); 1588 } 1589 } 1590 Err(HostError::new(HostErrorKind::TaskFailure)) 1591 } 1592 1593 fn lower_hex(bytes: &[u8]) -> String { 1594 use core::fmt::Write as _; 1595 1596 let mut encoded = String::with_capacity(bytes.len().saturating_mul(2)); 1597 for byte in bytes { 1598 write!(&mut encoded, "{byte:02x}").expect("writing to String cannot fail"); 1599 } 1600 encoded 1601 } 1602 1603 #[cfg(test)] 1604 mod tests { 1605 use super::*; 1606 1607 fn empty_status_context() -> RuntimeStatusContext { 1608 let time_entropy = RhiTimeEntropyAdapters::system(); 1609 RuntimeStatusContext { 1610 configuration_digest: "0".repeat(64), 1611 configuration_source: RhiStatusConfigurationSource::DerivedRepoLocal, 1612 schema_version: crate::RHI_STATE_SCHEMA_VERSION, 1613 generation: 1, 1614 source_count: 1, 1615 started_at: time_entropy.now_monotonic(), 1616 time_entropy, 1617 reconciliation: RhiReconciliationStatusV1::default(), 1618 publication: RhiPublicationStatusV1::default(), 1619 presence: RhiPresenceStatusV1::default(), 1620 } 1621 } 1622 1623 #[test] 1624 fn runtime_status_keeps_optional_operations_degradation_ready_and_reasoned() { 1625 let context = empty_status_context(); 1626 let observation = context 1627 .observation(RhiServicePhase::Degraded, true, false) 1628 .expect("degraded observation"); 1629 let (_, reader) = crate::rhi_status_cache( 1630 radroots_runtime_paths::InstanceId::new("primary").expect("instance"), 1631 observation, 1632 ) 1633 .expect("status cache"); 1634 let snapshot = reader.snapshot(); 1635 assert!(snapshot.is_ready()); 1636 let json = std::str::from_utf8(snapshot.detailed_status_json()).expect("status UTF-8"); 1637 assert!(json.contains("\"operations_listener_failed\"")); 1638 assert!(!json.contains("\"source_unavailable\"")); 1639 } 1640 1641 #[cfg(any(target_os = "linux", target_os = "macos"))] 1642 #[tokio::test] 1643 async fn runtime_status_refresh_reads_the_exact_empty_durable_work_summary() { 1644 use std::{fs, os::unix::fs::PermissionsExt as _}; 1645 1646 use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity}; 1647 use radroots_storage::event::SourceGeneration; 1648 1649 let root = tempfile::tempdir().expect("test root"); 1650 let root_text = root.path().to_str().expect("UTF-8 root"); 1651 let invocation = crate::parse_rhi_cli_v1_from([ 1652 "rhi", 1653 "--profile", 1654 "repo-local", 1655 "--instance", 1656 "primary", 1657 "--repo-local-root", 1658 root_text, 1659 "run", 1660 ]) 1661 .expect("invocation"); 1662 let runtime = crate::resolve_rhi_runtime_context( 1663 &crate::RadrootsPathResolver::new( 1664 crate::RadrootsPlatform::Linux, 1665 crate::RadrootsHostEnvironment::default(), 1666 ), 1667 &invocation, 1668 ) 1669 .expect("runtime"); 1670 fs::create_dir_all(runtime.context().paths().state()).expect("state root"); 1671 fs::set_permissions( 1672 runtime.context().paths().state(), 1673 fs::Permissions::from_mode(0o700), 1674 ) 1675 .expect("state mode"); 1676 let configuration = crate::parse_rhi_config_v1( 1677 include_bytes!("../contracts/services_hardening/config.v1.example.toml"), 1678 crate::RhiConfigProfile::RepoLocal, 1679 ) 1680 .expect("configuration"); 1681 let metadata = crate::RhiStateMetadata::new( 1682 &runtime, 1683 &configuration, 1684 SourceGeneration::new([0x41; 32]).expect("generation"), 1685 1_725_000_000_000, 1686 ) 1687 .expect("metadata"); 1688 let applied_at = MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("time"); 1689 let build = MigrationBuildIdentity::new( 1690 env!("CARGO_PKG_VERSION"), 1691 "1111111111111111111111111111111111111111", 1692 "053d0c750bf9cd683c6ea37cefe7e79617ba629f", 1693 "rustc-test", 1694 "test-target", 1695 "service-host", 1696 1, 1697 crate::RHI_STATE_SCHEMA_VERSION, 1698 1, 1699 1, 1700 1, 1701 ) 1702 .expect("build"); 1703 crate::initialize_rhi_state(&runtime, &metadata, applied_at, &build) 1704 .await 1705 .expect("initialize"); 1706 let state = crate::open_rhi_state_read_write(&runtime, &metadata, applied_at, &build) 1707 .await 1708 .expect("open"); 1709 let mut status = empty_status_context(); 1710 status.refresh_state(&state).await.expect("refresh"); 1711 assert_eq!(status.reconciliation, RhiReconciliationStatusV1::default()); 1712 assert_eq!(status.publication, RhiPublicationStatusV1::default()); 1713 assert_eq!(status.presence, RhiPresenceStatusV1::default()); 1714 state.close().await.expect("close"); 1715 } 1716 }