myc

Self-custodial remote signer for Radroots apps
git clone https://radroots.dev/git/myc.git
Log | Files | Refs | README | LICENSE

runtime_graph.rs (107189B)


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