rhi

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

commit 758ab117f9d94130389d2f0372117625a1ec8563
parent d622a3569eda032cc425bc4d853fb1ccdcf4631e
Author: triesap <tyson@radroots.org>
Date:   Mon, 24 Aug 2026 17:55:13 +0000

status: add passive RHI operations cache

- freeze bounded local status and service-specific lifecycle summaries
- expose only cached liveness, readiness, and metrics over optional TCP
- retain redacted root-only API and exact machine contracts

Diffstat:
MAGENTS.md | 5+++++
MCargo.toml | 2+-
MREADME | 23+++++++++++++++++++++++
Mcontracts/api_baselines/rhi.txt | 212+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acontracts/services_hardening/status_cache.v1.json | 48++++++++++++++++++++++++++++++++++++++++++++++++
Acontracts/services_hardening/tcp_operations.v1.json | 53+++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/lib.rs | 18++++++++++++++++++
Asrc/operations_v1.rs | 267+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Asrc/status_v1.rs | 1180+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/package_boundary.rs | 111++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Atests/services_hardening_operations.rs | 233+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Atests/services_hardening_status.rs | 501+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
12 files changed, 2651 insertions(+), 2 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -424,6 +424,11 @@ - Optional TCP operations expose only cached `/livez`, `/readyz`, and `/metrics`; requests must not perform SQLite, source, relay, DNS, identity, evidence, or credential probes. +- Step 212 owns the one-latest passive lifecycle/status cache and the optional + TCP adapter. Detailed status remains permissioned Unix-admin-only. The TCP + surface is exactly cached `/livez`, `/readyz`, and `/metrics`, with two fixed + RHI metric families and no route-registration extension, active probe, or + high-cardinality label authority. - Keep logs as safe structured stderr output. Keep result data on stdout and diagnostics on stderr. Use stable bounded public codes/messages, bounded metric labels, explicit redaction, and no trade/mutation/event/report IDs or diff --git a/Cargo.toml b/Cargo.toml @@ -81,7 +81,7 @@ url = "2" zeroize = { version = "1" } [dev-dependencies] -tokio = { version = "1", default-features = false, features = ["macros", "rt-multi-thread"] } +tokio = { version = "1", default-features = false, features = ["io-util", "macros", "net", "rt-multi-thread"] } [profile.release] lto = "thin" diff --git a/README b/README @@ -499,6 +499,29 @@ no signals, runtime, logger, or process-exit policy and does not claim the final supervised task graph. Its exact machine contract is [`runtime_foundation.v1.json`](contracts/services_hardening/runtime_foundation.v1.json). +## Passive lifecycle status and TCP operations + +`rhi_status_cache` retains exactly one latest immutable RHI lifecycle and +detailed-status publication behind a single non-cloneable publisher. Each +publication is bounded and validated before it replaces the prior snapshot. +Readers clone only passive cache authority; they perform no SQLite, +filesystem, evidence-source, relay, DNS, identity, credential, clock, or fresh +probe operation. Detailed JSON is reserved for the permissioned Unix-admin +surface and contains only the closed service identity, evidence transport, +reconciliation, publication, presence, persistence, configuration, build, and +stable reason vocabularies. Its exact contract is +[`status_cache.v1.json`](contracts/services_hardening/status_cache.v1.json). + +When the validated operations block is enabled, `RhiOperationsServer` exposes +exactly HTTP/1.1 `GET /livez`, `GET /readyz`, and `GET /metrics` on its admitted +TCP address. Requests read only the latest cached lifecycle/readiness and two +fixed bounded metric families; they cannot register another route or reach +detailed status. Requests perform no SQLite, filesystem, evidence-source, +relay, DNS, identity, credential, clock, or active health probe. The listener +is disabled by default and remains a supervisor-owned optional adapter. Its +exact contract is +[`tcp_operations.v1.json`](contracts/services_hardening/tcp_operations.v1.json). + ## Hardened v1 configuration contract The target service configuration is frozen by diff --git a/contracts/api_baselines/rhi.txt b/contracts/api_baselines/rhi.txt @@ -267,8 +267,26 @@ pub enum rhi::RhiIdentityRole pub rhi::RhiIdentityRole::Service impl rhi::RhiIdentityRole pub const fn rhi::RhiIdentityRole::as_str(self) -> &'static str +pub enum rhi::RhiIntegrityStateV1 +pub rhi::RhiIntegrityStateV1::Failed +pub rhi::RhiIntegrityStateV1::VerificationRequired +pub rhi::RhiIntegrityStateV1::Verified pub enum rhi::RhiMetricsCommandV1 pub rhi::RhiMetricsCommandV1::Snapshot +pub enum rhi::RhiOperationsErrorKind +pub rhi::RhiOperationsErrorKind::Accept +pub rhi::RhiOperationsErrorKind::Bind +pub rhi::RhiOperationsErrorKind::ConnectionTaskPanicked +pub rhi::RhiOperationsErrorKind::Disabled +pub rhi::RhiOperationsErrorKind::InvalidConfiguration +pub rhi::RhiOperationsErrorKind::LocalAddress +impl rhi::RhiOperationsErrorKind +pub const fn rhi::RhiOperationsErrorKind::code(self) -> &'static str +pub enum rhi::RhiPersistenceHealthV1 +pub rhi::RhiPersistenceHealthV1::ReadOnly +pub rhi::RhiPersistenceHealthV1::Ready +pub rhi::RhiPersistenceHealthV1::RepairRequired +pub rhi::RhiPersistenceHealthV1::Unavailable pub enum rhi::RhiPresenceAttemptOutcome pub rhi::RhiPresenceAttemptOutcome::Accepted pub rhi::RhiPresenceAttemptOutcome::AuthRequired @@ -350,6 +368,9 @@ impl rhi::RhiProcessResult pub const fn rhi::RhiProcessResult::code(self) -> &'static str pub fn rhi::RhiProcessResult::exit_code(self) -> std::process::ExitCode pub const fn rhi::RhiProcessResult::exit_code_u8(self) -> u8 +pub enum rhi::RhiProviderHealthV1 +pub rhi::RhiProviderHealthV1::Ready +pub rhi::RhiProviderHealthV1::Unavailable pub enum rhi::RhiPublicationAttemptEvidenceErrorKind pub rhi::RhiPublicationAttemptEvidenceErrorKind::InvalidAttemptNumber pub rhi::RhiPublicationAttemptEvidenceErrorKind::InvalidTargetOrdinal @@ -576,6 +597,13 @@ pub rhi::RhiRuntimeReadinessReason::SourceUnavailable pub rhi::RhiRuntimeReadinessReason::SubscriptionInactive impl rhi::RhiRuntimeReadinessReason pub const fn rhi::RhiRuntimeReadinessReason::as_str(self) -> &'static str +pub enum rhi::RhiServicePhase +pub rhi::RhiServicePhase::Degraded +pub rhi::RhiServicePhase::Failed +pub rhi::RhiServicePhase::Ready +pub rhi::RhiServicePhase::Starting +pub rhi::RhiServicePhase::Stopping +pub rhi::RhiServicePhase::Unready pub enum rhi::RhiSourcesCommandV1 pub rhi::RhiSourcesCommandV1::List pub enum rhi::RhiStateCatalogErrorKind @@ -658,6 +686,47 @@ pub rhi::RhiStateRepositoryWriteClass::CompareAndSwap pub rhi::RhiStateRepositoryWriteClass::Immutable impl rhi::RhiStateRepositoryWriteClass pub const fn rhi::RhiStateRepositoryWriteClass::code(self) -> &'static str +pub enum rhi::RhiStatusBuildMode +pub rhi::RhiStatusBuildMode::Development +pub rhi::RhiStatusBuildMode::Release +pub enum rhi::RhiStatusConfigurationSource +pub rhi::RhiStatusConfigurationSource::DerivedRepoLocal +pub rhi::RhiStatusConfigurationSource::ExplicitConfig +pub enum rhi::RhiStatusErrorKind +pub rhi::RhiStatusErrorKind::Encoding +pub rhi::RhiStatusErrorKind::InvalidBuildInfo +pub rhi::RhiStatusErrorKind::InvalidConfiguration +pub rhi::RhiStatusErrorKind::InvalidIdentityHealth +pub rhi::RhiStatusErrorKind::InvalidLifecycle +pub rhi::RhiStatusErrorKind::InvalidModel +pub rhi::RhiStatusErrorKind::InvalidPersistence +pub rhi::RhiStatusErrorKind::InvalidProviderState +pub rhi::RhiStatusErrorKind::InvalidReasonCode +pub rhi::RhiStatusErrorKind::InvalidTime +pub rhi::RhiStatusErrorKind::InvalidTransition +pub rhi::RhiStatusErrorKind::InvalidTransportState +pub rhi::RhiStatusErrorKind::PublisherDropped +pub rhi::RhiStatusErrorKind::ResponseTooLarge +pub rhi::RhiStatusErrorKind::TooManyReasonCodes +impl rhi::RhiStatusErrorKind +pub const fn rhi::RhiStatusErrorKind::code(self) -> &'static str +pub enum rhi::RhiStatusReasonCode +pub rhi::RhiStatusReasonCode::AdminListenerFailed +pub rhi::RhiStatusReasonCode::DatabaseLowDisk +pub rhi::RhiStatusReasonCode::DatabaseReadOnly +pub rhi::RhiStatusReasonCode::DatabaseSchemaMismatch +pub rhi::RhiStatusReasonCode::IdentityUnavailable +pub rhi::RhiStatusReasonCode::OperationsListenerFailed +pub rhi::RhiStatusReasonCode::PresenceStateUnavailable +pub rhi::RhiStatusReasonCode::PublicationRecoveryIncomplete +pub rhi::RhiStatusReasonCode::ReconciliationBacklogExceeded +pub rhi::RhiStatusReasonCode::RecoveryIncomplete +pub rhi::RhiStatusReasonCode::ShutdownInProgress +pub rhi::RhiStatusReasonCode::SourceUnavailable +pub rhi::RhiStatusReasonCode::SubscriptionInactive +impl rhi::RhiStatusReasonCode +pub const fn rhi::RhiStatusReasonCode::as_str(self) -> &'static str +pub fn rhi::RhiStatusReasonCode::new(impl core::convert::AsRef<str>) -> core::result::Result<Self, rhi::RhiStatusError> pub enum rhi::RhiTradeCommandV1 pub rhi::RhiTradeCommandV1::Projection pub rhi::RhiTradeCommandV1::ReportCurrent @@ -714,6 +783,10 @@ pub rhi::RhiTradeSourceIngestErrorKind::InvalidMode pub rhi::RhiTradeSourceIngestErrorKind::Storage impl rhi::RhiTradeSourceIngestErrorKind pub const fn rhi::RhiTradeSourceIngestErrorKind::code(self) -> &'static str +pub enum rhi::RhiTransportHealthV1 +pub rhi::RhiTransportHealthV1::Degraded +pub rhi::RhiTransportHealthV1::Ready +pub rhi::RhiTransportHealthV1::Unavailable pub enum rhi::TradeAgreementAttestationBackend pub rhi::TradeAgreementAttestationBackend::LocalStatementHash impl rhi::TradeAgreementAttestationBackend @@ -827,6 +900,12 @@ impl rhi::RhiBoundAdminServer pub async fn rhi::RhiBoundAdminServer::serve(self, rhi::RhiAdminCancellationToken) -> core::result::Result<(), rhi::RhiAdminServerError> impl core::fmt::Debug for rhi::RhiBoundAdminServer pub fn rhi::RhiBoundAdminServer::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiBoundOperationsServer +impl rhi::RhiBoundOperationsServer +pub fn rhi::RhiBoundOperationsServer::local_address(&self) -> core::net::socket_addr::SocketAddr +pub async fn rhi::RhiBoundOperationsServer::serve(self, rhi::RhiOperationsCancellationToken) -> core::result::Result<(), rhi::RhiOperationsError> +impl core::fmt::Debug for rhi::RhiBoundOperationsServer +pub fn rhi::RhiBoundOperationsServer::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub struct rhi::RhiCliExecutionPlanV1 impl rhi::RhiCliExecutionPlanV1 pub const fn rhi::RhiCliExecutionPlanV1::admin_operation(&self) -> core::option::Option<rhi::RhiCliAdminOperationV1> @@ -989,6 +1068,15 @@ impl rhi::RhiEvidencePolicyDigest pub const fn rhi::RhiEvidencePolicyDigest::as_bytes(&self) -> &[u8; 32] impl core::fmt::Debug for rhi::RhiEvidencePolicyDigest pub fn rhi::RhiEvidencePolicyDigest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiEvidenceTransportStatusV1 +impl rhi::RhiEvidenceTransportStatusV1 +pub const fn rhi::RhiEvidenceTransportStatusV1::configured_source_count(&self) -> u64 +pub const fn rhi::RhiEvidenceTransportStatusV1::health(&self) -> rhi::RhiTransportHealthV1 +pub fn rhi::RhiEvidenceTransportStatusV1::new(rhi::RhiTransportHealthV1, bool, bool, u64, u64, rhi::RhiStatusReasonCodes) -> core::result::Result<Self, rhi::RhiStatusError> +pub const fn rhi::RhiEvidenceTransportStatusV1::reachable_source_count(&self) -> u64 +pub const fn rhi::RhiEvidenceTransportStatusV1::reason_codes(&self) -> &rhi::RhiStatusReasonCodes +pub const fn rhi::RhiEvidenceTransportStatusV1::required_sources_ready(&self) -> bool +pub const fn rhi::RhiEvidenceTransportStatusV1::subscriber_active(&self) -> bool pub struct rhi::RhiExpectedPublicIdentity(_) impl rhi::RhiExpectedPublicIdentity pub fn rhi::RhiExpectedPublicIdentity::as_hex(&self) -> &str @@ -1009,6 +1097,12 @@ pub const fn rhi::RhiIdentityEnvelopeBinding::kind(&self) -> rhi::RhiIdentityPro pub const fn rhi::RhiIdentityEnvelopeBinding::role(&self) -> rhi::RhiIdentityRole impl core::fmt::Debug for rhi::RhiIdentityEnvelopeBinding pub fn rhi::RhiIdentityEnvelopeBinding::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiIdentityHealthV1 +impl rhi::RhiIdentityHealthV1 +pub const fn rhi::RhiIdentityHealthV1::is_available(&self) -> bool +pub const fn rhi::RhiIdentityHealthV1::is_configured(&self) -> bool +pub fn rhi::RhiIdentityHealthV1::new(bool, bool, rhi::RhiStatusReasonCodes) -> core::result::Result<Self, rhi::RhiStatusError> +pub const fn rhi::RhiIdentityHealthV1::reason_codes(&self) -> &rhi::RhiStatusReasonCodes pub struct rhi::RhiJitterBoundMilliseconds(_) impl rhi::RhiJitterBoundMilliseconds pub const fn rhi::RhiJitterBoundMilliseconds::get(self) -> u64 @@ -1028,6 +1122,33 @@ impl rhi::RhiNormalizedConfigDigest pub const fn rhi::RhiNormalizedConfigDigest::as_bytes(&self) -> &[u8; 32] impl core::fmt::Debug for rhi::RhiNormalizedConfigDigest pub fn rhi::RhiNormalizedConfigDigest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiOperationsCancellationToken +impl rhi::RhiOperationsCancellationToken +pub fn rhi::RhiOperationsCancellationToken::cancel(&self) +pub fn rhi::RhiOperationsCancellationToken::is_cancelled(&self) -> bool +pub fn rhi::RhiOperationsCancellationToken::new() -> Self +impl core::fmt::Debug for rhi::RhiOperationsCancellationToken +pub fn rhi::RhiOperationsCancellationToken::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiOperationsError +impl rhi::RhiOperationsError +pub const fn rhi::RhiOperationsError::code(self) -> &'static str +pub const fn rhi::RhiOperationsError::kind(self) -> rhi::RhiOperationsErrorKind +impl core::error::Error for rhi::RhiOperationsError +impl core::fmt::Debug for rhi::RhiOperationsError +pub fn rhi::RhiOperationsError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl core::fmt::Display for rhi::RhiOperationsError +pub fn rhi::RhiOperationsError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiOperationsServer +impl rhi::RhiOperationsServer +pub async fn rhi::RhiOperationsServer::bind(self) -> core::result::Result<rhi::RhiBoundOperationsServer, rhi::RhiOperationsError> +pub fn rhi::RhiOperationsServer::new(&rhi::RhiConfigDocumentV1, &rhi::RhiStatusReader) -> core::result::Result<Self, rhi::RhiOperationsError> +impl core::fmt::Debug for rhi::RhiOperationsServer +pub fn rhi::RhiOperationsServer::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPersistenceStatusV1 +impl rhi::RhiPersistenceStatusV1 +pub fn rhi::RhiPersistenceStatusV1::new(rhi::RhiPersistenceHealthV1, u32, u64, rhi::RhiIntegrityStateV1, rhi::RhiStatusReasonCodes) -> core::result::Result<Self, rhi::RhiStatusError> +impl core::fmt::Debug for rhi::RhiPersistenceStatusV1 +pub fn rhi::RhiPersistenceStatusV1::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub struct rhi::RhiPreparedPresenceAttempt impl rhi::RhiPreparedPresenceAttempt pub const fn rhi::RhiPreparedPresenceAttempt::attempt_id(&self) -> rhi::RhiPresenceAttemptId @@ -1150,6 +1271,11 @@ pub struct rhi::RhiPresenceRetryDelayMilliseconds(_) impl rhi::RhiPresenceRetryDelayMilliseconds pub const fn rhi::RhiPresenceRetryDelayMilliseconds::get(self) -> u64 pub fn rhi::RhiPresenceRetryDelayMilliseconds::new(u64) -> core::result::Result<Self, rhi::RhiPresencePublicationError> +pub struct rhi::RhiPresenceStatusV1 +impl rhi::RhiPresenceStatusV1 +pub const fn rhi::RhiPresenceStatusV1::new(u64, u64) -> Self +pub const fn rhi::RhiPresenceStatusV1::pending(self) -> u64 +pub const fn rhi::RhiPresenceStatusV1::unknown(self) -> u64 pub struct rhi::RhiPresenceTarget impl rhi::RhiPresenceTarget pub const fn rhi::RhiPresenceTarget::ordinal(&self) -> u8 @@ -1179,6 +1305,12 @@ pub const fn rhi::RhiProvenanceRepository<'_>::descriptor(&self) -> rhi::RhiStat pub const fn rhi::RhiProvenanceRepository<'_>::kind(&self) -> rhi::RhiStateRepositoryKind impl core::fmt::Debug for rhi::RhiProvenanceRepository<'_> pub fn rhi::RhiProvenanceRepository<'_>::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiProviderStatusV1 +impl rhi::RhiProviderStatusV1 +pub const fn rhi::RhiProviderStatusV1::health(&self) -> rhi::RhiProviderHealthV1 +pub const fn rhi::RhiProviderStatusV1::identity(&self) -> &rhi::RhiIdentityHealthV1 +pub fn rhi::RhiProviderStatusV1::new(rhi::RhiIdentityHealthV1, rhi::RhiStatusReasonCodes) -> core::result::Result<Self, rhi::RhiStatusError> +pub const fn rhi::RhiProviderStatusV1::reason_codes(&self) -> &rhi::RhiStatusReasonCodes pub struct rhi::RhiPublicationAttemptCommit impl rhi::RhiPublicationAttemptCommit pub const fn rhi::RhiPublicationAttemptCommit::attempt_id(self) -> rhi::RhiPublicationAttemptId @@ -1294,6 +1426,12 @@ pub const fn rhi::RhiPublicationRetryPolicy::attempt_deadline_milliseconds(self) pub const fn rhi::RhiPublicationRetryPolicy::initial_backoff_milliseconds(self) -> u64 pub const fn rhi::RhiPublicationRetryPolicy::maximum_attempts(self) -> u16 pub const fn rhi::RhiPublicationRetryPolicy::maximum_backoff_milliseconds(self) -> u64 +pub struct rhi::RhiPublicationStatusV1 +impl rhi::RhiPublicationStatusV1 +pub const fn rhi::RhiPublicationStatusV1::new(u64, u64, core::option::Option<rhi::RhiStatusUnixSeconds>) -> Self +pub const fn rhi::RhiPublicationStatusV1::oldest_pending_at_utc(self) -> core::option::Option<rhi::RhiStatusUnixSeconds> +pub const fn rhi::RhiPublicationStatusV1::pending(self) -> u64 +pub const fn rhi::RhiPublicationStatusV1::unknown(self) -> u64 pub struct rhi::RhiPublicationSubmissionError impl rhi::RhiPublicationSubmissionError pub const fn rhi::RhiPublicationSubmissionError::code(self) -> &'static str @@ -1629,6 +1767,13 @@ impl rhi::RhiReconciliationSourceSelectorDigest pub const fn rhi::RhiReconciliationSourceSelectorDigest::as_bytes(&self) -> &[u8; 32] impl core::fmt::Debug for rhi::RhiReconciliationSourceSelectorDigest pub fn rhi::RhiReconciliationSourceSelectorDigest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiReconciliationStatusV1 +impl rhi::RhiReconciliationStatusV1 +pub const fn rhi::RhiReconciliationStatusV1::exhausted(self) -> u64 +pub const fn rhi::RhiReconciliationStatusV1::leased(self) -> u64 +pub const fn rhi::RhiReconciliationStatusV1::new(u64, u64, u64, core::option::Option<rhi::RhiStatusUnixSeconds>) -> Self +pub const fn rhi::RhiReconciliationStatusV1::oldest_pending_at_utc(self) -> core::option::Option<rhi::RhiStatusUnixSeconds> +pub const fn rhi::RhiReconciliationStatusV1::pending(self) -> u64 pub struct rhi::RhiReconciliationUnixMilliseconds(_) impl rhi::RhiReconciliationUnixMilliseconds pub const fn rhi::RhiReconciliationUnixMilliseconds::get(self) -> u64 @@ -1865,6 +2010,65 @@ pub const fn rhi::RhiStateRepositoryDescriptor::kind(self) -> rhi::RhiStateRepos pub const fn rhi::RhiStateRepositoryDescriptor::write_class(self) -> rhi::RhiStateRepositoryWriteClass impl core::fmt::Debug for rhi::RhiStateRepositoryDescriptor pub fn rhi::RhiStateRepositoryDescriptor::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiStatusBuildInfoV1 +impl rhi::RhiStatusBuildInfoV1 +pub fn rhi::RhiStatusBuildInfoV1::new(rhi::RhiStatusBuildMode, core::option::Option<&str>, core::option::Option<&str>, core::option::Option<&str>, core::option::Option<&str>, core::option::Option<&str>, core::option::Option<&str>) -> core::result::Result<Self, rhi::RhiStatusError> +impl core::fmt::Debug for rhi::RhiStatusBuildInfoV1 +pub fn rhi::RhiStatusBuildInfoV1::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiStatusCommonV1 +impl rhi::RhiStatusCommonV1 +pub fn rhi::RhiStatusCommonV1::new(rhi::RhiServicePhase, bool, rhi::RhiStatusReasonCodes, u64, rhi::RhiStatusBuildInfoV1, rhi::RhiStatusConfigurationIdentityV1, rhi::RhiPersistenceStatusV1) -> core::result::Result<Self, rhi::RhiStatusError> +impl core::fmt::Debug for rhi::RhiStatusCommonV1 +pub fn rhi::RhiStatusCommonV1::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiStatusConfigurationIdentityV1 +impl rhi::RhiStatusConfigurationIdentityV1 +pub fn rhi::RhiStatusConfigurationIdentityV1::new(impl core::convert::AsRef<str>, rhi::RhiStatusConfigurationSource) -> core::result::Result<Self, rhi::RhiStatusError> +impl core::fmt::Debug for rhi::RhiStatusConfigurationIdentityV1 +pub fn rhi::RhiStatusConfigurationIdentityV1::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiStatusError +impl rhi::RhiStatusError +pub const fn rhi::RhiStatusError::code(self) -> &'static str +pub const fn rhi::RhiStatusError::kind(self) -> rhi::RhiStatusErrorKind +impl core::error::Error for rhi::RhiStatusError +impl core::fmt::Debug for rhi::RhiStatusError +pub fn rhi::RhiStatusError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl core::fmt::Display for rhi::RhiStatusError +pub fn rhi::RhiStatusError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiStatusObservationV1 +impl rhi::RhiStatusObservationV1 +pub fn rhi::RhiStatusObservationV1::new(rhi::RhiStatusCommonV1, rhi::RhiProviderStatusV1, rhi::RhiEvidenceTransportStatusV1, rhi::RhiReconciliationStatusV1, rhi::RhiPublicationStatusV1, rhi::RhiPresenceStatusV1) -> Self +impl core::fmt::Debug for rhi::RhiStatusObservationV1 +pub fn rhi::RhiStatusObservationV1::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiStatusPublisher +impl rhi::RhiStatusPublisher +pub fn rhi::RhiStatusPublisher::publish(&mut self, rhi::RhiStatusObservationV1) -> core::result::Result<(), rhi::RhiStatusError> +pub fn rhi::RhiStatusPublisher::subscribe(&self) -> rhi::RhiStatusReader +impl core::fmt::Debug for rhi::RhiStatusPublisher +pub fn rhi::RhiStatusPublisher::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiStatusReader +impl rhi::RhiStatusReader +pub async fn rhi::RhiStatusReader::changed(&mut self) -> core::result::Result<rhi::RhiStatusSnapshot, rhi::RhiStatusError> +pub fn rhi::RhiStatusReader::snapshot(&self) -> rhi::RhiStatusSnapshot +impl core::clone::Clone for rhi::RhiStatusReader +pub fn rhi::RhiStatusReader::clone(&self) -> Self +impl core::fmt::Debug for rhi::RhiStatusReader +pub fn rhi::RhiStatusReader::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiStatusReasonCodes(_) +impl rhi::RhiStatusReasonCodes +pub fn rhi::RhiStatusReasonCodes::as_slice(&self) -> &[rhi::RhiStatusReasonCode] +pub const fn rhi::RhiStatusReasonCodes::empty() -> Self +pub fn rhi::RhiStatusReasonCodes::new(impl core::iter::traits::collect::IntoIterator<Item = rhi::RhiStatusReasonCode>) -> core::result::Result<Self, rhi::RhiStatusError> +pub struct rhi::RhiStatusSnapshot +impl rhi::RhiStatusSnapshot +pub fn rhi::RhiStatusSnapshot::detailed_status_json(&self) -> &[u8] +pub fn rhi::RhiStatusSnapshot::is_ready(&self) -> bool +pub fn rhi::RhiStatusSnapshot::phase(&self) -> rhi::RhiServicePhase +impl core::fmt::Debug for rhi::RhiStatusSnapshot +pub fn rhi::RhiStatusSnapshot::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiStatusUnixSeconds(_) +impl rhi::RhiStatusUnixSeconds +pub const fn rhi::RhiStatusUnixSeconds::get(self) -> u64 +pub fn rhi::RhiStatusUnixSeconds::new(u64) -> core::result::Result<Self, rhi::RhiStatusError> pub struct rhi::RhiSupersessionRepository<'host> impl rhi::RhiSupersessionRepository<'_> pub const fn rhi::RhiSupersessionRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor @@ -2038,6 +2242,7 @@ pub const rhi::RHI_CONFIG_DOCUMENT_MAX_UTF8_BYTES: usize pub const rhi::RHI_CONFIG_EFFECTIVE_MAX_UTF8_BYTES: usize pub const rhi::RHI_CONFIG_SCHEMA: &str pub const rhi::RHI_CONFIG_SCHEMA_VERSION: u32 +pub const rhi::RHI_DETAILED_STATUS_MAX_UTF8_BYTES: usize pub const rhi::RHI_DOCTOR_CHECK_COUNT: usize pub const rhi::RHI_DOCTOR_CONTRACT_VERSION: u32 pub const rhi::RHI_DOCTOR_REPORT_MAX_UTF8_BYTES: usize @@ -2045,7 +2250,10 @@ pub const rhi::RHI_DOCTOR_SUMMARY_MAX_UTF8_BYTES: usize pub const rhi::RHI_ENCRYPTED_IDENTITY_BACKUP_INCLUDED: bool pub const rhi::RHI_ENCRYPTED_IDENTITY_ENVELOPE_CONTRACT_VERSION: u32 pub const rhi::RHI_ENCRYPTED_IDENTITY_ENVELOPE_MAX_BYTES: usize +pub const rhi::RHI_LIVEZ_PATH: &str +pub const rhi::RHI_METRICS_PATH: &str pub const rhi::RHI_MIGRATION_CATALOG_SHA256: [u8; 32] +pub const rhi::RHI_OPERATIONS_CONTRACT_VERSION: u32 pub const rhi::RHI_PRESENCE_DESIRED_CONTRACT_VERSION: u32 pub const rhi::RHI_PRESENCE_DESIRED_MAX_TARGETS: usize pub const rhi::RHI_PRESENCE_MAX_ATTEMPTS: u16 @@ -2060,6 +2268,7 @@ pub const rhi::RHI_PUBLICATION_MAX_ATTEMPTS: u16 pub const rhi::RHI_PUBLICATION_MAX_TARGETS: usize pub const rhi::RHI_PUBLICATION_SUBMISSION_CONTRACT_VERSION: u32 pub const rhi::RHI_PUBLICATION_TARGET_ORDINAL_MAXIMUM: u32 +pub const rhi::RHI_READYZ_PATH: &str pub const rhi::RHI_RECONCILIATION_ATTEMPT_CONTRACT_VERSION: u32 pub const rhi::RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES: usize pub const rhi::RHI_RECONCILIATION_ATTESTATION_CONTRACT_VERSION: u32 @@ -2112,7 +2321,9 @@ pub const rhi::RHI_STATE_SCHEMA_VERSION_8_SHA256: [u8; 32] pub const rhi::RHI_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256: [u8; 32] pub const rhi::RHI_STATE_SCHEMA_VERSION_9_OBJECT_COUNT: u32 pub const rhi::RHI_STATE_SCHEMA_VERSION_9_SHA256: [u8; 32] +pub const rhi::RHI_STATUS_CACHE_CONTRACT_VERSION: u32 pub const rhi::RHI_STATUS_CONTRACT_VERSION: u32 +pub const rhi::RHI_STATUS_REASON_CODE_COUNT: usize pub const rhi::RHI_TRADE_EVENT_EXTRA_FIELD_MAX_COUNT: usize pub const rhi::RHI_TRADE_EVENT_EXTRA_JSON_MAX_BYTES: usize pub const rhi::RHI_TRADE_EVENT_ID_MAX_BYTES: usize @@ -2167,6 +2378,7 @@ pub const fn rhi::rhi_doctor_check_definitions() -> &'static [rhi::RhiDoctorChec pub fn rhi::rhi_migration_catalog() -> core::result::Result<radroots_service_sqlite::migration::MigrationCatalog, rhi::RhiStateCatalogError> pub fn rhi::rhi_schema_catalog() -> core::result::Result<radroots_service_sqlite::integrity::catalog::SchemaCatalog, rhi::RhiStateCatalogError> pub const fn rhi::rhi_state_repository_descriptors() -> &'static [rhi::RhiStateRepositoryDescriptor; 21] +pub fn rhi::rhi_status_cache(radroots_runtime_paths::identifier::InstanceId, rhi::RhiStatusObservationV1) -> core::result::Result<(rhi::RhiStatusPublisher, rhi::RhiStatusReader), rhi::RhiStatusError> pub async fn rhi::run_rhi_doctor(&rhi::RhiRuntimeContext, &impl rhi::RhiDoctorProbe + ?core::marker::Sized) -> core::result::Result<rhi::RhiDoctorReport, rhi::RhiDoctorError> pub async fn rhi::stage_rhi_state_restore(&rhi::RhiRuntimeContext, &rhi::RhiStateMetadata, rhi::RhiVerifiedStateBackup) -> core::result::Result<rhi::RhiStagedStateRestore, rhi::RhiStateMaintenanceError> pub fn rhi::trade_mutation_subscription_kinds() -> alloc::vec::Vec<u32> diff --git a/contracts/services_hardening/status_cache.v1.json b/contracts/services_hardening/status_cache.v1.json @@ -0,0 +1,48 @@ +{ + "schema": "radroots.rhi.status-cache.v1", + "contract_version": 1, + "step": 212, + "shared_cache": "radroots_service_host::cached_service_state", + "publication": { + "authority": "single_non_clone_publisher", + "capacity": "one_latest_immutable_arc", + "replacement": "atomic", + "phase_transition": "shared_service_operational_state", + "failed_publication": "previous_snapshot_retained" + }, + "read": { + "kind": "passive_cached_snapshot", + "await_required": false, + "sqlite": false, + "provider": false, + "evidence_source": false, + "relay": false, + "credential": false, + "dns": false, + "clock": false, + "fresh_probe": false, + "task_spawn": false + }, + "detailed_status": { + "model": "service_status_v1", + "service": "rhi", + "local_transport": "permissioned_unix_admin_only", + "construction": "bounded_encode_before_atomic_publication", + "maximum_utf8_bytes": 1048576, + "identity_roles": ["service"], + "transport_count_keys": ["configured_source_count", "reachable_source_count"], + "reconciliation_count_keys": ["pending", "leased", "exhausted"], + "publication_count_keys": ["pending", "unknown"], + "presence_count_keys": ["pending", "unknown"], + "reason_codes": ["identity_unavailable", "database_schema_mismatch", "database_read_only", "database_low_disk", "source_unavailable", "subscription_inactive", "recovery_incomplete", "publication_recovery_incomplete", "presence_state_unavailable", "reconciliation_backlog_exceeded", "admin_listener_failed", "operations_listener_failed", "shutdown_in_progress"], + "oldest_pending_time": "optional_nonnegative_i64_range_unix_seconds_for_reconciliation_and_publication", + "protected_material": false, + "filesystem_paths": false, + "raw_errors": false + }, + "deferred": [ + "unix_admin_handler_binding", + "final_runtime_task_graph", + "signal_handling" + ] +} diff --git a/contracts/services_hardening/tcp_operations.v1.json b/contracts/services_hardening/tcp_operations.v1.json @@ -0,0 +1,53 @@ +{ + "schema": "radroots.rhi.tcp-operations.v1", + "contract_version": 1, + "step": 212, + "transport": "http_1_1_over_optional_tcp", + "configuration": "validated_rhi_config_v1_operations_block", + "route_registration_extension": false, + "routes": [ + { "method": "GET", "path": "/livez", "source": "cached_supervisor_state" }, + { "method": "GET", "path": "/readyz", "source": "cached_readiness_state" }, + { "method": "GET", "path": "/metrics", "source": "cached_bounded_metrics_snapshot" } + ], + "metrics": { + "format": "prometheus_text_0_0_4", + "families": [ + { "name": "radroots_rhi_service_phase", "kind": "gauge", "labels": ["phase"] }, + { "name": "radroots_rhi_service_ready", "kind": "gauge", "labels": [] } + ], + "high_cardinality_labels": false, + "arbitrary_labels": false + }, + "passive_read_forbidden_authority": [ + "sqlite", + "filesystem", + "provider", + "evidence_source", + "relay", + "credential", + "dns", + "fresh_probe" + ], + "local_only": [ + "detailed_status", + "configuration", + "paths", + "identities", + "sources", + "reconciliation", + "publication", + "presence", + "audit", + "mutations" + ], + "deferrals": [ + "daemon_composition", + "structured_logging_and_exit", + "integration_wave", + "rcld_promotion", + "nix", + "oci", + "terminal_consumer_convergence" + ] +} diff --git a/src/lib.rs b/src/lib.rs @@ -10,6 +10,7 @@ mod doctor_v1; mod features; mod identity_credential; mod identity_envelope; +mod operations_v1; mod presence_desired; mod presence_publication; mod process_result_v1; @@ -37,6 +38,7 @@ mod state_maintenance; mod state_metadata; mod state_repository; mod state_trade; +mod status_v1; mod trade_ingest; pub use adapters::nostr::event::NostrEventAdapter; @@ -89,6 +91,11 @@ pub use identity_envelope::{ RhiIdentityRole, RhiWrappingCredential, open_rhi_encrypted_identity, provision_rhi_encrypted_identity, }; +pub use operations_v1::{ + RHI_LIVEZ_PATH, RHI_METRICS_PATH, RHI_OPERATIONS_CONTRACT_VERSION, RHI_READYZ_PATH, + RhiBoundOperationsServer, RhiOperationsCancellationToken, RhiOperationsError, + RhiOperationsErrorKind, RhiOperationsServer, +}; pub use presence_desired::{ RHI_PRESENCE_DESIRED_CONTRACT_VERSION, RHI_PRESENCE_DESIRED_MAX_TARGETS, RhiPresenceDesiredAuthority, RhiPresenceDesiredCommitOutcome, RhiPresenceDesiredError, @@ -269,6 +276,17 @@ pub use state_trade::{ RhiTradeEvidencePersistenceErrorKind, RhiTradeEvidencePersistenceOutcome, RhiTradeSourceObservation, }; +pub use status_v1::{ + RHI_DETAILED_STATUS_MAX_UTF8_BYTES, RHI_STATUS_CACHE_CONTRACT_VERSION, + RHI_STATUS_REASON_CODE_COUNT, RhiEvidenceTransportStatusV1, RhiIdentityHealthV1, + RhiIntegrityStateV1, RhiPersistenceHealthV1, RhiPersistenceStatusV1, RhiPresenceStatusV1, + RhiProviderHealthV1, RhiProviderStatusV1, RhiPublicationStatusV1, RhiReconciliationStatusV1, + RhiServicePhase, RhiStatusBuildInfoV1, RhiStatusBuildMode, RhiStatusCommonV1, + RhiStatusConfigurationIdentityV1, RhiStatusConfigurationSource, RhiStatusError, + RhiStatusErrorKind, RhiStatusObservationV1, RhiStatusPublisher, RhiStatusReader, + RhiStatusReasonCode, RhiStatusReasonCodes, RhiStatusSnapshot, RhiStatusUnixSeconds, + RhiTransportHealthV1, rhi_status_cache, +}; pub use trade_ingest::{ RHI_TRADE_EVENT_EXTRA_FIELD_MAX_COUNT, RHI_TRADE_EVENT_EXTRA_JSON_MAX_BYTES, RHI_TRADE_EVENT_ID_MAX_BYTES, RHI_TRADE_EVENT_PUBLIC_KEY_MAX_BYTES, diff --git a/src/operations_v1.rs b/src/operations_v1.rs @@ -0,0 +1,267 @@ +//! Rhi-owned adapter for the exact passive TCP operations surface. + +use core::{fmt, time::Duration}; +use std::{error::Error, net::SocketAddr}; + +use radroots_service_host::{ + BoundOperationsServer as HostBoundOperationsServer, CancellationToken as HostCancellationToken, + OperationsBindPolicy as HostOperationsBindPolicy, + OperationsListenAddress as HostOperationsListenAddress, + OperationsListenerConfig as HostOperationsListenerConfig, + OperationsServer as HostOperationsServer, OperationsServerError as HostOperationsServerError, + OperationsTransportLimitValues as HostOperationsTransportLimitValues, + OperationsTransportLimits as HostOperationsTransportLimits, +}; +use serde_json::Value; + +use crate::{RhiConfigDocumentV1, RhiStatusReader}; + +/// Exact Rhi TCP operations contract version. +pub const RHI_OPERATIONS_CONTRACT_VERSION: u32 = 1; + +/// Exact liveness route exposed by the optional TCP listener. +pub const RHI_LIVEZ_PATH: &str = "/livez"; + +/// Exact readiness route exposed by the optional TCP listener. +pub const RHI_READYZ_PATH: &str = "/readyz"; + +/// Exact bounded metrics route exposed by the optional TCP listener. +pub const RHI_METRICS_PATH: &str = "/metrics"; + +/// Cloneable cooperative cancellation owned by the Rhi runtime supervisor. +#[derive(Clone, Default)] +pub struct RhiOperationsCancellationToken { + inner: HostCancellationToken, +} + +impl RhiOperationsCancellationToken { + #[must_use] + pub fn new() -> Self { + Self::default() + } + + /// Requests cancellation. Repeated requests have no additional effect. + pub fn cancel(&self) { + self.inner.cancel(); + } + + #[must_use] + pub fn is_cancelled(&self) -> bool { + self.inner.is_cancelled() + } +} + +impl fmt::Debug for RhiOperationsCancellationToken { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiOperationsCancellationToken") + .field("cancelled", &self.is_cancelled()) + .finish() + } +} + +/// Unbound exact-route Rhi TCP operations server. +pub struct RhiOperationsServer { + inner: HostOperationsServer, +} + +impl RhiOperationsServer { + /// Projects the already-validated Rhi configuration and passive status cache. + pub fn new( + config: &RhiConfigDocumentV1, + status: &RhiStatusReader, + ) -> Result<Self, RhiOperationsError> { + let listener = listener_config(config)?; + HostOperationsServer::new(listener, status.operations_cache()) + .map(|inner| Self { inner }) + .map_err(map_server_error) + } + + /// Binds the exact configured address without starting admission. + pub async fn bind(self) -> Result<RhiBoundOperationsServer, RhiOperationsError> { + self.inner + .bind() + .await + .map(|inner| RhiBoundOperationsServer { inner }) + .map_err(map_server_error) + } +} + +impl fmt::Debug for RhiOperationsServer { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiOperationsServer([sealed])") + } +} + +/// Successfully bound exact-route Rhi TCP operations server. +pub struct RhiBoundOperationsServer { + inner: HostBoundOperationsServer, +} + +impl RhiBoundOperationsServer { + #[must_use] + pub fn local_address(&self) -> SocketAddr { + self.inner.local_address() + } + + /// Serves until explicit supervisor cancellation, then drains owned work. + pub async fn serve( + self, + cancellation: RhiOperationsCancellationToken, + ) -> Result<(), RhiOperationsError> { + self.inner + .serve(cancellation.inner) + .await + .map_err(map_server_error) + } +} + +impl fmt::Debug for RhiBoundOperationsServer { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiBoundOperationsServer([sealed])") + } +} + +/// Stable source-free Rhi operations failure classification. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RhiOperationsErrorKind { + Disabled, + InvalidConfiguration, + Bind, + LocalAddress, + Accept, + ConnectionTaskPanicked, +} + +impl RhiOperationsErrorKind { + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::Disabled => "operations_disabled", + Self::InvalidConfiguration => "operations_configuration_invalid", + Self::Bind => "operations_bind_failed", + Self::LocalAddress => "operations_local_address_failed", + Self::Accept => "operations_accept_failed", + Self::ConnectionTaskPanicked => "operations_connection_task_panicked", + } + } +} + +/// One redacted source-free Rhi operations failure. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct RhiOperationsError { + kind: RhiOperationsErrorKind, +} + +impl RhiOperationsError { + const fn new(kind: RhiOperationsErrorKind) -> Self { + Self { kind } + } + + #[must_use] + pub const fn kind(self) -> RhiOperationsErrorKind { + self.kind + } + + #[must_use] + pub const fn code(self) -> &'static str { + self.kind.code() + } +} + +impl fmt::Debug for RhiOperationsError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiOperationsError") + .field("kind", &self.kind) + .finish() + } +} + +impl fmt::Display for RhiOperationsError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("Rhi TCP operations failed") + } +} + +impl Error for RhiOperationsError {} + +fn listener_config( + config: &RhiConfigDocumentV1, +) -> Result<HostOperationsListenerConfig, RhiOperationsError> { + let operations = config + .normalized() + .pointer("/operations") + .ok_or_else(invalid_configuration)?; + if !boolean(operations, "/enabled")? { + return Err(RhiOperationsError::new(RhiOperationsErrorKind::Disabled)); + } + let listen = string(operations, "/listen")? + .parse::<SocketAddr>() + .map_err(|_| invalid_configuration())?; + let listen = HostOperationsListenAddress::new(listen).map_err(|_| invalid_configuration())?; + let bind_policy = match string(operations, "/bind_policy")? { + "loopback_only" => HostOperationsBindPolicy::LoopbackOnly, + "explicit_public" => HostOperationsBindPolicy::Public, + _ => return Err(invalid_configuration()), + }; + let values = HostOperationsTransportLimitValues { + header_count: unsigned_u32(operations, "/limits/header_count")?, + header_bytes: unsigned_u32(operations, "/limits/header_bytes")?, + response_body_utf8_bytes: unsigned_u32(operations, "/limits/response_body_utf8_bytes")?, + concurrent_connections: unsigned_u32(operations, "/limits/concurrent_connections")?, + request_deadline: Duration::from_millis(unsigned( + operations, + "/limits/request_deadline_ms", + )?), + idle_timeout: Duration::from_millis(unsigned(operations, "/limits/idle_timeout_ms")?), + }; + let limits = HostOperationsTransportLimits::new(values).map_err(|_| invalid_configuration())?; + HostOperationsListenerConfig::enabled(listen, bind_policy, limits) + .map_err(|_| invalid_configuration()) +} + +fn boolean(value: &Value, pointer: &str) -> Result<bool, RhiOperationsError> { + value + .pointer(pointer) + .and_then(Value::as_bool) + .ok_or_else(invalid_configuration) +} + +fn string<'a>(value: &'a Value, pointer: &str) -> Result<&'a str, RhiOperationsError> { + value + .pointer(pointer) + .and_then(Value::as_str) + .ok_or_else(invalid_configuration) +} + +fn unsigned(value: &Value, pointer: &str) -> Result<u64, RhiOperationsError> { + value + .pointer(pointer) + .and_then(Value::as_u64) + .ok_or_else(invalid_configuration) +} + +fn unsigned_u32(value: &Value, pointer: &str) -> Result<u32, RhiOperationsError> { + u32::try_from(unsigned(value, pointer)?).map_err(|_| invalid_configuration()) +} + +const fn invalid_configuration() -> RhiOperationsError { + RhiOperationsError::new(RhiOperationsErrorKind::InvalidConfiguration) +} + +const fn map_server_error(error: HostOperationsServerError) -> RhiOperationsError { + let kind = match error { + HostOperationsServerError::Disabled => RhiOperationsErrorKind::Disabled, + HostOperationsServerError::HeaderLimitBelowParserFloor => { + RhiOperationsErrorKind::InvalidConfiguration + } + HostOperationsServerError::Bind { .. } => RhiOperationsErrorKind::Bind, + HostOperationsServerError::LocalAddress { .. } => RhiOperationsErrorKind::LocalAddress, + HostOperationsServerError::Accept { .. } => RhiOperationsErrorKind::Accept, + HostOperationsServerError::ConnectionTaskPanicked => { + RhiOperationsErrorKind::ConnectionTaskPanicked + } + }; + RhiOperationsError::new(kind) +} diff --git a/src/status_v1.rs b/src/status_v1.rs @@ -0,0 +1,1180 @@ +//! Passive, latest-value Rhi lifecycle and detailed-status publication. + +use core::fmt; +use std::{error::Error, sync::Arc, time::Duration}; + +use radroots_service_host::{ + BoundedMetricsSnapshot, BuildInfo as HostBuildInfo, + BuildInfoEnvironment as HostBuildInfoEnvironment, BuildMode as HostBuildMode, + CachedServiceState, CachedServiceStatePublisher, CachedServiceStateReader, CommonMetricGroup, + ConfigurationIdentity as HostConfigurationIdentity, + ConfigurationSource as HostConfigurationSource, ContractVersions as HostContractVersions, + InstanceId, IntegrityState as HostIntegrityState, MetricDescriptor, MetricKind, MetricLabel, + MetricLabelKey, MetricName, MetricSample, MetricValue, + PersistenceHealth as HostPersistenceHealth, PersistenceSummary as HostPersistenceSummary, + Readiness as HostReadiness, ReasonCode as HostReasonCode, ReasonCodes as HostReasonCodes, + ServiceId, ServiceOperationalState as HostServiceOperationalState, + ServicePhase as HostServicePhase, ServiceStatus, ServiceStatusDetail, + Sha256Digest as HostSha256Digest, StatusContractError, StatusEncodingError, StatusModelError, + UptimeMillis as HostUptimeMillis, cached_service_state, +}; +use serde::Serialize; + +/// Exact version of the passive Rhi status-cache contract. +pub const RHI_STATUS_CACHE_CONTRACT_VERSION: u32 = 1; + +/// Maximum encoded byte length of one detailed Rhi status response. +pub const RHI_DETAILED_STATUS_MAX_UTF8_BYTES: usize = + radroots_service_host::SERVICE_STATUS_MAX_UTF8_BYTES; + +/// Number of stable reason codes admitted by detailed Rhi status. +pub const RHI_STATUS_REASON_CODE_COUNT: usize = 13; + +/// One closed stable source-free status reason code. +#[derive(Clone, Copy, Debug, Hash, PartialEq, Eq, PartialOrd, Ord, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum RhiStatusReasonCode { + IdentityUnavailable, + DatabaseSchemaMismatch, + DatabaseReadOnly, + DatabaseLowDisk, + SourceUnavailable, + SubscriptionInactive, + RecoveryIncomplete, + PublicationRecoveryIncomplete, + PresenceStateUnavailable, + ReconciliationBacklogExceeded, + AdminListenerFailed, + OperationsListenerFailed, + ShutdownInProgress, +} + +impl RhiStatusReasonCode { + pub fn new(value: impl AsRef<str>) -> Result<Self, RhiStatusError> { + match value.as_ref() { + "identity_unavailable" => Ok(Self::IdentityUnavailable), + "database_schema_mismatch" => Ok(Self::DatabaseSchemaMismatch), + "database_read_only" => Ok(Self::DatabaseReadOnly), + "database_low_disk" => Ok(Self::DatabaseLowDisk), + "source_unavailable" => Ok(Self::SourceUnavailable), + "subscription_inactive" => Ok(Self::SubscriptionInactive), + "recovery_incomplete" => Ok(Self::RecoveryIncomplete), + "publication_recovery_incomplete" => Ok(Self::PublicationRecoveryIncomplete), + "presence_state_unavailable" => Ok(Self::PresenceStateUnavailable), + "reconciliation_backlog_exceeded" => Ok(Self::ReconciliationBacklogExceeded), + "admin_listener_failed" => Ok(Self::AdminListenerFailed), + "operations_listener_failed" => Ok(Self::OperationsListenerFailed), + "shutdown_in_progress" => Ok(Self::ShutdownInProgress), + _ => Err(RhiStatusError::new(RhiStatusErrorKind::InvalidReasonCode)), + } + } + + #[must_use] + pub const fn as_str(self) -> &'static str { + match self { + Self::IdentityUnavailable => "identity_unavailable", + Self::DatabaseSchemaMismatch => "database_schema_mismatch", + Self::DatabaseReadOnly => "database_read_only", + Self::DatabaseLowDisk => "database_low_disk", + Self::SourceUnavailable => "source_unavailable", + Self::SubscriptionInactive => "subscription_inactive", + Self::RecoveryIncomplete => "recovery_incomplete", + Self::PublicationRecoveryIncomplete => "publication_recovery_incomplete", + Self::PresenceStateUnavailable => "presence_state_unavailable", + Self::ReconciliationBacklogExceeded => "reconciliation_backlog_exceeded", + Self::AdminListenerFailed => "admin_listener_failed", + Self::OperationsListenerFailed => "operations_listener_failed", + Self::ShutdownInProgress => "shutdown_in_progress", + } + } +} + +/// Canonically ordered, unique, bounded status reasons. +#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize)] +#[serde(transparent)] +pub struct RhiStatusReasonCodes(Vec<RhiStatusReasonCode>); + +impl RhiStatusReasonCodes { + #[must_use] + pub const fn empty() -> Self { + Self(Vec::new()) + } + + pub fn new( + values: impl IntoIterator<Item = RhiStatusReasonCode>, + ) -> Result<Self, RhiStatusError> { + let mut bounded = Vec::with_capacity(RHI_STATUS_REASON_CODE_COUNT); + for value in values.into_iter().take(RHI_STATUS_REASON_CODE_COUNT + 1) { + if bounded.len() == RHI_STATUS_REASON_CODE_COUNT { + return Err(RhiStatusError::new(RhiStatusErrorKind::TooManyReasonCodes)); + } + bounded.push(value); + } + bounded.sort_unstable(); + bounded.dedup(); + Ok(Self(bounded)) + } + + #[must_use] + pub fn as_slice(&self) -> &[RhiStatusReasonCode] { + &self.0 + } + + fn into_host(self) -> Result<HostReasonCodes, RhiStatusError> { + let values = self + .0 + .into_iter() + .map(|value| HostReasonCode::new(value.as_str()).map_err(map_contract_error)) + .collect::<Result<Vec<_>, _>>()?; + HostReasonCodes::new(values).map_err(map_contract_error) + } +} + +/// Closed common lifecycle phase used by the Rhi public boundary. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum RhiServicePhase { + Starting, + Ready, + Degraded, + Unready, + Stopping, + Failed, +} + +impl RhiServicePhase { + const fn into_host(self) -> HostServicePhase { + match self { + Self::Starting => HostServicePhase::Starting, + Self::Ready => HostServicePhase::Ready, + Self::Degraded => HostServicePhase::Degraded, + Self::Unready => HostServicePhase::Unready, + Self::Stopping => HostServicePhase::Stopping, + Self::Failed => HostServicePhase::Failed, + } + } + + const fn from_host(value: HostServicePhase) -> Self { + match value { + HostServicePhase::Starting => Self::Starting, + HostServicePhase::Ready => Self::Ready, + HostServicePhase::Degraded => Self::Degraded, + HostServicePhase::Unready => Self::Unready, + HostServicePhase::Stopping => Self::Stopping, + HostServicePhase::Failed => Self::Failed, + } + } +} + +/// Build-metadata admission mode for Rhi status identity. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RhiStatusBuildMode { + Development, + Release, +} + +/// Complete deterministic build identity retained behind the Rhi boundary. +pub struct RhiStatusBuildInfoV1 { + inner: HostBuildInfo, +} + +impl RhiStatusBuildInfoV1 { + /// Validates the complete build/source-lock identity with fixed Rhi contracts. + pub fn new( + mode: RhiStatusBuildMode, + service_version: Option<&str>, + service_commit: Option<&str>, + lib_revision: Option<&str>, + rust_version: Option<&str>, + target: Option<&str>, + feature_profile: Option<&str>, + ) -> Result<Self, RhiStatusError> { + let contract_versions = HostContractVersions::new(1, 10, 1, 1, 1) + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::InvalidBuildInfo))?; + HostBuildInfo::from_compile_time( + match mode { + RhiStatusBuildMode::Development => HostBuildMode::Development, + RhiStatusBuildMode::Release => HostBuildMode::Release, + }, + HostBuildInfoEnvironment { + service_version, + service_commit, + lib_revision, + rust_version, + target, + feature_profile, + contract_versions, + }, + ) + .map(|inner| Self { inner }) + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::InvalidBuildInfo)) + } +} + +impl fmt::Debug for RhiStatusBuildInfoV1 { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiStatusBuildInfoV1([redacted])") + } +} + +/// Exact configuration-source vocabulary exposed by detailed status. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RhiStatusConfigurationSource { + ExplicitConfig, + DerivedRepoLocal, +} + +/// Safe configuration identity retained behind the Rhi boundary. +pub struct RhiStatusConfigurationIdentityV1 { + inner: HostConfigurationIdentity, +} + +impl RhiStatusConfigurationIdentityV1 { + pub fn new( + digest: impl AsRef<str>, + source: RhiStatusConfigurationSource, + ) -> Result<Self, RhiStatusError> { + let service = ServiceId::new("rhi") + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::InvalidConfiguration))?; + let digest = HostSha256Digest::new(digest) + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::InvalidConfiguration))?; + HostConfigurationIdentity::for_service( + &service, + digest, + match source { + RhiStatusConfigurationSource::ExplicitConfig => { + HostConfigurationSource::ExplicitConfig + } + RhiStatusConfigurationSource::DerivedRepoLocal => { + HostConfigurationSource::DerivedRepoLocal + } + }, + ) + .map(|inner| Self { inner }) + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::InvalidConfiguration)) + } +} + +impl fmt::Debug for RhiStatusConfigurationIdentityV1 { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiStatusConfigurationIdentityV1([redacted])") + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RhiPersistenceHealthV1 { + Ready, + ReadOnly, + RepairRequired, + Unavailable, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RhiIntegrityStateV1 { + Verified, + VerificationRequired, + Failed, +} + +/// Validated persistence summary retained behind the Rhi boundary. +pub struct RhiPersistenceStatusV1 { + inner: HostPersistenceSummary, + ready: bool, +} + +impl RhiPersistenceStatusV1 { + pub fn new( + health: RhiPersistenceHealthV1, + schema_version: u32, + generation: u64, + integrity: RhiIntegrityStateV1, + reason_codes: RhiStatusReasonCodes, + ) -> Result<Self, RhiStatusError> { + let reason_codes = reason_codes.into_host()?; + let ready = + health == RhiPersistenceHealthV1::Ready && integrity == RhiIntegrityStateV1::Verified; + HostPersistenceSummary::new( + match health { + RhiPersistenceHealthV1::Ready => HostPersistenceHealth::Ready, + RhiPersistenceHealthV1::ReadOnly => HostPersistenceHealth::ReadOnly, + RhiPersistenceHealthV1::RepairRequired => HostPersistenceHealth::RepairRequired, + RhiPersistenceHealthV1::Unavailable => HostPersistenceHealth::Unavailable, + }, + schema_version, + generation, + match integrity { + RhiIntegrityStateV1::Verified => HostIntegrityState::Verified, + RhiIntegrityStateV1::VerificationRequired => { + HostIntegrityState::VerificationRequired + } + RhiIntegrityStateV1::Failed => HostIntegrityState::Failed, + }, + reason_codes, + ) + .map(|inner| Self { inner, ready }) + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::InvalidPersistence)) + } +} + +impl fmt::Debug for RhiPersistenceStatusV1 { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiPersistenceStatusV1([redacted])") + } +} + +/// One validated Unix timestamp used only for an oldest pending work item. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize)] +#[serde(transparent)] +pub struct RhiStatusUnixSeconds(u64); + +impl RhiStatusUnixSeconds { + /// Constructs a timestamp representable by SQLite and the frozen wire contract. + pub fn new(value: u64) -> Result<Self, RhiStatusError> { + if value > i64::MAX as u64 { + return Err(RhiStatusError::new(RhiStatusErrorKind::InvalidTime)); + } + Ok(Self(value)) + } + + /// Returns exact whole Unix seconds. + #[must_use] + pub const fn get(self) -> u64 { + self.0 + } +} + +/// Passive availability of one configured Rhi identity role. +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] +pub struct RhiIdentityHealthV1 { + configured: bool, + available: bool, + reason_codes: RhiStatusReasonCodes, +} + +impl RhiIdentityHealthV1 { + /// Constructs one role observation, rejecting availability without configuration. + pub fn new( + configured: bool, + available: bool, + reason_codes: RhiStatusReasonCodes, + ) -> Result<Self, RhiStatusError> { + if available && !configured { + return Err(RhiStatusError::new( + RhiStatusErrorKind::InvalidIdentityHealth, + )); + } + Ok(Self { + configured, + available, + reason_codes, + }) + } + + #[must_use] + pub const fn is_configured(&self) -> bool { + self.configured + } + + #[must_use] + pub const fn is_available(&self) -> bool { + self.available + } + + #[must_use] + pub const fn reason_codes(&self) -> &RhiStatusReasonCodes { + &self.reason_codes + } +} + +/// Passive projection for the single configured RHI service identity. +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] +pub struct RhiProviderStatusV1 { + health: RhiProviderHealthV1, + identity: RhiIdentityHealthV1, + reason_codes: RhiStatusReasonCodes, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum RhiProviderHealthV1 { + Ready, + Unavailable, +} + +impl RhiProviderStatusV1 { + /// Derives provider health from the sole configured service identity. + pub fn new( + identity: RhiIdentityHealthV1, + reason_codes: RhiStatusReasonCodes, + ) -> Result<Self, RhiStatusError> { + if !identity.configured { + return Err(RhiStatusError::new( + RhiStatusErrorKind::InvalidProviderState, + )); + } + let health = if identity.available { + RhiProviderHealthV1::Ready + } else { + RhiProviderHealthV1::Unavailable + }; + Ok(Self { + health, + identity, + reason_codes, + }) + } + + #[must_use] + pub const fn health(&self) -> RhiProviderHealthV1 { + self.health + } + + #[must_use] + pub const fn identity(&self) -> &RhiIdentityHealthV1 { + &self.identity + } + + #[must_use] + pub const fn reason_codes(&self) -> &RhiStatusReasonCodes { + &self.reason_codes + } +} + +/// Passive evidence-source transport projection for detailed status. +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] +pub struct RhiEvidenceTransportStatusV1 { + health: RhiTransportHealthV1, + required_sources_ready: bool, + subscriber_active: bool, + configured_source_count: u64, + reachable_source_count: u64, + reason_codes: RhiStatusReasonCodes, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum RhiTransportHealthV1 { + Ready, + Degraded, + Unavailable, +} + +impl RhiEvidenceTransportStatusV1 { + /// Constructs one transport observation, rejecting contradictory ready state. + pub fn new( + health: RhiTransportHealthV1, + required_sources_ready: bool, + subscriber_active: bool, + configured_source_count: u64, + reachable_source_count: u64, + reason_codes: RhiStatusReasonCodes, + ) -> Result<Self, RhiStatusError> { + if reachable_source_count > configured_source_count + || (health == RhiTransportHealthV1::Ready + && (!required_sources_ready + || !subscriber_active + || configured_source_count == 0 + || reachable_source_count == 0)) + { + return Err(RhiStatusError::new( + RhiStatusErrorKind::InvalidTransportState, + )); + } + Ok(Self { + health, + required_sources_ready, + subscriber_active, + configured_source_count, + reachable_source_count, + reason_codes, + }) + } + + #[must_use] + pub const fn health(&self) -> RhiTransportHealthV1 { + self.health + } + + #[must_use] + pub const fn required_sources_ready(&self) -> bool { + self.required_sources_ready + } + + #[must_use] + pub const fn subscriber_active(&self) -> bool { + self.subscriber_active + } + + #[must_use] + pub const fn configured_source_count(&self) -> u64 { + self.configured_source_count + } + + #[must_use] + pub const fn reachable_source_count(&self) -> u64 { + self.reachable_source_count + } + + #[must_use] + pub const fn reason_codes(&self) -> &RhiStatusReasonCodes { + &self.reason_codes + } +} + +/// Passive bounded reconciliation-work summary. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)] +pub struct RhiReconciliationStatusV1 { + pending: u64, + leased: u64, + exhausted: u64, + #[serde(skip_serializing_if = "Option::is_none")] + oldest_pending_at_utc: Option<RhiStatusUnixSeconds>, +} + +impl RhiReconciliationStatusV1 { + #[must_use] + pub const fn new( + pending: u64, + leased: u64, + exhausted: u64, + oldest_pending_at_utc: Option<RhiStatusUnixSeconds>, + ) -> Self { + Self { + pending, + leased, + exhausted, + oldest_pending_at_utc, + } + } + + #[must_use] + pub const fn pending(self) -> u64 { + self.pending + } + + #[must_use] + pub const fn leased(self) -> u64 { + self.leased + } + + #[must_use] + pub const fn exhausted(self) -> u64 { + self.exhausted + } + + #[must_use] + pub const fn oldest_pending_at_utc(self) -> Option<RhiStatusUnixSeconds> { + self.oldest_pending_at_utc + } +} + +/// Passive publication-outbox summary derived before cache publication. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)] +pub struct RhiPublicationStatusV1 { + pending: u64, + unknown: u64, + #[serde(skip_serializing_if = "Option::is_none")] + oldest_pending_at_utc: Option<RhiStatusUnixSeconds>, +} + +impl RhiPublicationStatusV1 { + #[must_use] + pub const fn new( + pending: u64, + unknown: u64, + oldest_pending_at_utc: Option<RhiStatusUnixSeconds>, + ) -> Self { + Self { + pending, + unknown, + oldest_pending_at_utc, + } + } + + #[must_use] + pub const fn pending(self) -> u64 { + self.pending + } + + #[must_use] + pub const fn unknown(self) -> u64 { + self.unknown + } + + #[must_use] + pub const fn oldest_pending_at_utc(self) -> Option<RhiStatusUnixSeconds> { + self.oldest_pending_at_utc + } +} + +/// Passive desired-presence publication summary. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)] +pub struct RhiPresenceStatusV1 { + pending: u64, + unknown: u64, +} + +impl RhiPresenceStatusV1 { + #[must_use] + pub const fn new(pending: u64, unknown: u64) -> Self { + Self { pending, unknown } + } + + #[must_use] + pub const fn pending(self) -> u64 { + self.pending + } + + #[must_use] + pub const fn unknown(self) -> u64 { + self.unknown + } +} + +#[derive(Serialize)] +struct RhiStatusDetailV1 { + identity: RhiIdentityHealthV1, + reconciliation: RhiReconciliationStatusV1, + publication: RhiPublicationStatusV1, + presence: RhiPresenceStatusV1, +} + +impl ServiceStatusDetail for RhiStatusDetailV1 { + type Provider = RhiProviderStatusV1; + type Transport = RhiEvidenceTransportStatusV1; + + const FIELD_NAME: &'static str = "rhi"; +} + +/// Validated common fields shared by one detailed status publication. +/// +/// Construction and publication perform validation and bounded encoding only. +/// They do not query SQLite, providers, sources, relays, DNS, credentials, or the clock. +pub struct RhiStatusCommonV1 { + operational: HostServiceOperationalState, + uptime: HostUptimeMillis, + build: HostBuildInfo, + configuration: HostConfigurationIdentity, + persistence: HostPersistenceSummary, + persistence_ready: bool, +} + +impl RhiStatusCommonV1 { + /// Validates the common lifecycle and detailed-status envelope fields. + pub fn new( + phase: RhiServicePhase, + ready: bool, + reason_codes: RhiStatusReasonCodes, + uptime_millis: u64, + build: RhiStatusBuildInfoV1, + configuration: RhiStatusConfigurationIdentityV1, + persistence: RhiPersistenceStatusV1, + ) -> Result<Self, RhiStatusError> { + let operational = HostServiceOperationalState::new( + phase.into_host(), + if ready { + HostReadiness::READY + } else { + HostReadiness::NOT_READY + }, + reason_codes.into_host()?, + ) + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::InvalidLifecycle))?; + let uptime = HostUptimeMillis::from_duration(Duration::from_millis(uptime_millis)) + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::InvalidTime))?; + Ok(Self { + operational, + uptime, + build: build.inner, + configuration: configuration.inner, + persistence: persistence.inner, + persistence_ready: persistence.ready, + }) + } +} + +impl fmt::Debug for RhiStatusCommonV1 { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiStatusCommonV1([redacted])") + } +} + +/// One complete, already-observed status publication input. +pub struct RhiStatusObservationV1 { + common: RhiStatusCommonV1, + provider: RhiProviderStatusV1, + transport: RhiEvidenceTransportStatusV1, + reconciliation: RhiReconciliationStatusV1, + publication: RhiPublicationStatusV1, + presence: RhiPresenceStatusV1, +} + +impl RhiStatusObservationV1 { + #[must_use] + pub fn new( + common: RhiStatusCommonV1, + provider: RhiProviderStatusV1, + transport: RhiEvidenceTransportStatusV1, + reconciliation: RhiReconciliationStatusV1, + publication: RhiPublicationStatusV1, + presence: RhiPresenceStatusV1, + ) -> Self { + Self { + common, + provider, + transport, + reconciliation, + publication, + presence, + } + } + + fn into_cached(self, instance: &InstanceId) -> Result<PreparedRhiStatus, RhiStatusError> { + let operational = self.common.operational.clone(); + if operational.readiness().is_ready() + && (!self.common.persistence_ready + || self.provider.health != RhiProviderHealthV1::Ready + || !self.transport.required_sources_ready + || !self.transport.subscriber_active) + { + return Err(RhiStatusError::new(RhiStatusErrorKind::InvalidLifecycle)); + } + if operational.phase() == HostServicePhase::Ready + && self.transport.health != RhiTransportHealthV1::Ready + { + return Err(RhiStatusError::new(RhiStatusErrorKind::InvalidLifecycle)); + } + let operations_metrics = bounded_operations_metrics(&operational)?; + let detail = RhiStatusDetailV1 { + identity: self.provider.identity.clone(), + reconciliation: self.reconciliation, + publication: self.publication, + presence: self.presence, + }; + let service = ServiceId::new("rhi") + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::InvalidModel))?; + let status = ServiceStatus::new( + service, + instance.clone(), + self.common.operational, + self.common.uptime, + self.common.build, + self.common.configuration, + self.common.persistence, + self.provider, + self.transport, + detail, + ) + .map_err(map_model_error)?; + let json = status.to_bounded_json().map_err(map_encoding_error)?; + Ok(PreparedRhiStatus { + detail: CachedServiceState::new( + operational.clone(), + RhiCachedStatus { + json: json.into_boxed_slice(), + }, + ), + operations: CachedServiceState::new(operational, operations_metrics), + }) + } +} + +impl fmt::Debug for RhiStatusObservationV1 { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiStatusObservationV1([redacted])") + } +} + +struct RhiCachedStatus { + json: Box<[u8]>, +} + +struct PreparedRhiStatus { + detail: CachedServiceState<RhiCachedStatus>, + operations: CachedServiceState<BoundedMetricsSnapshot>, +} + +fn bounded_operations_metrics( + operational: &HostServiceOperationalState, +) -> Result<BoundedMetricsSnapshot, RhiStatusError> { + let phase_name = MetricName::new("radroots_rhi_service_phase").map_err(map_metrics_error)?; + let ready_name = MetricName::new("radroots_rhi_service_ready").map_err(map_metrics_error)?; + let descriptors = [ + MetricDescriptor::new( + CommonMetricGroup::Phase, + phase_name.clone(), + "Current cached Rhi service phase.", + MetricKind::Gauge, + [MetricLabelKey::Phase], + ) + .map_err(map_metrics_error)?, + MetricDescriptor::new( + CommonMetricGroup::Phase, + ready_name.clone(), + "Current cached Rhi readiness bit.", + MetricKind::Gauge, + [], + ) + .map_err(map_metrics_error)?, + ]; + let samples = [ + MetricSample::new( + phase_name, + MetricValue::Gauge(1), + [MetricLabel::phase(operational.phase())], + ) + .map_err(map_metrics_error)?, + MetricSample::new( + ready_name, + MetricValue::Gauge(i64::from(operational.readiness().is_ready())), + [], + ) + .map_err(map_metrics_error)?, + ]; + BoundedMetricsSnapshot::new(descriptors, samples).map_err(map_metrics_error) +} + +fn map_metrics_error(_: radroots_service_host::MetricsContractError) -> RhiStatusError { + RhiStatusError::new(RhiStatusErrorKind::InvalidModel) +} + +impl fmt::Debug for RhiCachedStatus { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiCachedStatus") + .field("json_utf8_bytes", &self.json.len()) + .finish() + } +} + +/// Sole publication authority for one process-local Rhi status cache. +/// +/// This type deliberately does not implement `Clone`. A successful publish +/// atomically replaces the one retained snapshot. A failed encoding or illegal +/// lifecycle transition leaves the previous snapshot unchanged. +pub struct RhiStatusPublisher { + instance: InstanceId, + inner: CachedServiceStatePublisher<RhiCachedStatus>, + operations: CachedServiceStatePublisher<BoundedMetricsSnapshot>, +} + +impl RhiStatusPublisher { + /// Encodes one observation, publishes its passive operations projection, + /// and then atomically replaces the detailed-status snapshot. + pub fn publish(&mut self, next: RhiStatusObservationV1) -> Result<(), RhiStatusError> { + let next = next.into_cached(&self.instance)?; + self.operations + .publish(next.operations) + .map_err(map_contract_error)?; + self.inner.publish(next.detail).map_err(map_contract_error) + } + + /// Creates another passive reader without sharing publication authority. + #[must_use] + pub fn subscribe(&self) -> RhiStatusReader { + RhiStatusReader { + inner: self.inner.subscribe(), + operations: self.operations.subscribe(), + } + } +} + +impl fmt::Debug for RhiStatusPublisher { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiStatusPublisher([sealed])") + } +} + +/// Cloneable passive reader of the latest Rhi lifecycle and detailed status. +pub struct RhiStatusReader { + inner: CachedServiceStateReader<RhiCachedStatus>, + operations: CachedServiceStateReader<BoundedMetricsSnapshot>, +} + +impl Clone for RhiStatusReader { + fn clone(&self) -> Self { + Self { + inner: self.inner.clone(), + operations: self.operations.clone(), + } + } +} + +impl RhiStatusReader { + /// Returns the latest immutable snapshot without awaiting or probing. + #[must_use] + pub fn snapshot(&self) -> RhiStatusSnapshot { + RhiStatusSnapshot { + inner: self.inner.snapshot(), + } + } + + /// Waits for a later publication and returns the newest retained value. + pub async fn changed(&mut self) -> Result<RhiStatusSnapshot, RhiStatusError> { + self.inner + .changed() + .await + .map(|inner| RhiStatusSnapshot { inner }) + .map_err(|_| RhiStatusError::new(RhiStatusErrorKind::PublisherDropped)) + } + + pub(crate) fn operations_cache(&self) -> CachedServiceStateReader<BoundedMetricsSnapshot> { + self.operations.clone() + } +} + +impl fmt::Debug for RhiStatusReader { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiStatusReader([passive])") + } +} + +/// One immutable point-in-time status snapshot backed by the retained cache `Arc`. +pub struct RhiStatusSnapshot { + inner: Arc<CachedServiceState<RhiCachedStatus>>, +} + +impl RhiStatusSnapshot { + #[must_use] + pub fn phase(&self) -> RhiServicePhase { + RhiServicePhase::from_host(self.inner.operational().phase()) + } + + #[must_use] + pub fn is_ready(&self) -> bool { + self.inner.operational().readiness().is_ready() + } + + /// Returns the already-bounded canonical detailed-status JSON bytes. + #[must_use] + pub fn detailed_status_json(&self) -> &[u8] { + &self.inner.metrics().json + } +} + +impl fmt::Debug for RhiStatusSnapshot { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiStatusSnapshot") + .field("phase", &self.phase()) + .field("ready", &self.is_ready()) + .field("json_utf8_bytes", &self.detailed_status_json().len()) + .finish() + } +} + +/// Creates the single-writer, one-latest-value Rhi status cache. +pub fn rhi_status_cache( + instance: InstanceId, + initial: RhiStatusObservationV1, +) -> Result<(RhiStatusPublisher, RhiStatusReader), RhiStatusError> { + let initial = initial.into_cached(&instance)?; + let (inner, reader) = cached_service_state(initial.detail); + let (operations, operations_reader) = cached_service_state(initial.operations); + Ok(( + RhiStatusPublisher { + instance, + inner, + operations, + }, + RhiStatusReader { + inner: reader, + operations: operations_reader, + }, + )) +} + +/// Stable source-free status failure category. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RhiStatusErrorKind { + InvalidReasonCode, + TooManyReasonCodes, + InvalidLifecycle, + InvalidBuildInfo, + InvalidConfiguration, + InvalidPersistence, + InvalidIdentityHealth, + InvalidProviderState, + InvalidTransportState, + InvalidTime, + InvalidModel, + Encoding, + ResponseTooLarge, + InvalidTransition, + PublisherDropped, +} + +impl RhiStatusErrorKind { + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::InvalidReasonCode => "status_reason_code_invalid", + Self::TooManyReasonCodes => "status_reason_count_exceeded", + Self::InvalidLifecycle => "status_lifecycle_invalid", + Self::InvalidBuildInfo => "status_build_info_invalid", + Self::InvalidConfiguration => "status_configuration_invalid", + Self::InvalidPersistence => "status_persistence_invalid", + Self::InvalidIdentityHealth => "status_identity_health_invalid", + Self::InvalidProviderState => "status_provider_state_invalid", + Self::InvalidTransportState => "status_transport_state_invalid", + Self::InvalidTime => "status_time_invalid", + Self::InvalidModel => "status_model_invalid", + Self::Encoding => "status_encoding_failed", + Self::ResponseTooLarge => "status_response_too_large", + Self::InvalidTransition => "status_transition_invalid", + Self::PublisherDropped => "status_publisher_dropped", + } + } + + const fn message(self) -> &'static str { + match self { + Self::InvalidReasonCode => "Rhi status reason code is invalid", + Self::TooManyReasonCodes => "Rhi status has too many reason codes", + Self::InvalidLifecycle => "Rhi lifecycle status is invalid", + Self::InvalidBuildInfo => "Rhi status build identity is invalid", + Self::InvalidConfiguration => "Rhi status configuration identity is invalid", + Self::InvalidPersistence => "Rhi persistence status is invalid", + Self::InvalidIdentityHealth => "Rhi identity health is invalid", + Self::InvalidProviderState => "Rhi provider status is invalid", + Self::InvalidTransportState => "Rhi transport status is invalid", + Self::InvalidTime => "Rhi status time is invalid", + Self::InvalidModel => "Rhi detailed status is invalid", + Self::Encoding => "Rhi detailed status encoding failed", + Self::ResponseTooLarge => "Rhi detailed status exceeds its byte limit", + Self::InvalidTransition => "Rhi lifecycle transition is invalid", + Self::PublisherDropped => "Rhi status publisher is unavailable", + } + } +} + +/// One redacted source-free status failure. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct RhiStatusError { + kind: RhiStatusErrorKind, +} + +impl RhiStatusError { + const fn new(kind: RhiStatusErrorKind) -> Self { + Self { kind } + } + + #[must_use] + pub const fn kind(self) -> RhiStatusErrorKind { + self.kind + } + + #[must_use] + pub const fn code(self) -> &'static str { + self.kind.code() + } +} + +impl fmt::Debug for RhiStatusError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiStatusError") + .field("kind", &self.kind) + .finish() + } +} + +impl fmt::Display for RhiStatusError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(self.kind.message()) + } +} + +impl Error for RhiStatusError {} + +const fn map_model_error(_error: StatusModelError) -> RhiStatusError { + RhiStatusError::new(RhiStatusErrorKind::InvalidModel) +} + +const fn map_encoding_error(error: StatusEncodingError) -> RhiStatusError { + match error { + StatusEncodingError::EncodingFailed => RhiStatusError::new(RhiStatusErrorKind::Encoding), + StatusEncodingError::ResponseTooLarge => { + RhiStatusError::new(RhiStatusErrorKind::ResponseTooLarge) + } + } +} + +const fn map_contract_error(_error: StatusContractError) -> RhiStatusError { + RhiStatusError::new(RhiStatusErrorKind::InvalidTransition) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn invalid_status_inputs_fail_with_safe_source_free_errors() { + assert_eq!( + RhiIdentityHealthV1::new(false, true, RhiStatusReasonCodes::empty()), + Err(RhiStatusError::new( + RhiStatusErrorKind::InvalidIdentityHealth + )) + ); + assert_eq!( + RhiProviderStatusV1::new( + RhiIdentityHealthV1::new(false, false, RhiStatusReasonCodes::empty()).unwrap(), + RhiStatusReasonCodes::empty(), + ), + Err(RhiStatusError::new( + RhiStatusErrorKind::InvalidProviderState + )) + ); + assert_eq!( + RhiEvidenceTransportStatusV1::new( + RhiTransportHealthV1::Ready, + false, + true, + 1, + 0, + RhiStatusReasonCodes::empty(), + ), + Err(RhiStatusError::new( + RhiStatusErrorKind::InvalidTransportState + )) + ); + assert_eq!( + RhiEvidenceTransportStatusV1::new( + RhiTransportHealthV1::Ready, + true, + true, + 1, + 2, + RhiStatusReasonCodes::empty(), + ), + Err(RhiStatusError::new( + RhiStatusErrorKind::InvalidTransportState + )) + ); + assert_eq!( + RhiStatusUnixSeconds::new(i64::MAX as u64 + 1), + Err(RhiStatusError::new(RhiStatusErrorKind::InvalidTime)) + ); + for kind in [ + RhiStatusErrorKind::InvalidReasonCode, + RhiStatusErrorKind::TooManyReasonCodes, + RhiStatusErrorKind::InvalidLifecycle, + RhiStatusErrorKind::InvalidBuildInfo, + RhiStatusErrorKind::InvalidConfiguration, + RhiStatusErrorKind::InvalidPersistence, + RhiStatusErrorKind::InvalidIdentityHealth, + RhiStatusErrorKind::InvalidProviderState, + RhiStatusErrorKind::InvalidTransportState, + RhiStatusErrorKind::InvalidTime, + RhiStatusErrorKind::InvalidModel, + RhiStatusErrorKind::Encoding, + RhiStatusErrorKind::ResponseTooLarge, + RhiStatusErrorKind::InvalidTransition, + RhiStatusErrorKind::PublisherDropped, + ] { + let error = RhiStatusError::new(kind); + assert!(!error.code().is_empty()); + assert!(Error::source(&error).is_none()); + assert!(!format!("{error} {error:?}").contains("source")); + } + } +} diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs @@ -7,9 +7,14 @@ const ROOT: &str = include_str!("../src/lib.rs"); const MAIN: &str = include_str!("../src/main.rs"); const ADMIN: &str = include_str!("../src/admin_v1.rs"); const DOCTOR: &str = include_str!("../src/doctor_v1.rs"); +const OPERATIONS: &str = include_str!("../src/operations_v1.rs"); +const STATUS: &str = include_str!("../src/status_v1.rs"); const PROCESS_RESULT: &str = include_str!("../src/process_result_v1.rs"); const OPERATOR_CONTRACT: &str = include_str!("../contracts/services_hardening/operator_contract.v1.json"); +const STATUS_CONTRACT: &str = include_str!("../contracts/services_hardening/status_cache.v1.json"); +const OPERATIONS_CONTRACT: &str = + include_str!("../contracts/services_hardening/tcp_operations.v1.json"); const ADMIN_IDENTITY_OFFLINE_CONTRACT: &str = include_str!("../contracts/services_hardening/admin_identity_offline.v1.json"); const ADMIN_WAVE_QUALIFICATION_CONTRACT: &str = @@ -87,6 +92,7 @@ const SOURCES: &[&str] = &[ include_str!("../src/features/trade_agreement_attestation.rs"), include_str!("../src/identity_credential.rs"), include_str!("../src/identity_envelope.rs"), + include_str!("../src/operations_v1.rs"), include_str!("../src/presence_desired.rs"), include_str!("../src/presence_publication.rs"), include_str!("../src/process_result_v1.rs"), @@ -114,6 +120,7 @@ const SOURCES: &[&str] = &[ include_str!("../src/state_metadata.rs"), include_str!("../src/state_repository.rs"), include_str!("../src/state_trade.rs"), + include_str!("../src/status_v1.rs"), include_str!("../src/trade_ingest.rs"), ]; @@ -171,6 +178,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { "features", "identity_credential", "identity_envelope", + "operations_v1", "presence_desired", "presence_publication", "process_result_v1", @@ -198,6 +206,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { "state_metadata", "state_repository", "state_trade", + "status_v1", "trade_ingest", ] { assert!( @@ -264,6 +273,20 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { "run_rhi_doctor", "rhi_doctor_check_definitions", "RhiProcessResult", + "RhiStatusPublisher", + "RhiStatusReader", + "RhiStatusSnapshot", + "RhiStatusObservationV1", + "RhiProviderStatusV1", + "RhiEvidenceTransportStatusV1", + "RhiReconciliationStatusV1", + "RhiPublicationStatusV1", + "RhiPresenceStatusV1", + "rhi_status_cache", + "RhiOperationsServer", + "RhiBoundOperationsServer", + "RhiOperationsCancellationToken", + "RhiOperationsErrorKind", "RhiPublicationErrorKind", "RhiPublicationMode", "RhiPublicationRetryPolicy", @@ -393,7 +416,93 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { assert!(!PUBLIC_API.contains("rhi::admin_v1::")); assert!(!PUBLIC_API.contains("rhi::cli_v1::")); assert!(!PUBLIC_API.contains("rhi::doctor_v1::")); + assert!(!PUBLIC_API.contains("rhi::operations_v1::")); assert!(!PUBLIC_API.contains("rhi::process_result_v1::")); + assert!(!PUBLIC_API.contains("rhi::status_v1::")); +} + +#[test] +fn step212_status_and_tcp_operations_are_passive_closed_and_dependency_neutral() { + let status_contract: serde_json::Value = + serde_json::from_str(STATUS_CONTRACT).expect("status contract"); + let operations_contract: serde_json::Value = + serde_json::from_str(OPERATIONS_CONTRACT).expect("operations contract"); + assert_eq!(status_contract["schema"], "radroots.rhi.status-cache.v1"); + assert_eq!(status_contract["step"], 212); + assert_eq!(status_contract["read"]["fresh_probe"], false); + assert_eq!( + operations_contract["schema"], + "radroots.rhi.tcp-operations.v1" + ); + assert_eq!(operations_contract["step"], 212); + assert_eq!(operations_contract["route_registration_extension"], false); + + for required in [ + "CachedServiceStatePublisher<RhiCachedStatus>", + "pub struct RhiStatusPublisher", + "pub struct RhiStatusReader", + "pub struct RhiStatusSnapshot", + "pub fn rhi_status_cache(", + "status.to_bounded_json()", + "reconciliation: RhiReconciliationStatusV1", + "publication: RhiPublicationStatusV1", + "presence: RhiPresenceStatusV1", + "RHI_STATUS_REASON_CODE_COUNT: usize = 13", + "radroots_rhi_service_phase", + "radroots_rhi_service_ready", + ] { + assert!(STATUS.contains(required), "status is missing {required}"); + } + for required in [ + "HostOperationsServer::new(listener, status.operations_cache())", + "RhiOperationsCancellationToken", + "HostOperationsTransportLimits::new(values)", + "HeaderLimitBelowParserFloor", + ] { + assert!( + OPERATIONS.contains(required), + "operations is missing {required}" + ); + } + for forbidden in [ + "sqlx::", + "std::fs::", + "tokio::spawn", + "spawn_blocking", + "SystemTime", + "std::env::", + "provider.execute", + "source.fetch", + "relay.connect", + "route(", + "Router", + ] { + assert!(!STATUS.contains(forbidden), "status gained {forbidden}"); + assert!( + !OPERATIONS.contains(forbidden), + "operations gained {forbidden}" + ); + } + assert!(PUBLIC_API.contains("impl core::clone::Clone for rhi::RhiStatusReader")); + assert!(!PUBLIC_API.contains("impl core::clone::Clone for rhi::RhiStatusPublisher")); + assert!(!PUBLIC_API.contains("impl core::clone::Clone for rhi::RhiStatusSnapshot")); + for forbidden in [ + "radroots_service_host::OperationsServer", + "radroots_service_host::BoundOperationsServer", + "radroots_service_host::CachedServiceState", + "radroots_service_host::BoundedMetricsSnapshot", + ] { + assert!(!PUBLIC_API.contains(forbidden)); + } + for required in [ + "## Passive lifecycle status and TCP operations", + "exactly HTTP/1.1 `GET /livez`, `GET /readyz`, and `GET /metrics`", + "Requests perform no SQLite", + "status_cache.v1.json", + "tcp_operations.v1.json", + ] { + assert!(README.contains(required), "README is missing {required}"); + } } #[test] @@ -1123,7 +1232,7 @@ fn public_errors_are_crate_owned_redacted_and_source_free() { .lines() .filter(|line| line.starts_with("pub struct rhi::") && line.ends_with("Error")) .count(); - assert_eq!(public_error_count, 36); + assert_eq!(public_error_count, 38); } #[test] diff --git a/tests/services_hardening_operations.rs b/tests/services_hardening_operations.rs @@ -0,0 +1,233 @@ +#![forbid(unsafe_code)] + +use std::error::Error; +use std::net::{Ipv4Addr, SocketAddrV4, TcpListener}; + +use rhi::{ + InstanceId, RHI_LIVEZ_PATH, RHI_METRICS_PATH, RHI_OPERATIONS_CONTRACT_VERSION, RHI_READYZ_PATH, + RhiConfigProfile, RhiEvidenceTransportStatusV1, RhiIdentityHealthV1, RhiIntegrityStateV1, + RhiOperationsCancellationToken, RhiOperationsErrorKind, RhiOperationsServer, + RhiPersistenceHealthV1, RhiPersistenceStatusV1, RhiPresenceStatusV1, RhiProviderStatusV1, + RhiPublicationStatusV1, RhiReconciliationStatusV1, RhiServicePhase, RhiStatusBuildInfoV1, + RhiStatusBuildMode, RhiStatusCommonV1, RhiStatusConfigurationIdentityV1, + RhiStatusConfigurationSource, RhiStatusObservationV1, RhiStatusReasonCodes, + RhiTransportHealthV1, parse_rhi_config_v1, rhi_status_cache, +}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpStream; + +const CONTRACT: &str = include_str!("../contracts/services_hardening/tcp_operations.v1.json"); +const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); +const SERVICE_REVISION: &str = "0123456789abcdef0123456789abcdef01234567"; +const LIB_REVISION: &str = "89abcdef0123456789abcdef0123456789abcdef"; + +fn build_info() -> RhiStatusBuildInfoV1 { + RhiStatusBuildInfoV1::new( + RhiStatusBuildMode::Release, + Some("0.1.0"), + Some(SERVICE_REVISION), + Some(LIB_REVISION), + Some("1.97.1"), + Some("x86_64-unknown-linux-gnu"), + Some("service-host"), + ) + .expect("build info") +} + +fn identity(configured: bool, available: bool) -> RhiIdentityHealthV1 { + RhiIdentityHealthV1::new(configured, available, RhiStatusReasonCodes::empty()) + .expect("identity health") +} + +fn observation(phase: RhiServicePhase, ready: bool) -> RhiStatusObservationV1 { + let configuration = RhiStatusConfigurationIdentityV1::new( + "a".repeat(64), + RhiStatusConfigurationSource::ExplicitConfig, + ) + .expect("configuration"); + let persistence = RhiPersistenceStatusV1::new( + RhiPersistenceHealthV1::Ready, + 10, + 42, + RhiIntegrityStateV1::Verified, + RhiStatusReasonCodes::empty(), + ) + .expect("persistence"); + let provider = RhiProviderStatusV1::new(identity(true, true), RhiStatusReasonCodes::empty()) + .expect("provider"); + let transport = RhiEvidenceTransportStatusV1::new( + RhiTransportHealthV1::Ready, + true, + true, + 2, + 2, + RhiStatusReasonCodes::empty(), + ) + .expect("transport"); + RhiStatusObservationV1::new( + RhiStatusCommonV1::new( + phase, + ready, + RhiStatusReasonCodes::empty(), + 1, + build_info(), + configuration, + persistence, + ) + .expect("common status"), + provider, + transport, + RhiReconciliationStatusV1::default(), + RhiPublicationStatusV1::new(0, 0, None), + RhiPresenceStatusV1::default(), + ) +} + +fn available_port() -> u16 { + let listener = + TcpListener::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)).expect("ephemeral listener"); + listener.local_addr().expect("local address").port() +} + +fn enabled_config(port: u16) -> String { + CONFIG.replacen( + "[operations]\nenabled = false", + &format!( + "[operations]\nenabled = true\nlisten = \"127.0.0.1:{port}\"\nbind_policy = \"loopback_only\"\n\n[operations.limits]\nheader_count = 16\nheader_bytes = 8192\nresponse_body_utf8_bytes = 4096\nconcurrent_connections = 4\nrequest_deadline_ms = 500\nidle_timeout_ms = 500" + ), + 1, + ) +} + +async fn raw_request(address: std::net::SocketAddr, request: &[u8]) -> Vec<u8> { + let mut stream = TcpStream::connect(address).await.expect("connect"); + stream.write_all(request).await.expect("request write"); + let mut response = Vec::new(); + stream + .read_to_end(&mut response) + .await + .expect("response read"); + response +} + +fn response_text(response: &[u8]) -> &str { + std::str::from_utf8(response).expect("response UTF-8") +} + +#[test] +fn machine_contract_freezes_exact_routes_metrics_and_non_authority() { + let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract"); + assert_eq!(contract["schema"], "radroots.rhi.tcp-operations.v1"); + assert_eq!( + contract["contract_version"], + RHI_OPERATIONS_CONTRACT_VERSION + ); + assert_eq!(contract["step"], 212); + assert_eq!( + contract["routes"] + .as_array() + .expect("routes") + .iter() + .map(|route| route["path"].as_str().expect("route path")) + .collect::<Vec<_>>(), + [RHI_LIVEZ_PATH, RHI_READYZ_PATH, RHI_METRICS_PATH] + ); + assert_eq!(contract["route_registration_extension"], false); + assert_eq!(contract["metrics"]["high_cardinality_labels"], false); + assert_eq!(contract["metrics"]["arbitrary_labels"], false); +} + +#[tokio::test] +async fn exact_tcp_routes_use_only_latest_cached_lifecycle_and_metrics() { + let config = parse_rhi_config_v1( + enabled_config(available_port()).as_bytes(), + RhiConfigProfile::Production, + ) + .expect("enabled config"); + let (mut publisher, reader) = rhi_status_cache( + InstanceId::new("primary").expect("instance"), + observation(RhiServicePhase::Ready, true), + ) + .expect("status cache"); + let detail_pointer = reader.snapshot().detailed_status_json().as_ptr(); + let bound = RhiOperationsServer::new(&config, &reader) + .expect("operations server") + .bind() + .await + .expect("bind"); + let address = bound.local_address(); + let cancellation = RhiOperationsCancellationToken::new(); + let task = tokio::spawn(bound.serve(cancellation.clone())); + + let live = raw_request(address, b"GET /livez HTTP/1.1\r\nhost: localhost\r\n\r\n").await; + let ready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await; + let metrics = raw_request(address, b"GET /metrics HTTP/1.1\r\nhost: localhost\r\n\r\n").await; + assert!(response_text(&live).starts_with("HTTP/1.1 200 OK\r\n")); + assert!(response_text(&live).ends_with("live\n")); + assert!(response_text(&ready).ends_with("ready\n")); + assert!(response_text(&metrics).contains("# TYPE radroots_rhi_service_phase gauge\n")); + assert!(response_text(&metrics).contains("radroots_rhi_service_phase{phase=\"ready\"} 1\n")); + assert!(response_text(&metrics).contains("radroots_rhi_service_ready 1\n")); + + for request in [ + &b"GET /status HTTP/1.1\r\nhost: localhost\r\n\r\n"[..], + &b"GET /readyz?probe=1 HTTP/1.1\r\nhost: localhost\r\n\r\n"[..], + &b"POST /metrics HTTP/1.1\r\nhost: localhost\r\ncontent-length: 0\r\n\r\n"[..], + &b"GET /v1/status HTTP/1.1\r\nhost: localhost\r\n\r\n"[..], + ] { + let rejected = raw_request(address, request).await; + assert!(response_text(&rejected).starts_with("HTTP/1.1 404 Not Found\r\n")); + } + + publisher + .publish(observation(RhiServicePhase::Unready, false)) + .expect("unready publication"); + let unready = raw_request(address, b"GET /readyz HTTP/1.1\r\nhost: localhost\r\n\r\n").await; + let metrics = raw_request(address, b"GET /metrics HTTP/1.1\r\nhost: localhost\r\n\r\n").await; + assert!(response_text(&unready).starts_with("HTTP/1.1 503 Service Unavailable\r\n")); + assert!(response_text(&unready).ends_with("unready\n")); + assert!(response_text(&metrics).contains("radroots_rhi_service_phase{phase=\"unready\"} 1\n")); + assert!(response_text(&metrics).contains("radroots_rhi_service_ready 0\n")); + assert_ne!( + reader.snapshot().detailed_status_json().as_ptr(), + detail_pointer + ); + + cancellation.cancel(); + assert_eq!(task.await.expect("serve task"), Ok(())); +} + +#[tokio::test] +async fn disabled_invalid_and_bind_failures_are_typed_source_free_and_redacted() { + let disabled = parse_rhi_config_v1(CONFIG.as_bytes(), RhiConfigProfile::Production) + .expect("disabled config"); + let (_, reader) = rhi_status_cache( + InstanceId::new("primary").expect("instance"), + observation(RhiServicePhase::Ready, true), + ) + .expect("status cache"); + let disabled_error = + RhiOperationsServer::new(&disabled, &reader).expect_err("disabled operations"); + assert_eq!(disabled_error.kind(), RhiOperationsErrorKind::Disabled); + assert_eq!(disabled_error.code(), "operations_disabled"); + assert!(Error::source(&disabled_error).is_none()); + + let below_floor = + enabled_config(available_port()).replace("header_bytes = 8192", "header_bytes = 8191"); + assert!(parse_rhi_config_v1(below_floor.as_bytes(), RhiConfigProfile::Production).is_err()); + + let occupied = + TcpListener::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)).expect("occupied listener"); + let port = occupied.local_addr().expect("occupied address").port(); + let config = parse_rhi_config_v1( + enabled_config(port).as_bytes(), + RhiConfigProfile::Production, + ) + .expect("enabled config"); + let server = RhiOperationsServer::new(&config, &reader).expect("server"); + assert!(!format!("{server:?}").contains(&port.to_string())); + let bind_error = server.bind().await.expect_err("occupied bind"); + assert_eq!(bind_error.kind(), RhiOperationsErrorKind::Bind); + assert!(Error::source(&bind_error).is_none()); + assert!(!format!("{bind_error:?} {bind_error}").contains(&port.to_string())); +} diff --git a/tests/services_hardening_status.rs b/tests/services_hardening_status.rs @@ -0,0 +1,501 @@ +#![forbid(unsafe_code)] + +use std::error::Error; + +use rhi::{ + InstanceId, RHI_DETAILED_STATUS_MAX_UTF8_BYTES, RHI_STATUS_CACHE_CONTRACT_VERSION, + RhiEvidenceTransportStatusV1, RhiIdentityHealthV1, RhiIntegrityStateV1, RhiPersistenceHealthV1, + RhiPersistenceStatusV1, RhiPresenceStatusV1, RhiProviderStatusV1, RhiPublicationStatusV1, + RhiReconciliationStatusV1, RhiServicePhase, RhiStatusBuildInfoV1, RhiStatusBuildMode, + RhiStatusCommonV1, RhiStatusConfigurationIdentityV1, RhiStatusConfigurationSource, + RhiStatusErrorKind, RhiStatusObservationV1, RhiStatusReasonCode, RhiStatusReasonCodes, + RhiStatusUnixSeconds, RhiTransportHealthV1, rhi_status_cache, +}; + +const CONTRACT: &str = include_str!("../contracts/services_hardening/status_cache.v1.json"); +const SERVICE_REVISION: &str = "0123456789abcdef0123456789abcdef01234567"; +const LIB_REVISION: &str = "89abcdef0123456789abcdef0123456789abcdef"; + +fn reasons(values: &[&str]) -> RhiStatusReasonCodes { + RhiStatusReasonCodes::new( + values + .iter() + .map(|value| RhiStatusReasonCode::new(value).expect("reason code")), + ) + .expect("reason codes") +} + +fn build_info() -> RhiStatusBuildInfoV1 { + RhiStatusBuildInfoV1::new( + RhiStatusBuildMode::Release, + Some("0.1.0"), + Some(SERVICE_REVISION), + Some(LIB_REVISION), + Some("1.97.1"), + Some("x86_64-unknown-linux-gnu"), + Some("service-host"), + ) + .expect("build info") +} + +fn identity(configured: bool, available: bool, reason: &[&str]) -> RhiIdentityHealthV1 { + RhiIdentityHealthV1::new(configured, available, reasons(reason)).expect("identity health") +} + +fn observation( + phase: RhiServicePhase, + ready: bool, + uptime_ms: u64, + pending_jobs: u64, +) -> RhiStatusObservationV1 { + let lifecycle_reasons = if phase == RhiServicePhase::Degraded { + reasons(&["source_unavailable"]) + } else { + RhiStatusReasonCodes::empty() + }; + let configuration = RhiStatusConfigurationIdentityV1::new( + "a".repeat(64), + RhiStatusConfigurationSource::ExplicitConfig, + ) + .expect("configuration"); + let persistence = RhiPersistenceStatusV1::new( + RhiPersistenceHealthV1::Ready, + 10, + 42, + RhiIntegrityStateV1::Verified, + RhiStatusReasonCodes::empty(), + ) + .expect("persistence"); + let provider = + RhiProviderStatusV1::new(identity(true, true, &[]), RhiStatusReasonCodes::empty()) + .expect("provider"); + let transport = RhiEvidenceTransportStatusV1::new( + if phase == RhiServicePhase::Degraded { + RhiTransportHealthV1::Degraded + } else { + RhiTransportHealthV1::Ready + }, + true, + true, + 2, + if phase == RhiServicePhase::Degraded { + 1 + } else { + 2 + }, + if phase == RhiServicePhase::Degraded { + reasons(&["source_unavailable"]) + } else { + RhiStatusReasonCodes::empty() + }, + ) + .expect("transport"); + RhiStatusObservationV1::new( + RhiStatusCommonV1::new( + phase, + ready, + lifecycle_reasons, + uptime_ms, + build_info(), + configuration, + persistence, + ) + .expect("common status"), + provider, + transport, + RhiReconciliationStatusV1::new( + pending_jobs, + 3, + 1, + Some(RhiStatusUnixSeconds::new(1_723_456_700).expect("job time")), + ), + RhiPublicationStatusV1::new( + 4, + 1, + Some(RhiStatusUnixSeconds::new(1_723_456_789).expect("publication time")), + ), + RhiPresenceStatusV1::new(2, 1), + ) +} + +#[test] +fn machine_contract_and_canonical_detailed_status_are_exact() { + let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("status contract"); + assert_eq!(contract["schema"], "radroots.rhi.status-cache.v1"); + assert_eq!( + contract["contract_version"], + RHI_STATUS_CACHE_CONTRACT_VERSION + ); + assert_eq!(contract["step"], 212); + assert_eq!( + contract["publication"]["capacity"], + "one_latest_immutable_arc" + ); + assert_eq!(contract["read"]["fresh_probe"], false); + assert_eq!( + contract["detailed_status"]["maximum_utf8_bytes"], + RHI_DETAILED_STATUS_MAX_UTF8_BYTES + ); + assert_eq!( + contract["detailed_status"]["reason_codes"] + .as_array() + .expect("reason inventory") + .iter() + .map(|value| value.as_str().expect("reason")) + .collect::<Vec<_>>(), + [ + RhiStatusReasonCode::IdentityUnavailable, + RhiStatusReasonCode::DatabaseSchemaMismatch, + RhiStatusReasonCode::DatabaseReadOnly, + RhiStatusReasonCode::DatabaseLowDisk, + RhiStatusReasonCode::SourceUnavailable, + RhiStatusReasonCode::SubscriptionInactive, + RhiStatusReasonCode::RecoveryIncomplete, + RhiStatusReasonCode::PublicationRecoveryIncomplete, + RhiStatusReasonCode::PresenceStateUnavailable, + RhiStatusReasonCode::ReconciliationBacklogExceeded, + RhiStatusReasonCode::AdminListenerFailed, + RhiStatusReasonCode::OperationsListenerFailed, + RhiStatusReasonCode::ShutdownInProgress, + ] + .map(RhiStatusReasonCode::as_str) + ); + + let (_publisher, reader) = rhi_status_cache( + InstanceId::new("primary").expect("instance"), + observation(RhiServicePhase::Ready, true, 120_000, 5), + ) + .expect("status cache"); + let snapshot = reader.snapshot(); + let wire = std::str::from_utf8(snapshot.detailed_status_json()).expect("status UTF-8"); + assert_eq!( + wire, + r#"{"contract_version":1,"service":"rhi","instance":"primary","phase":"ready","ready":true,"uptime_millis":120000,"reason_codes":[],"build_info":{"version":"0.1.0","revision":"0123456789abcdef0123456789abcdef01234567","toolchain":"1.97.1","contract_versions":{"config":1,"state":10,"admin":1,"status":1,"provider":1}},"configuration":{"schema":"radroots.rhi.config","schema_version":1,"digest":"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","source":"explicit_config"},"persistence":{"health":"ready","schema_version":10,"generation":42,"integrity":"verified","reason_codes":[]},"provider":{"health":"ready","identity":{"configured":true,"available":true,"reason_codes":[]},"reason_codes":[]},"transport":{"health":"ready","required_sources_ready":true,"subscriber_active":true,"configured_source_count":2,"reachable_source_count":2,"reason_codes":[]},"rhi":{"identity":{"configured":true,"available":true,"reason_codes":[]},"reconciliation":{"pending":5,"leased":3,"exhausted":1,"oldest_pending_at_utc":1723456700},"publication":{"pending":4,"unknown":1,"oldest_pending_at_utc":1723456789},"presence":{"pending":2,"unknown":1}}}"# + ); + assert!(wire.len() < RHI_DETAILED_STATUS_MAX_UTF8_BYTES); + for forbidden in [ + "secret", + "credential", + "private_key", + "password", + "filesystem_path", + "relay_url", + "raw_error", + ] { + assert!(!wire.contains(forbidden), "wire leaked `{forbidden}`"); + } +} + +#[tokio::test] +async fn latest_publication_is_atomic_passive_and_retains_old_snapshots() { + let (mut publisher, reader) = rhi_status_cache( + InstanceId::new("primary").expect("instance"), + observation(RhiServicePhase::Starting, false, 0, 9), + ) + .expect("status cache"); + let old = reader.snapshot(); + let old_pointer = old.detailed_status_json().as_ptr(); + for _ in 0..1_000 { + let same = reader.snapshot(); + assert_eq!(same.detailed_status_json().as_ptr(), old_pointer); + assert_eq!(same.phase(), RhiServicePhase::Starting); + } + + let mut changed = publisher.subscribe(); + publisher + .publish(observation(RhiServicePhase::Ready, true, 10, 2)) + .expect("ready publication"); + publisher + .publish(observation(RhiServicePhase::Degraded, true, 20, 1)) + .expect("degraded publication"); + + let latest = changed.changed().await.expect("latest publication"); + assert_eq!(latest.phase(), RhiServicePhase::Degraded); + assert!(latest.is_ready()); + assert!( + std::str::from_utf8(latest.detailed_status_json()) + .expect("status UTF-8") + .contains("\"uptime_millis\":20") + ); + assert_eq!(old.phase(), RhiServicePhase::Starting); + assert!( + std::str::from_utf8(old.detailed_status_json()) + .expect("old UTF-8") + .contains("\"uptime_millis\":0") + ); +} + +#[tokio::test] +async fn illegal_transition_and_publisher_drop_preserve_the_last_valid_value() { + let (mut publisher, reader) = rhi_status_cache( + InstanceId::new("primary").expect("instance"), + observation(RhiServicePhase::Starting, false, 1, 0), + ) + .expect("status cache"); + let before = reader.snapshot().detailed_status_json().to_vec(); + let error = publisher + .publish(observation(RhiServicePhase::Unready, false, 2, 0)) + .expect_err("illegal starting to unready transition"); + assert_eq!(error.kind(), RhiStatusErrorKind::InvalidTransition); + assert_eq!(reader.snapshot().detailed_status_json(), before); + assert!(Error::source(&error).is_none()); + + let mut dropped = publisher.subscribe(); + drop(publisher); + let dropped_error = dropped.changed().await.expect_err("publisher dropped"); + assert_eq!(dropped_error.kind(), RhiStatusErrorKind::PublisherDropped); + assert_eq!(reader.snapshot().detailed_status_json(), before); +} + +#[test] +fn closed_work_counts_time_and_safe_debug_bound_the_status_surface() { + let work = RhiReconciliationStatusV1::new(u64::MAX, 2, 3, None); + assert_eq!(work.pending(), u64::MAX); + assert_eq!(work.leased(), 2); + assert_eq!(work.exhausted(), 3); + assert_eq!(work.oldest_pending_at_utc(), None); + assert_eq!(RhiStatusUnixSeconds::new(0).expect("zero").get(), 0); + assert_eq!( + RhiStatusUnixSeconds::new(i64::MAX as u64) + .expect("maximum") + .get(), + i64::MAX as u64 + ); + assert_eq!( + RhiStatusUnixSeconds::new(i64::MAX as u64 + 1) + .expect_err("over maximum") + .kind(), + RhiStatusErrorKind::InvalidTime + ); + + let (publisher, reader) = rhi_status_cache( + InstanceId::new("primary").expect("instance"), + observation(RhiServicePhase::Ready, true, 5, 0), + ) + .expect("status cache"); + let rendered = format!("{publisher:?} {reader:?} {:?}", reader.snapshot()); + for forbidden in [ + SERVICE_REVISION, + LIB_REVISION, + "radroots.rhi.config", + "aaaaaaaaaaaaaaaa", + "oldest_pending_at_utc", + ] { + assert!(!rendered.contains(forbidden)); + } +} + +#[test] +fn model_boundaries_fail_closed_before_publication() { + assert_eq!( + RhiStatusReasonCode::new("") + .expect_err("empty reason") + .kind(), + RhiStatusErrorKind::InvalidReasonCode + ); + assert_eq!( + RhiStatusReasonCode::new("secret_canary_value") + .expect_err("unknown reason") + .kind(), + RhiStatusErrorKind::InvalidReasonCode + ); + let maximum = [ + RhiStatusReasonCode::IdentityUnavailable, + RhiStatusReasonCode::DatabaseSchemaMismatch, + RhiStatusReasonCode::DatabaseReadOnly, + RhiStatusReasonCode::DatabaseLowDisk, + RhiStatusReasonCode::SourceUnavailable, + RhiStatusReasonCode::SubscriptionInactive, + RhiStatusReasonCode::RecoveryIncomplete, + RhiStatusReasonCode::PublicationRecoveryIncomplete, + RhiStatusReasonCode::PresenceStateUnavailable, + RhiStatusReasonCode::ReconciliationBacklogExceeded, + RhiStatusReasonCode::AdminListenerFailed, + RhiStatusReasonCode::OperationsListenerFailed, + RhiStatusReasonCode::ShutdownInProgress, + ]; + assert_eq!( + RhiStatusReasonCodes::new(maximum) + .expect("maximum reasons") + .as_slice() + .len(), + 13 + ); + let mut infinite = std::iter::repeat(RhiStatusReasonCode::IdentityUnavailable); + assert_eq!( + RhiStatusReasonCodes::new(&mut infinite) + .expect_err("bounded infinite iterator") + .kind(), + RhiStatusErrorKind::TooManyReasonCodes + ); + assert_eq!( + infinite.next().expect("iterator retained"), + RhiStatusReasonCode::IdentityUnavailable + ); + + assert_eq!( + RhiStatusBuildInfoV1::new( + RhiStatusBuildMode::Release, + Some("0.1.0"), + None, + Some(LIB_REVISION), + Some("1.97.1"), + Some("x86_64-unknown-linux-gnu"), + Some("service-host"), + ) + .expect_err("release revision required") + .kind(), + RhiStatusErrorKind::InvalidBuildInfo + ); + assert_eq!( + RhiStatusConfigurationIdentityV1::new( + "A".repeat(64), + RhiStatusConfigurationSource::ExplicitConfig, + ) + .expect_err("lowercase digest required") + .kind(), + RhiStatusErrorKind::InvalidConfiguration + ); + assert_eq!( + RhiPersistenceStatusV1::new( + RhiPersistenceHealthV1::Ready, + 0, + 0, + RhiIntegrityStateV1::Verified, + RhiStatusReasonCodes::empty(), + ) + .expect_err("positive schema required") + .kind(), + RhiStatusErrorKind::InvalidPersistence + ); + + let error = RhiStatusCommonV1::new( + RhiServicePhase::Ready, + false, + RhiStatusReasonCodes::empty(), + 0, + build_info(), + RhiStatusConfigurationIdentityV1::new( + "a".repeat(64), + RhiStatusConfigurationSource::ExplicitConfig, + ) + .expect("configuration"), + RhiPersistenceStatusV1::new( + RhiPersistenceHealthV1::Ready, + 10, + 0, + RhiIntegrityStateV1::Verified, + RhiStatusReasonCodes::empty(), + ) + .expect("persistence"), + ) + .expect_err("ready phase requires readiness"); + assert_eq!(error.kind(), RhiStatusErrorKind::InvalidLifecycle); + + let inconsistent = RhiStatusObservationV1::new( + RhiStatusCommonV1::new( + RhiServicePhase::Ready, + true, + RhiStatusReasonCodes::empty(), + 1, + build_info(), + RhiStatusConfigurationIdentityV1::new( + "a".repeat(64), + RhiStatusConfigurationSource::ExplicitConfig, + ) + .expect("configuration"), + RhiPersistenceStatusV1::new( + RhiPersistenceHealthV1::Ready, + 10, + 1, + RhiIntegrityStateV1::Verified, + RhiStatusReasonCodes::empty(), + ) + .expect("persistence"), + ) + .expect("common"), + RhiProviderStatusV1::new( + identity(true, false, &["identity_unavailable"]), + reasons(&["identity_unavailable"]), + ) + .expect("provider"), + RhiEvidenceTransportStatusV1::new( + RhiTransportHealthV1::Ready, + true, + true, + 1, + 1, + RhiStatusReasonCodes::empty(), + ) + .expect("transport"), + RhiReconciliationStatusV1::default(), + RhiPublicationStatusV1::default(), + RhiPresenceStatusV1::default(), + ); + assert_eq!( + rhi_status_cache(InstanceId::new("primary").expect("instance"), inconsistent,) + .expect_err("ready status requires healthy critical dependencies") + .kind(), + RhiStatusErrorKind::InvalidLifecycle + ); +} + +#[test] +fn provider_transport_and_optional_oldest_time_are_deterministic() { + let ready = RhiProviderStatusV1::new(identity(true, true, &[]), RhiStatusReasonCodes::empty()) + .expect("ready provider"); + assert_eq!(ready.health(), rhi::RhiProviderHealthV1::Ready); + + let unavailable = RhiProviderStatusV1::new( + identity(true, false, &["identity_unavailable"]), + reasons(&["identity_unavailable"]), + ) + .expect("unavailable provider"); + assert_eq!(unavailable.health(), rhi::RhiProviderHealthV1::Unavailable); + + let (_publisher, reader) = rhi_status_cache( + InstanceId::new("primary").expect("instance"), + RhiStatusObservationV1::new( + RhiStatusCommonV1::new( + RhiServicePhase::Ready, + true, + RhiStatusReasonCodes::empty(), + 1, + build_info(), + RhiStatusConfigurationIdentityV1::new( + "a".repeat(64), + RhiStatusConfigurationSource::ExplicitConfig, + ) + .expect("configuration"), + RhiPersistenceStatusV1::new( + RhiPersistenceHealthV1::Ready, + 10, + 0, + RhiIntegrityStateV1::Verified, + RhiStatusReasonCodes::empty(), + ) + .expect("persistence"), + ) + .expect("common status"), + ready, + RhiEvidenceTransportStatusV1::new( + RhiTransportHealthV1::Ready, + true, + true, + 1, + 1, + RhiStatusReasonCodes::empty(), + ) + .expect("transport"), + RhiReconciliationStatusV1::default(), + RhiPublicationStatusV1::new(0, 0, None), + RhiPresenceStatusV1::default(), + ), + ) + .expect("cache"); + let wire = std::str::from_utf8(reader.snapshot().detailed_status_json()) + .expect("status UTF-8") + .to_owned(); + assert!(wire.contains("\"publication\":{\"pending\":0,\"unknown\":0}")); + assert!(!wire.contains("oldest_pending_at_utc")); +}