rhi

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

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 }